diff options
Diffstat (limited to 'src/9p_io.zig')
| -rw-r--r-- | src/9p_io.zig | 104 |
1 files changed, 98 insertions, 6 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig index 0b3db274..18c51ece 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -180,6 +180,10 @@ pub const Listener = struct { restore_writer: std.atomic.Value(bool) = .init(false), /// Clients turned away with every slot taken, since the log last said so. refused: std.atomic.Value(u32) = .init(0), + /// Writes being served right now, a retry of a parked one included: a + /// held Edit write a retry has taken out of its park is not flushed + /// (watchEdit). + writes_out: std.atomic.Value(u32) = .init(0), /// Connections accepted so far, and the count each slot's connection /// was accepted at: a Restore cuts those accepted before it (`reset`), /// and a new client in a freed slot has a later stamp. Written on the @@ -241,6 +245,10 @@ pub const Listener = struct { fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void { const l = of(ctx); + if (req.op == .write) _ = l.writes_out.fetchAdd(1, .acq_rel); + defer if (req.op == .write) { + _ = l.writes_out.fetchSub(1, .acq_rel); + }; // Once `stop` has begun the editor is tearing down and may never // rest again: answer without waiting on it. if (l.stopping.load(.acquire)) { @@ -268,7 +276,17 @@ pub const Listener = struct { const changes = pardes.ctlfs.changesPane(core, req); const fills = pardes.ctlfs.handsOff(core, req); const reply = core.serveFs(req); + const edit_started = core.fs.edit_started; + core.fs.edit_started = false; + // The held Edit write asked again by a retry, its Edit still + // running: parked again, as it was. + if (reply.status == .again and req.op == .write) return conn.reply(&reply, ""); if (req.op == .release and quiet) l.collectOs(); + // A close's held last line runs now when no step is out, so its err + // is logged before the close is answered, as before; with a step + // out (perhaps through a mount this close came by) it waits for + // that step, and the close is answered at once all the same. + if (req.op == .release and quiet) pardes.ctlfs.runClosedLines(core); // A read with nothing yet stays parked in the engine, and the core // keeps its ticket to answer it when what it waits on has something. if (reply.status == .again and req.op == .read) { @@ -308,14 +326,38 @@ pub const Listener = struct { // mount: 9ns answers nothing else on it while a clunk is out, and // the editor's step may be out opening a file through that mount. if (fills) return conn.reply(&reply, core.fsPayload(reply)); - // A refusal answers at once: its text may be in the core's one - // buffer for it, which a request run while this one waited would - // write over. + // A refusal that left saves going (Putall's one refused among the + // rest) is answered once they have landed, its words copied out of + // the core's one buffer for them, which a request run meanwhile + // could write over. + if (reply.status == .err and core.effects_len != 0 and !core.quit) { + var kept: [320]u8 = undefined; + const words = kept[0..@min(reply.ename.len, kept.len)]; + @memcpy(words, reply.ename[0..words.len]); + var waited = reply; + waited.ename = words; + pardes.turn.awaitSettled(epoch); + return conn.reply(&waited, ""); + } + // Any other refusal answers at once, for the same buffer's sake. if (reply.status == .err) return conn.reply(&reply, ""); // A write that quits the editor (Kill) is answered now: the editor // is on its way out and will settle nothing this could wait for, // and `deinit` lets the answer out before it cuts the connections. if (core.quit) return conn.reply(&reply, ""); + // An Edit whose `<`, `|` or `>` commands run off the loop: the write + // is held, as a read with nothing yet is, and answered when the Edit + // is done (answerHeld), so the connection serves on meanwhile -- a + // filter that reads this session through a mount among what it + // serves. A flush of it, or a hang-up, stops the Edit (watchEdit). + if (edit_started and core.pipe.edit_run != null) { + const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again }; + const ticket = conn.hold(&later) orelse return; + core.fs.edit_hold = .{ .asker = conn, .slot = ticket.slot, .seq = ticket.seq, .tag = req.tag, .node = req.node, .handle = req.handle, .written = reply.written }; + const watcher = std.Thread.spawn(.{}, watchEdit, .{ l, conn, ticket }) catch return; + watcher.detach(); + return; + } const restoring = core.restore_req != null; // Set with the turn still held, so `reset` sees it and waits. if (restoring) l.restore_writer.store(true, .release); @@ -395,6 +437,46 @@ pub const Listener = struct { const conn: *Runner.Conn = @ptrCast(@alignCast(held.asker)); if (!conn.answerWith(held.ticket, Held{ .core = core, .req = held.req }, Held.make)) o.held = null; } + // An Edit's held write, its Edit done. Not waiting: a retry has it + // out, and its asking answers it (Pardes.serveFs); or it was + // flushed, and watchEdit lets the record go. + if (core.fs.edit_hold) |h| if (h.done) { + const conn: *Runner.Conn = @ptrCast(@alignCast(h.asker)); + _ = conn.answerWith(.{ .slot = h.slot, .seq = h.seq }, core, struct { + fn make(c: *pardes.Pardes) ?Runner.Conn.Answer { + return .{ .reply = pardes.filesystem.EditHold.answer(&c.fs) }; + } + }.make); + }; + } + + /// While an Edit's write is held: when nobody waits for it any more -- + /// flushed (an interrupted writer), its fid clunked, its connection + /// hung up -- the Edit is stopped, its commands killed, nothing + /// changed (edit_cmd.cancel). The engine tells the backend nothing of a + /// flush, so it is looked for, every 100 ms; a park out for a retry is + /// not waiting either, so only three looks in a row with no write being + /// served count. + fn watchEdit(l: *Listener, conn: *Runner.Conn, ticket: cloud9.fs.Ticket) void { + const io = pardes.turn.io orelse return; + var misses: u8 = 0; + while (!l.stopping.load(.acquire)) { + io.sleep(.fromMilliseconds(100), .awake) catch return; + _ = pardes.turn.take(); + defer pardes.turn.give(); + const core = l.core; + const h = core.fs.edit_hold orelse return; + if (h.slot != ticket.slot or h.seq != ticket.seq or h.asker != @as(*anyopaque, conn)) return; + if (conn.waiting(ticket) or l.writes_out.load(.acquire) != 0) { + misses = 0; + continue; + } + misses += 1; + if (misses < 3) continue; + if (h.done) core.fs.edit_hold = null else @import("edit_cmd.zig").cancel(core); + l.kick(); + return; + } } /// A held read asked again, to make its answer while its park waits. @@ -1910,6 +1992,16 @@ pub fn start(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes) ?*Listene /// open to scripts and to mounted shells — while withholding `PARDES_PID`, /// which is the whole of what `--nested` means. pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void { + exportEnv(if (listener) |l| l.path() else null, serial, adopts); +} + +/// exportPaneEnv for a command run off the loop (an Edit's filter), with +/// this process's own socket: what a pane's command would be given. +pub fn exportSessionEnv(serial: u32, adopts: bool) void { + exportEnv(if (own_socket_len > 0) own_socket[0..own_socket_len] else null, serial, adopts); +} + +fn exportEnv(socket: ?[]const u8, serial: u32, adopts: bool) void { var announced = false; if (adopts) announcing: { var buf: [16]u8 = undefined; @@ -1919,13 +2011,13 @@ pub fn exportPaneEnv(listener: ?*const Listener, serial: u32, adopts: bool) void } if (!announced) _ = unsetenv("PARDES_PID"); - if (listener) |l| exporting: { + if (socket) |own| exporting: { var sock: [sun_path_len]u8 = undefined; - const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{l.path()}, 0) catch break :exporting; + const path = std.fmt.bufPrintSentinel(&sock, "{s}", .{own}, 0) catch break :exporting; var buf: [16]u8 = undefined; const id = std.fmt.bufPrintSentinel(&buf, "{d}", .{serial}, 0) catch break :exporting; if (setenv("PARDES_9P", path, 1) != 0) break :exporting; - exportMount(l.path()); + exportMount(own); if (setenv("PARDES_PANE", id, 1) == 0) return; } _ = unsetenv("PARDES_9P"); |
