summaryrefslogtreecommitdiff
path: root/src/9p_io.zig
diff options
context:
space:
mode:
authorGabriel Schneider <[email protected]>2026-09-28 13:51:15 -0300
committerGabriel Schneider <[email protected]>2026-10-01 00:12:15 -0300
commit1ad90ebaa016d35bd8659f345b6d8a83a9cd6bc9 (patch)
treedc24b567ce71e11d8404fc352f4087a843b4e314 /src/9p_io.zig
parent79b7505cba89a97153e2cfa2a3fef69bb4792faf (diff)
downloadpardes-1ad90ebaa016d35bd8659f345b6d8a83a9cd6bc9.tar.gz
pardes-1ad90ebaa016d35bd8659f345b6d8a83a9cd6bc9.zip
A Restore tells old connections from new by when they were accepted
The Restore marked the connections to cut in cloud9's Conn.user, which cloud9 also writes, clearing it for a connection accepted into that slot on another task: a client dialling in as the marks were made could be marked old and cut. Each connection is now stamped from a counter as it opens (cloud9's opened hook), and a Restore cuts those stamped before its own count; nothing is written from two tasks. Co-Authored-By: Claude Opus 5.5 <[email protected]>
Diffstat (limited to 'src/9p_io.zig')
-rw-r--r--src/9p_io.zig35
1 files changed, 25 insertions, 10 deletions
diff --git a/src/9p_io.zig b/src/9p_io.zig
index 393b6210..2ca11a87 100644
--- a/src/9p_io.zig
+++ b/src/9p_io.zig
@@ -176,6 +176,14 @@ pub const Listener = struct {
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),
+ /// 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
+ /// connection's task as it opens, read with the turn.
+ accepts: std.atomic.Value(u64) = .init(0),
+ accepted: [max_conns]std.atomic.Value(u64) = @splat(.init(0)),
+ /// Connections accepted before this count are being cut by a Restore.
+ cut_before: u64 = 0,
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,
@@ -209,6 +217,16 @@ pub const Listener = struct {
/// 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 onOpened(ctx: ?*anyopaque, conn: *Runner.Conn) void {
+ const l = of(ctx);
+ l.accepted[conn.index].store(l.accepts.fetchAdd(1, .acq_rel) + 1, .release);
+ }
+
+ /// A live connection accepted before the Restore under way.
+ fn cutting(l: *Listener, conn: *Runner.Conn) bool {
+ return conn.live() and l.accepted[conn.index].load(.acquire) <= l.cut_before and l.cut_before != 0;
+ }
+
fn onServe(ctx: ?*anyopaque, conn: *Runner.Conn, req: pardes.ctlfs.Req) void {
const l = of(ctx);
// Once `stop` has begun the editor is tearing down and may never
@@ -221,7 +239,7 @@ pub const Listener = struct {
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) {
+ if (l.cutting(conn)) {
const refused = pardes.ctlfs.Reply.fail(req.tag, pardes.ctlfs.E.IO);
return conn.reply(&refused, "");
}
@@ -516,18 +534,14 @@ pub const Listener = struct {
// 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;
- };
+ l.cut_before = l.accepts.load(.acquire);
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;
+ if (!l.cutting(conn)) continue;
conn.lock();
const n = conn.engine.output().len;
conn.unlock();
@@ -536,16 +550,17 @@ pub const Listener = struct {
if (!pending) break;
Client.nap(1);
}
- for (&l.runner.conns) |*conn| if (conn.user != null and conn.live()) conn.close();
+ for (&l.runner.conns) |*conn| if (l.cutting(conn)) conn.close();
const deadline = nowMs() +| 2000;
while (nowMs() < deadline) {
const left = for (&l.runner.conns) |*conn| {
- if (conn.user != null and conn.live()) break true;
+ if (l.cutting(conn)) break true;
} else false;
if (!left) break;
Client.nap(1);
}
pardes.turn.wake();
+ l.cut_before = 0;
}
l.collectOs();
for (&l.conns) |*conn| conn.accepted_ms = 0;
@@ -801,7 +816,7 @@ pub fn listen(io: std.Io, gpa: std.mem.Allocator, core: *pardes.Pardes, named: [
l.runner.init(.{
.io = io,
.root = pardes.ctlfs.root,
- .handler = .{ .ctx = l, .serve = Listener.onServe },
+ .handler = .{ .ctx = l, .serve = Listener.onServe, .opened = Listener.onOpened },
.greet_timeout_ms = Listener.greet_deadline_ms,
});
const entry_name = if (named.len != 0) named else fallback;