diff options
| author | Gabriel Schneider <[email protected]> | 2026-09-27 17:52:31 -0300 |
|---|---|---|
| committer | Gabriel Schneider <[email protected]> | 2026-10-01 00:12:14 -0300 |
| commit | f4412d6dfdb1e2bd8549a7edba785ee242fd38f2 (patch) | |
| tree | a48a67c60ce0ee55be8db3691f08f36c1986d7f0 /src/ninep | |
| parent | 542dd489149e33915d19878a277e23b3a9c0b070 (diff) | |
| download | pardes-f4412d6dfdb1e2bd8549a7edba785ee242fd38f2.tar.gz pardes-f4412d6dfdb1e2bd8549a7edba785ee242fd38f2.zip | |
Wake reads waiting on /log, event and pty/data; one reader per consuming file; pty/ctl reads back
A read parked on a queue was retried only when some unrelated write had
to wait, so a follower of /log (and a reader of event or pty/data) slept
until then. Pushing a record now marks the turn parked, and giving the
turn up wakes them (measured: stuck past 3 s before, 0 s after).
event and pty/data consume what they read, so a second open for reading
is refused with rio's "file in use" (EBUSY through 9ns); writers still
get in, and pty/data queues output only for an actual reader. pty/ctl
reads back "winsize C R", in the words it takes.
/log fixes from review: a record longer than a read comes in pieces (a
shell read loop failed on long lines), every repeated message is logged,
a record bigger than the ring is cut to fit instead of emptying it, the
ring is reserved at boot so recording never allocates, and panes present
at boot are recorded first.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
Diffstat (limited to 'src/ninep')
| -rw-r--r-- | src/ninep/events.zig | 90 | ||||
| -rw-r--r-- | src/ninep/pane.zig | 8 | ||||
| -rw-r--r-- | src/ninep/pty.zig | 33 | ||||
| -rw-r--r-- | src/ninep/screen.zig | 2 | ||||
| -rw-r--r-- | src/ninep/tree.zig | 33 |
5 files changed, 143 insertions, 23 deletions
diff --git a/src/ninep/events.zig b/src/ninep/events.zig index 472a114f..ebd1cd2c 100644 --- a/src/ninep/events.zig +++ b/src/ninep/events.zig @@ -50,6 +50,10 @@ pub const Queue = struct { q.buf.shrinkRetainingCapacity(q.buf.items.len - 4); return; }; + // A read parked on this queue is retried only when the turn is next + // given up with `parked` set; without it a follower sleeps until some + // unrelated request happens to park. + pardes.turn.parked = true; } pub fn peek(q: *const Queue) ?[]const u8 { @@ -163,7 +167,11 @@ fn pushLog(p: *Pardes, record: []u8) void { for (record[0 .. record.len - 1]) |*c| if (c.* < ' ') { c.* = ' '; }; - p.fs.log.push(p.gpa, record); + // One record larger than the ring would push every other out and then + // not fit itself; cut it to what fits instead. + const fit = @min(record.len, p.fs.log.cap - 4); + record[fit - 1] = '\n'; + p.fs.log.push(p.gpa, record[0..fit]); } /// An open freezes the ring's text, so `cat log` answers what happened lately @@ -208,6 +216,7 @@ pub fn readLog(p: *Pardes, req: Req) Reply { if (slot.next < q.dropped) { out.print(p.gpa, "lost {d}\n", .{q.dropped - slot.next}) catch return Reply.fail(req.tag, E.NOMEM); slot.next = q.dropped; + slot.part = 0; return .{ .tag = req.tag, .payload = .{ .staged = @intCast(out.items.len) } }; } // ponytail: walks the ring from its oldest record; it holds at most log_bytes. @@ -216,10 +225,17 @@ pub fn readLog(p: *Pardes, req: Req) Reply { while (at + 4 <= q.buf.items.len) : (seq += 1) { const len = std.mem.readInt(u32, q.buf.items[at..][0..4], .little); if (seq == slot.next) { - if (req.size < len) return Reply.fail(req.tag, E.INVAL); - out.appendSlice(p.gpa, q.buf.items[at + 4 ..][0..len]) catch return Reply.fail(req.tag, E.NOMEM); - slot.next += 1; - return .{ .tag = req.tag, .payload = .{ .staged = @intCast(len) } }; + // A short read (a shell's `read` takes a few bytes at a time) + // gets the record in pieces, as pty/data does. + const rest = q.buf.items[at + 4 + slot.part ..][0 .. len - slot.part]; + const n = @min(rest.len, req.size); + out.appendSlice(p.gpa, rest[0..n]) catch return Reply.fail(req.tag, E.NOMEM); + slot.part += @intCast(n); + if (slot.part == len) { + slot.next += 1; + slot.part = 0; + } + return .{ .tag = req.tag, .payload = .{ .staged = @intCast(n) } }; } at += 4 + len; } @@ -535,29 +551,39 @@ test "a write through the filesystem is reported once, attributed to the file it try testing.expectEqual(Status.again, rd(p, event, 0, 4096).reply.status); } -test "two event readers each count once, and the second closing leaves the first" { +test "one event reader at a time, and a writer still holds the pane" { const gpa = testing.allocator; const p = try withFile(gpa, "x\n"); defer p.deinit(); const event = Node.of(serialOf(p), .event); - _ = call(p, .{ .tag = 12, .op = .open, .node = event }); - _ = call(p, .{ .tag = 13, .op = .open, .node = event }); + const reader = call(p, .{ .tag = 12, .op = .open, .node = event }); + try testing.expectEqual(Status.ok, reader.reply.status); + // Reads consume, so a second reader would see half the clicks: refused. + const second = call(p, .{ .tag = 13, .op = .open, .node = event, .omode = 2 }); + try testing.expectEqual(Status.err, second.reply.status); + try testing.expectEqualStrings(tree.e_in_use, second.reply.ename); + const writer = call(p, .{ .tag = 13, .op = .open, .node = event, .omode = 1 }); + try testing.expectEqual(Status.ok, writer.reply.status); try testing.expectEqual(@as(u16, 2), p.fs.panes[0].readers); try testing.expectEqual(@as(u16, 2), p.fs.listeners); - _ = call(p, .{ .tag = 14, .op = .release, .node = event }); + _ = call(p, .{ .tag = 14, .op = .release, .node = event, .handle = writer.reply.handle }); try testing.expectEqual(@as(u16, 1), p.fs.panes[0].readers); try testing.expect(p.fs.scripted(0)); _ = call(p, .{ .tag = 15, .op = .release, .node = Node.of(serialOf(p), .body) }); try testing.expectEqual(@as(u16, 1), p.fs.listeners); - _ = call(p, .{ .tag = 16, .op = .release, .node = event }); + _ = call(p, .{ .tag = 16, .op = .release, .node = event, .handle = reader.reply.handle }); try testing.expectEqual(@as(u16, 0), p.fs.listeners); try testing.expect(!p.fs.scripted(0)); _ = call(p, .{ .tag = 17, .op = .release, .node = event }); try testing.expectEqual(@as(u16, 0), p.fs.listeners); + // Once it closed, the next reader gets in. + const next = call(p, .{ .tag = 18, .op = .open, .node = event }); + try testing.expectEqual(Status.ok, next.reply.status); + _ = call(p, .{ .tag = 19, .op = .release, .node = event, .handle = next.reply.handle }); } test "a pane deleted while its event file is open leaves no suppression behind" { @@ -569,7 +595,7 @@ test "a pane deleted while its event file is open leaves no suppression behind" _ = try th.newPane(p); const a = call(p, .{ .tag = 18, .op = .open, .node = event }); - const b = call(p, .{ .tag = 19, .op = .open, .node = event }); + const b = call(p, .{ .tag = 19, .op = .open, .node = event, .omode = 1 }); try testing.expectEqual(@as(u16, 2), p.fs.listeners); _ = wr(p, Node.of(serial, .exec), "Del\n"); @@ -687,6 +713,21 @@ test "the log records whether or not anyone reads, and an open that follows wait ); try testing.expectEqual(Status.again, rdf.next(p, log, fh, frozen).reply.status); + // A shell's `read` asks for a few bytes at a time: the record comes in + // pieces rather than failing. Every repeat is its own line. + const said_by = p.panes[0].?.serial; + p.setMessage(0, "same failure"); + p.setMessage(0, "same failure"); + var pieced: [64]u8 = undefined; + var got: usize = 0; + while (got == 0 or pieced[got - 1] != '\n') { + const piece = call(p, .{ .tag = 9, .op = .read, .node = log, .handle = fh, .off = frozen, .size = 3 }).bytes; + @memcpy(pieced[got..][0..piece.len], piece); + got += piece.len; + } + try testing.expectEqualStrings(try std.fmt.bufPrint(&expected, "msg {d} same failure\n", .{said_by}), pieced[0..got]); + try testing.expectEqualStrings(try std.fmt.bufPrint(&expected, "msg {d} same failure\n", .{said_by}), rdf.next(p, log, fh, frozen).bytes); + // A follower the ring outran hears how much it missed, then carries on. var filler: [200]u8 = @splat('x'); for (0..p.fs.log.cap / filler.len + 8) |i| { @@ -698,4 +739,31 @@ test "the log records whether or not anyone reads, and an open that follows wait try testing.expect(std.mem.startsWith(u8, rdf.next(p, log, fh, frozen).bytes, "msg ")); _ = call(p, .{ .tag = 10, .op = .release, .node = log, .handle = fh }); for (p.fs.snapshots) |slot| try testing.expect(slot.node == 0); + + // A record bigger than the whole ring is cut to fit, not dropped with + // everything else pushed out ahead of it. + p.fs.log.cap = 128; + var huge: [300]u8 = @splat('y'); + p.setMessage(0, &huge); + const cut = call(p, .{ .tag = 11, .op = .open, .node = log }); + const kept = call(p, .{ .tag = 12, .op = .read, .node = log, .handle = cut.reply.handle, .size = 512 }).bytes; + try testing.expectEqual(@as(usize, 124), kept.len); + try testing.expect(std.mem.startsWith(u8, kept, "msg ") and kept[kept.len - 1] == '\n'); + _ = call(p, .{ .tag = 13, .op = .release, .node = log, .handle = cut.reply.handle }); +} + +test "opens of the log share the snapshot slots, and a closed one frees its slot" { + const gpa = testing.allocator; + const p = try withFile(gpa, "x\n"); + defer p.deinit(); + const log = @intFromEnum(tree.TopFile.log); + var handles: [tree.screen.snapshot_slots]u32 = undefined; + for (&handles) |*h| h.* = call(p, .{ .tag = 1, .op = .open, .node = log }).reply.handle; + try testing.expectEqual(E.NFILE, call(p, .{ .tag = 2, .op = .open, .node = log }).errno()); + _ = call(p, .{ .tag = 3, .op = .release, .node = log, .handle = handles[5] }); + const again = call(p, .{ .tag = 4, .op = .open, .node = log }); + try testing.expectEqual(Status.ok, again.reply.status); + handles[5] = again.reply.handle; + for (handles) |h| _ = call(p, .{ .tag = 5, .op = .release, .node = log, .handle = h }); + for (p.fs.snapshots) |slot| try testing.expect(slot.node == 0); } diff --git a/src/ninep/pane.zig b/src/ninep/pane.zig index 4b34cf96..40d292c0 100644 --- a/src/ninep/pane.zig +++ b/src/ninep/pane.zig @@ -33,6 +33,8 @@ pub const State = struct { addr: Range = .{}, limit: ?Range = null, readers: u16 = 0, + /// One of `readers` reads `event`; a second reading open is refused. + event_reader: bool = false, events: events.Queue = .{}, nomark: bool = false, noscroll: bool = false, @@ -221,8 +223,9 @@ pub fn fileSize(p: *Pardes, id: usize, f: PaneFile) u64 { .look, .exec => ctl.resultsLen(p), .event => events.pending(&pf.events), .pty_status => pty.status_len, + .pty_ctl => pty.ctlLen(pane), .pty_data => events.pending(&pf.pty_out), - .dir, .errors, .pty, .pty_ctl => 0, + .dir, .errors, .pty => 0, }; } @@ -289,7 +292,8 @@ pub fn read(p: *Pardes, req: Req, id: usize, pane: *Pane, file: PaneFile) Reply .event => events.readQueue(p, req, &pf.events), .pty_status => pty.readStatus(p, req, id, pane), .pty_data => pty.readData(p, req, pf), - .dir, .errors, .pty, .pty_ctl => Reply.fail(req.tag, E.PERM), + .pty_ctl => pty.readCtl(p, req, pane), + .dir, .errors, .pty => Reply.fail(req.tag, E.PERM), }; } diff --git a/src/ninep/pty.zig b/src/ninep/pty.zig index c657d27d..b34bd59d 100644 --- a/src/ninep/pty.zig +++ b/src/ninep/pty.zig @@ -76,6 +76,18 @@ fn verb(p: *Pardes, id: usize, line: []const u8, apply: bool) bool { return true; } +/// ctl reads back its state in the words it takes, as a Plan 9 ctl does. +pub fn readCtl(p: *Pardes, req: Req, pane: *Pane) Reply { + p.fs.stage(p.gpa).print(p.gpa, "winsize {d} {d}\n", .{ pane.cols, pane.rows }) catch + return Reply.fail(req.tag, E.NOMEM); + return tree.stagedReply(p, req); +} + +pub fn ctlLen(pane: *const Pane) u64 { + var buf: [32]u8 = undefined; + return (std.fmt.bufPrint(&buf, "winsize {d} {d}\n", .{ pane.cols, pane.rows }) catch unreachable).len; +} + /// Three right-aligned fields: cols, rows and whether the host holds the tty. pub const status_len: u64 = 36; @@ -183,7 +195,7 @@ test "a terminal pane's pty/ holds exactly ctl, status and data" { const ctl = look_up(p, Node.of(serial, .pty), "ctl"); try testing.expectEqual(Node.of(serial, .pty_ctl), ctl.reply.attr.node); - try testing.expectEqual(@as(u16, 0o222), ctl.reply.attr.mode); + try testing.expectEqual(@as(u16, 0o666), ctl.reply.attr.mode); try testing.expectEqual(@as(u16, 0o444), look_up(p, Node.of(serial, .pty), "status").reply.attr.mode); try testing.expectEqual(E.NOENT, look_up(p, Node.of(serial, .pty), "body").errno()); try testing.expectEqual(E.NOENT, look_up(p, Node.of(serial, .pty), "pty").errno()); @@ -191,7 +203,7 @@ test "a terminal pane's pty/ holds exactly ctl, status and data" { try testing.expectEqual(E.NOTDIR, look_up(p, Node.of(serial, .pty_ctl), "x").errno()); try testing.expectEqual(E.NOTDIR, rdir(p, Node.of(serial, .pty_ctl), 0).errno()); try testing.expectEqual(E.PERM, rd(p, Node.of(serial, .pty), 0, 16).errno()); - try testing.expectEqual(E.PERM, rd(p, Node.of(serial, .pty_ctl), 0, 16).errno()); + try testing.expect(std.mem.startsWith(u8, rd(p, Node.of(serial, .pty_ctl), 0, 64).bytes, "winsize ")); try testing.expectEqual(E.PERM, wr(p, Node.of(serial, .pty_status), "x").errno()); } @@ -204,6 +216,8 @@ test "every pty/ctl verb, and every refusal" { const pane = p.panes[0].?; const cols = pane.cols; const rows = pane.rows; + var said: [32]u8 = undefined; + try testing.expectEqualStrings(try std.fmt.bufPrint(&said, "winsize {d} {d}\n", .{ cols, rows }), rd(p, ctl, 0, 64).bytes); const ws = wr(p, ctl, "winsize 132 44\n"); try testing.expectEqual(@as(u32, "winsize 132 44\n".len), ws.reply.written); try testing.expectEqual(@as(u16, 132), ws.winsize.?.cols); @@ -331,7 +345,12 @@ test "pty/data writes at the shell and reads the raw stream" { try testing.expectEqual(@as(usize, 0), p.fs.panes[0].pty_out.buf.items.len); try testing.expectEqual(Status.again, rd(p, data, 0, 64).reply.status); - _ = call(p, .{ .tag = 5, .op = .open, .node = data }); + // An open that only writes is not a reader: nothing queues for it. + const writer = call(p, .{ .tag = 4, .op = .open, .node = data, .omode = 1 }); + try testing.expectEqual(@as(u16, 0), p.fs.panes[0].pty_readers); + _ = call(p, .{ .tag = 4, .op = .release, .node = data, .handle = writer.reply.handle }); + + const reader = call(p, .{ .tag = 5, .op = .open, .node = data }); try testing.expectEqual(@as(u16, 1), p.fs.panes[0].pty_readers); try testing.expectEqual(@as(u16, 0), p.fs.listeners); try testing.expect(!p.fs.scripted(0)); @@ -351,14 +370,16 @@ test "pty/data writes at the shell and reads the raw stream" { p.update(.{ .output = .{ .pane = 0, .bytes = "orphan" } }); while (p.nextEffect()) |_| {} - _ = call(p, .{ .tag = 6, .op = .release, .node = data }); + _ = call(p, .{ .tag = 6, .op = .release, .node = data, .handle = reader.reply.handle }); try testing.expectEqual(@as(u16, 0), p.fs.panes[0].pty_readers); try testing.expectEqual(@as(usize, 0), p.fs.panes[0].pty_out.buf.capacity); try testing.expectEqual(Status.again, rd(p, data, 0, 64).reply.status); + // A second reader would take half the stream from the first: refused. _ = call(p, .{ .tag = 7, .op = .open, .node = data }); - _ = call(p, .{ .tag = 8, .op = .open, .node = data }); - _ = call(p, .{ .tag = 9, .op = .release, .node = data }); + const second = call(p, .{ .tag = 8, .op = .open, .node = data, .omode = 2 }); + try testing.expectEqual(Status.err, second.reply.status); + try testing.expectEqualStrings(tree.e_in_use, second.reply.ename); try testing.expectEqual(@as(u16, 1), p.fs.panes[0].pty_readers); p.update(.{ .output = .{ .pane = 0, .bytes = "still" } }); while (p.nextEffect()) |_| {} diff --git a/src/ninep/screen.zig b/src/ninep/screen.zig index 245a31ef..1743fbb6 100644 --- a/src/ninep/screen.zig +++ b/src/ninep/screen.zig @@ -17,6 +17,8 @@ pub const Snapshot = struct { /// /log only: this open waits for records after `bytes`, from `next`. follow: bool = false, next: u64 = 0, + /// How much of record `next` a short read already took. + part: u32 = 0, }; pub const snapshot_slots = 32; diff --git a/src/ninep/tree.zig b/src/ninep/tree.zig index f98429ab..a85a1873 100644 --- a/src/ninep/tree.zig +++ b/src/ninep/tree.zig @@ -61,6 +61,20 @@ pub fn failText(tag: u64, errno: u16, text: []const u8) Reply { pub const e_bad_addr = "bad address syntax"; pub const e_bad_ctl = "ill-formed control message"; pub const e_bad_event = "bad event syntax"; +/// A second open reading a file whose reads consume: rio's word for it +/// (windows/rio/xfid.c:25), which 9ns turns into EBUSY. +pub const e_in_use = "file in use"; + +/// Handles that mark an open as the one reading `event` or `pty/data`, so the +/// release knows to give the file up; other opens get `open_handle`. +const open_handle: u32 = 1; +// ponytail: EBUSY spelled locally; use E.BUSY once cloud9 names it and pardes bumps the pin. +const e_busy: u16 = 16; +const reader_handle: u32 = 2; + +fn reads(omode: u8) bool { + return omode & 3 != 1; // OWRITE is the one access mode that never reads +} /// Tremove closes a pane; nothing else in the tree is created or destroyed by /// the protocol, and wstat stays a truncation. Making a pane is an open of @@ -167,7 +181,7 @@ pub const PaneFile = enum(u5) { pub fn mode(f: PaneFile) u16 { return switch (f) { .dir, .pty => 0o755, - .errors, .pty_ctl => 0o222, + .errors => 0o222, .pty_status => 0o444, else => 0o666, }; @@ -507,16 +521,26 @@ fn open(p: *Pardes, req: Req, target: Target) Reply { if (t.file.inPty() and !pn.isTerminal()) return Reply.fail(req.tag, E.NOENT); switch (t.file) { .body => if (pn.isTerminal()) return screen.openSnapshot(p, req, false), + // Reads of these consume, so two readers would each see half + // the stream; the second is refused rather than robbed. .event => { + const reader = reads(req.omode); + if (reader and pf.event_reader) return failText(req.tag, e_busy, e_in_use); pf.readers +|= 1; p.fs.listeners +|= 1; + pf.event_reader = pf.event_reader or reader; + if (reader) return .{ .tag = req.tag, .handle = reader_handle }; + }, + .pty_data => if (reads(req.omode)) { + if (pf.pty_readers > 0) return failText(req.tag, e_busy, e_in_use); + pf.pty_readers = 1; + return .{ .tag = req.tag, .handle = reader_handle }; }, - .pty_data => pf.pty_readers +|= 1, else => {}, } }, } - return .{ .tag = req.tag, .handle = 1 }; + return .{ .tag = req.tag, .handle = open_handle }; } /// Tremove closes a pane. Nothing else in the tree can be removed, and this @@ -549,12 +573,13 @@ fn releaseHandle(p: *Pardes, req: Req) Reply { const id = p.paneBySerial(t.serial) orelse return .{ .tag = req.tag }; const pf = &p.fs.panes[id]; if (t.file == .pty_data) { - if (pf.pty_readers == 0) return .{ .tag = req.tag }; + if (req.handle != reader_handle or pf.pty_readers == 0) return .{ .tag = req.tag }; pf.pty_readers -= 1; if (pf.pty_readers == 0) pf.pty_out.clearAndFree(p.gpa); return .{ .tag = req.tag }; } if (pf.readers == 0) return .{ .tag = req.tag }; + if (req.handle == reader_handle) pf.event_reader = false; pf.readers -= 1; p.fs.listeners -|= 1; if (pf.readers == 0) pf.tag_snap.clearAndFree(p.gpa); |
