subc_daemon/
terminal_ring.rs1use std::{collections::VecDeque, sync::Arc};
2
3use crate::terminal_journal::{ring_history, JournalRead, TerminalJournal};
4
5use subc_control::{TerminalDisposition, TerminalExitKind};
6
7const DEFAULT_MAX_ENTRIES: usize = 32;
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub struct TerminalRingConfig {
12 max_entries: usize,
13}
14
15impl TerminalRingConfig {
16 pub const fn new(max_entries: usize) -> Self {
19 Self {
20 max_entries: if max_entries == 0 { 1 } else { max_entries },
21 }
22 }
23}
24
25impl Default for TerminalRingConfig {
26 fn default() -> Self {
27 Self::new(DEFAULT_MAX_ENTRIES)
28 }
29}
30
31#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
33pub struct TerminalRecord {
34 pub exit_code: Option<i32>,
35 pub exit_signal: Option<i32>,
36 pub at_ms: u64,
37 pub disposition: TerminalDisposition,
38 pub exit_kind: TerminalExitKind,
39 pub disposition_detail: Option<String>,
45}
46
47#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct TerminalHistorySnapshot {
50 pub daemon_started_at_ms: u64,
51 pub entries: Vec<TerminalRecord>,
52 pub dropped: u64,
53}
54
55pub(crate) enum DurableHistoryRead {
57 RingOnly(TerminalHistorySnapshot),
58 Journal(JournalRead),
59}
60
61impl DurableHistoryRead {
62 pub(crate) fn read(self, module_id: &str) -> subc_control::TerminalHistory {
64 match self {
65 Self::RingOnly(snapshot) => ring_history(snapshot, None),
66 Self::Journal(read) => read.read(module_id),
67 }
68 }
69}
70
71#[derive(Debug)]
76pub struct TerminalRing {
77 journal: Option<Arc<TerminalJournal>>,
78 config: TerminalRingConfig,
79 daemon_started_at_ms: u64,
80 start_clock: Option<crate::clock::StartClock>,
81 entries: VecDeque<TerminalRecord>,
82 dropped: u64,
83}
84
85impl TerminalRing {
86 pub fn new(config: TerminalRingConfig, daemon_started_at_ms: u64) -> Self {
87 Self {
88 journal: None,
89 config,
90 daemon_started_at_ms,
91 start_clock: None,
92 entries: VecDeque::new(),
93 dropped: 0,
94 }
95 }
96
97 pub(crate) fn with_start_clock(mut self, clock: crate::clock::StartClock) -> Self {
98 self.start_clock = Some(clock);
99 self
100 }
101
102 pub(crate) fn with_journal(mut self, journal: Option<Arc<TerminalJournal>>) -> Self {
103 self.journal = journal;
104 self
105 }
106
107 pub(crate) fn append_journal(&self, module_id: &str, entry: &TerminalRecord) {
108 if let Some(journal) = &self.journal {
109 journal.append(module_id, entry);
110 }
111 }
112
113 #[cfg(test)]
115 pub(crate) fn durable_history(&self, module_id: &str) -> subc_control::TerminalHistory {
116 self.capture_durable_history().read(module_id)
117 }
118
119 pub(crate) fn capture_durable_history(&self) -> DurableHistoryRead {
123 match &self.journal {
124 Some(journal) => DurableHistoryRead::Journal(journal.capture_read(self.snapshot())),
125 None => DurableHistoryRead::RingOnly(self.snapshot()),
126 }
127 }
128
129 pub fn push(&mut self, entry: TerminalRecord) {
130 self.entries.push_back(entry);
131 while self.entries.len() > self.config.max_entries {
132 self.entries.pop_front();
133 self.dropped = self.dropped.saturating_add(1);
134 }
135 }
136
137 pub fn snapshot(&self) -> TerminalHistorySnapshot {
138 TerminalHistorySnapshot {
139 daemon_started_at_ms: self
140 .start_clock
141 .map_or(self.daemon_started_at_ms, |clock| clock.started_at_ms()),
142 entries: self.entries.iter().cloned().collect(),
143 dropped: self.dropped,
144 }
145 }
146}
147
148#[cfg(test)]
149mod tests {
150 use super::{TerminalRecord, TerminalRing, TerminalRingConfig};
151 use subc_control::{TerminalDisposition, TerminalExitKind};
152
153 fn record(at_ms: u64) -> TerminalRecord {
154 TerminalRecord {
155 exit_code: Some(1),
156 exit_signal: None,
157 at_ms,
158 disposition: TerminalDisposition::Restarting,
159 exit_kind: TerminalExitKind::Crash,
160 disposition_detail: None,
161 }
162 }
163
164 #[test]
165 fn the_ring_evicts_oldest_exits_and_counts_them() {
166 let mut ring = TerminalRing::new(TerminalRingConfig::new(2), 10);
167 ring.push(record(11));
168 ring.push(record(12));
169 ring.push(record(13));
170
171 let snapshot = ring.snapshot();
172 assert_eq!(snapshot.daemon_started_at_ms, 10);
173 assert_eq!(snapshot.dropped, 1);
174 assert_eq!(
175 snapshot
176 .entries
177 .iter()
178 .map(|entry| entry.at_ms)
179 .collect::<Vec<_>>(),
180 vec![12, 13]
181 );
182 }
183
184 #[test]
185 fn an_incoherent_zero_capacity_keeps_one_terminal() {
186 let mut ring = TerminalRing::new(TerminalRingConfig::new(0), 10);
187 ring.push(record(11));
188 ring.push(record(12));
189
190 let snapshot = ring.snapshot();
191 assert_eq!(snapshot.dropped, 1);
192 assert_eq!(snapshot.entries.len(), 1);
193 assert_eq!(snapshot.entries[0].at_ms, 12);
194 }
195}