lib/coz/src/sampler.zig
daab053ee43316e1809a84551d573ddd1e5bf3d2
1 const std = @import("std");
2 const sys = @import("sys");
3
4 const experiment = @import("experiment.zig");
5 const perf = sys.perf;
6 const timer = sys.timer;
7
8 pub const default_sample_batch_size: u32 = @intCast(experiment.sample_batch_size);
9 pub const WaitFn = *const fn (u64) u64;
10
11 pub const LossCounter = union(enum) {
12 available: u64,
13 unsupported,
14 read_failed,
15 };
16
17 pub const Options = struct {
18 sample_period_ns: u64 = experiment.sample_period_ns,
19 sample_batch_size: u32 = default_sample_batch_size,
20 signal: timer.Signal = timer.sample_signal,
21 sample_type: u64 = perf.default_sample_type,
22 read_format: u64 = 0,
23 exclude_kernel: bool = true,
24 exclude_idle: bool = true,
25
26 pub fn perfOptions(self: Options) perf.SamplerOptions {
27 return .{
28 .sample_period_ns = self.sample_period_ns,
29 .sample_batch_size = self.sample_batch_size,
30 .sample_type = self.sample_type,
31 .read_format = self.read_format,
32 .exclude_kernel = self.exclude_kernel,
33 .exclude_idle = self.exclude_idle,
34 };
35 }
36 };
37
38 pub const Sampler = struct {
39 event: perf.Event = .{},
40 wake_timer: timer.Timer = .{},
41 wake_interval_ns: u64 = 0,
42 loss_counter_supported: bool = false,
43
44 pub fn openCurrentThread(options: Options) !Sampler {
45 try validateOptions(options);
46 const interval_ns = try wakeIntervalNs(options);
47
48 var opened = try openPerfEvent(options);
49 errdefer opened.event.close();
50
51 var wake_timer = try timer.Timer.createForCurrentThread(options.signal);
52 errdefer wake_timer.close();
53
54 return .{
55 .event = opened.event,
56 .wake_timer = wake_timer,
57 .wake_interval_ns = interval_ns,
58 .loss_counter_supported = opened.loss_counter_supported,
59 };
60 }
61
62 pub fn close(self: *Sampler) void {
63 self.wake_timer.close();
64 self.event.close();
65 self.wake_interval_ns = 0;
66 self.loss_counter_supported = false;
67 }
68
69 pub fn start(self: *Sampler) !void {
70 try self.event.start();
71 errdefer self.event.stop() catch {};
72 try self.wake_timer.startInterval(self.wake_interval_ns);
73 }
74
75 pub fn stop(self: *Sampler) !void {
76 var first_error: ?anyerror = null;
77 self.wake_timer.stop() catch |err| rememberError(&first_error, err);
78 self.event.stop() catch |err| rememberError(&first_error, err);
79 if (first_error) |err| return err;
80 }
81
82 pub fn ringReader(self: *Sampler) ?perf.RingReader {
83 return self.event.ringReader();
84 }
85
86 pub fn commitReader(self: *Sampler, reader: perf.RingReader) void {
87 self.event.commitReader(reader);
88 }
89
90 pub fn readLossCounter(self: *Sampler) LossCounter {
91 if (!self.loss_counter_supported) return .unsupported;
92 const result = self.event.countResult() catch return .read_failed;
93 return .{ .available = result.lost orelse return .read_failed };
94 }
95
96 pub fn pausingWait(self: *Sampler, wait_fn: WaitFn) PausingWait {
97 return .{ .sampler = self, .wait_fn = wait_fn };
98 }
99 };
100
101 const OpenedPerfEvent = struct {
102 event: perf.Event,
103 loss_counter_supported: bool,
104 };
105
106 fn openPerfEvent(options: Options) !OpenedPerfEvent {
107 if (comptime !perf.supported) return error.UnsupportedPlatform;
108 var perf_options = options.perfOptions();
109 if (perf_options.read_format & perf.ReadFormat.group != 0) {
110 return .{
111 .event = try perf.Event.openTaskClockSampler(perf_options),
112 .loss_counter_supported = false,
113 };
114 }
115
116 perf_options.read_format |= perf.ReadFormat.lost;
117 const event = perf.Event.openTaskClockSampler(perf_options) catch |err| switch (err) {
118 error.InvalidPerfEventOptions => fallback: {
119 perf_options.read_format &= ~perf.ReadFormat.lost;
120 break :fallback try perf.Event.openTaskClockSampler(perf_options);
121 },
122 else => return err,
123 };
124 return .{
125 .event = event,
126 .loss_counter_supported = event.config.read_format & perf.ReadFormat.lost != 0,
127 };
128 }
129
130 pub const PausingWait = struct {
131 sampler: *Sampler,
132 wait_fn: WaitFn,
133
134 pub fn wait(self: PausingWait, ns: u64) u64 {
135 self.sampler.stop() catch {};
136 defer self.sampler.start() catch {};
137 return self.wait_fn(ns);
138 }
139 };
140
141 pub fn validateOptions(options: Options) !void {
142 if (options.sample_period_ns == 0) return error.InvalidSamplePeriod;
143 if (options.sample_batch_size == 0) return error.InvalidSampleBatchSize;
144 _ = try wakeIntervalNs(options);
145 }
146
147 pub fn wakeIntervalNs(options: Options) !u64 {
148 return std.math.mul(u64, options.sample_period_ns, @as(u64, options.sample_batch_size));
149 }
150
151 pub fn unavailable(err: anyerror) bool {
152 return switch (err) {
153 error.UnsupportedPlatform,
154 error.PermissionDenied,
155 error.DeviceBusy,
156 error.ProcessResources,
157 error.EventRequiresUnsupportedCpuFeature,
158 error.TooManyBreakpoints,
159 error.SampleStackNotSupported,
160 error.EventNotSupported,
161 error.SampleMaxStackOverflow,
162 error.ProcessNotFound,
163 error.SystemResources,
164 error.TooBig,
165 => true,
166 else => false,
167 };
168 }
169
170 fn rememberError(first_error: *?anyerror, err: anyerror) void {
171 if (first_error.* == null) first_error.* = err;
172 }
173
174 var test_wait_calls: std.atomic.Value(u32) = .init(0);
175
176 fn countedWait(ns: u64) u64 {
177 _ = test_wait_calls.fetchAdd(1, .monotonic);
178 return ns + 1;
179 }
180
181 test "sampler options preserve upstream sampling cadence" {
182 const options: Options = .{};
183 const perf_options = options.perfOptions();
184
185 try std.testing.expectEqual(experiment.sample_period_ns, options.sample_period_ns);
186 try std.testing.expectEqual(default_sample_batch_size, options.sample_batch_size);
187 try std.testing.expectEqual(timer.sample_signal, options.signal);
188 try std.testing.expectEqual(perf.default_sample_type, options.sample_type);
189 try std.testing.expectEqual(experiment.sample_period_ns * experiment.sample_batch_size, try wakeIntervalNs(options));
190
191 try std.testing.expectEqual(options.sample_period_ns, perf_options.sample_period_ns);
192 try std.testing.expectEqual(options.sample_batch_size, perf_options.sample_batch_size);
193 try std.testing.expectEqual(options.sample_type, perf_options.sample_type);
194 try std.testing.expectEqual(options.read_format, perf_options.read_format);
195 try std.testing.expectEqual(options.exclude_kernel, perf_options.exclude_kernel);
196 try std.testing.expectEqual(options.exclude_idle, perf_options.exclude_idle);
197 }
198
199 test "sampler rejects invalid sampling cadence" {
200 try std.testing.expectError(error.InvalidSamplePeriod, validateOptions(.{ .sample_period_ns = 0 }));
201 try std.testing.expectError(error.InvalidSampleBatchSize, validateOptions(.{ .sample_batch_size = 0 }));
202 try std.testing.expectError(error.Overflow, wakeIntervalNs(.{
203 .sample_period_ns = std.math.maxInt(u64),
204 .sample_batch_size = 2,
205 }));
206 }
207
208 test "sampler close is idempotent without owned resources" {
209 var sample: Sampler = .{};
210
211 sample.close();
212 sample.close();
213
214 try std.testing.expectEqual(perf.invalid_fd, sample.event.fd);
215 try std.testing.expectEqual(timer.invalid_timer_id, sample.wake_timer.id);
216 try std.testing.expectEqual(@as(u64, 0), sample.wake_interval_ns);
217 try std.testing.expect(!sample.loss_counter_supported);
218 }
219
220 test "sampler start and stop reject unopened resources" {
221 var sample: Sampler = .{};
222
223 try std.testing.expectError(error.InvalidPerfEvent, sample.start());
224 try std.testing.expectError(error.UninitializedTimer, sample.stop());
225 }
226
227 test "sampler pausing wait delegates to wrapped wait function" {
228 var sample: Sampler = .{};
229 test_wait_calls.store(0, .monotonic);
230
231 const paused = sample.pausingWait(countedWait);
232
233 try std.testing.expectEqual(@as(u64, 8), paused.wait(7));
234 try std.testing.expectEqual(@as(u32, 1), test_wait_calls.load(.monotonic));
235 }
236
237 test "sampler opens starts stops and closes current thread resources" {
238 var sample = Sampler.openCurrentThread(.{
239 .sample_period_ns = std.time.ns_per_s * 60,
240 .sample_batch_size = 1,
241 }) catch |err| {
242 if (unavailable(err)) return error.SkipZigTest;
243 return err;
244 };
245 defer sample.close();
246
247 try std.testing.expect(sample.event.fd != perf.invalid_fd);
248 try std.testing.expect(sample.wake_timer.id != timer.invalid_timer_id);
249 try std.testing.expectEqual(std.time.ns_per_s * 60, sample.wake_interval_ns);
250
251 try sample.start();
252 try sample.stop();
253 try std.testing.expect(sample.readLossCounter() != .read_failed);
254 }
255
256 test "sampler reports unsupported perf platform before loss negotiation" {
257 if (comptime perf.supported) return;
258 try std.testing.expectError(error.UnsupportedPlatform, openPerfEvent(.{}));
259 }