lib/sql/src/history/recover.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

  1 const std = @import("std");
  2 const sys = @import("sys");
  3 const sql = @import("../root.zig");
  4 const conflict_mod = @import("conflict.zig");
  5 const record_mod = @import("record.zig");
  6 const store_mod = @import("store.zig");
  7 const materialize_mod = @import("materialize.zig");
  8 const instrumentation = @import("instrumentation.zig");
  9 const batch_mod = @import("batch.zig");
 10 const verify_mod = @import("verify.zig");
 11 const row = sql.row;
 12 const tree = sql.tree;
 13 const version = sql.version;
 14 
 15 const Allocator = std.mem.Allocator;
 16 const Error = store_mod.Error;
 17 const History = store_mod.History;
 18 const testing_io = std.Options.debug_io;
 19 
 20 const serial_batch_bytes: usize = 1 << 20;
 21 const parallel_batch_bytes: usize = 8 << 20;
 22 const parallel_min_bytes: usize = 16 << 20;
 23 
 24 const ReplayOptions = struct {
 25     batch_bytes: usize,
 26     hash_lanes: usize,
 27 };
 28 
 29 pub fn replayFile(
 30     self: *History,
 31     length: usize,
 32     recorder: ?*instrumentation.Recorder,
 33     control: sql.wal.Control,
 34 ) Error!usize {
 35     try control.check();
 36     const hash_lanes = if (length >= parallel_min_bytes and sys.thread.threadsSupported())
 37         sys.thread.cpuCount()
 38     else
 39         1;
 40     return replayFileWithOptions(self, length, .{
 41         .batch_bytes = if (hash_lanes == 1) serial_batch_bytes else parallel_batch_bytes,
 42         .hash_lanes = hash_lanes,
 43     }, recorder, control);
 44 }
 45 
 46 fn replayFileWithOptions(
 47     self: *History,
 48     length: usize,
 49     options: ReplayOptions,
 50     recorder: ?*instrumentation.Recorder,
 51     control: sql.wal.Control,
 52 ) Error!usize {
 53     const file = self.file orelse return error.InvalidHistory;
 54     var batch = batch_mod.Batch.init(self.allocator, self.io, file, length, options.batch_bytes);
 55     defer batch.deinit();
 56     var workers: verify_mod.Workers = .{};
 57     workers.start(options.hash_lanes);
 58     defer workers.deinit();
 59     var parse_arena = std.heap.ArenaAllocator.init(self.allocator);
 60     defer parse_arena.deinit();
 61 
 62     var offset: usize = 0;
 63     var batch_index: usize = 0;
 64     while (offset < length) {
 65         try control.check();
 66         const load_start = phaseStart(recorder);
 67         const terminator = try batch.load(offset);
 68         try control.check();
 69         const replay_bytes = batchReplayBytes(&batch);
 70         recordPhase(
 71             recorder,
 72             batch_index,
 73             .batch_load,
 74             load_start,
 75             replay_bytes,
 76             batch.records.items.len,
 77         );
 78         const verify_start = phaseStart(recorder);
 79         workers.verify(batch.bytes.items[0..batch.filled], batch.records.items);
 80         try control.check();
 81         recordPhase(
 82             recorder,
 83             batch_index,
 84             .hash_verification,
 85             verify_start,
 86             replay_bytes,
 87             batch.records.items.len,
 88         );
 89         const rebuild_start = phaseStart(recorder);
 90         for (batch.records.items) |record| {
 91             try control.check();
 92             const replayed: Error!bool = if (record.valid_hash)
 93                 replayVerifiedRecord(self, &batch, record, parse_arena.allocator())
 94             else
 95                 false;
 96             _ = parse_arena.reset(.retain_capacity);
 97             if (!try replayed) {
 98                 if (record.record_offset == 0) return error.InvalidHistory;
 99                 return record.record_offset;
100             }
101             offset = record.payload_offset + record.payload_len;
102         }
103         try control.check();
104         recordPhase(
105             recorder,
106             batch_index,
107             .index_rebuild,
108             rebuild_start,
109             replay_bytes,
110             batch.records.items.len,
111         );
112         batch_index += 1;
113         switch (terminator) {
114             .more => {},
115             .end, .tail => return offset,
116             .invalid => {
117                 if (offset == 0) return error.InvalidHistory;
118                 return offset;
119             },
120         }
121     }
122     return offset;
123 }
124 
125 fn batchReplayBytes(batch: *const batch_mod.Batch) usize {
126     if (batch.records.items.len == 0) return batch.filled;
127     const last = batch.records.items[batch.records.items.len - 1];
128     return last.payload_offset + last.payload_len - batch.start;
129 }
130 
131 fn phaseStart(recorder: ?*instrumentation.Recorder) i128 {
132     return if (recorder == null) 0 else sys.time.nanoTimestamp();
133 }
134 
135 fn recordPhase(
136     recorder: ?*instrumentation.Recorder,
137     batch: usize,
138     phase: instrumentation.Phase,
139     start_ns: i128,
140     bytes: usize,
141     records: usize,
142 ) void {
143     const active = recorder orelse return;
144     active.record(
145         batch,
146         phase,
147         instrumentation.elapsed(start_ns, sys.time.nanoTimestamp()),
148         bytes,
149         records,
150     );
151 }
152 
153 fn replayVerifiedRecord(self: *History, batch: *const batch_mod.Batch, record: batch_mod.Record, parse_allocator: Allocator) Error!bool {
154     const payload = batch.payload(record);
155     if (self.fast_forward != null and !isFastForwardRecord(record.kind)) {
156         return error.RecoveryRequired;
157     }
158     switch (record.kind) {
159         .row_chunk => {
160             const digest: version.Hash = payload[0..version.hash_bytes].*;
161             try indexRowChunk(self, digest, record.expected, record.payload_offset, record.payload_len);
162         },
163         .tree_nodes => indexTreeNodesPayload(self, record.expected, record.payload_offset, record.payload_len, payload) catch |err| switch (err) {
164             error.InvalidHistory => return false,
165             else => return err,
166         },
167         .relation_root => indexRelationRootPayload(self, parse_allocator, record.expected, record.payload_offset, record.payload_len, payload) catch |err| switch (err) {
168             error.InvalidHistory => return false,
169             else => return err,
170         },
171         .relation_spans => indexRelationSpansPayload(self, record.expected, record.payload_offset, record.payload_len, payload) catch |err| switch (err) {
172             error.InvalidHistory => return false,
173             else => return err,
174         },
175         .relation_rows => indexRelationRowsPayload(self, record.expected, record.payload_offset, record.payload_len, payload) catch |err| switch (err) {
176             error.InvalidHistory => return false,
177             else => return err,
178         },
179         .chunk_index_page => indexChunkIndexPagePayload(self, record.expected, record.payload_offset, record.payload_len, payload) catch |err| switch (err) {
180             error.InvalidHistory => return false,
181             else => return err,
182         },
183         .fast_forward_prepare => applyFastForwardPrepare(self, record.expected, payload) catch |err| switch (err) {
184             error.OutOfMemory => return err,
185             else => return error.RecoveryRequired,
186         },
187         .fast_forward_commit => applyFastForwardDecision(self, .target, payload) catch |err| switch (err) {
188             error.OutOfMemory => return err,
189             else => return error.RecoveryRequired,
190         },
191         .fast_forward_abort => applyFastForwardDecision(self, .baseline, payload) catch |err| switch (err) {
192             error.OutOfMemory => return err,
193             else => return error.RecoveryRequired,
194         },
195         .fast_forward_complete => applyFastForwardComplete(self, payload) catch |err| switch (err) {
196             error.OutOfMemory => return err,
197             else => return error.RecoveryRequired,
198         },
199         else => return try applyRecovered(self, record.kind, payload),
200     }
201     return true;
202 }
203 
204 pub fn applyRecovered(self: *History, kind: record_mod.RecordKind, payload: []const u8) Error!bool {
205     switch (kind) {
206         .database_root => applyDatabaseRootPayload(self, payload) catch |err| switch (err) {
207             error.InvalidHistory => return false,
208             else => return err,
209         },
210         .relation_root => unreachable,
211         .relation_rows => unreachable,
212         .commit => applyCommitPayload(self, payload) catch |err| switch (err) {
213             error.InvalidHistory => return false,
214             else => return err,
215         },
216         .ref => applyRefPayload(self, payload) catch |err| switch (err) {
217             error.InvalidHistory => return false,
218             else => return err,
219         },
220         .ref_delete => applyRefDeletePayload(self, payload) catch |err| switch (err) {
221             error.InvalidHistory => return false,
222             else => return err,
223         },
224         .conflict => applyConflictPayload(self, payload) catch |err| switch (err) {
225             error.InvalidHistory => return false,
226             else => return err,
227         },
228         .conflict_root => applyConflictRootPayload(self, payload) catch |err| switch (err) {
229             error.InvalidHistory => return false,
230             else => return err,
231         },
232         .relation_spans => unreachable,
233         .row_chunk => unreachable,
234         .chunk_index_page => unreachable,
235         .tree_nodes => unreachable,
236         .fast_forward_prepare,
237         .fast_forward_commit,
238         .fast_forward_abort,
239         .fast_forward_complete,
240         => unreachable,
241     }
242     return true;
243 }
244 
245 fn isFastForwardRecord(kind: record_mod.RecordKind) bool {
246     return switch (kind) {
247         .fast_forward_prepare,
248         .fast_forward_commit,
249         .fast_forward_abort,
250         .fast_forward_complete,
251         => true,
252         else => false,
253     };
254 }
255 
256 fn applyFastForwardPrepare(self: *History, id: version.Hash, payload: []const u8) Error!void {
257     if (self.fast_forward != null) return error.InvalidHistory;
258     var reader = record_mod.PayloadReader.init(payload);
259     const name_bytes = try reader.readBytes();
260     const expected = try reader.hash();
261     const target = try reader.hash();
262     try reader.finish();
263     if (self.findCommit(expected) == null or self.findCommit(target) == null) {
264         return error.InvalidHistory;
265     }
266     const ref_record = self.findRef(name_bytes) orelse return error.InvalidHistory;
267     if (!version.same(ref_record.ref.target, expected)) return error.InvalidHistory;
268     self.fast_forward = .{
269         .id = id,
270         .name = ref_record.name,
271         .expected = expected,
272         .target = target,
273     };
274 }
275 
276 fn applyFastForwardDecision(
277     self: *History,
278     decision: store_mod.FastForwardDecision,
279     payload: []const u8,
280 ) Error!void {
281     var reader = record_mod.PayloadReader.init(payload);
282     const id = try reader.hash();
283     try reader.finish();
284     const recovery = if (self.fast_forward) |*active| active else return error.InvalidHistory;
285     if (!version.same(recovery.id, id) or recovery.decision != .pending) {
286         return error.InvalidHistory;
287     }
288     const ref_record = self.findRef(recovery.name) orelse return error.InvalidHistory;
289     if (!version.same(ref_record.ref.target, recovery.expected)) return error.InvalidHistory;
290     if (decision == .target) ref_record.ref.target = recovery.target;
291     recovery.decision = decision;
292 }
293 
294 fn applyFastForwardComplete(self: *History, payload: []const u8) Error!void {
295     var reader = record_mod.PayloadReader.init(payload);
296     const id = try reader.hash();
297     try reader.finish();
298     const recovery = if (self.fast_forward) |*active| active else return error.InvalidHistory;
299     if (!version.same(recovery.id, id) or recovery.decision == .pending) {
300         return error.InvalidHistory;
301     }
302     const ref_record = self.findRef(recovery.name) orelse return error.InvalidHistory;
303     const selected = switch (recovery.decision) {
304         .pending => unreachable,
305         .baseline => recovery.expected,
306         .target => recovery.target,
307     };
308     if (!version.same(ref_record.ref.target, selected)) return error.InvalidHistory;
309     self.fast_forward = null;
310 }
311 
312 pub fn applyDatabaseRootPayload(self: *History, payload: []const u8) Error!void {
313     var reader = record_mod.PayloadReader.init(payload);
314     const conflicts = try reader.hash();
315     const entry_count = try reader.readU32();
316     const entries = try self.allocator.alloc(version.RelationEntry, entry_count);
317     defer self.allocator.free(entries);
318     for (entries) |*entry| {
319         const name = try reader.readBytes();
320         const hash = try reader.hash();
321         entry.* = .{
322             .name = name,
323             .hash = hash,
324         };
325     }
326     try reader.finish();
327     var root = try version.DatabaseRoot.initSorted(self.allocator, entries, .{ .hash = conflicts });
328     errdefer root.deinit();
329     if (self.findDatabaseRoot(root.hash) != null) {
330         root.deinit();
331         return;
332     }
333     try self.database_roots.ensureUnusedCapacity(self.allocator, 1);
334     try self.lookup.database_roots.ensureUnusedCapacity(self.allocator, 1);
335     self.lookup.database_roots.putAssumeCapacity(root.hash, self.database_roots.items.len);
336     self.database_roots.appendAssumeCapacity(.{ .root = root });
337 }
338 
339 pub fn indexRelationRootPayload(self: *History, allocator: Allocator, expected: version.Hash, offset: usize, payload_len: usize, payload: []const u8) Error!void {
340     var reader = record_mod.PayloadReader.init(payload);
341     const decoded = try materialize_mod.readRelationRootShallow(allocator, &reader);
342     var root = decoded.root;
343     defer root.deinit();
344     defer allocator.free(decoded.index_keys);
345     try reader.finish();
346     if (decoded.table_key) |key| {
347         if (self.findTreeNode(key) == null) return error.InvalidHistory;
348     }
349     for (decoded.index_keys) |index_key| {
350         const key = index_key orelse continue;
351         if (self.findTreeNode(key) == null) return error.InvalidHistory;
352     }
353     if (self.findRelationRoot(root.hash) != null) return;
354     try self.relation_roots.ensureUnusedCapacity(self.allocator, 1);
355     try self.lookup.relation_roots.ensureUnusedCapacity(self.allocator, 1);
356     indexRelationRootAssumeCapacity(self, root.hash, .{
357         .expected = expected,
358         .offset = offset,
359         .len = payload_len,
360     });
361 }
362 
363 pub fn indexRelationRootAssumeCapacity(self: *History, hash: version.Hash, location: record_mod.PayloadLocation) void {
364     self.lookup.relation_roots.putAssumeCapacity(hash, self.relation_roots.items.len);
365     self.relation_roots.appendAssumeCapacity(.{
366         .hash = hash,
367         .storage = .{ .indexed = location },
368     });
369 }
370 
371 pub fn indexRelationRowsPayload(self: *History, expected: version.Hash, offset: usize, payload_len: usize, payload: []const u8) Error!void {
372     var reader = record_mod.PayloadReader.init(payload);
373     const root = try reader.hash();
374     const page_count = try reader.readU32();
375     var page_index: u32 = 0;
376     while (page_index < page_count) : (page_index += 1) {
377         const digest = try reader.hash();
378         if (!self.hasIndexPage(digest)) return error.InvalidHistory;
379     }
380     try reader.finish();
381     if (self.hasRelationRows(root)) return;
382     try self.relation_rows.ensureUnusedCapacity(self.allocator, 1);
383     try self.lookup.relation_rows.ensureUnusedCapacity(self.allocator, 1);
384     self.lookup.relation_rows.putAssumeCapacity(root, self.relation_rows.items.len);
385     self.relation_rows.appendAssumeCapacity(.{
386         .root = root,
387         .storage = .{ .indexed = .{
388             .expected = expected,
389             .offset = offset,
390             .len = payload_len,
391         } },
392     });
393 }
394 
395 pub fn indexChunkIndexPagePayload(self: *History, expected: version.Hash, offset: usize, payload_len: usize, payload: []const u8) Error!void {
396     var reader = record_mod.PayloadReader.init(payload);
397     const digest = try reader.hash();
398     const chunk_count = try reader.readU32();
399     var chunk_index: u32 = 0;
400     while (chunk_index < chunk_count) : (chunk_index += 1) {
401         const chunk_digest = try reader.hash();
402         if (!self.hasRowChunk(chunk_digest)) return error.InvalidHistory;
403     }
404     try reader.finish();
405     if (self.hasIndexPage(digest)) return;
406     try self.index_pages.ensureUnusedCapacity(self.allocator, 1);
407     try self.lookup.index_pages.ensureUnusedCapacity(self.allocator, 1);
408     self.lookup.index_pages.putAssumeCapacity(digest, self.index_pages.items.len);
409     self.index_pages.appendAssumeCapacity(.{
410         .digest = digest,
411         .storage = .{ .indexed = .{
412             .expected = expected,
413             .offset = offset,
414             .len = payload_len,
415         } },
416     });
417 }
418 
419 pub fn indexTreeNodesPayload(self: *History, expected: version.Hash, payload_offset: usize, payload_len: usize, payload: []const u8) Error!void {
420     const initial_len = self.tree_nodes.items.len;
421     errdefer rollbackTreeNodeIndex(self, initial_len);
422     var reader = record_mod.PayloadReader.init(payload);
423     const count = try reader.readU32();
424     var applied: u32 = 0;
425     while (applied < count) : (applied += 1) {
426         const key = try reader.hash();
427         _ = try record_mod.nodeKind(try reader.readU8());
428         _ = try reader.readBytes();
429         _ = try reader.optionalBytes();
430         _ = try record_mod.readUsize(&reader);
431         _ = try record_mod.readSummary(&reader);
432         _ = try reader.hash();
433         const child_count = try reader.readU32();
434         var child_index: u32 = 0;
435         while (child_index < child_count) : (child_index += 1) {
436             const child = try reader.hash();
437             if (self.findTreeNode(child) == null) return error.InvalidHistory;
438         }
439         if (self.findTreeNode(key) != null) continue;
440         try self.tree_nodes.ensureUnusedCapacity(self.allocator, 1);
441         try self.lookup.tree_nodes.ensureUnusedCapacity(self.allocator, 1);
442         self.lookup.tree_nodes.putAssumeCapacity(key, self.tree_nodes.items.len);
443         self.tree_nodes.appendAssumeCapacity(.{
444             .key = key,
445             .payload = .{
446                 .expected = expected,
447                 .offset = payload_offset,
448                 .len = payload_len,
449             },
450         });
451     }
452     try reader.finish();
453 }
454 
455 pub fn rollbackTreeNodeIndex(self: *History, target_len: usize) void {
456     while (self.tree_nodes.items.len > target_len) {
457         const record = self.tree_nodes.pop().?;
458         _ = self.lookup.tree_nodes.orderedRemove(record.key);
459     }
460 }
461 
462 pub fn indexRowChunk(self: *History, digest: version.Hash, expected: version.Hash, offset: usize, payload_len: usize) Error!void {
463     if (self.findRowChunk(digest) != null) return;
464     try self.row_chunks.ensureUnusedCapacity(self.allocator, 1);
465     try self.lookup.row_chunks.ensureUnusedCapacity(self.allocator, 1);
466     self.lookup.row_chunks.putAssumeCapacity(digest, self.row_chunks.items.len);
467     self.row_chunks.appendAssumeCapacity(.{
468         .digest = digest,
469         .payload = .{
470             .expected = expected,
471             .offset = offset,
472             .len = payload_len,
473         },
474     });
475 }
476 
477 fn indexRelationSpansPayload(self: *History, expected: version.Hash, offset: usize, payload_len: usize, payload: []const u8) Error!void {
478     var reader = record_mod.PayloadReader.init(payload);
479     const root = try reader.hash();
480     const count = try reader.readU32();
481     const span_bytes = std.math.mul(usize, count, 2 * @sizeOf(u64)) catch return error.InvalidHistory;
482     if (span_bytes != payload.len - reader.cursor) return error.InvalidHistory;
483     if (self.lookup.relation_spans.get(root) != null) return;
484     try self.relation_spans.ensureUnusedCapacity(self.allocator, 1);
485     try self.lookup.relation_spans.ensureUnusedCapacity(self.allocator, 1);
486     self.lookup.relation_spans.putAssumeCapacity(root, self.relation_spans.items.len);
487     self.relation_spans.appendAssumeCapacity(.{
488         .root = root,
489         .storage = .{ .indexed = .{
490             .expected = expected,
491             .offset = offset,
492             .len = payload_len,
493         } },
494     });
495 }
496 
497 pub fn applyCommitPayload(self: *History, payload: []const u8) Error!void {
498     var reader = record_mod.PayloadReader.init(payload);
499     const root = try reader.hash();
500     const parent_count = try reader.readU32();
501     const parents = try self.allocator.alloc(version.Hash, parent_count);
502     var parent_index: usize = 0;
503     errdefer self.allocator.free(parents);
504     while (parent_index < parents.len) : (parent_index += 1) parents[parent_index] = try reader.hash();
505     try reader.finish();
506     const commit = version.Commit.init(root, parents);
507     if (self.findCommit(commit.hash) != null) {
508         self.allocator.free(parents);
509         return;
510     }
511     try self.commits.ensureUnusedCapacity(self.allocator, 1);
512     try self.lookup.commits.ensureUnusedCapacity(self.allocator, 1);
513     self.lookup.commits.putAssumeCapacity(commit.hash, self.commits.items.len);
514     self.commits.appendAssumeCapacity(.{
515         .parents = parents,
516         .commit = .{
517             .root = commit.root,
518             .parents = parents,
519             .hash = commit.hash,
520         },
521     });
522 }
523 
524 pub fn applyRefPayload(self: *History, payload: []const u8) Error!void {
525     var reader = record_mod.PayloadReader.init(payload);
526     const target = try reader.hash();
527     const name_bytes = try reader.readBytes();
528     try reader.finish();
529     if (self.findCommit(target) == null) return error.InvalidHistory;
530     if (self.findRef(name_bytes)) |record| {
531         record.ref.target = target;
532         return;
533     }
534     const name = try self.allocator.dupe(u8, name_bytes);
535     errdefer self.allocator.free(name);
536     try self.refs.append(self.allocator, .{
537         .name = name,
538         .ref = .{
539             .name = name,
540             .target = target,
541         },
542     });
543 }
544 
545 pub fn applyRefDeletePayload(self: *History, payload: []const u8) Error!void {
546     var reader = record_mod.PayloadReader.init(payload);
547     const name_bytes = try reader.readBytes();
548     try reader.finish();
549     const index = self.findRefIndex(name_bytes) orelse return;
550     var removed = self.refs.orderedRemove(index);
551     removed.deinit(self.allocator);
552 }
553 
554 pub fn applyConflictPayload(self: *History, payload: []const u8) Error!void {
555     var reader = record_mod.PayloadReader.init(payload);
556     const decoded = try conflict_mod.decodeConflictArtifactPayload(&reader);
557     try reader.finish();
558     if (self.findConflict(decoded.hash) != null) return;
559 
560     var artifact = try conflict_mod.cloneConflictArtifact(self.allocator, decoded);
561     errdefer conflict_mod.deinitConflictArtifact(self.allocator, &artifact);
562     try self.conflicts.ensureUnusedCapacity(self.allocator, 1);
563     try self.lookup.conflicts.ensureUnusedCapacity(self.allocator, 1);
564     self.lookup.conflicts.putAssumeCapacity(decoded.hash, self.conflicts.items.len);
565     self.conflicts.appendAssumeCapacity(.{ .artifact = artifact });
566 }
567 
568 pub fn applyConflictRootPayload(self: *History, payload: []const u8) Error!void {
569     var reader = record_mod.PayloadReader.init(payload);
570     const entries = try conflict_mod.readConflictEntries(self.allocator, &reader);
571     errdefer conflict_mod.deinitConflictEntries(self.allocator, entries);
572     try reader.finish();
573     self.validateConflictEntries(entries) catch {
574         return error.InvalidHistory;
575     };
576     const root = version.ConflictRoot.init(entries);
577     if (self.findConflictRoot(root.hash) != null) {
578         self.validateConflictRoot(root.hash) catch {
579             return error.InvalidHistory;
580         };
581         conflict_mod.deinitConflictEntries(self.allocator, entries);
582         return;
583     }
584     try self.conflict_roots.ensureUnusedCapacity(self.allocator, 1);
585     try self.lookup.conflict_roots.ensureUnusedCapacity(self.allocator, 1);
586     self.lookup.conflict_roots.putAssumeCapacity(root.hash, self.conflict_roots.items.len);
587     self.conflict_roots.appendAssumeCapacity(.{
588         .root = root,
589         .entries = entries,
590     });
591 }
592 
593 fn replayTestingHistory(dir: std.Io.Dir, path: []const u8, options: ReplayOptions) !History {
594     const length: usize = @intCast((try dir.statFile(testing_io, path, .{})).size);
595     var history = History{
596         .allocator = std.testing.allocator,
597         .io = testing_io,
598     };
599     errdefer history.deinit();
600     history.file = try dir.createFile(testing_io, path, .{ .read = true, .truncate = false });
601     history.bytes_written = try replayFileWithOptions(&history, length, options, null, .{});
602     return history;
603 }
604 
605 fn appendTestingCommitPayload(target: *std.ArrayList(u8), root: version.Hash) !void {
606     try record_mod.appendHash(std.testing.allocator, target, root);
607     try record_mod.appendU32(std.testing.allocator, target, 0);
608 }
609 
610 fn appendTestingRefPayload(target: *std.ArrayList(u8), target_hash: version.Hash, name: []const u8) !void {
611     try record_mod.appendHash(std.testing.allocator, target, target_hash);
612     try record_mod.appendBytes(std.testing.allocator, target, name);
613 }
614 
615 fn flipTestingByte(file: std.Io.File, offset: usize) !void {
616     var byte: [1]u8 = undefined;
617     try std.testing.expectEqual(@as(usize, 1), try file.readPositionalAll(testing_io, byte[0..], offset));
618     byte[0] ^= 0xff;
619     try file.writePositionalAll(testing_io, byte[0..], offset);
620 }
621 
622 test "history replay serial and parallel batches agree" {
623     if (comptime !sys.thread.threadsSupported()) return error.SkipZigTest;
624     var tmp = std.testing.tmpDir(.{});
625     defer tmp.cleanup();
626 
627     var file = try tmp.dir.createFile(testing_io, "parallel.history", .{ .read = true });
628     var payload: std.ArrayList(u8) = .empty;
629     defer payload.deinit(std.testing.allocator);
630     var roots: [8]version.Hash = undefined;
631     var offset: usize = 0;
632     for (&roots, 0..) |*root, index| {
633         var label_buffer: [64]u8 = undefined;
634         const label = try std.fmt.bufPrint(&label_buffer, "parallel-history-{d}", .{index});
635         root.* = version.emptyHash(label);
636         payload.clearRetainingCapacity();
637         try appendTestingCommitPayload(&payload, root.*);
638         offset = try store_mod.writeTestingRecord(file, offset, .commit, payload.items);
639     }
640     try file.setLength(testing_io, offset);
641     file.close(testing_io);
642 
643     for ([_]usize{ 1, 4 }) |hash_lanes| {
644         var history = try replayTestingHistory(tmp.dir, "parallel.history", .{
645             .batch_bytes = 170,
646             .hash_lanes = hash_lanes,
647         });
648         defer history.deinit();
649         try std.testing.expectEqual(offset, history.len());
650         for (roots) |root| try std.testing.expect(history.hasCommit(version.Commit.init(root, &.{}).hash));
651     }
652 }
653 
654 test "history replay keeps the earliest hash failure across worker order" {
655     if (comptime !sys.thread.threadsSupported()) return error.SkipZigTest;
656     var tmp = std.testing.tmpDir(.{});
657     defer tmp.cleanup();
658 
659     var file = try tmp.dir.createFile(testing_io, "hash-order.history", .{ .read = true });
660     var payload: std.ArrayList(u8) = .empty;
661     defer payload.deinit(std.testing.allocator);
662     const prefix_root = version.emptyHash("parallel-prefix");
663     try appendTestingCommitPayload(&payload, prefix_root);
664     const prefix = try store_mod.writeTestingRecord(file, 0, .commit, payload.items);
665 
666     const large_len: usize = 2 << 20;
667     const large = try std.testing.allocator.alloc(u8, large_len);
668     defer std.testing.allocator.free(large);
669     @memset(large, 0x5a);
670     const chunk_hash = version.emptyHash("parallel-large-chunk");
671     @memcpy(large[0..version.hash_bytes], chunk_hash[0..]);
672     const large_end = try store_mod.writeTestingRecord(file, prefix, .row_chunk, large);
673     try flipTestingByte(file, large_end - 1);
674 
675     payload.clearRetainingCapacity();
676     try appendTestingCommitPayload(&payload, version.emptyHash("parallel-later-hash"));
677     const later_end = try store_mod.writeTestingRecord(file, large_end, .commit, payload.items);
678     try flipTestingByte(file, later_end - 1);
679     try file.writePositionalAll(testing_io, "torn", later_end);
680     const original_len = later_end + "torn".len;
681     try file.setLength(testing_io, original_len);
682     file.close(testing_io);
683 
684     for ([_]usize{ 1, 4 }) |hash_lanes| {
685         var history = try replayTestingHistory(tmp.dir, "hash-order.history", .{
686             .batch_bytes = 4 << 20,
687             .hash_lanes = hash_lanes,
688         });
689         defer history.deinit();
690         try std.testing.expectEqual(prefix, history.len());
691         try std.testing.expect(history.hasCommit(version.Commit.init(prefix_root, &.{}).hash));
692         try std.testing.expect(!history.hasRowChunk(chunk_hash));
693         try std.testing.expectEqual(original_len, @as(usize, @intCast((try tmp.dir.statFile(testing_io, "hash-order.history", .{})).size)));
694     }
695 }
696 
697 test "history replay keeps semantic failure before later hashes" {
698     if (comptime !sys.thread.threadsSupported()) return error.SkipZigTest;
699     var tmp = std.testing.tmpDir(.{});
700     defer tmp.cleanup();
701 
702     var file = try tmp.dir.createFile(testing_io, "semantic-order.history", .{ .read = true });
703     var payload: std.ArrayList(u8) = .empty;
704     defer payload.deinit(std.testing.allocator);
705     const prefix_root = version.emptyHash("semantic-prefix");
706     try appendTestingCommitPayload(&payload, prefix_root);
707     const prefix = try store_mod.writeTestingRecord(file, 0, .commit, payload.items);
708 
709     payload.clearRetainingCapacity();
710     try appendTestingRefPayload(&payload, version.emptyHash("future-commit"), "main");
711     const semantic_end = try store_mod.writeTestingRecord(file, prefix, .ref, payload.items);
712     payload.clearRetainingCapacity();
713     try appendTestingCommitPayload(&payload, version.emptyHash("semantic-later-hash"));
714     const original_len = try store_mod.writeTestingRecord(file, semantic_end, .commit, payload.items);
715     try flipTestingByte(file, original_len - 1);
716     try file.setLength(testing_io, original_len);
717     file.close(testing_io);
718 
719     for ([_]usize{ 1, 4 }) |hash_lanes| {
720         var history = try replayTestingHistory(tmp.dir, "semantic-order.history", .{
721             .batch_bytes = 1 << 20,
722             .hash_lanes = hash_lanes,
723         });
724         defer history.deinit();
725         try std.testing.expectEqual(prefix, history.len());
726         try std.testing.expect((try history.ref("main")) == null);
727         try std.testing.expect(history.hasCommit(version.Commit.init(prefix_root, &.{}).hash));
728     }
729 }