summaryrefslogtreecommitdiff
path: root/src/ninep
diff options
context:
space:
mode:
Diffstat (limited to 'src/ninep')
-rw-r--r--src/ninep/events.zig90
-rw-r--r--src/ninep/pane.zig8
-rw-r--r--src/ninep/pty.zig33
-rw-r--r--src/ninep/screen.zig2
-rw-r--r--src/ninep/tree.zig33
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);