diff --git a/crates/edge-index/src/coverage.rs b/crates/edge-index/src/coverage.rs index 82e4aee..93f5635 100644 --- a/crates/edge-index/src/coverage.rs +++ b/crates/edge-index/src/coverage.rs @@ -42,7 +42,6 @@ impl Default for Coverage { #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct PromotionSignal { - pub live: bool, pub rev_micros: Option, pub now_micros: u64, pub skew_micros: u64, @@ -53,7 +52,7 @@ impl PromotionSignal { let Some(rev) = self.rev_micros else { return false; }; - self.live && self.now_micros.abs_diff(rev) <= self.skew_micros + self.now_micros.abs_diff(rev) <= self.skew_micros } } @@ -112,6 +111,19 @@ impl Coverage { warming => warming, } } + + pub const fn force_ready(self) -> Self { + match self { + Self::Ready { .. } => self, + Self::Warming { + events_processed, + last_cursor, + } => Self::Ready { + events_processed, + last_cursor, + }, + } + } } #[derive(Debug)] @@ -161,9 +173,8 @@ mod tests { const SKEW: u64 = 60_000_000; - fn signal(live: bool, rev: u64, now: u64) -> PromotionSignal { + fn signal(rev: u64, now: u64) -> PromotionSignal { PromotionSignal { - live, rev_micros: Some(rev), now_micros: now, skew_micros: SKEW, @@ -189,15 +200,14 @@ mod tests { } #[test] - fn promotion_requires_live_and_recent_rev() { + fn promotion_requires_recent_rev() { let now = 1_000_000_000; let recent = now - SKEW / 2; let stale = now - SKEW * 10; let c = Coverage::default().advance(HydrantCursor::new(1)); - assert!(!c.maybe_promote(signal(false, recent, now)).is_ready()); - assert!(!c.maybe_promote(signal(true, stale, now)).is_ready()); - assert!(c.maybe_promote(signal(true, recent, now)).is_ready()); + assert!(!c.maybe_promote(signal(stale, now)).is_ready()); + assert!(c.maybe_promote(signal(recent, now)).is_ready()); } #[test] @@ -205,7 +215,7 @@ mod tests { let now = 1_000_000_000; let c = Coverage::default() .advance(HydrantCursor::new(3)) - .maybe_promote(signal(true, now, now)) + .maybe_promote(signal(now, now)) .advance(HydrantCursor::new(4)); assert!(c.is_ready()); assert_eq!(c.events_processed(), 2); @@ -218,9 +228,9 @@ mod tests { let stale = now - SKEW * 10; let c = Coverage::default() .advance(HydrantCursor::new(1)) - .maybe_promote(signal(true, now, now)) + .maybe_promote(signal(now, now)) .advance(HydrantCursor::new(2)) - .maybe_promote(signal(false, stale, now)); + .maybe_promote(signal(stale, now)); assert!(c.is_ready()); } @@ -230,7 +240,7 @@ mod tests { let near_future = now + SKEW / 2; let c = Coverage::default() .advance(HydrantCursor::new(1)) - .maybe_promote(signal(true, near_future, now)); + .maybe_promote(signal(near_future, now)); assert!(c.is_ready()); } @@ -240,7 +250,7 @@ mod tests { let far_future = now + SKEW * 10; let c = Coverage::default() .advance(HydrantCursor::new(1)) - .maybe_promote(signal(true, far_future, now)); + .maybe_promote(signal(far_future, now)); assert!(!c.is_ready()); } @@ -248,7 +258,6 @@ mod tests { fn missing_rev_does_not_promote() { let c = Coverage::default().advance(HydrantCursor::new(1)); let s = PromotionSignal { - live: true, rev_micros: None, now_micros: 1_000_000_000, skew_micros: SKEW, diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index d5edd97..d898b9d 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -264,11 +264,15 @@ pub async fn run( runtime: IngestRuntime, ) -> Result<(), IngestError> { let metrics_dumper = spawn_metrics_dumper(&runtime); + let idle_promoter = spawn_idle_promoter(&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 Err(join) = idle_promoter.await { + warn!(?join, "idle promoter task panicked"); + } if let Some(handle) = warming_flusher && let Err(join) = handle.await { @@ -412,6 +416,41 @@ async fn flush_warming_buffer( ); } +const IDLE_PROMOTE_WINDOW: Duration = Duration::from_secs(15); +const IDLE_PROMOTE_MIN_EVENTS: u64 = 256; + +fn spawn_idle_promoter( + runtime: &IngestRuntime, +) -> tokio::task::JoinHandle<()> { + let rt = runtime.clone(); + tokio::spawn(async move { + let mut prev = rt.coverage.snapshot().events_processed(); + loop { + tokio::select! { + biased; + _ = rt.cancel.cancelled() => return, + _ = rt.clock.sleep(IDLE_PROMOTE_WINDOW) => {} + } + let snap = rt.coverage.snapshot(); + if snap.is_ready() { + return; + } + let processed = snap.events_processed(); + if processed >= IDLE_PROMOTE_MIN_EVENTS && processed == prev { + rt.coverage.update(|c| c.force_ready()); + info!( + target: "bobbin_ingest::coverage", + events_processed = processed, + last_cursor = snap.last_cursor().raw(), + "stream idle, promoting coverage to ready", + ); + return; + } + prev = processed; + } + }) +} + fn spawn_metrics_dumper( runtime: &IngestRuntime, ) -> tokio::task::JoinHandle<()> { @@ -1227,7 +1266,6 @@ async fn handle_frame( 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, @@ -1261,9 +1299,9 @@ async fn normalize_subjects( ) -> Vec { let warming = shadow.is_some() && !coverage.snapshot().is_ready(); futures::stream::iter(edges) - .then(|edge| async move { + .filter_map(|edge| async move { let Some((owner, rkey)) = parse_repo_subject(&edge.subject) else { - return edge; + return Some(edge); }; if warming && let Some(shadow) = shadow @@ -1274,11 +1312,22 @@ async fn normalize_subjects( .await; } match resolver.resolve(&owner, &rkey).await { - Resolution::Mapped(repo_did) => Edge { + Resolution::Mapped(repo_did) => Some(Edge { subject: owned_did_aturi(&repo_did), ..edge - }, - Resolution::NoRepoDid | Resolution::Unresolvable => edge, + }), + Resolution::NoRepoDid => { + warn!( + target: "bobbin_ingest::normalize", + kind = %edge.kind, + owner = owner.as_ref(), + rkey = rkey.as_ref(), + source = edge.source.as_ref(), + "dropping edge: target repo has no repoDid, no canonical DID subject available", + ); + None + } + Resolution::Unresolvable => Some(edge), } }) .collect() @@ -2291,7 +2340,7 @@ mod tests { } #[tokio::test] - async fn repo_without_repo_did_keeps_full_uri_subject() { + async fn repo_without_repo_did_drops_edge() { let (store, cov, resolver) = fresh(); let repo: HydrantFrame = parse_frame(json!({ "id": 1, @@ -2352,14 +2401,22 @@ mod tests { ) .await; - let key = bobbin_types::ids::EdgeKey::new( - Nsid::new_static("sh.tangled.feed.star").unwrap(), + let nsid = Nsid::new_static("sh.tangled.feed.star").unwrap(); + let uri_keyed = bobbin_types::ids::EdgeKey::new( + nsid.clone(), AtUri::new_owned("at://did:plc:nel/sh.tangled.repo/abcabcabcabcz").unwrap(), ); + let owner_keyed = + bobbin_types::ids::EdgeKey::new(nsid, AtUri::new_owned("at://did:plc:nel").unwrap()); assert_eq!( - store.count(&key), - 1, - "the at-uri is the only canonical identity for a repo with no repoDID", + store.count(&uri_keyed), + 0, + "no canonical DID exists for a repo without repoDID, so the edge must be dropped", + ); + assert_eq!( + store.count(&owner_keyed), + 0, + "the authoring DID is not the canonical repo identity", ); }