diff --git a/deploy/caddy/Dockerfile b/deploy/caddy/Dockerfile new file mode 100644 index 0000000..84ef1e6 --- /dev/null +++ b/deploy/caddy/Dockerfile @@ -0,0 +1,5 @@ +FROM caddy:2.10-builder AS builder +RUN xcaddy build --with github.com/mholt/caddy-ratelimit + +FROM caddy:2.10 +COPY --from=builder /usr/bin/caddy /usr/bin/caddy diff --git a/deploy/site/Caddyfile b/deploy/site/Caddyfile index 80f8841..04279c6 100644 --- a/deploy/site/Caddyfile +++ b/deploy/site/Caddyfile @@ -3,6 +3,52 @@ grafana.stream.waow.tech { } stream.89-167-122-160.sslip.io, stream.waow.tech { + # edge admission control (mholt/caddy-ratelimit, sliding window). + # per-IP keys grouped by /24 (v4) and /56 (v6) so address cycling + # inside one allocation shares a bucket. limits bound REQUEST rates + # only; concurrent-subscriber and cold-replay caps live in-process + # (the proxy cannot see what a connection costs the archive). + rate_limit { + # websocket connect attempts: a healthy consumer reconnects rarely; + # 30 per 5m per /24 tolerates flapping without permitting a spray + zone ws_connect { + match { + header Connection *Upgrade* + } + key {remote_host} + events 30 + window 5m + ipv4_prefix 24 + ipv6_prefix 56 + } + # keyed archive traffic: evelyn at 32 MB/s is ~150 getBlock/s + # (~214 KB avg block), so the ceiling sits well above a legit + # consumer; byte-rate fairness is the key ring's job, this only + # stops request-rate floods with a forged bearer header + zone xrpc_keyed { + match { + path /xrpc/* + header Authorization Bearer* + } + key {remote_host} + events 15000 + window 1m + } + # everything anonymous: static site, /metrics, /status, and + # unauthenticated xrpc probes (401 spray) share one modest bucket + zone anon_http { + match { + not header Connection *Upgrade* + not header Authorization Bearer* + } + key {remote_host} + events 240 + window 1m + ipv4_prefix 24 + ipv6_prefix 56 + } + } + root * /srv/site @home path / /index.html handle @home { diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 7db8dcf..c10fb2c 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -217,3 +217,28 @@ no restart). the fleet key remains valid and unmetered. rationale: metering protects compaction/live-delivery from uncapped archive drains on a ~200 MB/s volume; revocation avoids rotating the fleet credential per consumer. + +## in-process admission caps (2026-08-21, extension) + +upstream keeps connection admission at the edge: docs/README.md §7 puts +per-IP limits on subscriber counts, live-tail events, and segment +downloads at the CDN and origin proxy, explicitly calling in-process +limits "an implementation detail of the Bluesky-hosted instance" — the +binary itself only bounds per-request cost (filter caps, list limits, +byte-budgeted caches) plus one per-IP fixed-window limiter guarding +web repo actions. stream fronts with a single Caddy (no CDN), so two +caps that need process knowledge live in-process as opt-in extensions +(`--max-subscribers`, `--max-cold-readers`; 0/unset = unbounded, +preserving upstream behavior exactly): + +- max_subscribers: global subscriber ceiling; excess upgrades close + 1013 ("at capacity") before a Subscriber is allocated. +- max_cold_readers: concurrent cold-replay bound. only the process + knows a cursor is about to cost archive reads; over-limit passes + return .contended (20 ms pause + retry — backpressure, not an error; + never counts toward the rotation-seam stall disconnect). + +request-rate admission (per-IP/per-subnet) stays at the proxy layer per +upstream doctrine: deploy/caddy builds Caddy with mholt/caddy-ratelimit +and deploy/site/Caddyfile defines the zones (ws connects, keyed xrpc, +anonymous http). diff --git a/src/internal/runtime/cli.zig b/src/internal/runtime/cli.zig index a365187..3c53aff 100644 --- a/src/internal/runtime/cli.zig +++ b/src/internal/runtime/cli.zig @@ -206,6 +206,10 @@ pub const ServeOptions = struct { /// live tail and dictionary open, archive behind a token). archive_api_key: []const u8 = "", archive_api_keys_file: []const u8 = "", + /// stream extensions: in-process admission caps (0/absent = unbounded). + /// upstream keeps these at its CDN/proxy layer (docs/README.md §7). + max_subscribers: []const u8 = "", + max_cold_readers: []const u8 = "", timestamp_import_dir: []const u8 = "", pub const overrides = .{ @@ -324,6 +328,23 @@ pub fn effectiveSubscribeSlowWindowNanoseconds(raw: []const u8) !u64 { const configured = try parseDurationNanoseconds(raw); return if (configured == 0) default else configured; } +/// stream-extension admission caps (--max-subscribers/--max-cold-readers): +/// non-positive means unbounded, matching the pinned surface's sentinel +/// convention for optional integer controls. +fn parseAdmissionCap(raw: []const u8) !u32 { + const configured = try std.fmt.parseInt(i64, raw, 10); + if (configured <= 0) return 0; + return std.math.cast(u32, configured) orelse error.Overflow; +} + +test "admission caps use the non-positive unbounded sentinel" { + try std.testing.expectEqual(@as(u32, 0), try parseAdmissionCap("0")); + try std.testing.expectEqual(@as(u32, 0), try parseAdmissionCap("-1")); + try std.testing.expectEqual(@as(u32, 500), try parseAdmissionCap("500")); + try std.testing.expectError(error.Overflow, parseAdmissionCap("4294967296")); + try std.testing.expectError(error.InvalidCharacter, parseAdmissionCap("many")); +} + test "subscribe slow window uses the upstream non-positive sentinel" { const default = server_mod.default_slow_window_ns; try std.testing.expectEqual(default, try effectiveSubscribeSlowWindowNanoseconds("0")); @@ -471,6 +492,8 @@ pub const ServeConfig = struct { /// live tail and dictionary open, archive behind a token). archive_api_key: []const u8 = "", archive_api_keys_file: []const u8 = "", + max_subscribers: u32 = 0, + max_cold_readers: u32 = 0, timestamp_import_dir_arg: ?[]const u8 = null, repo_action_rate_limits: bool = true, store_fault_prefix: ?[]const u8 = null, @@ -679,6 +702,10 @@ pub fn parseServe( cfg.plan_config.whole_segment_threshold = try std.fmt.parseFloat(f64, arg["--plan-whole-segment-threshold=".len..]); } else if (std.mem.startsWith(u8, arg, "--archive-api-key=")) { cfg.archive_api_key = arg["--archive-api-key=".len..]; + } else if (std.mem.startsWith(u8, arg, "--max-subscribers=")) { + cfg.max_subscribers = try parseAdmissionCap(arg["--max-subscribers=".len..]); + } else if (std.mem.startsWith(u8, arg, "--max-cold-readers=")) { + cfg.max_cold_readers = try parseAdmissionCap(arg["--max-cold-readers=".len..]); } else if (std.mem.startsWith(u8, arg, "--archive-api-keys-file=")) { cfg.archive_api_keys_file = arg["--archive-api-keys-file=".len..]; } else if (std.mem.startsWith(u8, arg, "--timestamp-import-token=")) { diff --git a/src/internal/runtime/environment.zig b/src/internal/runtime/environment.zig index 4d96fee..4748185 100644 --- a/src/internal/runtime/environment.zig +++ b/src/internal/runtime/environment.zig @@ -51,6 +51,8 @@ pub const mappings = [_]Mapping{ .{ .env = "JETSTREAM_TIMESTAMP_IMPORT_DIR", .flag = "timestamp-import-dir" }, .{ .env = "JETSTREAM_ARCHIVE_API_KEY", .flag = "archive-api-key" }, .{ .env = "JETSTREAM_ARCHIVE_API_KEYS_FILE", .flag = "archive-api-keys-file" }, + .{ .env = "JETSTREAM_MAX_SUBSCRIBERS", .flag = "max-subscribers" }, + .{ .env = "JETSTREAM_MAX_COLD_READERS", .flag = "max-cold-readers" }, }; pub fn processEntries(allocator: std.mem.Allocator) ![]const []const u8 { @@ -134,7 +136,7 @@ fn isKnown(key: []const u8) bool { } test "unknown Jetstream variables are sorted, deduplicated, and scoped" { - try std.testing.expectEqual(@as(usize, 40), mappings.len); + try std.testing.expectEqual(@as(usize, 42), mappings.len); const entries = [_][]const u8{ "JETSTREAM_ZZZ=1", "JETSTREAM_ADDR=127.0.0.1:0", diff --git a/src/internal/runtime/metrics.zig b/src/internal/runtime/metrics.zig index e005826..d18b9d8 100644 --- a/src/internal/runtime/metrics.zig +++ b/src/internal/runtime/metrics.zig @@ -248,6 +248,8 @@ pub const Stats = struct { subscribe_hot_reads_total: Counter = .init(0), subscribe_cold_reads_total: Counter = .init(0), subscribe_cold_stalls_total: Counter = .init(0), + subscribe_cold_contended_total: Counter = .init(0), + subscribe_capacity_rejects_total: Counter = .init(0), subscribe_adversarial_drops_total: Counter = .init(0), subscribe_clean_disconnects_total: Counter = .init(0), subscribe_options_updates_total: Counter = .init(0), @@ -993,6 +995,10 @@ pub fn format( \\jetstream_subscribe_cold_reads_total {d} \\# TYPE stream_subscribe_cold_stalls_total counter \\stream_subscribe_cold_stalls_total {d} + \\# TYPE stream_subscribe_cold_contended_total counter + \\stream_subscribe_cold_contended_total {d} + \\# TYPE stream_subscribe_capacity_rejects_total counter + \\stream_subscribe_capacity_rejects_total {d} \\# TYPE jetstream_subscribe_adversarial_drops_total counter \\jetstream_subscribe_adversarial_drops_total {d} \\# TYPE jetstream_subscribe_clean_disconnects_total counter @@ -1007,6 +1013,8 @@ pub fn format( stats.subscribe_hot_reads_total.load(.monotonic), stats.subscribe_cold_reads_total.load(.monotonic), stats.subscribe_cold_stalls_total.load(.monotonic), + stats.subscribe_cold_contended_total.load(.monotonic), + stats.subscribe_capacity_rejects_total.load(.monotonic), stats.subscribe_adversarial_drops_total.load(.monotonic), stats.subscribe_clean_disconnects_total.load(.monotonic), stats.subscribe_options_updates_total.load(.monotonic), diff --git a/src/internal/serve/server.zig b/src/internal/serve/server.zig index 623626d..006ebfc 100644 --- a/src/internal/serve/server.zig +++ b/src/internal/serve/server.zig @@ -108,6 +108,13 @@ pub const Hub = struct { repo_action_rate_limits: bool = true, repo_action_limiter: repo_action_limiter.Limiter = .{}, active: std.atomic.Value(u32) = .init(0), + /// admission caps, 0 = unbounded. a recorded divergence: upstream keeps + /// connection admission at its CDN/proxy layer (docs/README.md §7); we + /// have no CDN, and only the process knows a cursor is about to cost + /// archive reads (docs/semantic-parity.md). + max_subscribers: u32 = 0, + max_cold_readers: u32 = 0, + cold_active: std.atomic.Value(u32) = .init(0), bootstrap_enabled: bool = false, compaction_enabled: bool = false, retry_enabled: bool = false, @@ -142,6 +149,23 @@ pub const Hub = struct { /// per-key archive metering ring, rendered into /metrics when set key_ring: ?*api_keys.KeyRing = null, + /// Bounded cold-replay concurrency. fetchAdd-then-undo keeps the pair + /// wait-free; transient overshoot is bounded by the number of racing + /// subscriber threads and never admits past the cap. + fn tryAcquireColdSlot(self: *Hub) bool { + if (self.max_cold_readers == 0) return true; + if (self.cold_active.fetchAdd(1, .acquire) >= self.max_cold_readers) { + _ = self.cold_active.fetchSub(1, .release); + return false; + } + return true; + } + + fn releaseColdSlot(self: *Hub) void { + if (self.max_cold_readers == 0) return; + _ = self.cold_active.fetchSub(1, .release); + } + pub fn deinit(self: *Hub) void { self.repo_action_limiter.deinit(self.allocator); } @@ -255,6 +279,14 @@ const Subscriber = struct { continue; }, .failed => return self.noteExit(.cold_failed), + .contended => { + // cold-reader slots full: brief pause and retry — + // deliberate backpressure, not a stall or an error + self.cold_stall_passes = 0; + self.hub.io.sleep(Io.Duration.fromMilliseconds(20), .awake) catch + return self.noteExit(.cold_failed); + continue; + }, } } @@ -381,7 +413,7 @@ const Subscriber = struct { return if (self.conn.compression != null) .deflate else .none; } - const ColdPass = enum { progressed, empty, failed }; + const ColdPass = enum { progressed, empty, failed, contended }; /// One stateless cold read from cursor_seq (upstream Tail.cold). No stop /// floor: overlapping the hot range is harmless because the cursor is @@ -391,6 +423,11 @@ const Subscriber = struct { const archive = hub.archive orelse return .failed; const reader = hub.cold_reader orelse return .failed; const from_seq = self.cursor_seq; + if (!hub.tryAcquireColdSlot()) { + _ = hub.stats.subscribe_cold_contended_total.fetchAdd(1, .monotonic); + return .contended; + } + defer hub.releaseColdSlot(); _ = hub.stats.subscribe_cold_reads_total.fetchAdd(1, .monotonic); const batch = reader.readSeqBatch( archive, @@ -751,6 +788,12 @@ pub const Handler = struct { fn startSubscriber(self: *Handler) !void { const hub = self.hub; + if (hub.max_subscribers > 0 and hub.active.load(.monotonic) >= hub.max_subscribers) { + _ = hub.stats.subscribe_capacity_rejects_total.fetchAdd(1, .monotonic); + self.conn.writeClose(.{ .code = 1013, .reason = "at capacity; retry later" }) catch {}; + self.noteCleanDisconnect(); + return error.AtCapacity; + } const plan = self.cursor_plan; const start_idx = plan.start_idx orelse hub.tail.publishedTip(); // cursor_seq == 0 means seq-unknown: live and time-indexed plans are @@ -2477,3 +2520,33 @@ test "accounts HTTP route verifies real archive state and rate limits by port-fr try testing.expect(std.mem.indexOf(u8, trusted.body, "repo action rate limit exceeded") == null); try testing.expectEqual(@as(u32, 15), identity_calls.load(.monotonic)); } + +test "cold-reader slots: cap admits, saturates, releases; zero is unbounded" { + const testing = std.testing; + var threaded: Io.Threaded = .init(testing.allocator, .{}); + defer threaded.deinit(); + const io = threaded.io(); + var tail = tail_mod.Tail.init(testing.allocator, io, 1 << 20); + defer tail.deinit(); + var stats: metrics.Stats = .{}; + var hub: Hub = .{ + .allocator = testing.allocator, + .io = io, + .tail = &tail, + .stats = &stats, + .max_cold_readers = 2, + }; + try testing.expect(hub.tryAcquireColdSlot()); + try testing.expect(hub.tryAcquireColdSlot()); + try testing.expect(!hub.tryAcquireColdSlot()); // saturated + hub.releaseColdSlot(); + try testing.expect(hub.tryAcquireColdSlot()); // freed slot readmits + hub.releaseColdSlot(); + hub.releaseColdSlot(); + try testing.expectEqual(@as(u32, 0), hub.cold_active.load(.monotonic)); + + hub.max_cold_readers = 0; // unbounded: never blocks, never counts + try testing.expect(hub.tryAcquireColdSlot()); + hub.releaseColdSlot(); + try testing.expectEqual(@as(u32, 0), hub.cold_active.load(.monotonic)); +} diff --git a/src/main.zig b/src/main.zig index cc5a5e5..362179e 100644 --- a/src/main.zig +++ b/src/main.zig @@ -518,6 +518,8 @@ pub fn main(init: std.process.Init.Minimal) !void { .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}", .{ diff --git a/tests/environment_contract.py b/tests/environment_contract.py index a2a68cb..0b00805 100644 --- a/tests/environment_contract.py +++ b/tests/environment_contract.py @@ -70,6 +70,8 @@ def exact_pinned_map_case() -> None: # Named, metered, file-revocable keys beside the fleet key — the # in-process approximation of the hosted gateway (semantic-parity.md) "JETSTREAM_ARCHIVE_API_KEYS_FILE": "archive-api-keys-file", + "JETSTREAM_MAX_SUBSCRIBERS": "max-subscribers", + "JETSTREAM_MAX_COLD_READERS": "max-cold-readers", } # Upstream variables Stream knowingly has NOT ported: the PDS-direct # fleet bootstrap (upstream specs/notes/2026-08-03-pds-direct-backfill-