atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609//! zat relay — AT Protocol firehose relay server//!//! crawls PDS instances directly via the Slurper (one subscriber per host),//! validates frames via DID resolution and signature verification, persists//! to disk with relay-assigned seq numbers, and rebroadcasts to downstream//! consumers over WebSocket.//!//! port 3000 (RELAY_PORT): WebSocket firehose + HTTP API (via httpFallback)//! /xrpc/com.atproto.sync.subscribeRepos — firehose WebSocket (supports ?cursor=N)//! /xrpc/com.atproto.sync.listRepos — paginated account listing//! /xrpc/com.atproto.sync.getRepoStatus — single account status//! /xrpc/com.atproto.sync.getLatestCommit — latest commit CID + rev//! /xrpc/com.atproto.sync.listReposByCollection — repos with records in a collection//! /xrpc/com.atproto.sync.listHosts — paginated active host listing//! /xrpc/com.atproto.sync.getHostStatus — single host status//! /xrpc/com.atproto.sync.requestCrawl — request PDS crawl (POST)//! /admin/hosts — list all hosts (GET, admin)//! /admin/hosts/block — block a host (POST, admin)//! /admin/hosts/unblock — unblock a host (POST, admin)//! /_health, /_stats — health, stats//!//! port 3001 (RELAY_METRICS_PORT): internal metrics + health//! /metrics — prometheus metrics//! /_health — liveness probe (DB check)
const std = @import("std");const Io = std.Io;const http = std.http;const websocket = @import("websocket");const broadcaster = @import("internal/broadcaster.zig");const validator_mod = @import("internal/validator.zig");const slurper_mod = @import("internal/slurper.zig");const event_log_mod = @import("internal/event_log.zig");const collection_index_mod = @import("internal/collection_index/index.zig");const backfill_mod = @import("internal/collection_index/backfill.zig");const cleaner_mod = @import("internal/collection_index/cleaner.zig");const resync_mod = @import("internal/collection_index/resync.zig");const host_ops_mod = @import("internal/host_ops.zig");const api = @import("internal/api/main.zig");const util = @import("internal/util/util.zig");const build_options = @import("build_options");
const getenv = util.getenv;const parseEnvInt = util.parseEnvInt;const malloc_trim: ?*const fn (pad: usize) callconv(.c) c_int = if (builtin.os.tag == .linux) @extern(*const fn (pad: usize) callconv(.c) c_int, .{ .name = "malloc_trim" })else null;
const log = std.log.scoped(.relay);
pub const default_stack_size = 8 * 1024 * 1024;
// -- Io backend selection --// Io.Threaded: one OS thread per io.concurrent() call (~2,800 for PDS subscribers).// Higher thread count than Evented (~35) but proven stable at 99%+ coverage on 0.15.//// Evented (io_uring fibers) shelved: 8 crash classes from cross-Io violations,// ReleaseSafe GPF from fiber context-switch codegen bug, persistent ~10-15% coverage// degradation, zig team marks Evented as experimental. See docs/evented-attempt.md.// Can revisit when zig's Evented runtime matures.const backend_config = @import("backend_config");const zio = @import("zio");
/// -Dbackend=zio builds the same source against zio's std.Io implementation./// zio is a hybrid: it wraps an inner Io.Threaded and only overrides sockets,/// sleep, and concurrency, so anything it does not implement degrades to a/// thread rather than breaking.const Backend = if (backend_config.use_zio) zio.ZioBackend else Io.Threaded;
var backend: Backend = undefined;var debug_threaded_io: Io.Threaded = undefined;
/// override single-threaded debug_io default — required for std.debug.print safety/// when multiple OS threads exist (frame worker pool, websocket server).pub const std_options_debug_threaded_io: ?*Io.Threaded = &debug_threaded_io;
var shutdown_flag: std.atomic.Value(bool) = .{ .raw = false };
/// metrics-only server on the internal portconst MetricsServer = struct { server: Io.net.Server, io: Io, stats: *broadcaster.Stats, validator: *validator_mod.Validator, data_dir: []const u8, persist: *event_log_mod.DiskPersist, bc: *broadcaster.Broadcaster, slurper: *slurper_mod.Slurper,
fn run(self: *MetricsServer) void { while (!shutdown_flag.load(.acquire)) { const stream = self.server.accept(self.io) catch |err| { if (shutdown_flag.load(.acquire)) return; log.debug("metrics accept error: {s}", .{@errorName(err)}); continue; }; self.handleMetricsConn(stream); } }
fn handleMetricsConn(self: *MetricsServer, stream: Io.net.Stream) void { defer stream.close(self.io);
var recv_buf: [4096]u8 = undefined; var send_buf: [4096]u8 = undefined; var connection_reader = stream.reader(self.io, &recv_buf); var connection_writer = stream.writer(self.io, &send_buf); var server = http.Server.init(&connection_reader.interface, &connection_writer.interface);
var request = server.receiveHead() catch return; const path = request.head.target;
if (std.mem.eql(u8, path, "/_healthz")) { // trivial liveness — constant-time, no dependencies request.respond("{\"status\":\"ok\"}", .{ .status = .ok, .keep_alive = false, .extra_headers = &.{ .{ .name = "content-type", .value = "application/json" }, .{ .name = "server", .value = "zlay (atproto-relay)" }, } }) catch {}; } else if (std.mem.eql(u8, path, "/_health") or std.mem.eql(u8, path, "/_readyz")) { // use atomic health flag — pg.Pool runs on Threaded, metrics server is Evented const db_ok = self.persist.isDbHealthy(); const status: http.Status = if (db_ok) .ok else .internal_server_error; const body = if (db_ok) "{\"status\":\"ok\"}" else "{\"status\":\"error\",\"msg\":\"database unavailable\"}"; request.respond(body, .{ .status = status, .keep_alive = false, .extra_headers = &.{ .{ .name = "content-type", .value = "application/json" }, .{ .name = "server", .value = "zlay (atproto-relay)" }, } }) catch {}; } else if (std.mem.eql(u8, path, "/metrics")) { const cache_entries = self.validator.cacheSize(); const attribution = broadcaster.AttributionMetrics{ .history_entries = self.bc.history.count(), .evtbuf_entries = self.persist.evtbufLen(), .did_cache_entries = self.persist.didCacheLen(), .host_authority_cache_entries = self.validator.hostAuthorityCacheSize(), .resolve_queue_len = self.validator.resolveQueueLen(), .resolve_queued_set_count = self.validator.resolveQueuedSetCount(), .validator_cache_map_cap = self.validator.cacheMapCapacity(), .did_cache_map_cap = self.persist.didCacheMapCap(), .queued_set_map_cap = self.validator.resolveQueuedSetCapacity(), .evtbuf_cap = self.persist.evtbufCap(), .outbuf_cap = self.persist.outbufCap(), .workers_count = self.slurper.workerCount(), .db_queue_depth = if (self.bc.db_queue) |q| q.depth() else 0, .db_queue_handled_total = if (self.bc.db_queue) |q| q.handled.load(.acquire) else 0, };
var metrics_buf: [65536]u8 = undefined; const base_body = broadcaster.formatPrometheusMetrics(self.stats, cache_entries, attribution, self.data_dir, &metrics_buf, self.io);
// zio runtime internals. a wedged loop, fibers piling up on a // lost wake, or an eroded poller table are otherwise invisible // from outside the process — the 2026-08-04 canary's failure had // to be inferred from what stopped answering. the scrape is the // one forensic channel that survives a pod dying, so put them // there. this handler runs as a fiber on the loop thread, which // is where debugStats() must be called from. var zio_buf: [2048]u8 = undefined; const body = if (comptime backend_config.use_zio) blk: { const d = backend.debugStats(); const extra = std.fmt.bufPrint(&zio_buf, \\# HELP zio_sched_switches_total fiber context switches (flat = loop wedged) \\# TYPE zio_sched_switches_total counter \\zio_sched_switches_total {d} \\# HELP zio_sched_fibers live fibers by state \\# TYPE zio_sched_fibers gauge \\zio_sched_fibers{{state="live"}} {d} \\zio_sched_fibers{{state="ready"}} {d} \\zio_sched_fibers{{state="parked"}} {d} \\# HELP zio_sched_futex_waiters fibers parked on futex words (climbing = lost wakes) \\# TYPE zio_sched_futex_waiters gauge \\zio_sched_futex_waiters {d} \\# HELP zio_sched_timers timer heap depth \\# TYPE zio_sched_timers gauge \\zio_sched_timers {d} \\# HELP zio_sched_inject_pending cross-thread messages awaiting the loop \\# TYPE zio_sched_inject_pending gauge \\zio_sched_inject_pending {d} \\# HELP zio_sched_free_fibers recycled fiber structs \\# TYPE zio_sched_free_fibers gauge \\zio_sched_free_fibers {d} \\# HELP zio_poller_registrations fds with an armed direction (collapse = erosion) \\# TYPE zio_poller_registrations gauge \\zio_poller_registrations {d} \\# HELP zio_dns_spill DNS spill pool depth and workers \\# TYPE zio_dns_spill gauge \\zio_dns_spill{{kind="queued"}} {d} \\zio_dns_spill{{kind="workers"}} {d} \\zio_dns_spill{{kind="idle"}} {d} \\ , .{ d.sched.switches, d.sched.live, d.sched.ready, d.sched.parked, d.sched.futex_waiters, d.sched.timers, d.sched.inject_pending, d.sched.free_list, d.poller_registrations, d.dns_queued, d.dns_workers, d.dns_idle, }) catch ""; if (base_body.len + extra.len <= metrics_buf.len) { @memcpy(metrics_buf[base_body.len..][0..extra.len], extra); break :blk metrics_buf[0 .. base_body.len + extra.len]; } break :blk base_body; } else base_body;
request.respond(body, .{ .status = .ok, .keep_alive = false, .extra_headers = &.{ .{ .name = "content-type", .value = "text/plain; version=0.0.4; charset=utf-8" }, .{ .name = "server", .value = "zlay (atproto-relay)" }, } }) catch {}; } else { request.respond("not found", .{ .status = .not_found, .keep_alive = false, .extra_headers = &.{ .{ .name = "content-type", .value = "text/plain" }, .{ .name = "server", .value = "zlay (atproto-relay)" }, } }) catch {}; } }};
pub fn main() !void { try runRelay(); // clean shutdown reached. runRelay's defers (incl. the event-log flush in // DiskPersist.deinit) have already run, so all durable state is persisted. // now skip glibc's exit-time teardown (atexit handlers / locale free): with // ~2,800 Io.Threaded backend threads still live (the backends are never // deinit'd), that teardown races a worker mid-getaddrinfo/strtold and GPFs on // every shutdown (libc.so.6+0x50f). the raw exit_group(2) syscall terminates // all threads atomically in-kernel before any of them can fault — note we // must NOT use std.process.exit here, which calls libc exit(3) (running the // very atexit/locale teardown we're avoiding) when linked against libc. see // docs/shutdown-gpf-analysis.md. if (builtin.os.tag == .linux) std.os.linux.exit_group(0) else std.process.exit(0);}
fn runRelay() !void { // exp-002: optional GPA wrapper for leak detection. // build with -Duse_gpa=true to enable. on clean shutdown (SIGTERM), // GPA logs every allocation that was never freed, with stack traces. var gpa: std.heap.DebugAllocator(.{ .stack_trace_frames = if (build_options.use_gpa) 8 else 0, // catch use-after-free at the use site with the free site's stack: // freed regions are kept unmapped-but-tracked instead of recycled. .retain_metadata = build_options.use_gpa, .never_unmap = build_options.use_gpa, }) = .init; defer if (build_options.use_gpa) { log.info("GPA: checking for leaks...", .{}); const status = gpa.deinit(); if (status == .leak) { log.err("GPA: leaks detected! see stderr for details", .{}); } else { log.info("GPA: no leaks detected", .{}); } }; const allocator = if (build_options.use_gpa) gpa.allocator() else std.heap.c_allocator;
// shared options for both debug and primary runtimes const io_opts: Io.Threaded.InitOptions = .{ .stack_size = default_stack_size, // 8MB (default is 16MB) .concurrent_limit = Io.Limit.limited(4096), // safety rail (steady-state ~2,800 hosts) };
// init debug io (for std.debug.print thread safety) debug_threaded_io = Io.Threaded.init(allocator, io_opts);
// init primary runtime (Evented: fibers for network orchestration) if (Backend == Io.Threaded) { backend = Io.Threaded.init(allocator, io_opts); } else { // zio forwards io_opts to its inner Threaded, then start() spawns the // loop thread. Order matters: nothing may touch io() before start(). try backend.init(allocator, io_opts); try backend.start(); } const io = backend.io();
// dedicated Threaded runtime for the frame worker pool. // 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, }); const pool_io = pool_io_backend.io();
log.info("io backend: {s}", .{if (Backend == Io.Threaded) "Threaded" else "zio"});
// parse config from env const port = parseEnvInt(u16, "RELAY_PORT", 3000); const metrics_port = parseEnvInt(u16, "RELAY_METRICS_PORT", 3001); const upstream = normalizeSeedHost(getenv("RELAY_UPSTREAM") orelse "bsky.network"); const data_dir = getenv("RELAY_DATA_DIR") orelse "data/events"; const retention_hours = parseEnvInt(u64, "RELAY_RETENTION_HOURS", 72); const max_events_gb = parseEnvInt(u64, "RELAY_MAX_EVENTS_GB", 100); const frame_workers = parseEnvInt(u16, "FRAME_WORKERS", 16); const frame_queue_capacity = parseEnvInt(u16, "FRAME_QUEUE_CAPACITY", 4096); const startup_batch_size = parseEnvInt(u16, "STARTUP_BATCH_SIZE", 50); const db_pool_size = parseEnvInt(u16, "DB_POOL_SIZE", 20);
// install signal handlers (including SIGPIPE ignore) installSignalHandlers();
// init components — pass io to network-facing modules var bc = broadcaster.Broadcaster.init(allocator, io, &shutdown_flag); defer bc.deinit();
// validator uses pool_io — its cache LRU and host resolvers are called from worker threads var val = validator_mod.Validator.init(allocator, &bc.stats, pool_io); defer val.deinit(); try val.start();
// init disk persistence (indigo-compatible diskpersist format + Postgres index) const database_url = getenv("DATABASE_URL") orelse "postgres://relay:relay@localhost:5432/relay"; var dp = event_log_mod.DiskPersist.init(allocator, data_dir, database_url, db_pool_size, pool_io) catch |err| { log.err("failed to init disk persist at {s}: {s}", .{ data_dir, @errorName(err) }); return err; }; defer dp.deinit(); dp.retention_hours = retention_hours; dp.max_dir_bytes = max_events_gb * 1024 * 1024 * 1024;
// 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, };
if (dp.lastSeq()) |last| { log.info("event log recovered: last_seq={d}", .{last}); }
// start flush thread try dp.start();
// spawn 2 DbRequestQueue worker threads on pool_io const db_worker_1 = std.Thread.spawn(.{}, event_log_mod.DbRequestQueue.run, .{ &db_queue, pool_io }) catch |err| { log.err("failed to start db queue worker 1: {s}", .{@errorName(err)}); return err; }; const db_worker_2 = std.Thread.spawn(.{}, event_log_mod.DbRequestQueue.run, .{ &db_queue, pool_io }) catch |err| { log.err("failed to start db queue worker 2: {s}", .{@errorName(err)}); return err; };
// wire persist into broadcaster for cursor replay and validator for migration checks bc.persist = &dp; bc.db_queue = &db_queue; val.persist = &dp;
// init collection index (RocksDB — inspired by lightrail/microcosm.blue) const ci_dir = getenv("COLLECTION_INDEX_DIR") orelse "data/collection-index"; var ci = collection_index_mod.CollectionIndex.open(allocator, ci_dir) catch |err| { log.err("failed to init collection index at {s}: {s}", .{ ci_dir, @errorName(err) }); return err; }; defer ci.deinit();
// init backfiller (collection index backfill from source relay) // uses pool_io (Threaded) — backfiller spawns std.Thread, DNS works var backfiller = backfill_mod.Backfiller.init(allocator, &ci, &dp, pool_io, &shutdown_flag);
// init cleaner (removes stale entries from collection index) // uses pool_io (Threaded) — cleaner spawns std.Thread, checks shutdown var cleaner = cleaner_mod.Cleaner.init(allocator, pool_io, &ci, &dp, &shutdown_flag);
// init resyncer (updates collection index on #sync events) // 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 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, .cursor_map = &cursor_map, .shutdown = &shutdown_flag, .max_consecutive_failures = 15, .bc = &bc, }; const host_ops_thread = std.Thread.spawn(.{}, host_ops_mod.HostOpsQueue.run, .{ &host_ops_queue, pool_io }) catch |err| { log.err("failed to start host ops thread: {s}", .{@errorName(err)}); return err; };
// init slurper (multi-host crawl manager) var slurper = slurper_mod.Slurper.init( allocator, &bc, &val, &dp, &shutdown_flag, .{ .seed_host = upstream, .max_message_size = 5 * 1024 * 1024, .frame_workers = frame_workers, .frame_queue_capacity = frame_queue_capacity, .startup_batch_size = startup_batch_size, }, io, pool_io, ); defer slurper.deinit(); slurper.collection_index = &ci; slurper.resyncer = &resyncer; slurper.host_ops = &host_ops_queue; slurper.cursor_map = &cursor_map; slurper.db_queue = &db_queue;
// start: loads active hosts from DB, spawns subscriber threads. // pullHosts runs on its own std.Thread (parallel, not gating spawnWorkers). try slurper.start();
// spawn pullHosts on a separate thread — runs in parallel with spawnWorkers, // uses pool_io for HTTP (DNS works), writes to DB via Threaded pool directly. const pull_hosts_thread = if (slurper.options.seed_host.len > 0) blk: { log.info("pulling hosts from {s} (background thread)...", .{slurper.options.seed_host}); break :blk std.Thread.spawn(.{}, slurper_mod.Slurper.pullHosts, .{&slurper}) catch |err| { log.warn("failed to spawn pullHosts thread: {s}", .{@errorName(err)}); break :blk null; }; } else blk: { log.info("no seed host configured, skipping bootstrap", .{}); break :blk null; };
// 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 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; };
// wire HTTP fallback into broadcaster (all API endpoints served on WS port) var http_context = api.HttpContext{ .stats = &bc.stats, .persist = &dp, .slurper = &slurper, .collection_index = &ci, .backfiller = &backfiller, .cleaner = &cleaner, .resyncer = &resyncer, .bc = &bc, .validator = &val, .host_ops = &host_ops_queue, .pool_io = pool_io, .io = io, .db_queue = &db_queue, .shutdown = &shutdown_flag, }; bc.http_fallback = api.handleHttpRequest; bc.http_fallback_ctx = @ptrCast(&http_context);
// start metrics-only server (internal port) const metrics_address = Io.net.Ip4Address.unspecified(metrics_port); var metrics_srv = MetricsServer{ .server = (Io.net.IpAddress{ .ip4 = metrics_address }).listen(io, .{ .reuse_address = true }) catch |err| { log.err("metrics server failed to listen on :{d}: {s}", .{ metrics_port, @errorName(err) }); return err; }, .io = io, .stats = &bc.stats, .validator = &val, .data_dir = data_dir, .persist = &dp, .bc = &bc, .slurper = &slurper, }; var metrics_future = try io.concurrent(MetricsServer.run, .{&metrics_srv}); defer _ = metrics_future.cancel(io);
// start downstream WebSocket server (also serves HTTP API via httpFallback) log.info("relay listening on :{d} (ws+http), :{d} (metrics)", .{ port, metrics_port }); log.info("seed host: {s}", .{upstream}); log.info("data dir: {s} (retention: {d}h, max: {d} GB)", .{ data_dir, retention_hours, max_events_gb });
var server = try websocket.Server(broadcaster.Handler).init(allocator, io, .{ .port = port, .address = "0.0.0.0", .max_conn = 4096, .max_message_size = 5 * 1024 * 1024, }); defer server.deinit();
// Io-native accept loop: fiber-based under Evented, thread-based under Threaded const ws_address = Io.net.Ip4Address.unspecified(port); var ws_listener = (Io.net.IpAddress{ .ip4 = ws_address }).listen(io, .{ .reuse_address = true }) catch |err| { log.err("websocket server failed to listen on :{d}: {s}", .{ port, @errorName(err) }); return err; }; var server_future = try io.concurrent(runWsServer, .{ &server, &ws_listener, &bc }); defer _ = server_future.cancel(io);
// wait for shutdown signal while (!shutdown_flag.load(.acquire)) { io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch break; }
log.info("shutdown signal received, stopping...", .{});
// stop WebSocket server: close listener to unblock accept, then cancel task ws_listener.deinit(io); server_future.cancel(io);
// join GC thread — it checks shutdown_flag and will exit its sleep loop. // must complete before dp.deinit() runs (dp is stack-owned). gc_thread.join();
// join pullHosts thread if running if (pull_hosts_thread) |t| t.join();
// join db queue workers — shutdown flag already set, they will drain and exit db_worker_1.join(); db_worker_2.join();
// join host ops thread — drains remaining ops before dp.deinit() host_ops_thread.join();
// join backfiller/cleaner threads if running — they touch dp and ci backfiller.waitForCompletion(); cleaner.waitForCompletion();
// cancel broadcaster fiber (shutdown flag already set, it will drain remaining) broadcast_future.cancel(io);
// close metrics listener to unblock accept(), then cancel task metrics_srv.server.deinit(io); metrics_future.cancel(io);
log.info("relay stopped cleanly", .{});}
const builtin = @import("builtin");
/// concrete wrapper for runIo — io.concurrent needs ArgsTuple, which can't handle anytypefn runWsServer(server: *websocket.Server(broadcaster.Handler), listener: *Io.net.Server, bc: *broadcaster.Broadcaster) void { server.runIo(listener, bc);}
fn gcLoop(dp: *event_log_mod.DiskPersist, io: Io) void { const gc_interval: u64 = 10 * 60; // 10 minutes in seconds while (!shutdown_flag.load(.acquire)) { // sleep in small increments to check shutdown var remaining: u64 = gc_interval; while (remaining > 0 and !shutdown_flag.load(.acquire)) { const chunk = @min(remaining, 1); io.sleep(Io.Duration.fromSeconds(@intCast(chunk)), .awake) catch return; remaining -= chunk; } if (shutdown_flag.load(.acquire)) return;
dp.gc() catch |err| { log.warn("event log GC failed: {s}", .{@errorName(err)}); };
// return freed pages to OS (glibc-specific, no-op on other platforms) if (comptime malloc_trim) |trim| { _ = trim(0); log.info("gc: malloc_trim complete", .{}); } }}
fn signalHandler(_: std.posix.SIG) callconv(.c) void { shutdown_flag.store(true, .release);}
fn installSignalHandlers() void { const act: std.posix.Sigaction = .{ .handler = .{ .handler = signalHandler }, .mask = std.posix.sigemptyset(), .flags = 0, }; std.posix.sigaction(std.posix.SIG.INT, &act, null); std.posix.sigaction(std.posix.SIG.TERM, &act, null);
// ignore SIGPIPE — writing to disconnected consumers must not crash the process const ignore_act: std.posix.Sigaction = .{ .handler = .{ .handler = std.posix.SIG.IGN }, .mask = std.posix.sigemptyset(), .flags = 0, }; std.posix.sigaction(std.posix.SIG.PIPE, &ignore_act, null);}
/// normalize seed host: empty string or "none" means disabled (no bootstrap).fn normalizeSeedHost(raw: []const u8) []const u8 { if (raw.len == 0 or std.mem.eql(u8, raw, "none")) return ""; return raw;}
test "normalizeSeedHost" { // normal hostnames pass through try std.testing.expectEqualStrings("bsky.network", normalizeSeedHost("bsky.network")); try std.testing.expectEqualStrings("relay.example.com", normalizeSeedHost("relay.example.com"));
// empty string and "none" disable bootstrap try std.testing.expectEqualStrings("", normalizeSeedHost("")); try std.testing.expectEqualStrings("", normalizeSeedHost("none"));}