Weave documentation
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(())
}

Source declarations