lib/tracy/src/flight.zig

daab053ee43316e1809a84551d573ddd1e5bf3d2

  1 const std = @import("std");
  2 const transport = @import("transport.zig");
  3 
  4 const assert = std.debug.assert;
  5 
  6 const schema = transport.schema;
  7 const OverflowPolicy = transport.OverflowPolicy;
  8 const State = transport.State;
  9 const Report = transport.Report;
 10 
 11 pub const FlightRecorder = struct {
 12     writer: std.Io.Writer,
 13     storage: []u8,
 14     event_storage: []u8,
 15     policy: OverflowPolicy,
 16     state: State = .accepting,
 17     header_seen: bool = false,
 18     header_len: usize = 0,
 19     ring_head: usize = 0,
 20     ring_len: usize = 0,
 21     ring_events: usize = 0,
 22     event_len: usize = 0,
 23     discarding_event: bool = false,
 24     observed_events: u64 = 0,
 25     stored_events: u64 = 0,
 26     overwritten_events: u64 = 0,
 27     dropped_events: u64 = 0,
 28     oversized_events: u64 = 0,
 29 
 30     pub fn init(
 31         storage: []u8,
 32         event_storage: []u8,
 33         writer_storage: []u8,
 34         policy: OverflowPolicy,
 35     ) FlightRecorder {
 36         assert(storage.len > 0);
 37         assert(event_storage.len > 0);
 38         assert(writer_storage.len > 0);
 39         assert(disjoint(storage, event_storage));
 40         assert(disjoint(storage, writer_storage));
 41         assert(disjoint(event_storage, writer_storage));
 42         var recorder: FlightRecorder = .{
 43             .writer = .{ .vtable = &writer_vtable, .buffer = writer_storage },
 44             .storage = storage,
 45             .event_storage = event_storage,
 46             .policy = policy,
 47         };
 48         recorder.assertValid();
 49         return recorder;
 50     }
 51 
 52     pub fn interface(self: *FlightRecorder) *std.Io.Writer {
 53         self.assertValid();
 54         return &self.writer;
 55     }
 56 
 57     pub fn reset(self: *FlightRecorder) void {
 58         self.writer.flush() catch unreachable;
 59         self.state = .accepting;
 60         self.header_seen = false;
 61         self.header_len = 0;
 62         self.ring_head = 0;
 63         self.ring_len = 0;
 64         self.ring_events = 0;
 65         self.event_len = 0;
 66         self.discarding_event = false;
 67         self.observed_events = 0;
 68         self.stored_events = 0;
 69         self.overwritten_events = 0;
 70         self.dropped_events = 0;
 71         self.oversized_events = 0;
 72         self.assertValid();
 73     }
 74 
 75     pub fn snapshot(self: *FlightRecorder, destination: *std.Io.Writer) !Report {
 76         if (destination == &self.writer) return error.InvalidSnapshotWriter;
 77         try self.writer.flush();
 78         self.assertValid();
 79         const result = self.report();
 80         if (self.header_len > 0) {
 81             try destination.writeAll(self.storage[0..self.header_len]);
 82         }
 83         const ring = self.ringStorage();
 84         if (self.ring_len > 0) {
 85             const first_len = @min(self.ring_len, ring.len - self.ring_head);
 86             try destination.writeAll(ring[self.ring_head..][0..first_len]);
 87             const second_len = self.ring_len - first_len;
 88             if (second_len > 0) try destination.writeAll(ring[0..second_len]);
 89         }
 90         try result.writeJsonl(destination);
 91         return result;
 92     }
 93 
 94     pub fn report(self: *FlightRecorder) Report {
 95         self.writer.flush() catch unreachable;
 96         self.assertValid();
 97         const retained_header: usize = if (self.header_len > 0) 1 else 0;
 98         return .{
 99             .policy = self.policy,
100             .state = self.state,
101             .capacity_bytes = self.storage.len,
102             .retained_bytes = self.header_len + self.ring_len,
103             .event_capacity_bytes = self.event_storage.len,
104             .writer_capacity_bytes = self.writer.buffer.len,
105             .observed_events = self.observed_events,
106             .stored_events = self.stored_events,
107             .retained_events = retained_header + self.ring_events,
108             .overwritten_events = self.overwritten_events,
109             .dropped_events = self.dropped_events,
110             .oversized_events = self.oversized_events,
111             .partial_event_bytes = self.event_len,
112             .discarding_oversized_event = self.discarding_event,
113         };
114     }
115 
116     fn accept(self: *FlightRecorder, bytes: []const u8) void {
117         for (bytes) |byte| {
118             if (self.discarding_event) {
119                 if (byte == '\n') {
120                     self.finishOversizedEvent();
121                 }
122                 continue;
123             }
124             if (byte == '\n') {
125                 self.commitEvent(self.event_storage[0..self.event_len]);
126                 self.event_len = 0;
127             } else if (self.event_len == self.event_storage.len) {
128                 self.event_len = 0;
129                 self.discarding_event = true;
130             } else {
131                 self.event_storage[self.event_len] = byte;
132                 self.event_len += 1;
133             }
134         }
135         self.assertValid();
136     }
137 
138     fn finishOversizedEvent(self: *FlightRecorder) void {
139         self.observed_events +|= 1;
140         self.oversized_events +|= 1;
141         self.discarding_event = false;
142         if (self.header_seen) return;
143         self.header_seen = true;
144         if (self.policy == .stop_when_full) self.state = .full;
145     }
146 
147     fn commitEvent(self: *FlightRecorder, payload: []const u8) void {
148         self.observed_events +|= 1;
149         const required = std.math.add(usize, payload.len, 1) catch {
150             self.oversized_events +|= 1;
151             return;
152         };
153         if (!self.header_seen) {
154             self.header_seen = true;
155             if (required > self.storage.len) {
156                 self.oversized_events +|= 1;
157                 if (self.policy == .stop_when_full) self.state = .full;
158                 return;
159             }
160             @memcpy(self.storage[0..payload.len], payload);
161             self.storage[payload.len] = '\n';
162             self.header_len = required;
163             self.stored_events +|= 1;
164             return;
165         }
166         self.commitRingEvent(payload, required);
167     }
168 
169     fn commitRingEvent(self: *FlightRecorder, payload: []const u8, required: usize) void {
170         if (self.state == .full) {
171             self.dropped_events +|= 1;
172             return;
173         }
174         const ring = self.ringStorage();
175         if (required > ring.len) {
176             self.oversized_events +|= 1;
177             if (self.policy == .stop_when_full) self.state = .full;
178             return;
179         }
180         if (self.policy == .stop_when_full and required > ring.len - self.ring_len) {
181             self.state = .full;
182             self.dropped_events +|= 1;
183             return;
184         }
185         while (required > ring.len - self.ring_len) self.evictOldest();
186         self.appendRing(payload);
187         self.stored_events +|= 1;
188         self.ring_events += 1;
189     }
190 
191     fn appendRing(self: *FlightRecorder, payload: []const u8) void {
192         const ring = self.ringStorage();
193         assert(payload.len + 1 <= ring.len - self.ring_len);
194         var tail = advance(self.ring_head, self.ring_len, ring.len);
195         const first_len = @min(payload.len, ring.len - tail);
196         @memcpy(ring[tail..][0..first_len], payload[0..first_len]);
197         const second_len = payload.len - first_len;
198         if (second_len > 0) @memcpy(ring[0..second_len], payload[first_len..]);
199         self.ring_len += payload.len;
200         tail = advance(self.ring_head, self.ring_len, ring.len);
201         ring[tail] = '\n';
202         self.ring_len += 1;
203     }
204 
205     fn evictOldest(self: *FlightRecorder) void {
206         const ring = self.ringStorage();
207         assert(self.ring_len > 0);
208         var event_bytes: usize = 1;
209         while (event_bytes <= self.ring_len) : (event_bytes += 1) {
210             const index = advance(self.ring_head, event_bytes - 1, ring.len);
211             if (ring[index] == '\n') break;
212         }
213         assert(event_bytes <= self.ring_len);
214         self.ring_head = advance(self.ring_head, event_bytes, ring.len);
215         self.ring_len -= event_bytes;
216         self.ring_events -= 1;
217         self.overwritten_events +|= 1;
218         if (self.ring_len == 0) self.ring_head = 0;
219     }
220 
221     fn ringStorage(self: *FlightRecorder) []u8 {
222         return self.storage[self.header_len..];
223     }
224 
225     fn assertValid(self: *FlightRecorder) void {
226         assert(self.storage.len > 0);
227         assert(self.event_storage.len > 0);
228         assert(self.writer.buffer.len > 0);
229         assert(self.header_len <= self.storage.len);
230         assert(self.ring_len <= self.storage.len - self.header_len);
231         assert(self.event_len <= self.event_storage.len);
232         if (self.discarding_event) assert(self.event_len == 0);
233         assert(self.ring_events <= self.ring_len);
234         if (self.state == .full) assert(self.policy == .stop_when_full);
235         if (self.ring_len == 0) {
236             assert(self.ring_head == 0);
237         } else {
238             assert(self.ringStorage().len > 0);
239             assert(self.ring_head < self.ringStorage().len);
240         }
241         const header_events: usize = if (self.header_len > 0) 1 else 0;
242         if (std.math.cast(u64, header_events + self.ring_events)) |retained| {
243             const accounted = std.math.add(u64, self.overwritten_events, retained) catch null;
244             if (accounted) |count| {
245                 if (self.stored_events != std.math.maxInt(u64)) {
246                     assert(self.stored_events == count);
247                 }
248             }
249         }
250         if (self.observed_events != std.math.maxInt(u64)) {
251             const accepted = std.math.add(u64, self.stored_events, self.dropped_events) catch null;
252             if (accepted) |count| {
253                 const accounted = std.math.add(u64, count, self.oversized_events) catch null;
254                 if (accounted) |total| assert(self.observed_events == total);
255             }
256         }
257     }
258 
259     const writer_vtable: std.Io.Writer.VTable = .{
260         .drain = drain,
261         .flush = std.Io.Writer.defaultFlush,
262         .rebase = std.Io.Writer.defaultRebase,
263     };
264 
265     fn drain(
266         writer: *std.Io.Writer,
267         data: []const []const u8,
268         splat: usize,
269     ) std.Io.Writer.Error!usize {
270         const self: *FlightRecorder = @alignCast(@fieldParentPtr("writer", writer));
271         assert(data.len > 0);
272         self.accept(writer.buffer[0..writer.end]);
273         writer.end = 0;
274         var consumed: usize = 0;
275         for (data[0 .. data.len - 1]) |bytes| {
276             self.accept(bytes);
277             consumed += bytes.len;
278         }
279         const last = data[data.len - 1];
280         for (0..splat) |_| {
281             self.accept(last);
282             consumed += last.len;
283         }
284         return consumed;
285     }
286 };
287 
288 fn disjoint(left: []const u8, right: []const u8) bool {
289     const left_start = @intFromPtr(left.ptr);
290     const right_start = @intFromPtr(right.ptr);
291     const left_end = std.math.add(usize, left_start, left.len) catch return false;
292     const right_end = std.math.add(usize, right_start, right.len) catch return false;
293     return left_end <= right_start or right_end <= left_start;
294 }
295 
296 fn advance(start: usize, amount: usize, capacity: usize) usize {
297     assert(capacity > 0);
298     assert(start < capacity);
299     assert(amount <= capacity);
300     const until_end = capacity - start;
301     if (amount < until_end) return start + amount;
302     return amount - until_end;
303 }
304 
305 fn expectSnapshot(
306     snapshot: []const u8,
307     expected_events: []const u8,
308     expected_report: Report,
309 ) !void {
310     try std.testing.expectEqual(expected_report.retained_bytes, expected_events.len);
311     try std.testing.expect(snapshot.len > expected_events.len);
312     try std.testing.expectEqualStrings(expected_events, snapshot[0..expected_events.len]);
313     const actual_report = try transport.parseLine(
314         std.testing.allocator,
315         snapshot[expected_events.len..],
316     );
317     try std.testing.expectEqualDeep(expected_report, actual_report);
318 }
319 
320 test "flight recorder preserves the header and latest complete events" {
321     var storage: [18]u8 = undefined;
322     var events: [8]u8 = undefined;
323     var writer_buffer: [5]u8 = undefined;
324     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
325     try recorder.interface().writeAll("header\none\ntwo22\nthree\n");
326 
327     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
328     defer snapshot.deinit();
329     const result = try recorder.snapshot(&snapshot.writer);
330     try expectSnapshot(snapshot.written(), "header\nthree\n", result);
331     try std.testing.expectEqual(@as(u64, 4), result.observed_events);
332     try std.testing.expectEqual(@as(u64, 2), result.overwritten_events);
333     try std.testing.expectEqual(@as(usize, 2), result.retained_events);
334 }
335 
336 test "flight recorder overwrite window matches an independent suffix model" {
337     const rows = [_][]const u8{ "h\n", "a\n", "b\n", "c\n", "d\n", "e\n" };
338     var storage: [16]u8 = undefined;
339     var events: [8]u8 = undefined;
340     var writer_buffer: [3]u8 = undefined;
341     var capacity: usize = 2;
342     while (capacity <= storage.len) : (capacity += 1) {
343         var recorder = FlightRecorder.init(
344             storage[0..capacity],
345             &events,
346             &writer_buffer,
347             .overwrite_oldest,
348         );
349         for (rows) |row| try recorder.interface().writeAll(row);
350         var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
351         defer snapshot.deinit();
352         const result = try recorder.snapshot(&snapshot.writer);
353 
354         var expected = std.Io.Writer.Allocating.init(std.testing.allocator);
355         defer expected.deinit();
356         try expected.writer.writeAll(rows[0]);
357         const suffix_count = @min(rows.len - 1, (capacity - rows[0].len) / 2);
358         for (rows[rows.len - suffix_count ..]) |row| try expected.writer.writeAll(row);
359         try expectSnapshot(snapshot.written(), expected.written(), result);
360         try std.testing.expect(result.retained_bytes <= capacity);
361         try std.testing.expectEqual(@as(u64, rows.len), result.observed_events);
362         try std.testing.expectEqual(1 + suffix_count, result.retained_events);
363     }
364 }
365 
366 test "flight recorder stop policy preserves the earliest complete window" {
367     var storage: [16]u8 = undefined;
368     var events: [8]u8 = undefined;
369     var writer_buffer: [5]u8 = undefined;
370     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .stop_when_full);
371     try recorder.interface().writeAll("h00\none\ntwo\ntri\nend\nmore\n");
372 
373     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
374     defer snapshot.deinit();
375     const result = try recorder.snapshot(&snapshot.writer);
376     try expectSnapshot(snapshot.written(), "h00\none\ntwo\ntri\n", result);
377     try std.testing.expectEqual(State.full, result.state);
378     try std.testing.expectEqual(@as(u64, 2), result.dropped_events);
379     try std.testing.expectEqual(@as(u64, 0), result.overwritten_events);
380 }
381 
382 test "flight recorder drops oversized rows without corrupting later rows" {
383     var storage: [16]u8 = undefined;
384     var events: [4]u8 = undefined;
385     var writer_buffer: [3]u8 = undefined;
386     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
387     try recorder.interface().writeAll("h\n12345\nok\n");
388 
389     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
390     defer snapshot.deinit();
391     const result = try recorder.snapshot(&snapshot.writer);
392     try expectSnapshot(snapshot.written(), "h\nok\n", result);
393     try std.testing.expectEqual(@as(u64, 3), result.observed_events);
394     try std.testing.expectEqual(@as(u64, 1), result.oversized_events);
395     try std.testing.expectEqual(@as(u64, 2), result.stored_events);
396 }
397 
398 test "flight recorder snapshots analyzable tracy jsonl" {
399     const event = @import("event.zig");
400     const record_mod = @import("record.zig");
401     var storage: [1024]u8 = undefined;
402     var events: [512]u8 = undefined;
403     var writer_buffer: [256]u8 = undefined;
404     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
405     try (event.TraceEvent{ .seq = 1, .kind = .start, .time_ns = 10, .thread = 1 })
406         .writeJsonLine(recorder.interface());
407     try (event.TraceEvent{ .seq = 2, .kind = .frame, .time_ns = 20, .thread = 1 })
408         .writeJsonLine(recorder.interface());
409 
410     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
411     defer snapshot.deinit();
412     _ = try recorder.snapshot(&snapshot.writer);
413     var lines = std.mem.tokenizeScalar(u8, snapshot.written(), '\n');
414     var event_count: usize = 0;
415     var report_count: usize = 0;
416     while (lines.next()) |line| {
417         var parsed = try record_mod.parseLine(std.testing.allocator, line);
418         defer parsed.deinit();
419         switch (parsed) {
420             .event => event_count += 1,
421             .flight => report_count += 1,
422         }
423     }
424     try std.testing.expectEqual(@as(usize, 2), event_count);
425     try std.testing.expectEqual(@as(usize, 1), report_count);
426 }
427 
428 test "flight recorder overwrite is visible in tracy capture integrity" {
429     const event = @import("event.zig");
430     const summary = @import("summary.zig");
431     var storage: [1024]u8 = undefined;
432     var events: [512]u8 = undefined;
433     var writer_buffer: [256]u8 = undefined;
434     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
435     try (event.TraceEvent{ .seq = 1, .kind = .start }).writeJsonLine(recorder.interface());
436     for (2..34) |sequence| {
437         try (event.TraceEvent{
438             .seq = sequence,
439             .kind = .message,
440             .name = "flight event",
441         }).writeJsonLine(recorder.interface());
442     }
443     try (event.TraceEvent{ .seq = 34, .kind = .stop }).writeJsonLine(recorder.interface());
444 
445     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
446     defer snapshot.deinit();
447     const snapshot_report = try recorder.snapshot(&snapshot.writer);
448     var analyzer = summary.Analyzer.init(std.testing.allocator);
449     defer analyzer.deinit();
450     try analyzer.ingestJsonlBytes(snapshot.written());
451     const integrity = analyzer.captureIntegrity();
452     try std.testing.expect(snapshot_report.overwritten_events > 0);
453     try std.testing.expectEqualDeep(snapshot_report, integrity.flight_report.?);
454     try std.testing.expectEqualStrings("sequence_gaps", integrity.status);
455     try std.testing.expect(integrity.missing_sequence_event_count > 0);
456     try std.testing.expectEqual(@as(u64, 1), integrity.start_event_count);
457     try std.testing.expectEqual(@as(u64, 1), integrity.stop_event_count);
458 }
459 
460 test "enabled tracy runtime records through the flight recorder" {
461     const build_options = @import("build_options");
462     if (!build_options.enabled) return;
463     const instrumentation = @import("instrumentation.zig");
464     const summary = @import("summary.zig");
465     var storage: [4096]u8 = undefined;
466     var events: [1024]u8 = undefined;
467     var writer_buffer: [512]u8 = undefined;
468     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
469     try std.testing.expect(try instrumentation.start(recorder.interface(), .{
470         .name = "flight test",
471     }));
472     defer instrumentation.stop();
473     const active_zone = instrumentation.zone("flight.test.zone");
474     active_zone.end();
475     instrumentation.stop();
476 
477     var snapshot = std.Io.Writer.Allocating.init(std.testing.allocator);
478     defer snapshot.deinit();
479     const snapshot_report = try recorder.snapshot(&snapshot.writer);
480     var analyzer = summary.Analyzer.init(std.testing.allocator);
481     defer analyzer.deinit();
482     try analyzer.ingestJsonlBytes(snapshot.written());
483     try std.testing.expectEqual(@as(u64, 1), analyzer.counters.completed_zones);
484     try std.testing.expectEqual(@as(u64, 4), snapshot_report.observed_events);
485     try std.testing.expectEqualDeep(snapshot_report, analyzer.captureIntegrity().flight_report.?);
486 }
487 
488 test "flight recorder report is machine readable and resettable" {
489     var storage: [16]u8 = undefined;
490     var events: [8]u8 = undefined;
491     var writer_buffer: [4]u8 = undefined;
492     var recorder = FlightRecorder.init(&storage, &events, &writer_buffer, .overwrite_oldest);
493     try recorder.interface().writeAll("head\nrow\npartial");
494     const before = recorder.report();
495     try std.testing.expectEqual(@as(usize, 7), before.partial_event_bytes);
496 
497     var jsonl = std.Io.Writer.Allocating.init(std.testing.allocator);
498     defer jsonl.deinit();
499     try before.writeJsonl(&jsonl.writer);
500     var parsed = try std.json.parseFromSlice(
501         std.json.Value,
502         std.testing.allocator,
503         jsonl.written(),
504         .{},
505     );
506     defer parsed.deinit();
507     try std.testing.expectEqualStrings(
508         schema,
509         parsed.value.object.get("schema").?.string,
510     );
511     recorder.reset();
512     const after = recorder.report();
513     try std.testing.expectEqual(@as(u64, 0), after.observed_events);
514     try std.testing.expectEqual(@as(usize, 0), after.retained_bytes);
515 }