From 516d85714b89f3cc13f84afd022bc84574a37cab Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 10 May 2026 10:18:27 +0300 Subject: [PATCH] feat(ingest): WarmingBuffer thru pipeline + bin Lewis: May this revision serve well! --- crates/bobbin/src/main.rs | 6 +- crates/ingest/examples/smoke.rs | 1 + crates/ingest/src/lib.rs | 363 ++++++++++++++++++++++++++++---- 3 files changed, 329 insertions(+), 41 deletions(-) diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index 34d129a..b1af496 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -7,7 +7,9 @@ use std::time::Duration; use anyhow::{Context, anyhow}; use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor}; -use bobbin_ingest::{IngestConfig, IngestRuntime, RepoIdResolver, run as run_ingest}; +use bobbin_ingest::{ + IngestConfig, IngestRuntime, RepoIdResolver, WarmingBuffer, run as run_ingest, +}; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; use bobbin_runtime::{Clock, OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; @@ -136,6 +138,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { )); let edges = Arc::new(EdgeStore::new(hasher.clone())); let coverage = Arc::new(CoverageWatch::new()); + let warming_buffer = Arc::new(WarmingBuffer::new(hasher.clone())); let knots = Arc::new(KnotProxy::new( KnotProxyConfig { allow_private_hosts: cfg.knot.allow_private, @@ -171,6 +174,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { cancel: cancel.clone(), disconnects: None, warming_shadow: None, + warming_buffer: Some(warming_buffer), }; let mut ingest_handle = tokio::spawn(run_ingest(ingest_cfg, ingest_runtime)); diff --git a/crates/ingest/examples/smoke.rs b/crates/ingest/examples/smoke.rs index e2ac8d9..e2e3019 100644 --- a/crates/ingest/examples/smoke.rs +++ b/crates/ingest/examples/smoke.rs @@ -43,6 +43,7 @@ async fn main() { cancel: cancel.clone(), disconnects: None, warming_shadow: None, + warming_buffer: None, }; let task = tokio::spawn(async move { let _ = run(cfg, runtime).await; diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 14d3d3a..d5edd97 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -1,14 +1,18 @@ -use std::collections::VecDeque; +use std::collections::{HashSet, VecDeque}; use std::num::NonZeroUsize; use std::sync::Arc; use std::time::Duration; -use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor, PromotionSignal}; +use bobbin_edge_index::{ + Coverage, CoverageWatch, EdgeStore, HydrantCursor, PromotionSignal, +}; use bobbin_record_lru::RecordStore; use bobbin_runtime::{ - Clock, Entropy, NetworkError, UnixMicros, WsConn, WsMessage, WsStream, WsTransport, + Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, + WsTransport, }; use bobbin_types::edges::{Edge, ExtractError, Record}; +use bobbin_types::ids::RepoIdent; use bobbin_types::record::RecordBody; use bobbin_types::search::{SearchSink, SearchableRecord}; use bytes::Bytes; @@ -29,10 +33,12 @@ use url::Url; mod frame; mod resolver; mod shadow; +mod warming; pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; use frame::HydrantStreamErrorFrame; pub use resolver::{RepoIdResolver, Resolution}; pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; +pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; const TANGLED_PREFIX: &str = "sh.tangled."; const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); @@ -46,6 +52,7 @@ 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); +const WARMING_FLUSH_PARALLELISM: usize = 64; pub const DEFAULT_INGEST_PARALLELISM: NonZeroUsize = match NonZeroUsize::new(16) { Some(n) => n, None => unreachable!(), @@ -206,6 +213,7 @@ pub struct IngestRuntime { pub cancel: CancellationToken, pub disconnects: Option>, pub warming_shadow: Option>, + pub warming_buffer: Option>, } impl Clone for IngestRuntime { @@ -222,19 +230,50 @@ impl Clone for IngestRuntime { cancel: self.cancel.clone(), disconnects: self.disconnects.clone(), warming_shadow: self.warming_shadow.clone(), + warming_buffer: self.warming_buffer.clone(), } } } +impl IngestRuntime { + fn pipeline_ctx(&self) -> PipelineCtx<'_, S> { + PipelineCtx { + resolver: &self.resolver, + store: &self.store, + coverage: &self.coverage, + records: &*self.records, + search: &self.search, + shadow: self.warming_shadow.as_deref(), + buffer: self.warming_buffer.as_deref(), + } + } +} + +struct PipelineCtx<'a, S: SearchSink + 'static> { + resolver: &'a RepoIdResolver, + store: &'a EdgeStore, + coverage: &'a CoverageWatch, + records: &'a dyn RecordStore, + search: &'a S, + shadow: Option<&'a WarmingShadowBuffer>, + buffer: Option<&'a WarmingBuffer>, +} + pub async fn run( config: IngestConfig, runtime: IngestRuntime, ) -> Result<(), IngestError> { let metrics_dumper = spawn_metrics_dumper(&runtime); + let warming_flusher = spawn_warming_flusher(&runtime); let result = run_inner(config, &runtime).await; if let Err(join) = metrics_dumper.await { warn!(?join, "metrics dumper task panicked"); } + if let Some(handle) = warming_flusher + && let Err(join) = handle.await + { + warn!(?join, "warming flusher task panicked"); + } result } @@ -284,6 +323,95 @@ async fn run_inner( } } +fn spawn_warming_flusher( + runtime: &IngestRuntime, +) -> Option> { + runtime.warming_buffer.as_ref()?; + let rt = runtime.clone(); + Some(tokio::spawn(async move { + let buffer = rt + .warming_buffer + .as_deref() + .expect("warming flusher only spawns when buffer is set"); + let mut rx = rt.coverage.subscribe(); + let reason = loop { + if rx.borrow_and_update().is_ready() { + break FlushReason::Ready; + } + tokio::select! { + biased; + _ = rt.cancel.cancelled() => break FlushReason::Cancelled, + res = rx.changed() => match res { + Ok(()) => continue, + Err(_) => break FlushReason::CoverageDropped, + }, + } + }; + flush_warming_buffer(&rt, buffer, reason).await; + })) +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum FlushReason { + Ready, + Cancelled, + CoverageDropped, +} + +impl FlushReason { + fn as_str(self) -> &'static str { + match self { + FlushReason::Ready => "ready", + FlushReason::Cancelled => "cancelled", + FlushReason::CoverageDropped => "coverage_dropped", + } + } +} + +async fn flush_warming_buffer( + runtime: &IngestRuntime, + buffer: &WarmingBuffer, + reason: FlushReason, +) { + let drained = buffer.drain_for_promote().await; + if drained.is_empty() { + return; + } + if reason == FlushReason::Ready { + let hasher = buffer.hasher().clone(); + let unique: HashSet = drained + .iter() + .flat_map(|(_, deps)| deps.iter().cloned()) + .fold(HashSet::with_hasher(hasher), |mut acc, dep| { + acc.insert(dep); + acc + }); + if !unique.is_empty() { + let resolver = runtime.resolver.clone(); + let _: Vec<()> = futures::stream::iter(unique) + .map(|key| { + let resolver = resolver.clone(); + async move { + let _ = resolver.resolve(&key.owner, &key.rkey).await; + } + }) + .buffer_unordered(WARMING_FLUSH_PARALLELISM) + .collect() + .await; + } + } + let upserts: Vec = drained.into_iter().map(|(u, _)| u).collect(); + let count = upserts.len(); + let ctx = runtime.pipeline_ctx(); + finalize_drained(&ctx, upserts).await; + info!( + target: "bobbin_ingest::warming", + flushed_entries = count, + reason = reason.as_str(), + "drained warming buffer", + ); +} + fn spawn_metrics_dumper( runtime: &IngestRuntime, ) -> tokio::task::JoinHandle<()> { @@ -641,6 +769,9 @@ enum PendingOp { cid: Option>, edges: Vec, }, + Parked { + nsid: Nsid, + }, Delete { source: AtUri, nsid: Nsid, @@ -665,6 +796,7 @@ fn pending_nsid(op: &PendingOp) -> Option<&Nsid> { match op { PendingOp::Upsert { nsid, .. } => Some(nsid), PendingOp::Delete { nsid, .. } => Some(nsid), + PendingOp::Parked { nsid, .. } => Some(nsid), PendingOp::Noop | PendingOp::ClearCache { .. } => None, } } @@ -682,13 +814,8 @@ async fn prep_stage( ) -> Prepared { let now = rt.clock.now_unix_micros(); let prepare_start = rt.clock.now_instant(); - let pending = prepare_frame( - frame, - &rt.resolver, - rt.warming_shadow.as_deref(), - now, - ) - .await; + let ctx = rt.pipeline_ctx(); + let pending = prepare_frame(frame, &ctx, now).await; let prepare_end = rt.clock.now_instant(); Prepared { pending, @@ -702,13 +829,8 @@ async fn resolve_stage( rt: IngestRuntime, ) -> Resolved { let resolve_start = rt.clock.now_instant(); - let pending = resolve_pending( - staged.pending, - &rt.resolver, - &rt.coverage, - rt.warming_shadow.as_deref(), - ) - .await; + let ctx = rt.pipeline_ctx(); + let pending = resolve_pending(staged.pending, &ctx).await; let resolve_end = rt.clock.now_instant(); Resolved { pending, @@ -755,10 +877,9 @@ async fn commit_stage( ); } -async fn prepare_frame( +async fn prepare_frame( frame: HydrantFrame, - resolver: &RepoIdResolver, - shadow: Option<&WarmingShadowBuffer>, + ctx: &PipelineCtx<'_, S>, now: UnixMicros, ) -> Pending { let cursor = HydrantCursor::new(frame.id); @@ -769,7 +890,7 @@ async fn prepare_frame( None => Regime::NonRecord, }; let op = match frame.kind { - FrameKind::Record => prepare_record(frame.record, resolver, shadow).await, + FrameKind::Record => prepare_record(frame.record, ctx).await, FrameKind::Identity | FrameKind::Account => PendingOp::Noop, FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); @@ -784,10 +905,9 @@ async fn prepare_frame( } } -async fn prepare_record( +async fn prepare_record( record: Option, - resolver: &RepoIdResolver, - shadow: Option<&WarmingShadowBuffer>, + ctx: &PipelineCtx<'_, S>, ) -> PendingOp { let Some(record) = record else { debug!("record-typed frame missing payload, skipping"); @@ -806,6 +926,7 @@ async fn prepare_record( }; match record.action { RecordAction::Create | RecordAction::Update => { + evict_from_buffer(ctx.buffer, &source).await; let Some(raw) = record.record else { debug!(collection = %nsid, "create/update missing record body, clearing cache"); return PendingOp::ClearCache { source }; @@ -823,16 +944,22 @@ async fn prepare_record( } }; if let Record::Repo(repo) = &parsed { - if let Some(shadow) = shadow { + if let Some(shadow) = ctx.shadow { shadow.note_observed(&record.did, &record.rkey).await; } - resolver + ctx.resolver .observe( record.did.clone(), record.rkey.clone(), repo.repo_did.clone(), ) .await; + if let Some(buffer) = ctx.buffer { + let drained = buffer.take_observed(&record.did, &record.rkey).await; + if !drained.is_empty() { + finalize_drained(ctx, drained).await; + } + } } let edges = match parsed.extract_edges(&source) { Ok(es) => es, @@ -841,6 +968,7 @@ async fn prepare_record( return PendingOp::ClearCache { source }; } }; + let _ = ctx.store.intern_source(&source); PendingOp::Upsert { source, nsid, @@ -850,7 +978,10 @@ async fn prepare_record( edges, } } - RecordAction::Delete => PendingOp::Delete { source, nsid }, + RecordAction::Delete => { + evict_from_buffer(ctx.buffer, &source).await; + PendingOp::Delete { source, nsid } + } RecordAction::Other => { debug!(collection = %nsid, "ignoring unknown record action"); PendingOp::Noop @@ -858,11 +989,17 @@ async fn prepare_record( } } -async fn resolve_pending( +async fn evict_from_buffer(buffer: Option<&WarmingBuffer>, source: &AtUri) { + if let Some(buffer) = buffer + && !buffer.is_sealed() + { + buffer.evict_source(source).await; + } +} + +async fn resolve_pending( pending: Pending, - resolver: &RepoIdResolver, - coverage: &CoverageWatch, - shadow: Option<&WarmingShadowBuffer>, + ctx: &PipelineCtx<'_, S>, ) -> Pending { let Pending { cursor, @@ -879,7 +1016,20 @@ async fn resolve_pending( cid, edges, } => { - let edges = normalize_subjects(edges, resolver, coverage, shadow).await; + let pieces = UpsertPieces { source, nsid, parsed, bytes, cid, edges }; + let pieces = match try_park_warming(ctx, cursor, pieces).await { + ParkOutcome::Parked { nsid } => { + return Pending { + cursor, + signal, + regime, + op: PendingOp::Parked { nsid }, + }; + } + ParkOutcome::Passthrough(pieces) => *pieces, + }; + let UpsertPieces { source, nsid, parsed, bytes, cid, edges } = pieces; + let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, ctx.shadow).await; PendingOp::Upsert { source, nsid, @@ -899,6 +1049,117 @@ async fn resolve_pending( } } +struct UpsertPieces { + source: AtUri, + nsid: Nsid, + parsed: Record, + bytes: Bytes, + cid: Option>, + edges: Vec, +} + +impl From for UpsertPieces { + fn from(u: ParkedUpsert) -> Self { + Self { + source: u.source, + nsid: u.nsid, + parsed: u.parsed, + bytes: u.bytes, + cid: u.cid, + edges: u.edges, + } + } +} + +enum ParkOutcome { + Parked { nsid: Nsid }, + Passthrough(Box), +} + +async fn try_park_warming( + ctx: &PipelineCtx<'_, S>, + cursor: HydrantCursor, + pieces: UpsertPieces, +) -> ParkOutcome { + let Some(buffer) = ctx.buffer else { + return ParkOutcome::Passthrough(Box::new(pieces)); + }; + if ctx.coverage.snapshot().is_ready() || buffer.is_sealed() { + return ParkOutcome::Passthrough(Box::new(pieces)); + } + let deps = collect_unresolved_deps(&pieces.edges, ctx.resolver).await; + if deps.is_empty() { + return ParkOutcome::Passthrough(Box::new(pieces)); + } + let nsid = pieces.nsid.clone(); + let upsert = ParkedUpsert { + cursor, + source: pieces.source, + nsid: pieces.nsid, + parsed: pieces.parsed, + bytes: pieces.bytes, + cid: pieces.cid, + edges: pieces.edges, + }; + let deps_for_shadow = ctx.shadow.is_some().then(|| deps.clone()); + match buffer.try_park(upsert, deps).await { + Ok(()) => { + if let Some((shadow, noted)) = ctx.shadow.zip(deps_for_shadow) { + let _ = futures::future::join_all(noted.into_iter().map(|dep| async move { + shadow.note_unresolved(dep.owner, dep.rkey).await; + })) + .await; + } + ParkOutcome::Parked { nsid } + } + Err(returned) => ParkOutcome::Passthrough(Box::new(returned.into())), + } +} + +async fn collect_unresolved_deps( + edges: &[Edge], + resolver: &RepoIdResolver, +) -> Vec { + futures::stream::iter(edges) + .fold(Vec::new(), |mut acc, edge| async move { + let Some((owner, rkey)) = parse_repo_subject(&edge.subject) else { + return acc; + }; + if resolver.cached_resolution(&owner, &rkey).await.is_some() { + return acc; + } + let candidate = RepoIdent::new(owner, rkey); + if !acc.contains(&candidate) { + acc.push(candidate); + } + acc + }) + .await +} + +async fn finalize_drained( + ctx: &PipelineCtx<'_, S>, + drained: Vec, +) { + for upsert in drained { + let ParkedUpsert { + cursor: _, + source, + nsid: _, + parsed, + bytes, + cid, + edges, + } = upsert; + let edges = normalize_subjects(edges, ctx.resolver, ctx.coverage, None).await; + cache_body(ctx.records, &source, cid, bytes); + ctx.store.upsert_source(&source, edges); + if let Some(searchable) = SearchableRecord::try_from_record(parsed) { + ctx.search.upsert(searchable.to_search_doc(&source)).await; + } + } +} + async fn commit_pending( pending: Pending, store: &EdgeStore, @@ -913,7 +1174,7 @@ async fn commit_pending( op, } = pending; match op { - PendingOp::Noop => {} + PendingOp::Noop | PendingOp::Parked { .. } => {} PendingOp::ClearCache { source } => records.remove(&source), PendingOp::Upsert { source, @@ -939,7 +1200,7 @@ async fn commit_pending( } #[cfg(test)] -async fn handle_frame( +async fn handle_frame( frame: HydrantFrame, store: &EdgeStore, coverage: &CoverageWatch, @@ -950,8 +1211,17 @@ async fn handle_frame( now: UnixMicros, ) { let _ = clock; - let pending = prepare_frame(frame, resolver, None, now).await; - let pending = resolve_pending(pending, resolver, coverage, None).await; + let ctx = PipelineCtx { + resolver, + store, + coverage, + records, + search, + shadow: None, + buffer: None, + }; + let pending = prepare_frame(frame, &ctx, now).await; + let pending = resolve_pending(pending, &ctx).await; commit_pending(pending, store, coverage, search, records).await; } @@ -1233,7 +1503,18 @@ mod tests { #[tokio::test] async fn prepare_frame_tags_regime_from_live_flag() { - let (_store, _cov, resolver) = fresh(); + let (store, cov, resolver) = fresh(); + let search = NoopSearchSink; + let records = NoopRecordStore; + let ctx = PipelineCtx { + resolver: &resolver, + store: &store, + coverage: &cov, + records: &records, + search: &search, + shadow: None, + buffer: None, + }; let mk = |live: bool| -> HydrantFrame { parse_frame(json!({ "id": 1, @@ -1253,16 +1534,16 @@ mod tests { } })) }; - let live_pending = prepare_frame(mk(true), &resolver, None, now()).await; + let live_pending = prepare_frame(mk(true), &ctx, now()).await; assert_eq!(live_pending.regime, Regime::Live); - let replay_pending = prepare_frame(mk(false), &resolver, None, now()).await; + let replay_pending = prepare_frame(mk(false), &ctx, now()).await; assert_eq!(replay_pending.regime, Regime::Replay); let identity: HydrantFrame = parse_frame(json!({ "id": 9, "type": "identity", })); - let id_pending = prepare_frame(identity, &resolver, None, now()).await; + let id_pending = prepare_frame(identity, &ctx, now()).await; assert_eq!(id_pending.regime, Regime::NonRecord); } @@ -1844,6 +2125,7 @@ mod tests { cancel, disconnects: None, warming_shadow: None, + warming_buffer: None, } } @@ -2354,6 +2636,7 @@ mod tests { cancel: CancellationToken::new(), disconnects: None, warming_shadow: None, + warming_buffer: None, }; let parallelism = 4usize; -- 2.51.2