diff --git a/CHANGELOG.md b/CHANGELOG.md index 63b6a1a..d78dd50 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,14 @@ Reconstructed from git history for everything up to `v0.1.1`; kept by hand from Two months of work since `v0.1.1` (2026-06-18). Headlines: +- **fix**: `applyWrites` indexes a batch in op order. It used to apply every delete before + every create, so a same-key create+delete in one batch left the record row present while + the MST — which mutates in op order — had dropped it. The reference PDS walks its writes + in one ordered loop (`actor-store/repo/transactor.ts` `indexWrites`) and tranquil resolves + each delete against the running tree; zds was the outlier and could reach a state neither + produces. Blob refs are now attached to their record by a staged span rather than matched + by URI, keeping the pass linear — `applyWrites` has a 5MB body cap but no op-count cap. + No measurable throughput change on the write benchmark. - **fix**: bump zat to `v0.4.3`, which makes `Mst.collectBlocks` emit the block for an empty MST node instead of mistaking it for an unloaded stub and skipping it. Two repos on pds.zat.dev carried empty subtree nodes minted by a pre-`v0.3.19` writer whose blocks diff --git a/src/storage/store.zig b/src/storage/store.zig index c1fd6c4..4afab4d 100644 --- a/src/storage/store.zig +++ b/src/storage/store.zig @@ -2117,17 +2117,25 @@ fn applyWritesMeasured(allocator: std.mem.Allocator, account: auth.Account, ops: var record_blocks: std.ArrayList(ImportedBlock) = .empty; var mst_blocks: std.ArrayList(zat.car.Block) = .empty; var blob_refs: std.ArrayList(BlobRef) = .empty; + // half-open [start, end) into blob_refs for each staged record, parallel to + // `records`. stageRecordWrite appends a record's refs contiguously, so the + // index pass can find them without rescanning the whole list per record. + var record_ref_spans: std.ArrayList(struct { start: usize, end: usize }) = .empty; const stage_start = monotonicNs(); for (resolved_ops) |op| switch (op) { .create => |create_op| { const rkey = create_op.rkey orelse return Error.InvalidRecordKey; + const refs_start = blob_refs.items.len; const record = try stageRecordWrite(allocator, &tree, account, create_op.collection, rkey, create_op.value, options.validate, rev, seq, &record_blocks, &blob_refs); try records.append(allocator, record); + try record_ref_spans.append(allocator, .{ .start = refs_start, .end = blob_refs.items.len }); }, .update => |update_op| { + const refs_start = blob_refs.items.len; const record = try stageRecordWrite(allocator, &tree, account, update_op.collection, update_op.rkey, update_op.value, options.validate, rev, seq, &record_blocks, &blob_refs); try records.append(allocator, record); + try record_ref_spans.append(allocator, .{ .start = refs_start, .end = blob_refs.items.len }); }, .delete => |delete_op| { const path = try repoPath(allocator, delete_op.collection, delete_op.rkey); @@ -2174,34 +2182,43 @@ fn applyWritesMeasured(allocator: std.mem.Allocator, account: auth.Account, ops: \\ repo_rev = COALESCE(repo_blocks.repo_rev, excluded.repo_rev) , .{ account.did, commit.cid, zqlite.blob(commit.data), rev }); + // Index the batch in op order, the way the tree was mutated. Applying every + // delete before every create would let a same-key create+delete leave the row + // behind while the MST says it is gone. The reference PDS walks its writes in + // one ordered loop (actor-store/repo/transactor.ts indexWrites) and tranquil + // resolves each delete against the running tree; this matches both. + var staged_idx: usize = 0; for (resolved_ops) |op| switch (op) { .delete => |delete_op| { const uri = try std.fmt.allocPrint(allocator, "at://{s}/{s}/{s}", .{ account.did, delete_op.collection, delete_op.rkey }); try conn.exec("DELETE FROM expected_blobs WHERE record_uri = ?", .{uri}); try conn.exec("DELETE FROM records WHERE did = ? AND collection = ? AND rkey = ?", .{ account.did, delete_op.collection, delete_op.rkey }); }, - else => {}, + else => { + const record = records.items[staged_idx]; + staged_idx += 1; + const uri = try record.uri(allocator); + try conn.exec( + \\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, + \\ rev = excluded.rev, + \\ seq = excluded.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}); + const span = record_ref_spans.items[staged_idx - 1]; + for (blob_refs.items[span.start..span.end]) |ref| { + try conn.exec( + \\INSERT INTO expected_blobs (blob_cid, record_uri) + \\VALUES (?, ?) + \\ON CONFLICT(blob_cid, record_uri) DO UPDATE SET blob_cid = excluded.blob_cid + , .{ ref.cid, uri }); + } + }, }; - for (records.items) |record| { - const uri = try record.uri(allocator); - try conn.exec( - \\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, - \\ rev = excluded.rev, - \\ seq = excluded.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| { - try conn.exec( - \\INSERT INTO expected_blobs (blob_cid, record_uri) - \\VALUES (?, ?) - \\ON CONFLICT(blob_cid, record_uri) DO UPDATE SET blob_cid = excluded.blob_cid - , .{ ref.cid, ref.uri }); - } + std.debug.assert(staged_idx == records.items.len); try conn.exec( \\INSERT INTO commits (seq, did, cid, rev, prev) \\VALUES (?, ?, ?, ?, ?) @@ -8604,3 +8621,111 @@ test "a create commit followed by a delete commit leaves no record behind" { defer db_mutex.unlock(store_io); try std.testing.expectEqual(@as(u64, 0), try scalarCountLocked("SELECT COUNT(*) FROM records WHERE did = ?", account.did)); } + +test "a same-key create and delete in one batch leaves no record, matching the reference" { + 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, "batch2.test", "batch2@test.com", "password", "did:plc:batchone", true); + const doc = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"io.atcr.manifest\"}", .{}); + defer doc.deinit(); + const rkey = "3jzfcijpj2z2b"; + + // the tree applies these in order, so the key is gone; the index must agree + _ = try applyWritesWithOptions(allocator, account, &.{ + .{ .create = .{ .collection = "io.atcr.manifest", .rkey = rkey, .value = doc.value } }, + .{ .delete = .{ .collection = "io.atcr.manifest", .rkey = rkey } }, + }, .{ .validate = .skip }); + + const root = blk: { + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try std.testing.expectEqual(@as(u64, 0), try scalarCountLocked("SELECT COUNT(*) FROM records WHERE did = ?", account.did)); + break :blk (try latestCommitRawLocked(allocator, account.did)).?; + }; + + // and the MST must agree with the index: no such key in the tree + var reader = RepoBlockReader{ .allocator = allocator, .did = account.did }; + var tree = try zat.mst.Mst.loadLazy(allocator, root.data_cid_raw, reader.reader()); + try std.testing.expect(tree.get("io.atcr.manifest/" ++ rkey) == null); +} + +test "a delete then create of the same key in one batch keeps the record" { + 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, "batch3.test", "batch3@test.com", "password", "did:plc:batchtwo", true); + const doc = try std.json.parseFromSlice(std.json.Value, allocator, "{\"$type\":\"io.atcr.manifest\"}", .{}); + defer doc.deinit(); + const rkey = "3jzfcijpj2z2c"; + + _ = try applyWritesWithOptions(allocator, account, &.{ + .{ .create = .{ .collection = "io.atcr.manifest", .rkey = rkey, .value = doc.value } }, + }, .{ .validate = .skip }); + // the reverse order: the create wins, and the row must survive the earlier delete + _ = try applyWritesWithOptions(allocator, account, &.{ + .{ .delete = .{ .collection = "io.atcr.manifest", .rkey = rkey } }, + .{ .create = .{ .collection = "io.atcr.manifest", .rkey = rkey, .value = doc.value } }, + }, .{ .validate = .skip }); + + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + try std.testing.expectEqual(@as(u64, 1), try scalarCountLocked("SELECT COUNT(*) FROM records WHERE did = ?", account.did)); +} + +test "a batch maps each record's blob refs to that record" { + var arena = std.heap.ArenaAllocator.init(std.testing.allocator); + defer arena.deinit(); + const allocator = arena.allocator(); + + var tmp = std.testing.tmpDir(.{}); + defer tmp.cleanup(); + var path_buf: [std.fs.max_path_bytes]u8 = undefined; + const path_len = try tmp.dir.realPath(std.testing.io, &path_buf); + blobstore.init(std.Options.debug_io, path_buf[0..path_len]); + + try init(std.Options.debug_io, ":memory:"); + defer close(); + const account = try createAccount(allocator, "refspan.test", "refspan@test.com", "password", "did:plc:refspan", true); + + const cid_a = try putBlob(allocator, std.Options.debug_io, account, "alpha!", "text/plain"); + const cid_b = try putBlob(allocator, std.Options.debug_io, account, "bravo!", "text/plain"); + + const body = struct { + fn make(a: std.mem.Allocator, cid: []const u8) ![]const u8 { + return std.fmt.allocPrint( + a, + "{{\"$type\":\"fm.example.blob\",\"media\":{{\"$type\":\"blob\",\"ref\":{{\"$link\":\"{s}\"}},\"mimeType\":\"text/plain\",\"size\":6}}}}", + .{cid}, + ); + } + }; + const doc_a = try std.json.parseFromSlice(std.json.Value, allocator, try body.make(allocator, cid_a), .{}); + defer doc_a.deinit(); + const doc_b = try std.json.parseFromSlice(std.json.Value, allocator, try body.make(allocator, cid_b), .{}); + defer doc_b.deinit(); + + // both records in ONE batch, so their ref spans must not be crossed + _ = try applyWritesWithOptions(allocator, account, &.{ + .{ .create = .{ .collection = "fm.example.blob", .rkey = "aaaaaaaaaaaaa", .value = doc_a.value } }, + .{ .create = .{ .collection = "fm.example.blob", .rkey = "bbbbbbbbbbbbb", .value = doc_b.value } }, + }, .{ .validate = .skip }); + + db_mutex.lockUncancelable(store_io); + defer db_mutex.unlock(store_io); + for ([_][2][]const u8{ .{ "aaaaaaaaaaaaa", cid_a }, .{ "bbbbbbbbbbbbb", cid_b } }) |want| { + const uri = try std.fmt.allocPrint(allocator, "at://{s}/fm.example.blob/{s}", .{ account.did, want[0] }); + const row = try conn.row("SELECT blob_cid FROM expected_blobs WHERE record_uri = ?", .{uri}); + try std.testing.expect(row != null); + defer row.?.deinit(); + try std.testing.expectEqualStrings(want[1], row.?.text(0)); + } +}