diff options
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); |
