summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/9p_io.zig63
-rw-r--r--src/detached/client.zig3
-rw-r--r--src/detached/server.zig8
-rw-r--r--src/dump.zig8
-rw-r--r--src/gui/gui.zig4
-rw-r--r--src/macos.zig2
-rw-r--r--src/ninep/events.zig23
-rw-r--r--src/pardes.zig7
-rw-r--r--src/tty/tty.zig3
9 files changed, 97 insertions, 24 deletions
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 <path>` when a Dump is written, `restore <path>` 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(