diff --git a/crates/bobbin/src/config.rs b/crates/bobbin/src/config.rs index ce58e9d..89d2463 100644 --- a/crates/bobbin/src/config.rs +++ b/crates/bobbin/src/config.rs @@ -15,6 +15,7 @@ const KNOWN_KEYS: &[&str] = &[ "server.shutdown_grace_secs", "hydrant.url", "hydrant.start_cursor", + "ingest.parallelism", "slingshot.url", "record_cache.lru_bytes", "search.heap_bytes", @@ -30,6 +31,7 @@ const KNOWN_ENVS: &[&str] = &[ "BOBBIN_SHUTDOWN_GRACE_SECS", "BOBBIN_HYDRANT_URL", "BOBBIN_START_CURSOR", + "BOBBIN_INGEST_PARALLELISM", "BOBBIN_SLINGSHOT_URL", "BOBBIN_RECORD_LRU_BYTES", "BOBBIN_SEARCH_HEAP_BYTES", @@ -47,6 +49,9 @@ pub struct BobbinConfig { #[config(nested)] pub hydrant: HydrantConfig, + #[config(nested)] + pub ingest: IngestConfig, + #[config(nested)] pub slingshot: SlingshotConfig, @@ -114,6 +119,16 @@ pub struct HydrantConfig { pub start_cursor: u64, } +#[derive(Debug, Config)] +pub struct IngestConfig { + /// Concurrent in-flight resolves during ingest. One slot per slingshot rtt. + /// Default 16 sits at the measured throughput + /// knee for cold replay against a healthy hydrant. The committer stays + /// serial so that like, cursor and Spur allocation order are preserved. + #[config(env = "BOBBIN_INGEST_PARALLELISM", default = 16)] + pub parallelism: usize, +} + #[derive(Debug, Config)] pub struct SlingshotConfig { /// Base URL of a slingshot instance. Used for record bodies and identity. diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index 719964e..ae56be3 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -1,4 +1,5 @@ use std::net::SocketAddr; +use std::num::NonZeroUsize; use std::path::PathBuf; use std::process::ExitCode; use std::sync::Arc; @@ -130,6 +131,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { let slingshot = SlingshotClient::with_default_http(cfg.slingshot.url.clone())?; let resolver = Arc::new(RepoIdResolver::with_slingshot( slingshot.clone(), + clock.clone(), hasher.clone(), )); let edges = Arc::new(EdgeStore::new(hasher.clone())); @@ -148,9 +150,12 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { .with_context(|| format!("search.heap_bytes {} exceeds usize", cfg.search.heap_bytes))?; let search = Arc::new(SearchIndex::new(search_heap, clock.clone())?); + let parallelism = NonZeroUsize::new(cfg.ingest.parallelism) + .ok_or_else(|| anyhow!("ingest.parallelism must be at least 1"))?; let ingest_cfg = IngestConfig { hydrant_base: cfg.hydrant.url.clone(), start_cursor: HydrantCursor::new(cfg.hydrant.start_cursor), + parallelism, }; let cancel = CancellationToken::new(); let ingest_coverage = coverage.clone(); diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 713b4d1..673979b 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -1,3 +1,5 @@ +use std::collections::VecDeque; +use std::num::NonZeroUsize; use std::sync::Arc; use std::time::Duration; @@ -14,10 +16,12 @@ use futures::StreamExt; use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; use jacquard_common::types::string::{AtStrError, AtUri, Cid}; use thiserror::Error; use tokio::time::Instant; +use tokio_stream::wrappers::ReceiverStream; use tokio_util::sync::CancellationToken; use tracing::{debug, info, warn}; use url::Url; @@ -25,6 +29,7 @@ use url::Url; mod frame; mod resolver; pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; +use frame::HydrantStreamErrorFrame; pub use resolver::{RepoIdResolver, Resolution}; const TANGLED_PREFIX: &str = "sh.tangled."; @@ -35,13 +40,20 @@ const PONG_TIMEOUT: Duration = Duration::from_secs(15); const READY_SKEW: Duration = Duration::from_secs(60); const FRAME_CHANNEL_DEPTH: usize = 256; const CONTROL_CHANNEL_DEPTH: usize = 16; -const NORMALIZE_CONCURRENCY: usize = 8; +const READER_HOLD_LIMIT: usize = 64; +const SEND_TIMEOUT: Duration = Duration::from_secs(10); const NORMAL_CLOSE: u16 = 1000; +const METRICS_DUMP_INTERVAL: Duration = Duration::from_secs(10); +pub const DEFAULT_INGEST_PARALLELISM: NonZeroUsize = match NonZeroUsize::new(16) { + Some(n) => n, + None => unreachable!(), +}; #[derive(Clone, Debug)] pub struct IngestConfig { pub hydrant_base: Url, pub start_cursor: HydrantCursor, + pub parallelism: NonZeroUsize, } impl IngestConfig { @@ -49,6 +61,7 @@ impl IngestConfig { Self { hydrant_base, start_cursor: HydrantCursor::new(0), + parallelism: DEFAULT_INGEST_PARALLELISM, } } @@ -88,6 +101,12 @@ pub enum IngestError { Extract(#[from] ExtractError), #[error("hydrant did not respond to ping within {0:?}")] PongTimeout(Duration), + #[error("websocket send blocked for at least {0:?}, treating link as dead")] + SendTimeout(Duration), + #[error("hydrant disconnected because bobbin's stream consumer fell behind: {message}")] + ConsumerTooSlow { message: String }, + #[error("hydrant signaled stream error {code}: {message}")] + HydrantStream { code: String, message: String }, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -133,11 +152,23 @@ impl Clone for IngestRuntime { pub async fn run( config: IngestConfig, runtime: IngestRuntime, +) -> Result<(), IngestError> { + let metrics_dumper = spawn_metrics_dumper(&runtime); + let result = run_inner(config, &runtime).await; + if let Err(join) = metrics_dumper.await { + warn!(?join, "metrics dumper task panicked"); + } + result +} + +async fn run_inner( + config: IngestConfig, + runtime: &IngestRuntime, ) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; loop { let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); - let SessionEnd { outcome, error } = run_session(&config, cursor, &runtime).await; + let SessionEnd { outcome, error } = run_session(&config, cursor, runtime).await; if runtime.cancel.is_cancelled() { info!( last_cursor = runtime.coverage.snapshot().last_cursor().raw(), @@ -167,6 +198,36 @@ pub async fn run( } } +fn spawn_metrics_dumper( + runtime: &IngestRuntime, +) -> tokio::task::JoinHandle<()> { + let rt = runtime.clone(); + tokio::spawn(async move { + loop { + tokio::select! { + biased; + _ = rt.cancel.cancelled() => break, + _ = rt.clock.sleep(METRICS_DUMP_INTERVAL) => { + let s = rt.resolver.stats(); + info!( + target: "bobbin_ingest::metrics", + resolver_hits = s.hits, + resolver_misses_mapped = s.misses_mapped, + resolver_misses_no_repo_did = s.misses_no_repo_did, + resolver_misses_unresolvable = s.misses_unresolvable, + resolver_misses_transient = s.misses_transient, + resolver_misses_no_client = s.misses_no_client, + resolver_miss_latency_micros_avg = s.miss_latency_micros_avg().unwrap_or(0), + resolver_miss_latency_micros_max = s.miss_latency_micros_max, + resolver_total = s.total(), + "resolver stats", + ); + } + } + } + }) +} + fn next_connect_cursor(snapshot: Coverage, start: HydrantCursor) -> HydrantCursor { if snapshot.events_processed() == 0 { start @@ -216,28 +277,25 @@ async fn run_session( } }; - let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); + let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); + let parallelism = config.parallelism.get(); let processor_runtime = runtime.clone(); let processor = tokio::spawn(async move { - loop { - tokio::select! { - biased; - _ = processor_runtime.cancel.cancelled() => break, - next = frame_rx.recv() => { - let Some(frame) = next else { break }; - let now = processor_runtime.clock.now_unix_micros(); - handle_frame( - frame, - &processor_runtime.store, - &processor_runtime.coverage, - &*processor_runtime.search, - &*processor_runtime.records, - &processor_runtime.resolver, - now, - ) - .await; - } - } + let cancel = processor_runtime.cancel.clone(); + let prep_rt = processor_runtime.clone(); + let resolve_rt = processor_runtime.clone(); + let commit_rt = processor_runtime; + + let pipeline = ReceiverStream::new(frame_rx) + .then(move |frame| prep_stage(frame, prep_rt.clone())) + .map(move |staged| resolve_stage(staged, resolve_rt.clone())) + .buffered(parallelism) + .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)); + + tokio::select! { + biased; + _ = cancel.cancelled() => {}, + _ = pipeline => {}, } }); @@ -253,14 +311,17 @@ async fn run_session( tokio::select! { biased; _ = runtime.cancel.cancelled() => { - let _ = ws_sink.send(WsMessage::Close { code: NORMAL_CLOSE, reason: "bobbin shutdown".to_owned() }).await; + let _ = timed_send( + &mut ws_sink, + WsMessage::Close { code: NORMAL_CLOSE, reason: "bobbin shutdown".to_owned() }, + ).await; break None; } _ = runtime.clock.sleep_until(next_ping) => { next_ping = runtime.clock.now_instant() + PING_INTERVAL; if pong_deadline.is_none() { - if let Err(e) = ws_sink.send(WsMessage::Ping(Bytes::new())).await { - break Some(IngestError::Network(e)); + if let Err(e) = timed_send(&mut ws_sink, WsMessage::Ping(Bytes::new())).await { + break Some(e); } pong_deadline = Some(runtime.clock.now_instant() + PONG_TIMEOUT); } @@ -272,8 +333,8 @@ async fn run_session( let Some(evt) = evt else { break None; }; match evt { WsEvent::IncomingPing(payload) => { - if let Err(e) = ws_sink.send(WsMessage::Pong(payload)).await { - break Some(IngestError::Network(e)); + if let Err(e) = timed_send(&mut ws_sink, WsMessage::Pong(payload)).await { + break Some(e); } } WsEvent::IncomingPong => { @@ -301,6 +362,17 @@ async fn run_session( } } +async fn timed_send( + sink: &mut Box, + msg: WsMessage, +) -> Result<(), IngestError> { + match tokio::time::timeout(SEND_TIMEOUT, sink.send(msg)).await { + Ok(Ok(())) => Ok(()), + Ok(Err(e)) => Err(IngestError::Network(e)), + Err(_elapsed) => Err(IngestError::SendTimeout(SEND_TIMEOUT)), + } +} + async fn reader_loop( mut ws_stream: Box, frame_tx: tokio::sync::mpsc::Sender, @@ -308,11 +380,55 @@ async fn reader_loop( cancel: CancellationToken, ) -> SessionEnd { let mut outcome = SessionOutcome::Empty; + let mut held: VecDeque = VecDeque::new(); + let error: Option = loop { - tokio::select! { - biased; - _ = cancel.cancelled() => break None, - msg = ws_stream.next() => { + let step = if held.is_empty() { + tokio::select! { + biased; + _ = cancel.cancelled() => ReaderStep::Cancelled, + msg = ws_stream.next() => ReaderStep::WsMessage(msg), + } + } else if held.len() < READER_HOLD_LIMIT { + tokio::select! { + biased; + _ = cancel.cancelled() => ReaderStep::Cancelled, + permit_result = frame_tx.reserve() => match permit_result { + Ok(permit) => { + let frame = held + .pop_front() + .expect("non-empty held when reserve succeeds"); + permit.send(frame); + ReaderStep::Sent + } + Err(_) => ReaderStep::FrameSinkClosed, + }, + msg = ws_stream.next() => ReaderStep::WsMessage(msg), + } + } else { + tokio::select! { + biased; + _ = cancel.cancelled() => ReaderStep::Cancelled, + permit_result = frame_tx.reserve() => match permit_result { + Ok(permit) => { + let frame = held + .pop_front() + .expect("held at limit when reserve succeeds"); + permit.send(frame); + ReaderStep::Sent + } + Err(_) => ReaderStep::FrameSinkClosed, + }, + } + }; + + match step { + ReaderStep::Cancelled => break None, + ReaderStep::FrameSinkClosed => break None, + ReaderStep::Sent => { + outcome = SessionOutcome::Progressed; + } + ReaderStep::WsMessage(msg) => { let Some(msg) = msg else { break None; }; let parsed = match msg { Ok(m) => m, @@ -320,18 +436,11 @@ async fn reader_loop( }; match parsed { WsMessage::Text(text) => { - let frame: HydrantFrame = match serde_json::from_str(&text) { + let frame = match classify_text_frame(&text) { Ok(f) => f, - Err(e) => break Some(IngestError::Decode(e)), + Err(e) => break Some(e), }; - tokio::select! { - biased; - _ = cancel.cancelled() => break None, - res = frame_tx.send(frame) => { - if res.is_err() { break None; } - } - } - outcome = SessionOutcome::Progressed; + held.push_back(frame); } WsMessage::Binary(_) => { debug!("hydrant sent unexpected binary frame, ignoring"); @@ -357,6 +466,45 @@ async fn reader_loop( SessionEnd { outcome, error } } +enum ReaderStep { + Cancelled, + FrameSinkClosed, + Sent, + WsMessage(Option>), +} + +fn classify_text_frame(text: &str) -> Result { + #[derive(serde::Deserialize)] + struct PeekType<'a> { + #[serde(rename = "type", borrow)] + kind: Option>, + } + let is_error_frame = serde_json::from_str::(text) + .ok() + .and_then(|p| p.kind) + .as_deref() + == Some("error"); + if is_error_frame { + return match serde_json::from_str::(text) { + Ok(err_frame) => Err(classify_hydrant_error(err_frame)), + Err(decode_err) => Err(IngestError::Decode(decode_err)), + }; + } + serde_json::from_str::(text).map_err(IngestError::Decode) +} + +fn classify_hydrant_error(frame: HydrantStreamErrorFrame) -> IngestError { + let HydrantStreamErrorFrame { error, message } = frame; + let message = message.unwrap_or_default(); + match error.as_str() { + "ConsumerTooSlow" => IngestError::ConsumerTooSlow { message }, + _ => IngestError::HydrantStream { + code: error, + message, + }, + } +} + #[derive(Debug)] enum WsEvent { IncomingPing(Bytes), @@ -370,81 +518,180 @@ async fn wait_until(deadline: Option, clock: &dyn Clock) { } } -async fn handle_frame( +struct Pending { + cursor: HydrantCursor, + signal: PromotionSignal, + op: PendingOp, +} + +enum PendingOp { + Noop, + ClearCache { + source: AtUri, + }, + Upsert { + source: AtUri, + nsid: Nsid, + parsed: Record, + bytes: Vec, + cid: Option>, + edges: Vec, + }, + Delete { + source: AtUri, + nsid: Nsid, + }, +} + +struct Prepared { + pending: Pending, + prepare_start: Instant, + prepare_end: Instant, +} + +struct Resolved { + pending: Pending, + prepare_start: Instant, + prepare_end: Instant, + resolve_start: Instant, + resolve_end: Instant, +} + +fn pending_nsid(op: &PendingOp) -> Option<&Nsid> { + match op { + PendingOp::Upsert { nsid, .. } => Some(nsid), + PendingOp::Delete { nsid, .. } => Some(nsid), + PendingOp::Noop | PendingOp::ClearCache { .. } => None, + } +} + +fn pending_edge_count(op: &PendingOp) -> u64 { + match op { + PendingOp::Upsert { edges, .. } => edges.len() as u64, + _ => 0, + } +} + +async fn prep_stage( + frame: HydrantFrame, + rt: IngestRuntime, +) -> Prepared { + let now = rt.clock.now_unix_micros(); + let prepare_start = rt.clock.now_instant(); + let pending = prepare_frame(frame, &rt.resolver, now).await; + let prepare_end = rt.clock.now_instant(); + Prepared { + pending, + prepare_start, + prepare_end, + } +} + +async fn resolve_stage( + staged: Prepared, + rt: IngestRuntime, +) -> Resolved { + let resolve_start = rt.clock.now_instant(); + let pending = resolve_pending(staged.pending, &rt.resolver).await; + let resolve_end = rt.clock.now_instant(); + Resolved { + pending, + prepare_start: staged.prepare_start, + prepare_end: staged.prepare_end, + resolve_start, + resolve_end, + } +} + +async fn commit_stage( + staged: Resolved, + rt: IngestRuntime, + parallelism: usize, +) { + let nsid = pending_nsid(&staged.pending.op).cloned(); + let edge_count = pending_edge_count(&staged.pending.op); + let Resolved { + pending, + prepare_start, + prepare_end, + resolve_start, + resolve_end, + } = staged; + let commit_start = rt.clock.now_instant(); + commit_pending(pending, &rt.store, &rt.coverage, &*rt.search, &*rt.records).await; + let commit_end = rt.clock.now_instant(); + tracing::trace!( + target: "bobbin_ingest::stage", + nsid = %nsid.as_ref().map(Nsid::as_str).unwrap_or(""), + edge_count, + prepare_us = prepare_end.duration_since(prepare_start).as_micros() as u64, + queue_resolve_wait_us = resolve_start.duration_since(prepare_end).as_micros() as u64, + resolve_us = resolve_end.duration_since(resolve_start).as_micros() as u64, + queue_commit_wait_us = commit_start.duration_since(resolve_end).as_micros() as u64, + commit_us = commit_end.duration_since(commit_start).as_micros() as u64, + total_us = commit_end.duration_since(prepare_start).as_micros() as u64, + parallelism, + "pipeline stage timings", + ); +} + +async fn prepare_frame( frame: HydrantFrame, - store: &EdgeStore, - coverage: &CoverageWatch, - search: &S, - records: &dyn RecordStore, resolver: &RepoIdResolver, now: UnixMicros, -) { +) -> Pending { let cursor = HydrantCursor::new(frame.id); let signal = promotion_signal(frame.record.as_ref(), now); - coverage.update(|c| c.advance(cursor).maybe_promote(signal)); - - match frame.kind { - FrameKind::Record => { - let Some(record) = frame.record else { - debug!( - id = frame.id, - "record-typed frame missing payload, skipping" - ); - return; - }; - if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { - return; - } - if let Err(err) = apply_record(record, store, search, records, resolver).await { - warn!(?err, "dropped record"); - } - } - FrameKind::Identity | FrameKind::Account => {} + let op = match frame.kind { + FrameKind::Record => prepare_record(frame.record, resolver).await, + FrameKind::Identity | FrameKind::Account => PendingOp::Noop, FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); + PendingOp::Noop } + }; + Pending { + cursor, + signal, + op, } } -fn promotion_signal(record: Option<&RecordFrame>, now: UnixMicros) -> PromotionSignal { - PromotionSignal { - live: record.is_some_and(|r| r.live), - rev_micros: record.map(|r| r.rev.timestamp()), - now_micros: now.raw(), - skew_micros: READY_SKEW.as_micros() as u64, +async fn prepare_record(record: Option, resolver: &RepoIdResolver) -> PendingOp { + let Some(record) = record else { + debug!("record-typed frame missing payload, skipping"); + return PendingOp::Noop; + }; + if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { + return PendingOp::Noop; } -} - -async fn apply_record( - record: RecordFrame, - store: &EdgeStore, - search: &S, - records: &dyn RecordStore, - resolver: &RepoIdResolver, -) -> Result<(), IngestError> { - let source = build_source_uri(&record)?; + let nsid = record.collection.clone(); + let source = match build_source_uri(&record) { + Ok(s) => s, + Err(e) => { + warn!(?e, "invalid frame source, dropping record"); + return PendingOp::Noop; + } + }; match record.action { RecordAction::Create | RecordAction::Update => { let Some(value) = record.record else { - records.remove(&source); - debug!(collection = %record.collection.as_ref(), "create/update missing record body, skipping"); - return Ok(()); + debug!(collection = %nsid, "create/update missing record body, clearing cache"); + return PendingOp::ClearCache { source }; }; let bytes = serde_json::to_vec(&value) .expect("serde_json::Value re-serialization is infallible"); let parsed = match Record::from_json_bytes(record.collection.as_ref(), &bytes) { Ok(r) => r, Err(ExtractError::UnknownCollection(name)) => { - records.remove(&source); - debug!(collection = %name, "unknown sh.tangled.* collection, skipping"); - return Ok(()); + debug!(collection = %name, "unknown sh.tangled.* collection, clearing cache"); + return PendingOp::ClearCache { source }; } Err(e) => { - records.remove(&source); - return Err(e.into()); + warn!(?e, "record decode failed, clearing cache"); + return PendingOp::ClearCache { source }; } }; - cache_body(records, &source, record.cid, bytes); if let Record::Repo(repo) = &parsed { resolver .observe( @@ -454,23 +701,118 @@ async fn apply_record( ) .await; } - let raw_edges = parsed.extract_edges(&source)?; - let edges = normalize_subjects(raw_edges, resolver).await; + let edges = match parsed.extract_edges(&source) { + Ok(es) => es, + Err(e) => { + warn!(?e, "edge extraction failed, clearing cache"); + return PendingOp::ClearCache { source }; + } + }; + PendingOp::Upsert { + source, + nsid, + parsed, + bytes, + cid: record.cid, + edges, + } + } + RecordAction::Delete => PendingOp::Delete { source, nsid }, + RecordAction::Other => { + debug!(collection = %nsid, "ignoring unknown record action"); + PendingOp::Noop + } + } +} + +async fn resolve_pending(pending: Pending, resolver: &RepoIdResolver) -> Pending { + let Pending { cursor, signal, op } = pending; + let op = match op { + PendingOp::Upsert { + source, + nsid, + parsed, + bytes, + cid, + edges, + } => { + let edges = normalize_subjects(edges, resolver).await; + PendingOp::Upsert { + source, + nsid, + parsed, + bytes, + cid, + edges, + } + } + other => other, + }; + Pending { + cursor, + signal, + op, + } +} + +async fn commit_pending( + pending: Pending, + store: &EdgeStore, + coverage: &CoverageWatch, + search: &S, + records: &dyn RecordStore, +) { + let Pending { cursor, signal, op } = pending; + match op { + PendingOp::Noop => {} + PendingOp::ClearCache { source } => records.remove(&source), + PendingOp::Upsert { + source, + nsid: _, + parsed, + bytes, + cid, + edges, + } => { + cache_body(records, &source, cid, bytes); store.upsert_source(&source, edges); if let Some(searchable) = SearchableRecord::try_from_record(parsed) { search.upsert(searchable.to_search_doc(&source)).await; } } - RecordAction::Delete => { + PendingOp::Delete { source, nsid: _ } => { store.remove_source(&source); - search.remove(&source).await; records.remove(&source); + search.remove(&source).await; } - RecordAction::Other => { - debug!(collection = %record.collection.as_ref(), "ignoring unknown record action"); - } } - Ok(()) + coverage.update(|c| c.advance(cursor).maybe_promote(signal)); +} + +#[cfg(test)] +async fn handle_frame( + frame: HydrantFrame, + store: &EdgeStore, + coverage: &CoverageWatch, + search: &S, + records: &dyn RecordStore, + resolver: &RepoIdResolver, + clock: &dyn Clock, + now: UnixMicros, +) { + let _ = clock; + let pending = prepare_frame(frame, resolver, now).await; + let pending = resolve_pending(pending, resolver).await; + commit_pending(pending, store, coverage, search, records).await; +} + +fn promotion_signal(record: Option<&RecordFrame>, now: UnixMicros) -> PromotionSignal { + PromotionSignal { + live: record.is_some_and(|r| r.live), + rev_micros: record.map(|r| r.rev.timestamp()), + now_micros: now.raw(), + skew_micros: READY_SKEW.as_micros() as u64, + } } fn cache_body( @@ -494,7 +836,7 @@ fn cache_body( async fn normalize_subjects(edges: Vec, resolver: &RepoIdResolver) -> Vec { futures::stream::iter(edges) - .map(|edge| async move { + .then(|edge| async move { let Some((owner, rkey)) = parse_repo_subject(&edge.subject) else { return edge; }; @@ -506,7 +848,6 @@ async fn normalize_subjects(edges: Vec, resolver: &RepoIdResolver) -> Vec< Resolution::NoRepoDid | Resolution::Unresolvable => edge, } }) - .buffered(NORMALIZE_CONCURRENCY) .collect() .await } @@ -563,6 +904,10 @@ mod tests { SystemClock::new().now_unix_micros() } + fn sys_clock() -> SystemClock { + SystemClock::new() + } + fn parse_frame(value: serde_json::Value) -> HydrantFrame { let text = serde_json::to_string(&value).expect("serialize fixture"); serde_json::from_str(&text).expect("deserialize fixture") @@ -595,6 +940,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -629,6 +975,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -658,6 +1005,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -693,6 +1041,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -703,6 +1052,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -742,7 +1092,7 @@ mod tests { })); let source = AtUri::new_owned("at://did:plc:olaren/sh.tangled.feed.star/abcabcabcabcz").unwrap(); - handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, now()).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now()).await; let cached = lru.get(&source).expect("hydrant cid must seed the lru"); assert_eq!(cached.cid.as_ref(), VALID_CID); let parsed: serde_json::Value = serde_json::from_slice(&cached.value).unwrap(); @@ -781,7 +1131,7 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, now()).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now()).await; assert!( lru.get(&source).is_none(), "missing cid means we cannot trust the body, so the lru must be cleared", @@ -816,7 +1166,7 @@ mod tests { "record": null } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, now()).await; + handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now()).await; assert!(lru.get(&source).is_none()); } @@ -848,6 +1198,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -883,6 +1234,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -890,6 +1242,284 @@ mod tests { assert!(matches!(cov.snapshot(), Coverage::Warming { .. })); } + #[test] + fn classify_consumer_too_slow_frame_returns_typed_variant() { + let text = r#"{"type":"error","error":"ConsumerTooSlow","message":"stream socket send blocked for at least 30 seconds"}"#; + match classify_text_frame(text) { + Err(IngestError::ConsumerTooSlow { message }) => { + assert!( + message.contains("30 seconds"), + "message field preserved verbatim, got: {message}" + ); + } + other => panic!("expected ConsumerTooSlow variant, got: {other:?}"), + } + } + + #[test] + fn classify_unknown_hydrant_error_falls_back_to_generic_variant() { + let text = + r#"{"type":"error","error":"NewFutureCode","message":"some new failure mode"}"#; + match classify_text_frame(text) { + Err(IngestError::HydrantStream { code, message }) => { + assert_eq!(code, "NewFutureCode"); + assert_eq!(message, "some new failure mode"); + } + other => panic!("expected HydrantStream variant, got: {other:?}"), + } + } + + #[test] + fn classify_error_frame_without_message_uses_empty_string() { + let text = r#"{"type":"error","error":"ConsumerTooSlow"}"#; + match classify_text_frame(text) { + Err(IngestError::ConsumerTooSlow { message }) => assert!(message.is_empty()), + other => panic!("expected ConsumerTooSlow with empty message, got: {other:?}"), + } + } + + #[test] + fn classify_normal_record_frame_unchanged() { + let text = r#"{"id":42,"type":"record"}"#; + let frame = classify_text_frame(text).expect("normal record frame must parse"); + assert_eq!(frame.id, 42); + assert_eq!(frame.kind, FrameKind::Record); + } + + #[test] + fn classify_garbage_object_returns_decode_error() { + let text = r#"{"random":"object","without":"required fields"}"#; + match classify_text_frame(text) { + Err(IngestError::Decode(_)) => {} + other => panic!("expected Decode error, got: {other:?}"), + } + } + + struct ScriptedWsStream { + messages: std::collections::VecDeque>, + } + + impl WsStream for ScriptedWsStream { + fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { + let msg = self.messages.pop_front(); + Box::pin(async move { msg }) + } + } + + #[tokio::test] + async fn reader_loop_surfaces_consumer_too_slow_from_error_frame() { + let mut messages = std::collections::VecDeque::new(); + messages.push_back(Ok(WsMessage::Text( + r#"{"type":"error","error":"ConsumerTooSlow","message":"stream delivery blocked"}"# + .to_owned(), + ))); + let stream: Box = Box::new(ScriptedWsStream { messages }); + let (frame_tx, _frame_rx) = + tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); + let (control_tx, _control_rx) = + tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); + let cancel = CancellationToken::new(); + + let end = reader_loop(stream, frame_tx, control_tx, cancel).await; + + match end.error { + Some(IngestError::ConsumerTooSlow { message }) => { + assert_eq!(message, "stream delivery blocked"); + } + other => panic!("expected ConsumerTooSlow SessionEnd error, got: {other:?}"), + } + assert_eq!(end.outcome, SessionOutcome::Empty); + } + + #[tokio::test] + async fn reader_loop_surfaces_unknown_hydrant_error_distinctly() { + let mut messages = std::collections::VecDeque::new(); + messages.push_back(Ok(WsMessage::Text( + r#"{"type":"error","error":"NewFutureCode","message":"new mode"}"#.to_owned(), + ))); + let stream: Box = Box::new(ScriptedWsStream { messages }); + let (frame_tx, _frame_rx) = + tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); + let (control_tx, _control_rx) = + tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); + let cancel = CancellationToken::new(); + + let end = reader_loop(stream, frame_tx, control_tx, cancel).await; + + match end.error { + Some(IngestError::HydrantStream { code, message }) => { + assert_eq!(code, "NewFutureCode"); + assert_eq!(message, "new mode"); + } + other => panic!("expected HydrantStream SessionEnd error, got: {other:?}"), + } + } + + struct ChannelWsStream { + rx: tokio::sync::mpsc::Receiver>, + } + + impl WsStream for ChannelWsStream { + fn next<'a>(&'a mut self) -> bobbin_runtime::WsMessageFuture<'a> { + Box::pin(async move { self.rx.recv().await }) + } + } + + fn star_frame_text(id: u64, rkey: &str) -> String { + json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": rkey, + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + }) + .to_string() + } + + #[tokio::test] + async fn pong_forwards_promptly_when_frame_channel_is_full() { + let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::>(8); + let stream: Box = Box::new(ChannelWsStream { rx: ws_rx }); + let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::(1); + let (control_tx, mut control_rx) = + tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); + let cancel = CancellationToken::new(); + + let prefill: HydrantFrame = parse_frame(json!({ + "id": 0, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "prefilrkey001", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })); + frame_tx + .try_send(prefill) + .expect("depth-1 frame channel must accept the prefill"); + + let reader_handle = tokio::spawn(reader_loop( + stream, + frame_tx.clone(), + control_tx.clone(), + cancel.clone(), + )); + + ws_tx + .send(Ok(WsMessage::Text(star_frame_text(1, "starrkeyaa001")))) + .await + .unwrap(); + ws_tx + .send(Ok(WsMessage::Pong(Bytes::new()))) + .await + .unwrap(); + + let pong_event = tokio::time::timeout(Duration::from_millis(500), control_rx.recv()) + .await + .expect("pong must surface inside 500ms even when frame_tx is saturated. A reader blocked on a frame send would never poll the next ws message") + .expect("control_tx was not closed"); + assert!( + matches!(pong_event, WsEvent::IncomingPong), + "first control event must be the pong, not a held text", + ); + + let too_slow = format!( + "{{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}}" + ); + ws_tx.send(Ok(WsMessage::Text(too_slow))).await.unwrap(); + + let end = tokio::time::timeout(Duration::from_secs(1), reader_handle) + .await + .expect("reader must exit promptly once ConsumerTooSlow is read") + .expect("reader task must not panic"); + match end.error { + Some(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "saturated"), + other => panic!("expected ConsumerTooSlow disconnect, got: {other:?}"), + } + + drop(frame_rx); + } + + #[tokio::test] + async fn reader_drains_held_frames_once_processor_catches_up() { + let (ws_tx, ws_rx) = tokio::sync::mpsc::channel::>(8); + let stream: Box = Box::new(ChannelWsStream { rx: ws_rx }); + let (frame_tx, mut frame_rx) = tokio::sync::mpsc::channel::(1); + let (control_tx, _control_rx) = + tokio::sync::mpsc::channel::(CONTROL_CHANNEL_DEPTH); + let cancel = CancellationToken::new(); + + let prefill: HydrantFrame = parse_frame(json!({ + "id": 0, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "prefilrkey001", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })); + frame_tx.try_send(prefill).unwrap(); + + let reader_handle = tokio::spawn(reader_loop( + stream, + frame_tx.clone(), + control_tx, + cancel.clone(), + )); + + ws_tx + .send(Ok(WsMessage::Text(star_frame_text(1, "heldrkeyaa001")))) + .await + .unwrap(); + ws_tx + .send(Ok(WsMessage::Text(star_frame_text(2, "heldrkeyaa002")))) + .await + .unwrap(); + + let _drained_prefill = frame_rx.recv().await.expect("prefilled frame drains"); + let first = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) + .await + .expect("first held frame must reach frame_rx after slot opens") + .expect("frame_tx still open"); + assert_eq!(first.id, 1); + let second = tokio::time::timeout(Duration::from_millis(500), frame_rx.recv()) + .await + .expect("second held frame must reach frame_rx after slot opens") + .expect("frame_tx still open"); + assert_eq!(second.id, 2); + + cancel.cancel(); + let _ = tokio::time::timeout(Duration::from_secs(1), reader_handle) + .await + .expect("reader must stop after cancel"); + } + #[test] fn first_connect_uses_configured_start_cursor() { let start = HydrantCursor::new(42); @@ -952,6 +1582,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -975,6 +1606,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -994,6 +1626,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1069,6 +1702,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1097,6 +1731,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1147,6 +1782,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1200,6 +1836,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1228,6 +1865,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1270,6 +1908,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1315,6 +1954,7 @@ mod tests { &NoopSearchSink, &NoopRecordStore, &resolver, + &sys_clock(), now(), ) .await; @@ -1343,4 +1983,236 @@ mod tests { .expect("cancel must short-circuit the hung ws connect"); assert!(matches!(outcome, Ok(Ok(()))), "got {outcome:?}"); } + + struct CloseOnConnectTransport { + used: std::sync::Mutex, + } + impl bobbin_runtime::WsTransport for CloseOnConnectTransport { + fn connect(&self, _url: Url) -> bobbin_runtime::WsConnectFuture { + let mut used = self.used.lock().unwrap(); + if *used { + return Box::pin(async move { + Err(NetworkError::Connect("only one connect allowed".to_owned())) + }); + } + *used = true; + Box::pin(async move { + let mut q = std::collections::VecDeque::new(); + q.push_back(Ok(WsMessage::Close { + code: 1000, + reason: "bye".to_owned(), + })); + let stream: Box = Box::new(ScriptedWsStream { messages: q }); + struct NoopSink; + impl bobbin_runtime::WsSink for NoopSink { + fn send<'a>( + &'a mut self, + _m: WsMessage, + ) -> bobbin_runtime::WsSendFuture<'a> { + Box::pin(async move { Ok(()) }) + } + } + let sink: Box = Box::new(NoopSink); + Ok(bobbin_runtime::WsConn { sink, stream }) + }) + } + } + + #[tokio::test] + async fn run_session_returns_after_remote_close_when_outer_cancel_unfired() { + let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").unwrap()); + let cancel = CancellationToken::new(); + let mut runtime = fresh_runtime(cancel.clone()); + runtime.ws = Arc::new(CloseOnConnectTransport { + used: std::sync::Mutex::new(false), + }); + + let res = tokio::time::timeout( + Duration::from_secs(3), + run_session(&cfg, HydrantCursor::new(0), &runtime), + ) + .await; + assert!( + res.is_ok(), + "run_session must return after a remote Close even when outer cancel never fires - regression for a session-scoped task hanging on the parent token", + ); + } + + struct HangingSink; + impl bobbin_runtime::WsSink for HangingSink { + fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { + Box::pin(std::future::pending()) + } + } + + #[tokio::test(start_paused = true)] + async fn timed_send_surfaces_send_timeout_when_sink_pends_forever() { + let mut sink: Box = Box::new(HangingSink); + let task = tokio::spawn(async move { + timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await + }); + tokio::time::advance(SEND_TIMEOUT + Duration::from_secs(1)).await; + let result = task.await.expect("task panicked"); + match result { + Err(IngestError::SendTimeout(d)) => assert_eq!(d, SEND_TIMEOUT), + other => panic!( + "expected SendTimeout, got {other:?}. A bare ws_sink.send.await would hang forever on a half-dead socket and starve the writer's pong-deadline arm", + ), + } + } + + struct OkSink; + impl bobbin_runtime::WsSink for OkSink { + fn send<'a>(&'a mut self, _: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { + Box::pin(async move { Ok(()) }) + } + } + + #[tokio::test] + async fn timed_send_returns_ok_when_sink_succeeds_promptly() { + let mut sink: Box = Box::new(OkSink); + let result = timed_send(&mut sink, WsMessage::Ping(Bytes::new())).await; + assert!(matches!(result, Ok(())), "got {result:?}"); + } + + #[test] + fn classify_error_frame_with_id_dispatches_as_error_not_unknown_kind() { + let text = r#"{"id":42,"type":"error","error":"ConsumerTooSlow","message":"slow"}"#; + match classify_text_frame(text) { + Err(IngestError::ConsumerTooSlow { message }) => assert_eq!(message, "slow"), + other => panic!( + "type=\"error\" must dispatch to the error path even when id is present, got: {other:?}" + ), + } + } + + #[tokio::test] + async fn buffered_pipeline_preserves_cursor_order_under_resolve_latency_skew() { + use bobbin_record_lru::RecordStore; + use bobbin_types::record::RecordBody; + use std::sync::Mutex; + + let server = wiremock::MockServer::start().await; + wiremock::Mock::given(wiremock::matchers::method("GET")) + .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) + .respond_with( + wiremock::ResponseTemplate::new(404) + .set_delay(Duration::from_millis(150)), + ) + .mount(&server) + .await; + + let client = bobbin_slingshot_client::SlingshotClient::with_default_http( + Url::parse(&server.uri()).unwrap(), + ) + .unwrap(); + let clock: Arc = Arc::new(SystemClock::new()); + let resolver = Arc::new(RepoIdResolver::with_slingshot( + client, + clock.clone(), + RuntimeHasher::default(), + )); + + let owner = Did::new_owned("did:plc:nel").unwrap(); + let abalone = Did::new_owned("did:plc:abalone").unwrap(); + let fast_rkeys = ["fastrkeyaa01", "fastrkeyaa02"]; + let slow_rkeys = ["slowrkeyaa01", "slowrkeyaa02"]; + for r in &fast_rkeys { + resolver + .observe( + owner.clone(), + Rkey::new_owned(r).unwrap(), + Some(abalone.clone()), + ) + .await; + } + + #[derive(Default)] + struct Capturing { + urls: Mutex>>, + } + impl RecordStore for Capturing { + fn get(&self, _uri: &AtUri) -> Option> { + None + } + fn put(&self, uri: AtUri, _body: Arc) { + self.urls.lock().unwrap().push(uri); + } + fn remove(&self, _uri: &AtUri) {} + } + let capturing = Arc::new(Capturing::default()); + + let runtime: IngestRuntime = IngestRuntime { + store: Arc::new(EdgeStore::new(RuntimeHasher::default())), + coverage: Arc::new(CoverageWatch::new()), + search: Arc::new(NoopSearchSink), + records: capturing.clone() as Arc, + resolver, + clock, + entropy: Arc::new(OsEntropy), + ws: TungsteniteWs::shared(), + cancel: CancellationToken::new(), + }; + + let parallelism = 4usize; + let (frame_tx, frame_rx) = tokio::sync::mpsc::channel::(64); + let pipeline_rt = runtime.clone(); + let pipeline = tokio::spawn(async move { + let prep_rt = pipeline_rt.clone(); + let resolve_rt = pipeline_rt.clone(); + let commit_rt = pipeline_rt; + ReceiverStream::new(frame_rx) + .then(move |frame| prep_stage(frame, prep_rt.clone())) + .map(move |staged| resolve_stage(staged, resolve_rt.clone())) + .buffered(parallelism) + .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)) + .await; + }); + + let mk = |id: u64, idx: u64, repo_rkey: &str| -> HydrantFrame { + parse_frame(json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": format!("starrkeya{idx:04}"), + "action": "create", + "cid": VALID_CID, + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": format!("at://did:plc:nel/sh.tangled.repo/{repo_rkey}"), + } + } + })) + }; + + frame_tx.send(mk(1, 1, slow_rkeys[0])).await.unwrap(); + frame_tx.send(mk(2, 2, fast_rkeys[0])).await.unwrap(); + frame_tx.send(mk(3, 3, slow_rkeys[1])).await.unwrap(); + frame_tx.send(mk(4, 4, fast_rkeys[1])).await.unwrap(); + drop(frame_tx); + + pipeline.await.unwrap(); + + let captured = capturing.urls.lock().unwrap(); + let rkeys: Vec = captured + .iter() + .filter_map(|uri| uri.rkey().map(|r| r.as_ref().to_owned())) + .collect(); + assert_eq!( + rkeys, + vec![ + "starrkeya0001".to_owned(), + "starrkeya0002".to_owned(), + "starrkeya0003".to_owned(), + "starrkeya0004".to_owned(), + ], + "buffered(N) must preserve cursor order even when resolves complete out of order, with ~150ms slow vs cache-hit fast as the latency skew here", + ); + assert_eq!(runtime.coverage.snapshot().last_cursor().raw(), 4); + } } diff --git a/example.toml b/example.toml index 46632ce..7baa285 100644 --- a/example.toml +++ b/example.toml @@ -32,6 +32,17 @@ # Default value: 0 #start_cursor = 0 +[ingest] +# Concurrent in-flight resolves during ingest. One slot per slingshot rtt. +# Default 16 sits at the measured throughput +# knee for cold replay against a healthy hydrant. The committer stays +# serial so that like, cursor and Spur allocation order are preserved. +# +# Can also be specified via environment variable `BOBBIN_INGEST_PARALLELISM`. +# +# Default value: 16 +#parallelism = 16 + [slingshot] # Base URL of a slingshot instance. Used for record bodies and identity. #