1use std::collections::BTreeMap;
7
8use serde::{Deserialize, Serialize};
9
10use crate::session::SessionId;
11use crate::vm::ObsEvent;
12
13fn default_trace_schema_version() -> u32 {
14 1
15}
16
17pub const TRACE_NORMALIZATION_SCHEMA_VERSION: u32 = 1;
19
20#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
22pub struct NormalizedTraceV1 {
23 #[serde(default = "default_trace_schema_version")]
25 pub schema_version: u32,
26 pub events: Vec<ObsEvent>,
28}
29
30#[must_use]
32pub fn obs_session(ev: &ObsEvent) -> Option<SessionId> {
33 match ev {
34 ObsEvent::Sent { session, .. }
35 | ObsEvent::Received { session, .. }
36 | ObsEvent::Opened { session, .. }
37 | ObsEvent::Closed { session, .. }
38 | ObsEvent::Acquired { session, .. }
39 | ObsEvent::Released { session, .. }
40 | ObsEvent::Transferred { session, .. }
41 | ObsEvent::Forked { session, .. }
42 | ObsEvent::Joined { session, .. }
43 | ObsEvent::Aborted { session, .. }
44 | ObsEvent::Tagged { session, .. }
45 | ObsEvent::Checked { session, .. } => Some(*session),
46 ObsEvent::Offered { edge, .. } | ObsEvent::Chose { edge, .. } => Some(edge.sid),
47 ObsEvent::EpochAdvanced { sid, .. } => Some(*sid),
48 ObsEvent::Halted { .. }
49 | ObsEvent::Invoked { .. }
50 | ObsEvent::Faulted { .. }
51 | ObsEvent::OutputConditionChecked { .. } => None,
52 }
53}
54
55#[must_use]
57pub fn with_tick(ev: &ObsEvent, tick: u64) -> ObsEvent {
58 let mut out = ev.clone();
59 set_obs_event_tick(&mut out, tick);
60 out
61}
62
63#[allow(clippy::too_many_lines)]
64fn set_obs_event_tick(out: &mut ObsEvent, tick: u64) {
65 match out {
66 ObsEvent::Sent {
67 tick: event_tick, ..
68 }
69 | ObsEvent::Received {
70 tick: event_tick, ..
71 }
72 | ObsEvent::Offered {
73 tick: event_tick, ..
74 }
75 | ObsEvent::Chose {
76 tick: event_tick, ..
77 }
78 | ObsEvent::Opened {
79 tick: event_tick, ..
80 }
81 | ObsEvent::Closed {
82 tick: event_tick, ..
83 }
84 | ObsEvent::EpochAdvanced {
85 tick: event_tick, ..
86 }
87 | ObsEvent::Halted {
88 tick: event_tick, ..
89 }
90 | ObsEvent::Invoked {
91 tick: event_tick, ..
92 }
93 | ObsEvent::Acquired {
94 tick: event_tick, ..
95 }
96 | ObsEvent::Released {
97 tick: event_tick, ..
98 }
99 | ObsEvent::Transferred {
100 tick: event_tick, ..
101 }
102 | ObsEvent::Forked {
103 tick: event_tick, ..
104 }
105 | ObsEvent::Joined {
106 tick: event_tick, ..
107 }
108 | ObsEvent::Aborted {
109 tick: event_tick, ..
110 }
111 | ObsEvent::Tagged {
112 tick: event_tick, ..
113 }
114 | ObsEvent::Checked {
115 tick: event_tick, ..
116 }
117 | ObsEvent::Faulted {
118 tick: event_tick, ..
119 }
120 | ObsEvent::OutputConditionChecked {
121 tick: event_tick, ..
122 } => *event_tick = tick,
123 }
124}
125
126#[must_use]
128pub fn normalize_trace(trace: &[ObsEvent]) -> Vec<ObsEvent> {
129 let mut counters: BTreeMap<SessionId, u64> = BTreeMap::new();
130 let mut out = Vec::with_capacity(trace.len());
131 for ev in trace {
132 if let Some(session) = obs_session(ev) {
133 let counter = counters.entry(session).or_insert(0);
134 let local_tick = *counter;
135 *counter += 1;
136 out.push(with_tick(ev, local_tick));
137 } else {
138 out.push(ev.clone());
139 }
140 }
141 out
142}
143
144#[must_use]
146pub fn strict_trace(trace: &[ObsEvent]) -> Vec<ObsEvent> {
147 trace.to_vec()
148}
149
150#[must_use]
152pub fn normalize_trace_v1(trace: &[ObsEvent]) -> NormalizedTraceV1 {
153 NormalizedTraceV1 {
154 schema_version: TRACE_NORMALIZATION_SCHEMA_VERSION,
155 events: normalize_trace(trace),
156 }
157}
158
159#[cfg(test)]
160mod tests {
161 use super::*;
162 use crate::session::Edge;
163
164 #[test]
165 #[allow(clippy::as_conversions)]
166 fn normalize_trace_memory_is_bounded_at_10k_events() {
167 let mut trace = Vec::with_capacity(10_000);
168 for i in 0..10_000usize {
169 let sid = i % 32;
170 let tick = i as u64;
171 if i % 2 == 0 {
172 trace.push(ObsEvent::Sent {
173 tick,
174 edge: Edge::new(sid, "A", "B"),
175 session: sid,
176 from: "A".to_string(),
177 to: "B".to_string(),
178 label: "m".to_string(),
179 });
180 } else {
181 trace.push(ObsEvent::Received {
182 tick,
183 edge: Edge::new(sid, "B", "A"),
184 session: sid,
185 from: "B".to_string(),
186 to: "A".to_string(),
187 label: "m".to_string(),
188 });
189 }
190 }
191
192 let normalized = normalize_trace(&trace);
193 assert_eq!(normalized.len(), trace.len());
194 assert!(
195 normalized.capacity() <= trace.len() + 1,
196 "normalize_trace should allocate O(n) space without capacity blow-up"
197 );
198 }
199}