From 2c7e10945c7dae2760f55ee4b410bf0c65ebaff2 Mon Sep 17 00:00:00 2001 From: dawn Date: Tue, 18 Aug 2026 01:07:36 +0900 Subject: [PATCH] bobbin/crates/ingest: warm minidoc cache for every did off the stream Signed-off-by: dawn --- bobbin/crates/ingest/src/lib.rs | 232 ++++++++++++++++++++++---------- 1 file changed, 159 insertions(+), 73 deletions(-) diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 695d461c..15c4df15 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -256,6 +256,7 @@ impl IngestRuntime { fn pipeline_ctx(&self) -> PipelineCtx<'_, S> { PipelineCtx { resolver: &self.resolver, + identity: &self.identity, store: &self.store, issue_states: &self.issue_states, pull_statuses: &self.pull_statuses, @@ -273,6 +274,7 @@ impl IngestRuntime { struct PipelineCtx<'a, S: SearchSink + 'static> { resolver: &'a RepoIdResolver, + identity: &'a IdentityResolver, store: &'a EdgeStore, issue_states: &'a StateIndex, pull_statuses: &'a StateIndex, @@ -1028,18 +1030,7 @@ async fn commit_stage( resolve_end, } = staged; let commit_start = rt.clock.now_instant(); - let settled = commit_pending( - pending, - &rt.store, - &rt.issue_states, - &rt.pull_statuses, - &rt.coverage, - &*rt.search, - &*rt.records, - &rt.resolver, - &rt.identity, - ) - .await; + let settled = commit_pending(pending, &rt.pipeline_ctx()).await; let commit_end = rt.clock.now_instant(); if let Some((watch, settled)) = rt.settlements.as_deref().zip(settled) { watch.settle(settled); @@ -1101,6 +1092,9 @@ async fn prepare_record( if !record.collection.as_ref().starts_with(TANGLED_PREFIX) { return PendingOp::Noop; } + // hydrant only announces an identity when it changes, so nothing tells us + // the handle of an author we first meet through a replayed record + ctx.identity.warm(&record.did); let nsid = record.collection.clone(); let source = match build_source_uri(&record) { Ok(s) => s, @@ -1470,18 +1464,13 @@ fn settle_cid(cid: Option<&Cid>) -> Option { }) } -#[allow(clippy::too_many_arguments)] -async fn commit_pending( +async fn commit_pending( pending: Pending, - store: &EdgeStore, - issue_states: &StateIndex, - pull_statuses: &StateIndex, - coverage: &CoverageWatch, - search: &S, - records: &dyn RecordStore, - resolver: &RepoIdResolver, - identity: &IdentityResolver, + ctx: &PipelineCtx<'_, S>, ) -> Option { + let (store, issue_states, pull_statuses) = (ctx.store, ctx.issue_states, ctx.pull_statuses); + let (coverage, search, records) = (ctx.coverage, ctx.search, ctx.records); + let (resolver, identity) = (ctx.resolver, ctx.identity); let Pending { cursor, signal, @@ -1602,8 +1591,10 @@ async fn handle_frame( now: UnixMicros, ) { let _ = clock; + let identity = IdentityResolver::detached(RuntimeHasher::default()); let ctx = PipelineCtx { resolver, + identity: &identity, store, issue_states, pull_statuses, @@ -1618,19 +1609,7 @@ async fn handle_frame( }; let pending = claim_pending(prepare_frame(frame, &ctx, now).await, &ctx).await; let pending = resolve_pending(pending, &ctx).await; - let identity = IdentityResolver::detached(RuntimeHasher::default()); - commit_pending( - pending, - store, - issue_states, - pull_statuses, - coverage, - search, - records, - resolver, - &identity, - ) - .await; + commit_pending(pending, &ctx).await; } fn log_unknown_state_variant(outcome: ApplyOutcome, source: &AtUri) { @@ -1763,6 +1742,7 @@ mod tests { use bobbin_record_lru::{CacheCapacity, LruRecordStore, NoopRecordStore, RecordStore}; use bobbin_runtime::{OsEntropy, RuntimeHasher, SystemClock, TungsteniteWs}; use bobbin_types::search::NoopSearchSink; + use jacquard_common::types::handle::Handle; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::tid::Tid; use serde_json::json; @@ -1806,6 +1786,10 @@ mod tests { SystemClock::new() } + fn detached_identity() -> IdentityResolver { + IdentityResolver::detached(RuntimeHasher::default()) + } + 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") @@ -1876,8 +1860,10 @@ mod tests { let registry = KnotRegistry::new(); let (store, issue_states, pull_statuses, cov, resolver) = fresh(); + let identity = detached_identity(); let ctx = PipelineCtx { resolver: &resolver, + identity: &identity, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, @@ -2067,8 +2053,10 @@ mod tests { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); let search = NoopSearchSink; let records = NoopRecordStore; + let identity = detached_identity(); let ctx = PipelineCtx { resolver: &resolver, + identity: &identity, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, @@ -2703,11 +2691,12 @@ mod tests { #[tokio::test] async fn identity_and_account_frames_drive_the_identity_resolver() { let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); - let identity = IdentityResolver::detached(RuntimeHasher::default()); + let identity = detached_identity(); let records = NoopRecordStore; let search = NoopSearchSink; let ctx = PipelineCtx { resolver: &resolver, + identity: &identity, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, @@ -2724,18 +2713,7 @@ mod tests { let apply = async |frame: HydrantFrame| { let pending = prepare_frame(frame, &ctx, now()).await; let pending = resolve_pending(pending, &ctx).await; - commit_pending( - pending, - &store, - &issue_states, - &pull_statuses, - &coverage, - &search, - &records, - &resolver, - &identity, - ) - .await; + commit_pending(pending, &ctx).await; }; apply(parse_frame(json!({ @@ -2779,6 +2757,135 @@ mod tests { ); } + #[tokio::test] + async fn replayed_records_warm_their_author_once() { + let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); + // a real clock, so the warm queue is live instead of inert + let identity = IdentityResolver::with_slingshot( + bobbin_slingshot_client::SlingshotClient::with_default_http( + url::Url::parse("http://127.0.0.1:1/").unwrap(), + ) + .unwrap(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ); + let records = NoopRecordStore; + let search = NoopSearchSink; + let ctx = PipelineCtx { + resolver: &resolver, + identity: &identity, + store: &store, + issue_states: &issue_states, + pull_statuses: &pull_statuses, + coverage: &coverage, + records: &records, + search: &search, + shadow: None, + buffer: None, + knot_registry: None, + knot_gate: None, + settlements: None, + }; + let star = |id: u64, rkey: &str| { + parse_frame(json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "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", + "subject": "at://did:plc:olaren/sh.tangled.repo/abcabcabcabcz" + } + } + })) + }; + let apply = async |frame: HydrantFrame| { + let pending = prepare_frame(frame, &ctx, now()).await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending(pending, &ctx).await; + }; + + apply(star(1, "abcabcabcabcz")).await; + assert!(store.key_count() > 0, "the record must land in the index"); + assert_eq!( + identity.stats().warm_queued, + 1, + "a replayed record is the only announcement of its author", + ); + + apply(star(2, "abcabcabcabcy")).await; + assert_eq!( + identity.stats().warm_queued, + 1, + "the queue dedups, so a busy author costs one lookup", + ); + } + + #[tokio::test] + async fn records_from_known_authors_queue_no_warm() { + let (store, issue_states, pull_statuses, coverage, resolver) = fresh(); + let identity = IdentityResolver::with_slingshot( + bobbin_slingshot_client::SlingshotClient::with_default_http( + url::Url::parse("http://127.0.0.1:1/").unwrap(), + ) + .unwrap(), + Arc::new(SystemClock::new()), + RuntimeHasher::default(), + ); + identity.observe( + Did::new_static("did:plc:nel").unwrap(), + Handle::new_static("nel.dev").unwrap(), + ); + let records = NoopRecordStore; + let search = NoopSearchSink; + let ctx = PipelineCtx { + resolver: &resolver, + identity: &identity, + store: &store, + issue_states: &issue_states, + pull_statuses: &pull_statuses, + coverage: &coverage, + records: &records, + search: &search, + shadow: None, + buffer: None, + knot_registry: None, + knot_gate: None, + settlements: None, + }; + let pending = prepare_frame( + parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:olaren/sh.tangled.repo/abcabcabcabcz" + } + } + })), + &ctx, + now(), + ) + .await; + let pending = resolve_pending(pending, &ctx).await; + commit_pending(pending, &ctx).await; + assert_eq!(identity.stats().warm_queued, 0); + } + #[tokio::test] async fn account_frame_advances_cursor_only() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); @@ -2959,11 +3066,12 @@ mod tests { #[tokio::test] async fn out_of_order_repo_preps_cannot_steal_a_rename() { let (store, issue_states, pull_statuses, cov, resolver) = fresh(); - let identity = IdentityResolver::detached(RuntimeHasher::default()); + let identity = detached_identity(); let search = NoopSearchSink; let records = NoopRecordStore; let ctx = PipelineCtx { resolver: &resolver, + identity: &identity, store: &store, issue_states: &issue_states, pull_statuses: &pull_statuses, @@ -3002,31 +3110,9 @@ mod tests { let old_pending = claim_pending(old_pending, &ctx).await; let new_pending = claim_pending(new_pending, &ctx).await; let old_pending = resolve_pending(old_pending, &ctx).await; - commit_pending( - old_pending, - &store, - &issue_states, - &pull_statuses, - &cov, - &search, - &records, - &resolver, - &identity, - ) - .await; + commit_pending(old_pending, &ctx).await; let new_pending = resolve_pending(new_pending, &ctx).await; - commit_pending( - new_pending, - &store, - &issue_states, - &pull_statuses, - &cov, - &search, - &records, - &resolver, - &identity, - ) - .await; + commit_pending(new_pending, &ctx).await; let owner = Did::new_owned("did:plc:bnuy").unwrap(); assert_eq!( -- 2.51.2