Skip to main content

performance_envelope/
performance_envelope.rs

1use std::{
2    error::Error,
3    time::{Duration, Instant},
4};
5
6use omena_reactive::{
7    ChangePolicyV0, ReactiveEngineV0, ReactiveGraphBuilderV0, ReactiveNodeIdV0, ReactiveStateV0,
8    ReactiveValueV0, StabilizeStatusV0,
9};
10
11const INPUT_COUNT: usize = 256;
12const EVENT_COUNT: usize = 512;
13const P95_CEILING: Duration = Duration::from_millis(250);
14const MAX_CEILING: Duration = Duration::from_millis(750);
15
16fn project_counter(state: &ReactiveStateV0) -> ReactiveStateV0 {
17    match state {
18        ReactiveStateV0::Available(ReactiveValueV0::Counter(value)) => {
19            ReactiveStateV0::available(ReactiveValueV0::Counter(value.rotate_left(7)))
20        }
21        _ => ReactiveStateV0::unavailable("counterRequired", "expected a counter deposit"),
22    }
23}
24
25fn settle(engine: &mut ReactiveEngineV0) -> Result<(), Box<dyn Error>> {
26    loop {
27        if matches!(
28            engine.stabilize_step(4_096)?,
29            StabilizeStatusV0::Settled { .. }
30        ) {
31            return Ok(());
32        }
33    }
34}
35
36fn graph() -> Result<(ReactiveEngineV0, Vec<ReactiveNodeIdV0>), Box<dyn Error>> {
37    let policy = ChangePolicyV0::exact("performanceSemanticValue");
38    let mut graph = ReactiveGraphBuilderV0::new();
39    let mut inputs = Vec::with_capacity(INPUT_COUNT);
40    let mut projections = Vec::with_capacity(INPUT_COUNT);
41    for index in 0..INPUT_COUNT {
42        let input = graph.add_input(
43            ReactiveStateV0::available(ReactiveValueV0::Counter(index as u64)),
44            policy,
45        );
46        inputs.push(input);
47        projections.push((
48            format!("projection-{index}"),
49            graph.add_map(input, project_counter, policy),
50        ));
51    }
52    let fold = graph.add_delta_fold(projections, policy)?;
53    let boundary = graph.add_effect_boundary(fold, "performance-observation", policy);
54    let mut engine = graph.build()?;
55    engine.observe(boundary)?;
56    settle(&mut engine)?;
57    engine.drain_effect_receipts();
58    Ok((engine, inputs))
59}
60
61fn measure(changes_input: bool) -> Result<Vec<Duration>, Box<dyn Error>> {
62    let (mut engine, inputs) = graph()?;
63    let mut samples = Vec::with_capacity(EVENT_COUNT);
64    for event in 0..EVENT_COUNT {
65        let index = event % INPUT_COUNT;
66        let value = if changes_input {
67            (index + event + 1) as u64
68        } else {
69            index as u64
70        };
71        let started = Instant::now();
72        engine.deposit(
73            inputs[index],
74            ReactiveStateV0::available(ReactiveValueV0::Counter(value)),
75        )?;
76        settle(&mut engine)?;
77        engine.drain_effect_receipts();
78        samples.push(started.elapsed());
79    }
80    Ok(samples)
81}
82
83fn percentile(samples: &[Duration], percentile: usize) -> Duration {
84    let mut sorted = samples.to_vec();
85    sorted.sort_unstable();
86    let index = (sorted.len() - 1) * percentile / 100;
87    sorted[index]
88}
89
90fn maximum(samples: &[Duration]) -> Duration {
91    samples.iter().copied().max().unwrap_or_default()
92}
93
94fn milliseconds(duration: Duration) -> f64 {
95    duration.as_secs_f64() * 1_000.0
96}
97
98fn main() -> Result<(), Box<dyn Error>> {
99    let changed_input = measure(true)?;
100    let unchanged_input = measure(false)?;
101    let changed_p95 = percentile(&changed_input, 95);
102    let changed_max = maximum(&changed_input);
103    let unchanged_p95 = percentile(&unchanged_input, 95);
104    let unchanged_max = maximum(&unchanged_input);
105
106    println!(
107        concat!(
108            "{{\"schemaVersion\":\"omena.reactive-engine-step-envelope.v0\",",
109            "\"inputCount\":{},\"eventCount\":{},",
110            "\"p95CeilingMs\":{:.3},\"maxCeilingMs\":{:.3},",
111            "\"changedInput\":{{\"p95Ms\":{:.3},\"maxMs\":{:.3}}},",
112            "\"unchangedInput\":{{\"p95Ms\":{:.3},\"maxMs\":{:.3}}}}}"
113        ),
114        INPUT_COUNT,
115        EVENT_COUNT,
116        milliseconds(P95_CEILING),
117        milliseconds(MAX_CEILING),
118        milliseconds(changed_p95),
119        milliseconds(changed_max),
120        milliseconds(unchanged_p95),
121        milliseconds(unchanged_max),
122    );
123
124    if changed_p95 > P95_CEILING
125        || changed_max > MAX_CEILING
126        || unchanged_p95 > P95_CEILING
127        || unchanged_max > MAX_CEILING
128    {
129        return Err("reactive engine stepping exceeded its latency envelope".into());
130    }
131    Ok(())
132}