Weave documentation
Nexus

Strand Inputs

Input Strand tracking is local progress state, not network discovery.

InputStrand::new(Arc<Strand>) holds the supplied Strand and a processed-length counter. Constructing it does not register with a Nexus. Use Nexus::add_input_strand(...).await to apply admission checks and process existing entries. The example only inspects the standalone progress wrapper; manually advancing its counter does not replicate or materialize entries.

This example is checked against source signatures; it was not compiled or run during this documentation pass.

use nexus::{InputStrand, StrandStats};
use std::sync::Arc;
use strand::{Strand, StrandConfig};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // Build an upstream writer Strand and append three events.
    let dir = tempfile::tempdir()?;
    let mut writer = Strand::new(StrandConfig::new().with_storage(dir.path())).await?;
    for payload in [b"e1".as_slice(), b"e2", b"e3"] {
        writer.append(payload).await?;
    }

    // Wrap it as a Nexus InputStrand. The Nexus tracks how far it has consumed
    // via an atomic processed_length so multi-strand merging is incremental.
    let input = InputStrand::new(Arc::new(writer));
    assert_eq!(input.processed_length(), 0);
    assert!(input.has_new_entries().await?);

    // Simulate the Nexus reader advancing across the first two events.
    input.set_processed_length(2);

    let stats = StrandStats::from_input_strand(&input).await?;
    println!(
        "writer_key_prefix={} total={} processed={} unprocessed={}",
        &stats.writer_key[..16.min(stats.writer_key.len())],
        stats.total_entries,
        stats.processed_entries,
        stats.unprocessed_entries,
    );
    Ok(())
}

Source declarations