diff --git a/examples/07_streamplace_chat.md b/examples/07_streamplace_chat.md index 3b7edda..d05d59f 100644 --- a/examples/07_streamplace_chat.md +++ b/examples/07_streamplace_chat.md @@ -41,31 +41,57 @@ logical Streamplace session: at://did:plc:.../place.stream.livestream/3ms73ig6jqr2e joined 2 contiguous livestream records; Streamplace rolled over the AT-URI mid-session. -Jetstream replay is retention-limited; the final count is observed, not authoritative for old cursors. -replaying this historical session's chat... +replaying from the archive (https://stream.waow.tech)... -[2026-08-03T17:37:53.900Z] did:plc:...: heloo hello -[2026-08-03T17:40:04.633Z] did:plc:...: oo this is groovy +[2026-08-03T18:07:18.906Z] did:plc:...: isnt BlueSky that one pixar subsidiary ... -replay complete: 652 observed messages across 2 segments +archive replay: 48 messages from 10209 blocks +finishing the session tail over the live socket... + +replay complete: 48 observed messages across 2 segments ``` The output uses author DIDs because Jetstream carries repo writes, not hydrated profiles. The implementation is in `examples/07_streamplace_chat.zig` and is -compile-checked by `zig build`. - -This example still relies on the Jetstream instance retaining the requested -cursor. Exact historical selection prevents chat from adjacent sessions from -leaking into the result, but a successful run is not proof of archival -completeness: an instance may clamp a cursor older than its replay window. - -For retention-independent current-state discovery, use the same shape as the +compile-checked by `zig build`. Build it with `-Doptimize=ReleaseSafe` — the +archive pass decompresses hundreds of zstd blocks, which is slow in Debug. + +Historical chat comes from `zat.ArchiveBackfill` against +[stream.waow.tech](https://stream.waow.tech)'s full-network archive, so replay +does not depend on Jetstream cursor retention. The session's time window maps +to plan seq bounds (`fetchSeqBounds`), so a recent broadcast plans only the +handful of blocks it touches. Sessions inside the archive's merged bootstrap +region (before 2026-08-04) still replay, but their rows are scattered per-DID, +so the plan spans more blocks — expect a bulk download, and watch the +`archive backfill: N blocks decoded` progress lines. The newest events sit in +the archive's unsealed tail, so the example always finishes over the live +socket from the covered boundary; the two passes overlap at the seam and the +printer dedups by rkey. + +Two honesty notes an archive consumer should internalize, both visible in the +sample output above: + +- **An archive replays what it holds, which is not always what happened.** + This very session was observed live at 652 messages on 2026-08-03, but the + archive replays 48: the broadcast fell inside a dated capture gap in + stream.waow.tech's bootstrap week (live capture was down from 2026-08-03 + 07:05Z into 2026-08-04; the relay replay after resume did not reach all the + way back). The 48 that do replay are `create_resync` rows — records + re-imported later by the archive's repo-healing — which is why they exist at + all. A different instance, or any session after steady state began + (2026-08-04 22:32Z), replays completely. +- **Delete-compaction is real.** stream runs jetstream's delete-compaction, so + a record deleted later has its create row physically folded out of sealed + history; the delete tombstone remains (delivered as a `.delete` event). An + archive is the network's memory, not a time machine for content its authors + withdrew. + +For current-state discovery without an archive, use the same shape as the [`collection_creators` flow](https://github.com/zzstoatzz/prefect-pack/blob/main/flows/collection_creators/main.py): enumerate every repo containing -`place.stream.chat.message`, resolve each repo's PDS, fetch its current records, -and filter them by `streamer` and time. That is substantially more expensive -than this Jetstream example and cannot recover deleted records. A durable chat -archive should instead consume continuously and persist both matching records -and its cursor. +`place.stream.chat.message`, resolve each repo's PDS, fetch its current +records, and filter them by `streamer` and time. That cannot recover deleted +records. A durable chat archive should still consume continuously and persist +both matching records and its cursor. Streamplace's `/api/websocket/` endpoint is useful when hydrated author profiles, moderation events, and viewer counts matter more than complete diff --git a/examples/07_streamplace_chat.zig b/examples/07_streamplace_chat.zig index 22e148f..2ce656d 100644 --- a/examples/07_streamplace_chat.zig +++ b/examples/07_streamplace_chat.zig @@ -8,6 +8,7 @@ const zat = @import("zat"); const rollover_tolerance_us: i64 = 2 * std.time.us_per_s; const ingestion_grace_us: i64 = 5 * std.time.us_per_min; +const archive_host = "https://stream.waow.tech"; const Segment = struct { uri: []u8, @@ -28,10 +29,20 @@ const Segment = struct { }; const ChatPrinter = struct { + allocator: std.mem.Allocator, streamer_did: []const u8, first_stream_rkey: []const u8, ended_at_us: ?i64, count: usize = 0, + /// the archive pass and the live replay overlap at the seam by design + /// (delivery is at-least-once); "did/rkey" keys dedup the seam + seen: std.StringHashMapUnmanaged(void) = .empty, + + fn deinit(self: *ChatPrinter) void { + var keys = self.seen.keyIterator(); + while (keys.next()) |key| self.allocator.free(key.*); + self.seen.deinit(self.allocator); + } pub fn onEvent(self: *ChatPrinter, event: zat.JetstreamEvent) void { const commit = switch (event) { @@ -54,6 +65,15 @@ const ChatPrinter = struct { if (created_at_us > end) return; } + var key_buf: [512]u8 = undefined; + const key = std.fmt.bufPrint(&key_buf, "{s}/{s}", .{ commit.did, commit.rkey }) catch return; + const entry = self.seen.getOrPut(self.allocator, key) catch return; + if (entry.found_existing) return; + entry.key_ptr.* = self.allocator.dupe(u8, key) catch { + _ = self.seen.remove(key); + return; + }; + const text = zat.json.getString(record, "text") orelse ""; self.count += 1; std.debug.print("[{s}] {s}: {s}\n", .{ created_at, commit.did, text }); @@ -167,28 +187,58 @@ pub fn main(init: std.process.Init) !void { if (segment_count > 1) { std.debug.print("joined {d} contiguous livestream records; Streamplace rolled over the AT-URI mid-session.\n", .{segment_count}); } - std.debug.print("Jetstream replay is retention-limited; the final count is observed, not authoritative for old cursors.\n{s}\n\n", .{ - if (replay_end == null) - "replaying chat, then following live..." - else - "replaying this historical session's chat...", - }); var printer = ChatPrinter{ + .allocator = allocator, .streamer_did = streamer_did, .first_stream_rkey = first_rkey, .ended_at_us = if (replay_end) |end| try parseUtcMicros(end) else null, }; + defer printer.deinit(); + // The livestream record's TID is its creation time in microseconds. + // Rewind one second so events at the boundary cannot be skipped. + const start_us = @as(i64, @intCast(first_tid.timestamp())) - std.time.us_per_s; const end_cursor = if (replay_end) |end| (try parseUtcMicros(end)) + ingestion_grace_us else null; + + // Historical chat comes from the archive, not from live-replay retention: + // stream.waow.tech archives every event, so any session since its live + // capture began replays completely — including messages later deleted. + // The witnessed-time window maps to plan seq bounds, which keeps the + // fetch to the handful of blocks the session actually touches. + const bounds = try zat.ArchiveBackfill.fetchSeqBounds(init.io, allocator, archive_host, start_us, end_cursor); + var archive_boundary_us: ?i64 = null; + if (bounds.covered) { + std.debug.print("replaying from the archive ({s})...\n\n", .{archive_host}); + const result = try zat.ArchiveBackfill.run(init.io, allocator, .{ + .host = archive_host, + .collections = &.{"place.stream.chat.message"}, + .after_seq = bounds.after_seq, + .before_seq = bounds.before_seq, + .concurrency = 6, + }, &printer); + archive_boundary_us = result.last_time_us; + std.debug.print("\narchive replay: {d} messages from {d} blocks\n", .{ printer.count, result.blocks_decoded }); + } else { + std.debug.print("session predates the archive's live capture; only live replay is available.\n", .{}); + } + + // The newest events live in the archive's unsealed tail, so a complete + // replay always finishes over the live socket from the covered boundary + // (the seam overlaps; the printer dedups). For a live session this is + // also the follow. + std.debug.print("{s}\n\n", .{ + if (replay_end == null) + "catching up over the live socket, then following live..." + else + "finishing the session tail over the live socket...", + }); var client = zat.JetstreamClient.init(init.io, allocator, .{ .wanted_collections = &.{"place.stream.chat.message"}, - // The livestream record's TID is its creation time in microseconds. - // Rewind one second so events at the boundary cannot be skipped. - .cursor = @as(i64, @intCast(first_tid.timestamp())) - std.time.us_per_s, + .cursor = if (archive_boundary_us) |t| t - std.time.us_per_s else start_us, // Stop close to the broadcast end, with room for delayed ingestion. // The idle timeout handles the case where no matching event lands // exactly beyond that boundary.