Something went wrong. Try again.
A Zig library for Bluesky activity and feed analytics, built on Zat.
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255const std = @import("std");const zat = @import("zat");const activity = @import("activity.zig");const Owned = @import("owned.zig").Owned;
pub const Format = enum { star, car };pub const Options = struct { format: Format, did: []const u8, window: activity.Window, observed_at: i64, retain_evidence: bool = false, collect_follows: bool = false, max_car_bytes: usize = 512 * 1024 * 1024,};pub const Diagnostics = struct { invalid_posts: usize = 0, invalid_reposts: usize = 0, invalid_likes: usize = 0, first_error: ?[]const u8 = null };pub const Analysis = struct { did: []const u8, revision: []const u8, window: activity.Window, observed_at: i64, records: usize, counts: activity.Counts, diagnostics: Diagnostics, evidence: ?activity.EvidenceSet, follows: []const []const u8,};
pub fn read(allocator: std.mem.Allocator, input: *std.Io.Reader, options: Options) !Owned(Analysis) { try options.window.validate(); if (zat.Did.parse(options.did) == null) return error.InvalidDid; var arena = std.heap.ArenaAllocator.init(allocator); errdefer arena.deinit(); const a = arena.allocator(); var scratch = std.heap.ArenaAllocator.init(allocator); defer scratch.deinit(); var context: Context = .{ .a = a, .scratch = &scratch, .options = options }; const revision = switch (options.format) { .star => try readStar(allocator, input, &context), .car => try readCar(allocator, input, &context), }; const coverage = context.counts.coverage; if (options.window.until > options.observed_at) context.counts.coverage = .{}; const follows = try context.follows.toOwnedSlice(a); std.mem.sort([]const u8, follows, {}, struct { fn less(_: void, left: []const u8, right: []const u8) bool { return std.mem.lessThan(u8, left, right); } }.less); return .{ .arena = arena, .value = .{ .did = try a.dupe(u8, options.did), .revision = revision, .window = options.window, .observed_at = options.observed_at, .records = context.records, .counts = context.counts, .diagnostics = context.diagnostics, .evidence = if (options.retain_evidence) .{ .records = try context.evidence.toOwnedSlice(a), .coverage = coverage, .observed_at = options.observed_at } else null, .follows = follows, } };}
fn readStar(a: std.mem.Allocator, input: *std.Io.Reader, context: *Context) ![]const u8 { var reader = try zat.star.PublicRepoReader.init(a, input, .{}); defer reader.deinit(); if (reader.subtree) return error.ExpectedFullRepository; const revision = try context.commit(reader.reader.metadata); while (try reader.next()) |entry| try context.add(entry.key, entry.record); return revision;}
fn readCar(a: std.mem.Allocator, input: *std.Io.Reader, context: *Context) ![]const u8 { var arena = std.heap.ArenaAllocator.init(a); defer arena.deinit(); const scratch = arena.allocator(); const bytes = try input.allocRemaining(scratch, .limited(context.options.max_car_bytes)); const car = try zat.car.readWithOptions(scratch, bytes, .{ .max_size = context.options.max_car_bytes, .max_blocks = 1000000 }); if (car.roots.len != 1) return error.InvalidRoots; const commit_bytes = zat.car.findBlock(car, car.roots[0].raw) orelse return error.MissingCommit; const commit = try zat.cbor.decodeAll(scratch, commit_bytes); const revision = try context.commit(commit); const root = commit.getCid("data") orelse return error.MissingRoot; var tree = try zat.mst.Mst.loadFromBlocks(scratch, car, root.raw); var walker: Walker = .{ .context = context, .car = car }; try tree.walk(.{ .ctx = &walker, .entryFn = Walker.entry }); return revision;}
const Walker = struct { context: *Context, car: zat.car.Car, fn entry(raw: *anyopaque, item: zat.mst.WalkEntry) anyerror!void { const self: *Walker = @ptrCast(@alignCast(raw)); const bytes = zat.car.findBlock(self.car, item.value.raw) orelse return error.MissingRecordBlock; try self.context.add(item.key, bytes); }};
const Context = struct { a: std.mem.Allocator, scratch: *std.heap.ArenaAllocator, options: Options, counts: activity.Counts = .{ .coverage = .all }, records: usize = 0, diagnostics: Diagnostics = .{}, evidence: std.ArrayList(activity.Evidence) = .empty, follows: std.ArrayList([]const u8) = .empty, seen: std.StringHashMapUnmanaged(void) = .empty,
fn commit(self: *Context, value: zat.cbor.Value) ![]const u8 { if (!std.mem.eql(u8, value.getString("did") orelse return error.MissingDid, self.options.did)) return error.IdentityMismatch; const revision = value.getString("rev") orelse return error.MissingRevision; if (zat.Tid.parse(revision) == null or value.getInt("version") != 3) return error.InvalidCommit; return self.a.dupe(u8, revision); }
fn add(self: *Context, key: []const u8, bytes: []const u8) !void { self.records += 1; const slash = std.mem.indexOfScalar(u8, key, '/') orelse return error.InvalidRecordKey; const nsid = key[0..slash]; const follow = self.options.collect_follows and std.mem.eql(u8, nsid, "app.bsky.graph.follow"); const index: ?usize = for (activity.collections, 0..) |collection, i| { if (std.mem.eql(u8, nsid, collection)) break i; } else null; if (!follow and index == null) return; _ = self.scratch.reset(.retain_capacity); const a = self.scratch.allocator(); const record = zat.cbor.decodeAll(a, bytes) catch |err| { if (err == error.OutOfMemory or follow) return err; self.invalidate(index.?, @errorName(err)); return; }; if (follow) { if (!std.mem.eql(u8, record.getString("$type") orelse return error.InvalidFollow, nsid)) return error.InvalidFollow; const subject = record.getString("subject") orelse return error.InvalidFollow; if (zat.Did.parse(subject) == null) return error.InvalidFollow; if (!self.seen.contains(subject)) { const owned = try self.a.dupe(u8, subject); try self.seen.put(self.a, owned, {}); try self.follows.append(self.a, owned); } return; } var evidence = activity.classify(record, nsid, "") catch |err| { self.invalidate(index.?, @errorName(err)); return; }; if (evidence.micros > self.options.observed_at) { self.invalidate(index.?, "FutureTimestamp"); return; } if (self.options.window.contains(evidence.micros)) self.counts.add(evidence.kind); if (self.options.retain_evidence) { evidence.uri = try std.fmt.allocPrint(self.a, "at://{s}/{s}", .{ self.options.did, key }); if (evidence.subject) |subject| evidence.subject = try self.a.dupe(u8, subject); if (evidence.root) |root| evidence.root = try self.a.dupe(u8, root); try self.evidence.append(self.a, evidence); } }
fn invalidate(self: *Context, index: usize, reason: []const u8) void { if (self.diagnostics.first_error == null) self.diagnostics.first_error = reason; switch (index) { 0 => { self.counts.coverage.post = false; self.counts.coverage.reply = false; self.diagnostics.invalid_posts += 1; }, 1 => { self.counts.coverage.repost = false; self.diagnostics.invalid_reposts += 1; }, 2 => { self.counts.coverage.like = false; self.diagnostics.invalid_likes += 1; }, else => unreachable, } }};
test "CAR and STAR produce the same owned evidence and offline answers" { const fixture = @import("test_repo.zig"); var input = try fixture.make(std.testing.allocator, false); defer input.deinit(); const until = zat.Datetime.parse("2026-09-02T00:00:00Z").?.micros; for ([_]Format{ .star, .car }) |format| { var reader: std.Io.Reader = .fixed(if (format == .star) input.value.star else input.value.car); var result = try read(std.testing.allocator, &reader, .{ .format = format, .did = fixture.did, .window = .{ .since = 0, .until = until }, .observed_at = until, .retain_evidence = true, .collect_follows = true }); defer result.deinit(); const value = result.value; try std.testing.expect(value.records == 5 and value.evidence.?.records.len == 4); try std.testing.expect(value.counts.complete()); try std.testing.expectEqual(@as(usize, 1), value.counts.post); try std.testing.expectEqual(@as(usize, 1), value.counts.reply); try std.testing.expectEqual(@as(usize, 1), value.counts.repost); try std.testing.expectEqual(@as(usize, 1), value.counts.like); try std.testing.expectEqualStrings("did:plc:other", value.follows[0]); const targets = try activity.interactions(std.testing.allocator, value.evidence.?, value.window); defer targets.deinit(std.testing.allocator); try std.testing.expect(targets.rows.len == 1 and targets.rows[0].likes == 1 and targets.rows[0].replies == 1 and targets.rows[0].reposts == 1); const later = try activity.count(value.evidence.?, .{ .since = until, .until = until }); try std.testing.expect(activity.pattern(later) == .no_retained_activity); }}
test "bad records lose only their own coverage, wrong identity and incomplete archives fail" { const fixture = @import("test_repo.zig"); var input = try fixture.make(std.testing.allocator, true); defer input.deinit(); for ([_]Format{ .star, .car }) |format| { const bytes = if (format == .star) input.value.star else input.value.car; var reader: std.Io.Reader = .fixed(bytes); const options: Options = .{ .format = format, .did = fixture.did, .window = .{ .since = 0, .until = std.math.maxInt(i64) }, .observed_at = std.math.maxInt(i64) }; var result = try read(std.testing.allocator, &reader, options); defer result.deinit(); try std.testing.expect(result.value.counts.coverage.post and result.value.counts.coverage.reply and result.value.counts.coverage.repost); try std.testing.expect(!result.value.counts.coverage.like); try std.testing.expect(result.value.diagnostics.invalid_likes == 1); var wrong = options; wrong.did = "did:plc:other"; reader = .fixed(bytes); try std.testing.expectError(error.IdentityMismatch, read(std.testing.allocator, &reader, wrong)); reader = .fixed(bytes[0 .. bytes.len - 1]); if (read(std.testing.allocator, &reader, options)) |owned| { var unexpected = owned; unexpected.deinit(); return error.ExpectedFailure; } else |_| {} const corrupted = try std.testing.allocator.dupe(u8, bytes); defer std.testing.allocator.free(corrupted); corrupted[corrupted.len - 1] ^= 1; reader = .fixed(corrupted); if (read(std.testing.allocator, &reader, options)) |owned| { var unexpected = owned; unexpected.deinit(); return error.ExpectedFailure; } else |_| {} }}
test "archive results clean up at every allocation failure" { const fixture = @import("test_repo.zig"); var input = try fixture.make(std.testing.allocator, false); defer input.deinit(); const Check = struct { fn run(a: std.mem.Allocator, bytes: []const u8) !void { var reader: std.Io.Reader = .fixed(bytes); var result = try read(a, &reader, .{ .format = .star, .did = fixture.did, .window = .{ .since = 0, .until = std.math.maxInt(i64) }, .observed_at = std.math.maxInt(i64), .retain_evidence = true, .collect_follows = true }); defer result.deinit(); } }; try std.testing.checkAllAllocationFailures(std.testing.allocator, Check.run, .{input.value.star});}