diff --git a/bench/README.md b/bench/README.md index 2679bc5..627545d 100644 --- a/bench/README.md +++ b/bench/README.md @@ -84,26 +84,27 @@ This table only includes operation/count pairs measured in both implementations. The operations are matched to Tranquil's metastore bench: write one record, look up one record CID, and list records from a seeded repo. -Summary: Tranquil is faster on CID lookups, by about 1.6x at 10 callers and -1.6x at 100 callers. ZDS is slower on 10-caller writes, at about 0.7x -Tranquil's throughput, but faster on 100-caller writes, at about 1.6x -Tranquil's throughput. ZDS is faster on list-record reads in these runs, by -about 1.4x at both 10 and 100 callers. The next read-path question is why -Tranquil's CID lookup is still ahead while ZDS's list path is already ahead. +Summary: Tranquil is faster on the read rows after ZDS moved record materialize +paths to decode DAG-CBOR from `repo_blocks`: about 1.5-1.8x faster for CID +lookup and about 1.6x faster for list records. Tranquil is also faster on +10-caller writes, by about 2.8x. ZDS is faster on the 100-caller write row, by +about 1.3x, which is the main sign that the sharded write lane is helping under +contention. | operation | callers | ops | zds | tranquil | |---|---:|---:|---:|---:| -| apply/write | 10 | 10000 | 299 ops/s, p95 84 ms | 436 ops/s, p95 36.0 ms | -| apply/write | 100 | 20000 | 1502 ops/s, p95 56 ms | 944 ops/s, p95 119 ms | -| get record CID | 10 | 10000 | 355k ops/s, p95 5 us | 571k ops/s, p95 31 us | -| get record CID | 100 | 20000 | 349k ops/s, p95 6 us | 571k ops/s, p95 244 us | -| list records | 10 | 10000 | 52k ops/s, p95 20 us | 38k ops/s, p95 309 us | -| list records | 100 | 20000 | 51k ops/s, p95 22 us | 37k ops/s, p95 3.0 ms | +| apply/write | 10 | 10000 | 282 ops/s, p95 93.8 ms | 780 ops/s, p95 14.1 ms | +| apply/write | 100 | 20000 | 1357 ops/s, p95 66.7 ms | 1007 ops/s, p95 109 ms | +| get record CID | 10 | 10000 | 296k ops/s, p95 6 us | 455k ops/s, p95 43 us | +| get record CID | 100 | 20000 | 294k ops/s, p95 405 us | 537k ops/s, p95 233 us | +| list records | 10 | 10000 | 22k ops/s, p95 52 us | 36k ops/s, p95 310 us | +| list records | 100 | 20000 | 23k ops/s, p95 49 us | 36k ops/s, p95 2.9 ms | Tranquil's read row is `get_record_cid`, so ZDS reports `get-cid` for that comparison. Full `com.atproto.repo.getRecord` is a separate local probe because -it materializes record JSON and response fields that Tranquil's metastore bench -does not request. +it joins the record index to `repo_blocks`, decodes DAG-CBOR, and renders the +XRPC JSON response shape. Tranquil's metastore bench intentionally measures a +metadata lookup for the current CID, not record body materialization. ## local probes @@ -111,8 +112,8 @@ These rows are useful for ZDS tuning but are not direct Tranquil comparisons. | operation | callers | ops | zds | |---|---:|---:|---:| -| full getRecord | 10 | 10000 | 56k ops/s, p95 65 us | -| full getRecord | 100 | 20000 | 52k ops/s, p95 7.2 ms | +| full getRecord | 10 | 10000 | 40k ops/s, p95 352 us | +| full getRecord | 100 | 20000 | 39k ops/s, p95 8.2 ms | ## write profile diff --git a/src/atproto/repo.zig b/src/atproto/repo.zig index 6c1c2dd..a616087 100644 --- a/src/atproto/repo.zig +++ b/src/atproto/repo.zig @@ -461,7 +461,6 @@ fn collectImportedRecords( for (node.entries.items) |entry| { const block = zat.car.findBlock(repo_car, entry.value.raw) orelse return error.MissingRecordBlock; const record_value = try zat.cbor.decodeAll(allocator, block); - const record_json = try cborToJson(allocator, record_value); var blobs: std.ArrayList([]const u8) = .empty; try collectBlobCids(allocator, record_value, &blobs); const slash = std.mem.indexOfScalar(u8, entry.key, '/') orelse return error.InvalidRepoPath; @@ -469,7 +468,6 @@ fn collectImportedRecords( .collection = try allocator.dupe(u8, entry.key[0..slash]), .rkey = try allocator.dupe(u8, entry.key[slash + 1 ..]), .cid = try cidString(allocator, entry.value.raw), - .value_json = record_json, .blob_cids = try blobs.toOwnedSlice(allocator), }); try collectFromChild(allocator, repo_car, entry.right, out); @@ -489,47 +487,6 @@ fn collectFromChild( } } -fn cborToJson(allocator: std.mem.Allocator, value: zat.cbor.Value) ![]const u8 { - var out: std.Io.Writer.Allocating = .init(allocator); - defer out.deinit(); - try writeCborJson(allocator, &out.writer, value); - return out.toOwnedSlice(); -} - -fn writeCborJson(allocator: std.mem.Allocator, writer: anytype, value: zat.cbor.Value) !void { - switch (value) { - .unsigned => |v| try writer.print("{d}", .{v}), - .negative => |v| try writer.print("{d}", .{v}), - .bytes => |bytes| { - const encoded = try allocator.alloc(u8, std.base64.standard.Encoder.calcSize(bytes.len)); - defer allocator.free(encoded); - _ = std.base64.standard.Encoder.encode(encoded, bytes); - try writer.print("{{\"$bytes\":{f}}}", .{std.json.fmt(encoded, .{})}); - }, - .text => |text| try writer.print("{f}", .{std.json.fmt(text, .{})}), - .array => |items| { - try writer.writeByte('['); - for (items, 0..) |item, idx| { - if (idx != 0) try writer.writeByte(','); - try writeCborJson(allocator, writer, item); - } - try writer.writeByte(']'); - }, - .map => |entries| { - try writer.writeByte('{'); - for (entries, 0..) |entry, idx| { - if (idx != 0) try writer.writeByte(','); - try writer.print("{f}:", .{std.json.fmt(entry.key, .{})}); - try writeCborJson(allocator, writer, entry.value); - } - try writer.writeByte('}'); - }, - .boolean => |v| try writer.print("{}", .{v}), - .null => try writer.writeAll("null"), - .cid => |cid| try writer.print("{{\"$link\":{f}}}", .{std.json.fmt(try cidString(allocator, cid.raw), .{})}), - } -} - fn collectBlobCids(allocator: std.mem.Allocator, value: zat.cbor.Value, out: *std.ArrayList([]const u8)) !void { switch (value) { .map => |entries| { diff --git a/src/internal/cbor_json.zig b/src/internal/cbor_json.zig new file mode 100644 index 0000000..d97aadb --- /dev/null +++ b/src/internal/cbor_json.zig @@ -0,0 +1,46 @@ +const std = @import("std"); +const zat = @import("zat"); + +pub fn writeAlloc(allocator: std.mem.Allocator, value: zat.cbor.Value) ![]const u8 { + var out: std.Io.Writer.Allocating = .init(allocator); + defer out.deinit(); + try writeValue(allocator, &out.writer, value); + return out.toOwnedSlice(); +} + +fn writeValue(allocator: std.mem.Allocator, writer: anytype, value: zat.cbor.Value) !void { + switch (value) { + .unsigned => |v| try writer.print("{d}", .{v}), + .negative => |v| try writer.print("{d}", .{v}), + .bytes => |bytes| { + const encoded = try allocator.alloc(u8, std.base64.standard.Encoder.calcSize(bytes.len)); + defer allocator.free(encoded); + _ = std.base64.standard.Encoder.encode(encoded, bytes); + try writer.print("{{\"$bytes\":{f}}}", .{std.json.fmt(encoded, .{})}); + }, + .text => |text| try writer.print("{f}", .{std.json.fmt(text, .{})}), + .array => |items| { + try writer.writeByte('['); + for (items, 0..) |item, idx| { + if (idx != 0) try writer.writeByte(','); + try writeValue(allocator, writer, item); + } + try writer.writeByte(']'); + }, + .map => |entries| { + try writer.writeByte('{'); + for (entries, 0..) |entry, idx| { + if (idx != 0) try writer.writeByte(','); + try writer.print("{f}:", .{std.json.fmt(entry.key, .{})}); + try writeValue(allocator, writer, entry.value); + } + try writer.writeByte('}'); + }, + .boolean => |v| try writer.print("{}", .{v}), + .null => try writer.writeAll("null"), + .cid => |cid| try writer.print( + "{{\"$link\":{f}}}", + .{std.json.fmt(try zat.multibase.base32lower.encode(allocator, cid.raw), .{})}, + ), + } +} diff --git a/src/storage/store.zig b/src/storage/store.zig index 06cdefc..0aa8dfd 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -2,6 +2,7 @@ const std = @import("std"); const atid = @import("../core/atid.zig"); const auth = @import("../auth/tokens.zig"); const blobstore = @import("blobstore.zig"); +const cbor_json = @import("../internal/cbor_json.zig"); const eventlog = @import("eventlog.zig"); const sharded_locks = @import("../internal/sharded_locks.zig"); const zat = @import("zat"); @@ -13,6 +14,7 @@ pub const Error = error{ InvalidRecordKey, InvalidRecordType, MissingRecord, + MissingRecordBlock, RepoNotFound, InvalidRepoPath, InvalidDagCbor, @@ -115,7 +117,6 @@ pub const ImportedRecord = struct { collection: []const u8, rkey: []const u8, cid: []const u8, - value_json: []const u8, blob_cids: []const []const u8, }; @@ -309,15 +310,17 @@ pub fn profileAvatarCid(allocator: std.mem.Allocator, did: []const u8) !?[]const try requireInitialized(); const row = try conn.row( - \\SELECT value_json - \\FROM records - \\WHERE did = ? AND collection = 'app.bsky.actor.profile' AND rkey = 'self' + \\SELECT rb.data + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.did = ? AND r.collection = 'app.bsky.actor.profile' AND r.rkey = 'self' \\LIMIT 1 , .{did}); if (row == null) return null; defer row.?.deinit(); - var parsed = std.json.parseFromSlice(std.json.Value, allocator, row.?.text(0), .{}) catch return null; + const profile_json = recordJsonFromBlock(allocator, row.?.nullableBlob(0) orelse "") catch return null; + var parsed = std.json.parseFromSlice(std.json.Value, allocator, profile_json, .{}) catch return null; defer parsed.deinit(); if (parsed.value != .object) return null; const avatar = parsed.value.object.get("avatar") orelse return null; @@ -823,15 +826,14 @@ fn applyWritesMeasured(allocator: std.mem.Allocator, account: auth.Account, ops: for (records.items) |record| { const uri = try record.uri(allocator); try conn.exec( - \\INSERT INTO records (did, collection, rkey, uri, cid, value_json, rev, seq) - \\VALUES (?, ?, ?, ?, ?, ?, ?, ?) + \\INSERT INTO records (did, collection, rkey, uri, cid, rev, seq) + \\VALUES (?, ?, ?, ?, ?, ?, ?) \\ON CONFLICT(did, collection, rkey) DO UPDATE SET \\ uri = excluded.uri, \\ cid = excluded.cid, - \\ value_json = excluded.value_json, \\ rev = excluded.rev, \\ seq = excluded.seq - , .{ record.did, record.collection, record.rkey, uri, record.cid, record.value_json, record.rev, @as(i64, @intCast(seq)) }); + , .{ record.did, record.collection, record.rkey, uri, record.cid, record.rev, @as(i64, @intCast(seq)) }); try conn.exec("DELETE FROM expected_blobs WHERE record_uri = ?", .{uri}); } for (blob_refs.items) |ref| { @@ -945,9 +947,10 @@ pub fn get(did: []const u8, collection: []const u8, rkey: []const u8) ?Record { requireInitialized() catch return null; const row = conn.row( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE did = ? AND collection = ? AND rkey = ? + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.did = ? AND r.collection = ? AND r.rkey = ? , .{ did, collection, rkey }) catch return null; if (row == null) return null; defer row.?.deinit(); @@ -982,9 +985,10 @@ pub fn listRecentRecords(allocator: std.mem.Allocator, limit: usize) ![]Record { try requireInitialized(); var rows = try conn.rows( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\ORDER BY seq DESC + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\ORDER BY r.seq DESC \\LIMIT ? , .{@as(i64, @intCast(if (limit == 0) 100 else limit))}); defer rows.deinit(); @@ -1008,10 +1012,11 @@ pub fn listRecords( try requireInitialized(); var rows = try conn.rows( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE did = ? AND collection = ? - \\ORDER BY seq DESC + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.did = ? AND r.collection = ? + \\ORDER BY r.seq DESC \\LIMIT ? , .{ did, collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit(); @@ -1035,17 +1040,21 @@ pub fn listRecordsContaining( try requireInitialized(); var rows = try conn.rows( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE collection = ? AND instr(value_json, ?) - \\ORDER BY seq DESC + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.collection = ? + \\ORDER BY r.seq DESC \\LIMIT ? - , .{ collection, needle, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); + , .{ collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit(); var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { - try records.append(allocator, try recordFromRow(row, allocator)); + const record = try recordFromRow(row, allocator); + if (std.mem.indexOf(u8, record.value_json, needle) != null) { + try records.append(allocator, record); + } } if (rows.err) |err| return err; return records.toOwnedSlice(allocator); @@ -1063,17 +1072,21 @@ pub fn listRecordsByDidContaining( try requireInitialized(); var rows = try conn.rows( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE did = ? AND collection = ? AND instr(value_json, ?) - \\ORDER BY seq DESC + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.did = ? AND r.collection = ? + \\ORDER BY r.seq DESC \\LIMIT ? - , .{ did, collection, needle, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); + , .{ did, collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit(); var records: std.ArrayList(Record) = .empty; while (rows.next()) |row| { - try records.append(allocator, try recordFromRow(row, allocator)); + const record = try recordFromRow(row, allocator); + if (std.mem.indexOf(u8, record.value_json, needle) != null) { + try records.append(allocator, record); + } } if (rows.err) |err| return err; return records.toOwnedSlice(allocator); @@ -1093,9 +1106,10 @@ fn getByUriLocked(uri: []const u8) ?Record { fn getByStoredUriLocked(uri: []const u8) ?Record { const row = conn.row( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE uri = ? + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.uri = ? , .{uri}) catch return null; if (row == null) return null; defer row.?.deinit(); @@ -1116,10 +1130,23 @@ pub fn countSubject(collection: []const u8, subject_did: []const u8) usize { db_mutex.lockUncancelable(store_io); defer db_mutex.unlock(store_io); requireInitialized() catch return 0; - const row = conn.row("SELECT count(*) FROM records WHERE collection = ? AND instr(value_json, ?) > 0", .{ collection, subject_did }) catch return 0; - if (row == null) return 0; - defer row.?.deinit(); - return @intCast(row.?.int(0)); + var rows = conn.rows( + \\SELECT rb.data + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.collection = ? + , .{collection}) catch return 0; + defer rows.deinit(); + + var count_result: usize = 0; + var arena = std.heap.ArenaAllocator.init(std.heap.page_allocator); + defer arena.deinit(); + while (rows.next()) |row| { + _ = arena.reset(.retain_capacity); + const value = recordJsonFromBlock(arena.allocator(), row.nullableBlob(0) orelse "") catch continue; + if (std.mem.indexOf(u8, value, subject_did) != null) count_result += 1; + } + return count_result; } pub fn listCollectionsJson(allocator: std.mem.Allocator, did: []const u8) ![]const u8 { @@ -1390,9 +1417,9 @@ pub fn importRepo( const uri = try std.fmt.allocPrint(std.heap.page_allocator, "at://{s}/{s}/{s}", .{ account.did, record.collection, record.rkey }); defer std.heap.page_allocator.free(uri); try conn.exec( - \\INSERT INTO records (did, collection, rkey, uri, cid, value_json, rev, seq) - \\VALUES (?, ?, ?, ?, ?, ?, ?, ?) - , .{ account.did, record.collection, record.rkey, uri, record.cid, record.value_json, rev, @as(i64, @intCast(seq)) }); + \\INSERT INTO records (did, collection, rkey, uri, cid, rev, seq) + \\VALUES (?, ?, ?, ?, ?, ?, ?) + , .{ account.did, record.collection, record.rkey, uri, record.cid, rev, @as(i64, @intCast(seq)) }); for (record.blob_cids) |blob_cid| { try conn.exec( \\INSERT INTO expected_blobs (blob_cid, record_uri) @@ -1587,10 +1614,11 @@ pub fn writeListJson( try requireInitialized(); var rows = try conn.rows( - \\SELECT did, collection, rkey, cid, value_json, rev, seq - \\FROM records - \\WHERE did = ? AND collection = ? - \\ORDER BY seq DESC + \\SELECT r.did, r.collection, r.rkey, r.cid, rb.data, r.rev, r.seq + \\FROM records r + \\JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE r.did = ? AND r.collection = ? + \\ORDER BY r.seq DESC \\LIMIT ? , .{ did, collection, @as(i64, @intCast(if (limit == 0) 100 else limit)) }); defer rows.deinit(); @@ -1631,6 +1659,59 @@ fn migrate() !void { try migrateBlobTable(); try migrateAppPreferencesTable(); inline for (post_schema_statements) |sql| try conn.execNoArgs(sql); + try migrateRecordsTable(); + try validateRecordBlocksPresent(); +} + +fn migrateRecordsTable() !void { + if (!try recordsTableHasValueJsonColumn()) return; + try conn.execNoArgs("DROP TABLE IF EXISTS records_next"); + try conn.execNoArgs( + \\CREATE TABLE IF NOT EXISTS records_next ( + \\ did TEXT NOT NULL REFERENCES accounts(did) ON DELETE CASCADE, + \\ collection TEXT NOT NULL, + \\ rkey TEXT NOT NULL, + \\ uri TEXT NOT NULL UNIQUE, + \\ cid TEXT NOT NULL, + \\ rev TEXT NOT NULL, + \\ seq INTEGER NOT NULL, + \\ PRIMARY KEY (did, collection, rkey) + \\) + ); + try conn.execNoArgs( + \\INSERT INTO records_next (did, collection, rkey, uri, cid, rev, seq) + \\SELECT did, collection, rkey, uri, cid, rev, seq + \\FROM records + ); + try conn.execNoArgs("PRAGMA foreign_keys = OFF"); + errdefer conn.execNoArgs("PRAGMA foreign_keys = ON") catch {}; + try conn.execNoArgs("DROP TABLE records"); + try conn.execNoArgs("ALTER TABLE records_next RENAME TO records"); + try conn.execNoArgs("PRAGMA foreign_keys = ON"); + try conn.execNoArgs("CREATE INDEX IF NOT EXISTS records_collection_idx ON records (did, collection, seq DESC)"); + try conn.execNoArgs("CREATE INDEX IF NOT EXISTS records_cid_idx ON records (cid)"); +} + +fn recordsTableHasValueJsonColumn() !bool { + var rows = try conn.rows("PRAGMA table_info(records)", .{}); + defer rows.deinit(); + while (rows.next()) |row| { + if (std.mem.eql(u8, row.text(1), "value_json")) return true; + } + if (rows.err) |err| return err; + return false; +} + +fn validateRecordBlocksPresent() !void { + const row = try conn.row( + \\SELECT COUNT(*) + \\FROM records r + \\LEFT JOIN repo_blocks rb ON rb.did = r.did AND rb.cid = r.cid + \\WHERE rb.cid IS NULL + , .{}); + if (row == null) return; + defer row.?.deinit(); + if (row.?.int(0) != 0) return Error.MissingRecordBlock; } fn migrateAppPreferencesTable() !void { @@ -2316,18 +2397,24 @@ fn cidText(allocator: std.mem.Allocator, raw: []const u8) ![]const u8 { } fn recordFromRow(row: zqlite.Row, allocator: std.mem.Allocator) !Record { + const record_bytes = row.nullableBlob(4) orelse return Error.MissingRecordBlock; return .{ .did = try allocator.dupe(u8, row.text(0)), .collection = try allocator.dupe(u8, row.text(1)), .rkey = try allocator.dupe(u8, row.text(2)), .cid = try allocator.dupe(u8, row.text(3)), - .value_json = try allocator.dupe(u8, row.text(4)), + .value_json = try recordJsonFromBlock(allocator, record_bytes), .validation_status = validationStatusForRecord(row.text(1)), .rev = try allocator.dupe(u8, row.text(5)), .seq = @intCast(row.int(6)), }; } +fn recordJsonFromBlock(allocator: std.mem.Allocator, data: []const u8) ![]const u8 { + const value = zat.cbor.decodeAll(allocator, data) catch return Error.InvalidDagCbor; + return cbor_json.writeAlloc(allocator, value); +} + fn oauthRequestFromRow(row: zqlite.Row, allocator: std.mem.Allocator) !OAuthRequest { return .{ .request_id = try allocator.dupe(u8, row.text(0)), @@ -2712,7 +2799,6 @@ const schema_statements = [_][*:0]const u8{ \\ rkey TEXT NOT NULL, \\ uri TEXT NOT NULL UNIQUE, \\ cid TEXT NOT NULL, - \\ value_json BLOB NOT NULL, \\ rev TEXT NOT NULL, \\ seq INTEGER NOT NULL, \\ PRIMARY KEY (did, collection, rkey) @@ -2832,11 +2918,23 @@ test "persists records in sqlite" { const record = try create(allocator, account, "app.bsky.feed.post", "3ztest", parsed.value); try std.testing.expectEqualStrings("3ztest", record.rkey); + try std.testing.expect(!try recordsTableHasValueJsonColumn()); + + const block_row = try conn.row( + "SELECT count(*) FROM repo_blocks WHERE did = ? AND cid = ?", + .{ account.did, record.cid }, + ); + try std.testing.expect(block_row != null); + defer block_row.?.deinit(); + try std.testing.expectEqual(@as(i64, 1), block_row.?.int(0)); const fetched = get(account.did, "app.bsky.feed.post", "3ztest").?; try std.testing.expectEqualStrings(record.cid, fetched.cid); try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "hello") != null); + const listed = try writeListJson(allocator, account.did, "app.bsky.feed.post", 10); + try std.testing.expect(std.mem.indexOf(u8, listed, "hello") != null); + const repo_car = try writeRepoCar(allocator, account.did); const loaded = try zat.loadCommitFromCAR(allocator, repo_car); try std.testing.expectEqualStrings(account.did, loaded.commit.did); @@ -2846,6 +2944,51 @@ test "persists records in sqlite" { const tree = try zat.mst.Mst.loadFromBlocks(allocator, loaded.repo_car, loaded.commit.data_cid); const found = tree.get("app.bsky.feed.post/3ztest") orelse return error.MissingRecord; try std.testing.expectEqualStrings(record.cid, try cidText(allocator, found.raw)); + + try conn.exec( + "DELETE FROM repo_blocks WHERE did = ? AND cid = ?", + .{ account.did, record.cid }, + ); + try std.testing.expect(get(account.did, "app.bsky.feed.post", "3ztest") == null); +} + +test "put updates record index while content comes from repo 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, + "alice.test", + "alice@test.com", + "password", + "did:plc:cmadossymmii3izkabdbp5en", + true, + ); + + var first = try std.json.parseFromSlice(std.json.Value, allocator, "{\"text\":\"before\"}", .{}); + defer first.deinit(); + _ = try create(allocator, account, "app.bsky.feed.post", "3zput", first.value); + + var second = try std.json.parseFromSlice(std.json.Value, allocator, "{\"text\":\"after\"}", .{}); + defer second.deinit(); + const updated = try put(allocator, account, "app.bsky.feed.post", "3zput", second.value); + + const fetched = get(account.did, "app.bsky.feed.post", "3zput").?; + try std.testing.expectEqualStrings(updated.cid, fetched.cid); + try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "after") != null); + try std.testing.expect(std.mem.indexOf(u8, fetched.value_json, "before") == null); + + const block_row = try conn.row( + "SELECT count(*) FROM repo_blocks WHERE did = ? AND cid = ?", + .{ account.did, updated.cid }, + ); + try std.testing.expect(block_row != null); + defer block_row.?.deinit(); + try std.testing.expectEqual(@as(i64, 1), block_row.?.int(0)); } test "stores blob metadata in sqlite and bytes in disk blobstore" {