From 80eca7859d884ed6f5cf6ab44a92084f7df8818b Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 10 Apr 2026 00:54:08 -0500 Subject: [PATCH] clean up stale Evented-era comments across codebase MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit now that Backend is Io.Threaded, remove references to cross-Io constraints, Evented fiber context requirements, and Uring thread warnings that no longer apply. the historical context is preserved in docs/evented-attempt.md and docs/notes.md. no behavioral changes — comments only. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/api/admin.zig | 4 ++-- src/api/xrpc.zig | 2 +- src/broadcaster.zig | 6 +++--- src/event_log.zig | 10 +++++----- src/host_ops.zig | 4 ++-- src/main.zig | 34 +++++++++++++++------------------- src/resync.zig | 6 ++---- src/slurper.zig | 2 +- src/subscriber.zig | 4 ++-- 9 files changed, 33 insertions(+), 39 deletions(-) diff --git a/src/api/admin.zig b/src/api/admin.zig index bd98d9a..7dce94b 100644 --- a/src/api/admin.zig +++ b/src/api/admin.zig @@ -3,8 +3,8 @@ //! all handlers require Bearer token auth against RELAY_ADMIN_PASSWORD. //! includes host blocking/unblocking, account bans, and backfill control. //! -//! DB-accessing handlers use DbRequest + DbRequestQueue to route queries through -//! pool_io (Threaded) workers. +//! DB-accessing handlers use DbRequest + DbRequestQueue to route queries +//! through pool_io workers. const std = @import("std"); const Io = std.Io; diff --git a/src/api/xrpc.zig b/src/api/xrpc.zig index 18458d4..c6db9ce 100644 --- a/src/api/xrpc.zig +++ b/src/api/xrpc.zig @@ -5,7 +5,7 @@ //! listHosts, getHostStatus, requestCrawl //! //! DB-accessing handlers use DbRequest + DbRequestQueue to route queries through -//! pool_io (Threaded) workers, avoiding the broken Evented pg.Pool. +//! pool_io workers. const std = @import("std"); const Io = std.Io; diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 5f058ac..7701d73 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -646,12 +646,12 @@ pub const Broadcaster = struct { /// disk playback runs on pool_io (Threaded) via request/reply — playback() /// holds the DiskPersist mutex and reads files, all of which require Threaded Io. pub fn replayTo(self: *Broadcaster, consumer: *Consumer, cursor: u64) void { - // phase 1: disk replay from diskpersist via cross-Io request/reply + // phase 1: disk replay from diskpersist via request/reply if (self.persist) |dp| { var req: event_log_mod.PlaybackRequest = .{ .since = cursor, .allocator = self.allocator }; dp.enqueuePlayback(&req); - // poll until pool_io worker completes the request (yields to Evented scheduler). + // poll until pool_io worker completes the request. // SAFETY: req is stack-local — we MUST wait for the worker to finish before // returning, otherwise the stack frame unwinds while the worker still holds &req. // if sleep fails (shutdown/io error), fall back to spin-wait. the host_ops worker @@ -720,7 +720,7 @@ pub const Broadcaster = struct { return self.consumers.items.len; } - /// broadcast loop — runs as an Evented fiber, drains the broadcast queue + /// broadcast loop — drains the broadcast queue, fans out to consumers /// and calls broadcast() for each item. this is the only path that touches /// consumers_mutex / consumer.mutex / consumer.cond. pub fn runBroadcastLoop(self: *Broadcaster) void { diff --git a/src/event_log.zig b/src/event_log.zig index 7e65814..9f9b6c1 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -77,7 +77,7 @@ const PersistJob = struct { // --- disk persistence --- -/// cross-Io playback request — Evented fiber posts, pool_io worker executes +/// playback request — caller posts, pool_io worker executes pub const PlaybackRequest = struct { since: u64, allocator: Allocator, @@ -87,7 +87,7 @@ pub const PlaybackRequest = struct { next: std.atomic.Value(?*PlaybackRequest) = .{ .raw = null }, }; -/// cross-Io DB request — Evented fiber posts, pool_io worker executes. +/// DB request — caller posts, pool_io worker executes. /// callers define typed structs embedding DbRequest + @fieldParentPtr. pub const DbRequest = struct { callback: *const fn (*DbRequest, *DiskPersist) void, @@ -114,7 +114,7 @@ pub const DbRequest = struct { }; /// MPSC FIFO ring buffer for general DB traffic. -/// multiple producers (Evented fibers) push via spinlock, +/// multiple producers push via spinlock, /// multiple consumers (pool_io worker threads) pop via CAS. pub const DbRequestQueue = struct { const CAPACITY = 4096; @@ -210,10 +210,10 @@ pub const DiskPersist = struct { io: Io, /// last successful DB interaction (epoch seconds, set by Threaded workers). - /// read by metrics server to report health without cross-Io pg.Pool access. + /// read by metrics server to report health without direct pg.Pool access. last_db_success: std.atomic.Value(i64) = .{ .raw = 0 }, - /// MPSC queue for cross-Io playback requests (Evented → pool_io) + /// MPSC queue for playback requests (caller → pool_io worker) playback_head: std.atomic.Value(?*PlaybackRequest) = .{ .raw = null }, /// current evtbuf entry count (for metrics — non-blocking, returns 0 if lock is contended) diff --git a/src/host_ops.zig b/src/host_ops.zig index 0b05b2c..0dd1448 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -1,4 +1,4 @@ -//! host ops — cross-Io host bookkeeping without pg.Pool from Evented fibers +//! host ops — background thread for host bookkeeping (cursor flush, failure tracking) //! //! two mechanisms: //! @@ -236,7 +236,7 @@ pub const HostOpsQueue = struct { drained += 1; } - // drain playback requests (cross-Io: Evented fibers post, we execute) + // drain playback requests drained += self.drainPlaybackRequests(); // periodic cursor sweep diff --git a/src/main.zig b/src/main.zig index 385c898..4eadc51 100644 --- a/src/main.zig +++ b/src/main.zig @@ -182,9 +182,9 @@ pub fn main() !void { const io = backend.io(); // dedicated Threaded runtime for the frame worker pool. - // worker threads are plain std.Thread — they cannot use Evented io - // (Evented futex calls ev.yield() which requires fiber context). - // this io is used for: persist ordering mutex, timestamps, validator cache, + // historically needed to isolate Threaded work from Evented fibers (cross-Io + // crash class). now redundant since Backend is also Threaded, but harmless. + // used for: persist ordering mutex, timestamps, validator cache, // DID resolution HTTP, and thread pool internal sync. var pool_io_backend = Io.Threaded.init(allocator, .{ .stack_size = default_stack_size, @@ -227,8 +227,9 @@ pub fn main() !void { dp.retention_hours = retention_hours; dp.max_dir_bytes = max_events_gb * 1024 * 1024 * 1024; - // DbRequestQueue — general DB traffic from Evented fibers routed to pool_io workers. - // replaces the broken ev_db (Evented pg.Pool) approach. + // DbRequestQueue — routes DB requests through pool_io worker threads. + // originally needed to bridge Evented fibers to Threaded pg.Pool; now + // redundant under all-Threaded but still functional. cleanup candidate. var db_queue: event_log_mod.DbRequestQueue = .{ .shutdown = &shutdown_flag, .persist = &dp, @@ -273,18 +274,16 @@ pub fn main() !void { var cleaner = cleaner_mod.Cleaner.init(allocator, pool_io, &ci, &dp, &shutdown_flag); // init resyncer (updates collection index on #sync events) - // runs entirely on pool_io (Threaded) — enqueue() is called from frame worker - // threads, and the background worker is a plain std.Thread. MUST NOT use - // Evented io (Threaded futex on Evented fiber blocks Uring thread → deadlock; - // Evented futex on plain thread → NULL threadlocal → SIGSEGV). + // runs on pool_io — enqueue() is called from frame worker threads, + // background worker is a plain std.Thread. var resyncer = resync_mod.Resyncer.init(allocator, pool_io, &ci); try resyncer.start(); defer resyncer.deinit(); - // init cursor map + host ops queue — subscriber fibers (Evented) write - // cursor seqs to the coalescing map (atomic store, no lock) and push rare - // ops (failures, status) to the MPSC queue. a background thread (Threaded) - // sweeps cursors every 5s and drains rare ops immediately. + // init cursor map + host ops queue — subscriber threads write cursor seqs + // to the coalescing map (atomic store, no lock) and push rare ops (failures, + // status) to the MPSC queue. a background thread sweeps cursors every 5s + // and drains rare ops immediately. var cursor_map: host_ops_mod.CursorMap = .{}; var host_ops_queue: host_ops_mod.HostOpsQueue = .{ .persist = &dp, @@ -339,15 +338,12 @@ pub fn main() !void { break :blk null; }; - // start broadcaster fiber — drains broadcast queue, owns all consumer state. - // this is the Evented-side sequencer: frame workers push results to the queue, - // this fiber does the actual fan-out to downstream consumers. + // start broadcast loop — drains broadcast queue, owns all consumer state. + // frame workers push results to the queue, this thread does the fan-out. var broadcast_future = try io.concurrent(broadcaster.Broadcaster.runBroadcastLoop, .{&bc}); defer _ = broadcast_future.cancel(io); - // start GC loop on a plain thread — dp.gc() uses pool_io (Threaded) mutex - // and pg.Pool. MUST NOT run as Evented fiber: Threaded futex on Evented - // fiber dereferences NULL Thread.current() threadlocal → heap corruption. + // start GC loop on a plain thread — dp.gc() uses pool_io mutex and pg.Pool. const gc_thread = std.Thread.spawn(.{}, gcLoop, .{ &dp, pool_io }) catch |err| { log.err("failed to start GC thread: {s}", .{@errorName(err)}); return err; diff --git a/src/resync.zig b/src/resync.zig index 1d48f5f..df8fc67 100644 --- a/src/resync.zig +++ b/src/resync.zig @@ -5,10 +5,8 @@ //! describeRepo from the PDS to get the current collection list, then //! replaces the index entries for that DID. //! -//! runs entirely on pool_io (Threaded) — enqueue() is called from frame worker -//! threads (plain std.Thread), and the background worker is also a plain thread. -//! MUST NOT use Evented io: Threaded futex on an Evented fiber blocks the Uring -//! thread, and Evented futex on a plain thread dereferences a NULL threadlocal. +//! runs on pool_io — enqueue() is called from frame worker threads, +//! background worker is a plain std.Thread. const std = @import("std"); const Io = std.Io; diff --git a/src/slurper.zig b/src/slurper.zig index 964de23..9c3e99f 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -494,7 +494,7 @@ pub const Slurper = struct { log.warn("host {s}: failed to spawn check thread", .{hostname}); return; }; - // wait for check to complete (Evented fiber yields) + // wait for check to complete while (!check_req.done.load(.acquire)) { if (self.shutdown.load(.acquire)) return; self.io.sleep(Io.Duration.fromMicroseconds(100), .awake) catch { diff --git a/src/subscriber.zig b/src/subscriber.zig index bb4c1b8..3dc904d 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -211,7 +211,7 @@ pub const Subscriber = struct { pool: ?*frame_worker_mod.FramePool = null, /// dedicated Threaded io for frame workers — safe from plain OS threads pool_io: ?Io = null, - /// host ops queue — pushes rare DB ops to a background thread (avoids cross-Io pg.Pool access) + /// host ops queue — pushes rare DB ops to a background thread host_ops: ?*host_ops_mod.HostOpsQueue = null, /// coalescing cursor map — subscriber writes latest seq atomically, worker sweeps every 5s cursor_map: ?*host_ops_mod.CursorMap = null, @@ -262,7 +262,7 @@ pub const Subscriber = struct { var backoff: u64 = 1; const max_backoff: u64 = 60; - // cursor is set at spawn time by slurper (avoids cross-Io pg.Pool access) + // cursor is set at spawn time by slurper if (self.last_upstream_seq) |seq| { log.info("host {s}: resuming from cursor {d}", .{ self.options.hostname, seq }); } -- 2.51.2