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 }