jetstream client
atproto jetstream client
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145//! loopback e2e: two in-process fake v2 instances drive the REAL unified//! client — websocket.zig live tail, HttpTransport archive sweep — through//! the failure modes that motivated multi-host failover (docs/failover.md)://! a primary that dies mid-live, and a primary that goes connected-but-//! silent (the 2026-08-16 bsky.network front wedge). the fixtures serve//! the actual wire: listSegments/planSnapshot JSON, raw-block zstd frames//! over getBlock, and subscribeEvents lexicon frames over a hand-rolled//! websocket upgrade.
const std = @import("std");const zat = @import("zat");const archive = @import("archive_backfill.zig");const client_mod = @import("client.zig");const livedecode = @import("livedecode.zig");
const Io = std.Io;const Allocator = std.mem.Allocator;const testing = std.testing;
/// witnessed base: 2026-07-13T00:00:00Z (the livedecode fixture's day)const epoch_s: i64 = 1783900800;
fn witnessedUs(offset_s: i64) i64 { return (epoch_s + offset_s) * std.time.us_per_s;}
const Row = struct { seq: u64, /// seconds past epoch_s (0..59: rendered into one RFC3339 minute) t: i64, did: []const u8 = "did:plc:loopback", collection: []const u8 = "app.bsky.feed.post",};
const Instance = struct { io: Io, allocator: Allocator, listener: Io.net.Server, port: u16, sealed: []const Row, live: []const Row, /// close the websocket after emitting `live` (a server that dies mid-tail) live_then_close: bool = false, /// accept the websocket, emit nothing, hold it open (the silent wedge) silent_ws: bool = false, /// after this many websocket sessions, close every new connection /// immediately (the host is down) max_ws_sessions: u32 = std.math.maxInt(u32), ws_sessions: u32 = 0, wire_cursors: [32]?u64 = @splat(null), archive_calls: usize = 0, reject_archive: bool = false, reject_ws_count: u32 = 0, first_session_rows: ?usize = null, retention_floor_us: ?u64 = null, cursor_too_old: bool = false, page_size: usize = 0, page_rows: ?[]const Row = null, plan_calls: usize = 0, last_before_seq: ?u64 = null, first_tip: ?u64 = null, stalled_plan: bool = false, archive_fault: enum { none, fetch, decode, segment } = .none, compression: enum { none, valid, rotate, same, invalid } = .none, bad_compressed_first: bool = false, include_cursor_boundary: bool = false, dict_calls: usize = 0, compressed_sessions: usize = 0, reject_plan: bool = false, retry_block_calls: ?usize = null, live_authorized: bool = false, archive_authorized: bool = false, archive_missing_auth: bool = false,
fn init(io: Io, allocator: Allocator, sealed: []const Row, live: []const Row) !Instance { var listener = try (Io.net.IpAddress{ .ip4 = .loopback(0) }).listen(io, .{ .reuse_address = true }); errdefer listener.deinit(io); return .{ .io = io, .allocator = allocator, .listener = listener, .port = listener.socket.address.getPort(), .sealed = sealed, .live = live, }; }
fn deinit(self: *Instance) void { self.listener.deinit(self.io); }
fn hostUrl(self: *const Instance, buf: []u8) []const u8 { return std.fmt.bufPrint(buf, "http://127.0.0.1:{d}", .{self.port}) catch unreachable; }
/// sequential accept loop: every handler terminates promptly — the /// client under test closes its previous socket before reconnecting, /// so a held websocket never blocks the next accept for long. exits /// via the /quit poison pill so no thread is ever inside accept when /// the listener closes (Io.Threaded panics on that race in Debug). fn serve(self: *Instance) !void { while (true) { var stream = self.listener.accept(self.io) catch return; const stop_requested = self.handleConn(&stream) catch false; stream.close(self.io); if (stop_requested) return; } }
/// unblock and stop the serve loop from the test thread fn quit(self: *Instance) void { var transport = zat.HttpTransport.init(self.io, self.allocator); defer transport.deinit(); // A shutdown request must not retry after the server has exited. transport.dead_connection_attempts = 1; var url_buf: [64]u8 = undefined; const url = std.fmt.bufPrint(&url_buf, "http://127.0.0.1:{d}/quit", .{self.port}) catch unreachable; var response = transport.fetch(.{ .url = url, .deadline_ns = std.time.ns_per_s }) catch return; response.deinit(self.allocator); }
fn handleConn(self: *Instance, stream: *Io.net.Stream) !bool { var rbuf: [8192]u8 = undefined; var wbuf: [16384]u8 = undefined; var reader = stream.reader(self.io, &rbuf); var writer = stream.writer(self.io, &wbuf); const r = &reader.interface; const w = &writer.interface;
var head: [4096]u8 = undefined; var n: usize = 0; while (std.mem.indexOf(u8, head[0..n], "\r\n\r\n") == null) { if (n == head.len) return error.HeadTooLarge; head[n] = try r.takeByte(); n += 1; } const request = head[0..n]; const line_end = std.mem.indexOf(u8, request, "\r\n") orelse return error.BadRequest; var line = std.mem.tokenizeScalar(u8, request[0..line_end], ' '); const method = line.next() orelse return error.BadRequest; const target = line.next() orelse return error.BadRequest;
var arena_state = std.heap.ArenaAllocator.init(self.allocator); defer arena_state.deinit(); const arena = arena_state.allocator();
if (std.mem.startsWith(u8, target, "/quit")) { try respondBytes(w, "text/plain", ""); return true; } if (std.mem.indexOf(u8, target, "subscribeEvents") != null) { self.live_authorized = headerValue(request, "Authorization") != null; try self.handleWebsocket(request, target, r, w); return false; } if (std.mem.indexOf(u8, target, "getZstdDictionary") != null) { self.dict_calls += 1; if (self.compression == .none) { try w.writeAll("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); } else if (self.compression == .invalid) { try respondBytes(w, "application/octet-stream", "invalid dictionary"); } else { var dict = @embedFile("testdata/live.dict").*; if (self.compression == .rotate and self.dict_calls > 1) { const old = std.mem.readInt(u32, dict[4..8], .little); std.mem.writeInt(u32, dict[4..8], old + 1, .little); } try respondBytes(w, "application/octet-stream", &dict); } return false; } self.archive_calls += 1; if (self.reject_archive) { try w.writeAll("HTTP/1.1 401 Unauthorized\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); return false; } const authorized = if (headerValue(request, "Authorization")) |value| std.mem.eql(u8, value, "Bearer test-archive-key") else false; self.archive_authorized = self.archive_authorized or authorized; self.archive_missing_auth = self.archive_missing_auth or !authorized; if (std.mem.indexOf(u8, target, "listSegments") != null) { try respondJson(w, try self.segmentListJson(arena)); return false; } if (std.mem.eql(u8, method, "POST") and std.mem.indexOf(u8, target, "planSnapshot") != null) { const body = try arena.alloc(u8, contentLength(request) orelse 0); try r.readSliceAll(body); self.plan_calls += 1; if (self.reject_plan) { try w.writeAll("HTTP/1.1 400 Bad Request\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); return false; } if (self.page_size > 0) { const parsed = try std.json.parseFromSlice(std.json.Value, arena, body, .{}); const fields = parsed.value.object; const after: u64 = if (fields.get("afterSeq")) |v| @intCast(v.integer) else 0; self.last_before_seq = if (fields.get("beforeSeq")) |v| @intCast(v.integer) else null; var start: usize = 0; while (start < self.sealed.len and self.sealed[start].seq <= after) : (start += 1) {} var end = @min(start + self.page_size, self.sealed.len); if (self.last_before_seq) |before| { while (end > start and self.sealed[end - 1].seq > before) : (end -= 1) {} } self.page_rows = self.sealed[start..end]; } try respondJson(w, try self.planJson(arena)); return false; } if (std.mem.indexOf(u8, target, "getSegment") != null and self.archive_fault == .segment) { try respondBytes(w, "application/octet-stream", "invalid segment header"); return false; } if (std.mem.indexOf(u8, target, "getBlock") != null) { if (self.retry_block_calls) |count| { self.retry_block_calls = count + 1; try w.writeAll("HTTP/1.1 503 Service Unavailable\r\nRetry-After: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); return false; } if (self.archive_fault != .none) { const second = std.mem.indexOf(u8, target, "segment=second") != null; const block: u64 = if (std.mem.indexOf(u8, target, "blockIndex=1") != null) 1 else if (std.mem.indexOf(u8, target, "blockIndex=2") != null) 2 else 0; if (!second and block == 1) { if (self.archive_fault == .fetch) { try w.writeAll("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); } else try respondBytes(w, "application/octet-stream", "invalid zstd frame"); return false; } var fixture = self.*; fixture.sealed = &.{.{ .seq = if (second) 100 else block + 1, .t = if (second) 10 else @intCast(block + 1) }}; fixture.page_rows = null; try respondBytes(w, "application/octet-stream", try fixture.blockZstd(arena)); return false; } try respondBytes(w, "application/octet-stream", try self.blockZstd(arena)); return false; } try w.writeAll("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"); try w.flush(); return false; }
fn segmentListJson(self: *const Instance, arena: Allocator) ![]const u8 { if (self.sealed.len == 0) return "{\"segments\": []}"; return std.fmt.allocPrint(arena, \\{{"segments": [{{"name": "seg_0000000001", "minSeq": {d}, "maxSeq": {d}, "minWitnessedAt": {d}, "maxWitnessedAt": {d}}}]}} , .{ self.sealed[0].seq, self.sealed[self.sealed.len - 1].seq, witnessedUs(self.sealed[0].t), witnessedUs(self.sealed[self.sealed.len - 1].t), }); }
fn planJson(self: *const Instance, arena: Allocator) ![]const u8 { if (self.archive_fault == .segment) return \\{"plannedThroughSeq":100,"sealedTipSeq":100,"segments":[ \\{"name":"first","index":0,"mode":"segment"}, \\{"name":"second","index":1,"mode":"blocks","blocks":[{"first":0,"last":0}]}]} ; if (self.archive_fault != .none) return \\{"plannedThroughSeq":100,"sealedTipSeq":100,"segments":[ \\{"name":"first","index":0,"mode":"blocks","blocks":[{"first":0,"last":2}]}, \\{"name":"second","index":1,"mode":"blocks","blocks":[{"first":0,"last":0}]}]} ; if (self.sealed.len == 0) return "{\"plannedThroughSeq\": 0, \"sealedTipSeq\": 0, \"segments\": []}"; const tip = self.last_before_seq orelse (if (self.plan_calls == 1) self.first_tip else null) orelse self.sealed[self.sealed.len - 1].seq; if (self.stalled_plan) return std.fmt.allocPrint(arena, "{{\"plannedThroughSeq\":0,\"sealedTipSeq\":{d},\"segments\":[]}}", .{tip}); const rows = self.page_rows orelse self.sealed; if (rows.len == 0) return std.fmt.allocPrint(arena, "{{\"plannedThroughSeq\":{d},\"sealedTipSeq\":{d},\"segments\":[]}}", .{ tip, tip }); return std.fmt.allocPrint(arena, \\{{"plannedThroughSeq": {d}, "sealedTipSeq": {d}, "segments": [ \\ {{"name": "seg_0000000001", "index": 0, "checksum": "0", "minSeq": {d}, "maxSeq": {d}, \\ "mode": "blocks", "blocks": [{{"first": 0, "last": 0}}]}}]}} , .{ rows[rows.len - 1].seq, tip, rows[0].seq, rows[rows.len - 1].seq }); }
/// the sealed rows as one columnar jss block wrapped in a raw-block /// zstd frame (no compressor needed: magic, single-segment FHD with /// 4-byte FCS, one raw block) fn blockZstd(self: *const Instance, arena: Allocator) ![]u8 { const source_rows = self.page_rows orelse self.sealed; const rows = try arena.alloc(BlockRow, source_rows.len); for (source_rows, rows) |row, *out| out.* = .{ .seq = row.seq, .witnessed_at = witnessedUs(row.t), .kind = 1, // create .collection = row.collection, .did = row.did, .rkey = "rkey1", .rev = "rev1", .payload = &archive.test_record_cbor, }; const block = try archive.buildTestBlock(arena, rows); var out: std.Io.Writer.Allocating = .init(arena); const w = &out.writer; try w.writeAll(&.{ 0x28, 0xb5, 0x2f, 0xfd, 0xa0 }); // magic + FHD try w.writeInt(u32, @intCast(block.len), .little); // frame content size try w.writeInt(u24, @intCast( (block.len << 3) | 0b001, ), .little); // last raw block try w.writeAll(block); return out.written(); }
const BlockRow = @typeInfo(@typeInfo(@TypeOf(archive.buildTestBlock)).@"fn".param_types[1].?).pointer.child;
fn handleWebsocket(self: *Instance, request: []const u8, target: []const u8, r: *Io.Reader, w: *Io.Writer) !void { if (self.ws_sessions < self.wire_cursors.len) self.wire_cursors[self.ws_sessions] = queryCursor(target); self.ws_sessions += 1; if (self.ws_sessions <= self.reject_ws_count) return; if (self.cursor_too_old) { const body = "{\"error\":\"CursorTooOld\"}"; try w.print("HTTP/1.1 400 Bad Request\r\nContent-Type: application/json\r\nContent-Length: {d}\r\nConnection: close\r\n\r\n{s}", .{ body.len, body }); try w.flush(); return; } if (self.ws_sessions > self.max_ws_sessions) return; // down: slam the door
const compressed = std.mem.indexOf(u8, target, "zstdDictionary=") != null; if (compressed) self.compressed_sessions += 1; if (compressed and self.ws_sessions == 1 and (self.compression == .rotate or self.compression == .same)) { const body = "{\"error\":\"UnknownZstdDictionary\"}"; try w.print("HTTP/1.1 400 Bad Request\r\nContent-Type: application/json\r\nContent-Length: {d}\r\nConnection: close\r\n\r\n{s}", .{ body.len, body }); try w.flush(); return; } const key = headerValue(request, "Sec-WebSocket-Key") orelse return error.NoKey; var sha = std.crypto.hash.Sha1.init(.{}); sha.update(key); sha.update("258EAFA5-E914-47DA-95CA-C5AB0DC85B11"); var digest: [20]u8 = undefined; sha.final(&digest); var accept_buf: [28]u8 = undefined; const accept = std.base64.standard.Encoder.encode(&accept_buf, &digest); try w.print("HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {s}\r\nSec-WebSocket-Protocol: xrpc.v1.json\r\n\r\n", .{accept}); try w.flush();
const cursor = queryCursor(target); const clamped = if (cursor) |c| if (self.retention_floor_us) |floor| c >= 1_000_000_000_000_000 and c < floor else false else false; if (clamped) try writeTextFrame(w, \\{"$type":"message","payload":{"$type":"network.bsky.jetstream.subscribeEvents#info","name":"OutdatedCursor","message":"resumed at retention floor"}} ); if (compressed and self.bad_compressed_first) try writeFrame(w, "invalid zstd", 0x82); if (!self.silent_ws) { const rows = if (self.ws_sessions == 1 and self.first_session_rows != null) self.live[0..self.first_session_rows.?] else self.live; for (rows) |row| { if (self.retention_floor_us) |floor| if (witnessedUs(row.t) < floor) continue; if (!self.include_cursor_boundary) { if (cursor) |c| { if (c >= 1_000_000_000_000_000) { if (witnessedUs(row.t) < c) continue; } else if (row.seq < c) continue; } } var frame_buf: [512]u8 = undefined; const frame = try std.fmt.bufPrint(&frame_buf, \\{{"$type":"message","payload":{{"$type":"network.bsky.jetstream.subscribeEvents#commit","collection":"{s}","did":"{s}","operation":"create","record":{{"$type":"app.bsky.feed.post","text":"hi"}},"rev":"rev1","rkey":"rkey1","seq":{d},"time":"2026-07-13T00:00:{d:0>2}.000000Z"}}}} , .{ row.collection, row.did, row.seq, @as(u64, @intCast(row.t)) }); if (compressed) { var buffer: [600]u8 = undefined; var output: Io.Writer = .fixed(&buffer); try output.writeAll(&.{ 0x28, 0xb5, 0x2f, 0xfd, 0xa3 }); const id = std.mem.readInt(u32, @embedFile("testdata/live.dict")[4..8], .little) + @as(u32, if (self.compression == .rotate and self.dict_calls > 1) 1 else 0); try output.writeInt(u32, id, .little); try output.writeInt(u32, @intCast(frame.len), .little); try output.writeInt(u24, @intCast((frame.len << 3) | 1), .little); try output.writeAll(frame); try writeFrame(w, output.buffered(), 0x82); } else try writeTextFrame(w, frame); } if (self.live_then_close or (self.ws_sessions == 1 and self.first_session_rows != null)) return; } // hold the tail open; the client closing its socket (stop, stall, // rotation, reconnect) ends this read var discard: [256]u8 = undefined; while (true) { const read = r.readSliceShort(&discard) catch return; if (read == 0) return; } }};
fn writeTextFrame(w: *Io.Writer, payload: []const u8) !void { return writeFrame(w, payload, 0x81);}
fn writeFrame(w: *Io.Writer, payload: []const u8, opcode: u8) !void { if (payload.len < 126) { try w.writeAll(&.{ opcode, @as(u8, @intCast(payload.len)) }); } else { try w.writeAll(&.{ opcode, 126 }); try w.writeInt(u16, @intCast(payload.len), .big); } try w.writeAll(payload); try w.flush();}
fn headerValue(request: []const u8, name: []const u8) ?[]const u8 { var lines = std.mem.tokenizeSequence(u8, request, "\r\n"); while (lines.next()) |l| { if (l.len > name.len + 2 and std.ascii.startsWithIgnoreCase(l, name) and l[name.len] == ':') return std.mem.trim(u8, l[name.len + 1 ..], " "); } return null;}
fn contentLength(request: []const u8) ?usize { const v = headerValue(request, "Content-Length") orelse return null; return std.fmt.parseInt(usize, v, 10) catch null;}
fn queryCursor(target: []const u8) ?u64 { const idx = std.mem.indexOf(u8, target, "cursor=") orelse return null; var end = idx + "cursor=".len; while (end < target.len and std.ascii.isDigit(target[end])) end += 1; return std.fmt.parseInt(u64, target[idx + "cursor=".len .. end], 10) catch null;}
fn respondJson(w: *Io.Writer, body: []const u8) !void { try respondBytes(w, "application/json", body);}
fn respondBytes(w: *Io.Writer, content_type: []const u8, body: []const u8) !void { try w.print("HTTP/1.1 200 OK\r\nContent-Type: {s}\r\nContent-Length: {d}\r\nConnection: close\r\n\r\n", .{ content_type, body.len }); try w.writeAll(body); try w.flush();}
const Collected = struct { seq: u64, t: i64 };
fn CollectingHandler(comptime stop_after: usize) type { return struct { events: [32]Collected = undefined, len: usize = 0,
/// commit rows missing canonical record bytes (both halves must /// populate record_cbor for create/update — upstream RecordCBOR) missing_record_cbor: usize = 0,
pub fn onEvent(self: *@This(), event: livedecode.Event) bool { if (self.len == self.events.len) return false; if (event.payload == .commit and event.payload.commit.operation != .delete and event.payload.commit.record_cbor.len == 0) self.missing_record_cbor += 1; self.events[self.len] = .{ .seq = event.seq, .t = @divFloor(event.time_us, std.time.us_per_s) - epoch_s, }; self.len += 1; return self.len < stop_after; }
fn slice(self: *const @This()) []const Collected { return self.events[0..self.len]; } };}
fn subscribeWithin(io: Io, options: client_mod.Options, handler: anytype) !void { const Runner = struct { fn run(task_io: Io, task_options: client_mod.Options, task_handler: @TypeOf(handler), done: *std.atomic.Value(bool)) !void { defer done.store(true, .release); try client_mod.subscribe(task_io, testing.allocator, task_options, task_handler); } }; var done: std.atomic.Value(bool) = .init(false); var future = try io.concurrent(Runner.run, .{ io, options, handler, &done }); defer _ = future.cancel(io) catch {}; const start = Io.Timestamp.now(io, .awake).toMicroseconds(); while (!done.load(.acquire)) { if (Io.Timestamp.now(io, .awake).toMicroseconds() - start >= 10 * std.time.us_per_s) return error.TestTimeout; try io.sleep(Io.Duration.fromMilliseconds(10), .awake); } try future.await(io);}
test "e2e timestamp failover needs no archive and reconnects with fallback sequences" { try timestampFailover(false);}
test "e2e timestamp failover preserves resume through a refused connection" { try timestampFailover(true);}
fn timestampFailover(reject_first: bool) !void { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var primary = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1001, .t = 20 }, .{ .seq = 1002, .t = 21 } }); primary.live_then_close = true; primary.max_ws_sessions = 1; primary.reject_archive = true; var fallback = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 10 }, .{ .seq = 2, .t = 11 }, .{ .seq = 3, .t = 20 }, .{ .seq = 4, .t = 22 } }); fallback.reject_archive = true; fallback.reject_ws_count = if (reject_first) 1 else 0; fallback.first_session_rows = if (reject_first) null else 3; var primary_future = try io.concurrent(Instance.serve, .{&primary}); defer primary.deinit(); defer _ = primary_future.cancel(io) catch {}; defer primary.quit(); var fallback_future = try io.concurrent(Instance.serve, .{&fallback}); defer fallback.deinit(); defer _ = fallback_future.cancel(io) catch {}; defer fallback.quit(); const Handler = struct { collected: CollectingHandler(5) = .{}, failures: usize = 0, pub fn onEvent(self: *@This(), event: livedecode.Event) bool { return self.collected.onEvent(event); } pub fn onError(self: *@This(), _: anyerror) bool { self.failures += 1; return self.failures < 16; } }; var handler = Handler{}; var primary_url: [64]u8 = undefined; var fallback_url: [64]u8 = undefined; try subscribeWithin(io, .{ .hosts = &.{ primary.hostUrl(&primary_url), fallback.hostUrl(&fallback_url) }, .zstd_compression = false, .failover_after_errors = 2, .failover_stall_ns = 300 * std.time.ns_per_ms, .backoff_min_ns = std.time.ns_per_ms, .backoff_max_ns = 4 * std.time.ns_per_ms, }, &handler); const expected = [_]Collected{ .{ .seq = 1001, .t = 20 }, .{ .seq = 1002, .t = 21 }, .{ .seq = 2, .t = 11 }, .{ .seq = 3, .t = 20 }, .{ .seq = 4, .t = 22 } }; try testing.expectEqualSlices(Collected, &expected, handler.collected.slice()); try testing.expectEqual(@as(usize, 0), primary.archive_calls + fallback.archive_calls); try testing.expectEqual(@as(?u64, @intCast(witnessedUs(11))), fallback.wire_cursors[0]); try testing.expectEqual(@as(?u64, if (reject_first) @intCast(witnessedUs(11)) else 3), fallback.wire_cursors[1]);}
test "e2e initial timestamp survives a dead primary without widening the replay window" { try initialCursorFailover(@intCast(witnessedUs(11)), false);}
test "e2e initial sequence cannot transfer from a dead primary" { try initialCursorFailover(42, true);}
fn initialCursorFailover(cursor: u64, expect_unhealthy: bool) !void { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var primary = try Instance.init(io, testing.allocator, &.{}, &.{}); primary.max_ws_sessions = 0; var fallback = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 10 }, .{ .seq = 2, .t = 11 } }); var primary_future = try io.concurrent(Instance.serve, .{&primary}); defer primary.deinit(); defer _ = primary_future.cancel(io) catch {}; defer primary.quit(); var fallback_future = try io.concurrent(Instance.serve, .{&fallback}); defer fallback.deinit(); defer _ = fallback_future.cancel(io) catch {}; defer fallback.quit(); var primary_url: [64]u8 = undefined; var fallback_url: [64]u8 = undefined; var handler = CollectingHandler(1){}; const result = subscribeWithin(io, .{ .hosts = &.{ primary.hostUrl(&primary_url), fallback.hostUrl(&fallback_url) }, .live_cursor = cursor, .failover_after_errors = 1, .zstd_compression = false, .backoff_min_ns = std.time.ns_per_ms, }, &handler); if (expect_unhealthy) { try testing.expectError(error.HostUnhealthy, result); try testing.expectEqual(@as(usize, 0), handler.len); try testing.expectEqual(@as(u32, 0), fallback.ws_sessions); } else { try result; try testing.expectEqualSlices(Collected, &.{.{ .seq = 2, .t = 11 }}, handler.slice()); try testing.expectEqual(@as(?u64, cursor), fallback.wire_cursors[0]); } try testing.expectEqual(@as(usize, 0), primary.archive_calls + fallback.archive_calls);}
test "e2e expired timestamp reports the retention gap before delivering retained events" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 10 }, .{ .seq = 2, .t = 20 } }); host.retention_floor_us = @intCast(witnessedUs(20)); var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { collected: CollectingHandler(1) = .{}, outdated: usize = 0, warned_before_event: bool = false, pub fn onEvent(self: *@This(), event: livedecode.Event) bool { self.warned_before_event = self.outdated == 1; return self.collected.onEvent(event); } pub fn onInfo(self: *@This(), info: livedecode.Info) void { if (std.mem.eql(u8, info.name, "OutdatedCursor")) self.outdated += 1; } }; var handler = Handler{}; var url: [64]u8 = undefined; try subscribeWithin(io, .{ .hosts = &.{host.hostUrl(&url)}, .live_cursor = @intCast(witnessedUs(10)), .zstd_compression = false, }, &handler); try testing.expect(handler.warned_before_event); try testing.expectEqualSlices(Collected, &.{.{ .seq = 2, .t = 20 }}, handler.collected.slice()); try testing.expectEqual(@as(usize, 0), host.archive_calls);}
test "e2e expired live sequence returns CursorTooOld without skipping to tip" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{}); host.cursor_too_old = true; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(1){}; try testing.expectError(error.CursorTooOld, subscribeWithin(io, .{ .hosts = &.{host.hostUrl(&url)}, .live_cursor = 42, .zstd_compression = false, }, &handler)); try testing.expectEqual(@as(u32, 1), host.ws_sessions); try testing.expectEqual(@as(usize, 0), host.archive_calls); try testing.expectEqual(@as(usize, 0), handler.slice().len);}
test "e2e failover: archive replay completes before live timestamp failover" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io();
// primary: sealed seqs 1-2 (t=1,2), live 3-4 (t=3,4), then the host dies var primary = try Instance.init(io, testing.allocator, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 } }, &.{ .{ .seq = 3, .t = 3 }, .{ .seq = 4, .t = 4 } }); primary.live_then_close = true; primary.max_ws_sessions = 1;
// fallback: a DIFFERENT seq space. floor = last_time(4s) - rewind(1s) // = 3s: seq 1001 (t=2) is trimmed, 1002 (t=3, boundary duplicate) // and 1003 (t=5) deliver, then live 1004 (t=6) var fallback = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1001, .t = 2 }, .{ .seq = 1002, .t = 3 }, .{ .seq = 1003, .t = 5 }, .{ .seq = 1004, .t = 6 } });
fallback.reject_archive = true;
// LIFO: request shutdown, join the serve task, then close the listener. var primary_future = try io.concurrent(Instance.serve, .{&primary}); defer primary.deinit(); defer _ = primary_future.cancel(io) catch {}; defer primary.quit(); var fallback_future = try io.concurrent(Instance.serve, .{&fallback}); defer fallback.deinit(); defer _ = fallback_future.cancel(io) catch {}; defer fallback.quit();
var primary_url: [64]u8 = undefined; var fallback_url: [64]u8 = undefined; var handler = CollectingHandler(7){}; try subscribeWithin(io, .{ .hosts = &.{ primary.hostUrl(&primary_url), fallback.hostUrl(&fallback_url) }, .after_seq = 0, .failover_after_errors = 2, .failover_rewind_us = 1 * std.time.us_per_s, .backoff_min_ns = 1 * std.time.ns_per_ms, .backoff_max_ns = 4 * std.time.ns_per_ms, }, &handler);
const expected = [_]Collected{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, .{ .seq = 4, .t = 4 }, .{ .seq = 1002, .t = 3 }, .{ .seq = 1003, .t = 5 }, .{ .seq = 1004, .t = 6 }, }; try testing.expectEqualSlices(Collected, &expected, handler.slice()); try testing.expectEqual(@as(usize, 0), handler.missing_record_cbor); try testing.expect(primary.plan_calls > 0); try testing.expectEqual(@as(usize, 0), fallback.archive_calls);}
test "e2e stall: connected-but-silent primary rotates; fresh tip on the fallback" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io();
var primary = try Instance.init(io, testing.allocator, &.{}, &.{}); primary.silent_ws = true; primary.max_ws_sessions = 1; // one silent session, then down
var fallback = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 501, .t = 10 }, .{ .seq = 502, .t = 11 } });
var primary_future = try io.concurrent(Instance.serve, .{&primary}); defer primary.deinit(); defer _ = primary_future.cancel(io) catch {}; defer primary.quit(); var fallback_future = try io.concurrent(Instance.serve, .{&fallback}); defer fallback.deinit(); defer _ = fallback_future.cancel(io) catch {}; defer fallback.quit();
var primary_url: [64]u8 = undefined; var fallback_url: [64]u8 = undefined; var handler = CollectingHandler(2){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{ primary.hostUrl(&primary_url), fallback.hostUrl(&fallback_url) }, .failover_after_errors = 1, .failover_stall_ns = 300 * std.time.ns_per_ms, .backoff_min_ns = 1 * std.time.ns_per_ms, .backoff_max_ns = 4 * std.time.ns_per_ms, }, &handler);
const expected = [_]Collected{ .{ .seq = 501, .t = 10 }, .{ .seq = 502, .t = 11 } }; try testing.expectEqualSlices(Collected, &expected, handler.slice()); try testing.expectEqual(@as(usize, 0), handler.missing_record_cbor);}
test "e2e live-phase clean stop: onEvent=false ends subscribe instead of reconnect-looping" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io();
var host = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 } }); var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit();
var url: [64]u8 = undefined; var handler = CollectingHandler(1){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .backoff_min_ns = 1 * std.time.ns_per_ms, .backoff_max_ns = 4 * std.time.ns_per_ms, }, &handler);
try testing.expectEqual(@as(usize, 1), handler.len); try testing.expectEqual(@as(u64, 1), handler.slice()[0].seq);}
// Upstream engine_test.go: TestEngineMultiPageBackfillCutover and// TestEnginePinnedBeforeSeqAcrossPages (289b032 and 58c4d7f).test "conformance: snapshot consumes every plan page and pins the sealed tip" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, .{ .seq = 4, .t = 4 }, }, &.{}); host.page_size = 1; host.first_tip = 3; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(32){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .after_seq = 0, .snapshot_only = true, .concurrency = 1, }, &handler); try testing.expectEqualSlices(Collected, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, }, handler.slice()); try testing.expectEqual(@as(usize, 3), host.plan_calls); try testing.expectEqual(@as(?u64, 3), host.last_before_seq);}
// Upstream TestEngineLiveOnlyAppliesCollectionFilter and// TestMatcherDIDFilterAllKinds: the server deliberately ignores filters.test "conformance: live delivery reapplies collection and DID filters" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 1, .collection = "app.bsky.feed.like" }, .{ .seq = 2, .t = 2, .did = "did:plc:other" }, .{ .seq = 3, .t = 3 }, }); var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { collected: CollectingHandler(32) = .{}, pub fn onEvent(self: *@This(), event: livedecode.Event) bool { _ = self.collected.onEvent(event); return event.seq < 3; } }; var handler = Handler{}; var url: [64]u8 = undefined; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .collections = &.{"app.bsky.feed.post"}, .dids = &.{"did:plc:loopback"}, }, &handler); try testing.expectEqualSlices(Collected, &.{.{ .seq = 3, .t = 3 }}, handler.collected.slice());}
// Upstream sweepSealedArchive rejects a non-advancing continuation.test "conformance: a stalled plan fails instead of looping" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{.{ .seq = 3, .t = 3 }}, &.{}); host.stalled_plan = true; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(32){}; try testing.expectError(error.PlanStalled, client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .after_seq = 0, .snapshot_only = true, }, &handler)); try testing.expectEqual(@as(usize, 1), host.plan_calls);}
test "archive block budget applies across plan pages" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 }, .{ .seq = 3, .t = 3 }, }, &.{}); host.page_size = 1; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { count: usize = 0, pub fn onRow(self: *@This(), _: livedecode.Event) bool { self.count += 1; return true; } }; var handler = Handler{}; var url: [64]u8 = undefined; const result = try archive.run(io, testing.allocator, .{ .host = host.hostUrl(&url), .after_seq = 0, .max_blocks = 2, .concurrency = 1, }, &handler); try testing.expect(result.truncated); try testing.expectEqual(@as(u64, 2), result.blocks_decoded); try testing.expectEqual(@as(usize, 2), handler.count); try testing.expectEqual(@as(usize, 2), host.plan_calls); try testing.expectEqual(@as(u64, 3), result.sealed_tip_seq);}
// Upstream TestDownloadBlocksErrorStopsEntryNotPlan: good prefix, error,// then the next entry; a decode error instead leaves later blocks readable.test "archive errors are ordered and the caller chooses whether to continue" { for ([_]bool{ true, false }) |accept_error| { for ([_]usize{ 1, 4 }) |concurrency| { for ([_]bool{ false, true }) |decode_failure| { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{}); host.archive_fault = if (decode_failure) .decode else .fetch; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { items: [8]u64 = undefined, len: usize = 0, accept: bool, err: ?anyerror = null, pub fn onEvent(self: *@This(), event: livedecode.Event) bool { self.items[self.len] = event.seq; self.len += 1; return true; } pub fn onError(self: *@This(), err: anyerror) bool { self.err = err; self.items[self.len] = 0; self.len += 1; return self.accept; } }; var handler = Handler{ .accept = accept_error }; var url: [64]u8 = undefined; // If an archive failure incorrectly rotates, port 1 cannot // supply the next entry and the expected sequence is lost. try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{ host.hostUrl(&url), "http://127.0.0.1:1" }, .after_seq = 0, .snapshot_only = true, .concurrency = concurrency, }, &handler); try testing.expectEqual(@as(?anyerror, if (decode_failure) error.BlockDecompressFailed else error.FetchFailed), handler.err); const expected: []const u64 = if (!accept_error) &.{ 1, 0 } else if (decode_failure) &.{ 1, 0, 3, 100 } else &.{ 1, 0, 100 }; try testing.expectEqualSlices(u64, expected, handler.items[0..handler.len]); } } }}
test "archive plan failures stay terminal even when the caller accepts errors" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{}); host.reject_plan = true; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { errors: usize = 0, pub fn onEvent(_: *@This(), _: livedecode.Event) bool { return true; } pub fn onError(self: *@This(), _: anyerror) bool { self.errors += 1; return true; } }; var handler = Handler{}; var url: [64]u8 = undefined; try testing.expectError(error.PlanFailed, client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .after_seq = 0, .snapshot_only = true, }, &handler)); try testing.expectEqual(@as(usize, 0), handler.errors);}
test "a whole segment error can be accepted or declined before the next entry" { for ([_]bool{ true, false }) |accept| { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{}); host.archive_fault = .segment; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { accept: bool, errors: usize = 0, last: u64 = 0, pub fn onRow(self: *@This(), event: livedecode.Event) bool { self.last = event.seq; return true; } pub fn onError(self: *@This(), err: anyerror) bool { if (err != error.MalformedSegment) @panic("unexpected error"); self.errors += 1; return self.accept; } }; var handler = Handler{ .accept = accept }; var url: [64]u8 = undefined; const result = try archive.run(io, testing.allocator, .{ .host = host.hostUrl(&url) }, &handler); try testing.expectEqual(@as(usize, 1), handler.errors); try testing.expectEqual(@as(u64, if (accept) 100 else 0), handler.last); try testing.expectEqual(!accept, result.stopped); }}
test "archive key authenticates archive requests but not live handshakes" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{.{ .seq = 1, .t = 1 }}, &.{.{ .seq = 2, .t = 2 }}); var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(2){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .after_seq = 0, .api_key = "test-archive-key", .concurrency = 1, }, &handler); try testing.expectEqual(@as(usize, 2), handler.len); try testing.expect(host.archive_authorized); try testing.expect(!host.archive_missing_auth); try testing.expect(!host.live_authorized);}
test "archive block retries stop after three HTTP attempts" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{.{ .seq = 1, .t = 1 }}, &.{}); host.retry_block_calls = 0; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(2){}; try testing.expectError(error.FetchRetriesExhausted, client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .after_seq = 0, .snapshot_only = true, .concurrency = 1, }, &handler)); try testing.expectEqual(@as(?usize, 3), host.retry_block_calls); try testing.expectEqual(@as(usize, 0), handler.len);}
test "default compression negotiates and recovers dictionary rotation" { inline for (.{ .valid, .rotate, .same, .invalid, .none }) |mode| { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{.{ .seq = 1, .t = 1 }}); host.compression = mode; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(1){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)} }, &handler); try testing.expectEqual(@as(usize, 1), handler.len); try testing.expectEqual(@as(usize, if (mode == .rotate or mode == .same) 2 else 1), host.dict_calls); try testing.expectEqual(@as(usize, if (mode == .rotate) 2 else if (mode == .valid or mode == .same) 1 else 0), host.compressed_sessions); }}
test "explicitly disabled compression makes no dictionary request" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{.{ .seq = 1, .t = 1 }}); host.compression = .valid; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(1){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .zstd_compression = false }, &handler); try testing.expectEqual(@as(usize, 1), handler.len); try testing.expectEqual(@as(usize, 0), host.dict_calls); try testing.expectEqual(@as(usize, 0), host.compressed_sessions);}
test "malformed compressed live frame follows caller continue or stop policy" { inline for (.{ true, false }) |accept| { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{.{ .seq = 1, .t = 1 }}); host.compression = .valid; host.bad_compressed_first = true; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { errors: usize = 0, events: usize = 0, pub fn onError(self: *@This(), err: anyerror) bool { testing.expectEqual(error.DecompressFailed, err) catch @panic("unexpected live error"); self.errors += 1; return accept; } pub fn onEvent(self: *@This(), _: livedecode.Event) bool { self.events += 1; return false; } }; var handler = Handler{}; var url: [64]u8 = undefined; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)} }, &handler); try testing.expectEqual(@as(usize, 1), handler.errors); try testing.expectEqual(@as(usize, if (accept) 1 else 0), handler.events); try testing.expectEqual(@as(u32, 1), host.ws_sessions); }}
test "live resume suppresses boundary sent by server" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{ .{ .seq = 1, .t = 1 }, .{ .seq = 2, .t = 2 } }); host.include_cursor_boundary = true; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); var url: [64]u8 = undefined; var handler = CollectingHandler(1){}; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)}, .live_cursor = 1, .zstd_compression = false }, &handler); try testing.expectEqual(@as(u64, 2), handler.events[0].seq);}
test "dictionary rejection reports recovery and honors caller stop" { inline for (.{ true, false }) |accept| { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); const io = threaded.io(); var host = try Instance.init(io, testing.allocator, &.{}, &.{.{ .seq = 1, .t = 1 }}); host.compression = .rotate; var future = try io.concurrent(Instance.serve, .{&host}); defer host.deinit(); defer _ = future.cancel(io) catch {}; defer host.quit(); const Handler = struct { errors: usize = 0, events: usize = 0, pub fn onError(self: *@This(), _: anyerror) bool { self.errors += 1; return accept; } pub fn onEvent(self: *@This(), _: livedecode.Event) bool { self.events += 1; return false; } }; var handler = Handler{}; var url: [64]u8 = undefined; try client_mod.subscribe(io, testing.allocator, .{ .hosts = &.{host.hostUrl(&url)} }, &handler); try testing.expectEqual(@as(usize, 1), handler.errors); try testing.expectEqual(@as(usize, if (accept) 1 else 0), handler.events); try testing.expectEqual(@as(u32, if (accept) 2 else 1), host.ws_sessions); try testing.expectEqual(@as(usize, 2), host.dict_calls); }}