From 8fae81d576b9fe20fddcf89b7c54a59d2c37a261 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Mon, 29 Jun 2026 13:23:05 -0500 Subject: [PATCH] Fix full repo export reachability --- GETREPO_NOTES.md | 25 +++++--- README.md | 1 + bench/README.md | 14 +++++ build.zig.zon | 4 +- docs/operations.md | 4 ++ src/atproto/sync.zig | 16 +++++ src/core/config.zig | 9 +++ src/internal/cli.zig | 14 ++++- src/main.zig | 1 + src/storage/store.zig | 141 ++++++++++++++++++++++++++++++++++-------- 10 files changed, 191 insertions(+), 38 deletions(-) diff --git a/GETREPO_NOTES.md b/GETREPO_NOTES.md index 61d9230..411423f 100644 --- a/GETREPO_NOTES.md +++ b/GETREPO_NOTES.md @@ -7,21 +7,26 @@ ZDS is directly relevant because it serves `com.atproto.sync.getRepo` from `src/ Current shape: - ZDS serves raw CAR bytes and does not appear to add application-layer gzip, so it avoids the Node compression issue described in the operator thread. -- Full repo export is built from SQLite `repo_blocks`, ordered by CID. +- As of the 2026-06-29 getRepo pass, full repo export starts at the latest commit block and uses `zat.mst.collectReachableBlocks` to include only current MST and record blocks reachable from the commit's `data` CID. - Incremental export with `since` filters by `repo_rev`. - Repo writes already use lazy MST loading backed by `repo_blocks`, so ZDS is closer to a Hubble-style block/CAR-serving shape than an eager full-repo rebuild path. -- `subscribeRepos` has a connection cap, but `getRepo` does not appear to have route-specific backpressure. +- `subscribeRepos` has a connection cap. Full `getRepo` now has its own lightweight route-local concurrency cap so backup/backfill exports cannot consume unbounded handler slots. Operators can tune it with `ZDS_MAX_CONCURRENT_REPO_EXPORTS`; the default is `4`, matching Tranquil's current full-export default. Open questions: -- Should full `getRepo` walk the current commit/MST and include only reachable blocks, instead of exporting every retained block for the DID? -- Do retained old MST or record blocks make the current full export disclose deleted or stale record data? -- Should `getRepo` have route-specific concurrency or rate limits separate from normal public reads? -- Should the benchmark suite compare the current range-scan export against a reachability-based export for small, medium, and large repos? +- Should `HEAD com.atproto.sync.getRepo` stay semantically cheap, or should it compute/export-sized metadata? Today it is treated as a probe and still shares the same handler path. +- Should `getRepo` backpressure eventually become per-account/IP rate limiting in addition to process-local concurrency? +- Should the benchmark suite keep a historical range-scan comparison around, now that production full export is reachability-based? + +Done in this pass: + +1. Adopted `zat` 0.3.9 for `collectReachableBlocks`. +2. Added a regression test for update/delete followed by full `getRepo`; stale record CIDs are excluded from full export. +3. Changed full export from whole-DID `repo_blocks` scan to current-root reachability. +4. Added lightweight route backpressure for full `getRepo`, configurable with `ZDS_MAX_CONCURRENT_REPO_EXPORTS`. +5. Ran `just bench repo-size`; reachable full export measured about 1.4 ms for 100 records, 16.8 ms for 1,000 records, and 95.7 ms for 5,000 records on the local synthetic benchmark. Suggested next pass: -1. Add a regression test for delete/update followed by full `getRepo`, checking whether stale record blocks are exported. -2. If stale blocks are visible, change full export to current-root reachability rather than whole-DID block scan. -3. Add lightweight route backpressure for full `getRepo`. -4. Re-run `bench repo-size` and keep raw range-scan numbers separate from reachable-export numbers. +1. Decide whether `GETREPO_NOTES.md` should stay as a root handoff note or move into `docs/`. +2. Add a larger deliberate `just bench repo-large 100000` run when we want stress numbers comparable to large public repos. diff --git a/README.md b/README.md index 97232b9..8585198 100644 --- a/README.md +++ b/README.md @@ -79,6 +79,7 @@ ZDS_BLOB_UPLOAD_LIMIT=100000000 \ ZDS_BLOBSTORE_PATH=/var/lib/zds/blobs \ ZDS_HANDLE_DOMAINS='.example.com,example.com' \ ZDS_CRAWLERS='https://bsky.network,https://vsky.network' \ +ZDS_MAX_CONCURRENT_REPO_EXPORTS=4 \ ZDS_PLC_ROTATION_KEY='64-hex-secp256k1-secret-or-private-multikey' \ ZDS_RECOVERY_DID_KEY='did:key:optionalRecoveryKey' \ ZDS_JWT_SECRET='at-least-32-random-bytes-here' \ diff --git a/bench/README.md b/bench/README.md index 8d20123..9ad31e1 100644 --- a/bench/README.md +++ b/bench/README.md @@ -107,6 +107,20 @@ magnitude for large-repo stress. Live network corpora should be captured into fixtures before they become stable benchmarks. Do not make routine CI or pre-commit checks depend on live PDS availability. +June 29, 2026 local reachable-export check after switching full +`com.atproto.sync.getRepo` from whole-DID `repo_blocks` scans to current-root +reachability via `zat.mst.collectReachableBlocks`: + +| repo tier | records | list | full CAR export | +|---|---:|---:|---:| +| small | 100 | 13.4k ops/s, 15.0 ms | 3.6k ops/s, 72.1 MB/s, 1.4 ms | +| medium | 1,000 | 12.4k ops/s, 16.1 ms | 297 ops/s, 59.6 MB/s, 16.8 ms | +| large | 5,000 | 11.9k ops/s, 16.8 ms | 52 ops/s, 53.1 MB/s, 95.7 ms | + +These are synthetic local trend numbers, not cross-PDS comparisons. The old +whole-DID range-scan export path is intentionally no longer the production full +export behavior because it could include retained stale blocks. + ## methodology - Build ZDS with `-Doptimize=ReleaseFast`. diff --git a/build.zig.zon b/build.zig.zon index 343d5c7..b6c64cc 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,8 +5,8 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.3.8.tar.gz", - .hash = "zat-0.3.8-5PuC7glgCgCVn18R0k6AuXiNok-jjqDRGgaBFj5JsI_N", + .url = "https://tangled.org/zat.dev/zat/archive/v0.3.9.tar.gz", + .hash = "zat-0.3.9-5PuC7v6aCgDbfEdNEUxbSbbpq27V5upmAsa6J-kqdCXx", }, .zqlite = .{ .url = "git+https://github.com/karlseguin/zqlite.zig?ref=master#05a88d6758753e1c63fdd45b211dde2057094b0c", diff --git a/docs/operations.md b/docs/operations.md index 28a0f64..2e79af6 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -105,6 +105,10 @@ shape, and production troubleshooting notes. - `ZDS_BLOB_UPLOAD_LIMIT`: upload body limit. Default: `100000000`. - `ZDS_BLOBSTORE_PATH`: disk blobstore root. - `ZDS_CRAWLERS`: comma-separated relay crawl targets. +- `ZDS_MAX_CONCURRENT_REPO_EXPORTS`: maximum concurrent full + `com.atproto.sync.getRepo` exports. Default: `4`. Incremental + `getRepo?since=...` exports are not intended to consume this full-export + backpressure budget. - `ZDS_ADMIN_TOKEN`: bearer token for privileged local PDS administration, including invite-code minting. - `ZDS_INVITE_REQUIRED`: set to `true` to require invite codes for account diff --git a/src/atproto/sync.zig b/src/atproto/sync.zig index 62ed753..deae482 100644 --- a/src/atproto/sync.zig +++ b/src/atproto/sync.zig @@ -13,6 +13,7 @@ const max_subscribe_repos_connections = 32; const notify_threshold_ms = 20 * 60 * 1000; var subscribe_repos_connections: usize = 0; +var get_repo_connections: usize = 0; var crawler_last_notified_ms: std.atomic.Value(i64) = .init(0); var crawler_notify_in_flight: std.atomic.Value(bool) = .init(false); @@ -207,6 +208,21 @@ pub fn getRepo(request: *http_api.Request) !void { var since_buf: [256]u8 = undefined; const since = http_api.queryParam(request.url.raw, "since", &since_buf); try requirePublicRepoAvailable(request, did); + + const full_export = since == null; + if (full_export) { + const active = @atomicRmw(usize, &get_repo_connections, .Add, 1, .monotonic); + const max_connections = config.maxConcurrentRepoExports(); + if (active >= max_connections) { + _ = @atomicRmw(usize, &get_repo_connections, .Sub, 1, .monotonic); + log.debug("sync getRepo rejected too_many_connections active={d} max={d} did={s}\n", .{ active + 1, max_connections, did }); + return http_api.xrpcError(request, .too_many_requests, "RateLimitExceeded", "too many full getRepo exports"); + } + } + defer { + if (full_export) _ = @atomicRmw(usize, &get_repo_connections, .Sub, 1, .monotonic); + } + const body = store.writeRepoCarSince(allocator, did, since) catch { return http_api.xrpcError(request, .not_found, "RepoNotFound", "Repo not found"); }; diff --git a/src/core/config.zig b/src/core/config.zig index c0224e1..e3d1ce0 100644 --- a/src/core/config.zig +++ b/src/core/config.zig @@ -14,6 +14,7 @@ var blob_upload_limit_value: usize = 100_000_000; var blobstore_path_value: []const u8 = "dev/blobs"; var handle_domains_value: []const u8 = ".test"; var crawlers_value: []const u8 = "https://bsky.network,https://vsky.network"; +var max_concurrent_repo_exports_value: usize = 4; var proxy_service_did_value: []const u8 = "did:web:api.bsky.app"; var proxy_service_id_value: []const u8 = "bsky_appview"; var proxy_service_url_value: []const u8 = "https://api.bsky.app"; @@ -91,6 +92,10 @@ pub fn crawlers() []const u8 { return crawlers_value; } +pub fn maxConcurrentRepoExports() usize { + return max_concurrent_repo_exports_value; +} + pub fn proxyServiceDid() []const u8 { return proxy_service_did_value; } @@ -187,6 +192,10 @@ pub fn setCrawlers(value: []const u8) void { crawlers_value = value; } +pub fn setMaxConcurrentRepoExports(value: usize) void { + max_concurrent_repo_exports_value = value; +} + pub fn setProxyServiceDid(value: []const u8) void { proxy_service_did_value = value; } diff --git a/src/internal/cli.zig b/src/internal/cli.zig index 77b9260..f1686ca 100644 --- a/src/internal/cli.zig +++ b/src/internal/cli.zig @@ -18,6 +18,7 @@ pub const Options = struct { blobstore_path: ?[]const u8 = null, handle_domains: ?[]const u8 = null, crawlers: ?[]const u8 = null, + max_concurrent_repo_exports: ?usize = null, proxy_service_did: ?[]const u8 = null, proxy_service_id: ?[]const u8 = null, proxy_service_url: ?[]const u8 = null, @@ -35,6 +36,7 @@ pub const ParseError = error{ MissingBlobUploadLimit, MissingBlobstorePath, MissingCrawlers, + MissingMaxConcurrentRepoExports, MissingProxyServiceDid, MissingProxyServiceId, MissingProxyServiceUrl, @@ -78,6 +80,7 @@ pub fn parse(init: std.process.Init) ParseError!Options { .blobstore_path = env("ZDS_BLOBSTORE_PATH"), .handle_domains = env("ZDS_HANDLE_DOMAINS"), .crawlers = env("ZDS_CRAWLERS"), + .max_concurrent_repo_exports = try envUsize("ZDS_MAX_CONCURRENT_REPO_EXPORTS"), .proxy_service_did = env("ZDS_PROXY_SERVICE_DID"), .proxy_service_id = env("ZDS_PROXY_SERVICE_ID"), .proxy_service_url = env("ZDS_PROXY_SERVICE_URL"), @@ -108,7 +111,8 @@ pub fn usage() void { \\ [--plc-rotation-key KEY] [--recovery-did-key DIDKEY] \\ [--mail-provider comail|resend] [--email-from ADDRESS] \\ [--blob-upload-limit BYTES] [--blobstore-path PATH] - \\ [--handle-domains DOMAINS] [--crawlers URLS] [--jwt-secret SECRET] + \\ [--handle-domains DOMAINS] [--crawlers URLS] + \\ [--max-concurrent-repo-exports N] [--jwt-secret SECRET] \\ [--dpop-secret SECRET] \\ [--proxy-service-did DID] [--proxy-service-id ID] [--proxy-service-url URL] \\ [--admin-token TOKEN] [--invite-required] @@ -199,6 +203,10 @@ fn parseSplitArg(options: *Options, arg: []const u8, args: *std.process.Args.Ite options.crawlers = args.next() orelse return error.MissingCrawlers; return true; } + if (std.mem.eql(u8, arg, "--max-concurrent-repo-exports")) { + options.max_concurrent_repo_exports = try std.fmt.parseInt(usize, args.next() orelse return error.MissingMaxConcurrentRepoExports, 10); + return true; + } if (std.mem.eql(u8, arg, "--proxy-service-did")) { options.proxy_service_did = args.next() orelse return error.MissingProxyServiceDid; return true; @@ -245,6 +253,10 @@ fn parseJoinedArg(options: *Options, arg: []const u8) ParseError!bool { options.blob_upload_limit = try std.fmt.parseInt(usize, arg["--blob-upload-limit=".len..], 10); return true; } + if (std.mem.startsWith(u8, arg, "--max-concurrent-repo-exports=")) { + options.max_concurrent_repo_exports = try std.fmt.parseInt(usize, arg["--max-concurrent-repo-exports=".len..], 10); + return true; + } return false; } diff --git a/src/main.zig b/src/main.zig index 8d12732..7da352b 100644 --- a/src/main.zig +++ b/src/main.zig @@ -39,6 +39,7 @@ pub fn main(init: std.process.Init) !void { if (options.blobstore_path) |value| zds.core.config.setBlobstorePath(value); if (options.handle_domains) |value| zds.core.config.setHandleDomains(value); if (options.crawlers) |value| zds.core.config.setCrawlers(value); + if (options.max_concurrent_repo_exports) |value| zds.core.config.setMaxConcurrentRepoExports(value); if (options.proxy_service_did) |value| zds.core.config.setProxyServiceDid(value); if (options.proxy_service_id) |value| zds.core.config.setProxyServiceId(value); if (options.proxy_service_url) |value| zds.core.config.setProxyServiceUrl(value); diff --git a/src/storage/store.zig b/src/storage/store.zig index e4cfc54..c7a91db 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -427,6 +427,7 @@ const BlobRef = struct { const RepoBlockReader = struct { allocator: std.mem.Allocator, did: []const u8, + locked: bool = false, fn reader(self: *RepoBlockReader) zat.mst.BlockReader { return .{ @@ -439,24 +440,32 @@ const RepoBlockReader = struct { const self: *RepoBlockReader = @ptrCast(@alignCast(ctx)); const cid = try cidText(self.allocator, cid_raw); + if (self.locked) { + return repoBlockDataLocked(self.allocator, self.did, cid); + } + db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); - const row = try conn.row( - \\SELECT data - \\FROM repo_blocks - \\WHERE did = ? AND cid = ? - \\LIMIT 1 - , .{ self.did, cid }); - if (row == null) return null; - defer row.?.deinit(); - - const data = row.?.nullableBlob(0) orelse return null; - return try self.allocator.dupe(u8, data); + return repoBlockDataLocked(self.allocator, self.did, cid); } }; +fn repoBlockDataLocked(allocator: std.mem.Allocator, did: []const u8, cid: []const u8) !?[]const u8 { + const row = try conn.row( + \\SELECT data + \\FROM repo_blocks + \\WHERE did = ? AND cid = ? + \\LIMIT 1 + , .{ did, cid }); + if (row == null) return null; + defer row.?.deinit(); + + const data = row.?.nullableBlob(0) orelse return null; + return try allocator.dupe(u8, data); +} + var conn: zqlite.Conn = undefined; var initialized = false; var store_io: Io = undefined; @@ -2684,6 +2693,8 @@ pub fn writeRepoCar(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { } pub fn writeRepoCarSince(allocator: std.mem.Allocator, did: []const u8, since: ?[]const u8) ![]const u8 { + if (since == null) return writeRepoCarFull(allocator, did); + db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); try requireInitialized(); @@ -2692,20 +2703,13 @@ pub fn writeRepoCarSince(allocator: std.mem.Allocator, did: []const u8, since: ? const root_raw = try zat.multibase.base32lower.decode(allocator, root.cid[1..]); const car_root = zat.cbor.Cid{ .raw = root_raw }; - var rows = if (since) |since_rev| - try conn.rows( - \\SELECT cid, data - \\FROM repo_blocks - \\WHERE did = ? AND (repo_rev IS NULL OR repo_rev > ?) - \\ORDER BY repo_rev DESC, cid DESC - , .{ did, since_rev }) - else - try conn.rows( - \\SELECT cid, data - \\FROM repo_blocks - \\WHERE did = ? - \\ORDER BY cid ASC - , .{did}); + const since_rev = since.?; + var rows = try conn.rows( + \\SELECT cid, data + \\FROM repo_blocks + \\WHERE did = ? AND (repo_rev IS NULL OR repo_rev > ?) + \\ORDER BY repo_rev DESC, cid DESC + , .{ did, since_rev }); defer rows.deinit(); var blocks: std.ArrayList(zat.car.Block) = .empty; @@ -2726,6 +2730,37 @@ pub fn writeRepoCarSince(allocator: std.mem.Allocator, did: []const u8, since: ? return zat.car.writeAlloc(allocator, c); } +fn writeRepoCarFull(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try requireInitialized(); + + const root = try latestCommitRawLocked(allocator, did) orelse return Error.RepoNotFound; + var blocks: std.ArrayList(zat.car.Block) = .empty; + try blocks.append(allocator, .{ + .cid_raw = root.commit_cid_raw, + .data = root.commit_data, + }); + + var repo_block_reader = RepoBlockReader{ + .allocator = allocator, + .did = did, + .locked = true, + }; + try zat.mst.collectReachableBlocks( + allocator, + root.data_cid_raw, + repo_block_reader.reader(), + &blocks, + .{ .include_records = true }, + ); + + return zat.car.writeAlloc(allocator, .{ + .roots = &.{.{ .raw = root.commit_cid_raw }}, + .blocks = blocks.items, + }); +} + pub fn writeRepoListJson(allocator: std.mem.Allocator, cursor: ?[]const u8, limit: usize) ![]const u8 { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); @@ -5878,6 +5913,62 @@ test "getRepo since filters repo blocks by revision" { try std.testing.expectEqualStrings(second_result.commit.cid, try cidText(allocator, empty_diff_car.roots[0].raw)); } +fn carContainsCid(allocator: std.mem.Allocator, c: zat.car.Car, cid: []const u8) !bool { + for (c.blocks) |block| { + if (std.mem.eql(u8, cid, try cidText(allocator, block.cid_raw))) return true; + } + return false; +} + +test "full getRepo exports only current reachable record blocks" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + try init(std.Options.debug_io, ":memory:"); + defer close(); + + const account = try createAccount( + allocator, + "reachable.test", + "reachable@test.com", + "password", + "did:plc:reachablecheck", + true, + ); + + const first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"first\"}", .{}); + defer first.deinit(); + const first_result = try applyWritesWithOptions(allocator, account, &.{.{ .create = .{ + .collection = "app.bsky.feed.post", + .rkey = "3zreach", + .value = first.value, + } }}, .{}); + const first_cid = first_result.records[0].cid; + + const second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"app.bsky.feed.post\",\"text\":\"second\"}", .{}); + defer second.deinit(); + const second_result = try applyWritesWithOptions(allocator, account, &.{.{ .update = .{ + .collection = "app.bsky.feed.post", + .rkey = "3zreach", + .value = second.value, + } }}, .{}); + const second_cid = second_result.records[0].cid; + + const updated_car = try zat.car.read(allocator, try writeRepoCar(allocator, account.did)); + try std.testing.expect(!try carContainsCid(allocator, updated_car, first_cid)); + try std.testing.expect(try carContainsCid(allocator, updated_car, second_cid)); + + _ = try applyWritesWithOptions(allocator, account, &.{.{ .delete = .{ + .collection = "app.bsky.feed.post", + .rkey = "3zreach", + } }}, .{}); + + const deleted_car = try zat.car.read(allocator, try writeRepoCar(allocator, account.did)); + try std.testing.expect(!try carContainsCid(allocator, deleted_car, first_cid)); + try std.testing.expect(!try carContainsCid(allocator, deleted_car, second_cid)); +} + test "permissioned spaces store self-owned records outside public repo" { var arena = std.heap.ArenaAllocator.init(std.testing.allocator); defer arena.deinit(); -- 2.51.2