1use std::collections::{BTreeMap, HashMap, HashSet};
8use std::hash::{DefaultHasher, Hash, Hasher};
9use std::path::{Path, PathBuf};
10use std::sync::Arc;
11
12use futures::stream::{self, BoxStream};
13use futures::StreamExt;
14use harn_session_store::{
15 ListFilter, ReadRange, SessionEventKind, SessionMeta, SessionStore, StoredEvent, MAX_READ_BATCH,
16};
17use serde::{Deserialize, Serialize};
18
19use crate::agent_sessions::event_facts as facts;
20use crate::agent_sessions::event_facts::{semantic_string, semantic_value};
21use crate::event_log::{AnyEventLog, EventId, EventLog, LogError, LogEvent, Topic};
22use crate::orchestration::{load_run_record, RunRecord, RunTraceSpanRecord};
23use crate::redact::{current_policy, RedactionPolicy};
24
25pub const SESSION_TIMELINE_SCHEMA_VERSION: u32 = 2;
26pub const SESSION_TIMELINE_QUERY_METHOD: &str = "harn.session_timeline.query";
27pub const SESSION_TIMELINE_SUBSCRIBE_METHOD: &str = "harn.session_timeline.subscribe";
28pub const SESSION_TIMELINE_UNSUBSCRIBE_METHOD: &str = "harn.session_timeline.unsubscribe";
29pub const SESSION_TIMELINE_UPDATE_METHOD: &str = "harn.session_timeline.update";
30
31const DEFAULT_QUERY_LIMIT: usize = 1024;
32const READ_BATCH_SIZE: usize = 256;
33
34#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
35#[serde(default, rename_all = "camelCase")]
36pub struct SessionTimelineQuery {
37 #[serde(alias = "session_id")]
38 pub session_id: Option<String>,
39 #[serde(alias = "run_id")]
40 pub run_id: Option<String>,
41 #[serde(alias = "run_path")]
42 pub run_path: Option<String>,
43 #[serde(alias = "project_id")]
44 pub project_id: Option<String>,
45 #[serde(alias = "from_cursor")]
46 pub from_cursor: SessionTimelineCursor,
47 pub limit: Option<usize>,
48}
49
50impl SessionTimelineQuery {
51 pub fn for_session(session_id: impl Into<String>) -> Self {
52 Self {
53 session_id: Some(session_id.into()),
54 ..Self::default()
55 }
56 }
57
58 fn limit(&self) -> usize {
59 self.limit.unwrap_or(DEFAULT_QUERY_LIMIT).max(1)
60 }
61
62 fn topics(&self) -> Vec<Topic> {
63 let mut topics = Vec::new();
64 if let Some(session_id) = self.session_id.as_deref() {
65 topics.push(agent_events_topic(session_id));
66 }
67 topics.push(static_topic(crate::channels::CHANNEL_TRANSCRIPT_TOPIC));
68 topics.push(static_topic(crate::channels::CHANNEL_AUDIT_TOPIC));
69 topics
70 }
71}
72
73#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
74#[serde(default)]
75pub struct SessionTimelineCursor {
76 pub topics: BTreeMap<String, EventId>,
77}
78
79impl SessionTimelineCursor {
80 pub fn event_id_for(&self, topic: &Topic) -> Option<EventId> {
81 self.topics.get(topic.as_str()).copied()
82 }
83
84 fn bump(&mut self, topic: &str, event_id: EventId) {
85 self.topics
86 .entry(topic.to_string())
87 .and_modify(|cursor| *cursor = (*cursor).max(event_id))
88 .or_insert(event_id);
89 }
90}
91
92#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
93#[serde(rename_all = "camelCase")]
94pub struct SessionTimelineSnapshot {
95 pub schema_version: u32,
96 pub query: SessionTimelineQuery,
97 pub cursor: SessionTimelineCursor,
98 #[serde(default)]
99 pub coverage: SessionTimelineCoverage,
100 pub nodes: Vec<SessionTimelineNode>,
101}
102
103#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
108#[serde(default, rename_all = "camelCase")]
109pub struct SessionTimelineCoverage {
110 pub returned: usize,
111 pub available: Option<usize>,
112 pub truncated: bool,
113}
114
115#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
116#[serde(rename_all = "camelCase")]
117pub struct SessionTimelineUpdate {
118 pub schema_version: u32,
119 pub cursor: SessionTimelineCursor,
120 pub node: SessionTimelineNode,
121}
122
123#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
124#[serde(rename_all = "camelCase")]
125pub struct SessionTimelineNode {
126 pub id: String,
127 #[serde(skip_serializing_if = "Option::is_none")]
128 pub parent_id: Option<String>,
129 #[serde(default)]
130 pub children: Vec<String>,
131 pub category: String,
132 pub kind: String,
133 pub name: String,
134 pub status: String,
135 #[serde(skip_serializing_if = "Option::is_none")]
136 pub trace_id: Option<String>,
137 #[serde(skip_serializing_if = "Option::is_none")]
138 pub span_id: Option<String>,
139 #[serde(skip_serializing_if = "Option::is_none")]
140 pub occurred_at_ms: Option<i64>,
141 #[serde(skip_serializing_if = "Option::is_none")]
142 pub start_ms: Option<u64>,
143 #[serde(skip_serializing_if = "Option::is_none")]
144 pub duration_ms: Option<u64>,
145 #[serde(default)]
146 pub attributes: serde_json::Value,
147 #[serde(default)]
148 pub references: Vec<SessionTimelineReference>,
149 #[serde(default)]
150 pub links: Vec<SessionTimelineLink>,
151 pub order: u64,
152}
153
154#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
155#[serde(rename_all = "camelCase")]
156pub struct SessionTimelineReference {
157 pub kind: String,
158 #[serde(skip_serializing_if = "Option::is_none")]
159 pub id: Option<String>,
160 #[serde(skip_serializing_if = "Option::is_none")]
161 pub topic: Option<String>,
162 #[serde(skip_serializing_if = "Option::is_none")]
163 pub event_id: Option<EventId>,
164}
165
166#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
167#[serde(rename_all = "camelCase")]
168pub struct SessionTimelineLink {
169 pub kind: String,
170 #[serde(skip_serializing_if = "Option::is_none")]
171 pub target_id: Option<String>,
172 #[serde(skip_serializing_if = "Option::is_none")]
173 pub trace_id: Option<String>,
174 #[serde(skip_serializing_if = "Option::is_none")]
175 pub span_id: Option<String>,
176 #[serde(skip_serializing_if = "Option::is_none")]
177 pub event_id: Option<String>,
178}
179
180#[derive(Debug)]
181pub enum SessionTimelineError {
182 EventLog(LogError),
183 RunRecord(String),
184 SessionStore(String),
185}
186
187impl std::fmt::Display for SessionTimelineError {
188 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
189 match self {
190 Self::EventLog(error) => error.fmt(f),
191 Self::RunRecord(message) => f.write_str(message),
192 Self::SessionStore(message) => f.write_str(message),
193 }
194 }
195}
196
197impl std::error::Error for SessionTimelineError {}
198
199impl From<LogError> for SessionTimelineError {
200 fn from(error: LogError) -> Self {
201 Self::EventLog(error)
202 }
203}
204
205#[derive(Clone)]
206struct TimelineDraft {
207 sort_ms: i128,
208 sequence: u64,
209 node: SessionTimelineNode,
210}
211
212pub fn agent_events_topic(session_id: &str) -> Topic {
213 Topic::new(format!(
214 "observability.agent_events.{}",
215 crate::event_log::sanitize_topic_component(session_id)
216 ))
217 .expect("sanitized session id should produce a valid topic")
218}
219
220pub fn timeline_from_run_record(
221 run: &RunRecord,
222 query: SessionTimelineQuery,
223) -> SessionTimelineSnapshot {
224 let policy = current_policy();
225 let mut builder = TimelineBuilder::new(query.clone());
226 if run_matches_query(run, &query) {
227 builder.add_run_spans(run, &policy);
228 }
229 builder.finish()
230}
231
232pub async fn query_session_timeline(
233 log: Option<&AnyEventLog>,
234 run: Option<&RunRecord>,
235 query: SessionTimelineQuery,
236) -> Result<SessionTimelineSnapshot, SessionTimelineError> {
237 let policy = current_policy();
238 let mut builder = TimelineBuilder::new(query.clone());
239 if let Some(run) = run.filter(|run| run_matches_query(run, &query)) {
240 builder.add_run_spans(run, &policy);
241 } else if run.is_none() {
242 if let Some(run) = load_run_for_timeline(&query)? {
243 if run_matches_query(&run, &query) {
244 builder.add_run_spans(&run, &policy);
245 }
246 }
247 }
248 if let Some(log) = log {
249 builder.add_event_log(log, &policy).await?;
250 }
251 Ok(builder.finish())
252}
253
254pub async fn query_persisted_session_timeline(
258 project_root: &Path,
259 query: SessionTimelineQuery,
260) -> Result<Option<SessionTimelineSnapshot>, SessionTimelineError> {
261 let Some(store) = crate::stdlib::session_store::open_existing_canonical_store(project_root)
262 .map_err(|error| SessionTimelineError::SessionStore(error.to_string()))?
263 else {
264 return Ok(None);
265 };
266 query_session_store_timeline(&store, query).await
267}
268
269pub async fn query_session_store_timeline(
275 store: &dyn SessionStore,
276 query: SessionTimelineQuery,
277) -> Result<Option<SessionTimelineSnapshot>, SessionTimelineError> {
278 let Some(session_id) = query.session_id.as_deref() else {
279 return Ok(None);
280 };
281
282 let topic = canonical_session_topic(session_id);
283 let mut from = query.from_cursor.topics.get(&topic).copied();
284 let mut builder = TimelineBuilder::new(query.clone());
285 let mut saw_event = false;
286 loop {
287 let remaining = query.limit().saturating_sub(builder.nodes.len()).max(1);
288 let page = match store
289 .read(
290 session_id,
291 ReadRange {
292 from_event_id: from,
293 limit: Some(remaining.min(MAX_READ_BATCH)),
294 ..ReadRange::default()
295 },
296 )
297 .await
298 {
299 Ok(page) => page,
300 Err(harn_session_store::StoreError::NotFound(_)) => return Ok(None),
301 Err(error) => return Err(SessionTimelineError::SessionStore(error.to_string())),
302 };
303 saw_event |= !page.events.is_empty();
304 for event in page.events {
305 builder.add_stored_event(&topic, event);
306 }
307 if page.next_cursor.is_none() {
308 break;
309 }
310 if builder.nodes.len() >= query.limit() {
311 builder.exhaustive = false;
312 break;
313 }
314 from = page.next_cursor;
315 }
316 if !saw_event {
317 match store.describe(session_id).await {
318 Ok(_) => {}
319 Err(harn_session_store::StoreError::NotFound(_)) => return Ok(None),
320 Err(error) => return Err(SessionTimelineError::SessionStore(error.to_string())),
321 }
322 }
323 Ok(Some(builder.finish_in_source_order()))
324}
325
326pub async fn list_persisted_sessions(
329 project_root: &Path,
330 limit: usize,
331) -> Result<Vec<SessionMeta>, SessionTimelineError> {
332 let Some(store) = crate::stdlib::session_store::open_existing_canonical_store(project_root)
333 .map_err(|error| SessionTimelineError::SessionStore(error.to_string()))?
334 else {
335 return Ok(Vec::new());
336 };
337 store
338 .list(ListFilter {
339 project_scope: Some(project_root.to_string_lossy().into_owned()),
340 limit: Some(limit),
341 ..ListFilter::default()
342 })
343 .await
344 .map_err(|error| SessionTimelineError::SessionStore(error.to_string()))
345}
346
347pub async fn subscribe_session_timeline(
348 log: Arc<AnyEventLog>,
349 query: SessionTimelineQuery,
350) -> Result<
351 BoxStream<'static, Result<SessionTimelineUpdate, SessionTimelineError>>,
352 SessionTimelineError,
353> {
354 let policy = current_policy();
355 let mut streams = Vec::new();
356 for topic in query.topics() {
357 let topic_name = topic.as_str().to_string();
358 let from_cursor = query.from_cursor.event_id_for(&topic);
359 let events = log.clone().subscribe(&topic, from_cursor).await?;
360 let query = query.clone();
361 let policy = policy.clone();
362 streams.push(Box::pin(events.filter_map(move |item| {
363 let topic_name = topic_name.clone();
364 let query = query.clone();
365 let policy = policy.clone();
366 async move {
367 match item {
368 Ok((event_id, event)) => {
369 event_update(&query, &policy, &topic_name, event_id, event).map(Ok)
370 }
371 Err(error) => Some(Err(SessionTimelineError::EventLog(error))),
372 }
373 }
374 }))
375 as BoxStream<
376 'static,
377 Result<SessionTimelineUpdate, SessionTimelineError>,
378 >);
379 }
380 Ok(Box::pin(stream::select_all(streams)))
381}
382
383struct TimelineBuilder {
384 query: SessionTimelineQuery,
385 cursor: SessionTimelineCursor,
386 nodes: Vec<TimelineDraft>,
387 exhaustive: bool,
388 tool_positions: HashMap<u64, usize>,
389 collided_tool_positions: HashMap<String, usize>,
390}
391
392impl TimelineBuilder {
393 fn new(query: SessionTimelineQuery) -> Self {
394 let capacity = query.limit().min(10_000);
395 Self {
396 cursor: query.from_cursor.clone(),
397 query,
398 nodes: Vec::with_capacity(capacity),
399 exhaustive: true,
400 tool_positions: HashMap::with_capacity(capacity / 2),
401 collided_tool_positions: HashMap::new(),
402 }
403 }
404
405 fn push(&mut self, draft: TimelineDraft) {
406 self.nodes.push(draft);
407 }
408
409 fn register_tool_position(&mut self, hash: u64, index: usize) {
410 const COLLISION: usize = usize::MAX;
411 match self.tool_positions.get(&hash).copied() {
412 None => {
413 self.tool_positions.insert(hash, index);
414 }
415 Some(COLLISION) => {
416 self.collided_tool_positions
417 .insert(self.nodes[index].node.id.clone(), index);
418 }
419 Some(existing) => {
420 self.tool_positions.insert(hash, COLLISION);
421 self.collided_tool_positions
422 .insert(self.nodes[existing].node.id.clone(), existing);
423 self.collided_tool_positions
424 .insert(self.nodes[index].node.id.clone(), index);
425 }
426 }
427 }
428
429 fn stored_tool_position(&self, hash: u64, event: &StoredEvent) -> Option<usize> {
430 const COLLISION: usize = usize::MAX;
431 let index = self.tool_positions.get(&hash).copied()?;
432 if index == COLLISION {
433 let id = stored_tool_node_id(event)?;
434 return self.collided_tool_positions.get(&id).copied();
435 }
436 stored_tool_node_matches(&self.nodes[index].node, event).then_some(index)
437 }
438
439 fn add_run_spans(&mut self, run: &RunRecord, policy: &RedactionPolicy) {
440 for span in &run.trace_spans {
441 if !span_matches_query(span, &self.query) {
442 continue;
443 }
444 let node = span_node(span, policy);
445 self.push(TimelineDraft {
446 sort_ms: i128::from(span.start_ms),
447 sequence: span.span_id,
448 node,
449 });
450 }
451 }
452
453 fn add_stored_event(&mut self, topic: &str, event: StoredEvent) {
454 self.cursor.bump(topic, event.event_id);
455 let is_tool_result = matches!(&event.kind, SessionEventKind::ToolResult);
456 let tool_hash = stored_tool_node_hash(&event);
457 let sequence = event.event_id;
458 let event_ts_ms = event.ts_ms;
459 let sort_ms = i128::from(event_ts_ms);
460 if let Some(index) = tool_hash.and_then(|hash| self.stored_tool_position(hash, &event)) {
461 if is_tool_result {
462 merge_stored_tool_result(&mut self.nodes[index].node, event);
463 } else {
464 let mut node = stored_event_node(event);
465 merge_missing_attributes(&mut node.attributes, &self.nodes[index].node.attributes);
466 self.nodes[index].node = node;
467 }
468 return;
469 }
470 let mut node = stored_event_node(event);
471 if is_tool_result {
472 node.duration_ms = node.start_ms.and_then(|start| {
473 nonnegative_u64(event_ts_ms).map(|end| end.saturating_sub(start))
474 });
475 }
476 let index = self.nodes.len();
477 self.push(TimelineDraft {
478 sort_ms,
479 sequence,
480 node,
481 });
482 if let Some(hash) = tool_hash {
483 self.register_tool_position(hash, index);
484 }
485 }
486
487 async fn add_event_log(
488 &mut self,
489 log: &AnyEventLog,
490 policy: &RedactionPolicy,
491 ) -> Result<(), SessionTimelineError> {
492 if self.nodes.len() > self.query.limit() {
493 self.exhaustive = false;
494 return Ok(());
495 }
496 for topic in self.query.topics() {
497 let topic_name = topic.as_str().to_string();
498 let mut from = self.query.from_cursor.event_id_for(&topic);
499 loop {
500 let batch = log.read_range(&topic, from, READ_BATCH_SIZE).await?;
501 let batch_len = batch.len();
502 for (event_id, event) in batch {
503 from = Some(event_id);
504 self.cursor.bump(&topic_name, event_id);
505 if let Some(node) =
506 event_node(&self.query, policy, &topic_name, event_id, event)
507 {
508 let sort_ms = node
509 .occurred_at_ms
510 .map(i128::from)
511 .or_else(|| node.start_ms.map(i128::from))
512 .unwrap_or(i128::from(event_id));
513 self.push(TimelineDraft {
514 sort_ms,
515 sequence: event_id,
516 node,
517 });
518 if self.nodes.len() > self.query.limit() {
519 self.exhaustive = false;
520 return Ok(());
521 }
522 }
523 }
524 if batch_len < READ_BATCH_SIZE {
525 break;
526 }
527 }
528 }
529 Ok(())
530 }
531
532 fn finish(self) -> SessionTimelineSnapshot {
533 self.finish_with_ordering(true)
534 }
535
536 fn finish_in_source_order(self) -> SessionTimelineSnapshot {
537 self.finish_with_ordering(false)
538 }
539
540 fn finish_with_ordering(mut self, sort: bool) -> SessionTimelineSnapshot {
541 if sort {
542 self.nodes.sort_by(|left, right| {
543 left.sort_ms
544 .cmp(&right.sort_ms)
545 .then_with(|| left.sequence.cmp(&right.sequence))
546 .then_with(|| left.node.id.cmp(&right.node.id))
547 });
548 }
549 let available = self.exhaustive.then_some(self.nodes.len());
550 let truncated = !self.exhaustive || self.nodes.len() > self.query.limit();
551 self.nodes.truncate(self.query.limit());
552
553 let mut children_by_parent: BTreeMap<String, Vec<String>> = BTreeMap::new();
554 if self
555 .nodes
556 .iter()
557 .any(|draft| draft.node.parent_id.is_some())
558 {
559 let visible_ids: HashSet<&str> = self
560 .nodes
561 .iter()
562 .map(|draft| draft.node.id.as_str())
563 .collect();
564 for draft in &self.nodes {
565 let Some(parent_id) = draft.node.parent_id.as_ref() else {
566 continue;
567 };
568 if visible_ids.contains(parent_id.as_str()) {
569 children_by_parent
570 .entry(parent_id.clone())
571 .or_default()
572 .push(draft.node.id.clone());
573 }
574 }
575 }
576
577 let nodes: Vec<_> = self
578 .nodes
579 .into_iter()
580 .enumerate()
581 .map(|(index, mut draft)| {
582 draft.node.order = index as u64;
583 draft.node.children = children_by_parent
584 .remove(&draft.node.id)
585 .unwrap_or_default();
586 draft.node
587 })
588 .collect();
589
590 SessionTimelineSnapshot {
591 schema_version: SESSION_TIMELINE_SCHEMA_VERSION,
592 query: self.query,
593 cursor: self.cursor,
594 coverage: SessionTimelineCoverage {
595 returned: nodes.len(),
596 available,
597 truncated,
598 },
599 nodes,
600 }
601 }
602}
603
604fn canonical_session_topic(session_id: &str) -> String {
605 format!("session-store:{session_id}")
606}
607
608fn stored_tool_node_id(event: &StoredEvent) -> Option<String> {
609 if !matches!(
610 &event.kind,
611 SessionEventKind::ToolCall | SessionEventKind::ToolResult
612 ) {
613 return None;
614 }
615 let tool_call_id = event.headers.get("tool_call_id")?;
616 let run_id = event.headers.get("run_id")?;
617 let turn_id = event.headers.get("turn_id")?;
618 Some(format!(
619 "session:{}:run:{run_id}:turn:{turn_id}:tool:{tool_call_id}",
620 event.session_id
621 ))
622}
623
624fn stored_tool_node_hash(event: &StoredEvent) -> Option<u64> {
625 if !matches!(
626 &event.kind,
627 SessionEventKind::ToolCall | SessionEventKind::ToolResult
628 ) {
629 return None;
630 }
631 let mut hasher = DefaultHasher::new();
632 event.session_id.hash(&mut hasher);
633 event.headers.get("run_id")?.hash(&mut hasher);
634 event.headers.get("turn_id")?.hash(&mut hasher);
635 event.headers.get("tool_call_id")?.hash(&mut hasher);
636 Some(hasher.finish())
637}
638
639fn stored_tool_node_matches(node: &SessionTimelineNode, event: &StoredEvent) -> bool {
640 let reference_session = node
641 .references
642 .iter()
643 .find(|reference| reference.kind == "session_event")
644 .and_then(|reference| reference.id.as_deref());
645 let link_target = |kind: &str| {
646 node.links
647 .iter()
648 .find(|link| link.kind == kind)
649 .and_then(|link| link.target_id.as_deref())
650 };
651 reference_session == Some(event.session_id.as_str())
652 && link_target("run") == event.headers.get("run_id").map(String::as_str)
653 && link_target("turn") == event.headers.get("turn_id").map(String::as_str)
654 && link_target("tool_call") == event.headers.get("tool_call_id").map(String::as_str)
655}
656
657fn merge_stored_tool_result(existing: &mut SessionTimelineNode, mut event: StoredEvent) {
658 debug_assert!(matches!(&event.kind, SessionEventKind::ToolResult));
659 let source_event_id = event.headers.remove("source_event_id");
660 let message_id = event.headers.remove("message_id");
661 let tool_name = semantic_string(&event.payload, &facts::TOOL_NAME_ANY);
662 let role = semantic_string(&event.payload, &facts::ROLE);
663 let output = semantic_value(&event.payload, &facts::TOOL_OUTPUT_ANY);
664 let is_error = facts::bool_at(&event.payload, facts::TOOL_IS_ERROR);
665 let end_ms = nonnegative_u64(event.ts_ms);
666 let mut attributes = event.payload;
667 if let serde_json::Value::Object(attributes) = &mut attributes {
668 attributes.insert("revision".to_string(), event.event_id.into());
669 attributes.insert(
670 "recordHash".to_string(),
671 std::mem::take(&mut event.record_hash).into(),
672 );
673 if let Some(role) = role {
674 attributes.insert("role".to_string(), role.into());
675 }
676 if let Some(output) = output {
677 attributes.insert("output".to_string(), output);
678 }
679 attributes.insert("isError".to_string(), is_error.into());
680 }
681 let previous_attributes = std::mem::take(&mut existing.attributes);
682 merge_missing_attributes_owned(&mut attributes, previous_attributes);
683
684 if let Some(tool_name) = tool_name {
685 existing.name = tool_name;
686 }
687 let start_ms = existing
688 .start_ms
689 .or_else(|| existing.occurred_at_ms.and_then(nonnegative_u64));
690 existing.kind.clear();
691 existing.kind.push_str(event.kind.discriminator());
692 existing.status.clear();
693 existing
694 .status
695 .push_str(if is_error { "failed" } else { "completed" });
696 existing.occurred_at_ms = Some(event.ts_ms);
697 existing.start_ms = start_ms;
698 existing.duration_ms = start_ms
699 .zip(end_ms)
700 .map(|(start, end)| end.saturating_sub(start));
701 existing.attributes = attributes;
702 if let Some(reference) = existing
703 .references
704 .iter_mut()
705 .find(|reference| reference.kind == "session_event")
706 {
707 reference.event_id = Some(event.event_id);
708 }
709 let mut previous_links = std::mem::take(&mut existing.links);
710 let mut links = Vec::with_capacity(previous_links.len().max(5));
711 for kind in ["run", "turn"] {
712 if let Some(link) = take_timeline_link(&mut previous_links, kind) {
713 links.push(link);
714 }
715 }
716 links.extend(
717 [("source_event", source_event_id), ("message", message_id)]
718 .into_iter()
719 .filter_map(|(kind, target_id)| {
720 target_id.map(|target_id| SessionTimelineLink {
721 kind: kind.to_string(),
722 target_id: Some(target_id),
723 trace_id: None,
724 span_id: None,
725 event_id: None,
726 })
727 }),
728 );
729 if let Some(link) = take_timeline_link(&mut previous_links, "tool_call") {
730 links.push(link);
731 }
732 existing.links = links;
733}
734
735fn take_timeline_link(
736 links: &mut Vec<SessionTimelineLink>,
737 kind: &str,
738) -> Option<SessionTimelineLink> {
739 let index = links.iter().position(|link| link.kind == kind)?;
740 Some(links.remove(index))
741}
742
743fn stored_event_node(mut event: StoredEvent) -> SessionTimelineNode {
744 let source_event_id = event.headers.remove("source_event_id");
745 let message_id = event.headers.remove("message_id");
746 let tool_call_id = event.headers.remove("tool_call_id");
747 let run_id = event.headers.remove("run_id");
748 let turn_id = event.headers.remove("turn_id");
749 let id = tool_call_id
750 .as_ref()
751 .zip(run_id.as_ref())
752 .zip(turn_id.as_ref())
753 .filter(|_| {
754 matches!(
755 &event.kind,
756 SessionEventKind::ToolCall | SessionEventKind::ToolResult
757 )
758 })
759 .map(|((tool_call_id, run_id), turn_id)| {
760 format!(
761 "session:{}:run:{run_id}:turn:{turn_id}:tool:{tool_call_id}",
762 event.session_id
763 )
764 })
765 .or_else(|| {
766 source_event_id
767 .as_ref()
768 .map(|id| format!("session:{}:source:{id}", event.session_id))
769 })
770 .unwrap_or_else(|| format!("session:{}:event:{}", event.session_id, event.event_id));
771 let category = match &event.kind {
772 SessionEventKind::Message => "message",
773 SessionEventKind::ToolCall | SessionEventKind::ToolResult => "tool",
774 SessionEventKind::Plan => "plan",
775 SessionEventKind::Compaction => "compaction",
776 SessionEventKind::PermissionDecision => "permission",
777 SessionEventKind::Receipt => "receipt",
778 _ => "event",
779 }
780 .to_string();
781 let status = match &event.kind {
782 SessionEventKind::ToolCall => "running",
783 SessionEventKind::ToolResult => tool_result_status(&event),
784 SessionEventKind::Custom { custom_type } if custom_type == "agent_run_terminal" => event
785 .payload
786 .pointer(facts::FINAL_STATUS)
787 .and_then(serde_json::Value::as_str)
788 .unwrap_or("completed"),
789 _ => "completed",
790 }
791 .to_string();
792 let tool_name = || semantic_string(&event.payload, &facts::TOOL_NAME_ANY);
793 let visible_text = || semantic_string(&event.payload, &facts::TEXT);
794 let name = match &event.kind {
795 SessionEventKind::ToolCall => tool_name().or_else(visible_text),
796 SessionEventKind::ToolResult => tool_name(),
800 _ => visible_text(),
801 }
802 .unwrap_or_else(|| event.kind.discriminator().to_string());
803 let links = [
804 ("run", run_id),
805 ("turn", turn_id),
806 ("source_event", source_event_id),
807 ("message", message_id),
808 ("tool_call", tool_call_id),
809 ]
810 .into_iter()
811 .filter_map(|(kind, target_id)| {
812 target_id.map(|target_id| SessionTimelineLink {
813 kind: kind.to_string(),
814 target_id: Some(target_id),
815 trace_id: None,
816 span_id: None,
817 event_id: None,
818 })
819 })
820 .collect();
821 let role = semantic_string(&event.payload, &facts::ROLE);
822 let semantic_attribute = match &event.kind {
823 SessionEventKind::ToolCall => {
824 semantic_value(&event.payload, &facts::TOOL_INPUT_ANY).map(|value| ("input", value))
825 }
826 SessionEventKind::ToolResult => {
827 semantic_value(&event.payload, &facts::TOOL_OUTPUT_ANY).map(|value| ("output", value))
828 }
829 _ => None,
830 };
831 let is_error = matches!(&event.kind, SessionEventKind::ToolResult)
832 .then(|| facts::bool_at(&event.payload, facts::TOOL_IS_ERROR));
833 let mut attributes = event.payload;
834 if let serde_json::Value::Object(attributes) = &mut attributes {
835 attributes.insert("sessionId".to_string(), event.session_id.clone().into());
836 attributes.insert("revision".to_string(), event.event_id.into());
837 attributes.insert("recordHash".to_string(), event.record_hash.into());
838 if let Some(role) = role {
839 attributes.insert("role".to_string(), role.into());
840 }
841 if let Some((key, value)) = semantic_attribute {
842 attributes.insert(key.to_string(), value);
843 }
844 if let Some(is_error) = is_error {
845 attributes.insert("isError".to_string(), is_error.into());
846 }
847 }
848 let start_ms = matches!(&event.kind, SessionEventKind::ToolCall)
849 .then(|| nonnegative_u64(event.ts_ms))
850 .flatten();
851 let session_topic = canonical_session_topic(&event.session_id);
852 SessionTimelineNode {
853 id,
854 parent_id: None,
855 children: Vec::new(),
856 category,
857 kind: event.kind.discriminator().to_string(),
858 name,
859 status,
860 trace_id: None,
861 span_id: None,
862 occurred_at_ms: Some(event.ts_ms),
863 start_ms,
864 duration_ms: None,
865 attributes,
866 references: vec![SessionTimelineReference {
867 kind: "session_event".to_string(),
868 id: Some(event.session_id),
869 topic: Some(session_topic),
870 event_id: Some(event.event_id),
871 }],
872 links,
873 order: 0,
874 }
875}
876
877fn nonnegative_u64(value: i64) -> Option<u64> {
878 u64::try_from(value).ok()
879}
880
881fn merge_missing_attributes(current: &mut serde_json::Value, previous: &serde_json::Value) {
882 let (serde_json::Value::Object(current), serde_json::Value::Object(previous)) =
883 (current, previous)
884 else {
885 return;
886 };
887 for (key, value) in previous {
888 current.entry(key.clone()).or_insert_with(|| value.clone());
889 }
890}
891
892fn merge_missing_attributes_owned(current: &mut serde_json::Value, previous: serde_json::Value) {
893 let (serde_json::Value::Object(current), serde_json::Value::Object(previous)) =
894 (current, previous)
895 else {
896 return;
897 };
898 for (key, value) in previous {
899 current.entry(key).or_insert(value);
900 }
901}
902
903fn tool_result_status(event: &StoredEvent) -> &'static str {
904 if facts::bool_at(&event.payload, facts::TOOL_IS_ERROR) {
905 "failed"
906 } else {
907 "completed"
908 }
909}
910
911fn event_update(
912 query: &SessionTimelineQuery,
913 policy: &RedactionPolicy,
914 topic: &str,
915 event_id: EventId,
916 event: LogEvent,
917) -> Option<SessionTimelineUpdate> {
918 let mut node = event_node(query, policy, topic, event_id, event)?;
919 node.order = 0;
920 let mut cursor = SessionTimelineCursor::default();
921 cursor.bump(topic, event_id);
922 Some(SessionTimelineUpdate {
923 schema_version: SESSION_TIMELINE_SCHEMA_VERSION,
924 cursor,
925 node,
926 })
927}
928
929fn event_node(
930 query: &SessionTimelineQuery,
931 policy: &RedactionPolicy,
932 topic: &str,
933 event_id: EventId,
934 mut event: LogEvent,
935) -> Option<SessionTimelineNode> {
936 event.redact_in_place(policy);
937 if topic.starts_with("observability.agent_events.") {
938 return agent_event_node(query, topic, event_id, event);
939 }
940 if topic == crate::channels::CHANNEL_TRANSCRIPT_TOPIC {
941 return channel_lifecycle_node(query, topic, event_id, event);
942 }
943 if topic == crate::channels::CHANNEL_AUDIT_TOPIC {
944 return channel_audit_node(query, topic, event_id, event);
945 }
946 None
947}
948
949fn span_node(span: &RunTraceSpanRecord, policy: &RedactionPolicy) -> SessionTimelineNode {
950 let mut attributes = serde_json::json!(span.metadata);
951 policy.redact_json_in_place(&mut attributes);
952 let status = attributes
953 .get("status")
954 .and_then(serde_json::Value::as_str)
955 .unwrap_or("completed")
956 .to_string();
957 SessionTimelineNode {
958 id: span_node_id(&span.trace_id, span.span_id),
959 parent_id: span
960 .parent_id
961 .map(|parent| span_node_id(&span.trace_id, parent)),
962 children: Vec::new(),
963 category: "span".to_string(),
964 kind: span.kind.clone(),
965 name: span.name.clone(),
966 status,
967 trace_id: Some(span.trace_id.clone()),
968 span_id: Some(span.span_id.to_string()),
969 occurred_at_ms: None,
970 start_ms: Some(span.start_ms),
971 duration_ms: Some(span.duration_ms),
972 attributes,
973 references: vec![SessionTimelineReference {
974 kind: "run_trace_span".to_string(),
975 id: Some(span.span_id.to_string()),
976 topic: None,
977 event_id: None,
978 }],
979 links: span
980 .links
981 .iter()
982 .map(|link| SessionTimelineLink {
983 kind: link
984 .attributes
985 .get("harn.link.kind")
986 .cloned()
987 .unwrap_or_else(|| "span_link".to_string()),
988 target_id: Some(format!("span:{}:{}", link.trace_id, link.span_id)),
989 trace_id: Some(link.trace_id.clone()),
990 span_id: Some(link.span_id.clone()),
991 event_id: None,
992 })
993 .collect(),
994 order: 0,
995 }
996}
997
998fn agent_event_node(
999 query: &SessionTimelineQuery,
1000 topic: &str,
1001 event_id: EventId,
1002 event: LogEvent,
1003) -> Option<SessionTimelineNode> {
1004 if !event_matches_query(
1005 query,
1006 &event.payload,
1007 Some(&event.headers),
1008 &["session_id"],
1009 &[],
1010 ) {
1011 return None;
1012 }
1013 let event_value = event.payload.get("event").unwrap_or(&event.payload);
1014 let event_type = event_value
1015 .get("type")
1016 .and_then(serde_json::Value::as_str)
1017 .unwrap_or(event.kind.as_str());
1018 let status = event_status(event_value).unwrap_or("observed").to_string();
1019 Some(SessionTimelineNode {
1020 id: format!("event:{topic}:{event_id}"),
1021 parent_id: None,
1022 children: Vec::new(),
1023 category: "agent_event".to_string(),
1024 kind: event.kind.clone(),
1025 name: event_type.to_string(),
1026 status,
1027 trace_id: None,
1028 span_id: None,
1029 occurred_at_ms: Some(event.occurred_at_ms),
1030 start_ms: None,
1031 duration_ms: duration_ms(event_value),
1032 attributes: event.payload,
1033 references: vec![event_ref(topic, event_id)],
1034 links: Vec::new(),
1035 order: 0,
1036 })
1037}
1038
1039fn channel_lifecycle_node(
1040 query: &SessionTimelineQuery,
1041 topic: &str,
1042 event_id: EventId,
1043 event: LogEvent,
1044) -> Option<SessionTimelineNode> {
1045 if !event_matches_query(
1046 query,
1047 &event.payload,
1048 Some(&event.headers),
1049 &["session_id", "matched_in_session_id"],
1050 &["pipeline_id"],
1051 ) {
1052 return None;
1053 }
1054 let channel_event_id = string_field(&event.payload, "event_id");
1055 let trigger_id = string_field(&event.payload, "trigger_id");
1056 let is_match = event.kind == crate::channels::CHANNEL_MATCH_TRANSCRIPT_KIND;
1057 let id = if is_match {
1058 format!(
1059 "channel:{}:match:{}",
1060 channel_event_id.as_deref().unwrap_or("unknown"),
1061 trigger_id.as_deref().unwrap_or("unknown")
1062 )
1063 } else {
1064 format!(
1065 "channel:{}:emit",
1066 channel_event_id.as_deref().unwrap_or("unknown")
1067 )
1068 };
1069 let mut links: Vec<SessionTimelineLink> = if is_match {
1070 channel_event_id
1071 .as_ref()
1072 .map(|event_id| SessionTimelineLink {
1073 kind: "channel_emit".to_string(),
1074 target_id: Some(format!("channel:{event_id}:emit")),
1075 trace_id: None,
1076 span_id: None,
1077 event_id: Some(event_id.clone()),
1078 })
1079 .into_iter()
1080 .collect()
1081 } else {
1082 Vec::new()
1083 };
1084 if is_match {
1085 links.extend(channel_batch_links(&event.payload));
1086 }
1087 Some(SessionTimelineNode {
1088 id,
1089 parent_id: None,
1090 children: Vec::new(),
1091 category: "channel".to_string(),
1092 kind: event.kind.clone(),
1093 name: string_field(&event.payload, "name_resolved")
1094 .or_else(|| string_field(&event.payload, "name"))
1095 .unwrap_or_else(|| event.kind.clone()),
1096 status: if event
1097 .payload
1098 .get("duplicate")
1099 .and_then(serde_json::Value::as_bool)
1100 .unwrap_or(false)
1101 {
1102 "duplicate".to_string()
1103 } else {
1104 "observed".to_string()
1105 },
1106 trace_id: None,
1107 span_id: string_field(&event.payload, "span_id"),
1108 occurred_at_ms: Some(event.occurred_at_ms),
1109 start_ms: None,
1110 duration_ms: None,
1111 attributes: event.payload,
1112 references: vec![event_ref(topic, event_id)],
1113 links,
1114 order: 0,
1115 })
1116}
1117
1118fn channel_audit_node(
1119 query: &SessionTimelineQuery,
1120 topic: &str,
1121 event_id: EventId,
1122 event: LogEvent,
1123) -> Option<SessionTimelineNode> {
1124 if !event_matches_query(
1125 query,
1126 &event.payload,
1127 Some(&event.headers),
1128 &["session_id", "matched_in_session_id"],
1129 &["pipeline_id", "run_id"],
1130 ) {
1131 return None;
1132 }
1133 let channel_event_id = string_field(&event.payload, "event_id");
1134 let trigger_id = string_field(&event.payload, "trigger_id");
1135 let is_match = event.kind == crate::channels::CHANNEL_MATCH_RECEIPT_KIND;
1136 let id = if is_match {
1137 format!(
1138 "channel_receipt:{}:match:{}",
1139 channel_event_id.as_deref().unwrap_or("unknown"),
1140 trigger_id.as_deref().unwrap_or("unknown")
1141 )
1142 } else {
1143 format!(
1144 "channel_receipt:{}:emit",
1145 channel_event_id.as_deref().unwrap_or("unknown")
1146 )
1147 };
1148 let mut links: Vec<SessionTimelineLink> = if is_match {
1149 channel_event_id
1150 .as_ref()
1151 .map(|event_id| SessionTimelineLink {
1152 kind: "channel_emit".to_string(),
1153 target_id: Some(format!("channel_receipt:{event_id}:emit")),
1154 trace_id: None,
1155 span_id: None,
1156 event_id: Some(event_id.clone()),
1157 })
1158 .into_iter()
1159 .collect()
1160 } else {
1161 Vec::new()
1162 };
1163 if is_match {
1164 links.extend(channel_batch_links(&event.payload));
1165 }
1166 Some(SessionTimelineNode {
1167 id,
1168 parent_id: None,
1169 children: Vec::new(),
1170 category: "channel_audit".to_string(),
1171 kind: event.kind.clone(),
1172 name: string_field(&event.payload, "name_resolved").unwrap_or_else(|| event.kind.clone()),
1173 status: event
1174 .payload
1175 .get("handler_result")
1176 .and_then(|value| value.get("status"))
1177 .and_then(serde_json::Value::as_str)
1178 .or_else(|| {
1179 event.payload.get("inserted").and_then(|inserted| {
1180 if inserted.as_bool() == Some(false) {
1181 Some("duplicate")
1182 } else {
1183 None
1184 }
1185 })
1186 })
1187 .unwrap_or("recorded")
1188 .to_string(),
1189 trace_id: None,
1190 span_id: string_field(&event.payload, "span_id"),
1191 occurred_at_ms: Some(event.occurred_at_ms),
1192 start_ms: None,
1193 duration_ms: None,
1194 attributes: event.payload,
1195 references: vec![event_ref(topic, event_id)],
1196 links,
1197 order: 0,
1198 })
1199}
1200
1201fn event_matches_query(
1202 query: &SessionTimelineQuery,
1203 payload: &serde_json::Value,
1204 headers: Option<&BTreeMap<String, String>>,
1205 session_keys: &[&str],
1206 run_keys: &[&str],
1207) -> bool {
1208 field_query_matches(query.session_id.as_deref(), payload, headers, session_keys)
1209 && field_query_matches(query.run_id.as_deref(), payload, headers, run_keys)
1210 && field_query_matches(
1211 query.project_id.as_deref(),
1212 payload,
1213 headers,
1214 &["project_id", "projectId", "workspace_id", "workspaceId"],
1215 )
1216}
1217
1218fn field_query_matches(
1219 expected: Option<&str>,
1220 payload: &serde_json::Value,
1221 headers: Option<&BTreeMap<String, String>>,
1222 keys: &[&str],
1223) -> bool {
1224 let Some(expected) = expected else {
1225 return true;
1226 };
1227 if expected.is_empty() {
1228 return true;
1229 }
1230 if keys.is_empty() {
1231 return true;
1232 }
1233 keys.iter().any(|key| {
1234 payload
1235 .get(*key)
1236 .and_then(serde_json::Value::as_str)
1237 .is_some_and(|value| value == expected)
1238 || payload
1239 .get("event")
1240 .and_then(|event| event.get(*key))
1241 .and_then(serde_json::Value::as_str)
1242 .is_some_and(|value| value == expected)
1243 || headers
1244 .and_then(|headers| headers.get(*key))
1245 .is_some_and(|value| value == expected)
1246 })
1247}
1248
1249fn span_matches_query(span: &RunTraceSpanRecord, query: &SessionTimelineQuery) -> bool {
1250 if let Some(session_id) = query.session_id.as_deref() {
1251 let has_session_attr = span.metadata.contains_key("session_id")
1252 || span.metadata.contains_key("agent_session_id");
1253 if has_session_attr
1254 && !metadata_matches(
1255 &span.metadata,
1256 &["session_id", "agent_session_id"],
1257 session_id,
1258 )
1259 {
1260 return false;
1261 }
1262 }
1263 true
1264}
1265
1266fn run_matches_query(run: &RunRecord, query: &SessionTimelineQuery) -> bool {
1267 if let Some(run_id) = query.run_id.as_deref() {
1268 if run.id != run_id {
1269 return false;
1270 }
1271 }
1272 if let Some(project_id) = query.project_id.as_deref() {
1273 if !metadata_matches(&run.metadata, &["project_id", "projectId"], project_id) {
1274 return false;
1275 }
1276 }
1277 true
1278}
1279
1280fn metadata_matches(
1281 metadata: &BTreeMap<String, serde_json::Value>,
1282 keys: &[&str],
1283 expected: &str,
1284) -> bool {
1285 keys.iter().any(|key| {
1286 metadata
1287 .get(*key)
1288 .and_then(serde_json::Value::as_str)
1289 .is_some_and(|value| value == expected)
1290 })
1291}
1292
1293fn event_status(value: &serde_json::Value) -> Option<&str> {
1294 value
1295 .get("status")
1296 .and_then(serde_json::Value::as_str)
1297 .or_else(|| value.get("verdict").and_then(serde_json::Value::as_str))
1298}
1299
1300fn duration_ms(value: &serde_json::Value) -> Option<u64> {
1301 value
1302 .get("duration_ms")
1303 .or_else(|| value.get("judge_duration_ms"))
1304 .and_then(serde_json::Value::as_u64)
1305}
1306
1307fn string_field(value: &serde_json::Value, key: &str) -> Option<String> {
1308 let value = value.get(key)?;
1309 if let Some(text) = value.as_str() {
1310 if !text.is_empty() {
1311 return Some(text.to_string());
1312 }
1313 return None;
1314 }
1315 value.as_u64().map(|number| number.to_string())
1316}
1317
1318fn channel_batch_links(payload: &serde_json::Value) -> Vec<SessionTimelineLink> {
1319 payload
1320 .get("batch")
1321 .and_then(|batch| batch.get("constituent_event_ids"))
1322 .and_then(serde_json::Value::as_array)
1323 .into_iter()
1324 .flatten()
1325 .filter_map(|value| value.as_str())
1326 .map(|event_id| SessionTimelineLink {
1327 kind: "channel_batch_member".to_string(),
1328 target_id: None,
1329 trace_id: None,
1330 span_id: None,
1331 event_id: Some(event_id.to_string()),
1332 })
1333 .collect()
1334}
1335
1336fn span_node_id(trace_id: &str, span_id: u64) -> String {
1337 format!("span:{trace_id}:{span_id}")
1338}
1339
1340fn event_ref(topic: &str, event_id: EventId) -> SessionTimelineReference {
1341 SessionTimelineReference {
1342 kind: "event_log".to_string(),
1343 id: None,
1344 topic: Some(topic.to_string()),
1345 event_id: Some(event_id),
1346 }
1347}
1348
1349fn static_topic(topic: &str) -> Topic {
1350 Topic::new(topic).expect("static session timeline topic should be valid")
1351}
1352
1353fn load_run_for_timeline(
1354 query: &SessionTimelineQuery,
1355) -> Result<Option<RunRecord>, SessionTimelineError> {
1356 if let Some(path) = query
1357 .run_path
1358 .as_deref()
1359 .map(str::trim)
1360 .filter(|path| !path.is_empty())
1361 {
1362 return load_run_record_for_timeline(Path::new(path), true);
1363 }
1364
1365 let Some(run_id) = query
1366 .run_id
1367 .as_deref()
1368 .map(str::trim)
1369 .filter(|run_id| !run_id.is_empty())
1370 else {
1371 return Ok(None);
1372 };
1373 let path = default_run_record_path(run_id)?;
1374 load_run_record_for_timeline(&path, false)
1375}
1376
1377fn load_run_record_for_timeline(
1378 path: &Path,
1379 explicit: bool,
1380) -> Result<Option<RunRecord>, SessionTimelineError> {
1381 if !path.exists() {
1382 if explicit {
1383 return Err(SessionTimelineError::RunRecord(format!(
1384 "session timeline run record not found: {}",
1385 path.display()
1386 )));
1387 }
1388 return Ok(None);
1389 }
1390 load_run_record(path).map(Some).map_err(|error| {
1391 SessionTimelineError::RunRecord(format!(
1392 "failed to load session timeline run record {}: {error}",
1393 path.display()
1394 ))
1395 })
1396}
1397
1398fn default_run_record_path(run_id: &str) -> Result<PathBuf, SessionTimelineError> {
1399 if run_id == "." || run_id == ".." || run_id.contains('/') || run_id.contains('\\') {
1400 return Err(SessionTimelineError::RunRecord(format!(
1401 "session timeline runId is not a valid default run-record filename: {run_id}"
1402 )));
1403 }
1404 let base = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
1405 Ok(crate::runtime_paths::run_root(&base).join(format!("{run_id}.json")))
1406}
1407
1408#[cfg(test)]
1409#[path = "session_timeline_tests.rs"]
1410mod session_timeline_tests;