summaryrefslogtreecommitdiff
path: root/web/mux.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-20 03:53:17 -0300
committerGabriel Schneider <[email protected]>2026-09-20 03:53:17 -0300
commitf1b53c1533539aecbf16ad19fd9156deae091f92 (patch)
tree38ff00810ef5e0ad271c26e218cf6af4578c71ce /web/mux.zig
parentba7ec40782ba7020d82a56896a5eb1b52578d6aa (diff)
downloadcloud9-f1b53c1533539aecbf16ad19fd9156deae091f92.tar.gz
cloud9-f1b53c1533539aecbf16ad19fd9156deae091f92.zip
9web: multiplexer, HTTP view of the tree, live streams, richer page
One upstream 9P connection now serves any number of browser WebSocket sessions and, with --serve, plain 9P clients over TCP or Unix (a 9pserve- style frame remux: tags and fids remapped, Tflush forwarded, reconnect on upstream loss). /fs/<path> maps HTTP onto the tree: GET file or directory (JSON or HTML), Range, HEAD with 9P headers, PUT (create, truncate, append, trailing slash makes a directory), DELETE, and ?follow=1 or text/event-stream turning a blocking read into server-sent events with Tflush on disconnect. The page gains a lazy tree, stat panel, create, rename, delete, upload and follow mode. --probe embeds 9proc for self-introspection. New mux-test and http-fs-test steps; e2e still passes. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Diffstat (limited to 'web/mux.zig')
-rw-r--r--web/mux.zig812
1 files changed, 812 insertions, 0 deletions
diff --git a/web/mux.zig b/web/mux.zig
new file mode 100644
index 0000000..23f2e6c
--- /dev/null
+++ b/web/mux.zig
@@ -0,0 +1,812 @@
+//! A transparent 9P2000 multiplexer: one shared upstream `cloud9.Client`
+//! connection (all 16 tags) fanned out to any number of downstream sessions.
+//!
+//! Each downstream request is forwarded upstream on a remapped tag with its
+//! fids remapped into the shared upstream fid space; each upstream reply is
+//! routed back to the originating downstream by tag. Tversion is answered
+//! locally (the upstream session is negotiated once at connect); Tattach is
+//! forwarded so every downstream gets its own upstream tree root; Tflush is
+//! forwarded as Tflush. Downstream fid spaces are isolated by construction:
+//! two downstreams never share an upstream fid.
+//!
+//! The upstream is driven by one reader task and any number of downstream
+//! forwarder tasks, all serialized by `lock`. Nothing here allocates per
+//! request: the tag/fid/pending tables are comptime-sized. On an upstream I/O
+//! failure the reader reconnects, bumps `generation`, and every downstream's
+//! fids from an older generation are answered EIO until it re-attaches.
+//!
+//! This is deliberately *not* an `fs.Server`/`serve.Runner` backend: the
+//! engine re-decomposes each request into filesystem operations and re-encodes
+//! directories with synthetic stats and an entry-index cursor, which loses the
+//! upstream's real directory stats, iounit and qids and complicates the
+//! readdir byte offset. A frame-level remux forwards requests unchanged, which
+//! is exactly the "forward each request upstream" contract and matches the way
+//! plan9port's 9pserve multiplexes. The HTTP `/fs` view (see httpfs.zig) uses
+//! the same `Mux` at the fid level directly.
+const std = @import("std");
+const c9 = @import("cloud9");
+const Io = std.Io;
+const wire = c9.wire;
+const transport = c9.transport;
+
+const notag = wire.notag;
+const nofid = wire.nofid;
+
+pub const Limits = struct {
+ /// Largest 9P frame on either side; sizes the shared upstream buffers.
+ msize: u32 = 65536,
+ /// Upstream fids the shared connection may hold across all downstreams.
+ upstream_fids: usize = 4096,
+ /// Fids one downstream session may hold at once.
+ fids_per_conn: usize = 512,
+ /// Downstream sessions routed at once (for introspection only).
+ downstreams: usize = 256,
+};
+
+/// How the mux hands an upstream reply back to a downstream. `send` is called
+/// off the mux lock and must serialize writes on that downstream itself.
+pub const Sink = struct {
+ ctx: *anyopaque,
+ send: *const fn (ctx: *anyopaque, frame: []const u8) anyerror!void,
+};
+
+pub const Address = union(enum) {
+ network: transport.Address,
+ /// An already configured duplex device (serial), one session at a time.
+ file: []const u8,
+};
+
+pub const Error = error{
+ Disconnected,
+ TooManyFids,
+ UnknownFid,
+ BadRequest,
+ Upstream,
+};
+
+pub fn Mux(comptime limits: Limits) type {
+ if (limits.msize < 256) @compileError("mux msize too small");
+ return struct {
+ const Self = @This();
+ pub const msize = limits.msize;
+
+ io: Io,
+ address: Address,
+ user: []const u8,
+ tree: []const u8,
+
+ client: c9.Client = undefined,
+ in: [msize]u8 = undefined,
+ out: [msize]u8 = undefined,
+ rbuf: [msize]u8 = undefined,
+ stage: [msize]u8 = undefined,
+ /// One send buffer per upstream tag (16 ordinary + 1 flush).
+ frames: [17][msize]u8 = undefined,
+ /// A scratch buffer to serialize one submitted frame out of the client.
+ wbuf: [msize]u8 = undefined,
+
+ stream: ?Io.net.Stream = null,
+ file: ?Io.File = null,
+ reader: Io.net.Stream.Reader = undefined,
+ writer: Io.net.Stream.Writer = undefined,
+ freader: Io.File.Reader = undefined,
+ fwriter: Io.File.Writer = undefined,
+ ureader: *Io.Reader = undefined,
+ uwriter: *Io.Writer = undefined,
+
+ mutex: Io.Mutex = .init,
+ wmutex: Io.Mutex = .init,
+ cond: Io.Condition = .init,
+ negotiated: u32 = 0,
+ version: [16]u8 = undefined,
+ version_len: usize = 0,
+ generation: u32 = 1,
+ connected: bool = false,
+ stopping: std.atomic.Value(bool) = .init(0 != 0),
+
+ /// Counts for the /probe introspection.
+ downstreams: std.atomic.Value(u32) = .init(0),
+ upstream_up: std.atomic.Value(u32) = .init(0),
+ reconnects: std.atomic.Value(u32) = .init(0),
+
+ fid_used: [limits.upstream_fids]bool = @splat(false),
+ fid_next: usize = 0,
+
+ pending: [17]Pending = @splat(.{}),
+
+ const Kind = enum { plain, walk, clunk, flush };
+ /// A blocking fid-level RPC used by the HTTP `/fs` view. The reply
+ /// frame is copied into `buf` before the reader advances, so its
+ /// borrowed data stays valid until the waiter consumes it.
+ pub const Rpc = struct {
+ buf: []u8,
+ len: usize = 0,
+ ready: bool = false,
+ failed: bool = false,
+ /// The upstream tag this call holds while in flight (for `flushRpc`).
+ tag: u16 = 0,
+ event: Io.Event = .unset,
+ };
+ const Pending = struct {
+ active: bool = false,
+ kind: Kind = .plain,
+ sink: Sink = undefined,
+ conn: ?*Conn = null,
+ rpc: ?*Rpc = null,
+ down_tag: u16 = 0,
+ /// walk: the downstream/upstream newfid and whether it was in place.
+ newfid_down: u32 = 0,
+ newfid_up: u32 = 0,
+ walk_names: u16 = 0,
+ inplace: bool = false,
+ /// clunk/remove: the downstream fid to drop on completion.
+ clunk_down: u32 = 0,
+ /// flush: the upstream tag it cancels.
+ flush_up: u16 = 0,
+ };
+
+ /// A downstream session: its own fid map and the generation it belongs
+ /// to. `sink` routes replies; `write` on the transport must serialize.
+ pub const Conn = struct {
+ mux: *Self,
+ sink: Sink,
+ generation: u32 = 0,
+ alive: bool = true,
+ fids: [limits.fids_per_conn]FidMap = @splat(.{}),
+
+ const FidMap = struct { used: bool = false, down: u32 = 0, up: u32 = 0 };
+
+ pub fn init(m: *Self, sink: Sink) Conn {
+ _ = m.downstreams.fetchAdd(1, .monotonic);
+ return .{ .mux = m, .sink = sink, .generation = m.generation };
+ }
+
+ /// Drops every upstream fid this session holds and cancels its
+ /// outstanding requests, then unregisters it.
+ pub fn deinit(c: *Conn) void {
+ const m = c.mux;
+ m.lock();
+ c.alive = false;
+ // Orphan any in-flight replies bound for this session.
+ for (&m.pending) |*p| if (p.active and p.conn == c) {
+ p.conn = null;
+ };
+ // Best-effort clunk of every live upstream fid.
+ if (c.generation == m.generation and m.connected) {
+ for (&c.fids) |*e| if (e.used) {
+ m.clunkUpstreamLocked(e.up);
+ e.used = false;
+ };
+ m.flushOutput();
+ }
+ m.unlock();
+ _ = m.downstreams.fetchSub(1, .monotonic);
+ }
+
+ fn mapFind(c: *Conn, down: u32) ?*FidMap {
+ for (&c.fids) |*e| if (e.used and e.down == down) return e;
+ return null;
+ }
+ fn mapAdd(c: *Conn, down: u32, up: u32) void {
+ for (&c.fids) |*e| if (!e.used) {
+ e.* = .{ .used = true, .down = down, .up = up };
+ return;
+ };
+ unreachable; // caller checked capacity via free upstream fid
+ }
+ };
+
+ pub fn init(m: *Self, io: Io, address: Address, user: []const u8, tree: []const u8) void {
+ m.* = .{ .io = io, .address = address, .user = user, .tree = tree };
+ }
+
+ fn lock(m: *Self) void {
+ m.mutex.lockUncancelable(m.io);
+ }
+ fn unlock(m: *Self) void {
+ m.mutex.unlock(m.io);
+ }
+
+ pub fn upstreamState(m: *Self) []const u8 {
+ return if (m.upstream_up.load(.acquire) != 0) "connected" else "disconnected";
+ }
+ pub fn downstreamCount(m: *Self) u32 {
+ return m.downstreams.load(.acquire);
+ }
+ pub fn reconnectCount(m: *Self) u32 {
+ return m.reconnects.load(.acquire);
+ }
+ pub fn negotiatedMsize(m: *Self) u32 {
+ return m.negotiated;
+ }
+
+ // -- upstream connection ------------------------------------------
+
+ fn dial(m: *Self) !void {
+ switch (m.address) {
+ .network => |addr| {
+ const s = try transport.connect(m.io, addr);
+ m.stream = s;
+ m.reader = s.reader(m.io, &m.rbuf);
+ m.writer = s.writer(m.io, &m.wbuf);
+ m.ureader = &m.reader.interface;
+ m.uwriter = &m.writer.interface;
+ },
+ .file => |path| {
+ const f = try Io.Dir.cwd().openFile(m.io, path, .{ .mode = .read_write });
+ m.file = f;
+ m.freader = f.readerStreaming(m.io, &m.rbuf);
+ m.fwriter = f.writerStreaming(m.io, &m.wbuf);
+ m.ureader = &m.freader.interface;
+ m.uwriter = &m.fwriter.interface;
+ },
+ }
+ }
+
+ fn closeUpstream(m: *Self) void {
+ if (m.stream) |s| {
+ s.close(m.io);
+ m.stream = null;
+ }
+ if (m.file) |f| {
+ f.close(m.io);
+ m.file = null;
+ }
+ }
+
+ /// Connects and negotiates the shared session. Call once before the
+ /// reader task runs; returns an error if the upstream is unreachable.
+ pub fn connect(m: *Self) !void {
+ try m.dial();
+ errdefer m.closeUpstream();
+ m.client = .init(.{ .in = &m.in, .out = &m.out });
+ try m.handshake();
+ m.connected = true;
+ m.upstream_up.store(1, .release);
+ }
+
+ fn handshake(m: *Self) !void {
+ _ = m.client.submit(.{ .version = .{ .msize = msize } }) catch return error.Upstream;
+ try m.sendClientOutput();
+ const done = try m.readOne();
+ if (done.op != .version) return error.Upstream;
+ m.negotiated = done.result.version.msize;
+ const v = done.result.version.version;
+ m.version_len = @min(v.len, m.version.len);
+ @memcpy(m.version[0..m.version_len], v[0..m.version_len]);
+ if (!std.mem.eql(u8, m.version[0..m.version_len], "9P2000")) return error.Upstream;
+ }
+
+ /// Writes whatever the client has staged (used only during handshake).
+ fn sendClientOutput(m: *Self) !void {
+ const bytes = m.client.output();
+ m.uwriter.writeAll(bytes) catch return error.Upstream;
+ m.uwriter.flush() catch return error.Upstream;
+ m.client.wrote(bytes.len);
+ }
+
+ fn readOne(m: *Self) !c9.Client.Done {
+ while (true) {
+ if (m.client.take()) |done| return done;
+ if (m.client.dead) return error.Upstream;
+ const frame = transport.readFrame(m.ureader, &m.stage, msize) catch return error.Upstream;
+ if (m.client.push(frame) != frame.len) return error.Upstream;
+ }
+ }
+
+ // -- upstream fid pool --------------------------------------------
+
+ fn allocFid(m: *Self) ?u32 {
+ var i: usize = 0;
+ while (i < limits.upstream_fids) : (i += 1) {
+ const idx = (m.fid_next + i) % limits.upstream_fids;
+ if (!m.fid_used[idx]) {
+ m.fid_used[idx] = true;
+ m.fid_next = (idx + 1) % limits.upstream_fids;
+ return @intCast(idx);
+ }
+ }
+ return null;
+ }
+ fn freeFid(m: *Self, fid: u32) void {
+ if (fid < limits.upstream_fids) m.fid_used[fid] = false;
+ }
+
+ // -- forwarding ----------------------------------------------------
+
+ /// Handles one downstream frame. Tversion is answered locally through
+ /// the sink; every other message is remapped and forwarded upstream,
+ /// its reply delivered later by the reader task. On a synchronous
+ /// failure it answers the downstream with an Rerror.
+ pub fn forward(m: *Self, c: *Conn, frame: []const u8) void {
+ const got = wire.decode(frame) catch {
+ c.sink.send(c.sink.ctx, m.errorFrame(&m.wbuf, 0, "protocol botch")) catch {};
+ return;
+ };
+ if (!wire.isT(got.msg.msgType())) {
+ m.rerror(c, got.tag, "protocol botch");
+ return;
+ }
+ switch (got.msg) {
+ .tversion => |v| {
+ // The shared upstream session is already negotiated. Answer
+ // locally and reset this downstream's fid space.
+ m.lock();
+ for (&c.fids) |*e| if (e.used and c.generation == m.generation and m.connected) {
+ m.clunkUpstreamLocked(e.up);
+ e.used = false;
+ } else {
+ e.used = false;
+ };
+ m.flushOutput();
+ const use: u32 = @min(@min(v.msize, m.negotiated), msize);
+ c.generation = m.generation;
+ m.unlock();
+ const reply = wire.encode(.{ .rversion = .{ .msize = use, .version = "9P2000" } }, got.tag, &m.wbuf) catch return;
+ c.sink.send(c.sink.ctx, reply) catch {};
+ },
+ else => m.forwardRequest(c, got),
+ }
+ }
+
+ fn rerror(m: *Self, c: *Conn, tag: u16, text: []const u8) void {
+ var buf: [wire.header_len + 2 + 128]u8 = undefined;
+ const frame = m.errorFrame(&buf, tag, text);
+ c.sink.send(c.sink.ctx, frame) catch {};
+ }
+
+ fn errorFrame(m: *Self, buf: []u8, tag: u16, text: []const u8) []const u8 {
+ _ = m;
+ const t = text[0..@min(text.len, 128)];
+ return wire.encode(.{ .rerror = .{ .ename = t } }, tag, buf) catch buf[0..0];
+ }
+
+ fn forwardRequest(m: *Self, c: *Conn, got: wire.Decoded) void {
+ m.lock();
+ defer m.unlock();
+
+ if (!m.connected) {
+ m.unlock();
+ m.rerror(c, got.tag, "Transport endpoint is not connected");
+ m.lock();
+ return;
+ }
+ if (c.generation != m.generation) {
+ // Stale session after a reconnect: fids are gone. Reset and,
+ // unless this is a fresh attach, answer EIO.
+ for (&c.fids) |*e| e.used = false;
+ c.generation = m.generation;
+ if (got.msg != .tattach) {
+ m.unlock();
+ m.rerror(c, got.tag, "Input/output error");
+ m.lock();
+ return;
+ }
+ }
+
+ // Translate to a client request with upstream fids.
+ var pend: Pending = .{ .active = true, .sink = c.sink, .conn = c, .down_tag = got.tag };
+ const req: c9.Client.Request = switch (got.msg) {
+ .tattach => |a| blk: {
+ if (c.mapFind(a.fid) != null) return m.syncErr(c, got.tag, "fid already in use");
+ const up = m.allocFid() orelse return m.syncErr(c, got.tag, "Too many open files in system");
+ pend.kind = .walk; // reuse walk bookkeeping to bind on success
+ pend.newfid_down = a.fid;
+ pend.newfid_up = up;
+ pend.walk_names = 0;
+ pend.inplace = false;
+ break :blk .{ .attach = .{ .fid = up, .uname = a.uname, .aname = a.aname } };
+ },
+ .twalk => |w| blk: {
+ const src = c.mapFind(w.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ var up_new: u32 = src.up;
+ const inplace = w.newfid == w.fid;
+ if (!inplace) {
+ if (c.mapFind(w.newfid) != null) return m.syncErr(c, got.tag, "fid already in use");
+ up_new = m.allocFid() orelse return m.syncErr(c, got.tag, "Too many open files in system");
+ }
+ pend.kind = .walk;
+ pend.newfid_down = w.newfid;
+ pend.newfid_up = up_new;
+ pend.walk_names = w.nwname;
+ pend.inplace = inplace;
+ break :blk .{ .walk = .{ .fid = src.up, .newfid = up_new, .names = w.wname[0..w.nwname] } };
+ },
+ .topen => |o| blk: {
+ const e = c.mapFind(o.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .open = .{ .fid = e.up, .mode = o.mode } };
+ },
+ .tcreate => |cr| blk: {
+ const e = c.mapFind(cr.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .create = .{ .fid = e.up, .name = cr.name, .perm = cr.perm, .mode = cr.mode } };
+ },
+ .tread => |r| blk: {
+ const e = c.mapFind(r.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .read = .{ .fid = e.up, .offset = r.offset, .count = r.count } };
+ },
+ .twrite => |w| blk: {
+ const e = c.mapFind(w.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .write = .{ .fid = e.up, .offset = w.offset, .data = w.data } };
+ },
+ .tclunk => |cl| blk: {
+ const e = c.mapFind(cl.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ pend.kind = .clunk;
+ pend.clunk_down = cl.fid;
+ break :blk .{ .clunk = .{ .fid = e.up } };
+ },
+ .tremove => |rm| blk: {
+ const e = c.mapFind(rm.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ pend.kind = .clunk;
+ pend.clunk_down = rm.fid;
+ break :blk .{ .remove = .{ .fid = e.up } };
+ },
+ .tstat => |s| blk: {
+ const e = c.mapFind(s.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .stat = .{ .fid = e.up } };
+ },
+ .twstat => |s| blk: {
+ const e = c.mapFind(s.fid) orelse return m.syncErr(c, got.tag, "fid unknown or out of range");
+ break :blk .{ .wstat = .{ .fid = e.up, .stat = s.stat } };
+ },
+ .tflush => |f| blk: {
+ const up = m.findUpTag(c, f.oldtag);
+ if (up == null) {
+ m.unlock();
+ var buf: [wire.header_len]u8 = undefined;
+ const rf = wire.encode(.rflush, got.tag, &buf) catch buf[0..0];
+ c.sink.send(c.sink.ctx, rf) catch {};
+ m.lock();
+ return;
+ }
+ pend.kind = .flush;
+ pend.flush_up = up.?;
+ break :blk .{ .flush = .{ .oldtag = up.? } };
+ },
+ .tauth => return m.syncErr(c, got.tag, "authentication not required"),
+ else => return m.syncErr(c, got.tag, "protocol botch"),
+ };
+
+ const up_tag = m.submitLocked(req) catch |err| {
+ if (pend.kind == .walk and !pend.inplace) m.freeFid(pend.newfid_up);
+ m.unlock();
+ m.rerror(c, got.tag, switch (err) {
+ error.Disconnected => "Transport endpoint is not connected",
+ else => "Input/output error",
+ });
+ m.lock();
+ return;
+ };
+ m.pending[up_tag] = pend;
+ }
+
+ /// A synchronous error while holding the lock: drops the lock to send,
+ /// then reacquires so the deferred unlock stays balanced.
+ fn syncErr(m: *Self, c: *Conn, tag: u16, text: []const u8) void {
+ m.unlock();
+ m.rerror(c, tag, text);
+ m.lock();
+ }
+
+ fn findUpTag(m: *Self, c: *Conn, down_tag: u16) ?u16 {
+ for (&m.pending, 0..) |*p, i| {
+ if (p.active and p.conn == c and p.down_tag == down_tag and p.kind != .flush) return @intCast(i);
+ }
+ return null;
+ }
+
+ /// Submits a request, waiting for a free tag; sends it upstream. The
+ /// caller holds `lock`. Returns the upstream tag.
+ fn submitLocked(m: *Self, req: c9.Client.Request) !u16 {
+ while (true) {
+ if (!m.connected) return error.Disconnected;
+ const tag = m.client.submit(req) catch |err| switch (err) {
+ error.NoTags => {
+ m.cond.wait(m.io, &m.mutex) catch return error.Disconnected;
+ continue;
+ },
+ else => return error.BadRequest,
+ };
+ // The client's output buffer is disjoint from the stream
+ // writer's buffer, so write it out directly under the lock.
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch return error.Disconnected;
+ m.client.wrote(bytes.len);
+ return tag;
+ }
+ }
+
+ fn writeUpstream(m: *Self, bytes: []const u8) !void {
+ m.wmutex.lockUncancelable(m.io);
+ defer m.wmutex.unlock(m.io);
+ m.uwriter.writeAll(bytes) catch return error.Upstream;
+ m.uwriter.flush() catch return error.Upstream;
+ }
+
+ fn clunkUpstreamLocked(m: *Self, up: u32) void {
+ // Fire-and-forget clunk to reclaim an upstream fid on disconnect.
+ const tag = m.client.submit(.{ .clunk = .{ .fid = up } }) catch {
+ m.freeFid(up);
+ return;
+ };
+ m.pending[tag] = .{ .active = true, .kind = .clunk, .conn = null, .clunk_down = 0 };
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch {};
+ m.client.wrote(bytes.len);
+ m.freeFid(up);
+ }
+
+ fn flushOutput(m: *Self) void {
+ _ = m;
+ }
+
+ // -- synchronous fid-level RPC (for the HTTP /fs view) -------------
+
+ /// Allocates an upstream fid from the shared pool. Returns null when
+ /// exhausted or the upstream is down.
+ pub fn takeFid(m: *Self) ?u32 {
+ m.lock();
+ defer m.unlock();
+ if (!m.connected) return null;
+ return m.allocFid();
+ }
+
+ pub fn dropFid(m: *Self, fid: u32) void {
+ m.lock();
+ defer m.unlock();
+ m.freeFid(fid);
+ }
+
+ /// Blocks until the shared upstream answers `request`, copying the
+ /// reply into `r.buf`. Returns the decoded reply. Fids in `request`
+ /// are upstream fids (from `takeFid`); the caller owns their lifetime.
+ pub fn rpc(m: *Self, request: c9.Client.Request, r: *Rpc) !wire.Decoded {
+ r.* = .{ .buf = r.buf };
+ m.lock();
+ if (!m.connected) {
+ m.unlock();
+ return error.Disconnected;
+ }
+ const tag = m.submitLocked(request) catch |err| {
+ m.unlock();
+ return err;
+ };
+ r.tag = tag;
+ m.pending[tag] = .{ .active = true, .kind = .plain, .conn = null, .rpc = r, .down_tag = 0 };
+ m.unlock();
+ r.event.wait(m.io) catch {
+ // The waiting fiber was cancelled (HTTP timeout, follow
+ // teardown). Its `r` is about to be freed, so make sure the
+ // mux stops referencing it before returning: Tflush the
+ // upstream op and wait, uncancelably, until its pending clears.
+ m.cancelAndDrain(r);
+ return error.Canceled;
+ };
+ if (!r.ready or r.len == 0) return error.Upstream;
+ return wire.decode(r.buf[0..r.len]) catch return error.Upstream;
+ }
+
+ /// Guarantees `r` is no longer referenced by any pending slot before
+ /// the caller frees it. Best-effort Tflush of `r`'s upstream tag (the
+ /// flush's own reply clears `r`'s pending in `route`); then an
+ /// uncancelable wait for that to happen. A late reply or a disconnect
+ /// also clears it, so this returns even if the flush cannot be sent.
+ fn cancelAndDrain(m: *Self, r: *Rpc) void {
+ m.lock();
+ var live = false;
+ for (&m.pending) |*p| if (p.active and p.rpc == r) {
+ live = true;
+ };
+ if (!live) {
+ m.unlock();
+ return;
+ }
+ r.event.reset();
+ if (m.connected) {
+ if (m.client.submit(.{ .flush = .{ .oldtag = r.tag } })) |ftag| {
+ m.pending[ftag] = .{ .active = true, .kind = .flush, .conn = null, .flush_up = r.tag };
+ const bytes = m.client.output();
+ m.writeUpstream(bytes) catch {};
+ m.client.wrote(bytes.len);
+ } else |_| {}
+ }
+ m.unlock();
+ while (true) {
+ m.lock();
+ var still = false;
+ for (&m.pending) |*p| if (p.active and p.rpc == r) {
+ still = true;
+ };
+ m.unlock();
+ if (!still) return;
+ r.event.waitUncancelable(m.io);
+ r.event.reset();
+ }
+ }
+
+ // -- reader task ---------------------------------------------------
+
+ /// Reads upstream replies forever, routing each to its downstream.
+ /// Reconnects on failure. Runs on its own task; ended by `stop`.
+ pub fn readerLoop(m: *Self) void {
+ while (!m.stopping.load(.acquire)) {
+ if (!m.connected) {
+ m.reconnect();
+ if (m.stopping.load(.acquire)) return;
+ continue;
+ }
+ const frame = transport.readFrame(m.ureader, &m.stage, msize) catch {
+ m.onDisconnect();
+ if (m.stopping.load(.acquire)) return;
+ m.reconnect();
+ continue;
+ };
+ m.deliver(frame);
+ }
+ }
+
+ const Outgoing = struct { sink: Sink, tag: u8, len: usize };
+
+ fn deliver(m: *Self, frame: []const u8) void {
+ var outs: [17]Outgoing = undefined;
+ var nouts: usize = 0;
+ m.lock();
+ if (m.client.push(frame) != frame.len) {
+ m.unlock();
+ m.onDisconnect();
+ m.reconnect();
+ return;
+ }
+ while (m.client.take()) |done| {
+ const p = &m.pending[done.tag];
+ if (!p.active) continue;
+ if (p.rpc) |r| {
+ const built = m.route(done, p, r.buf);
+ p.active = false;
+ r.len = built orelse 0;
+ r.failed = done.result == .fail;
+ r.ready = true;
+ r.event.set(m.io);
+ continue;
+ }
+ const built = m.route(done, p, &m.frames[done.tag]);
+ const conn = p.conn;
+ p.active = false;
+ if (built) |len| if (conn != null) {
+ outs[nouts] = .{ .sink = p.sink, .tag = @intCast(done.tag), .len = len };
+ nouts += 1;
+ };
+ }
+ m.cond.broadcast(m.io);
+ m.unlock();
+ // Send outside the lock; a slow downstream cannot stall the client.
+ for (outs[0..nouts]) |o| o.sink.send(o.sink.ctx, m.frames[o.tag][0..o.len]) catch {};
+ }
+
+ /// Builds the downstream reply frame into `buf`, updating fid state.
+ /// Returns the encoded length, or null if nothing should be sent.
+ fn route(m: *Self, done: c9.Client.Done, p: *Pending, buf: []u8) ?usize {
+ // Fid bookkeeping first (independent of whether we send).
+ switch (p.kind) {
+ .walk => {
+ const full = done.result == .attach or
+ (done.result == .walk and done.result.walk.nwqid == p.walk_names);
+ if (p.conn) |c| {
+ if (full) {
+ if (!p.inplace) c.mapAdd(p.newfid_down, p.newfid_up);
+ } else if (!p.inplace) {
+ m.freeFid(p.newfid_up);
+ }
+ } else if (!p.inplace and !full) {
+ m.freeFid(p.newfid_up);
+ } else if (p.conn == null and full and !p.inplace) {
+ // Session gone: reclaim the fid the server just bound.
+ m.clunkUpstreamLocked(p.newfid_up);
+ }
+ },
+ .clunk => {
+ if (p.conn) |c| {
+ if (c.mapFind(p.clunk_down)) |e| {
+ m.freeFid(e.up);
+ e.used = false;
+ }
+ }
+ // fire-and-forget clunk (conn==null): the fid was already freed.
+ },
+ .flush => {
+ // The flushed request gets no reply after Rflush: clear its
+ // pending and wake any blocking rpc waiting on it.
+ const fp = &m.pending[p.flush_up];
+ if (fp.active) {
+ fp.active = false;
+ if (fp.rpc) |rr| {
+ rr.failed = true;
+ rr.ready = false;
+ rr.event.set(m.io);
+ }
+ }
+ },
+ .plain => {},
+ }
+ if (p.conn == null and p.rpc == null) return null;
+
+ const reply: wire.Msg = switch (done.result) {
+ .fail => |ename| .{ .rerror = .{ .ename = ename } },
+ .version => return null,
+ .auth => |q| .{ .rauth = .{ .aqid = q } },
+ .attach => |q| .{ .rattach = .{ .qid = q } },
+ .walk => |w| .{ .rwalk = .{ .nwqid = w.nwqid, .wqid = w.wqid } },
+ .open => |o| .{ .ropen = .{ .qid = o.qid, .iounit = o.iounit } },
+ .create => |cr| .{ .rcreate = .{ .qid = cr.qid, .iounit = cr.iounit } },
+ .read => |data| .{ .rread = .{ .data = data } },
+ .write => |n| .{ .rwrite = .{ .count = n } },
+ .clunk => .rclunk,
+ .remove => .rremove,
+ .stat => |s| .{ .rstat = .{ .stat = s } },
+ .wstat => .rwstat,
+ .flush => .rflush,
+ };
+ const encoded = wire.encode(reply, p.down_tag, buf) catch {
+ return (wire.encode(.{ .rerror = .{ .ename = "Invalid argument" } }, p.down_tag, buf) catch return null).len;
+ };
+ return encoded.len;
+ }
+
+ fn onDisconnect(m: *Self) void {
+ var outs: [17]Outgoing = undefined;
+ var nouts: usize = 0;
+ m.lock();
+ m.connected = false;
+ m.upstream_up.store(0, .release);
+ m.client.hangup();
+ // Fail every outstanding request so no downstream hangs.
+ for (&m.pending, 0..) |*p, i| if (p.active) {
+ p.active = false;
+ if (p.rpc) |r| {
+ r.failed = true;
+ r.ready = false;
+ r.event.set(m.io);
+ } else if (p.conn != null) {
+ const frame = m.errorFrame(&m.frames[i], p.down_tag, "Transport endpoint is not connected");
+ outs[nouts] = .{ .sink = p.sink, .tag = @intCast(i), .len = frame.len };
+ nouts += 1;
+ }
+ };
+ for (&m.fid_used) |*u| u.* = false;
+ m.fid_next = 0;
+ m.closeUpstream();
+ m.cond.broadcast(m.io);
+ m.unlock();
+ for (outs[0..nouts]) |o| o.sink.send(o.sink.ctx, m.frames[o.tag][0..o.len]) catch {};
+ }
+
+ fn reconnect(m: *Self) void {
+ while (!m.stopping.load(.acquire)) {
+ m.io.sleep(.fromMilliseconds(200), .awake) catch return;
+ m.dial() catch continue;
+ m.client = .init(.{ .in = &m.in, .out = &m.out });
+ m.handshake() catch {
+ m.closeUpstream();
+ continue;
+ };
+ m.lock();
+ m.generation +%= 1;
+ if (m.generation == 0) m.generation = 1;
+ m.connected = true;
+ m.upstream_up.store(1, .release);
+ _ = m.reconnects.fetchAdd(1, .monotonic);
+ m.cond.broadcast(m.io);
+ m.unlock();
+ return;
+ }
+ }
+
+ pub fn stop(m: *Self) void {
+ m.stopping.store(true, .release);
+ m.lock();
+ m.closeUpstream();
+ m.connected = false;
+ m.cond.broadcast(m.io);
+ m.unlock();
+ }
+ };
+}