summaryrefslogtreecommitdiff
path: root/test/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 /test/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 'test/web/mux.zig')
-rw-r--r--test/web/mux.zig310
1 files changed, 310 insertions, 0 deletions
diff --git a/test/web/mux.zig b/test/web/mux.zig
new file mode 100644
index 0000000..ab95281
--- /dev/null
+++ b/test/web/mux.zig
@@ -0,0 +1,310 @@
+//! Unit tests for the 9web multiplexer (web/mux.zig): one shared upstream 9P
+//! connection fanned out to several downstream sessions. The upstream is an
+//! in-process `serve.Runner` over an `fs.Server` backend with a blocking
+//! "event" file (its reads park until a write to "data" wakes them). The tests
+//! drive the mux with hand-built 9P frames and check fid isolation, a parked
+//! read released by a 9P-side write, Tflush forwarding, and reconnection with
+//! downstream fids invalidated.
+const std = @import("std");
+const c9 = @import("cloud9");
+const mux = @import("mux");
+const wire = c9.wire;
+const fs = c9.fs;
+const serve = c9.serve;
+const transport = c9.transport;
+const Io = std.Io;
+const testing = std.testing;
+
+const msize = 8192;
+const M = mux.Mux(.{ .msize = msize, .upstream_fids = 64, .fids_per_conn = 32, .downstreams = 8 });
+
+// -- the upstream backend: root/{event,data} ---------------------------------
+
+const root_node = 1;
+const event_node = 2;
+const data_node = 3;
+
+const EventFs = struct {
+ pub const Req = fs.Req;
+ pub const Reply = fs.Reply;
+
+ mutex: Io.Mutex = .init,
+ posted: ?[]const u8 = null,
+ store: [256]u8 = undefined,
+
+ fn attrOf(node: u64) fs.Attr {
+ return switch (node) {
+ root_node => .{ .name = "/", .node = root_node, .dir = true, .mode = 0o500 },
+ event_node => .{ .name = "event", .node = event_node, .mode = 0o400 },
+ data_node => .{ .name = "data", .node = data_node, .mode = 0o600 },
+ else => unreachable,
+ };
+ }
+};
+
+const Runner = serve.Runner(EventFs, .{ .fid_capacity = 32, .slot_capacity = 8 }, .{ .msize = msize, .connections = 2, .listeners = 1 });
+
+fn serveFn(ctx: ?*anyopaque, conn: *Runner.Conn, req: fs.Req) void {
+ const st: *EventFs = @ptrCast(@alignCast(ctx.?));
+ st.mutex.lockUncancelable(testing.io);
+ defer st.mutex.unlock(testing.io);
+ const fail: fs.Reply = .fail(req.tag, fs.E.NOENT);
+ switch (req.op) {
+ .lookup => {
+ if (std.mem.eql(u8, req.data, "..")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(root_node) }, "");
+ if (std.mem.eql(u8, req.data, "event")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(event_node) }, "");
+ if (std.mem.eql(u8, req.data, "data")) return conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(data_node) }, "");
+ return conn.reply(&fail, "");
+ },
+ .getattr, .setattr => conn.reply(&.{ .tag = req.tag, .attr = EventFs.attrOf(req.node) }, ""),
+ .open => conn.reply(&.{ .tag = req.tag, .handle = 7 }, ""),
+ .release => conn.reply(&.{ .tag = req.tag }, ""),
+ .readdir => conn.reply(&.{ .tag = req.tag }, ""),
+ .read => {
+ if (req.node == event_node) {
+ if (st.posted) |bytes| {
+ st.posted = null;
+ return conn.reply(&.{ .tag = req.tag }, bytes);
+ }
+ return conn.reply(&.{ .tag = req.tag, .status = .again }, "");
+ }
+ conn.reply(&.{ .tag = req.tag }, ""); // data reads EOF
+ },
+ .write => {
+ const n = @min(req.data.len, st.store.len);
+ @memcpy(st.store[0..n], req.data[0..n]);
+ st.posted = st.store[0..n];
+ conn.reply(&.{ .tag = req.tag, .written = @intCast(n) }, "");
+ conn.wake(); // retry the parked event read on this connection
+ },
+ }
+}
+
+// -- a downstream session that captures replies -------------------------------
+
+const Down = struct {
+ m: *M,
+ conn: M.Conn = undefined,
+ mutex: Io.Mutex = .init,
+ cond: Io.Condition = .init,
+ ring: [8][msize]u8 = undefined,
+ len: [8]usize = @splat(0),
+ head: usize = 0,
+ tail: usize = 0,
+ count: usize = 0,
+
+ fn init(d: *Down, m: *M) void {
+ d.* = .{ .m = m };
+ d.conn = M.Conn.init(m, .{ .ctx = d, .send = sink });
+ }
+ fn deinit(d: *Down) void {
+ d.conn.deinit();
+ }
+ fn sink(ctx: *anyopaque, frame: []const u8) anyerror!void {
+ const d: *Down = @ptrCast(@alignCast(ctx));
+ d.mutex.lockUncancelable(testing.io);
+ defer d.mutex.unlock(testing.io);
+ const i = d.tail;
+ const n = @min(frame.len, msize);
+ @memcpy(d.ring[i][0..n], frame[0..n]);
+ d.len[i] = n;
+ d.tail = (i + 1) % d.ring.len;
+ d.count += 1;
+ d.cond.signal(testing.io);
+ }
+ fn send(d: *Down, frame: []const u8) void {
+ d.m.forward(&d.conn, frame);
+ }
+ fn recv(d: *Down) !wire.Decoded {
+ d.mutex.lockUncancelable(testing.io);
+ defer d.mutex.unlock(testing.io);
+ var tries: usize = 0;
+ while (d.count == 0) : (tries += 1) {
+ if (tries > 4000) return error.Timeout;
+ d.cond.wait(testing.io, &d.mutex) catch return error.Canceled;
+ }
+ const i = d.head;
+ d.head = (i + 1) % d.ring.len;
+ d.count -= 1;
+ return wire.decode(d.ring[i][0..d.len[i]]);
+ }
+};
+
+// -- frame builders -----------------------------------------------------------
+
+var enc_buf: [msize]u8 = undefined;
+fn enc(msg: wire.Msg, tag: u16) []const u8 {
+ return wire.encode(msg, tag, &enc_buf) catch unreachable;
+}
+
+fn handshake(d: *Down) !void {
+ d.send(enc(.{ .tversion = .{ .msize = msize, .version = "9P2000" } }, wire.notag));
+ const v = try d.recv();
+ try testing.expectEqual(wire.Type.rversion, v.msg.msgType());
+ d.send(enc(.{ .tattach = .{ .fid = 0, .afid = wire.nofid, .uname = "u", .aname = "" } }, 1));
+ const a = try d.recv();
+ try testing.expectEqual(wire.Type.rattach, a.msg.msgType());
+}
+
+fn walk1(d: *Down, newfid: u32, name: []const u8) !void {
+ var wn: [wire.max_welem][]const u8 = @splat("");
+ wn[0] = name;
+ d.send(enc(.{ .twalk = .{ .fid = 0, .newfid = newfid, .nwname = 1, .wname = wn } }, 2));
+ const r = try d.recv();
+ try testing.expectEqual(wire.Type.rwalk, r.msg.msgType());
+ try testing.expectEqual(@as(u16, 1), r.msg.rwalk.nwqid);
+}
+
+// -- the rig ------------------------------------------------------------------
+
+const Rig = struct {
+ dir: testing.TmpDir,
+ path_buf: [std.fs.max_path_bytes]u8 = undefined,
+ sock_buf: [transport.sun_path_len]u8 = undefined,
+ sock: [:0]const u8 = undefined,
+ fsx: EventFs = .{},
+ runner: *Runner = undefined,
+ m: *M = undefined,
+ reader: Io.Future(void) = undefined,
+
+ fn start(r: *Rig) !void {
+ const io = testing.io;
+ r.dir = testing.tmpDir(.{});
+ const plen = try r.dir.dir.realPath(io, &r.path_buf);
+ r.sock = try std.fmt.bufPrintSentinel(&r.sock_buf, "{s}/up", .{r.path_buf[0..plen]}, 0);
+ r.runner = try testing.allocator.create(Runner);
+ r.runner.init(.{ .io = io, .root = root_node, .handler = .{ .serve = serveFn, .ctx = &r.fsx } });
+ _ = try r.runner.listen(.{ .unix = r.sock }, 4);
+ r.m = try testing.allocator.create(M);
+ r.m.init(io, .{ .network = .{ .unix = r.sock } }, "u", "");
+ try r.m.connect();
+ r.reader = try io.concurrent(M.readerLoop, .{r.m});
+ }
+
+ fn stopUpstream(r: *Rig) void {
+ r.runner.stop();
+ Io.Dir.cwd().deleteFile(testing.io, r.sock) catch {};
+ }
+
+ fn restartUpstream(r: *Rig) !void {
+ r.runner.init(.{ .io = testing.io, .root = root_node, .handler = .{ .serve = serveFn, .ctx = &r.fsx } });
+ _ = try r.runner.listen(.{ .unix = r.sock }, 4);
+ }
+
+ fn end(r: *Rig) void {
+ r.m.stop();
+ r.reader.cancel(testing.io);
+ r.runner.stop();
+ testing.allocator.destroy(r.runner);
+ testing.allocator.destroy(r.m);
+ r.dir.cleanup();
+ }
+};
+
+test "mux: two downstreams share one upstream with isolated fids" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&a);
+ try handshake(&b);
+ // Both use downstream fid 1; the mux maps them to distinct upstream fids.
+ try walk1(&a, 1, "event");
+ try walk1(&b, 1, "data");
+ // A opens its event; B opens its data. Independent handles.
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ try testing.expectEqual(wire.Type.ropen, (try a.recv()).msg.msgType());
+ b.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.owrite } }, 3));
+ try testing.expectEqual(wire.Type.ropen, (try b.recv()).msg.msgType());
+}
+
+test "mux: a parked read is released by another downstream's write" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&a);
+ try handshake(&b);
+ try walk1(&a, 1, "event");
+ try walk1(&b, 1, "data");
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ _ = try a.recv();
+ b.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.owrite } }, 3));
+ _ = try b.recv();
+ // A's read parks upstream (no reply yet).
+ a.send(enc(.{ .tread = .{ .fid = 1, .offset = 0, .count = 128 } }, 4));
+ try testing.io.sleep(.fromMilliseconds(30), .awake);
+ // B writes "hello", which posts the event and wakes A's parked read.
+ b.send(enc(.{ .twrite = .{ .fid = 1, .offset = 0, .data = "hello" } }, 4));
+ try testing.expectEqual(wire.Type.rwrite, (try b.recv()).msg.msgType());
+ const rr = try a.recv();
+ try testing.expectEqual(wire.Type.rread, rr.msg.msgType());
+ try testing.expectEqualStrings("hello", rr.msg.rread.data);
+}
+
+test "mux: Tflush is forwarded and cancels a parked read" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ try handshake(&a);
+ try walk1(&a, 1, "event");
+ a.send(enc(.{ .topen = .{ .fid = 1, .mode = c9.oread } }, 3));
+ _ = try a.recv();
+ a.send(enc(.{ .tread = .{ .fid = 1, .offset = 0, .count = 128 } }, 5));
+ try testing.io.sleep(.fromMilliseconds(30), .awake);
+ a.send(enc(.{ .tflush = .{ .oldtag = 5 } }, 6));
+ // Expect the interrupted read (Rerror) and the Rflush, in either order.
+ var saw_err = false;
+ var saw_flush = false;
+ for (0..2) |_| {
+ const got = try a.recv();
+ switch (got.msg.msgType()) {
+ .rerror => saw_err = true,
+ .rflush => saw_flush = true,
+ else => return error.Unexpected,
+ }
+ }
+ try testing.expect(saw_err and saw_flush);
+}
+
+test "mux: an upstream drop errors downstream, and it reconnects" {
+ var rig: Rig = .{ .dir = undefined };
+ try rig.start();
+ defer rig.end();
+ var a: Down = undefined;
+ a.init(rig.m);
+ defer a.deinit();
+ try handshake(&a);
+ try walk1(&a, 1, "event");
+ // Drop the upstream. An in-flight/next request must error, not hang.
+ rig.stopUpstream();
+ try testing.io.sleep(.fromMilliseconds(50), .awake);
+ a.send(enc(.{ .tstat = .{ .fid = 1 } }, 7));
+ const err = try a.recv();
+ try testing.expectEqual(wire.Type.rerror, err.msg.msgType());
+ // Bring the upstream back; the mux reconnects and a fresh attach works.
+ try rig.restartUpstream();
+ var tries: usize = 0;
+ while (tries < 200) : (tries += 1) {
+ try testing.io.sleep(.fromMilliseconds(20), .awake);
+ if (std.mem.eql(u8, rig.m.upstreamState(), "connected")) break;
+ }
+ var b: Down = undefined;
+ b.init(rig.m);
+ defer b.deinit();
+ try handshake(&b); // a new session on the reconnected upstream
+ try testing.expect(rig.m.reconnectCount() >= 1);
+}