gate4agent_runtime_native/tick_profile.rs
1//! Fixed-window latency profiling for the phases inside
2//! [`crate::NativeRuntime::tick`].
3//!
4//! Before this module existed, "the node burns CPU with zero live PTY
5//! sessions" had no number attached to it anywhere in the process --
6//! `tick()` ran six distinct pieces of work every call (drain observations,
7//! drain ingress, reduce the control plane, dispatch effects, publish the
8//! step, and walk the provider supervisors) and none of them were timed.
9//! This module gives each of those six phases its own always-on reading.
10//!
11//! Two contracts every series here honours, mirrored from
12//! `hatchery-tui`'s `profile.rs` (that crate is a different workspace and
13//! is not a dependency of this one, so the types are re-derived here rather
14//! than imported -- see that module's own doc comment for the same
15//! reasoning spelled out for the TUI's redraw loop):
16//! - **Distributions, not an average.** A mean spreads one expensive tick
17//! across 255 idle ones and reports a number nobody can act on. Every
18//! phase keeps the last [`SAMPLE_WINDOW`] samples in a fixed ring
19//! ([`RingStats`]) and reports p50/p95/max plus the sample count,
20//! computed on demand rather than tracked incrementally.
21//! - **Always-on cheap, no allocation after construction.** `push` is an
22//! array write and an index bump. The O(N log N) sort behind `stats()`
23//! only runs when something actually reads a snapshot (an HTTP `/metrics`
24//! request), never once per tick -- so this stays cheap enough to leave
25//! enabled unconditionally, which is the only way it can ever catch the
26//! stall it exists to find.
27use std::time::Duration;
28
29/// Ring capacity shared by every phase below -- "the last 256 ticks," never
30/// an average across the process lifetime. At the ~10ms drive-loop cadence
31/// this is a little over 2.5 seconds of tick history, enough to catch a
32/// transient stall without growing unbounded.
33pub const SAMPLE_WINDOW: usize = 256;
34
35/// One windowed phase's current reading: nearest-rank p50/p95/max over
36/// whatever [`RingStats`] currently holds, plus `count` -- the denominator
37/// a reader must check alongside the three numbers, since `count` below
38/// `SAMPLE_WINDOW` means "the process hasn't produced a full window yet,"
39/// not "the window is smaller than advertised."
40#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
41pub struct Distribution {
42 pub p50: u32,
43 pub p95: u32,
44 pub max: u32,
45 pub count: usize,
46}
47
48/// Fixed-size ring buffer of `u32` samples. `push` overwrites the oldest
49/// entry once full -- the only state this holds is `values`/`len`/`next`,
50/// all stack-sized by the const generic `N`, so a profiler holding several
51/// of these never allocates past its own construction.
52#[derive(Clone, Debug)]
53pub struct RingStats<const N: usize> {
54 values: [u32; N],
55 len: usize,
56 next: usize,
57}
58
59impl<const N: usize> Default for RingStats<N> {
60 fn default() -> Self {
61 Self { values: [0; N], len: 0, next: 0 }
62 }
63}
64
65impl<const N: usize> RingStats<N> {
66 pub fn push(&mut self, value: u32) {
67 self.values[self.next] = value;
68 self.next = (self.next + 1) % N;
69 self.len = (self.len + 1).min(N);
70 }
71
72 /// Nearest-rank percentiles over whatever the ring currently holds.
73 /// Percentiles do not care about chronological order, so this sorts a
74 /// stack-local COPY of the valid slice (`values` is `[u32; N]`, `Copy`
75 /// because `u32` is `Copy` -- never a heap allocation) rather than the
76 /// ring itself, which must keep its own write position intact for the
77 /// next `push`.
78 pub fn stats(&self) -> Distribution {
79 if self.len == 0 {
80 return Distribution::default();
81 }
82 let mut sorted = self.values;
83 sorted[..self.len].sort_unstable();
84 let rank = |percentile: usize| sorted[(self.len * percentile / 100).min(self.len - 1)];
85 Distribution { p50: rank(50), p95: rank(95), max: sorted[self.len - 1], count: self.len }
86 }
87}
88
89/// Clamped `Duration` -> microsecond sample: a tick phase measured in hours
90/// would mean the process already hung far worse than this profiler needs
91/// to describe, so this saturates at `u32::MAX` rather than panicking or
92/// widening every sample to a `u128`.
93pub fn duration_micros(elapsed: Duration) -> u32 {
94 elapsed.as_micros().min(u128::from(u32::MAX)) as u32
95}
96
97/// One snapshot of every phase's current distribution -- what
98/// [`crate::NativeRuntime::tick_profile_snapshot`] returns.
99#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
100pub struct TickProfileSnapshot {
101 pub drain_observations_us: Distribution,
102 pub drain_ingress_us: Distribution,
103 pub step_control_plane_us: Distribution,
104 pub dispatch_effects_us: Distribution,
105 pub publish_step_us: Distribution,
106 pub provider_supervisors_us: Distribution,
107}
108
109/// Owns the six per-phase rings `NativeRuntime::tick` writes into every
110/// call. Lives on `NativeRuntime` itself (not a side channel) so the timing
111/// and the work it describes can never drift apart.
112#[derive(Clone, Debug, Default)]
113pub struct TickPhaseProfiler {
114 drain_observations_us: RingStats<SAMPLE_WINDOW>,
115 drain_ingress_us: RingStats<SAMPLE_WINDOW>,
116 step_control_plane_us: RingStats<SAMPLE_WINDOW>,
117 dispatch_effects_us: RingStats<SAMPLE_WINDOW>,
118 publish_step_us: RingStats<SAMPLE_WINDOW>,
119 provider_supervisors_us: RingStats<SAMPLE_WINDOW>,
120}
121
122impl TickPhaseProfiler {
123 /// `effects.drain_observations` -- pulling completed native-provider
124 /// observations (and any terminal frames riding along with them) off
125 /// the effect dispatcher before the kernel reduces this tick.
126 pub fn record_drain_observations(&mut self, elapsed: Duration) {
127 self.drain_observations_us.push(duration_micros(elapsed));
128 }
129
130 /// `port.drain_ingress` -- pulling queued operator/control commands off
131 /// the bounded control-plane port.
132 pub fn record_drain_ingress(&mut self, elapsed: Duration) {
133 self.drain_ingress_us.push(duration_micros(elapsed));
134 }
135
136 /// `kernel.step_control_plane` -- the synchronous kernel reduction
137 /// itself: ingress and observations in, effects and a fresh backend
138 /// snapshot out. This phase is the leading suspect for idle-tick cost,
139 /// since `step_control_plane` builds a full `KernelStep` (including a
140 /// backend snapshot clone) unconditionally on every call, whether or
141 /// not ingress/observations carried anything to reduce.
142 pub fn record_step_control_plane(&mut self, elapsed: Duration) {
143 self.step_control_plane_us.push(duration_micros(elapsed));
144 }
145
146 /// The loop dispatching this tick's own kernel effects to per-instance
147 /// native workers (`effects.dispatch`).
148 pub fn record_dispatch_effects(&mut self, elapsed: Duration) {
149 self.dispatch_effects_us.push(duration_micros(elapsed));
150 }
151
152 /// `port.publish_step` -- publishing the step's snapshot/events back
153 /// out through the control-plane port.
154 pub fn record_publish_step(&mut self, elapsed: Duration) {
155 self.publish_step_us.push(duration_micros(elapsed));
156 }
157
158 /// Walking every installed `ProviderSupervisor::tick` plus draining
159 /// their exit-ack/fault queues (`collect_provider_supervisor_events`).
160 pub fn record_provider_supervisors(&mut self, elapsed: Duration) {
161 self.provider_supervisors_us.push(duration_micros(elapsed));
162 }
163
164 pub fn snapshot(&self) -> TickProfileSnapshot {
165 TickProfileSnapshot {
166 drain_observations_us: self.drain_observations_us.stats(),
167 drain_ingress_us: self.drain_ingress_us.stats(),
168 step_control_plane_us: self.step_control_plane_us.stats(),
169 dispatch_effects_us: self.dispatch_effects_us.stats(),
170 publish_step_us: self.publish_step_us.stats(),
171 provider_supervisors_us: self.provider_supervisors_us.stats(),
172 }
173 }
174}
175
176#[cfg(test)]
177mod tests {
178 use super::*;
179
180 #[test]
181 fn ring_stats_reports_nearest_rank_percentiles() {
182 let mut ring: RingStats<8> = RingStats::default();
183 for value in [10, 20, 30, 40, 50, 60, 70, 80] {
184 ring.push(value);
185 }
186 let stats = ring.stats();
187 assert_eq!(stats.count, 8);
188 assert_eq!(stats.max, 80);
189 assert_eq!(stats.p50, 50);
190 }
191
192 #[test]
193 fn ring_stats_evicts_oldest_once_full() {
194 let mut ring: RingStats<4> = RingStats::default();
195 for value in 1..=6u32 {
196 ring.push(value);
197 }
198 let stats = ring.stats();
199 // Only the last 4 pushes (3, 4, 5, 6) survive.
200 assert_eq!(stats.count, 4);
201 assert_eq!(stats.max, 6);
202 assert_eq!(stats.p50, 5);
203 }
204
205 #[test]
206 fn empty_ring_reports_zeroed_distribution() {
207 let ring: RingStats<SAMPLE_WINDOW> = RingStats::default();
208 assert_eq!(ring.stats(), Distribution::default());
209 }
210
211 #[test]
212 fn tick_phase_profiler_snapshot_reflects_every_recorded_phase() {
213 let mut profiler = TickPhaseProfiler::default();
214 profiler.record_drain_observations(Duration::from_micros(10));
215 profiler.record_drain_ingress(Duration::from_micros(20));
216 profiler.record_step_control_plane(Duration::from_micros(1_500));
217 profiler.record_dispatch_effects(Duration::from_micros(30));
218 profiler.record_publish_step(Duration::from_micros(40));
219 profiler.record_provider_supervisors(Duration::from_micros(50));
220 let snapshot = profiler.snapshot();
221 assert_eq!(snapshot.drain_observations_us.max, 10);
222 assert_eq!(snapshot.drain_ingress_us.max, 20);
223 assert_eq!(snapshot.step_control_plane_us.max, 1_500);
224 assert_eq!(snapshot.dispatch_effects_us.max, 30);
225 assert_eq!(snapshot.publish_step_us.max, 40);
226 assert_eq!(snapshot.provider_supervisors_us.max, 50);
227 for distribution in [
228 snapshot.drain_observations_us,
229 snapshot.drain_ingress_us,
230 snapshot.step_control_plane_us,
231 snapshot.dispatch_effects_us,
232 snapshot.publish_step_us,
233 snapshot.provider_supervisors_us,
234 ] {
235 assert_eq!(distribution.count, 1);
236 }
237 }
238}