Skip to main content

subc_daemon/
terminal_ring.rs

1use 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/// Fixed-size retention policy for one module's terminal exits.
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub struct TerminalRingConfig {
12    max_entries: usize,
13}
14
15impl TerminalRingConfig {
16    /// A zero-sized history cannot answer whether an exit was observed, so clamp it
17    /// to one record instead of constructing a ring that lies by omission.
18    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/// One observed child exit and the disposition chosen by its supervisor.
32#[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    /// Why the supervisor chose this disposition, when the disposition alone
40    /// does not say. `failed` records the exhausted crash budget here, naming
41    /// the limit AND the window it was counted over, because a module stopped
42    /// by three crashes in ten minutes and one stopped by three crashes in a
43    /// week are the same `failed` and call for different reactions.
44    pub disposition_detail: Option<String>,
45}
46
47/// The retained terminal suffix for one module.
48#[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
55/// A durable history read captured under the ring lock and finished after it.
56pub(crate) enum DurableHistoryRead {
57    RingOnly(TerminalHistorySnapshot),
58    Journal(JournalRead),
59}
60
61impl DurableHistoryRead {
62    /// Blocking when a journal is configured: it reads journal files.
63    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/// Bounded terminal-exit history for one supervised module.
72///
73/// This belongs to the module rather than an individual child so a replacement
74/// process retains the exits that caused it to exist.
75#[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    /// Capture and read in one step, for tests that own the ring directly.
114    #[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    /// Pin this ring's snapshot and the journal files a history read will
120    /// see. Cheap; the caller may then release the ring lock before the
121    /// file reading in [`DurableHistoryRead::read`].
122    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}