//! slice sync: everything one account (or one app's collections) ever did, //! replayed from the whole-network archive in seconds — the local-first //! pattern from the Jetstream v2 launch thread ("backfill all of my Semble //! data into the browser: 12 seconds, 383 requests, 190kb over the wire"). //! //! the planSnapshot request carries the did/collection filters, so the //! server prunes to just the blocks containing matching rows — the client //! never downloads the archive, only the slice. //! //! run: zig build example-slice-sync -- [collection,collection,...] //! zig build example-slice-sync -- zzstoatzz.io app.bsky.feed.post const std = @import("std"); const zat = @import("zat"); const jetstream_sdk = @import("jetstream"); const SliceStats = struct { rows: usize = 0, first_time_us: i64 = 0, last_time_us: i64 = 0, by_collection: std.StringHashMapUnmanaged(usize) = .empty, arena: std.mem.Allocator, pub fn onRow(self: *@This(), event: jetstream_sdk.Event) bool { self.rows += 1; if (self.first_time_us == 0) self.first_time_us = event.time_us; self.last_time_us = event.time_us; if (event.payload == .commit) { const gop = self.by_collection.getOrPut(self.arena, event.payload.commit.collection) catch return true; if (!gop.found_existing) { gop.key_ptr.* = self.arena.dupe(u8, event.payload.commit.collection) catch return true; gop.value_ptr.* = 0; } gop.value_ptr.* += 1; } return true; } }; pub fn main(init: std.process.Init) !void { const allocator = init.gpa; const args = try init.minimal.args.toSlice(init.arena.allocator()); if (args.len < 2) { std.debug.print("usage: zig build example-slice-sync -- [collections,...]\n", .{}); return error.MissingActor; } // "-" as the actor means no DID filter: a network-wide collection slice var did: []const u8 = args[1]; if (std.mem.eql(u8, did, "-")) { did = ""; } else if (!std.mem.startsWith(u8, did, "did:")) { var resolver = zat.HandleResolver.init(init.io, allocator); defer resolver.deinit(); const resolved = try resolver.resolve(zat.Handle.parse(args[1]) orelse return error.InvalidHandle); defer allocator.free(resolved); did = try init.arena.allocator().dupe(u8, resolved); } // token-gated archives (stream.waow.tech is one) need a bearer key const api_key: ?[]const u8 = if (std.c.getenv("JETSTREAM_API_KEY")) |v| std.mem.span(v) else null; var collections: std.ArrayList([]const u8) = .empty; if (args.len > 2) { var it = std.mem.tokenizeScalar(u8, args[2], ','); while (it.next()) |c| try collections.append(init.arena.allocator(), c); } std.debug.print("slicing the archive for {s} ({d} collection filters)\n", .{ if (did.len > 0) did else "the whole network", collections.items.len }); const started = std.Io.Timestamp.now(init.io, .awake).nanoseconds; var stats = SliceStats{ .arena = init.arena.allocator() }; const result = try jetstream_sdk.ArchiveBackfill.run(init.io, allocator, .{ .host = "https://stream.waow.tech", .after_seq = 0, .dids = if (did.len > 0) &.{did} else &.{}, .collections = collections.items, // jim's sparse-backfill recipe: kind=commit changes the PLAN, not // just client filtering — without it, marker/sentinel rows admit // nearly every block (measured 83x more blocks for a network-wide // collection slice) .kinds = &.{.commit}, .api_key = api_key, }, &stats); const elapsed_ms = @divFloor(std.Io.Timestamp.now(init.io, .awake).nanoseconds - started, std.time.ns_per_ms); std.debug.print("\n{d} rows in {d}ms — {d} block fetches, {d} KB over the wire\n", .{ stats.rows, elapsed_ms, result.blocks_decoded, result.bytes_fetched / 1024 }); var it = stats.by_collection.iterator(); while (it.next()) |entry| { std.debug.print(" {s}: {d}\n", .{ entry.key_ptr.*, entry.value_ptr.* }); } if (stats.rows > 0) { const span_days = @divFloor(stats.last_time_us - stats.first_time_us, std.time.us_per_s * 86_400); std.debug.print(" spanning {d} days of history\n", .{span_days}); } }