Skip to main content

harn_vm/
session_timeline.rs

1//! Session timeline projection for client-facing observability.
2//!
3//! This module does not own persistence. It projects the existing run record
4//! spans and event-log topics into one stable, redacted shape that clients can
5//! query or subscribe to without learning Harn's storage internals.
6
7use 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/// States whether a bounded snapshot covers every matching semantic node.
104///
105/// `available` is exact when Harn exhausted every selected source. It is `None`
106/// when Harn stopped after proving that the requested limit omitted evidence.
107#[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
254/// Project durable canonical transcript rows into the same timeline contract
255/// used by live observability. Returns `None` when the project has no canonical
256/// store or the requested session does not exist.
257pub 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
269/// Project one [`SessionStore`] through Harn's canonical semantic timeline.
270///
271/// Transport adapters use this seam so HTTP, ACP, and embedded callers receive
272/// identical node identity, grouping, status, references, links, and redaction
273/// without learning how stored event envelopes map to product chronology.
274pub 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
326/// List canonical sessions for one project without exposing SQLite layout or
327/// schema policy to product clients.
328pub 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        // Result text is output, not the tool's product label. The
797        // discriminator sentinel makes the revision merge retain the call
798        // node's name when this envelope omits tool_name.
799        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;