Something went wrong. Try again.
A Zig library for Bluesky activity and feed analytics, built on Zat.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162const std = @import("std");const Client = @import("client.zig").Client;const archive = @import("archive.zig");const Owned = @import("owned.zig").Owned;const store = @import("store.zig");
const Info = struct { did: []const u8, pds: ?[]const u8 = null, accountStatus: struct { active: bool }, archive: struct { available: bool }, syncState: struct { state: []const u8, rev: ?[]const u8 = null },};
pub const Checkpoint = struct { did: []const u8, source: []const u8, revision: []const u8, path: []const u8, format: archive.Format, observed_at: i64, records: usize,};
pub const Result = struct { checkpoint: Checkpoint, checked_at: i64, reused: bool, window: @import("activity.zig").Window, counts: @import("activity.zig").Counts, diagnostics: archive.Diagnostics };
pub const Options = struct { window: @import("activity.zig").Window, max_bytes: usize = 512 * 1024 * 1024 };
pub fn refresh(client: *Client, dir: std.Io.Dir, did: []const u8, options: Options) !Owned(Result) { try options.window.validate(); if (client.session != null) return error.AuthenticatedArchiveClient; var arena = std.heap.ArenaAllocator.init(client.allocator); errdefer arena.deinit(); const a = arena.allocator(); const response = client.query(Info, "blue.microcosm.hubble.getRepoInfo", &.{.{ .name = "did", .value = did }}) catch |err| { if (client.last_fault) |fault| { if (std.mem.eql(u8, fault.@"error", "RepoGone")) return error.RepoGone; if (std.mem.eql(u8, fault.@"error", "RepoNotFound")) return error.RepoNotFound; } return err; }; defer response.deinit(); const info = response.value; try validate(info, did); const checked_at = std.Io.Timestamp.now(client.io, .real).toMicroseconds(); const revision = info.syncState.rev.?; const key = try cacheKey(a, client.rpc.host, did, revision, archive.Format.car); const path = try std.fmt.allocPrint(a, "{s}.{t}", .{ key, archive.Format.car }); const metadata = try std.fmt.allocPrint(a, "{s}.json", .{key}); const cached = dir.readFileAlloc(client.io, metadata, a, .limited(16 * 1024)) catch |err| switch (err) { error.FileNotFound => null, else => return err, }; if (cached) |bytes| { const checkpoint = try std.json.parseFromSliceLeaky(Checkpoint, a, bytes, .{ .allocate = .alloc_always }); if (!std.mem.eql(u8, checkpoint.did, did) or !std.mem.eql(u8, checkpoint.source, client.rpc.host) or !std.mem.eql(u8, checkpoint.revision, revision) or !std.mem.eql(u8, checkpoint.path, path) or checkpoint.format != .car) return error.CheckpointMismatch; var file = try dir.openFile(client.io, path, .{}); defer file.close(client.io); var summary = try analyze(client.io, client.allocator, file, checkpoint, options.window, checked_at); defer summary.deinit(); if (!std.mem.eql(u8, summary.value.revision, revision)) return error.RevisionMismatch; return .{ .arena = arena, .value = .{ .checkpoint = checkpoint, .checked_at = checked_at, .reused = true, .window = options.window, .counts = summary.value.counts, .diagnostics = summary.value.diagnostics } }; } var staged = try dir.createFileAtomic(client.io, path, .{ .permissions = .fromMode(0o600), .replace = true }); defer staged.deinit(client.io); var buffer: [64 * 1024]u8 = undefined; var writer = staged.file.writer(client.io, &buffer); _ = try client.download("com.atproto.sync.getRepo", &.{.{ .name = "did", .value = did }}, "application/vnd.ipld.car", &writer.interface, options.max_bytes); try writer.interface.flush(); const observed_at = std.Io.Timestamp.now(client.io, .real).toMicroseconds(); const staged_path = std.fmt.hex(staged.file_basename_hex); const readable = try staged.dir.openFile(client.io, &staged_path, .{}); defer readable.close(client.io); var checkpoint: Checkpoint = .{ .did = try a.dupe(u8, did), .source = try a.dupe(u8, client.rpc.host), .revision = try a.dupe(u8, revision), .path = path, .format = archive.Format.car, .observed_at = observed_at, .records = 0 }; var summary = try analyze(client.io, client.allocator, readable, checkpoint, options.window, observed_at); defer summary.deinit(); if (!std.mem.eql(u8, summary.value.revision, revision)) return error.RevisionChanged; checkpoint.records = summary.value.records; try staged.file.sync(client.io); try staged.replace(client.io); try store.save(client.io, a, dir, metadata, checkpoint, .replace); return .{ .arena = arena, .value = .{ .checkpoint = checkpoint, .checked_at = checked_at, .reused = false, .window = options.window, .counts = summary.value.counts, .diagnostics = summary.value.diagnostics } };}
pub fn analyze(io: std.Io, allocator: std.mem.Allocator, file: std.Io.File, checkpoint: Checkpoint, window: @import("activity.zig").Window, checked_at: i64) !Owned(archive.Analysis) { var buffer: [64 * 1024]u8 = undefined; var reader = file.reader(io, &buffer); var result = try archive.read(allocator, &reader.interface, .{ .format = checkpoint.format, .did = checkpoint.did, .window = window, .observed_at = checked_at }); errdefer result.deinit(); if (!std.mem.eql(u8, result.value.revision, checkpoint.revision)) return error.RevisionMismatch; return result;}
fn validate(info: Info, did: []const u8) !void { if (!std.mem.eql(u8, info.did, did)) return error.IdentityMismatch; if (!info.accountStatus.active) return error.InactiveAccount; if (!std.mem.eql(u8, info.syncState.state, "synchronized") or info.syncState.rev == null) return error.MirrorIncomplete; if (!info.archive.available) return error.ArchiveUnavailable;}
fn cacheKey(a: std.mem.Allocator, source: []const u8, did: []const u8, revision: []const u8, format: archive.Format) ![]const u8 { var hash = std.crypto.hash.sha2.Sha256.init(.{}); for ([_][]const u8{ source, did, revision, @tagName(format), "v1" }) |part| { hash.update(part); hash.update(&.{0}); } var digest: [32]u8 = undefined; hash.final(&digest); return std.fmt.allocPrint(a, "{s}", .{std.fmt.bytesToHex(digest, .lower)});}
test "mirror coverage and identity are required before archive acquisition" { var info: Info = .{ .did = "did:plc:a", .accountStatus = .{ .active = true }, .archive = .{ .available = true }, .syncState = .{ .state = "synchronized", .rev = "rev" } }; try validate(info, "did:plc:a"); try std.testing.expectError(error.IdentityMismatch, validate(info, "did:plc:b")); info.syncState.state = "desynchronized"; try std.testing.expectError(error.MirrorIncomplete, validate(info, "did:plc:a")); info.archive.available = false; info.syncState.state = "synchronized"; try std.testing.expectError(error.ArchiveUnavailable, validate(info, "did:plc:a"));}
test "failed verification publishes nothing and retry reuses a verified archive" { const fixture = @import("test_repo.zig"); const origin_module = @import("test_origin.zig"); const io = std.testing.io; var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); const a = arena.allocator(); var repo = try fixture.make(std.testing.allocator, false); defer repo.deinit(); const damaged = try a.dupe(u8, repo.value.car); damaged[damaged.len - 1] ^= 1; const info = try std.json.Stringify.valueAlloc(a, .{ .did = fixture.did, .accountStatus = .{ .active = true }, .archive = .{ .available = true }, .syncState = .{ .state = "synchronized", .rev = fixture.revision } }, .{}); const response = try origin_module.response(a, info); const replies: []const []const u8 = &.{ response, try origin_module.response(a, damaged), response, try origin_module.response(a, repo.value.car), response }; var origin = try origin_module.Origin.init(io, replies); defer origin.deinit(); var serving = try io.concurrent(origin_module.Origin.serve, .{&origin}); defer _ = serving.cancel(io) catch {}; var tmp = std.testing.tmpDir(.{}); defer tmp.cleanup(); const host = try origin.url(a); var client = Client.init(io, std.testing.allocator, host); defer client.deinit(); const options: Options = .{ .window = .{ .since = 0, .until = 1_790_000_000_000_000 } }; try std.testing.expectError(error.BadBlockHash, refresh(&client, tmp.dir, fixture.did, options)); const key = try cacheKey(a, host, fixture.did, fixture.revision, .car); const name = try std.fmt.allocPrint(a, "{s}.json", .{key}); try std.testing.expectError(error.FileNotFound, tmp.dir.openFile(io, name, .{})); var fresh = try refresh(&client, tmp.dir, fixture.did, options); defer fresh.deinit(); var reused = try refresh(&client, tmp.dir, fixture.did, options); defer reused.deinit(); try std.testing.expect(!fresh.value.reused and reused.value.reused); try std.testing.expectEqual(fresh.value.checkpoint.observed_at, reused.value.checkpoint.observed_at); try std.testing.expect(fresh.value.counts.complete() and fresh.value.counts.reply == 1); try serving.await(io); try std.testing.expectEqual(@as(usize, 5), origin.requests);}