pub struct AsyncSink { /* private fields */ }Expand description
Opt-in non-blocking wrapper around one sink.
emit rejects immediately with LogError::Capacity when either bound is
exhausted. Call Self::flush or Self::shutdown at an explicit
lifecycle boundary; dropping the last owner also shuts down and joins the
worker.
Implementations§
Source§impl AsyncSink
impl AsyncSink
Sourcepub fn new(
config: AsyncSinkConfig,
sink: Arc<dyn LogSink>,
) -> Result<Self, LogError>
pub fn new( config: AsyncSinkConfig, sink: Arc<dyn LogSink>, ) -> Result<Self, LogError>
Starts one bounded worker around sink.
Examples found in repository?
examples/async_file.rs (lines 33-39)
19fn main() -> Result<(), LogError> {
20 let directory = std::path::PathBuf::from("target/appcore-log-example");
21 std::fs::create_dir_all(&directory).map_err(|_| LogError::Io)?;
22
23 let path = directory.join("async.jsonl");
24 let file = Arc::new(FileSink::new(FileSinkConfig {
25 path: path.clone(),
26 max_bytes: LOG_SIZE_8_MIB,
27 sync_each_write: true,
28 retention: 2,
29 archive: None,
30 })?);
31
32 // This queue retains at most 256 events and 1 MiB, including active I/O.
33 let asynchronous = Arc::new(AsyncSink::new(
34 AsyncSinkConfig {
35 max_events: 256,
36 max_bytes: 1024 * 1024,
37 },
38 file,
39 )?);
40 let dispatcher = LogDispatcher::new(LogPolicy::default(), vec![asynchronous.clone()]);
41 let log = dispatcher.event(0, "application");
42
43 log.info("application started");
44 log.warn("storage response is slow");
45
46 // The lifecycle owner drains durable writes before exiting.
47 asynchronous.shutdown()?;
48 let path = std::fs::canonicalize(path).map_err(|_| LogError::Io)?;
49 println!("log written to {}", path.display());
50 Ok(())
51}Sourcepub fn flush(&self) -> Result<(), LogError>
pub fn flush(&self) -> Result<(), LogError>
Waits until every event admitted before this call has been delivered.
Sourcepub fn shutdown(&self) -> Result<(), LogError>
pub fn shutdown(&self) -> Result<(), LogError>
Drains admitted events, terminates the worker and rejects future emits.
This waits for the wrapped sink. Use only with a sink whose own I/O has a bounded completion contract; Rust cannot terminate arbitrary blocking sink code safely.
Examples found in repository?
examples/async_file.rs (line 47)
19fn main() -> Result<(), LogError> {
20 let directory = std::path::PathBuf::from("target/appcore-log-example");
21 std::fs::create_dir_all(&directory).map_err(|_| LogError::Io)?;
22
23 let path = directory.join("async.jsonl");
24 let file = Arc::new(FileSink::new(FileSinkConfig {
25 path: path.clone(),
26 max_bytes: LOG_SIZE_8_MIB,
27 sync_each_write: true,
28 retention: 2,
29 archive: None,
30 })?);
31
32 // This queue retains at most 256 events and 1 MiB, including active I/O.
33 let asynchronous = Arc::new(AsyncSink::new(
34 AsyncSinkConfig {
35 max_events: 256,
36 max_bytes: 1024 * 1024,
37 },
38 file,
39 )?);
40 let dispatcher = LogDispatcher::new(LogPolicy::default(), vec![asynchronous.clone()]);
41 let log = dispatcher.event(0, "application");
42
43 log.info("application started");
44 log.warn("storage response is slow");
45
46 // The lifecycle owner drains durable writes before exiting.
47 asynchronous.shutdown()?;
48 let path = std::fs::canonicalize(path).map_err(|_| LogError::Io)?;
49 println!("log written to {}", path.display());
50 Ok(())
51}Sourcepub fn stats(&self) -> AsyncSinkStats
pub fn stats(&self) -> AsyncSinkStats
Returns bounded queue and delivery counters without waiting for I/O.
Trait Implementations§
Auto Trait Implementations§
impl !Freeze for AsyncSink
impl !RefUnwindSafe for AsyncSink
impl !UnwindSafe for AsyncSink
impl Send for AsyncSink
impl Sync for AsyncSink
impl Unpin for AsyncSink
impl UnsafeUnpin for AsyncSink
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more