1use std::collections::{BTreeMap, BTreeSet};
4
5use serde::{Deserialize, Serialize};
6
7use crate::event::ObservabilityEvent;
8
9#[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
13pub struct EventFilter {
14 #[serde(skip_serializing_if = "Option::is_none")]
16 pub conversation_id: Option<String>,
17 #[serde(skip_serializing_if = "Option::is_none")]
19 pub kind: Option<String>,
20 #[serde(skip_serializing_if = "Option::is_none")]
22 pub min_tick: Option<u64>,
23 #[serde(skip_serializing_if = "Option::is_none")]
25 pub max_tick: Option<u64>,
26 #[serde(skip_serializing_if = "Option::is_none")]
28 pub kernel_id: Option<String>,
29 #[serde(skip_serializing_if = "Option::is_none")]
31 pub tool_name: Option<String>,
32 #[serde(skip_serializing_if = "Option::is_none")]
34 pub call_id: Option<String>,
35 #[serde(skip_serializing_if = "Option::is_none")]
37 pub skill_id: Option<String>,
38 #[serde(skip_serializing_if = "Option::is_none")]
40 pub model: Option<String>,
41}
42
43impl EventFilter {
44 pub fn new() -> Self {
46 Self::default()
47 }
48
49 pub fn conversation_id(mut self, conversation_id: impl Into<String>) -> Self {
51 self.conversation_id = Some(conversation_id.into());
52 self
53 }
54
55 pub fn kind(mut self, kind: impl Into<String>) -> Self {
57 self.kind = Some(kind.into());
58 self
59 }
60
61 pub fn min_tick(mut self, tick: u64) -> Self {
63 self.min_tick = Some(tick);
64 self
65 }
66
67 pub fn max_tick(mut self, tick: u64) -> Self {
69 self.max_tick = Some(tick);
70 self
71 }
72
73 pub fn kernel_id(mut self, kernel_id: impl Into<String>) -> Self {
75 self.kernel_id = Some(kernel_id.into());
76 self
77 }
78
79 pub fn tool_name(mut self, tool_name: impl Into<String>) -> Self {
81 self.tool_name = Some(tool_name.into());
82 self
83 }
84
85 pub fn call_id(mut self, call_id: impl Into<String>) -> Self {
87 self.call_id = Some(call_id.into());
88 self
89 }
90
91 pub fn skill_id(mut self, skill_id: impl Into<String>) -> Self {
93 self.skill_id = Some(skill_id.into());
94 self
95 }
96
97 pub fn model(mut self, model: impl Into<String>) -> Self {
99 self.model = Some(model.into());
100 self
101 }
102
103 pub fn matches(&self, event: &ObservabilityEvent) -> bool {
105 if self
106 .conversation_id
107 .as_ref()
108 .is_some_and(|expected| expected != &event.conversation_id)
109 {
110 return false;
111 }
112 if self
113 .kind
114 .as_ref()
115 .is_some_and(|expected| expected != event.kind.discriminant())
116 {
117 return false;
118 }
119 if self.min_tick.is_some_and(|min_tick| event.tick < min_tick) {
120 return false;
121 }
122 if self.max_tick.is_some_and(|max_tick| event.tick > max_tick) {
123 return false;
124 }
125
126 let fields = event.kind.scalar_fields();
127 if self
128 .kernel_id
129 .as_ref()
130 .is_some_and(|expected| expected != fields.kernel_id)
131 {
132 return false;
133 }
134 if self
135 .tool_name
136 .as_ref()
137 .is_some_and(|expected| expected != fields.tool_name)
138 {
139 return false;
140 }
141 if self
142 .call_id
143 .as_ref()
144 .is_some_and(|expected| expected != fields.call_id)
145 {
146 return false;
147 }
148 if self
149 .skill_id
150 .as_ref()
151 .is_some_and(|expected| expected != fields.skill_id)
152 {
153 return false;
154 }
155 if self
156 .model
157 .as_ref()
158 .is_some_and(|expected| expected != fields.model)
159 {
160 return false;
161 }
162
163 true
164 }
165}
166
167#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)]
173pub struct EventQuery {
174 events: Vec<ObservabilityEvent>,
175}
176
177impl EventQuery {
178 pub fn new(events: Vec<ObservabilityEvent>) -> Self {
180 Self { events }
181 }
182
183 pub fn all(&self) -> &[ObservabilityEvent] {
185 &self.events
186 }
187
188 pub fn len(&self) -> usize {
190 self.events.len()
191 }
192
193 pub fn is_empty(&self) -> bool {
195 self.events.is_empty()
196 }
197
198 pub fn filter(&self, filter: &EventFilter) -> Vec<ObservabilityEvent> {
200 self.events
201 .iter()
202 .filter(|event| filter.matches(event))
203 .cloned()
204 .collect()
205 }
206
207 pub fn latest(&self, limit: usize) -> Vec<ObservabilityEvent> {
209 let mut events = self
210 .events
211 .iter()
212 .rev()
213 .take(limit)
214 .cloned()
215 .collect::<Vec<_>>();
216 events.reverse();
217 events
218 }
219
220 pub fn count_by_kind(&self) -> BTreeMap<String, usize> {
222 let mut counts = BTreeMap::new();
223 for event in &self.events {
224 let count = counts
225 .entry(event.kind.discriminant().to_string())
226 .or_insert(0);
227 *count += 1;
228 }
229 counts
230 }
231
232 pub fn conversations(&self) -> Vec<String> {
234 self.events
235 .iter()
236 .map(|event| event.conversation_id.clone())
237 .collect::<BTreeSet<_>>()
238 .into_iter()
239 .collect()
240 }
241}
242
243impl From<Vec<ObservabilityEvent>> for EventQuery {
244 fn from(events: Vec<ObservabilityEvent>) -> Self {
245 Self::new(events)
246 }
247}
248
249#[cfg(test)]
250#[allow(
251 clippy::unwrap_used,
252 clippy::panic,
253 clippy::indexing_slicing,
254 clippy::expect_used
255)]
256mod tests {
257 use super::*;
258 use crate::event::{EventKind, SCHEMA_VERSION};
259
260 fn event(tick: u64, conversation_id: &str, kind: EventKind) -> ObservabilityEvent {
261 ObservabilityEvent {
262 version: SCHEMA_VERSION,
263 occurred_at_millis: 1_715_000_000_000 + tick,
264 tick,
265 conversation_id: conversation_id.into(),
266 kind,
267 }
268 }
269
270 #[test]
271 fn filter_matches_conversation_kind_and_tick_window() {
272 let query = EventQuery::new(vec![
273 event(
274 1,
275 "a",
276 EventKind::PromptStarted {
277 model: "m".into(),
278 messages_in: 1,
279 },
280 ),
281 event(
282 2,
283 "a",
284 EventKind::ToolCompleted {
285 tool_name: "search".into(),
286 provider_call_id: None,
287 call_id: "call-1".into(),
288 result: "ok".into(),
289 truncated: false,
290 },
291 ),
292 event(
293 3,
294 "b",
295 EventKind::ToolCompleted {
296 tool_name: "search".into(),
297 provider_call_id: None,
298 call_id: "call-2".into(),
299 result: "ok".into(),
300 truncated: false,
301 },
302 ),
303 ]);
304
305 let matches = query.filter(
306 &EventFilter::new()
307 .conversation_id("a")
308 .kind("tool.completed")
309 .min_tick(2)
310 .max_tick(3),
311 );
312
313 assert_eq!(matches.len(), 1);
314 assert_eq!(matches[0].tick, 2);
315 }
316
317 #[test]
318 fn filter_matches_scalar_fields() {
319 let query = EventQuery::new(vec![event(
320 1,
321 "thread",
322 EventKind::ComposeSkillResolved {
323 kernel_id: "kernel".into(),
324 skill_id: "retrieval".into(),
325 applies: true,
326 delta: Some(0.2),
327 confidence: Some(0.8),
328 },
329 )]);
330
331 let matches = query.filter(&EventFilter::new().kernel_id("kernel").skill_id("retrieval"));
332
333 assert_eq!(matches.len(), 1);
334 assert!(
335 query
336 .filter(&EventFilter::new().tool_name("search"))
337 .is_empty()
338 );
339 }
340
341 #[test]
342 fn query_summarizes_conversations_and_kinds() {
343 let query = EventQuery::new(vec![
344 event(
345 1,
346 "b",
347 EventKind::PromptStarted {
348 model: "m".into(),
349 messages_in: 1,
350 },
351 ),
352 event(
353 2,
354 "a",
355 EventKind::PromptStarted {
356 model: "m".into(),
357 messages_in: 2,
358 },
359 ),
360 event(
361 3,
362 "b",
363 EventKind::ContextSampled {
364 message_count: 1,
365 byte_size: 2,
366 token_estimate: None,
367 },
368 ),
369 ]);
370
371 assert_eq!(query.conversations(), vec!["a", "b"]);
372 assert_eq!(query.count_by_kind().get("prompt.started"), Some(&2));
373 assert_eq!(
374 query
375 .latest(2)
376 .iter()
377 .map(|event| event.tick)
378 .collect::<Vec<_>>(),
379 vec![2, 3]
380 );
381 }
382}