1use 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;
21pub const MAX_ASYNC_LOG_EVENTS: usize = 65_536;
23pub const MAX_ASYNC_LOG_BYTES: usize = 64 * 1024 * 1024;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct AsyncSinkConfig {
29 pub max_events: usize,
31 pub max_bytes: usize,
33}
34
35#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
37pub struct AsyncSinkStats {
38 pub retained_events: usize,
40 pub retained_bytes: usize,
42 pub delivered: u64,
44 pub failures: u64,
46 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
81pub 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 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 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 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 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}