Skip to main content

structfs_handles/
tail.rs

1//! `TailLog`: an append-only event log with atomic tail reads.
2//!
3//! Streaming stores conventionally serve events at `events/from/{seq}` and
4//! terminal state at a separate path. Because those are two reads, every
5//! consumer needs a "close-out drain" for events landing between the tail
6//! read and the terminal check. `TailLog::read_from` returns items *and*
7//! terminal status in one atomic operation, so that race cannot exist.
8
9use std::sync::Mutex;
10
11use structfs_core_store::Value;
12
13use crate::gate::{CancelToken, Cancelled, Gate};
14
15/// One page of a tail read.
16#[derive(Debug, Clone, PartialEq)]
17pub struct TailPage<T> {
18    /// Events from the requested cursor to the end of the log.
19    pub items: Vec<T>,
20    /// The cursor to pass to the next read.
21    pub next: u64,
22    /// Whether the log is finished. When true, no more events will ever
23    /// arrive; the consumer can stop without a close-out drain.
24    pub done: bool,
25}
26
27impl<T: Into<Value>> TailPage<T> {
28    /// Encode as the conventional tail envelope:
29    /// `{items: [...], next: N, status: "open" | "done"}`.
30    pub fn into_value(self) -> Value {
31        let mut map = std::collections::BTreeMap::new();
32        map.insert(
33            "items".to_string(),
34            Value::Array(self.items.into_iter().map(Into::into).collect()),
35        );
36        map.insert("next".to_string(), Value::Integer(self.next as i64));
37        map.insert(
38            "status".to_string(),
39            Value::String(if self.done { "done" } else { "open" }.to_string()),
40        );
41        Value::Map(map)
42    }
43}
44
45struct TailState<T> {
46    events: Vec<T>,
47    done: bool,
48}
49
50/// An append-only event log with terminal state and parked tail reads.
51pub struct TailLog<T> {
52    state: Mutex<TailState<T>>,
53    gate: Gate,
54}
55
56impl<T: Clone> Default for TailLog<T> {
57    fn default() -> Self {
58        Self::new()
59    }
60}
61
62impl<T: Clone> TailLog<T> {
63    /// Create an empty, open log.
64    pub fn new() -> Self {
65        Self {
66            state: Mutex::new(TailState {
67                events: Vec::new(),
68                done: false,
69            }),
70            gate: Gate::new(),
71        }
72    }
73
74    fn lock(&self) -> std::sync::MutexGuard<'_, TailState<T>> {
75        self.state.lock().unwrap_or_else(|e| e.into_inner())
76    }
77
78    /// Append an event and wake parked readers.
79    ///
80    /// Returns false (dropping the event) if the log is already finished.
81    pub fn push(&self, event: T) -> bool {
82        {
83            let mut state = self.lock();
84            if state.done {
85                return false;
86            }
87            state.events.push(event);
88        }
89        self.gate.notify();
90        true
91    }
92
93    /// Mark the log finished and wake parked readers. Idempotent.
94    pub fn finish(&self) {
95        self.lock().done = true;
96        self.gate.notify();
97    }
98
99    /// Whether the log is finished.
100    pub fn is_done(&self) -> bool {
101        self.lock().done
102    }
103
104    /// Number of events appended so far.
105    pub fn len(&self) -> u64 {
106        self.lock().events.len() as u64
107    }
108
109    /// Whether no events have been appended.
110    pub fn is_empty(&self) -> bool {
111        self.len() == 0
112    }
113
114    fn page_from(state: &TailState<T>, seq: u64) -> TailPage<T> {
115        // Clamp an out-of-range cursor instead of panicking: a stale or
116        // corrupt cursor yields an empty page at the log's end.
117        let start = (seq as usize).min(state.events.len());
118        TailPage {
119            items: state.events[start..].to_vec(),
120            next: state.events.len() as u64,
121            done: state.done,
122        }
123    }
124
125    /// Non-blocking snapshot of the tail from `seq`.
126    pub fn snapshot_from(&self, seq: u64) -> TailPage<T> {
127        Self::page_from(&self.lock(), seq)
128    }
129
130    /// Atomic tail read: park until there are events past `seq` or the log
131    /// is finished, then return them together with the terminal status.
132    pub async fn read_from(&self, seq: u64) -> TailPage<T> {
133        self.gate
134            .wait_until(|| {
135                let state = self.lock();
136                if (state.events.len() as u64) > seq || state.done {
137                    Some(Self::page_from(&state, seq))
138                } else {
139                    None
140                }
141            })
142            .await
143    }
144
145    /// [`TailLog::read_from`], cancellable.
146    pub async fn read_from_cancellable(
147        &self,
148        seq: u64,
149        token: &CancelToken,
150    ) -> Result<TailPage<T>, Cancelled> {
151        self.gate
152            .wait_until_cancellable(token, || {
153                let state = self.lock();
154                if (state.events.len() as u64) > seq || state.done {
155                    Some(Self::page_from(&state, seq))
156                } else {
157                    None
158                }
159            })
160            .await
161    }
162}
163
164#[cfg(test)]
165mod tests {
166    use super::*;
167    use std::sync::Arc;
168
169    #[tokio::test]
170    async fn tail_read_returns_existing_events() {
171        let log = TailLog::new();
172        log.push(1);
173        log.push(2);
174
175        let page = log.read_from(0).await;
176        assert_eq!(page.items, vec![1, 2]);
177        assert_eq!(page.next, 2);
178        assert!(!page.done);
179    }
180
181    #[tokio::test]
182    async fn tail_read_parks_until_push() {
183        let log = Arc::new(TailLog::new());
184        let reader = {
185            let log = log.clone();
186            tokio::spawn(async move { log.read_from(0).await })
187        };
188        tokio::task::yield_now().await;
189        log.push("event");
190
191        let page = reader.await.unwrap();
192        assert_eq!(page.items, vec!["event"]);
193    }
194
195    #[tokio::test]
196    async fn terminal_status_arrives_with_items() {
197        // The close-out race: events pushed immediately before finish must
198        // arrive in the same page that reports done.
199        let log = Arc::new(TailLog::new());
200        log.push(1);
201
202        let reader = {
203            let log = log.clone();
204            tokio::spawn(async move {
205                let first = log.read_from(0).await;
206                let second = log.read_from(first.next).await;
207                (first, second)
208            })
209        };
210        tokio::task::yield_now().await;
211        log.push(2);
212        log.finish();
213
214        let (first, second) = reader.await.unwrap();
215        assert_eq!(first.items, vec![1]);
216        // Second read gets the racing event AND the terminal flag together.
217        assert_eq!(second.items, vec![2]);
218        assert!(second.done);
219    }
220
221    #[tokio::test]
222    async fn finish_unblocks_empty_tail() {
223        let log: Arc<TailLog<i32>> = Arc::new(TailLog::new());
224        let reader = {
225            let log = log.clone();
226            tokio::spawn(async move { log.read_from(0).await })
227        };
228        tokio::task::yield_now().await;
229        log.finish();
230
231        let page = reader.await.unwrap();
232        assert!(page.items.is_empty());
233        assert!(page.done);
234    }
235
236    #[tokio::test]
237    async fn stale_cursor_clamps() {
238        let log = TailLog::new();
239        log.push(1);
240        log.finish();
241
242        let page = log.read_from(999).await;
243        assert!(page.items.is_empty());
244        assert_eq!(page.next, 1);
245        assert!(page.done);
246    }
247
248    #[test]
249    fn push_after_finish_is_dropped() {
250        let log = TailLog::new();
251        assert!(log.push(1));
252        log.finish();
253        assert!(!log.push(2));
254        assert_eq!(log.len(), 1);
255    }
256
257    #[tokio::test]
258    async fn cancellation_fails_parked_read() {
259        let log: Arc<TailLog<i32>> = Arc::new(TailLog::new());
260        let token = CancelToken::new();
261        let reader = {
262            let log = log.clone();
263            let token = token.clone();
264            tokio::spawn(async move { log.read_from_cancellable(0, &token).await })
265        };
266        tokio::task::yield_now().await;
267        token.cancel();
268
269        assert!(reader.await.unwrap().is_err());
270    }
271
272    #[test]
273    fn page_value_encoding() {
274        let page = TailPage {
275            items: vec![Value::from(1i64)],
276            next: 1,
277            done: true,
278        };
279        let value = page.into_value();
280        let map = match value {
281            Value::Map(m) => m,
282            _ => panic!("expected map"),
283        };
284        assert_eq!(map.get("next"), Some(&Value::Integer(1)));
285        assert_eq!(map.get("status"), Some(&Value::from("done")));
286        assert!(matches!(map.get("items"), Some(Value::Array(a)) if a.len() == 1));
287    }
288}