diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-20 03:53:17 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-09-20 03:53:17 -0300 |
| commit | f1b53c1533539aecbf16ad19fd9156deae091f92 (patch) | |
| tree | 38ff00810ef5e0ad271c26e218cf6af4578c71ce /web/mux.zig | |
| parent | ba7ec40782ba7020d82a56896a5eb1b52578d6aa (diff) | |
| download | cloud9-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.zig | 812 |
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(); + } + }; +} |
