diff options
| -rw-r--r-- | .agents/skills/pardes-9p/SKILL.md | 2 | ||||
| -rw-r--r-- | docs/fs.md | 9 | ||||
| -rw-r--r-- | src/fs-help.txt | 2 | ||||
| -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 | ||||
| -rw-r--r-- | src/pardes.zig | 9 | ||||
| -rw-r--r-- | test/fs.py | 20 |
10 files changed, 176 insertions, 32 deletions
diff --git a/.agents/skills/pardes-9p/SKILL.md b/.agents/skills/pardes-9p/SKILL.md index 2867ed18..81853c3e 100644 --- a/.agents/skills/pardes-9p/SKILL.md +++ b/.agents/skills/pardes-9p/SKILL.md @@ -47,7 +47,7 @@ $m/status pid, version, panes $m/look write a line = a right click on it at the active pane $m/exec write a line = a middle click: an editor command word, or a shell line $m/log recent events, then EOF: new|del|rename|save <serial> <name>, msg <serial|-> <text> - (exec 3<>$m/log; echo follow >&3; cat <&3 waits for new ones) + (exec 3<>$m/log; echo follow >&3; cat <&3 waits for new ones; tail -f does not) $m/screen the rendered screen as JSON, frozen per open $m/listeners this session's dial addresses $m/pane/new open it to make a pane, read names it; rmdir $m/pane/<n> closes it @@ -186,12 +186,17 @@ an open renders its frame; stat the entry. `/log` is one ring (64 KiB) that records whether or not anyone reads it: `new`, `del`, `rename` and `save <serial> <name>`, and `msg <serial|-> <text>` -for every line the editor says, builtins announcing themselves included. +for every line the editor says, repeats included (with `verbose` on, that +includes each builtin announcing itself as it runs). A `msg` said while a +pane is being made can precede that pane's `new`; panes present at boot are +recorded before anything else. Control characters in a record become spaces, so a record is one line. An open freezes the ring's text the way `/screen` freezes a frame: reads walk it and end. Writing `follow` to that same open makes reads past it park for the next record, one per read; a follower the ring outran reads `lost N` first. -Closing the open is the only way back, as with rio's `consctl`. `/screen` returns JSON with +Closing the open is the only way back, as with rio's `consctl`. A record +longer than a read comes in pieces, so a shell's `read` loop works. `tail -f` +never writes `follow`, so it sees nothing new: use the follow open instead. `/screen` returns JSON with `cols`, `rows`, `cursor`, a `styles` table, and row-major `cells` of `[grapheme, style_index]`. Each open freezes one frame until close. A terminal `body` freezes its history on the first read of each open handle; diff --git a/src/fs-help.txt b/src/fs-help.txt index dca2c14e..db046946 100644 --- a/src/fs-help.txt +++ b/src/fs-help.txt @@ -6,7 +6,7 @@ status pid, version and pane count look write a line: a right click on it at the active pane; read: the serials it touched exec write a line: a middle click, an editor command word or a shell line; read the same log recent events, one a line: new/del/rename/save <serial> <name>, and - msg <serial|-> <text> for what the editor said; write follow to wait + msg <serial|-> <text>, what the editor said; write follow to wait (tail -f won't) screen rendered screen as JSON, frozen from open to close listeners the session's dial addresses pane/new open it to make a pane; the read answers that pane's serial 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); diff --git a/src/pardes.zig b/src/pardes.zig index 6d294230..06bc81ed 100644 --- a/src/pardes.zig +++ b/src/pardes.zig @@ -6604,6 +6604,9 @@ pub const Pardes = struct { // A config init that opens a file with `Look …` armed the pulse before // anyone touched anything. Nobody asked for that, so boot is silent. _ = p.takeHaptic(); + // The ring is reserved whole now, so recording never allocates later. + try p.fs.log.buf.ensureTotalCapacityPrecise(p.gpa, p.fs.log.cap); + ctlfs.events.announce(p); // /log starts with the panes it booted with return p; } @@ -7438,6 +7441,9 @@ pub const Pardes = struct { if (text.len == 0) return; const serial: u32 = if (id < MAX_PANES) if (p.panes[id]) |pane| pane.serial else 0 else 0; const kept = text[0..@min(text.len, LoggedMessage.cap)]; + // /log hears every one: a client that retried and failed the same way + // is waiting on that second line. Only the +Messages view collapses. + ctlfs.events.noteMessage(p, serial, kept); if (p.messages_len > 0) { const last = &p.messages[(p.messages_head + limits.message_log - 1) % limits.message_log]; if (last.serial == serial and @@ -7451,7 +7457,6 @@ pub const Pardes = struct { return; } } - ctlfs.events.noteMessage(p, serial, kept); const slot = &p.messages[p.messages_head]; slot.len = @intCast(kept.len); @memcpy(slot.text[0..slot.len], kept); @@ -13635,6 +13640,8 @@ pub const Pardes = struct { p.sync(); p.presentation.enabled = true; _ = p.takeHaptic(); // see init: a restored session is not a gesture + try p.fs.log.buf.ensureTotalCapacityPrecise(p.gpa, p.fs.log.cap); + ctlfs.events.announce(p); return p; } @@ -142,7 +142,7 @@ def walk_tree(client, path='/', skip=('/log', '/os', '/screen')): def discovery(binary, embedded=False): with tempfile.TemporaryDirectory(prefix='pardes-discovery-') as directory: root = Path(directory) - with session(binary, root, 'discovery') as (client, _): + with session(binary, root, 'discovery') as (client, address): top = client.list('/') assert top[:10] == ['README', 'index', 'status', 'look', 'exec', 'log', 'screen', 'listeners', 'pane', 'os'], top @@ -197,14 +197,26 @@ def discovery(binary, embedded=False): assert 'new' in client.list('/pane'), 'new should be visible to ls' client.stat('/pane/new') assert newest(client) == second, 'a stat of new made a pane' - # What the editor said is in the stream too, and the fixture's own - # `new` lands on whichever side of the open its first update did. + # What the editor said is in the stream too. def pane_event(): - while (record := client.read_fid(log, frozen)).startswith((b'msg ', f'new {fixture} '.encode())): + while (record := client.read_fid(log, frozen)).startswith(b'msg '): pass return record assert pane_event() == f'new {first} {root}/+New\n'.encode() assert pane_event() == f'new {second} {root}/+New\n'.encode() + # A follower that is waiting wakes when the editor logs, not when + # some unrelated request next happens to wait. + with Client(address) as follower: + waiting = follower.open('/log', 2) + start = len(follower.read_fid(waiting)) + follower.rpc(118, struct.pack('<IQI', waiting, 0, 7) + b'follow\n') + woke = [] + reader = threading.Thread(target=lambda: woke.append(follower.read_fid(waiting, start)), daemon=True) + reader.start() + time.sleep(.2) + client.write('/exec', b'Msg woken\n') + reader.join(5) + assert woke and woke[0].endswith(b' woken\n'), woke assert set(client.list('/pane')) == {'new', str(fixture), str(first), str(second)} client.write(f'/pane/{first}/body', b'first pane', truncate=True) assert client.read(f'/pane/{first}/body') == b'first pane' |
