Skip to main content

appcore_log/
async_sink.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: async_sink.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/09/07 00:00:00 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/09/07 00:00:00 by dnettoRaw
8//      ###########      S: 1.0.2-rc
9// =============================================================================
10
11//! Explicit bounded asynchronous delivery for callers that accept queueing.
12
13use crate::{LogError, LogEvent, LogSink};
14use parking_lot::Mutex;
15use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
16use std::sync::mpsc::{self, Receiver, SyncSender, TrySendError};
17use std::sync::Arc;
18use std::thread::{self, JoinHandle};
19
20const ASYNC_LOG_THREAD_STACK_BYTES: usize = 256 * 1024;
21/// Maximum events accepted by one asynchronous sink configuration.
22pub const MAX_ASYNC_LOG_EVENTS: usize = 65_536;
23/// Maximum estimated retained bytes accepted by one asynchronous sink.
24pub const MAX_ASYNC_LOG_BYTES: usize = 64 * 1024 * 1024;
25
26/// Count and retained-byte ceilings for one asynchronous sink.
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct AsyncSinkConfig {
29    /// Maximum events retained across the queue and active delivery.
30    pub max_events: usize,
31    /// Maximum estimated event bytes retained across queue and active delivery.
32    pub max_bytes: usize,
33}
34
35/// Lock-free observation of one asynchronous sink.
36#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
37pub struct AsyncSinkStats {
38    /// Events currently retained by the queue or worker.
39    pub retained_events: usize,
40    /// Estimated event bytes currently retained by the queue or worker.
41    pub retained_bytes: usize,
42    /// Events delivered successfully.
43    pub delivered: u64,
44    /// Inner sink delivery failures.
45    pub failures: u64,
46    /// Events rejected because either queue ceiling was reached.
47    pub rejected: u64,
48}
49
50enum Message {
51    Event(Box<LogEvent>, usize),
52    Flush(SyncSender<()>),
53    Shutdown,
54}
55
56struct Counters {
57    retained_events: AtomicUsize,
58    retained_bytes: AtomicUsize,
59    delivered: AtomicU64,
60    failures: AtomicU64,
61    rejected: AtomicU64,
62}
63
64impl Counters {
65    fn new() -> Self {
66        Self {
67            retained_events: AtomicUsize::new(0),
68            retained_bytes: AtomicUsize::new(0),
69            delivered: AtomicU64::new(0),
70            failures: AtomicU64::new(0),
71            rejected: AtomicU64::new(0),
72        }
73    }
74}
75
76struct Lifecycle {
77    sender: SyncSender<Message>,
78    worker: Option<JoinHandle<()>>,
79}
80
81/// Opt-in non-blocking wrapper around one sink.
82///
83/// `emit` rejects immediately with [`LogError::Capacity`] when either bound is
84/// exhausted. Call [`Self::flush`] or [`Self::shutdown`] at an explicit
85/// lifecycle boundary; dropping the last owner also shuts down and joins the
86/// worker.
87pub struct AsyncSink {
88    config: AsyncSinkConfig,
89    lifecycle: Mutex<Lifecycle>,
90    counters: Arc<Counters>,
91    closed: AtomicBool,
92    accepts_sensitive: bool,
93}
94
95impl AsyncSink {
96    /// Starts one bounded worker around `sink`.
97    pub fn new(config: AsyncSinkConfig, sink: Arc<dyn LogSink>) -> Result<Self, LogError> {
98        if config.max_events == 0
99            || config.max_events > MAX_ASYNC_LOG_EVENTS
100            || config.max_bytes == 0
101            || config.max_bytes > MAX_ASYNC_LOG_BYTES
102        {
103            return Err(LogError::Capacity);
104        }
105        let (sender, receiver) = mpsc::sync_channel(config.max_events);
106        let counters = Arc::new(Counters::new());
107        let worker_counters = Arc::clone(&counters);
108        let accepts_sensitive = sink.accepts_sensitive();
109        let worker = thread::Builder::new()
110            .name("appcore-log-sink".to_string())
111            .stack_size(ASYNC_LOG_THREAD_STACK_BYTES)
112            .spawn(move || worker_loop(receiver, sink, &worker_counters))
113            .map_err(|_| LogError::Io)?;
114        Ok(Self {
115            config,
116            lifecycle: Mutex::new(Lifecycle {
117                sender,
118                worker: Some(worker),
119            }),
120            counters,
121            closed: AtomicBool::new(false),
122            accepts_sensitive,
123        })
124    }
125
126    /// Waits until every event admitted before this call has been delivered.
127    pub fn flush(&self) -> Result<(), LogError> {
128        if self.closed.load(Ordering::Acquire) {
129            return Err(LogError::Io);
130        }
131        let lifecycle = self.lifecycle.lock();
132        let (sender, receiver) = mpsc::sync_channel(0);
133        lifecycle
134            .sender
135            .send(Message::Flush(sender))
136            .map_err(|_| LogError::Io)?;
137        receiver.recv().map_err(|_| LogError::Io)
138    }
139
140    /// Drains admitted events, terminates the worker and rejects future emits.
141    ///
142    /// This waits for the wrapped sink. Use only with a sink whose own I/O has
143    /// a bounded completion contract; Rust cannot terminate arbitrary blocking
144    /// sink code safely.
145    pub fn shutdown(&self) -> Result<(), LogError> {
146        if self.closed.swap(true, Ordering::AcqRel) {
147            return Ok(());
148        }
149        let (sent, worker) = {
150            let mut lifecycle = self.lifecycle.lock();
151            (
152                lifecycle.sender.send(Message::Shutdown).is_ok(),
153                lifecycle.worker.take(),
154            )
155        };
156        let joined = worker.ok_or(LogError::Io)?.join().is_ok();
157        if sent && joined {
158            Ok(())
159        } else {
160            Err(LogError::Io)
161        }
162    }
163
164    /// Returns bounded queue and delivery counters without waiting for I/O.
165    pub fn stats(&self) -> AsyncSinkStats {
166        AsyncSinkStats {
167            retained_events: self.counters.retained_events.load(Ordering::Relaxed),
168            retained_bytes: self.counters.retained_bytes.load(Ordering::Relaxed),
169            delivered: self.counters.delivered.load(Ordering::Relaxed),
170            failures: self.counters.failures.load(Ordering::Relaxed),
171            rejected: self.counters.rejected.load(Ordering::Relaxed),
172        }
173    }
174
175    fn reserve(&self, bytes: usize) -> Result<(), LogError> {
176        if bytes > self.config.max_bytes
177            || !reserve_bounded(&self.counters.retained_events, 1, self.config.max_events)
178        {
179            increment(&self.counters.rejected);
180            return Err(LogError::Capacity);
181        }
182        if !reserve_bounded(&self.counters.retained_bytes, bytes, self.config.max_bytes) {
183            self.counters
184                .retained_events
185                .fetch_sub(1, Ordering::Relaxed);
186            increment(&self.counters.rejected);
187            return Err(LogError::Capacity);
188        }
189        Ok(())
190    }
191
192    fn release(&self, bytes: usize) {
193        self.counters
194            .retained_events
195            .fetch_sub(1, Ordering::Relaxed);
196        self.counters
197            .retained_bytes
198            .fetch_sub(bytes, Ordering::Relaxed);
199    }
200}
201
202impl LogSink for AsyncSink {
203    fn emit(&self, event: &LogEvent) -> Result<(), LogError> {
204        if self.closed.load(Ordering::Acquire) {
205            return Err(LogError::Io);
206        }
207        let bytes = event.retained_bytes();
208        self.reserve(bytes)?;
209        let lifecycle = self.lifecycle.lock();
210        if self.closed.load(Ordering::Acquire) {
211            self.release(bytes);
212            return Err(LogError::Io);
213        }
214        match lifecycle
215            .sender
216            .try_send(Message::Event(Box::new(event.clone()), bytes))
217        {
218            Ok(()) => Ok(()),
219            Err(TrySendError::Full(_)) => {
220                self.release(bytes);
221                increment(&self.counters.rejected);
222                Err(LogError::Capacity)
223            }
224            Err(TrySendError::Disconnected(_)) => {
225                self.release(bytes);
226                Err(LogError::Io)
227            }
228        }
229    }
230
231    fn accepts_sensitive(&self) -> bool {
232        self.accepts_sensitive
233    }
234
235    fn name(&self) -> &'static str {
236        "async"
237    }
238}
239
240impl Drop for AsyncSink {
241    fn drop(&mut self) {
242        self.closed.store(true, Ordering::Release);
243        let _ = self.lifecycle.get_mut().sender.try_send(Message::Shutdown);
244    }
245}
246
247fn reserve_bounded(counter: &AtomicUsize, amount: usize, maximum: usize) -> bool {
248    counter
249        .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
250            current.checked_add(amount).filter(|next| *next <= maximum)
251        })
252        .is_ok()
253}
254
255fn increment(counter: &AtomicU64) {
256    let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
257        Some(current.saturating_add(1))
258    });
259}
260
261fn worker_loop(receiver: Receiver<Message>, sink: Arc<dyn LogSink>, counters: &Counters) {
262    while let Ok(message) = receiver.recv() {
263        match message {
264            Message::Event(event, bytes) => {
265                let delivered =
266                    std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| sink.emit(&event)));
267                if matches!(delivered, Ok(Ok(()))) {
268                    increment(&counters.delivered);
269                } else {
270                    increment(&counters.failures);
271                }
272                counters.retained_events.fetch_sub(1, Ordering::Relaxed);
273                counters.retained_bytes.fetch_sub(bytes, Ordering::Relaxed);
274            }
275            Message::Flush(completed) => {
276                let _ = completed.send(());
277            }
278            Message::Shutdown => break,
279        }
280    }
281}
282
283#[cfg(test)]
284mod tests {
285    use super::*;
286    use crate::{Severity, Verbosity};
287    use std::sync::{Condvar, Mutex as StdMutex};
288
289    struct GatedSink {
290        state: StdMutex<(bool, bool)>,
291        changed: Condvar,
292    }
293
294    impl GatedSink {
295        fn new() -> Self {
296            Self {
297                state: StdMutex::new((false, false)),
298                changed: Condvar::new(),
299            }
300        }
301
302        fn wait_until_entered(&self) {
303            let mut state = self.state.lock().unwrap();
304            while !state.0 {
305                state = self.changed.wait(state).unwrap();
306            }
307        }
308
309        fn release(&self) {
310            let mut state = self.state.lock().unwrap();
311            state.1 = true;
312            self.changed.notify_all();
313        }
314    }
315
316    impl LogSink for GatedSink {
317        fn emit(&self, _event: &LogEvent) -> Result<(), LogError> {
318            let mut state = self.state.lock().unwrap();
319            state.0 = true;
320            self.changed.notify_all();
321            while !state.1 {
322                state = self.changed.wait(state).unwrap();
323            }
324            Ok(())
325        }
326    }
327
328    #[test]
329    fn active_delivery_remains_inside_both_bounds() {
330        let inner = Arc::new(GatedSink::new());
331        let sink = AsyncSink::new(
332            AsyncSinkConfig {
333                max_events: 1,
334                max_bytes: 4096,
335            },
336            inner.clone(),
337        )
338        .unwrap();
339        let event = LogEvent::new(1, Severity::Info, Verbosity::V4, "test", "message");
340
341        sink.emit(&event).unwrap();
342        inner.wait_until_entered();
343        assert_eq!(sink.emit(&event), Err(LogError::Capacity));
344        assert_eq!(sink.stats().retained_events, 1);
345
346        inner.release();
347        sink.flush().unwrap();
348        assert_eq!(
349            sink.stats(),
350            AsyncSinkStats {
351                delivered: 1,
352                rejected: 1,
353                ..AsyncSinkStats::default()
354            }
355        );
356        sink.shutdown().unwrap();
357        assert_eq!(sink.emit(&event), Err(LogError::Io));
358    }
359}