//! 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(); } }; }