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 }