performance_envelope/
performance_envelope.rs1use 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}