lib/trace/src/session.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

  1 const std = @import("std");
  2 const event = @import("event.zig");
  3 const checkpoint = @import("checkpoint.zig");
  4 
  5 const Allocator = std.mem.Allocator;
  6 
  7 pub const ReplayProgress = struct {
  8     cursor: u64,
  9     event_count: u64,
 10     remaining_count: u64,
 11 
 12     pub fn complete(self: ReplayProgress) bool {
 13         return self.remaining_count == 0;
 14     }
 15 };
 16 
 17 pub const Session = struct {
 18     allocator: Allocator,
 19     mode: event.Mode,
 20     sink: ?event.Sink = null,
 21     source: ?event.Source = null,
 22     epoch: u64 = 0,
 23     next_seq: u64 = 1,
 24     sequence_exhausted: bool = false,
 25     replay_cursor: u64 = 0,
 26     last_checkpoint: ?event.Event = null,
 27 
 28     pub fn initOff(allocator: Allocator) Session {
 29         return .{ .allocator = allocator, .mode = .off };
 30     }
 31 
 32     pub fn initRecord(allocator: Allocator, sink: event.Sink) Session {
 33         return .{ .allocator = allocator, .mode = .record, .sink = sink };
 34     }
 35 
 36     pub fn initReplay(allocator: Allocator, source: event.Source) Session {
 37         return .{ .allocator = allocator, .mode = .replay, .source = source };
 38     }
 39 
 40     pub fn deinit(self: *Session) void {
 41         if (self.last_checkpoint) |*item| item.deinit(self.allocator);
 42         self.* = undefined;
 43     }
 44 
 45     pub fn replayProgress(self: *const Session) ReplayProgress {
 46         const event_count = if (self.source) |source| source.count() else 0;
 47         std.debug.assert(self.replay_cursor <= event_count);
 48         return .{
 49             .cursor = self.replay_cursor,
 50             .event_count = event_count,
 51             .remaining_count = event_count - self.replay_cursor,
 52         };
 53     }
 54 
 55     pub fn sessionStart(
 56         self: *Session,
 57         thread_id: event.ThreadId,
 58         label: []const u8,
 59     ) !event.Timepoint {
 60         const timepoint = try self.prepareTimepoint(thread_id);
 61         return try self.handleEvent(event.Event.sessionStart(timepoint, label));
 62     }
 63 
 64     pub fn sessionEnd(
 65         self: *Session,
 66         thread_id: event.ThreadId,
 67         status: i64,
 68     ) !event.Timepoint {
 69         const timepoint = try self.prepareTimepoint(thread_id);
 70         return try self.handleEvent(event.Event.sessionEnd(timepoint, status));
 71     }
 72 
 73     pub fn functionEnter(
 74         self: *Session,
 75         thread_id: event.ThreadId,
 76         site: event.Safepoint,
 77     ) !event.Timepoint {
 78         const timepoint = try self.prepareTimepoint(thread_id);
 79         return try self.handleEvent(event.Event.functionEnter(timepoint, site));
 80     }
 81 
 82     pub fn functionExit(
 83         self: *Session,
 84         thread_id: event.ThreadId,
 85         site: event.Safepoint,
 86     ) !event.Timepoint {
 87         const timepoint = try self.prepareTimepoint(thread_id);
 88         return try self.handleEvent(event.Event.functionExit(timepoint, site));
 89     }
 90 
 91     pub fn safepoint(
 92         self: *Session,
 93         thread_id: event.ThreadId,
 94         site: event.Safepoint,
 95     ) !event.Timepoint {
 96         const timepoint = try self.prepareTimepoint(thread_id);
 97         return try self.handleEvent(event.Event.safepointReached(timepoint, site));
 98     }
 99 
100     pub fn allocation(
101         self: *Session,
102         thread_id: event.ThreadId,
103         object_id: u64,
104         size: u64,
105         alignment: u32,
106         label: ?[]const u8,
107     ) !event.Timepoint {
108         const timepoint = try self.prepareTimepoint(thread_id);
109         return try self.handleEvent(event.Event.allocation(
110             timepoint,
111             object_id,
112             size,
113             alignment,
114             label,
115         ));
116     }
117 
118     pub fn free(
119         self: *Session,
120         thread_id: event.ThreadId,
121         object_id: u64,
122     ) !event.Timepoint {
123         const timepoint = try self.prepareTimepoint(thread_id);
124         return try self.handleEvent(event.Event.free(timepoint, object_id));
125     }
126 
127     pub fn boundaryBytesAlloc(
128         self: *Session,
129         thread_id: event.ThreadId,
130         operation: []const u8,
131         recorded_value: ?[]const u8,
132     ) ![]u8 {
133         const timepoint = try self.prepareTimepoint(thread_id);
134         switch (self.mode) {
135             .off => {
136                 const source_value = recorded_value orelse return error.MissingBoundaryValue;
137                 const value = try self.allocator.dupe(u8, source_value);
138                 self.commitTimepoint(timepoint);
139                 return value;
140             },
141             .record => {
142                 const value = recorded_value orelse return error.MissingBoundaryValue;
143                 const owned = try self.allocator.dupe(u8, value);
144                 errdefer self.allocator.free(owned);
145                 try self.sink.?.append(event.Event.boundaryBytes(timepoint, operation, value));
146                 self.commitTimepoint(timepoint);
147                 return owned;
148             },
149             .replay => {
150                 const source = self.source.?;
151                 const actual = (try source.peek()) orelse return error.MissingReplayEvent;
152                 const recorded = actual.data orelse return error.ReplayEventMismatch;
153                 const expected = event.Event.boundaryBytes(timepoint, operation, recorded);
154                 if (!expected.eqlForReplay(actual.*)) return error.ReplayEventMismatch;
155                 const next_replay_cursor = try self.prepareReplayCursor();
156                 const value = try self.allocator.dupe(u8, recorded);
157                 source.advance();
158                 self.replay_cursor = next_replay_cursor;
159                 self.commitTimepoint(timepoint);
160                 return value;
161             },
162         }
163     }
164 
165     pub fn checkpointCapture(
166         self: *Session,
167         thread_id: event.ThreadId,
168         label: []const u8,
169         provider: checkpoint.Provider,
170     ) !event.Timepoint {
171         const timepoint = try self.prepareTimepoint(thread_id);
172         const bytes = try provider.captureAlloc(self.allocator);
173         defer self.allocator.free(bytes);
174         return try self.handleEvent(event.Event.checkpoint(timepoint, label, bytes));
175     }
176 
177     pub fn checkpointRestore(
178         self: *Session,
179         thread_id: event.ThreadId,
180         label: []const u8,
181         provider: checkpoint.Provider,
182     ) !event.Timepoint {
183         const checkpoint_event = self.last_checkpoint orelse return error.MissingCheckpoint;
184         const checkpoint_bytes = checkpoint_event.data orelse return error.MissingCheckpoint;
185         const timepoint = try self.prepareTimepoint(thread_id);
186         const expected = event.Event.checkpointRestore(timepoint, label);
187         const prepared_replay = if (self.mode == .replay)
188             try self.prepareReplay(expected)
189         else
190             null;
191         const replay_timepoint = if (prepared_replay) |prepared|
192             prepared.actual.timepoint
193         else
194             timepoint;
195 
196         provider.prepareRestoreBytes(checkpoint_bytes) catch |err| {
197             provider.cancelRestore();
198             return err;
199         };
200         errdefer provider.cancelRestore();
201 
202         if (self.mode == .record) try self.sink.?.append(expected);
203 
204         provider.commitRestore();
205         if (prepared_replay) |prepared| self.commitReplay(prepared);
206         self.commitTimepoint(timepoint);
207         return replay_timepoint;
208     }
209 
210     pub fn userEvent(
211         self: *Session,
212         thread_id: event.ThreadId,
213         label: []const u8,
214         bytes: []const u8,
215     ) !event.Timepoint {
216         const timepoint = try self.prepareTimepoint(thread_id);
217         return try self.handleEvent(event.Event.user(timepoint, label, bytes));
218     }
219 
220     pub fn verifyReplayComplete(self: *Session) !void {
221         if (self.mode != .replay) return;
222         const source = self.source.?;
223         if (self.replay_cursor != source.count()) return error.ReplayNotComplete;
224         if (try source.peek() != null) return error.ReplayNotComplete;
225     }
226 
227     fn prepareTimepoint(self: *const Session, thread_id: event.ThreadId) !event.Timepoint {
228         if (self.sequence_exhausted) return error.SequenceExhausted;
229         return event.Timepoint{
230             .epoch = self.epoch,
231             .thread_id = thread_id,
232             .seq = self.next_seq,
233         };
234     }
235 
236     fn commitTimepoint(self: *Session, timepoint: event.Timepoint) void {
237         std.debug.assert(!self.sequence_exhausted);
238         std.debug.assert(timepoint.epoch == self.epoch);
239         std.debug.assert(timepoint.seq == self.next_seq);
240         if (self.next_seq == std.math.maxInt(u64)) {
241             self.sequence_exhausted = true;
242         } else {
243             self.next_seq += 1;
244         }
245     }
246 
247     fn handleEvent(self: *Session, expected: event.Event) !event.Timepoint {
248         var owned_checkpoint: ?event.Event = null;
249         errdefer if (owned_checkpoint) |*owned| owned.deinit(self.allocator);
250 
251         var replay_timepoint: ?event.Timepoint = null;
252         switch (self.mode) {
253             .off => {
254                 if (expected.kind == .checkpoint) {
255                     owned_checkpoint = try expected.cloneAlloc(self.allocator);
256                 }
257             },
258             .record => {
259                 if (expected.kind == .checkpoint) {
260                     owned_checkpoint = try expected.cloneAlloc(self.allocator);
261                 }
262                 try self.sink.?.append(expected);
263             },
264             .replay => {
265                 const prepared = try self.prepareReplay(expected);
266                 replay_timepoint = prepared.actual.timepoint;
267                 if (prepared.actual.kind == .checkpoint) {
268                     owned_checkpoint = try prepared.actual.cloneAlloc(self.allocator);
269                 }
270                 self.commitReplay(prepared);
271             },
272         }
273 
274         self.commitTimepoint(expected.timepoint);
275         if (owned_checkpoint) |owned| {
276             self.replaceCheckpoint(owned);
277             owned_checkpoint = null;
278         }
279         return replay_timepoint orelse expected.timepoint;
280     }
281 
282     const PreparedReplay = struct {
283         source: event.Source,
284         actual: *const event.Event,
285         next_cursor: u64,
286     };
287 
288     fn prepareReplay(self: *Session, expected: event.Event) !PreparedReplay {
289         const source = self.source.?;
290         const actual = (try source.peek()) orelse return error.MissingReplayEvent;
291         if (!expected.eqlForReplay(actual.*)) return error.ReplayEventMismatch;
292         return .{
293             .source = source,
294             .actual = actual,
295             .next_cursor = try self.prepareReplayCursor(),
296         };
297     }
298 
299     fn prepareReplayCursor(self: *const Session) !u64 {
300         return std.math.add(u64, self.replay_cursor, 1) catch
301             return error.ReplayCursorExhausted;
302     }
303 
304     fn commitReplay(self: *Session, prepared: PreparedReplay) void {
305         prepared.source.advance();
306         self.replay_cursor = prepared.next_cursor;
307     }
308 
309     fn replaceCheckpoint(self: *Session, owned: event.Event) void {
310         if (self.last_checkpoint) |*previous| previous.deinit(self.allocator);
311         self.last_checkpoint = owned;
312     }
313 };
314 
315 const SliceStream = struct {
316     events: []const event.Event,
317     cursor: usize = 0,
318 
319     fn source(self: *SliceStream) event.Source {
320         return .{ .context = self, .peekFn = peek, .advanceFn = advance, .countFn = count };
321     }
322 
323     fn peek(context: *anyopaque) !?*const event.Event {
324         const self: *SliceStream = @ptrCast(@alignCast(context));
325         if (self.cursor == self.events.len) return null;
326         return &self.events[self.cursor];
327     }
328 
329     fn advance(context: *anyopaque) void {
330         const self: *SliceStream = @ptrCast(@alignCast(context));
331         self.cursor += 1;
332     }
333 
334     fn count(context: *anyopaque) u64 {
335         const self: *SliceStream = @ptrCast(@alignCast(context));
336         return self.events.len;
337     }
338 };
339 
340 const CaptureSink = struct {
341     timepoints: [8]event.Timepoint = undefined,
342     len: usize = 0,
343 
344     fn sink(self: *CaptureSink) event.Sink {
345         return .{ .context = self, .appendFn = append };
346     }
347 
348     fn append(context: *anyopaque, item: event.Event) !void {
349         const self: *CaptureSink = @ptrCast(@alignCast(context));
350         if (self.len == self.timepoints.len) return error.TestSinkFull;
351         self.timepoints[self.len] = item.timepoint;
352         self.len += 1;
353     }
354 };
355 
356 const FailOnceSink = struct {
357     timepoints: [8]event.Timepoint = undefined,
358     len: usize = 0,
359     fail_next: bool = true,
360 
361     fn sink(self: *FailOnceSink) event.Sink {
362         return .{ .context = self, .appendFn = append };
363     }
364 
365     fn append(context: *anyopaque, item: event.Event) !void {
366         const self: *FailOnceSink = @ptrCast(@alignCast(context));
367         if (self.fail_next) {
368             self.fail_next = false;
369             return error.InjectedSinkFailure;
370         }
371         if (self.len == self.timepoints.len) return error.TestSinkFull;
372         self.timepoints[self.len] = item.timepoint;
373         self.len += 1;
374     }
375 };
376 
377 const FailOnceStream = struct {
378     events: []const event.Event,
379     cursor: usize = 0,
380     fail_next: bool = true,
381 
382     fn source(self: *FailOnceStream) event.Source {
383         return .{ .context = self, .peekFn = peek, .advanceFn = advance, .countFn = count };
384     }
385 
386     fn peek(context: *anyopaque) !?*const event.Event {
387         const self: *FailOnceStream = @ptrCast(@alignCast(context));
388         if (self.fail_next) {
389             self.fail_next = false;
390             return error.InjectedSourceFailure;
391         }
392         if (self.cursor == self.events.len) return null;
393         return &self.events[self.cursor];
394     }
395 
396     fn advance(context: *anyopaque) void {
397         const self: *FailOnceStream = @ptrCast(@alignCast(context));
398         self.cursor += 1;
399     }
400 
401     fn count(context: *anyopaque) u64 {
402         const self: *FailOnceStream = @ptrCast(@alignCast(context));
403         return self.events.len;
404     }
405 };
406 
407 const ProviderProbe = struct {
408     value: [32]u8 = undefined,
409     value_len: usize = 0,
410     staged: [32]u8 = undefined,
411     staged_len: usize = 0,
412     prepared: bool = false,
413     prepare_calls: u64 = 0,
414     commit_calls: u64 = 0,
415     cancel_calls: u64 = 0,
416     fail_prepare: bool = false,
417 
418     fn init(bytes: []const u8) ProviderProbe {
419         var self: ProviderProbe = .{};
420         self.setValue(bytes);
421         return self;
422     }
423 
424     fn setValue(self: *ProviderProbe, bytes: []const u8) void {
425         std.debug.assert(bytes.len <= self.value.len);
426         @memcpy(self.value[0..bytes.len], bytes);
427         self.value_len = bytes.len;
428     }
429 
430     fn current(self: *const ProviderProbe) []const u8 {
431         return self.value[0..self.value_len];
432     }
433 
434     fn provider(self: *ProviderProbe) checkpoint.Provider {
435         return .{
436             .context = self,
437             .capture = capture,
438             .prepare_restore = prepareRestore,
439             .commit_restore = commitRestore,
440             .cancel_restore = cancelRestore,
441         };
442     }
443 
444     fn capture(context: *anyopaque, allocator: Allocator) ![]u8 {
445         const self: *ProviderProbe = @ptrCast(@alignCast(context));
446         return try allocator.dupe(u8, self.current());
447     }
448 
449     fn prepareRestore(context: *anyopaque, bytes: []const u8) !void {
450         const self: *ProviderProbe = @ptrCast(@alignCast(context));
451         self.prepare_calls += 1;
452         if (self.fail_prepare) return error.InjectedProviderFailure;
453         if (bytes.len > self.staged.len) return error.TestProviderCapacityExceeded;
454         @memcpy(self.staged[0..bytes.len], bytes);
455         self.staged_len = bytes.len;
456         self.prepared = true;
457     }
458 
459     fn commitRestore(context: *anyopaque) void {
460         const self: *ProviderProbe = @ptrCast(@alignCast(context));
461         std.debug.assert(self.prepared);
462         @memcpy(self.value[0..self.staged_len], self.staged[0..self.staged_len]);
463         self.value_len = self.staged_len;
464         self.prepared = false;
465         self.commit_calls += 1;
466     }
467 
468     fn cancelRestore(context: *anyopaque) void {
469         const self: *ProviderProbe = @ptrCast(@alignCast(context));
470         self.prepared = false;
471         self.cancel_calls += 1;
472     }
473 };
474 
475 test "recording streams events without retaining copies" {
476     var capture: CaptureSink = .{};
477     var session = Session.initRecord(std.testing.allocator, capture.sink());
478     defer session.deinit();
479     const start = try session.sessionStart(1, "test");
480     const enter = try session.functionEnter(1, .{ .function_id = 10, .site_id = 1 });
481     const exit = try session.functionExit(1, .{ .function_id = 10, .site_id = 2 });
482     try std.testing.expectEqual(@as(u64, 1), start.seq);
483     try std.testing.expectEqual(@as(u64, 2), enter.seq);
484     try std.testing.expectEqual(@as(u64, 3), exit.seq);
485     try std.testing.expectEqual(@as(usize, 3), capture.len);
486 }
487 
488 test "session failure atomicity retries a failed record append without a sequence gap" {
489     var capture: FailOnceSink = .{};
490     var session = Session.initRecord(std.testing.allocator, capture.sink());
491     defer session.deinit();
492 
493     try std.testing.expectError(
494         error.InjectedSinkFailure,
495         session.sessionStart(1, "test"),
496     );
497     try std.testing.expectEqual(@as(usize, 0), capture.len);
498 
499     const retried = try session.sessionStart(1, "test");
500     try std.testing.expectEqual(@as(u64, 1), retried.seq);
501     try std.testing.expectEqual(@as(usize, 1), capture.len);
502     try std.testing.expectEqual(@as(u64, 1), capture.timepoints[0].seq);
503 }
504 
505 test "session failure atomicity retries a failed source peek" {
506     const events = [_]event.Event{
507         event.Event.sessionStart(.{ .thread_id = 1, .seq = 1 }, "test"),
508     };
509     var stream = FailOnceStream{ .events = &events };
510     var session = Session.initReplay(std.testing.allocator, stream.source());
511     defer session.deinit();
512 
513     try std.testing.expectError(
514         error.InjectedSourceFailure,
515         session.sessionStart(1, "test"),
516     );
517     try std.testing.expectEqual(@as(usize, 0), stream.cursor);
518     try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
519 
520     const retried = try session.sessionStart(1, "test");
521     try std.testing.expectEqual(@as(u64, 1), retried.seq);
522     try session.verifyReplayComplete();
523 }
524 
525 test "replay validates a borrowed stream and exposes numeric progress" {
526     const events = [_]event.Event{
527         event.Event.sessionStart(.{ .thread_id = 1, .seq = 1 }, "test"),
528         event.Event.safepointReached(.{ .thread_id = 1, .seq = 2 }, .{ .function_id = 20, .site_id = 5 }),
529     };
530     var stream = SliceStream{ .events = &events };
531     var session = Session.initReplay(std.testing.allocator, stream.source());
532     defer session.deinit();
533     var progress = session.replayProgress();
534     try std.testing.expectEqual(@as(u64, 0), progress.cursor);
535     try std.testing.expectEqual(@as(u64, 2), progress.remaining_count);
536     _ = try session.sessionStart(1, "test");
537     progress = session.replayProgress();
538     try std.testing.expectEqual(@as(u64, 1), progress.cursor);
539     try std.testing.expectEqual(@as(u64, 1), progress.remaining_count);
540     _ = try session.safepoint(1, .{ .function_id = 20, .site_id = 5 });
541     try session.verifyReplayComplete();
542     try std.testing.expect(session.replayProgress().complete());
543 }
544 
545 test "session failure atomicity allocates a record boundary before append" {
546     var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{
547         .fail_index = 0,
548     });
549     var capture: CaptureSink = .{};
550     var session = Session.initRecord(failing.allocator(), capture.sink());
551     defer session.deinit();
552 
553     try std.testing.expectError(
554         error.OutOfMemory,
555         session.boundaryBytesAlloc(1, "env.TEST", "value"),
556     );
557     try std.testing.expectEqual(@as(usize, 0), capture.len);
558 
559     failing.fail_index = std.math.maxInt(usize);
560     failing.resize_fail_index = std.math.maxInt(usize);
561     const retried = try session.boundaryBytesAlloc(1, "env.TEST", "value");
562     defer failing.allocator().free(retried);
563     try std.testing.expectEqualStrings("value", retried);
564     try std.testing.expectEqual(@as(usize, 1), capture.len);
565     try std.testing.expectEqual(@as(u64, 1), capture.timepoints[0].seq);
566 }
567 
568 test "session failure atomicity retries boundary replay after allocation failure" {
569     const events = [_]event.Event{
570         event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "value"),
571     };
572     var stream = SliceStream{ .events = &events };
573     var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{
574         .fail_index = 0,
575     });
576     var session = Session.initReplay(failing.allocator(), stream.source());
577     defer session.deinit();
578 
579     try std.testing.expectError(
580         error.OutOfMemory,
581         session.boundaryBytesAlloc(1, "env.TEST", null),
582     );
583     try std.testing.expectEqual(@as(usize, 0), stream.cursor);
584     try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
585 
586     failing.fail_index = std.math.maxInt(usize);
587     failing.resize_fail_index = std.math.maxInt(usize);
588     const retried = try session.boundaryBytesAlloc(1, "env.TEST", null);
589     defer failing.allocator().free(retried);
590     try std.testing.expectEqualStrings("value", retried);
591     try session.verifyReplayComplete();
592 }
593 
594 test "session failure atomicity rejects noncanonical boundary fields" {
595     var events = [_]event.Event{
596         event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "value"),
597     };
598     events[0].label = "unexpected";
599     var stream = SliceStream{ .events = &events };
600     var session = Session.initReplay(std.testing.allocator, stream.source());
601     defer session.deinit();
602 
603     try std.testing.expectError(
604         error.ReplayEventMismatch,
605         session.boundaryBytesAlloc(1, "env.TEST", null),
606     );
607     try std.testing.expectEqual(@as(usize, 0), stream.cursor);
608     try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
609 
610     events[0].label = null;
611     const retried = try session.boundaryBytesAlloc(1, "env.TEST", null);
612     defer std.testing.allocator.free(retried);
613     try std.testing.expectEqualStrings("value", retried);
614     try session.verifyReplayComplete();
615 }
616 
617 test "session failure atomicity requires a live boundary value while off" {
618     var session = Session.initOff(std.testing.allocator);
619     defer session.deinit();
620 
621     try std.testing.expectError(
622         error.MissingBoundaryValue,
623         session.boundaryBytesAlloc(1, "env.TEST", null),
624     );
625     const retried = try session.boundaryBytesAlloc(1, "env.TEST", "value");
626     defer std.testing.allocator.free(retried);
627     try std.testing.expectEqualStrings("value", retried);
628     const next = try session.sessionEnd(1, 0);
629     try std.testing.expectEqual(@as(u64, 2), next.seq);
630 }
631 
632 test "replay mismatch leaves the session retryable" {
633     const events = [_]event.Event{
634         event.Event.safepointReached(.{ .thread_id = 1, .seq = 1 }, .{ .function_id = 20, .site_id = 5 }),
635     };
636     var stream = SliceStream{ .events = &events };
637     var session = Session.initReplay(std.testing.allocator, stream.source());
638     defer session.deinit();
639     try std.testing.expectError(
640         error.ReplayEventMismatch,
641         session.safepoint(1, .{ .function_id = 21, .site_id = 5 }),
642     );
643     try std.testing.expectEqual(@as(usize, 0), stream.cursor);
644     try std.testing.expectEqual(@as(u64, 0), session.replayProgress().cursor);
645     const retried = try session.safepoint(1, .{ .function_id = 20, .site_id = 5 });
646     try std.testing.expectEqual(@as(u64, 1), retried.seq);
647     try session.verifyReplayComplete();
648 }
649 
650 test "boundary replay copies borrowed bytes before advancing" {
651     const events = [_]event.Event{
652         event.Event.boundaryBytes(.{ .thread_id = 1, .seq = 1 }, "env.TEST", "hello"),
653     };
654     var stream = SliceStream{ .events = &events };
655     var session = Session.initReplay(std.testing.allocator, stream.source());
656     defer session.deinit();
657     const replayed = try session.boundaryBytesAlloc(1, "env.TEST", null);
658     defer std.testing.allocator.free(replayed);
659     try std.testing.expectEqualStrings("hello", replayed);
660     try session.verifyReplayComplete();
661 }
662 
663 test "session failure atomicity validates restore before staging provider state" {
664     const events = [_]event.Event{
665         event.Event.checkpoint(.{ .thread_id = 1, .seq = 1 }, "saved", "checkpoint"),
666         event.Event.checkpointRestore(.{ .thread_id = 1, .seq = 2 }, "restore"),
667     };
668     var stream = SliceStream{ .events = &events };
669     var provider_probe = ProviderProbe.init("checkpoint");
670     var session = Session.initReplay(std.testing.allocator, stream.source());
671     defer session.deinit();
672 
673     _ = try session.checkpointCapture(1, "saved", provider_probe.provider());
674     provider_probe.setValue("live");
675 
676     try std.testing.expectError(
677         error.ReplayEventMismatch,
678         session.checkpointRestore(1, "wrong", provider_probe.provider()),
679     );
680     try std.testing.expectEqualStrings("live", provider_probe.current());
681     try std.testing.expectEqual(@as(u64, 0), provider_probe.prepare_calls);
682     try std.testing.expectEqual(@as(usize, 1), stream.cursor);
683     try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
684 
685     provider_probe.fail_prepare = true;
686     try std.testing.expectError(
687         error.InjectedProviderFailure,
688         session.checkpointRestore(1, "restore", provider_probe.provider()),
689     );
690     try std.testing.expectEqualStrings("live", provider_probe.current());
691     try std.testing.expect(!provider_probe.prepared);
692     try std.testing.expectEqual(@as(u64, 1), provider_probe.cancel_calls);
693     try std.testing.expectEqual(@as(usize, 1), stream.cursor);
694     try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
695 
696     provider_probe.fail_prepare = false;
697     const restored = try session.checkpointRestore(1, "restore", provider_probe.provider());
698     try std.testing.expectEqual(@as(u64, 2), restored.seq);
699     try std.testing.expectEqualStrings("checkpoint", provider_probe.current());
700     try std.testing.expectEqual(@as(u64, 2), provider_probe.prepare_calls);
701     try std.testing.expectEqual(@as(u64, 1), provider_probe.commit_calls);
702     try session.verifyReplayComplete();
703 }
704 
705 test "session failure atomicity preserves a checkpoint when replacement allocation fails" {
706     const events = [_]event.Event{
707         event.Event.checkpoint(.{ .thread_id = 1, .seq = 1 }, "first", "one"),
708         event.Event.checkpoint(.{ .thread_id = 1, .seq = 2 }, "second", "two"),
709     };
710     var stream = SliceStream{ .events = &events };
711     var failing = std.testing.FailingAllocator.init(std.testing.allocator, .{});
712     var provider_probe = ProviderProbe.init("one");
713     var session = Session.initReplay(failing.allocator(), stream.source());
714     defer session.deinit();
715 
716     _ = try session.checkpointCapture(1, "first", provider_probe.provider());
717     try std.testing.expectEqualStrings("one", session.last_checkpoint.?.data.?);
718     provider_probe.setValue("two");
719     failing.fail_index = failing.alloc_index + 1;
720 
721     try std.testing.expectError(
722         error.OutOfMemory,
723         session.checkpointCapture(1, "second", provider_probe.provider()),
724     );
725     try std.testing.expectEqualStrings("one", session.last_checkpoint.?.data.?);
726     try std.testing.expectEqual(@as(usize, 1), stream.cursor);
727     try std.testing.expectEqual(@as(u64, 1), session.replayProgress().cursor);
728 
729     failing.fail_index = std.math.maxInt(usize);
730     failing.resize_fail_index = std.math.maxInt(usize);
731     const retried = try session.checkpointCapture(1, "second", provider_probe.provider());
732     try std.testing.expectEqual(@as(u64, 2), retried.seq);
733     try std.testing.expectEqualStrings("two", session.last_checkpoint.?.data.?);
734     try session.verifyReplayComplete();
735 }
736 
737 test "session failure atomicity cancels a staged restore when recording fails" {
738     var capture: FailOnceSink = .{ .fail_next = false };
739     var provider_probe = ProviderProbe.init("checkpoint");
740     var session = Session.initRecord(std.testing.allocator, capture.sink());
741     defer session.deinit();
742 
743     _ = try session.checkpointCapture(1, "saved", provider_probe.provider());
744     provider_probe.setValue("live");
745     capture.fail_next = true;
746 
747     try std.testing.expectError(
748         error.InjectedSinkFailure,
749         session.checkpointRestore(1, "restore", provider_probe.provider()),
750     );
751     try std.testing.expectEqualStrings("live", provider_probe.current());
752     try std.testing.expect(!provider_probe.prepared);
753     try std.testing.expectEqual(@as(u64, 1), provider_probe.cancel_calls);
754     try std.testing.expectEqual(@as(usize, 1), capture.len);
755 
756     const restored = try session.checkpointRestore(1, "restore", provider_probe.provider());
757     try std.testing.expectEqual(@as(u64, 2), restored.seq);
758     try std.testing.expectEqualStrings("checkpoint", provider_probe.current());
759     try std.testing.expectEqual(@as(usize, 2), capture.len);
760     try std.testing.expectEqual(@as(u64, 2), capture.timepoints[1].seq);
761 }
762 
763 test "checkpoint provider owns committed restore after checkpoint replacement" {
764     var provider_probe = ProviderProbe.init("one");
765     {
766         var session = Session.initOff(std.testing.allocator);
767         defer session.deinit();
768 
769         _ = try session.checkpointCapture(1, "first", provider_probe.provider());
770         provider_probe.setValue("live");
771         _ = try session.checkpointRestore(1, "restore", provider_probe.provider());
772         try std.testing.expectEqualStrings("one", provider_probe.current());
773         _ = try session.checkpointCapture(1, "second", provider_probe.provider());
774         try std.testing.expectEqualStrings("one", provider_probe.current());
775     }
776     try std.testing.expectEqualStrings("one", provider_probe.current());
777 }
778 
779 test "session sequence capacity admits the final timepoint exactly once" {
780     var session = Session.initOff(std.testing.allocator);
781     defer session.deinit();
782     session.next_seq = std.math.maxInt(u64);
783 
784     const final = try session.sessionStart(1, "last");
785     try std.testing.expectEqual(std.math.maxInt(u64), final.seq);
786     try std.testing.expectError(
787         error.SequenceExhausted,
788         session.sessionEnd(1, 0),
789     );
790 }