Locus
Streams and Watchers
Memory-backed streams and drive-wide broadcast notifications.
ReadStream wraps in-memory Bytes; WriteStream buffers bytes and optionally invokes a synchronous flush callback. Neither constructor opens a Locus file or commits bytes to BlobStore. Locus::watch() registers a drive-wide watcher; subscribe to it and filter paths in your consumer. There is no path-filter parameter. Broadcast receivers can lag and must handle receive errors.
This example is checked against source signatures; it was not compiled or run during this documentation pass.
use bytes::Bytes;
use locus::{ReadStream, WatchEvent, Watcher, WriteStream};
use tokio::io::{AsyncReadExt, AsyncWriteExt};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
// ReadStream is an AsyncRead over an in-memory Bytes payload.
let mut reader = ReadStream::new(Bytes::from_static(b"hello locus stream"));
let mut buf = Vec::new();
reader.read_to_end(&mut buf).await?;
// WriteStream is an AsyncWrite with an in-memory buffer up to `max_size`,
// optionally draining via a flush callback (commit to blob storage, etc).
let mut writer = WriteStream::new(64 * 1024);
writer.write_all(b"locus write").await?;
writer.shutdown().await?;
let buffered: Bytes = writer.into_bytes();
// Watchers fan out drive events to broadcast subscribers; the engine
// wires them in via Locus::watch(). Here we drive one manually to inspect
// the WatchEvent variants a subscriber will receive.
let watcher = Watcher::new();
let mut rx = watcher.subscribe();
watcher.emit_created("/docs/notes.md".into(), false);
watcher.emit_modified("/docs/notes.md".into(), false);
let evt = rx.recv().await?;
assert!(matches!(evt, WatchEvent::Created { .. }));
println!(
"read={} bytes buffered={} bytes watcher_id={:?}",
buf.len(),
buffered.len(),
watcher.id(),
);
Ok(())
}