summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig505
1 files changed, 217 insertions, 288 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig
index 08deae7a..7e09fb27 100644
--- a/src/9p_io.zig
+++ b/src/9p_io.zig
@@ -62,7 +62,7 @@ pub fn ensureSocketDir(dir: [:0]const u8) bool {
const prefix = "pardes-9p-";
-pub const msize: u32 = 8192;
+pub const msize = ninep.msize;
pub const max_conns = 4;
@@ -143,76 +143,17 @@ const Srv = ninep.Server(pardes.ctlfs, ninep.editor);
/// The Unix and TCP listeners run on cloud9's `std.Io` runner: it accepts,
/// reads frames, steps each connection's engine and writes replies on its
-/// own tasks, and hands every backend request to the editor's thread
-/// through `Slot` (see `tick`). QUIC keeps the poll loop below.
+/// own tasks, and every backend request is answered right there, on the
+/// connection's task, by taking the editor's turn (`pardes.turn`): while the
+/// editor waits for input, or while a step of it is out in a syscall --
+/// which is what lets the editor read its own tree through a mount. QUIC
+/// keeps the poll loop below, on the editor's thread.
const Runner = if (supported) cloud9.serve.Runner(pardes.ctlfs, ninep.editor, .{
.msize = msize,
.connections = max_conns,
.listeners = 2,
}) else void;
-/// Answers one backend request on the editor's thread.
-fn step(core: *pardes.Pardes, srv: *Srv, req: pardes.ctlfs.Req) void {
- core.update(.{ .fs_req = req });
- var answered = false;
- while (core.nextEffect()) |effect| {
- if (effect == .fs_reply) {
- const reply = effect.fs_reply;
- if (reply.tag == req.tag) answered = true;
- srv.reply(&reply, core.fsPayload(reply));
- } else core.perform(effect);
- }
- if (!answered) {
- const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO);
- srv.reply(&reply, "");
- }
-}
-
-/// Pays the backend a request whose connection is gone; replies go nowhere.
-fn pay(core: *pardes.Pardes, req: pardes.ctlfs.Req) void {
- core.update(.{ .fs_req = req });
- while (core.nextEffect()) |effect| if (effect != .fs_reply) core.perform(effect);
-}
-
-/// One runner slot's mailbox: the requests its connection task handed over,
-/// answered by the editor's thread in `tick`. A hangup's releases (at most
-/// one per fid) and one outstanding request are the most it ever holds,
-/// because a closing connection waits for its entries to be paid before
-/// the slot is reused.
-const Slot = struct {
- const capacity = ninep.editor.fid_capacity + 2;
- const Entry = struct { gen: u32, req: pardes.ctlfs.Req };
-
- mutex: std.Io.Mutex = .init,
- queue: [capacity]Entry = undefined,
- head: usize = 0,
- len: usize = 0,
- /// Bumped when a connection takes the slot; entries carry the value.
- gen: u32 = 0,
- closing: bool = false,
- drained: std.Io.Event = .unset,
- /// Requests refused for want of room; never expected.
- refused: u32 = 0,
-
- fn push(s: *Slot, req: pardes.ctlfs.Req) bool {
- if (s.len == s.queue.len) {
- s.refused += 1;
- return false;
- }
- s.queue[(s.head + s.len) % capacity] = .{ .gen = s.gen, .req = req };
- s.len += 1;
- return true;
- }
-
- fn pop(s: *Slot) ?Entry {
- if (s.len == 0) return null;
- const e = s.queue[s.head];
- s.head = (s.head + 1) % capacity;
- s.len -= 1;
- return e;
- }
-};
-
const quic_slots = if (quic_enabled) max_conns else 0;
/// A QUIC connection on the poll path.
@@ -229,12 +170,10 @@ const accept_pause_ms: i64 = 100;
pub const Listener = struct {
io: std.Io,
+ /// The core that answers; `reset` puts a replacement in.
+ core: *pardes.Pardes,
runner: Runner = undefined,
- slots: [max_conns]Slot = @splat(.{}),
stopping: std.atomic.Value(bool) = .init(false),
- /// While `reset` runs: the generation each slot is paying off; anything
- /// newer waits for the replacement core.
- resetting: ?[max_conns]u32 = null,
tcp_address: ?std.Io.net.IpAddress = null,
quic: if (quic_enabled) ?quic.Listener else void = if (quic_enabled) null else {},
quic_address: ?std.Io.net.IpAddress = null,
@@ -264,47 +203,59 @@ pub const Listener = struct {
return @ptrCast(@alignCast(ctx.?));
}
- /// The runner's handler: on the connection's task.
- fn onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void {
- const l = of(ctx);
- const s = &l.slots[conn.index];
- s.mutex.lockUncancelable(l.io);
- s.gen +%= 1;
- s.mutex.unlock(l.io);
- }
-
+ /// The runner's handler, on the connection's task: takes the turn and
+ /// answers. A request that would change a pane while the editor is out
+ /// in a syscall mid-step is parked in the engine instead, and retried
+ /// when the turn is next given up between steps (`wakeParked`).
fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void {
const l = of(ctx);
- const s = &l.slots[conn.index];
- s.mutex.lockUncancelable(l.io);
- const queued = s.push(req);
- s.mutex.unlock(l.io);
- if (!queued) {
- const reply = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO);
- conn.reply(&reply, "");
- return;
+ // Once `stop` has begun the editor is tearing down and may never
+ // rest again: answer without waiting on it.
+ if (l.stopping.load(.acquire)) {
+ const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO);
+ return conn.reply(&refused, "");
+ }
+ const quiet = pardes.turn.take();
+ defer pardes.turn.give();
+ if (!quiet and pardes.ctlfs.needsQuiet(req)) {
+ pardes.turn.parked = true;
+ const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again };
+ return conn.reply(&later, "");
}
+ const core = l.core;
+ const epoch = pardes.turn.epoch;
+ const restores = pardes.turn.restores;
+ const reply = core.serveFs(req);
+ if (req.op == .release and quiet) l.collectOs();
+ if (!pardes.ctlfs.changesPane(req)) return conn.reply(&reply, core.fsPayload(reply));
+ // The editor draws the change and performs what it asked for; when
+ // it asked for something -- a save, a shell, a watch -- the answer
+ // waits until that is done, so `echo Save > exec` returns with the
+ // file written. A change carries no payload, so the reply is still
+ // whole after the wait.
l.kick();
+ const restoring = core.restore_req != null;
+ if (core.effects_len != 0) pardes.turn.awaitSettled(epoch);
+ // Once a wait returns `core` may be gone: a Restore meanwhile put a
+ // replacement in (`reset`) and is hanging this connection up, which
+ // is a Restore's answer. Only one that failed leaves this connection
+ // here to answer.
+ if (l.core != core) return;
+ if (restoring) pardes.turn.awaitRestored(restores);
+ if (l.core != core) return;
+ conn.reply(&reply, "");
}
- /// Holds the slot until the editor's thread has paid what the
- /// connection still owes, so a successor never inherits its entries.
- fn onClosed(ctx: ?*anyopaque, conn: *Runner.Conn) void {
+ /// The turn was given up quiet with a request parked: every connection
+ /// retries what it parked.
+ fn wakeParked(ctx: ?*anyopaque) void {
const l = of(ctx);
- const s = &l.slots[conn.index];
if (l.stopping.load(.acquire)) return;
- s.mutex.lockUncancelable(l.io);
- if (s.len == 0) {
- s.mutex.unlock(l.io);
- return;
- }
- s.closing = true;
- s.drained.reset();
- s.mutex.unlock(l.io);
- l.kick();
- s.drained.wait(l.io) catch {};
+ if (comptime supported) l.runner.wakeAll();
}
+ /// A request changed the core: the editor wakes to draw it and perform
+ /// what it asked for.
fn kick(l: *Listener) void {
if (l.stopping.load(.acquire)) return;
if (l.wake) |f| f(l.wake_ctx);
@@ -411,68 +362,46 @@ pub const Listener = struct {
c.draining = true;
}
- pub const Drained = struct { count: usize = 0, pending: bool = false };
+ /// What one QUIC poll pass did: `pending` says there is more to do at
+ /// once.
+ pub const Polled = struct { count: usize = 0, pending: bool = false };
- /// Answers what every connection has asked since the last call, retries
- /// their parked reads against the editor's current state and has the
- /// runner send the replies. Call it after each editor update, on the
- /// editor's thread; `pending` says QUIC has more to do right away.
- pub fn drain(l: *Listener, core: *pardes.Pardes) Drained {
- core.fs.socket_path = l.path();
- core.fs.tcp_address = l.tcp_address;
- core.fs.quic_address = l.quic_address;
- l.expire();
- var result: Drained = .{};
- if (comptime supported) {
- for (&l.slots, 0..) |*s, i| {
- const conn = &l.runner.conns[i];
- s.mutex.lockUncancelable(l.io);
- const only: ?u32 = if (l.resetting) |gens| gens[i] else null;
- var count: usize = 0;
- if (s.len != 0 or conn.live()) {
- conn.lock();
- while (s.len != 0) {
- if (only) |gen| if (s.queue[s.head].gen != gen) break;
- const e = s.pop().?;
- if (e.gen == s.gen) step(core, &conn.engine, e.req) else if (e.req.op == .release) pay(core, e.req);
- count += 1;
- }
- if (conn.live() and (only == null or only.? == s.gen)) {
- while (conn.engine.retry()) |req| {
- step(core, &conn.engine, req);
- count += 1;
- }
- while (conn.engine.next()) |req| {
- step(core, &conn.engine, req);
- count += 1;
- }
- }
- const owed = conn.engine.output().len != 0;
- conn.unlock();
- if (owed) conn.flush();
- }
- if (s.closing and s.len == 0) {
- s.closing = false;
- s.drained.set(l.io);
- }
- s.mutex.unlock(l.io);
- result.count += count;
- if (count != 0) l.collectOs(core);
- }
- }
+ /// Answers one QUIC request, on the editor's thread.
+ fn step(l: *Listener, srv: *Srv, req: pardes.ctlfs.Req) void {
+ const reply = l.core.serveFs(req);
+ srv.reply(&reply, l.core.fsPayload(reply));
+ if (req.op == .release) l.collectOs();
+ }
+
+ /// QUIC's events, accepts, reads, requests and replies, on the editor's
+ /// thread after each of its steps; `pending` says there is more to do
+ /// right away. Unix and TCP need none of this.
+ pub fn tick(l: *Listener) Polled {
+ var result: Polled = .{};
+ // Between the editor's steps, which is when the table may be swept:
+ // the editor's own directory listings leave nodes in it too.
+ l.collectOs();
if (comptime quic_enabled) {
+ if (l.quic) |*listener| listener.events() catch |err| {
+ log.warn("QUIC listener stopped: {s}", .{@errorName(err)});
+ for (&l.conns, 0..) |*conn, i| if (conn.quic != null) l.drop(@intCast(i));
+ listener.deinit();
+ l.quic = null;
+ l.quic_address = null;
+ };
+ l.accept();
+ for (0..quic_slots) |i| if (l.live(@intCast(i))) l.fill(@intCast(i));
+ l.expire();
for (&l.conns, 0..) |*conn, i| {
if (!l.live(@intCast(i)) and !conn.draining) continue;
var count: usize = 0;
while (conn.srv.retry()) |req| {
- step(core, &conn.srv, req);
- l.collectOs(core);
+ l.step(&conn.srv, req);
count += 1;
}
while (count < 64) {
const req = conn.srv.next() orelse break;
- step(core, &conn.srv, req);
- l.collectOs(core);
+ l.step(&conn.srv, req);
count += 1;
}
result.count += count;
@@ -482,64 +411,44 @@ pub const Listener = struct {
} else l.flush(@intCast(i));
if (conn.quic) |*connection| result.pending = result.pending or connection.pending();
}
+ l.arm();
}
return result;
}
- /// `drain` plus the QUIC listener's events, accepts and reads.
- pub fn tick(l: *Listener, core: *pardes.Pardes) Drained {
- if (comptime quic_enabled) {
- if (l.quic) |*listener| listener.events() catch |err| {
- log.warn("QUIC listener stopped: {s}", .{@errorName(err)});
- for (&l.conns, 0..) |*conn, i| if (conn.quic != null) l.drop(@intCast(i));
- listener.deinit();
- l.quic = null;
- l.quic_address = null;
- };
- l.accept();
- for (0..quic_slots) |i| if (l.live(@intCast(i))) l.fill(@intCast(i));
- }
- const result = l.drain(core);
- l.arm();
- return result;
- }
-
- /// Hangs every connection up and pays `core` what they owed, so that a
- /// replacement core starts with no client holding anything. A client
- /// that connects meanwhile is not answered until the caller has put
- /// the replacement in and calls `tick` again.
- pub fn reset(l: *Listener, core: *pardes.Pardes) void {
+ /// Hangs every connection up, so that `replacement` starts with no
+ /// client holding anything. From here requests are answered from the
+ /// replacement -- a client that connects meanwhile is a client of the
+ /// replacement -- and every task waiting on the old core is let go: it
+ /// finds the core changed and answers nothing, its connection being
+ /// hung up. The connections pay their releases on their own tasks, so
+ /// the editor rests while they do.
+ pub fn reset(l: *Listener, replacement: *pardes.Pardes) void {
for (0..quic_slots) |i| l.drop(@intCast(i));
+ if (comptime quic_enabled) while (l.tick().pending) {};
+ l.core = replacement;
+ replacement.fs.socket_path = l.path();
+ replacement.fs.tcp_address = l.tcp_address;
+ replacement.fs.quic_address = l.quic_address;
if (comptime supported) {
- var gens: [max_conns]u32 = undefined;
- for (&l.slots, &gens) |*s, *gen| {
- s.mutex.lockUncancelable(l.io);
- gen.* = s.gen;
- s.mutex.unlock(l.io);
- }
- l.resetting = gens;
- defer l.resetting = null;
l.runner.closeAll();
+ pardes.turn.settle();
+ pardes.turn.restoreSettled();
+ pardes.turn.rest();
const deadline = nowMs() +| 2000;
- while (true) {
- const drained = l.drain(core);
- var idle = !drained.pending;
- for (&l.slots, &gens, 0..) |*s, gen, i| {
- s.mutex.lockUncancelable(l.io);
- if (s.len != 0 and s.queue[s.head].gen == gen) idle = false;
- if (s.gen == gen and l.runner.conns[i].live()) idle = false;
- s.mutex.unlock(l.io);
- }
- if (idle or nowMs() >= deadline) break;
- Client.nap(1);
- }
- } else while (l.drain(core).pending) {}
- l.collectOs(core);
+ while (l.runner.count() != 0 and nowMs() < deadline) Client.nap(1);
+ pardes.turn.wake();
+ }
+ l.collectOs();
for (&l.conns) |*conn| conn.accepted_ms = 0;
l.arm();
}
- fn collectOs(l: *Listener, core: *pardes.Pardes) void {
+ /// Forgets the host paths no connection names any more. With the turn,
+ /// and only between the editor's steps: a step out in a syscall may be
+ /// reading one of these entries.
+ fn collectOs(l: *Listener) void {
+ const core = l.core;
var i: usize = 0;
while (i < core.fs.os_paths.items.len) {
const entry = core.fs.os_paths.items[i];
@@ -552,17 +461,10 @@ pub const Listener = struct {
}
}
- /// Whether any connection, or a release still to be paid, names `node`.
+ /// Whether any connection names `node`.
fn references(l: *Listener, node: u64) bool {
if (comptime supported) {
- for (&l.slots, 0..) |*s, i| {
- const conn = &l.runner.conns[i];
- s.mutex.lockUncancelable(l.io);
- defer s.mutex.unlock(l.io);
- var k: usize = 0;
- while (k < s.len) : (k += 1) {
- if (s.queue[(s.head + k) % Slot.capacity].req.node == node) return true;
- }
+ for (&l.runner.conns) |*conn| {
if (!conn.live()) continue;
conn.lock();
defer conn.unlock();
@@ -576,16 +478,12 @@ pub const Listener = struct {
return false;
}
- /// Where the runner's tasks report that a request waits for `tick`.
- pub fn setWake(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) void {
+ /// Where the runner's tasks report that a request changed the core,
+ /// and for QUIC a thread that watches its socket and calls `wake` when
+ /// `tick` has something to read.
+ pub fn watch(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void {
l.wake_ctx = ctx;
l.wake = wake;
- }
-
- /// `setWake`, and for QUIC a thread that watches its socket and calls
- /// `wake` when `tick` has something to read.
- pub fn watch(l: *Listener, ctx: ?*anyopaque, wake: *const fn (?*anyopaque) void) !void {
- l.setWake(ctx, wake);
if (comptime !quic_enabled) return;
if (l.quic == null or l.watcher != null) return;
if (libc.pipe(&l.control) != 0) return error.PipeFailed;
@@ -735,7 +633,15 @@ pub const Listener = struct {
if (comptime quic_enabled) {
if (l.quic) |*listener| listener.deinit();
}
- if (comptime supported) l.runner.stop();
+ if (comptime supported) {
+ // A connection task may be waiting for the turn; it answers
+ // nothing more once stopping, but the wait itself must end.
+ pardes.turn.rest();
+ l.runner.stop();
+ pardes.turn.wake();
+ }
+ pardes.turn.wake_parked = null;
+ pardes.turn.stop();
if (l.path_len != 0) {
var z: [sun_path_len:0]u8 = undefined;
@memcpy(z[0..l.path_len], l.path_buf[0..l.path_len]);
@@ -752,7 +658,9 @@ pub fn socketPath(buf: *[sun_path_len]u8, dir: []const u8, name: []const u8) ?[:
return std.fmt.bufPrintSentinel(buf, "{s}/" ++ prefix ++ "{s}.sock", .{ dir, name }, 0) catch null;
}
-pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener {
+/// Serves `core`. From here on the calling thread has the editor's turn
+/// (`pardes.turn`) and must give it up while it waits.
+pub fn listen(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes, named: []const u8, fallback: []const u8, tcp_dial: ?[]const u8, quic_dial: ?[]const u8) ?*Listener {
if (comptime !supported) return null;
var dir_buf: [sun_path_len:0]u8 = undefined;
const dir = socketDir(&dir_buf) orelse {
@@ -761,11 +669,14 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [
};
if (!ensureSocketDir(dir)) return null;
const l = gpa.create(Listener) catch return null;
- l.* = .{ .io = io };
+ l.* = .{ .io = io, .core = core };
+ pardes.turn.start(io);
+ pardes.turn.wake_parked = Listener.wakeParked;
+ pardes.turn.wake_ctx = l;
l.runner.init(.{
.io = io,
.root = pardes.ctlfs.root,
- .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened, .closed = Listener.onClosed },
+ .handler = .{ .ctx = l, .serve = Listener.onServe },
.greet_timeout_ms = Listener.greet_deadline_ms,
});
const entry_name = if (named.len != 0) named else fallback;
@@ -838,6 +749,9 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, named: []const u8, fallback: [
}
}
log.info("serving 9P2000 on {s}", .{p});
+ core.fs.socket_path = l.path();
+ core.fs.tcp_address = l.tcp_address;
+ core.fs.quic_address = l.quic_address;
return l;
}
@@ -1160,33 +1074,6 @@ test "TCP addresses are numeric and normalize mapped IPv4" {
} else try testing.expectError(error.QuicUnavailable, Client.validateDial("quic!127.0.0.1!5640"));
}
-test "same-session TCP mounts compare canonical endpoints and local wildcard destinations" {
- if (comptime !supported) return error.SkipZigTest;
- const Case = struct { bound: []const u8, dial: []const u8, same: bool };
- for ([_]Case{
- .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!::ffff:127.0.0.1!5640", .same = true },
- .{ .bound = "tcp!::ffff:127.0.0.1!5640", .dial = "tcp!127.0.0.1!5640", .same = true },
- .{ .bound = "tcp!::1!5640", .dial = "tcp!0:0:0:0:0:0:0:1!5640", .same = true },
- .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!0.0.0.0!5640", .same = true },
- .{ .bound = "tcp!::1!5640", .dial = "tcp!::!5640", .same = true },
- .{ .bound = "tcp!0.0.0.0!5640", .dial = "tcp!127.0.0.2!5640", .same = true },
- .{ .bound = "tcp!::!5640", .dial = "tcp!::1!5640", .same = true },
- .{ .bound = "tcp!127.0.0.1!5640", .dial = "tcp!127.0.0.1!5641", .same = false },
- .{ .bound = "tcp!0.0.0.0!5640", .dial = "tcp!192.0.2.1!5640", .same = false },
- .{ .bound = "tcp!::!5640", .dial = "tcp!2001:db8::1!5640", .same = false },
- .{ .bound = "tcp!::!5640", .dial = "tcp!127.0.0.1!5640", .same = false },
- }) |c| try testing.expectEqual(c.same, Client.sameSession(c.dial, "", try networkAddress(c.bound, true), null));
- try testing.expect(Client.sameSession("/tmp/pardes-owned.sock", "/tmp/pardes-owned.sock", null, null));
- try testing.expect(Client.sameSession("unix!/tmp/pardes-owned.sock", "/tmp/pardes-owned.sock", null, null));
- try testing.expect(!Client.sameSession("tcp!127.0.0.1!5640", "/tmp/pardes-owned.sock", null, null));
- if (quic_enabled) {
- const endpoint = try networkAddress("quic!127.0.0.1!5640", false);
- try testing.expect(Client.sameSession("quic!::ffff:127.0.0.1!5640", "", null, endpoint));
- try testing.expect(!Client.sameSession("tcp!127.0.0.1!5640", "", null, endpoint));
- try testing.expect(!Client.sameSession("quic!127.0.0.1!5640", "", endpoint, null));
- }
-}
-
extern "c" fn mkdtemp(template: [*:0]u8) ?[*:0]u8;
extern "c" fn rmdir(path: [*:0]const u8) c_int;
@@ -1225,7 +1112,9 @@ test "a listening editor posts itself into the 9P registry and unposts on stop"
var link: [sun_path_len]u8 = undefined;
{
- const l = listen(testing.io, gpa, "unit", "", null, null) orelse return error.ListenFailed;
+ const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 });
+ defer p.deinit();
+ const l = listen(testing.io, gpa, p, "unit", "", null, null) orelse return error.ListenFailed;
defer l.deinit(gpa);
// The socket stays exactly where pardes has always bound it:
// adopting the registry moves nothing, it only advertises.
@@ -1312,7 +1201,7 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r
var bind_buf: [64]u8 = undefined;
const bind = try std.fmt.bufPrint(&bind_buf, "{s}!{s}!0", .{ protocol, host });
const tcp = std.mem.eql(u8, protocol, "tcp");
- const l = listen(testing.io, gpa, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed;
+ const l = listen(testing.io, gpa, p, "roundtrip", "", if (tcp) bind else null, if (tcp) null else bind) orelse return error.ListenFailed;
defer {
l.reset(p);
l.deinit(gpa);
@@ -1331,39 +1220,108 @@ test "Unix TCP and QUIC share one listener through reads writes reconnects and r
var worker: Worker = .{ .dial = dial, .body_path = body, .expected = if (attempt == 0) "initial\n" else replacement, .replacement = replacement };
const thread = try std.Thread.spawn(.{}, Worker.run, .{&worker});
defer thread.join();
+ // The editor rests while the client works: its requests are
+ // answered on the runner's tasks, not by this thread.
const deadline = Client.nowMs() + 3 * Client.budget_ms;
+ pardes.turn.rest();
while (!worker.done.load(.acquire) and Client.nowMs() < deadline) {
- _ = l.tick(p);
+ if (quic_enabled) {
+ pardes.turn.wake();
+ _ = l.tick();
+ pardes.turn.rest();
+ }
Client.nap(1);
}
+ pardes.turn.wake();
try testing.expect(worker.done.load(.acquire));
if (worker.failure) |err| return err;
l.reset(p);
try testing.expectEqual(@as(usize, 0), l.runner.count());
- for (&l.slots) |*slot| try testing.expectEqual(@as(usize, 0), slot.len);
for (&l.conns, 0..) |conn, i| try testing.expect(!l.live(@intCast(i)) and !conn.draining);
for (p.fs.snapshots) |snapshot| try testing.expect(snapshot.node == 0);
}
};
}
+test "a change waits while the editor is out mid-step, a read does not, and the change lands when the editor rests" {
+ if (comptime !supported) return error.SkipZigTest;
+ const gpa = testing.allocator;
+ var directory: [64:0]u8 = undefined;
+ _ = try std.fmt.bufPrintSentinel(&directory, "/tmp/pardes-park-XXXXXX", .{}, 0);
+ if (mkdtemp(&directory) == null) return error.TempDirectoryFailed;
+ defer _ = rmdir(&directory);
+ const old_runtime = if (libc.getenv("XDG_RUNTIME_DIR")) |v| try gpa.dupeZ(u8, std.mem.span(v)) else null;
+ defer {
+ if (old_runtime) |v| {
+ _ = setenv("XDG_RUNTIME_DIR", v, 1);
+ gpa.free(v);
+ } else _ = unsetenv("XDG_RUNTIME_DIR");
+ }
+ try testing.expectEqual(@as(c_int, 0), setenv("XDG_RUNTIME_DIR", &directory, 1));
+
+ const p = try pardes.Pardes.init(gpa, .{ .tty_only = true, .cols = 40, .rows = 12 });
+ defer p.deinit();
+ const pane = try p.setTestFile("before\n");
+ while (p.nextEffect()) |_| {}
+ // This thread is the editor's: it has the turn from here on.
+ const l = listen(testing.io, gpa, p, "park", "", null, null) orelse return error.ListenFailed;
+ defer {
+ l.reset(p);
+ l.deinit(gpa);
+ }
+ var body_buf: [64]u8 = undefined;
+ const body = try std.fmt.bufPrint(&body_buf, "/pane/{d}/body", .{pane.serial});
+
+ const Writer = struct {
+ dial: []const u8,
+ path: []const u8,
+ done: std.atomic.Value(bool) = .init(false),
+ failure: ?anyerror = null,
+ fn run(w: *@This()) void {
+ defer w.done.store(true, .release);
+ Client.write(testing.allocator, w.dial, w.path, "after\n") catch |err| {
+ w.failure = err;
+ };
+ }
+ };
+ var writer: Writer = .{ .dial = l.path(), .path = body };
+
+ // The editor goes out mid-step -- a Look resolving a path, say -- and a
+ // write arrives meanwhile. It must not land: the step still holds
+ // pointers into the pane. A read is answered all the same.
+ pardes.turn.yield();
+ const thread = try std.Thread.spawn(.{}, Writer.run, .{&writer});
+ defer thread.join();
+ Client.nap(150);
+ try testing.expect(!writer.done.load(.acquire));
+ const seen = try Client.readLimit(gpa, l.path(), body, body, 64);
+ defer gpa.free(seen);
+ try testing.expectEqualStrings("before\n", seen);
+ try testing.expect(!writer.done.load(.acquire));
+ pardes.turn.back();
+ try testing.expectEqualStrings("before\n", pane.file.?.content);
+
+ // Resting is what lets it in: the parked write is retried, served, and
+ // the client hears back.
+ pardes.turn.rest();
+ const deadline = Client.nowMs() + Client.budget_ms;
+ while (!writer.done.load(.acquire) and Client.nowMs() < deadline) Client.nap(1);
+ pardes.turn.wake();
+ try testing.expect(writer.done.load(.acquire));
+ if (writer.failure) |err| return err;
+ try testing.expectEqualStrings("after\n", pane.file.?.content);
+}
+
extern "c" fn setenv(name: [*:0]const u8, value: [*:0]const u8, overwrite: c_int) c_int;
extern "c" fn unsetenv(name: [*:0]const u8) c_int;
pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listener {
var name: [16]u8 = undefined;
const fallback = std.fmt.bufPrint(&name, "{d}", .{@as(u32, @intCast(libc.getpid()))}) catch unreachable;
- const listener = listen(io, gpa, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse {
+ return listen(io, gpa, core, core.opts.ninep_name, fallback, core.opts.ninep_tcp, core.opts.ninep_quic) orelse {
core.reportError(0, "9p listener", error.ListenFailed);
return null;
};
- core.fs.socket_path = listener.path();
- // So the editor can recognise its own tree by name and refuse to walk into
- // it through the mount, which would deadlock the loop that serves it.
- pardes.filesystem.noteOwnSocket(listener.path());
- core.fs.tcp_address = listener.tcp_address;
- core.fs.quic_address = listener.quic_address;
- return listener;
}
/// What a pane shell is told about the editor above it, set into this
@@ -1423,7 +1381,7 @@ test "9P shell environment states being inside pardes apart from how to reach it
const listener = try testing.allocator.create(Listener);
defer testing.allocator.destroy(listener);
- listener.* = .{ .io = testing.io };
+ listener.* = .{ .io = testing.io, .core = undefined }; // only its path is read
const path = "/tmp/pardes-example.sock";
@memcpy(listener.path_buf[0..path.len], path);
listener.path_len = path.len;
@@ -1549,35 +1507,6 @@ pub const Client = struct {
_ = try resolve(&buf, dial);
}
- pub fn sameSession(dial: []const u8, socket_path: []const u8, tcp_address: ?std.Io.net.IpAddress, quic_address: ?std.Io.net.IpAddress) bool {
- if (comptime !supported) return false;
- var buf: [sun_path_len]u8 = undefined;
- const address = resolve(&buf, dial) catch return false;
- switch (address) {
- .unix => |path| return socket_path.len != 0 and std.mem.eql(u8, path, socket_path),
- .tcp, .quic => |destination| {
- var ip = destination;
- switch (ip) {
- .ip4 => |v4| if (std.mem.allEqual(u8, &v4.bytes, 0)) {
- ip = .{ .ip4 = .loopback(v4.port) };
- },
- .ip6 => |v6| if (std.mem.allEqual(u8, &v6.bytes, 0)) {
- ip = .{ .ip6 = .loopback(v6.port) };
- },
- }
- const bound = canonicalIp((if (address == .tcp) tcp_address else quic_address) orelse return false);
- if (ip.getPort() != bound.getPort() or @as(std.Io.net.IpAddress.Family, ip) != @as(std.Io.net.IpAddress.Family, bound)) return false;
- if (ip.eql(&bound)) return true;
- const wildcard = switch (bound) {
- .ip4 => |v4| std.mem.allEqual(u8, &v4.bytes, 0),
- .ip6 => |v6| std.mem.allEqual(u8, &v6.bytes, 0),
- };
- if (wildcard) return localIp(ip);
- return false;
- },
- }
- }
-
const Session = struct {
fd: c_int,
quic: if (quic_enabled) ?quic.Connection else void = if (quic_enabled) null else {},