1use std::sync::Mutex;
10
11use structfs_core_store::Value;
12
13use crate::gate::{CancelToken, Cancelled, Gate};
14
15#[derive(Debug, Clone, PartialEq)]
17pub struct TailPage<T> {
18 pub items: Vec<T>,
20 pub next: u64,
22 pub done: bool,
25}
26
27impl<T: Into<Value>> TailPage<T> {
28 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
50pub 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 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 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 pub fn finish(&self) {
95 self.lock().done = true;
96 self.gate.notify();
97 }
98
99 pub fn is_done(&self) -> bool {
101 self.lock().done
102 }
103
104 pub fn len(&self) -> u64 {
106 self.lock().events.len() as u64
107 }
108
109 pub fn is_empty(&self) -> bool {
111 self.len() == 0
112 }
113
114 fn page_from(state: &TailState<T>, seq: u64) -> TailPage<T> {
115 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 pub fn snapshot_from(&self, seq: u64) -> TailPage<T> {
127 Self::page_from(&self.lock(), seq)
128 }
129
130 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 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 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 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}