summaryrefslogtreecommitdiff
path: root/src/ninep/events.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-27 17:52:31 -0300
committerGabriel Schneider <[email protected]>2026-10-01 00:12:14 -0300
commitf4412d6dfdb1e2bd8549a7edba785ee242fd38f2 (patch)
treea48a67c60ce0ee55be8db3691f08f36c1986d7f0 /src/ninep/events.zig
parent542dd489149e33915d19878a277e23b3a9c0b070 (diff)
downloadpardes-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/events.zig')
-rw-r--r--src/ninep/events.zig90
1 files changed, 79 insertions, 11 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);
}