Skip to main content

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}