Skip to main content

starweaver_cli/local_store/
replay.rs

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