Skip to main content

starweaver_cli/local_store/
replay.rs

1//! Local display replay windows over shared stream cursor contracts.
2
3use starweaver_stream::{
4    DisplayMessage, ReplayCursor, ReplayCursorFamily, ReplayEvent, ReplayEventKind, ReplayScope,
5};
6
7use super::LocalStore;
8use crate::{CliError, CliResult};
9
10/// Display replay events and the next sequence for live tail continuation.
11#[derive(Clone, Debug)]
12pub struct DisplayReplayWindow {
13    /// Replay scope used by the events.
14    pub scope: ReplayScope,
15    /// Replay events after the requested cursor.
16    pub events: Vec<ReplayEvent>,
17    /// Next sequence number for live events in this scope.
18    pub next_sequence: usize,
19}
20
21impl LocalStore {
22    /// Replay display messages as scoped replay events.
23    pub fn replay_display_window(
24        &self,
25        session_id: &str,
26        run_id: Option<&str>,
27        cursor: Option<&ReplayCursor>,
28    ) -> CliResult<DisplayReplayWindow> {
29        let scope = run_id.map_or_else(|| ReplayScope::session(session_id), ReplayScope::run);
30        if let Some(cursor) = cursor {
31            cursor
32                .validate(ReplayCursorFamily::Display, &scope)
33                .map_err(|error| CliError::Usage(error.to_string()))?;
34        }
35        let messages = self.replay_display(session_id, run_id, None)?;
36        let (events, next_sequence) = if run_id.is_some() {
37            let next_sequence = messages
38                .last()
39                .map_or(0, |message| message.sequence.saturating_add(1));
40            (
41                messages
42                    .into_iter()
43                    .map(|message| display_replay_event(&scope, message.sequence, message))
44                    .collect::<Vec<_>>(),
45                next_sequence,
46            )
47        } else {
48            let next_sequence = messages.len();
49            (
50                messages
51                    .into_iter()
52                    .enumerate()
53                    .map(|(sequence, message)| display_replay_event(&scope, sequence, message))
54                    .collect::<Vec<_>>(),
55                next_sequence,
56            )
57        };
58        Ok(DisplayReplayWindow {
59            scope,
60            events: filter_replay_events(events, cursor),
61            next_sequence,
62        })
63    }
64}
65
66fn display_replay_event(
67    scope: &ReplayScope,
68    sequence: usize,
69    message: DisplayMessage,
70) -> ReplayEvent {
71    ReplayEvent::new(
72        scope.clone(),
73        sequence,
74        ReplayEventKind::DisplayMessage(Box::new(message)),
75    )
76}
77
78fn filter_replay_events(
79    events: Vec<ReplayEvent>,
80    cursor: Option<&ReplayCursor>,
81) -> Vec<ReplayEvent> {
82    events
83        .into_iter()
84        .filter(|event| cursor.is_none_or(|cursor| event.sequence > cursor.sequence))
85        .collect()
86}