From 6fc9f416938b59b665b1fddc052402770a21cc3f Mon Sep 17 00:00:00 2001 From: Gabriel Schneider Date: Mon, 28 Sep 2026 12:44:57 -0300 Subject: A Restore answers its writer before hanging up, and the log records dumps and restores A client that wrote Restore saw its connection cut with no answer, and could not tell a Restore from a crash. The listener now lets the writer's answer out before the cut, and cuts only the old editor's connections, refusing their requests meanwhile; a client that dials during it is the new editor's and stays. Dump logs 'dump ' and the restored editor's log 'restore '. Keeping connections across a Restore was weighed and left: the fids name the old editor's panes and opens, so it would mean carrying serials and open records into the new one, where acme's Load only adds windows. tty's Restore also closed its shells' ptys without reaping them; it retires them now. Co-Authored-By: Claude Opus 5.5 --- src/9p_io.zig | 63 +++++++++++++++++++++++++++++++++++++++++-------- src/detached/client.zig | 3 ++- src/detached/server.zig | 8 +++---- src/dump.zig | 8 +++++-- src/gui/gui.zig | 4 ++-- src/macos.zig | 2 +- src/ninep/events.zig | 23 ++++++++++++++++++ src/pardes.zig | 7 +++--- src/tty/tty.zig | 3 ++- 9 files changed, 97 insertions(+), 24 deletions(-) (limited to 'src') diff --git a/src/9p_io.zig b/src/9p_io.zig index d6eef9d4..393b6210 100644 --- a/src/9p_io.zig +++ b/src/9p_io.zig @@ -174,6 +174,8 @@ pub const Listener = struct { core: *pardes.Pardes, runner: Runner = undefined, stopping: std.atomic.Value(bool) = .init(false), + /// A connection wrote the Restore and has yet to answer it. + restore_writer: std.atomic.Value(bool) = .init(false), 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, @@ -217,6 +219,12 @@ pub const Listener = struct { } const quiet = pardes.turn.take(); defer pardes.turn.give(); + // A connection a Restore is cutting (`reset`) is the old editor's + // client: its fids name the old editor's panes and opens. + if (conn.user != null) { + const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO); + return conn.reply(&refused, ""); + } if (!quiet and pardes.ctlfs.needsQuiet(req)) { pardes.turn.parked = true; const later: pardes.ctlfs.Reply = .{ .tag = req.tag, .status = .again }; @@ -267,13 +275,20 @@ pub const Listener = struct { // and `deinit` lets the answer out before it cuts the connections. if (core.quit) return conn.reply(&reply, ""); 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); + defer if (restoring) l.restore_writer.store(false, .release); if (core.effects_len != 0) pardes.turn.awaitSettled(epoch); + // The Restore's own write is answered once it is done, and before + // `reset` hangs every connection up, this one too, so the writer + // hears that it happened rather than a cut; a failed one leaves the + // connection and says why on the message row and in the log. + if (restoring) { + if (l.core == core) pardes.turn.awaitRestored(restores); + return conn.reply(&reply, ""); + } // 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); + // replacement in and is hanging this connection up. if (l.core != core) return; conn.reply(&reply, ""); } @@ -484,9 +499,9 @@ pub const Listener = struct { /// 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 + /// replacement -- and every task waiting on the old core is let go: the + /// Restore's own writer answers, the rest find the core changed and + /// answer nothing, their connections 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)); @@ -496,12 +511,40 @@ pub const Listener = struct { replacement.fs.tcp_address = l.tcp_address; replacement.fs.quic_address = l.quic_address; if (comptime supported) { - l.runner.closeAll(); + // The Restore's writer answers as it wakes (`onServe`), and its + // answer gets 200 ms to leave before the cut; every other request + // from the old editor's clients meanwhile is refused, since the + // fids it names are the old editor's. A client that connects now + // is the replacement's, and stays. + // `user` marks the connections to cut: cloud9 leaves it to us and + // clears it for a new connection in the slot. + for (&l.runner.conns) |*conn| if (conn.live()) { + conn.user = l; + }; pardes.turn.settle(); pardes.turn.restoreSettled(); pardes.turn.rest(); + const flush_by = nowMs() +| 200; + while (nowMs() < flush_by) { + const pending = l.restore_writer.load(.acquire) or for (&l.runner.conns) |*conn| { + if (conn.user == null or !conn.live()) continue; + conn.lock(); + const n = conn.engine.output().len; + conn.unlock(); + if (n != 0) break true; + } else false; + if (!pending) break; + Client.nap(1); + } + for (&l.runner.conns) |*conn| if (conn.user != null and conn.live()) conn.close(); const deadline = nowMs() +| 2000; - while (l.runner.count() != 0 and nowMs() < deadline) Client.nap(1); + while (nowMs() < deadline) { + const left = for (&l.runner.conns) |*conn| { + if (conn.user != null and conn.live()) break true; + } else false; + if (!left) break; + Client.nap(1); + } pardes.turn.wake(); } l.collectOs(); diff --git a/src/detached/client.zig b/src/detached/client.zig index 4468c961..ba533751 100644 --- a/src/detached/client.zig +++ b/src/detached/client.zig @@ -838,6 +838,7 @@ test "detached Restore keeps attached frontends and follows queued frames with a const fd = h.session.clients[c.slot].fd; try testing.expectError(error.BadDumpMagic, h.session.restore( ".{ .magic = \"not-a-pardes-dump\", .theme = \"dark\", .screen = .{ .cols = 60, .rows = 16 } }", + "test", )); try testing.expectEqual(before, h.session.core); try testing.expectEqual(fd, h.session.clients[c.slot].fd); @@ -848,7 +849,7 @@ test "detached Restore keeps attached frontends and follows queued frames with a defer testing.allocator.free(saved); while (h.core.nextEffect()) |_| {} _ = try h.core.setTestFile("changed after dump\n"); - try h.session.restore(saved); + try h.session.restore(saved, "test"); h.core = h.session.core; try testing.expect(h.core != before); try testing.expectEqual(fd, h.session.clients[c.slot].fd); diff --git a/src/detached/server.zig b/src/detached/server.zig index 5c0e1487..b088d833 100644 --- a/src/detached/server.zig +++ b/src/detached/server.zig @@ -233,8 +233,8 @@ pub const Session = struct { return batch.len != 0 or batch.status != null; } - pub fn restore(s: *Session, bytes: []const u8) !void { - const replacement = try dump.restore(s.core, bytes); + pub fn restore(s: *Session, bytes: []const u8, from: []const u8) !void { + const replacement = try dump.restore(s.core, bytes, from); s.cancelWorkers(); for (0..s.ptys.len) |pane| s.closePty(@intCast(pane)); s.harvest(); @@ -1183,7 +1183,7 @@ pub fn run(init: std.process.Init, opts: pardes.Options, name: []const u8) !void break :restore; }; defer gpa.free(bytes); - session.restore(bytes) catch |err| session.core.reportError(session.core.active, "Restore", err); + session.restore(bytes, path) catch |err| session.core.reportError(session.core.active, "Restore", err); } } } @@ -1223,7 +1223,7 @@ test "detached queued results preserve current requests and are discarded before try std.testing.expect(s.pipe_tasks.add(.{ .id = 77, .future = .{ .any_future = null, .result = {} } })); s.mailbox.post(.{ .pipe = .{ .id = 77, .success = true, .outputs = outputs } }); Session.lspStatus(&s, "old status"); - try s.restore(saved); + try s.restore(saved, "test"); try std.testing.expect(s.lsp_task == null); try std.testing.expectEqual(@as(usize, 0), s.pipe_tasks.len); try std.testing.expect(!s.drainCompletions(true)); diff --git a/src/dump.zig b/src/dump.zig index 508e7300..3d0f0f0e 100644 --- a/src/dump.zig +++ b/src/dump.zig @@ -671,11 +671,15 @@ pub fn dumpState(p: *Pardes) !void { p.emit(.write_dump); } -pub fn restore(p: *Pardes, zon_bytes: []const u8) !*Pardes { +/// The replacement core for a Restore of the dump at `from`, whose log says +/// so after its panes' `new`s, for a client that reconnects to read. +pub fn restore(p: *Pardes, zon_bytes: []const u8, from: []const u8) !*Pardes { var opts = p.opts; opts.cols = p.screen_w; opts.rows = p.screen_h; - return initDump(p.gpa, opts, zon_bytes, p); + const replacement = try initDump(p.gpa, opts, zon_bytes, p); + pardes.ctlfs.events.notePath(replacement, "restore", from); + return replacement; } pub fn initFromDump(gpa: std.mem.Allocator, opts: Options, zon_bytes: []const u8) !*Pardes { diff --git a/src/gui/gui.zig b/src/gui/gui.zig index 6978284b..9ada305b 100644 --- a/src/gui/gui.zig +++ b/src/gui/gui.zig @@ -1322,7 +1322,7 @@ test "GUI PTY Restore joins a real reader waiting for queue space" { try std.testing.expect(host_io.writeFd(child.file.handle, "printf 'after-full'; exit\n")); try PtyTests.waitBlocked(&queue); const old_serial = core.panes[0].?.serial; - const replacement = try dump.restore(core, core.dump_out.?); + const replacement = try dump.restore(core, core.dump_out.?, "test"); shell.stopPtys(); core.deinit(); core = replacement; @@ -2561,7 +2561,7 @@ fn localSession( break :blk; }; defer gpa.free(bytes); - const nc = dump.restore(core, bytes) catch |err| { + const nc = dump.restore(core, bytes, rp) catch |err| { core.reportError(core.active, "Restore", err); break :blk; }; diff --git a/src/macos.zig b/src/macos.zig index 62858887..7bfa6c49 100644 --- a/src/macos.zig +++ b/src/macos.zig @@ -1345,7 +1345,7 @@ fn restoreCore(st: *State) bool { return false; }; defer st.gpa.free(bytes); - const replacement = dump.restore(st.core, bytes) catch |err| { + const replacement = dump.restore(st.core, bytes, path) catch |err| { st.core.reportError(st.core.active, "Restore", err); return false; }; diff --git a/src/ninep/events.zig b/src/ninep/events.zig index 69dc271d..1cab5605 100644 --- a/src/ninep/events.zig +++ b/src/ninep/events.zig @@ -7,6 +7,7 @@ const look = @import("../look.zig"); const cloud9 = @import("cloud9"); const tree = @import("tree.zig"); const pane_files = @import("pane.zig"); +const dump = @import("../dump.zig"); const Pardes = pardes.Pardes; const Pane = pardes.Pane; @@ -155,6 +156,13 @@ pub fn noteMessage(p: *Pardes, serial: u32, text: []const u8) void { /// kernel mapped the reply to (`Invalid argument`), and here is the reason. /// The log is the one place for it, as acme's `errors` file takes text and /// answers nothing: a per-pane readable error file would be a second. +/// `dump ` when a Dump is written, `restore ` first in a +/// Restore's replacement. +pub fn notePath(p: *Pardes, what: []const u8, path: []const u8) void { + var buf: [pardes.memory.limits.host_path_cap + 16]u8 = undefined; + pushLog(p, std.fmt.bufPrint(&buf, "{s} {s}\n", .{ what, path }) catch return); +} + pub fn noteError(p: *Pardes, req: Req, reply: Reply) void { const why = if (reply.ename.len > 0) reply.ename else cloud9.fs.errString(reply.errno); var name: [16]u8 = undefined; @@ -870,6 +878,21 @@ test "a refused or failed write is an err record in the log, saying which file a _ = call(p, .{ .tag = 9, .op = .release, .node = log, .handle = g }); } +test "a Dump written and a Restore made are in the log, with their files" { + const p = try withFile(testing.allocator, "one\n"); + defer p.deinit(); + p.setLastDump("/tmp/pardes.dump.zon"); + try dump.dumpState(p); + const restored = try dump.restore(p, p.dump_out.?, "/tmp/pardes.dump.zon"); + defer restored.deinit(); + for ([_]*Pardes{ p, restored }, [_][]const u8{ "dump", "restore" }) |core, what| { + const log = try freezeLog(core); + defer core.gpa.free(log.bytes); + var want: [64]u8 = undefined; + try testing.expect(std.mem.endsWith(u8, log.bytes, try std.fmt.bufPrint(&want, "{s} /tmp/pardes.dump.zon\n", .{what}))); + } +} + test "opens of the log share the open records, and a closed one frees its record" { const gpa = testing.allocator; const p = try withFile(gpa, "x\n"); diff --git a/src/pardes.zig b/src/pardes.zig index 6e78380a..c049fa39 100644 --- a/src/pardes.zig +++ b/src/pardes.zig @@ -1770,7 +1770,7 @@ test "owned cwd restore either copies the directory or preserves the old core" { for (0..256) |failure| { allocator.has_induced_failure = false; allocator.fail_index = allocator.alloc_index + failure; - const replacement = dump.restore(p, p.dump_out.?) catch { + const replacement = dump.restore(p, p.dump_out.?, "test") catch { try std.testing.expect(allocator.has_induced_failure); try std.testing.expectEqual(original, p.panes[0].?); try std.testing.expectEqualStrings("/retained/terminal/directory", original.cwdSlice()); @@ -2722,14 +2722,14 @@ test "editable workspace and column tags are typed into and persist" { try std.testing.expect(std.mem.startsWith(u8, tagline.columnTag(p, 0), "Grep New")); try std.testing.expectEqualStrings("untouched\n", p.panes[0].?.file.?.content); try dump.dumpState(p); - const restored = try dump.restore(p, p.dump_out.?); + const restored = try dump.restore(p, p.dump_out.?, "test"); defer restored.deinit(); try std.testing.expectEqualStrings(p.global_tag.own.?, restored.global_tag.own.?); try std.testing.expectEqualStrings(tagline.columnTag(p, 0), tagline.columnTag(restored, 0)); p.gpa.free(p.col_tags[0].own.?); p.col_tags[0].own = try p.gpa.dupe(u8, ""); try dump.dumpState(p); - const empty = try dump.restore(p, p.dump_out.?); + const empty = try dump.restore(p, p.dump_out.?, "test"); defer empty.deinit(); try std.testing.expectEqualStrings("", tagline.columnTag(empty, 0)); } @@ -4095,6 +4095,7 @@ pub const Pardes = struct { const copy = p.gpa.dupe(u8, path) catch return; if (p.last_dump) |old| p.gpa.free(old); p.last_dump = copy; + ctlfs.events.notePath(p, "dump", path); } /// the shell polls this each frame: a pending Restore's dump path, or null diff --git a/src/tty/tty.zig b/src/tty/tty.zig index 32f9138f..4c1e2db7 100644 --- a/src/tty/tty.zig +++ b/src/tty/tty.zig @@ -803,7 +803,7 @@ fn localSession( break :blk; }; defer gpa.free(bytes); - const nc = dump.restore(core, bytes) catch |err| { + const nc = dump.restore(core, bytes, rp) catch |err| { core.reportError(core.active, "Restore", err); break :blk; }; @@ -816,6 +816,7 @@ fn localSession( for (&sh.ptys) |*slot| if (slot.*) |*pt| { pt.reader.cancel(io) catch {}; _ = libc.close(pt.file.handle); + host_io.retireShell(pt.pid); slot.* = null; }; for (0..pardes.MAX_PANES) |wid| file_watch.watchPane( -- cgit v1.3