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