jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151//! stream — a zig jetstream. entrypoint: parse args, wire components.
const std = @import("std");const zat = @import("zat");const websocket = @import("websocket");const archive_mod = @import("internal/storage/archive.zig");const lifecycle = @import("internal/ingest/backfill/lifecycle.zig");const repo_store = @import("internal/ingest/backfill/repo_store.zig");const phase_mod = @import("internal/ingest/backfill/phase.zig");const retry_mod = @import("internal/ingest/backfill/retry.zig");const backfill_engine = @import("internal/ingest/backfill/engine.zig");const crashpoint = @import("internal/runtime/crashpoint.zig");const cursor_mod = @import("internal/ingest/cursor.zig");const ingest = @import("internal/ingest/ingest.zig");const verify = @import("internal/ingest/verify.zig");const xrpcapi = @import("internal/serve/xrpcapi.zig");const api_keys_mod = @import("internal/serve/api_keys.zig");const metrics = @import("internal/runtime/metrics.zig");const meta_store = @import("internal/storage/meta_store.zig");const pipeline_mod = @import("internal/ingest/pipeline.zig");const repair_mod = @import("internal/ingest/repair.zig");const server_mod = @import("internal/serve/server.zig");const segment_io = @import("internal/storage/segment_io.zig");const compact_pass = @import("internal/compact/pass.zig");const cold = @import("internal/storage/cold.zig");const steady_mod = @import("internal/compact/steady.zig");const rebloom_mod = @import("internal/compact/rebloom.zig");const tail_mod = @import("internal/serve/tail.zig");const timestamp_rules = @import("internal/timestamp/rules.zig");const timestamp_jobs = @import("internal/timestamp/jobs.zig");const timestamp_manager = @import("internal/timestamp/manager.zig");const operational = @import("internal/runtime/operational.zig");const inspect_all = @import("internal/runtime/inspect_all.zig");const logging = @import("internal/runtime/logging.zig");const environment = @import("internal/runtime/environment.zig");const cli = @import("internal/runtime/cli.zig");const observability = @import("internal/runtime/observability.zig");const build_options = @import("build_options");
const Io = std.Io;const log = std.log.scoped(.stream);
pub const std_options: std.Options = .{ // Runtime filtering happens in logging.logFn so ReleaseSafe retains the // canonical --log-level=debug surface instead of compiling it away. .log_level = .debug, .logFn = logging.logFn,};
const default_subscribe_read_log_retention_bytes: usize = 256 * 1024 * 1024;
/// 8 MB: ReleaseSafe inlining makes TLS/CBOR/crypto call chains deep enough/// to overflow 2-4 MB stacks; zig's 16 MB default maps needless VM at scale/// (docs/lessons-from-zlay.md #1)const default_stack_size: usize = 8 * 1024 * 1024;
var shutdown_requested: std.atomic.Value(bool) = .init(false);
const LifecycleRun = struct { allocator: std.mem.Allocator, io: Io, options: lifecycle.Options, archive: *archive_mod.Archive, meta: *meta_store.Store, cursor_store: *cursor_mod.Store, verifier: ?*verify.Verifier, stats: *metrics.Stats, trace_context: observability.Context, done: std.atomic.Value(bool) = .init(false), result: ?anyerror = null,};
fn signalHandler(_: std.posix.SIG) callconv(.c) void { // Atomic storage is the only work performed in signal context. The main // task observes it and owns cancellation, pipeline drain, durable flush, // and all allocator/I/O cleanup. shutdown_requested.store(true, .release);}
fn installSignalHandlers() void { const action: std.posix.Sigaction = .{ .handler = .{ .handler = signalHandler }, .mask = std.posix.sigemptyset(), .flags = 0, }; std.posix.sigaction(std.posix.SIG.INT, &action, null); std.posix.sigaction(std.posix.SIG.TERM, &action, null);
const ignore_pipe: std.posix.Sigaction = .{ .handler = .{ .handler = std.posix.SIG.IGN }, .mask = std.posix.sigemptyset(), .flags = 0, }; std.posix.sigaction(std.posix.SIG.PIPE, &ignore_pipe, null);}
fn parseBackfillRepos(allocator: std.mem.Allocator, raw: []const u8) ![]const []const u8 { if (std.mem.trim(u8, raw, " \t\r\n").len == 0) return &.{}; var out: std.ArrayList([]const u8) = .empty; errdefer out.deinit(allocator); var it = std.mem.splitScalar(u8, raw, ','); while (it.next()) |part| { const did = std.mem.trim(u8, part, " \t\r\n"); if (did.len == 0 or zat.Did.parse(did) == null) return error.InvalidBackfillRepo; for (out.items) |prior| if (std.mem.eql(u8, prior, did)) return error.DuplicateBackfillRepo; try out.append(allocator, did); } if (out.items.len == 0) return error.InvalidBackfillRepo; return out.toOwnedSlice(allocator);}
fn effectiveSkipMergeDiscovery(requested: bool, max_repos: usize, selected_repos: usize) bool { return requested or max_repos > 0 or selected_repos > 0;}
const RelayUrls = struct { http: []u8, websocket: []u8,
fn deinit(self: RelayUrls, allocator: std.mem.Allocator) void { allocator.free(self.http); allocator.free(self.websocket); }};
fn normalizeRelayUrls(allocator: std.mem.Allocator, relay_url: []const u8) !RelayUrls { const secure = std.mem.startsWith(u8, relay_url, "https://"); const prefix_len: usize = if (secure) "https://".len else if (std.mem.startsWith(u8, relay_url, "http://")) "http://".len else return error.InvalidRelayUrl; const rest = relay_url[prefix_len..]; const authority_len = std.mem.indexOfAny(u8, rest, "/?#") orelse rest.len; if (authority_len == 0 or std.mem.indexOfAny(u8, rest[0..authority_len], " \t\r\n") != null) return error.InvalidRelayUrl; const authority = rest[0..authority_len]; const http = try std.fmt.allocPrint(allocator, "{s}://{s}", .{ if (secure) "https" else "http", authority }); errdefer allocator.free(http); return .{ .http = http, .websocket = try std.fmt.allocPrint(allocator, "{s}://{s}", .{ if (secure) "wss" else "ws", authority }), };}
test "canonical relay URL normalizes HTTP and live websocket endpoints" { const testing = std.testing; const secure = try normalizeRelayUrls(testing.allocator, "https://relay.example/ignored?x=1#fragment"); defer secure.deinit(testing.allocator); try testing.expectEqualStrings("https://relay.example", secure.http); try testing.expectEqualStrings("wss://relay.example", secure.websocket); const local = try normalizeRelayUrls(testing.allocator, "http://127.0.0.1:7777"); defer local.deinit(testing.allocator); try testing.expectEqualStrings("http://127.0.0.1:7777", local.http); try testing.expectEqualStrings("ws://127.0.0.1:7777", local.websocket); try testing.expectError(error.InvalidRelayUrl, normalizeRelayUrls(testing.allocator, "ws://relay.example")); try testing.expectError(error.InvalidRelayUrl, normalizeRelayUrls(testing.allocator, "https:///missing"));}
test "explicit backfill repo parsing matches upstream selection rules" { const testing = std.testing; const parsed = try parseBackfillRepos(testing.allocator, "did:plc:aaa, did:web:example.com"); defer testing.allocator.free(parsed); try testing.expectEqual(@as(usize, 2), parsed.len); try testing.expectEqualStrings("did:plc:aaa", parsed[0]); try testing.expectEqualStrings("did:web:example.com", parsed[1]); try testing.expectEqual(@as(usize, 0), (try parseBackfillRepos(testing.allocator, " \t ")).len); try testing.expectError(error.DuplicateBackfillRepo, parseBackfillRepos(testing.allocator, "did:plc:aaa,did:plc:aaa")); try testing.expectError(error.InvalidBackfillRepo, parseBackfillRepos(testing.allocator, "did:plc:aaa,")); try testing.expectError(error.InvalidBackfillRepo, parseBackfillRepos(testing.allocator, "not-a-did"));}
test "merge discovery skip is explicit or automatic for either debug selection" { const testing = std.testing; try testing.expect(!effectiveSkipMergeDiscovery(false, 0, 0)); try testing.expect(effectiveSkipMergeDiscovery(true, 0, 0)); try testing.expect(effectiveSkipMergeDiscovery(false, 1, 0)); try testing.expect(effectiveSkipMergeDiscovery(false, 0, 1));}
/// Parse the non-negative subset of Go's time.ParseDuration used by upstream's/// duration flags. Like Go, fractional values are truncated below 1 ns./// Full signed Go duration surface used by urfave/cli DurationFlag. The/// non-negative parser above remains the implementation shared by flags whose/// downstream types cannot represent negative durations./// net/http Server.New applies this default when its configured timeout is/// non-positive; the CLI's websocket client-drain duration does not.fn writeVersion(io: Io, args: []const []const u8) !void { if (args.len != 0) return error.BadArgs; var buffer: [1024]u8 = undefined; var stdout = Io.File.stdout().writer(io, &buffer); try operational.renderVersion(&stdout.interface); try stdout.interface.flush();}
fn runInspectSegment(allocator: std.mem.Allocator, io: Io, args: []const []const u8) !void { var mode: operational.BlocksMode = .table; var truncate: usize = 100; var path: ?[]const u8 = null; var i: usize = 0; while (i < args.len) : (i += 1) { const arg = args[i]; if (std.mem.startsWith(u8, arg, "--blocks=")) { mode = parseBlocksMode(arg["--blocks=".len..]) orelse return error.BadArgs; } else if (std.mem.eql(u8, arg, "--blocks")) { i += 1; if (i >= args.len) return error.BadArgs; mode = parseBlocksMode(args[i]) orelse return error.BadArgs; } else if (std.mem.startsWith(u8, arg, "--blocks-truncate=")) { truncate = try cli.parseNonNegativeUsize(arg["--blocks-truncate=".len..]); } else if (std.mem.eql(u8, arg, "--blocks-truncate")) { i += 1; if (i >= args.len) return error.BadArgs; truncate = try cli.parseNonNegativeUsize(args[i]); } else if (std.mem.startsWith(u8, arg, "--")) { return error.BadArgs; } else if (path != null) { return error.BadArgs; } else { path = arg; } } var inspection = try operational.inspectSegment(allocator, io, path orelse return error.BadArgs); defer inspection.deinit(); var buffer: [16 * 1024]u8 = undefined; var stdout = Io.File.stdout().writer(io, &buffer); try operational.renderInspection(&stdout.interface, &inspection, mode, truncate); try stdout.interface.flush();}
fn runInspectAll(allocator: std.mem.Allocator, io: Io, args: []const []const u8) !void { var data_dir: []const u8 = "./data"; var skip_unsealed = false; var truncate: usize = 100; var i: usize = 0; while (i < args.len) : (i += 1) { const arg = args[i]; if (std.mem.startsWith(u8, arg, "--data-dir=")) { data_dir = arg["--data-dir=".len..]; } else if (std.mem.eql(u8, arg, "--data-dir")) { i += 1; if (i >= args.len) return error.BadArgs; data_dir = args[i]; } else if (std.mem.eql(u8, arg, "--skip-unsealed")) { skip_unsealed = true; } else if (std.mem.startsWith(u8, arg, "--collections-truncate=")) { truncate = try cli.parseNonNegativeUsize(arg["--collections-truncate=".len..]); } else if (std.mem.eql(u8, arg, "--collections-truncate")) { i += 1; if (i >= args.len) return error.BadArgs; truncate = try cli.parseNonNegativeUsize(args[i]); } else { return error.BadArgs; } } var aggregate = try inspect_all.inspectAll(allocator, io, data_dir, .{ .skip_unsealed = skip_unsealed }); defer aggregate.deinit(); var buffer: [16 * 1024]u8 = undefined; var stdout = Io.File.stdout().writer(io, &buffer); try inspect_all.renderInspectAll( &stdout.interface, data_dir, &aggregate, Io.Timestamp.now(io, .real).toMicroseconds(), truncate, ); try stdout.interface.flush();}
fn parseBlocksMode(raw: []const u8) ?operational.BlocksMode { if (std.mem.eql(u8, raw, "summary")) return .summary; if (std.mem.eql(u8, raw, "table")) return .table; if (std.mem.eql(u8, raw, "full")) return .full; return null;}
const BindAddress = struct { address: Io.net.IpAddress, port: u16,};
fn splitBindAddress(raw: []const u8) !struct { host: []const u8, port: u16 } { if (raw.len == 0) return error.InvalidBindAddress; if (raw[0] == '[') { const close = std.mem.indexOfScalar(u8, raw, ']') orelse return error.InvalidBindAddress; if (close + 1 >= raw.len or raw[close + 1] != ':') return error.InvalidBindAddress; return .{ .host = raw[1..close], .port = std.fmt.parseInt(u16, raw[close + 2 ..], 10) catch return error.InvalidBindAddress, }; } const colon = std.mem.lastIndexOfScalar(u8, raw, ':') orelse return error.InvalidBindAddress; if (std.mem.indexOfScalar(u8, raw[0..colon], ':') != null) return error.InvalidBindAddress; return .{ .host = raw[0..colon], .port = std.fmt.parseInt(u16, raw[colon + 1 ..], 10) catch return error.InvalidBindAddress, };}
fn resolveBindAddress(io: Io, raw: []const u8) !BindAddress { const split = try splitBindAddress(raw); if (split.host.len == 0) return .{ .address = .{ .ip4 = .unspecified(split.port) }, .port = split.port, }; if (Io.net.IpAddress.parse(split.host, split.port)) |address| return .{ .address = address, .port = split.port, } else |_| {}
const host = Io.net.HostName.init(split.host) catch return error.InvalidBindAddress; var results_buffer: [16]Io.net.HostName.LookupResult = undefined; var results: Io.Queue(Io.net.HostName.LookupResult) = .init(&results_buffer); var lookup = io.async(Io.net.HostName.lookup, .{ host, io, &results, .{ .port = split.port } }); defer lookup.cancel(io) catch {}; var first: ?Io.net.IpAddress = null; while (true) { const result = results.getOne(io) catch |err| switch (err) { error.Closed => break, error.Canceled => return error.Canceled, }; switch (result) { .address => |address| if (first == null) { first = address; }, .canonical_name => {}, } } try lookup.await(io); return .{ .address = first orelse return error.InvalidBindAddress, .port = split.port };}
test "canonical bind address parsing covers upstream listener forms" { const testing = std.testing; const wildcard = try splitBindAddress(":8080"); try testing.expectEqualStrings("", wildcard.host); try testing.expectEqual(@as(u16, 8080), wildcard.port); const ipv4 = try splitBindAddress("127.0.0.1:6060"); try testing.expectEqualStrings("127.0.0.1", ipv4.host); try testing.expectEqual(@as(u16, 6060), ipv4.port); const ipv6 = try splitBindAddress("[::1]:0"); try testing.expectEqualStrings("::1", ipv6.host); try testing.expectEqual(@as(u16, 0), ipv6.port); try testing.expectError(error.InvalidBindAddress, splitBindAddress("::1:8080")); try testing.expectError(error.InvalidBindAddress, splitBindAddress("localhost")); try testing.expectError(error.InvalidBindAddress, splitBindAddress("127.0.0.1:nope"));}
pub fn main(init: std.process.Init.Minimal) !void { const allocator = std.heap.c_allocator;
var threaded: Io.Threaded = .init(allocator, .{ .stack_size = default_stack_size }); defer threaded.deinit(); const io = threaded.io();
const env_entries = try environment.processEntries(allocator); defer allocator.free(env_entries); const unknown_env = try environment.unknown(allocator, env_entries); defer allocator.free(unknown_env); if (unknown_env.len == 1) { log.err("unrecognized JETSTREAM_ environment variable {s}", .{unknown_env[0]}); return error.UnknownJetstreamEnvironmentVariable; } else if (unknown_env.len > 1) { var joined: std.Io.Writer.Allocating = .init(allocator); defer joined.deinit(); for (unknown_env, 0..) |name, i| { if (i > 0) try joined.writer.writeAll(", "); try joined.writer.writeAll(name); } log.err("unrecognized JETSTREAM_ environment variables: {s}", .{joined.written()}); return error.UnknownJetstreamEnvironmentVariable; }
var arg_it = init.args.iterate(); _ = arg_it.next(); // program name var log_level: []const u8 = environment.value(env_entries, "JETSTREAM_LOG_LEVEL") orelse "info"; var log_format: []const u8 = environment.value(env_entries, "JETSTREAM_LOG_FORMAT") orelse "json"; var first_arg = arg_it.next(); // urfave/cli root flags are persistent: accept them before every command, // with either of the spellings supported by the upstream CLI. while (first_arg) |arg| { if (std.mem.startsWith(u8, arg, "--log-level=")) { log_level = arg["--log-level=".len..]; } else if (std.mem.eql(u8, arg, "--log-level")) { log_level = arg_it.next() orelse return error.BadArgs; } else if (std.mem.startsWith(u8, arg, "--log-format=")) { log_format = arg["--log-format=".len..]; } else if (std.mem.eql(u8, arg, "--log-format")) { log_format = arg_it.next() orelse return error.BadArgs; } else break; first_arg = arg_it.next(); } var explicit_serve = false; if (first_arg) |command| { if (std.mem.eql(u8, command, "version") or std.mem.eql(u8, command, "inspect-segment") or std.mem.eql(u8, command, "inspect-all")) { var command_args: std.ArrayList([]const u8) = .empty; defer command_args.deinit(allocator); var inspect_env_arg: ?[]u8 = null; defer if (inspect_env_arg) |owned| allocator.free(owned); if (std.mem.eql(u8, command, "inspect-all")) { if (environment.value(env_entries, "JETSTREAM_DATA_DIR")) |configured| { inspect_env_arg = try std.fmt.allocPrint(allocator, "--data-dir={s}", .{configured}); try command_args.append(allocator, inspect_env_arg.?); } } while (arg_it.next()) |arg| { if (std.mem.startsWith(u8, arg, "--log-level=")) { log_level = arg["--log-level=".len..]; } else if (std.mem.eql(u8, arg, "--log-level")) { log_level = arg_it.next() orelse return error.BadArgs; } else if (std.mem.startsWith(u8, arg, "--log-format=")) { log_format = arg["--log-format=".len..]; } else if (std.mem.eql(u8, arg, "--log-format")) { log_format = arg_it.next() orelse return error.BadArgs; } else try command_args.append(allocator, arg); } try logging.configure(log_level, log_format); if (std.mem.eql(u8, command, "version")) { try writeVersion(io, command_args.items); } else if (std.mem.eql(u8, command, "inspect-segment")) { try runInspectSegment(allocator, io, command_args.items); } else { try runInspectAll(allocator, io, command_args.items); } return; } // Like upstream's explicit `serve` subcommand, while retaining the // historical flag-only invocation for existing Stream deployments. if (std.mem.eql(u8, command, "serve")) { explicit_serve = true; first_arg = arg_it.next(); } }
shutdown_requested.store(false, .release); installSignalHandlers();
var parsed = try cli.parseServe(allocator, env_entries, first_arg, &arg_it, explicit_serve, log_level, log_format); defer parsed.deinit(allocator); const cfg = &parsed.config; explicit_serve = parsed.explicit_serve; log_level = parsed.log_level; log_format = parsed.log_format; try logging.configure(log_level, log_format); log.debug("logging configured: level={s} format={s}", .{ log_level, log_format }); if (cfg.legacy_port_set and cfg.public_addr_raw != null) { log.err("--port cannot be combined with upstream --addr/--debug-addr or explicit serve mode", .{}); return error.BadArgs; } if (!(cfg.plan_config.whole_segment_threshold > 0 and cfg.plan_config.whole_segment_threshold <= 1)) { log.err("--plan-whole-segment-threshold must be > 0 and <= 1", .{}); return error.BadArgs; } if (cfg.store_fault_prefix) |prefix| { if (!meta_store.armFault(prefix, cfg.store_fault_ordinal)) { log.err("invalid store fault: prefix must be non-empty and ordinal must be positive", .{}); return error.BadArgs; } } if (cfg.segment_fault_spec) |spec| { if (!segment_io.armSpec(spec)) { log.err("invalid segment fault; expected op:positive-ordinal:kind", .{}); return error.BadArgs; } } const backfill_repos: []const []const u8 = if (cfg.backfill_repos_raw) |raw| try parseBackfillRepos(allocator, raw) else &.{}; defer if (backfill_repos.len > 0) allocator.free(backfill_repos); if (backfill_repos.len > 0 and cfg.backfill_max_repos > 0) return error.ConflictingBackfillSelection; // Pinned cfg.upstream automatically enables the debug short circuit for // either partial selection mode, in addition to accepting it explicitly. cfg.skip_merge_discovery = effectiveSkipMergeDiscovery(cfg.skip_merge_discovery, cfg.backfill_max_repos, backfill_repos.len);
var tracing = observability.Runtime.init(allocator, .{ .io = io, .environment = env_entries, .service_name = cfg.otel_service_name, .service_version = build_options.version, .shutdown_timeout_ms = shutdownTimeoutMillis(cfg.shutdown_timeout), }) catch |err| { log.err("setup tracing failed: {s}", .{@errorName(err)}); return @errorCast(err); }; defer tracing.deinit(); if (tracing.enabled()) log.info("OTLP tracing enabled (service.name={s})", .{cfg.otel_service_name});
var orchestrator_trace = observability.Span.start( allocator, "ingest/orchestrator", "Run", observability.root_context, &.{}, ); defer orchestrator_trace.deinit(); defer orchestrator_trace.succeed(); try orchestrator_trace.recordError(serve(allocator, io, cfg, backfill_repos, &tracing, &orchestrator_trace));}
fn serve( allocator: std.mem.Allocator, io: Io, cfg: *cli.ServeConfig, backfill_repos: []const []const u8, tracing: *observability.Runtime, orchestrator_trace: *observability.Span,) !void { var hot_tail = tail_mod.Tail.init(allocator, io, cfg.subscribe_read_log_retention_bytes); defer hot_tail.deinit(); var stats: metrics.Stats = .{ .start_time_s = @divTrunc(Io.Timestamp.now(io, .real).toMicroseconds(), std.time.us_per_s), }; var hub: server_mod.Hub = .{ .allocator = allocator, .io = io, .tail = &hot_tail, .stats = &stats, .bootstrap_enabled = cfg.do_backfill or cfg.backfill_max_repos > 0 or backfill_repos.len > 0, .compaction_enabled = cfg.compaction_interval_ns > 0, .retry_enabled = cfg.retry_interval_ns > 0, .plc_url = cfg.plc_url, .data_dir = cfg.data_dir, .repo_action_rate_limits = cfg.repo_action_rate_limits, .cursor_lookback_ns = cfg.cursor_lookback_ns, .subscribe_slow_window_ns = cfg.subscribe_slow_window_ns, .subscribe_slow_min_rate = cfg.subscribe_slow_min_rate, .subscribe_read_batch = cfg.subscribe_read_batch, .max_subscribers = cfg.max_subscribers, .max_cold_readers = cfg.max_cold_readers, }; defer hub.deinit(); log.info("subscribe controls: read-log-retention={d} bytes block-cache={d} bytes read-batch={d} slow-window={d}ns slow-min-rate={d}", .{ cfg.subscribe_read_log_retention_bytes, cfg.subscribe_block_cache_bytes, cfg.subscribe_read_batch, cfg.subscribe_slow_window_ns, cfg.subscribe_slow_min_rate, });
// Upstream starts its server and registers every route before storage // recovery, gating data paths with 503 until ready (jetstreamd // runtime.go: lifecycle.IsSteadyState on subscribe, the xrpcapi Ready // hook on the archive surface). Same ordering here: bind and accept // immediately — with both gates closed — so a deploy never presents // connection-refused while segment recovery and rocksdb open run // (measured ~2.5 min on 2026-08-16). The gates open below once the // stores are published and the durable lifecycle is resolved. Cleanup // for the server is deliberately registered AFTER storage init so // shutdown still cancels the accept loops before the stores they read // are deinitialized; an error return in between exits the process, // which reclaims the sockets. hub.serving.store(false, .release); hub.storage_ready.store(false, .release); const public_bind = if (cfg.public_addr_raw) |raw| resolveBindAddress(io, raw) catch |err| { log.err("invalid public bind address {s}: {s}", .{ raw, @errorName(err) }); return @errorCast(err); } else BindAddress{ .address = .{ .ip4 = .unspecified(cfg.port) }, .port = cfg.port }; const debug_bind: ?BindAddress = if (cfg.debug_addr_raw) |raw| resolveBindAddress(io, raw) catch |err| { log.err("invalid debug bind address {s}: {s}", .{ raw, @errorName(err) }); return @errorCast(err); } else null; const canonical_listeners = cfg.public_addr_raw != null; var public_context: server_mod.ListenerContext = .{ .hub = &hub, .role = if (canonical_listeners) .public else .legacy, }; var debug_context: server_mod.ListenerContext = .{ .hub = &hub, .role = .debug };
var ws_server = try websocket.Server(server_mod.Handler).init(allocator, io, .{ .port = public_bind.port, .address = "0.0.0.0", .max_conn = 4096, .max_message_size = 10_000_000, .websocket_shutdown_grace = cfg.client_drain_timeout, .http_shutdown_grace = cli.effectiveHttpShutdownDuration(cfg.shutdown_timeout), // Match cfg.upstream v1's coder/websocket defaults exactly. The handler // declines this per connection for v2 and custom-zstd subscribers. .compression = .{ .write_threshold = 128, .client_no_context_takeover = false, .server_no_context_takeover = false, }, // The websocket package defaults to a 1 KiB / 10-header handshake, // which rejects ordinary modern browser HTTP requests before the // homepage or status fallback can run. This remains tightly bounded // while accommodating browser and reverse-proxy headers. .handshake = .{ .max_size = 16 * 1024, .max_headers = 64 }, }); var listener = public_bind.address.listen(io, .{ .reuse_address = true }) catch |err| { log.err("failed to bind public listener {s}: {s}", .{ cfg.public_addr_raw orelse "legacy --port", @errorName(err) }); return @errorCast(err); }; var server_future = try io.concurrent(runWsServer, .{ &ws_server, &listener, &public_context }); // An error anywhere before the normal server defers register (debug // listener setup, storage init) unwinds through here; without these the // accept loops would never be canceled and the runtime would block // joining them instead of exiting (the f0f84f1 gate caught exactly that: // --segment-fault=write:1:enospc hung for 300s where it must exit). // Disarmed once the normal defers take over so nothing runs twice. var early_server_cleanup = true; errdefer if (early_server_cleanup) { hub.listeners_ready.store(false, .release); _ = server_future.cancel(io); listener.deinit(io); ws_server.deinit(); }; var debug_ws_server: ?websocket.Server(server_mod.Handler) = if (debug_bind != null) try websocket.Server(server_mod.Handler).init(allocator, io, .{ .port = debug_bind.?.port, .address = "0.0.0.0", .max_conn = 256, .max_message_size = 1024, .websocket_shutdown_grace = cfg.client_drain_timeout, .http_shutdown_grace = cli.effectiveHttpShutdownDuration(cfg.shutdown_timeout), .compression = null, .handshake = .{ .max_size = 16 * 1024, .max_headers = 64 }, }) else null; var debug_listener: ?Io.net.Server = if (debug_bind) |bind| bind.address.listen(io, .{ .reuse_address = true }) catch |err| { log.err("failed to bind debug listener {s}: {s}", .{ cfg.debug_addr_raw.?, @errorName(err) }); return @errorCast(err); } else null; var debug_future: ?Io.Future(void) = if (debug_listener != null) try io.concurrent(runWsServer, .{ &debug_ws_server.?, &debug_listener.?, &debug_context }) else null; errdefer if (early_server_cleanup) { if (debug_future) |*future| _ = future.cancel(io); if (debug_listener) |*debug| debug.deinit(io); if (debug_ws_server) |*server| server.deinit(); }; hub.listeners_ready.store(true, .release); if (canonical_listeners) { log.info("public listener {s}; debug listener {s} (gates closed until storage is ready)", .{ cfg.public_addr_raw.?, cfg.debug_addr_raw orelse "disabled" }); } else { log.info("legacy combined listener on :{d} (gates closed until storage is ready)", .{cfg.port}); }
var archive = archive_mod.Archive.initWithStats(allocator, io, cfg.data_dir, &stats) catch |err| { logPersistenceFailure(cfg.data_dir, err); return @errorCast(err); }; var archive_closed = false; defer { archive.deinit(); } archive.max_segment_bytes = cfg.max_segment_bytes; var rule_store = try timestamp_rules.Store.open(allocator, io, cfg.data_dir); defer rule_store.deinit(); archive.stamper_ctx = &rule_store; archive.stamper = timestamp_rules.Store.stampOpaque; hot_tail.enableDurabilityGate(archive.committed_seq.load(.acquire)); log.info("archive at {s}/segments (next seq {d})", .{ cfg.data_dir, archive.next_seq }); hub.archive = &archive; var cold_reader = cold.ColdReader.init(allocator, io, cfg.subscribe_block_cache_bytes); defer cold_reader.deinit(); archive.manifest.setInvalidator(&cold_reader, cold.ColdReader.invalidateOpaque); defer archive.manifest.setInvalidator(null, null); hub.cold_reader = &cold_reader; var xrpc_api: xrpcapi.Api = .{ .allocator = allocator, .io = io, .archive = &archive, .segment_cache_max_age_s = cfg.segment_cache_max_age_s, .plan = cfg.plan_config, }; hub.xrpc = &xrpc_api;
var meta = try meta_store.Store.openWithStats(allocator, io, cfg.data_dir, &stats); defer meta.deinit(); try meta.migrateLegacy(io, cfg.data_dir); var repo_diagnostics = try repo_store.Store.init(allocator, io, &meta); defer repo_diagnostics.deinit(); hub.repo_status = &repo_diagnostics;
// Resolve the durable lifecycle before opening the gates. The listener // is already accepting but closed (serving=false, storage_ready=false // since before bind), so the old bootstrap race — a client observing the // default serving=true before the gate shut — cannot occur: serving only // ever flips true after this read, and stays false when mid-lifecycle. const stored_phase = blk: { var ps = phase_mod.Store.init(&meta); defer ps.deinit(io); break :blk try ps.read(io); }; if (stored_phase) |phase| stats.setOrchestratorPhase(switch (phase) { .bootstrap => .bootstrap, .merging => .merging, .steady_state => .steady_state, }); const mid_lifecycle = if (stored_phase) |ph| ph != .steady_state else false; const bootstrap_requested = cfg.do_backfill or cfg.backfill_max_repos > 0 or backfill_repos.len > 0;
var import_dir_buffer: [Io.Dir.max_path_bytes]u8 = undefined; const import_dir = cfg.timestamp_import_dir_arg orelse try std.fmt.bufPrint(&import_dir_buffer, "{s}/imports", .{cfg.data_dir}); try Io.Dir.cwd().createDirPath(io, import_dir); var import_scratch_buffer: [Io.Dir.max_path_bytes]u8 = undefined; const import_scratch = try std.fmt.bufPrint(&import_scratch_buffer, "{s}/timestamp-import-jobs", .{cfg.data_dir}); var import_job_store = timestamp_jobs.Store.init(allocator, &meta); var import_ready: std.atomic.Value(bool) = .init(false); var import_manager = try timestamp_manager.Manager.init(.{ .allocator = allocator, .io = io, .import_dir = import_dir, .scratch_root = import_scratch, .jobs = &import_job_store, .runner = .{ .allocator = allocator, .io = io, .scratch_root = import_scratch, .archive = &archive, .rules = &rule_store, .jobs = &import_job_store, .stats = &stats, }, .ready = &import_ready, }); var import_manager_err_owned = true; errdefer if (import_manager_err_owned) import_manager.deinit(); xrpc_api.import = &import_manager; xrpc_api.import_token = cfg.timestamp_import_token; xrpc_api.archive_api_key = cfg.archive_api_key; var key_ring: ?api_keys_mod.KeyRing = if (cfg.archive_api_keys_file.len > 0) api_keys_mod.KeyRing.init(allocator, cfg.archive_api_keys_file) else null; defer if (key_ring) |*ring| ring.deinit(); if (key_ring) |*ring| { xrpc_api.key_ring = ring; hub.key_ring = ring; log.info("archive key ring enabled: {s} (revocation = edit the file)", .{cfg.archive_api_keys_file}); }
// Killed getRepo downloads can leave incomplete scratch files. Bootstrap // owns backfill/, while steady repair has a separate root so a legitimate // live repair can never resurrect lifecycle state after cutover. const fetched_car = @import("internal/ingest/backfill/fetched_car.zig"); var bootstrap_scratch_buf: [Io.Dir.max_path_bytes]u8 = undefined; const bootstrap_scratch = try std.fmt.bufPrint(&bootstrap_scratch_buf, "{s}/backfill/repo-scratch", .{cfg.data_dir}); try fetched_car.resetScratch(io, bootstrap_scratch); var repair_scratch_root_buf: [Io.Dir.max_path_bytes]u8 = undefined; const repair_scratch_root = try std.fmt.bufPrint(&repair_scratch_root_buf, "{s}/repair-scratch", .{cfg.data_dir}); try fetched_car.resetScratch(io, repair_scratch_root); var live_repair_scratch_buf: [Io.Dir.max_path_bytes]u8 = undefined; const live_repair_scratch = try std.fmt.bufPrint(&live_repair_scratch_buf, "{s}/live", .{repair_scratch_root});
// The stores the request paths read are all published above; the release // store makes them visible to any handler that acquires storage_ready. // serving additionally requires the durable lifecycle to be steady — a // resumed bootstrap keeps it closed until cutover, exactly as before. hub.storage_ready.store(true, .release); if (!(mid_lifecycle or bootstrap_requested)) hub.serving.store(true, .release);
// Shutdown cleanup for the early-started server, registered here so it // runs BEFORE the storage defers above it in this function (defers are // LIFO): accept loops cancel and listeners close while the archive and // meta store are still alive. The early errdefer stands down now. early_server_cleanup = false; defer ws_server.deinit(); defer listener.deinit(io); defer _ = server_future.cancel(io); defer if (debug_ws_server) |*server| server.deinit(); defer if (debug_listener) |*debug| debug.deinit(io); defer if (debug_future) |*future| { _ = future.cancel(io); }; defer hub.listeners_ready.store(false, .release);
var cursor_store = cursor_mod.Store.init(&meta); defer cursor_store.deinit(io);
var verifier: ?verify.Verifier = if (cfg.do_verify) verify.Verifier.init(allocator, io, cfg.plc_url, &meta) else null; if (verifier) |*v| v.stats = &stats; defer if (verifier) |*v| v.deinit(); if (cfg.do_verify) log.info("signature verification on (plc: {s})", .{cfg.plc_url});
// lifecycle: a stored non-steady phase means a crash mid-bootstrap or // mid-merge — always resume it. otherwise only --backfill (full crawl) // or --max-backfill-repos (debug) start one; serving is 503-gated // until the lifecycle commits steady_state. if (mid_lifecycle or bootstrap_requested) { var lifecycle_run: LifecycleRun = .{ .allocator = allocator, .io = io, .options = .{ .data_dir = cfg.data_dir, .upstream = cfg.upstream, .relay_http = cfg.relay_http, .plc_url = cfg.plc_url, .max_repos = cfg.backfill_max_repos, .selected_repos = backfill_repos, .workers = cfg.backfill_workers, .batch_size = cfg.backfill_batch_size, .async_flush_workers = cfg.backfill_async_flush_workers, .skip_merge_discovery = cfg.skip_merge_discovery, .compaction_enabled = cfg.compaction_interval_ns != 0, }, .archive = &archive, .meta = &meta, .cursor_store = &cursor_store, .verifier = if (verifier) |*v| v else null, .stats = &stats, .trace_context = orchestrator_trace.context(), }; var lifecycle_future = try io.concurrent(runLifecycle, .{&lifecycle_run}); while (!shutdown_requested.load(.acquire) and !lifecycle_run.done.load(.acquire)) Io.sleep(io, Io.Duration.fromMilliseconds(25), .awake) catch |err| switch (err) { error.Canceled => {}, }; if (shutdown_requested.load(.acquire)) { log.info("shutdown requested during bootstrap; canceling lifecycle and flushing its durable prefix", .{}); hub.draining.store(true, .release); hub.listeners_ready.store(false, .release);
var listener_shutdowns: Io.Group = .init; errdefer listener_shutdowns.cancel(io); try listener_shutdowns.concurrent(io, cancelWsServer, .{ &server_future, io }); if (debug_future) |*future| try listener_shutdowns.concurrent(io, cancelWsServer, .{ future, io });
lifecycle_future.cancel(io); archive.close(); archive_closed = true; try listener_shutdowns.await(io); orchestrator_trace.succeed(); try shutdownTracing(tracing, cfg.shutdown_timeout); return; } lifecycle_future.await(io); if (lifecycle_run.result) |err| { logPersistenceFailure(cfg.data_dir, err); return @errorCast(err); } hub.serving.store(true, .release); } if (cfg.cursor == null) { cfg.cursor = try cursor_store.load(io); if (cfg.cursor) |c| log.info("resuming from persisted cursor {d}", .{c}); }
var steady_trace = observability.Span.start( allocator, "ingest/orchestrator", "runSteadyState", orchestrator_trace.context(), &.{}, ); defer steady_trace.deinit(); defer steady_trace.succeed(); try steady_trace.recordError(runSteadyState( allocator, io, cfg, tracing, orchestrator_trace, &steady_trace, &hub, &hot_tail, &stats, &archive, &archive_closed, &meta, &cursor_store, &verifier, &import_manager, &import_manager_err_owned, &import_ready, live_repair_scratch, &server_future, &debug_future, ));}
fn runSteadyState( allocator: std.mem.Allocator, io: Io, cfg: *cli.ServeConfig, tracing: *observability.Runtime, orchestrator_trace: *observability.Span, steady_trace: *observability.Span, hub: *server_mod.Hub, hot_tail: *tail_mod.Tail, stats: *metrics.Stats, archive: *archive_mod.Archive, archive_closed: *bool, meta: *meta_store.Store, cursor_store: *cursor_mod.Store, verifier: *?verify.Verifier, import_manager: *timestamp_manager.Manager, import_manager_err_owned: *bool, import_ready: *std.atomic.Value(bool), live_repair_scratch: []const u8, server_future: *Io.Future(void), debug_future: *?Io.Future(void),) !void { import_manager.setTraceContext(steady_trace.context());
// steady compactor: live tombstone set fed from the archive's append // hook, periodic + cap-triggered passes. interval 0 disables all of it // (run() returns immediately, no hook, no rebuild). var compactor = steady_mod.Compactor.init(allocator, io, archive, meta, cfg.data_dir); defer compactor.deinit(); compactor.interval_ns = cfg.compaction_interval_ns; compactor.tombstone_cap = cfg.compaction_tombstone_cap; compactor.rewrite_workers = cfg.compaction_rewrite_workers; compactor.stats = stats; compactor.trace_context = steady_trace.context(); // The archive outlives this frame; its hooks must not. defer { archive.on_append = null; archive.on_append_ctx = null; } if (cfg.compaction_interval_ns != 0) { try compactor.rebuild(); archive.on_append_ctx = &compactor; archive.on_append = steady_mod.Compactor.onAppend; log.info("steady compactor on (interval {d}ns, tombstone cap {d}, rewrite workers {d})", .{ cfg.compaction_interval_ns, cfg.compaction_tombstone_cap, cfg.compaction_rewrite_workers, }); } var compactor_future = try io.concurrent(steady_mod.Compactor.run, .{&compactor}); defer { compactor.requestStop(); _ = compactor_future.cancel(io); }
// one-shot legacy rebloom sweep: right-size fat per-block bloom regions // in already-sealed segments. Sealed-file rewriting is single-writer, so // the sweep refuses to coexist with an enabled compactor. if (cfg.rebloom_sweep and cfg.compaction_interval_ns != 0) log.err("--rebloom-sweep requires --compaction-interval=0 (sealed rewrites are single-writer); sweep not started", .{}); var rebloom_sweep: rebloom_mod.Sweep = .{ .allocator = allocator, .io = io, .archive = archive, .data_dir = cfg.data_dir, .stats = stats, .enabled = cfg.rebloom_sweep and cfg.compaction_interval_ns == 0, }; var rebloom_future = try io.concurrent(rebloom_mod.Sweep.run, .{&rebloom_sweep}); defer { rebloom_sweep.requestStop(); _ = rebloom_future.cancel(io); }
// failed-repo retry loop: periodic resync of repos whose backfill // download failed (or that discovery queued). appends ride the // archive writer mutex beside the live consumer. var retry = retry_mod.Retry.init(allocator, io, archive, meta, cfg.data_dir, cfg.relay_http); defer retry.deinit(); retry.interval_ns = cfg.retry_interval_ns; retry.workers = cfg.retry_workers; retry.host_workers = cfg.retry_host_workers; retry.max_delay_ns = cfg.retry_max_delay_ns; retry.stats = stats; retry.trace_context = steady_trace.context(); retry.shutdown = &shutdown_requested; if (cfg.retry_interval_ns != 0) log.info("failed-repo retry loop on (interval {d}ns, workers {d}, host workers {d}, max delay {d}ns)", .{ cfg.retry_interval_ns, cfg.retry_workers, cfg.retry_host_workers, cfg.retry_max_delay_ns, }); var retry_future = try io.concurrent(retry_mod.Retry.run, .{&retry}); defer { retry.requestStop(); _ = retry_future.cancel(io); }
var consumer: ingest.Consumer = .{ .allocator = allocator, .io = io, .tail = hot_tail, .upstream = cfg.upstream, .cursor = cfg.cursor, .to_stdout = cfg.to_stdout, .verifier = if (verifier.*) |*v| v else null, .cursor_store = cursor_store, .archive = archive, .stats = stats, .shutdown = &shutdown_requested, .slow_upstream = .{ .min_rate = cfg.upstream_slow_min_rate }, }; defer consumer.deinit(); defer { archive.on_durable = null; archive.on_durable_ctx = null; } // The archive's durability hook points into consumer. Keep the fallback // close above consumer's deinit in the unwind stack so every error path // flushes while the callback context is still alive. defer if (!archive_closed.*) { archive.close(); archive_closed.* = true; }; archive.on_durable_ctx = &consumer; archive.on_durable = ingest.Consumer.onArchiveDurable; var pipe: pipeline_mod.Pipeline = .{ .allocator = allocator, .io = io, .verifier = if (verifier.*) |*v| v else null, .stats = stats, .emit_ctx = &consumer, .emit = ingest.Consumer.emitVerified, .emit_repair = ingest.Consumer.emitRepair, .emit_cursor = ingest.Consumer.emitCursor, .cursor_watermark = cfg.cursor orelse 0, .trace_context = steady_trace.context(), }; var repair_coordinator: ?repair_mod.Coordinator = if (verifier.*) |*v| .{ .allocator = allocator, .io = io, .verifier = v, .relay_http = cfg.relay_http, .scratch_dir = live_repair_scratch, .stats = stats, .publish_ctx = &pipe, .publish = pipeline_mod.Pipeline.publishRepair, } else null; if (repair_coordinator) |*coordinator| { try coordinator.start(); pipe.repair = coordinator; } defer if (repair_coordinator) |*coordinator| coordinator.deinit(); try pipe.start(); defer pipe.stop(); consumer.pipeline = &pipe; defer import_manager.deinit(); import_manager_err_owned.* = false; import_ready.store(true, .release); if (try import_manager.resumeIncomplete()) log.info("resumed incomplete timestamp import", .{});
var consumer_done: std.atomic.Value(bool) = .init(false); var consumer_future = try io.concurrent(runConsumer, .{ &consumer, &consumer_done }); defer { if (consumer_future.any_future != null) _ = consumer_future.cancel(io) catch {}; } while (!shutdown_requested.load(.acquire) and !consumer_done.load(.acquire)) Io.sleep(io, Io.Duration.fromMilliseconds(25), .awake) catch |err| switch (err) { error.Canceled => {}, };
if (shutdown_requested.load(.acquire)) { log.info("shutdown requested; draining listeners, subscribers, and live pipeline", .{}); hub.draining.store(true, .release); hub.listeners_ready.store(false, .release);
// Upstream begins ordinary HTTP shutdown and WebSocket client drain // from the same canceled root context. Give each listener future to // one helper task so Future.cancel's single-threaded ownership rule is // preserved while storage/pipeline shutdown proceeds concurrently. var listener_shutdowns: Io.Group = .init; errdefer listener_shutdowns.cancel(io); try listener_shutdowns.concurrent(io, cancelWsServer, .{ server_future, io }); if (debug_future.*) |*future| try listener_shutdowns.concurrent(io, cancelWsServer, .{ future, io });
var shutdown_err: ?anyerror = null; _ = consumer_future.cancel(io) catch |err| switch (err) { error.Canceled => {}, else => shutdown_err = err, }; pipe.drainAndStop() catch |err| { shutdown_err = shutdown_err orelse err; }; // The durability hook points into consumer, so the final flush must // happen here while consumer is still alive—not in archive's older // outer defer after consumer teardown. archive.close(); archive_closed.* = true; try listener_shutdowns.await(io); if (shutdown_err) |err| return @errorCast(err); steady_trace.succeed(); orchestrator_trace.succeed(); try shutdownTracing(tracing, cfg.shutdown_timeout); return; } try consumer_future.await(io); try pipe.drainAndStop(); archive.close(); archive_closed.* = true; steady_trace.succeed(); orchestrator_trace.succeed(); try shutdownTracing(tracing, cfg.shutdown_timeout);}
fn shutdownTracing(runtime: *observability.Runtime, configured: Io.Duration) !void { const result = runtime.shutdown(shutdownTimeoutMillis(configured)); if (!result.isSuccess()) { log.err("tracer shutdown failed: {s}", .{@tagName(result)}); return error.TracerShutdownFailed; }}
fn shutdownTimeoutMillis(configured: Io.Duration) u64 { const timeout = cli.effectiveHttpShutdownDuration(configured); return @intCast(@divFloor(timeout.nanoseconds + std.time.ns_per_ms - 1, std.time.ns_per_ms));}
fn logPersistenceFailure(data_dir: []const u8, err: anyerror) void { if (err == error.NoSpaceLeft) { log.err("fatal persistence error: disk full under {s}; free space, then restart stream", .{data_dir}); }}
fn runConsumer(consumer: *ingest.Consumer, done: *std.atomic.Value(bool)) !void { defer done.store(true, .release); try consumer.run();}
fn runLifecycle(run: *LifecycleRun) void { defer run.done.store(true, .release); lifecycle.runTraced( run.allocator, run.io, run.options, run.archive, run.meta, run.cursor_store, run.verifier, run.stats, run.trace_context, ) catch |err| { run.result = err; };}
fn runWsServer(server: *websocket.Server(server_mod.Handler), listener: *Io.net.Server, context: *server_mod.ListenerContext) void { server.runIo(listener, context);}
fn cancelWsServer(future: *Io.Future(void), io: Io) void { _ = future.cancel(io);}