diff --git a/crates/bobbin-sim/src/determinism.rs b/crates/bobbin-sim/src/determinism.rs index 34126b7..e153473 100644 --- a/crates/bobbin-sim/src/determinism.rs +++ b/crates/bobbin-sim/src/determinism.rs @@ -110,10 +110,7 @@ impl LeakRunResult { } } -pub fn run_leak_check( - config: LeakRunConfig, - workload_factory: F, -) -> LeakRunResult +pub fn run_leak_check(config: LeakRunConfig, workload_factory: F) -> LeakRunResult where F: Fn() -> Box, { @@ -125,12 +122,7 @@ where let (first_report, first_lines) = run_once(&config, workload_factory()); let (second_report, second_lines) = run_once(&config, workload_factory()); - let outcome = compare_runs( - &first_report, - &first_lines, - &second_report, - &second_lines, - ); + let outcome = compare_runs(&first_report, &first_lines, &second_report, &second_lines); LeakRunResult { config, @@ -141,10 +133,7 @@ where } } -fn run_once( - config: &LeakRunConfig, - workload: Box, -) -> (SimReport, Vec) { +fn run_once(config: &LeakRunConfig, workload: Box) -> (SimReport, Vec) { let capture = TraceCapture::new(); let subscriber = Registry::default().with(capture.layer()); diff --git a/crates/bobbin-sim/src/main.rs b/crates/bobbin-sim/src/main.rs index a264249..7657b92 100644 --- a/crates/bobbin-sim/src/main.rs +++ b/crates/bobbin-sim/src/main.rs @@ -5,13 +5,12 @@ use std::time::Duration; use bobbin_sim::workloads::{ CancelMidHydration, CancelMidHydrationConfig, ColdStartUnderLiveLoad, - ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, - ConcurrentReadsDuringReplayConfig, FrameBurst, HydrantDisconnectBarrage, - HydrantDisconnectBarrageConfig, SlingshotFlap, SlingshotFlapConfig, + ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, ConcurrentReadsDuringReplayConfig, + FrameBurst, HydrantDisconnectBarrage, HydrantDisconnectBarrageConfig, SlingshotFlap, + SlingshotFlapConfig, }; use bobbin_sim::{ - LeakOutcome, LeakRunConfig, Sim, SimConfig, SimOutcome, SimReport, Workload, - run_leak_check, + LeakOutcome, LeakRunConfig, Sim, SimConfig, SimOutcome, SimReport, Workload, run_leak_check, }; use clap::{Parser, ValueEnum}; @@ -230,9 +229,9 @@ fn build_workload_from(snap: &CliSnapshot) -> Box { omit_target_repos: false, emit_live_promotion_frame: false, })), - WorkloadName::CancelMidHydration => Box::new(CancelMidHydration::new( - CancelMidHydrationConfig::default(), - )), + WorkloadName::CancelMidHydration => { + Box::new(CancelMidHydration::new(CancelMidHydrationConfig::default())) + } WorkloadName::HydrantDisconnectBarrage => Box::new(HydrantDisconnectBarrage::new( HydrantDisconnectBarrageConfig::default(), )), diff --git a/crates/bobbin-sim/src/runtime.rs b/crates/bobbin-sim/src/runtime.rs index 037c23e..b1291d2 100644 --- a/crates/bobbin-sim/src/runtime.rs +++ b/crates/bobbin-sim/src/runtime.rs @@ -10,8 +10,8 @@ use bobbin_ingest::{ }; use bobbin_record_lru::NoopRecordStore; use bobbin_runtime::{ - Clock, DEFAULT_MEM_WS_CAPACITY, MemHttpTransport, MemWsTransport, RuntimeHasher, - SeededEntropy, SimClock, UnixMicros, + Clock, DEFAULT_MEM_WS_CAPACITY, MemHttpTransport, MemWsTransport, RuntimeHasher, SeededEntropy, + SimClock, UnixMicros, }; use bobbin_slingshot_client::SlingshotClient; use bobbin_types::search::NoopSearchSink; @@ -96,8 +96,7 @@ impl Sim { }; let hooks = self.workload.build(ctx); - let slingshot_http = - MemHttpTransport::shared(hooks.slingshot.clone(), clock.clone()); + let slingshot_http = MemHttpTransport::shared(hooks.slingshot.clone(), clock.clone()); let slingshot_client = SlingshotClient::new(slingshot_base, slingshot_http) .expect("slingshot base url is valid"); let resolver = Arc::new(RepoIdResolver::with_slingshot( diff --git a/crates/bobbin-sim/src/workload.rs b/crates/bobbin-sim/src/workload.rs index 64c3ef3..2d51b82 100644 --- a/crates/bobbin-sim/src/workload.rs +++ b/crates/bobbin-sim/src/workload.rs @@ -21,8 +21,7 @@ pub struct WorkloadCtx { pub consumer_too_slow_count: Arc, } -pub type WorkloadScript = - Pin + Send + 'static>>; +pub type WorkloadScript = Pin + Send + 'static>>; pub struct WorkloadHooks { pub slingshot: Arc, diff --git a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs index 87cb372..6c885da 100644 --- a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs +++ b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs @@ -104,23 +104,18 @@ impl Workload for CancelMidHydration { clock.sleep(Duration::from_secs(2)).await; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let slingshot_after_drain = slingshot_calls_probe.load(Ordering::Relaxed); - let post_observe_events = - snap.events_processed().saturating_sub(events_after_observe); + let post_observe_events = snap.events_processed().saturating_sub(events_after_observe); let post_observe_slingshot = slingshot_after_drain.saturating_sub(slingshot_after_observe); - let outcome = if !matches!(cancel_outcome, SimOutcome::Passed) { - SimOutcome::Failed - } else if post_observe_slingshot != 0 { - SimOutcome::Failed - } else if post_observe_events != 0 { - SimOutcome::Failed - } else if snap.events_processed() > frames_total { - SimOutcome::Failed - } else if snap.events_processed() < cancel_at { + let failed = !matches!(cancel_outcome, SimOutcome::Passed) + || post_observe_slingshot != 0 + || post_observe_events != 0 + || snap.events_processed() > frames_total + || snap.events_processed() < cancel_at; + let outcome = if failed { SimOutcome::Failed } else { SimOutcome::Passed diff --git a/crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs b/crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs index ef6e2c2..2c3f902 100644 --- a/crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs +++ b/crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs @@ -2,9 +2,7 @@ use std::sync::Arc; use std::sync::Mutex; use std::time::Duration; -use bobbin_runtime::{ - Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage, -}; +use bobbin_runtime::{Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; use tokio::sync::mpsc; use url::Url; @@ -89,9 +87,8 @@ impl Workload for ColdStartUnderLiveLoad { } }; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let stray_slingshot = slingshot_probe.calls(); let outcome = match outcome { SimOutcome::Passed if stray_slingshot != 0 => SimOutcome::Failed, diff --git a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs index 42a8ffa..b328958 100644 --- a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs +++ b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs @@ -4,9 +4,7 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; use bobbin_edge_index::{PageCursor, PageLimit}; -use bobbin_runtime::{ - Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage, -}; +use bobbin_runtime::{Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; use bobbin_types::ids::{EdgeKey, SubjectRef}; use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; @@ -134,9 +132,8 @@ impl Workload for ConcurrentReadsDuringReplay { let _ = reader_task.await; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let violations = monotonicity_violations.load(Ordering::Relaxed); let observed = max_observed.load(Ordering::Relaxed); let iters = read_iterations.load(Ordering::Relaxed); diff --git a/crates/bobbin-sim/src/workloads/frame_burst.rs b/crates/bobbin-sim/src/workloads/frame_burst.rs index e314f6c..807d398 100644 --- a/crates/bobbin-sim/src/workloads/frame_burst.rs +++ b/crates/bobbin-sim/src/workloads/frame_burst.rs @@ -69,9 +69,8 @@ impl Workload for FrameBurst { } }; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let stray_slingshot = slingshot_probe.calls(); let session_total = sessions.load(Ordering::Relaxed); let outcome = match outcome { diff --git a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs index 19050ae..3545630 100644 --- a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs +++ b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs @@ -3,9 +3,7 @@ use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; -use bobbin_runtime::{ - MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage, -}; +use bobbin_runtime::{MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage}; use tokio::sync::mpsc; use url::Url; @@ -94,9 +92,8 @@ impl Workload for HydrantDisconnectBarrage { } }; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let dc = disconnects.load(Ordering::Relaxed); let ss = sessions.load(Ordering::Relaxed); let stray_slingshot = slingshot_probe.calls(); diff --git a/crates/bobbin-sim/src/workloads/slingshot_flap.rs b/crates/bobbin-sim/src/workloads/slingshot_flap.rs index f0130de..e82f0f5 100644 --- a/crates/bobbin-sim/src/workloads/slingshot_flap.rs +++ b/crates/bobbin-sim/src/workloads/slingshot_flap.rs @@ -125,9 +125,8 @@ impl Workload for SlingshotFlap { } }; let snap = coverage.snapshot(); - let virtual_runtime = Duration::from_micros( - clock.now_unix_micros().raw().saturating_sub(started.raw()), - ); + let virtual_runtime = + Duration::from_micros(clock.now_unix_micros().raw().saturating_sub(started.raw())); let cts = cts_counter.load(Ordering::Relaxed); let sessions = sessions_counter.load(Ordering::Relaxed); let failure_reason = match outcome { @@ -229,7 +228,12 @@ fn build_frame_script(cfg: &SlingshotFlapConfig, started_unix: UnixMicros) -> Ve } } if cfg.emit_live_promotion_frame { - let prior = cfg.cross_did_stars + if cfg.omit_target_repos { 0 } else { cfg.cross_did_stars }; + let prior = cfg.cross_did_stars + + if cfg.omit_target_repos { + 0 + } else { + cfg.cross_did_stars + }; let id = (prior + 1) as u64; let promoter_name = format!("periwinkle-{prior}"); let promoter_did = format!("did:plc:{promoter_name}"); diff --git a/crates/bobbin-sim/tests/determinism_leak.rs b/crates/bobbin-sim/tests/determinism_leak.rs index 73106ff..738feef 100644 --- a/crates/bobbin-sim/tests/determinism_leak.rs +++ b/crates/bobbin-sim/tests/determinism_leak.rs @@ -3,9 +3,9 @@ use std::time::Duration; use bobbin_sim::workloads::{ CancelMidHydration, CancelMidHydrationConfig, ColdStartUnderLiveLoad, - ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, - ConcurrentReadsDuringReplayConfig, FrameBurst, HydrantDisconnectBarrage, - HydrantDisconnectBarrageConfig, SlingshotFlap, SlingshotFlapConfig, + ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, ConcurrentReadsDuringReplayConfig, + FrameBurst, HydrantDisconnectBarrage, HydrantDisconnectBarrageConfig, SlingshotFlap, + SlingshotFlapConfig, }; use bobbin_sim::{LeakOutcome, LeakRunConfig, Workload, run_leak_check}; @@ -90,10 +90,12 @@ fn hydrant_disconnect_barrage_is_byte_deterministic() { warming_buffer_enabled: true, }; let factory = || -> Box { - Box::new(HydrantDisconnectBarrage::new(HydrantDisconnectBarrageConfig { - frames: 256, - disconnect_after_frames_per_session: vec![50, 80, 60], - })) + Box::new(HydrantDisconnectBarrage::new( + HydrantDisconnectBarrageConfig { + frames: 256, + disconnect_after_frames_per_session: vec![50, 80, 60], + }, + )) }; let result = run_leak_check(config.clone(), factory); assert!( diff --git a/crates/bobbin-sim/tests/warming_buffer.rs b/crates/bobbin-sim/tests/warming_buffer.rs index 8bf2a01..4ecde47 100644 --- a/crates/bobbin-sim/tests/warming_buffer.rs +++ b/crates/bobbin-sim/tests/warming_buffer.rs @@ -43,8 +43,7 @@ fn lever_b_absorbs_sustained_brownout_without_disconnect() { assert_eq!( report.disconnect_count, 0, "expected zero disconnects under sustained brownout once Lever B parks the cohort, got {} (last={:?})", - report.disconnect_count, - report.last_disconnect, + report.disconnect_count, report.last_disconnect, ); let b = report.warming_buffer; @@ -93,7 +92,12 @@ fn shadow_and_buffer_observe_the_same_population() { .expect("build current_thread runtime"); let report = runtime.block_on(Sim::new(sim_config, workload).run()); - assert_eq!(report.outcome, SimOutcome::Passed, "{:?}", report.failure_reason); + assert_eq!( + report.outcome, + SimOutcome::Passed, + "{:?}", + report.failure_reason + ); let s = report.warming_shadow; let b = report.warming_buffer; @@ -175,7 +179,8 @@ fn warming_to_ready_promote_drains_residual_via_parallel_slingshot_wave() { report.resolver_misses, ); assert_eq!( - report.edge_count, stars * 2 + 1, + report.edge_count, + stars * 2 + 1, "each drained star contributes a primary edge plus a sh.tangled.feed.star.by mirror edge; promoter repo adds one. got {} for {stars} stars", report.edge_count, ); @@ -213,8 +218,18 @@ fn buffer_enabled_and_disabled_produce_identical_edge_index() { let with_buffer = run_with_buffer(true); let without_buffer = run_with_buffer(false); - assert_eq!(with_buffer.outcome, SimOutcome::Passed, "{:?}", with_buffer.failure_reason); - assert_eq!(without_buffer.outcome, SimOutcome::Passed, "{:?}", without_buffer.failure_reason); + assert_eq!( + with_buffer.outcome, + SimOutcome::Passed, + "{:?}", + with_buffer.failure_reason + ); + assert_eq!( + without_buffer.outcome, + SimOutcome::Passed, + "{:?}", + without_buffer.failure_reason + ); assert_eq!( with_buffer.events_processed, without_buffer.events_processed, "events_processed must match across buffer modes", diff --git a/crates/edge-index/src/lib.rs b/crates/edge-index/src/lib.rs index 29a231c..f188173 100644 --- a/crates/edge-index/src/lib.rs +++ b/crates/edge-index/src/lib.rs @@ -236,16 +236,18 @@ impl EdgeStore { let Some((_, entries)) = self.reverse.remove_sync(&id) else { return; }; - entries.into_iter().for_each(|ReverseEntry { key, sort_micros }| { - self.forward.update_sync(&key, |_, bucket| { - bucket.sources.remove(&(sort_micros, id.raw())); - if let Some(author) = author { - drop_ref(&mut bucket.author_refs, author); - } + entries + .into_iter() + .for_each(|ReverseEntry { key, sort_micros }| { + self.forward.update_sync(&key, |_, bucket| { + bucket.sources.remove(&(sort_micros, id.raw())); + if let Some(author) = author { + drop_ref(&mut bucket.author_refs, author); + } + }); + self.forward + .remove_if_sync(&key, |bucket| bucket.sources.is_empty()); }); - self.forward - .remove_if_sync(&key, |bucket| bucket.sources.is_empty()); - }); } fn intern_author(&self, source: &str) -> Option { @@ -589,7 +591,13 @@ mod tests { #[test] fn cursor_decode_rejects_malformed() { - let bad = ["", "deadbeef0", "no-hex!!", "1234567", "way-too-long-not-a-tid"]; + let bad = [ + "", + "deadbeef0", + "no-hex!!", + "1234567", + "way-too-long-not-a-tid", + ]; bad.into_iter().for_each(|s| { assert!( matches!(PageToken::decode_token(s), Err(CursorParseError::Malformed)), diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index a1fe7d6..d4e5759 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -3,15 +3,13 @@ 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_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes}; use bobbin_runtime::{ Clock, Entropy, NetworkError, RuntimeHasher, UnixMicros, WsConn, WsMessage, WsStream, WsTransport, }; -use bobbin_resolver::{NormalizeRepoRefs, decode_canon_or_upgrade_bytes}; use bobbin_types::edges::{Edge, ExtractError, Record}; use bobbin_types::ids::{RepoIdent, SubjectRef}; use bobbin_types::record::RecordBody; @@ -35,8 +33,8 @@ mod frame; mod resolver; mod shadow; mod warming; -pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; use frame::HydrantStreamErrorFrame; +pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; pub use resolver::{RepoIdResolver, Resolution}; pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; pub use warming::{ParkedUpsert, WarmingBuffer, WarmingBufferSnapshot}; @@ -170,10 +168,7 @@ impl DisconnectSink { } pub fn record(&self, snap: DisconnectSnapshot) { - *self - .last - .lock() - .expect("disconnect sink mutex poisoned") = Some(snap); + *self.last.lock().expect("disconnect sink mutex poisoned") = Some(snap); self.count .fetch_add(1, std::sync::atomic::Ordering::Relaxed); } @@ -382,28 +377,35 @@ async fn flush_warming_buffer( 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; - } + if reason != FlushReason::Ready { + info!( + target: "bobbin_ingest::warming", + abandoned_entries = drained.len(), + reason = reason.as_str(), + "abandoning parked items on non-ready flush", + ); + return; + } + 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(); @@ -541,7 +543,8 @@ async fn run_session( let commit_rt = processor_runtime; let pipeline = ReceiverStream::new(frame_rx) - .then(move |frame| prep_stage(frame, prep_rt.clone())) + .map(move |frame| prep_stage(frame, prep_rt.clone())) + .buffered(parallelism) .map(move |staged| resolve_stage(staged, resolve_rt.clone())) .buffered(parallelism) .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)); @@ -683,7 +686,9 @@ async fn reader_loop( outcome = SessionOutcome::Progressed; } ReaderStep::WsMessage(msg) => { - let Some(msg) = msg else { break None; }; + let Some(msg) = msg else { + break None; + }; let parsed = match msg { Ok(m) => m, Err(e) => break Some(IngestError::Network(e)), @@ -700,7 +705,11 @@ async fn reader_loop( debug!("hydrant sent unexpected binary frame, ignoring"); } WsMessage::Ping(payload) => { - if control_tx.send(WsEvent::IncomingPing(payload)).await.is_err() { + if control_tx + .send(WsEvent::IncomingPing(payload)) + .await + .is_err() + { break None; } } @@ -804,7 +813,7 @@ enum PendingOp { Upsert { source: AtUri, nsid: Nsid, - parsed: Record, + parsed: Box, bytes: Bytes, cid: Option>, edges: Vec, @@ -1046,7 +1055,7 @@ async fn prepare_record( PendingOp::Upsert { source, nsid, - parsed, + parsed: Box::new(parsed), bytes, cid: record.cid, edges, @@ -1093,7 +1102,14 @@ async fn resolve_pending( cid, edges, } => { - let pieces = UpsertPieces { source, nsid, parsed, bytes, cid, edges }; + let pieces = UpsertPieces { + source, + nsid, + parsed: *parsed, + bytes, + cid, + edges, + }; let pieces = match try_park_warming(ctx, cursor, pieces).await { ParkOutcome::Parked { nsid } => { return Pending { @@ -1105,12 +1121,19 @@ async fn resolve_pending( } ParkOutcome::Passthrough(pieces) => *pieces, }; - let UpsertPieces { source, nsid, parsed, bytes, cid, edges } = 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, - parsed, + parsed: Box::new(parsed), bytes, cid, edges, @@ -1193,10 +1216,7 @@ async fn try_park_warming( } } -async fn collect_unresolved_deps( - edges: &[Edge], - resolver: &RepoIdResolver, -) -> Vec { +async fn collect_unresolved_deps(edges: &[Edge], resolver: &RepoIdResolver) -> Vec { futures::stream::iter(edges) .fold(Vec::new(), |mut acc, edge| async move { let Some(uri) = edge.subject.as_uri() else { @@ -1265,7 +1285,7 @@ async fn commit_pending( } => { cache_body(records, &source, cid, bytes); store.upsert_source(&source, edges); - index_search(search, resolver, &source, parsed).await; + index_search(search, resolver, &source, *parsed).await; } PendingOp::Delete { source, nsid: _ } => { store.remove_source(&source); @@ -1292,6 +1312,7 @@ async fn index_search( } #[cfg(test)] +#[allow(clippy::too_many_arguments)] async fn handle_frame( frame: HydrantFrame, store: &EdgeStore, @@ -1685,13 +1706,22 @@ 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, &sys_clock(), 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(); assert_eq!( - parsed["subject"]["did"], - "did:plc:abalone", + parsed["subject"]["did"], "did:plc:abalone", "legacy wire is upgraded to canon shape before caching so downstream readers see canonical fields" ); } @@ -1728,7 +1758,17 @@ mod tests { } } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), 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", @@ -1763,7 +1803,17 @@ mod tests { "record": null } })); - handle_frame(frame, &store, &cov, &NoopSearchSink, &lru, &resolver, &sys_clock(), now()).await; + handle_frame( + frame, + &store, + &cov, + &NoopSearchSink, + &lru, + &resolver, + &sys_clock(), + now(), + ) + .await; assert!(lru.get(&source).is_none()); } @@ -1855,8 +1905,7 @@ mod tests { #[test] fn classify_unknown_hydrant_error_falls_back_to_generic_variant() { - let text = - r#"{"type":"error","error":"NewFutureCode","message":"some new failure mode"}"#; + 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"); @@ -1911,8 +1960,7 @@ mod tests { .to_owned(), ))); let stream: Box = Box::new(ScriptedWsStream { messages }); - let (frame_tx, _frame_rx) = - tokio::sync::mpsc::channel::(FRAME_CHANNEL_DEPTH); + 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(); @@ -1935,8 +1983,7 @@ mod tests { 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 (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(); @@ -2024,10 +2071,7 @@ mod tests { .send(Ok(WsMessage::Text(star_frame_text(1, "starrkeyaa001")))) .await .unwrap(); - ws_tx - .send(Ok(WsMessage::Pong(Bytes::new()))) - .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 @@ -2038,9 +2082,9 @@ mod tests { "first control event must be the pong, not a held text", ); - let too_slow = format!( - "{{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}}" - ); + let too_slow = + "{\"type\":\"error\",\"error\":\"ConsumerTooSlow\",\"message\":\"saturated\"}" + .to_string(); ws_tx.send(Ok(WsMessage::Text(too_slow))).await.unwrap(); let end = tokio::time::timeout(Duration::from_secs(1), reader_handle) @@ -2711,10 +2755,7 @@ mod tests { 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> { + fn send<'a>(&'a mut self, _m: WsMessage) -> bobbin_runtime::WsSendFuture<'a> { Box::pin(async move { Ok(()) }) } } @@ -2754,9 +2795,8 @@ mod tests { #[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 - }); + 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 { @@ -2802,8 +2842,7 @@ mod tests { 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)), + wiremock::ResponseTemplate::new(404).set_delay(Duration::from_millis(150)), ) .mount(&server) .await; diff --git a/crates/ingest/src/shadow.rs b/crates/ingest/src/shadow.rs index 18bacb8..5dc27e9 100644 --- a/crates/ingest/src/shadow.rs +++ b/crates/ingest/src/shadow.rs @@ -58,8 +58,7 @@ impl WarmingShadowBuffer { if count > 0 { self.drained_via_observe_total .fetch_add(count, Ordering::Relaxed); - self.current_concurrent - .fetch_sub(count, Ordering::Relaxed); + self.current_concurrent.fetch_sub(count, Ordering::Relaxed); } } diff --git a/crates/ingest/src/warming.rs b/crates/ingest/src/warming.rs index 3f6dbc9..63aa454 100644 --- a/crates/ingest/src/warming.rs +++ b/crates/ingest/src/warming.rs @@ -111,7 +111,10 @@ impl WarmingBuffer { } async fn park(&self, upsert: ParkedUpsert, deps: Vec) { - debug_assert!(!deps.is_empty(), "park requires at least one unresolved dep"); + debug_assert!( + !deps.is_empty(), + "park requires at least one unresolved dep" + ); let dep_count = deps.len() as u64; let source = upsert.source.clone(); let handle: EntryHandle = Arc::new(Mutex::new(Some(EntryState { @@ -144,7 +147,8 @@ impl WarmingBuffer { self.dep_enqueued_total .fetch_add(dep_count, Ordering::Relaxed); let cur = self.current_entries.fetch_add(1, Ordering::Relaxed) + 1; - self.max_concurrent_entries.fetch_max(cur, Ordering::Relaxed); + self.max_concurrent_entries + .fetch_max(cur, Ordering::Relaxed); } async fn register_dep(&self, dep: RepoIdent, handle: EntryHandle) { @@ -201,7 +205,9 @@ impl WarmingBuffer { None => None, } }; - let Some((matched, drain_now)) = outcome else { return; }; + let Some((matched, drain_now)) = outcome else { + return; + }; if matched { dep_matches += 1; } @@ -363,11 +369,13 @@ mod tests { buf.park( make_upsert("at://did:plc:starer1/sh.tangled.feed.star/aaaaaaaaaaaaa"), vec![dep_a.clone()], - ).await; + ) + .await; buf.park( make_upsert("at://did:plc:starer2/sh.tangled.feed.star/bbbbbbbbbbbbb"), vec![dep_b.clone()], - ).await; + ) + .await; let s = buf.snapshot(); assert_eq!(s.enqueued_total, 2); assert_eq!(s.distinct_keys_seen, 2); @@ -383,8 +391,11 @@ mod tests { buf.park( make_upsert("at://did:plc:starer/sh.tangled.feed.star/zzzzzzzzzzzzz"), vec![RepoIdent::new(d("did:plc:nel"), r("abcabcabcabcz"))], - ).await; - let drained = buf.take_observed(&d("did:plc:olaren"), &r("nopenopenopep")).await; + ) + .await; + let drained = buf + .take_observed(&d("did:plc:olaren"), &r("nopenopenopep")) + .await; assert!(drained.is_empty()); let s = buf.snapshot(); assert_eq!(s.current_entries, 1); @@ -418,7 +429,8 @@ mod tests { let dep_a = RepoIdent::new(d("did:plc:nel"), r("abcabcabcabcz")); let dep_b = RepoIdent::new(d("did:plc:olaren"), r("abcabcabcabd1")); let source = "at://did:plc:starer/sh.tangled.feed.star/zzzzzzzzzzzzz"; - buf.park(make_upsert(source), vec![dep_a.clone(), dep_b.clone()]).await; + buf.park(make_upsert(source), vec![dep_a.clone(), dep_b.clone()]) + .await; assert!(buf.evict_source(&at(source)).await); let drained_a = buf.take_observed(&dep_a.owner, &dep_a.rkey).await; let drained_b = buf.take_observed(&dep_b.owner, &dep_b.rkey).await; @@ -469,7 +481,8 @@ mod tests { buf.park( make_upsert("at://did:plc:starer/sh.tangled.feed.star/zzzzzzzzzzzzz"), vec![dep], - ).await; + ) + .await; let residual = buf.drain_for_promote().await; assert_eq!(residual.len(), 1); let s = buf.snapshot(); @@ -500,7 +513,8 @@ mod tests { buf.park( make_upsert("at://did:plc:starer/sh.tangled.feed.star/zzzzzzzzzzzzz"), vec![dep.clone()], - ).await; + ) + .await; let first = buf.take_observed(&dep.owner, &dep.rkey).await; let second = buf.take_observed(&dep.owner, &dep.rkey).await; assert_eq!(first.len(), 1); diff --git a/crates/knot-proxy/src/breaker.rs b/crates/knot-proxy/src/breaker.rs index 4b6f6b6..7b0dcb9 100644 --- a/crates/knot-proxy/src/breaker.rs +++ b/crates/knot-proxy/src/breaker.rs @@ -199,7 +199,9 @@ mod tests { let after = t0 + Duration::from_millis(60); b.record_failure_at(t0); assert!(b.try_acquire_at(t0).is_err()); - let trial = b.try_acquire_at(after).expect("trial admitted after cooldown"); + let trial = b + .try_acquire_at(after) + .expect("trial admitted after cooldown"); assert!( b.try_acquire_at(after).is_err(), "second concurrent half-open call must be rejected", @@ -228,7 +230,9 @@ mod tests { let t0 = Instant::now(); let after = t0 + Duration::from_millis(60); b.record_failure_at(t0); - b.try_acquire_at(after).expect("trial admitted").record_success(); + b.try_acquire_at(after) + .expect("trial admitted") + .record_success(); b.try_acquire_at(after).expect("closed").record_success(); b.try_acquire_at(after).expect("closed").record_success(); } diff --git a/crates/resolver/src/legacy_upgrade.rs b/crates/resolver/src/legacy_upgrade.rs index fb2374a..e958ebf 100644 --- a/crates/resolver/src/legacy_upgrade.rs +++ b/crates/resolver/src/legacy_upgrade.rs @@ -12,8 +12,8 @@ use jacquard_common::DefaultStr; use jacquard_common::types::did::Did; use jacquard_common::types::string::AtUri; -use crate::{RepoIdResolver, Resolution}; use crate::normalize::{is_repo_at_uri, resolve_repo_uri}; +use crate::{RepoIdResolver, Resolution}; use jacquard_common::IntoStatic; use jacquard_common::types::ident::AtIdentifier; use jacquard_common::types::recordkey::Rkey; @@ -76,8 +76,7 @@ pub async fn decode_canon_or_upgrade_bytes<'a>( "{nsid}: legacy upgrade failed" ))); }; - let canon_bytes = - serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson)?; + let canon_bytes = serialize_canon_variant(&canon).map_err(ExtractError::DecodeJson)?; Ok((canon, alloc::borrow::Cow::Owned(canon_bytes))) } } @@ -94,10 +93,7 @@ fn serialize_canon_variant(record: &Record) -> Result, serde } } -pub async fn upgrade( - legacy: LegacyRecord, - resolver: &RepoIdResolver, -) -> Option { +pub async fn upgrade(legacy: LegacyRecord, resolver: &RepoIdResolver) -> Option { match legacy { LegacyRecord::Issue(l) => upgrade_issue(l, resolver).await.map(Record::Issue), LegacyRecord::Pull(l) => upgrade_pull(l, resolver).await.map(Record::Pull), @@ -272,7 +268,10 @@ mod tests { let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; let decoded = DecodedRecord::try_decode("sh.tangled.repo.issue", json) .expect("legacy issue must decode"); - assert!(matches!(decoded, DecodedRecord::Legacy(LegacyRecord::Issue(_)))); + assert!(matches!( + decoded, + DecodedRecord::Legacy(LegacyRecord::Issue(_)) + )); } #[test] @@ -300,8 +299,7 @@ mod tests { async fn upgrade_issue_uses_repo_did_directly() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.issue","repoDid":"did:plc:abalone","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = - LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), @@ -318,8 +316,7 @@ mod tests { .observe(owner.clone(), key.clone(), Some(did("did:plc:abalone"))) .await; let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = - LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Issue(i) => assert_eq!(i.repo, did("did:plc:abalone")), @@ -331,8 +328,7 @@ mod tests { async fn upgrade_issue_drops_when_resolver_cannot_map() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.issue","repo":"at://did:plc:nel/sh.tangled.repo/abcabcabcabcz","title":"t","createdAt":"2026-05-01T00:00:00Z"}"#; - let legacy = - LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.issue", json).expect("decode"); assert!( upgrade(legacy, &resolver).await.is_none(), "no resolver entry and no repoDid means the canon Did cannot be constructed", @@ -343,8 +339,7 @@ mod tests { async fn upgrade_pull_propagates_target_resolution() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"}}"#; - let legacy = - LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { @@ -359,8 +354,7 @@ mod tests { async fn upgrade_pull_source_repo_resolution_is_independent_of_target() { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let json = br#"{"$type":"sh.tangled.repo.pull","title":"t","createdAt":"2026-05-01T00:00:00Z","rounds":[],"target":{"branch":"main","repoDid":"did:plc:abalone"},"source":{"branch":"feat","repo":"at://did:plc:nel/sh.tangled.repo/missing"}}"#; - let legacy = - LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); + let legacy = LegacyRecord::from_json_bytes("sh.tangled.repo.pull", json).expect("decode"); let canon = upgrade(legacy, &resolver).await.expect("upgrade"); match canon { Record::Pull(p) => { @@ -448,7 +442,8 @@ mod tests { let resolver = RepoIdResolver::detached(RuntimeHasher::default()); let with_did = br#"{"$type":"sh.tangled.repo.collaborator","createdAt":"2026-05-01T00:00:00Z","subject":"did:plc:lyna","repoDid":"did:plc:abalone"}"#; let canon = upgrade( - LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", with_did).expect("decode"), + LegacyRecord::from_json_bytes("sh.tangled.repo.collaborator", with_did) + .expect("decode"), &resolver, ) .await diff --git a/crates/resolver/src/lib.rs b/crates/resolver/src/lib.rs index 69f9708..8b10a9c 100644 --- a/crates/resolver/src/lib.rs +++ b/crates/resolver/src/lib.rs @@ -10,6 +10,8 @@ use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; +use tokio::time::Instant; + use bobbin_runtime::{Clock, RuntimeHasher}; use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; use bobbin_types::edges::{ExtractError, Record}; @@ -19,9 +21,11 @@ use jacquard_common::types::did::Did; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::recordkey::Rkey; use scc::HashMap as SccMap; +use tokio::sync::OnceCell; use tracing::warn; const REPO_COLLECTION: &str = "sh.tangled.repo"; +const TRANSIENT_TTL: Duration = Duration::from_secs(60); #[derive(Clone, Debug, Eq, PartialEq)] pub enum Resolution { @@ -49,6 +53,7 @@ impl AuthoritativeResolution { enum CacheEntry { Authoritative(AuthoritativeResolution), Provisional(Resolution), + Transient { expires_at: Instant }, } impl CacheEntry { @@ -57,8 +62,13 @@ impl CacheEntry { Self::Authoritative(AuthoritativeResolution::Mapped(did)) => Resolution::Mapped(did), Self::Authoritative(AuthoritativeResolution::NoRepoDid) => Resolution::NoRepoDid, Self::Provisional(r) => r, + Self::Transient { .. } => Resolution::Unresolvable, } } + + fn is_expired_transient(&self, now: Instant) -> bool { + matches!(self, Self::Transient { expires_at } if *expires_at <= now) + } } #[derive(Default)] @@ -158,6 +168,7 @@ struct SlingshotProbe { pub struct RepoIdResolver { cache: SccMap, by_repo_did: SccMap, RepoIdent, RuntimeHasher>, + in_flight: SccMap>, RuntimeHasher>, probe: Option, stats: ResolverStats, } @@ -170,7 +181,8 @@ impl RepoIdResolver { ) -> Self { Self { cache: SccMap::with_hasher(hasher.clone()), - by_repo_did: SccMap::with_hasher(hasher), + by_repo_did: SccMap::with_hasher(hasher.clone()), + in_flight: SccMap::with_hasher(hasher), probe: Some(SlingshotProbe { client, clock }), stats: ResolverStats::default(), } @@ -179,7 +191,8 @@ impl RepoIdResolver { pub fn detached(hasher: RuntimeHasher) -> Self { Self { cache: SccMap::with_hasher(hasher.clone()), - by_repo_did: SccMap::with_hasher(hasher), + by_repo_did: SccMap::with_hasher(hasher.clone()), + in_flight: SccMap::with_hasher(hasher), probe: None, stats: ResolverStats::default(), } @@ -195,10 +208,14 @@ impl RepoIdResolver { rkey: &Rkey, ) -> Option { let key = RepoIdent::new(owner.clone(), rkey.clone()); - self.cache - .get_async(&key) - .await - .map(|entry| entry.get().clone().into_resolution()) + let entry = self.cache.get_async(&key).await?; + let now = self.probe.as_ref().map(|p| p.clock.now_instant()); + if let Some(now) = now + && entry.get().is_expired_transient(now) + { + return None; + } + Some(entry.get().clone().into_resolution()) } pub async fn observe( @@ -216,9 +233,7 @@ impl RepoIdResolver { .and_modify(|existing| *existing = entry.clone()) .or_insert(entry); - let Some(repo_did) = repo_did else { - return None; - }; + let repo_did = repo_did?; let mut prior: Option = None; self.by_repo_did .entry_async(repo_did) @@ -261,16 +276,68 @@ impl RepoIdResolver { .or_insert(entry); } + async fn fill_transient(&self, key: RepoIdent, expires_at: Instant) { + let entry = CacheEntry::Transient { expires_at }; + self.cache + .entry_async(key) + .await + .and_modify(|existing| { + if matches!(existing, CacheEntry::Authoritative(_)) { + return; + } + *existing = entry.clone(); + }) + .or_insert(entry); + } + pub async fn resolve(&self, owner: &Did, rkey: &Rkey) -> Resolution { let key = RepoIdent::new(owner.clone(), rkey.clone()); - if let Some(entry) = self.cache.get_async(&key).await { - self.stats.record_hit(); - return entry.get().clone().into_resolution(); - } + let Some(probe) = self.probe.as_ref() else { + if let Some(entry) = self.cache.get_async(&key).await { + self.stats.record_hit(); + return entry.get().clone().into_resolution(); + } self.stats.record_miss(MissKind::NoClient, None); return Resolution::Unresolvable; }; + + let now = probe.clock.now_instant(); + if let Some(entry) = self.cache.get_async(&key).await + && !entry.get().is_expired_transient(now) + { + self.stats.record_hit(); + return entry.get().clone().into_resolution(); + } + + let cell: Arc> = self + .in_flight + .entry_async(key.clone()) + .await + .or_insert_with(|| Arc::new(OnceCell::new())) + .get() + .clone(); + + let result = cell + .get_or_init(|| async { self.fetch_repo_did(owner, rkey, &key).await }) + .await + .clone(); + + self.in_flight.remove_async(&key).await; + + result + } + + async fn fetch_repo_did( + &self, + owner: &Did, + rkey: &Rkey, + key: &RepoIdent, + ) -> Resolution { + let probe = self + .probe + .as_ref() + .expect("fetch_repo_did is only called when a probe is present"); let started = probe.clock.now_instant(); let nsid: Nsid = nsid_static(REPO_COLLECTION); let provisional = match probe.client.get_record(owner, &nsid, rkey).await { @@ -309,11 +376,12 @@ impl RepoIdResolver { error = ?e, owner = owner.as_ref(), rkey = rkey.as_ref(), - "slingshot transient failure during repoDID lookup, will retry", + "caching transient slingshot failure for repoDID lookup under short TTL", ); let elapsed = probe.clock.now_instant().duration_since(started); - self.stats - .record_miss(MissKind::Transient, Some(elapsed)); + self.stats.record_miss(MissKind::Transient, Some(elapsed)); + let expires_at = probe.clock.now_instant() + TRANSIENT_TTL; + self.fill_transient(key.clone(), expires_at).await; return Resolution::Unresolvable; } }; @@ -324,7 +392,8 @@ impl RepoIdResolver { Resolution::Unresolvable => MissKind::Unresolvable, }; self.stats.record_miss(kind, Some(elapsed)); - self.fill_provisional(key, provisional.clone()).await; + self.fill_provisional(key.clone(), provisional.clone()) + .await; provisional } } @@ -564,7 +633,8 @@ mod tests { let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -590,7 +660,8 @@ mod tests { let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -621,7 +692,8 @@ mod tests { let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -648,7 +720,8 @@ mod tests { let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -663,25 +736,30 @@ mod tests { } #[tokio::test] - async fn slingshot_transport_error_does_not_cache() { + async fn slingshot_transport_error_caches_with_short_ttl() { 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(503)) - .expect(2) + .expect(1) .mount(&server) .await; let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); let first = resolver.resolve(&owner, &key).await; let second = resolver.resolve(&owner, &key).await; assert_eq!(first, Resolution::Unresolvable); - assert_eq!(second, Resolution::Unresolvable); + assert_eq!( + second, + Resolution::Unresolvable, + "transient TTL must suppress immediate re-hammering of a sick upstream", + ); } #[tokio::test] @@ -690,13 +768,14 @@ mod tests { wiremock::Mock::given(wiremock::matchers::method("GET")) .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) .respond_with(wiremock::ResponseTemplate::new(503)) - .expect(2) + .expect(1) .mount(&server) .await; let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -705,9 +784,10 @@ mod tests { let snap = resolver.stats(); assert_eq!( - snap.misses_transient, 2, - "transport-error retries must count as transient, not as canonical unresolvable", + snap.misses_transient, 1, + "second resolve must hit the short-TTL cache instead of re-firing the transient miss", ); + assert_eq!(snap.hits, 1, "second call hits cached transient entry"); assert_eq!( snap.misses_unresolvable, 0, "canonical unresolvable counter is reserved for cached terminal answers", @@ -716,7 +796,62 @@ mod tests { snap.miss_latency_micros_sum > 0, "transient misses still have latency contributions", ); - assert_eq!(snap.miss_count(), 2); + assert_eq!(snap.miss_count(), 1); + } + + #[tokio::test] + async fn slingshot_in_flight_requests_coalesce() { + let server = wiremock::MockServer::start().await; + let body = serde_json::json!({ + "uri": "at://did:plc:nel/sh.tangled.repo/abcabcabcabcz", + "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", + "value": {"$type": "sh.tangled.repo", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": "did:plc:limpet"} + }); + wiremock::Mock::given(wiremock::matchers::method("GET")) + .and(wiremock::matchers::path("/xrpc/com.atproto.repo.getRecord")) + .respond_with( + wiremock::ResponseTemplate::new(200) + .set_body_json(body) + .set_delay(Duration::from_millis(200)), + ) + .expect(1) + .mount(&server) + .await; + + let client = + SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); + let resolver = Arc::new(RepoIdResolver::with_slingshot( + client, + test_clock(), + RuntimeHasher::default(), + )); + + let owner = did("did:plc:nel"); + let key = rkey("abcabcabcabcz"); + let r0 = resolver.clone(); + let r1 = resolver.clone(); + let r2 = resolver.clone(); + let o0 = owner.clone(); + let o1 = owner.clone(); + let o2 = owner.clone(); + let k0 = key.clone(); + let k1 = key.clone(); + let k2 = key.clone(); + let (a, b, c) = tokio::join!( + tokio::spawn(async move { r0.resolve(&o0, &k0).await }), + tokio::spawn(async move { r1.resolve(&o1, &k1).await }), + tokio::spawn(async move { r2.resolve(&o2, &k2).await }), + ); + let expected = Resolution::Mapped(did("did:plc:limpet")); + assert_eq!(a.unwrap(), expected); + assert_eq!(b.unwrap(), expected); + assert_eq!(c.unwrap(), expected); + + let snap = resolver.stats(); + assert_eq!( + snap.misses_mapped, 1, + "only the winning task pays the slingshot RTT", + ); } #[tokio::test] @@ -729,7 +864,8 @@ mod tests { .await; let client = SlingshotClient::with_default_http(url::Url::parse(&server.uri()).unwrap()).unwrap(); - let resolver = RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); + let resolver = + RepoIdResolver::with_slingshot(client, test_clock(), RuntimeHasher::default()); let owner = did("did:plc:nel"); let key = rkey("abcabcabcabcz"); @@ -737,11 +873,17 @@ mod tests { resolver.resolve(&owner, &key).await; let snap = resolver.stats(); - assert_eq!(snap.misses_unresolvable, 1, "first call is the slingshot miss"); + assert_eq!( + snap.misses_unresolvable, 1, + "first call is the slingshot miss" + ); assert_eq!(snap.hits, 1, "second call hits the unresolvable cache"); assert_eq!(snap.miss_count(), 1); assert_eq!(snap.total(), 2); - assert!(snap.miss_latency_micros_sum > 0, "latency recorded for slingshot miss"); + assert!( + snap.miss_latency_micros_sum > 0, + "latency recorded for slingshot miss" + ); assert!(snap.miss_latency_micros_avg().unwrap() > 0); } diff --git a/crates/runtime/src/entropy.rs b/crates/runtime/src/entropy.rs index 7f0367d..679bad0 100644 --- a/crates/runtime/src/entropy.rs +++ b/crates/runtime/src/entropy.rs @@ -96,7 +96,10 @@ mod tests { let e = SeededEntropy::new(0); let mut seen = std::collections::HashSet::new(); for _ in 0..1024 { - assert!(seen.insert(e.next_u64()), "splitmix collided inside 1024 draws"); + assert!( + seen.insert(e.next_u64()), + "splitmix collided inside 1024 draws" + ); } } diff --git a/crates/runtime/src/mem_network.rs b/crates/runtime/src/mem_network.rs index 8b95cc3..0491914 100644 --- a/crates/runtime/src/mem_network.rs +++ b/crates/runtime/src/mem_network.rs @@ -32,7 +32,10 @@ pub struct MemHttpBody { impl MemHttpBody { pub fn ok_json(body: Bytes) -> Self { let mut headers = HeaderMap::new(); - headers.insert(http::header::CONTENT_TYPE, "application/json".parse().unwrap()); + headers.insert( + http::header::CONTENT_TYPE, + "application/json".parse().unwrap(), + ); Self { status: StatusCode::OK, headers, @@ -64,7 +67,10 @@ impl MemHttpTransport { Self { responder, clock } } - pub fn shared(responder: Arc, clock: Arc) -> Arc { + pub fn shared( + responder: Arc, + clock: Arc, + ) -> Arc { Arc::new(Self::new(responder, clock)) } } @@ -203,14 +209,13 @@ mod tests { fn respond(&self, _: &HttpRequest) -> MemHttpResponse { let i = self.cursor.fetch_add(1, Ordering::Relaxed); let mut guard = self.responses.lock().unwrap(); - let next = std::mem::replace( + std::mem::replace( &mut guard[i], MemHttpResponse { latency: Duration::ZERO, result: Err(NetworkError::Transport("script consumed".into())), }, - ); - next + ) } } @@ -231,9 +236,8 @@ mod tests { #[tokio::test(start_paused = true)] async fn http_returns_scripted_body_after_injected_latency() { let clock: Arc = Arc::new(SimClock::at(UnixMicros::new(0))); - let responder: Arc = Arc::new(ScriptedHttp::new(vec![ - ok_response("{\"hello\":1}", 10), - ])); + let responder: Arc = + Arc::new(ScriptedHttp::new(vec![ok_response("{\"hello\":1}", 10)])); let transport = MemHttpTransport::new(responder, clock); let request = HttpRequest { url: Url::parse("http://oyster.cafe/xrpc/x").unwrap(), @@ -345,7 +349,10 @@ mod tests { .connect(Url::parse("ws://oyster.cafe/").unwrap()) .await .unwrap(); - conn.sink.send(WsMessage::Text("ping".into())).await.unwrap(); + conn.sink + .send(WsMessage::Text("ping".into())) + .await + .unwrap(); let echoed = conn.stream.next().await.unwrap().unwrap(); assert!(matches!(echoed, WsMessage::Text(t) if t == "ping")); } diff --git a/crates/search/src/lib.rs b/crates/search/src/lib.rs index 7f10ea3..64d50ad 100644 --- a/crates/search/src/lib.rs +++ b/crates/search/src/lib.rs @@ -533,7 +533,12 @@ mod tests { idx.flush().await; let page = idx - .search("barnacle", SearchFilters::default(), SearchCursor::Start, 10) + .search( + "barnacle", + SearchFilters::default(), + SearchCursor::Start, + 10, + ) .await .unwrap(); assert_eq!(page.hits.len(), 1); diff --git a/crates/types/src/edges.rs b/crates/types/src/edges.rs index 14dcf88..cd5c855 100644 --- a/crates/types/src/edges.rs +++ b/crates/types/src/edges.rs @@ -273,15 +273,33 @@ const MIRROR_KINDS: &[(&str, &str)] = &[ ("sh.tangled.knot.member", "sh.tangled.knot.member.by"), ("sh.tangled.label.op", "sh.tangled.label.op.by"), ("sh.tangled.pipeline", "sh.tangled.pipeline.by"), - ("sh.tangled.pipeline.status", "sh.tangled.pipeline.status.by"), + ( + "sh.tangled.pipeline.status", + "sh.tangled.pipeline.status.by", + ), ("sh.tangled.repo.artifact", "sh.tangled.repo.artifact.by"), - ("sh.tangled.repo.collaborator", "sh.tangled.repo.collaborator.by"), + ( + "sh.tangled.repo.collaborator", + "sh.tangled.repo.collaborator.by", + ), ("sh.tangled.repo.issue", "sh.tangled.repo.issue.by"), - ("sh.tangled.repo.issue.comment", "sh.tangled.repo.issue.comment.by"), - ("sh.tangled.repo.issue.state", "sh.tangled.repo.issue.state.by"), + ( + "sh.tangled.repo.issue.comment", + "sh.tangled.repo.issue.comment.by", + ), + ( + "sh.tangled.repo.issue.state", + "sh.tangled.repo.issue.state.by", + ), ("sh.tangled.repo.pull", "sh.tangled.repo.pull.by"), - ("sh.tangled.repo.pull.comment", "sh.tangled.repo.pull.comment.by"), - ("sh.tangled.repo.pull.status", "sh.tangled.repo.pull.status.by"), + ( + "sh.tangled.repo.pull.comment", + "sh.tangled.repo.pull.comment.by", + ), + ( + "sh.tangled.repo.pull.status", + "sh.tangled.repo.pull.status.by", + ), ("sh.tangled.spindle.member", "sh.tangled.spindle.member.by"), ]; diff --git a/crates/types/src/ids.rs b/crates/types/src/ids.rs index 04fbd59..3146e52 100644 --- a/crates/types/src/ids.rs +++ b/crates/types/src/ids.rs @@ -113,7 +113,10 @@ mod tests { #[test] fn owner_did_rejects_handle_authority() { - assert_eq!(owner_did_from_aturi("at://oyster.cafe/sh.tangled.repo/r1"), None); + assert_eq!( + owner_did_from_aturi("at://oyster.cafe/sh.tangled.repo/r1"), + None + ); } #[test] @@ -129,6 +132,10 @@ mod tests { s.insert(SubjectRef::Uri( AtUri::new_owned("at://did:plc:abalone").unwrap(), )); - assert_eq!(s.len(), 2, "did variant and uri variant must hash distinctly"); + assert_eq!( + s.len(), + 2, + "did variant and uri variant must hash distinctly" + ); } } diff --git a/crates/types/src/legacy.rs b/crates/types/src/legacy.rs index 23133e7..ebafaa6 100644 --- a/crates/types/src/legacy.rs +++ b/crates/types/src/legacy.rs @@ -2,7 +2,7 @@ use alloc::collections::BTreeMap; use alloc::vec::Vec; use jacquard_common::deps::smol_str::SmolStr; -use jacquard_common::types::string::{Datetime, Did, AtUri}; +use jacquard_common::types::string::{AtUri, Datetime, Did}; use jacquard_common::types::value::Data; use jacquard_common::{BosStr, DefaultStr}; use serde::Deserialize; @@ -35,7 +35,10 @@ pub struct LegacyIssue { } #[derive(Debug, Deserialize)] -#[serde(rename_all = "camelCase", bound(deserialize = "S: Deserialize<'de> + BosStr"))] +#[serde( + rename_all = "camelCase", + bound(deserialize = "S: Deserialize<'de> + BosStr") +)] pub struct LegacyTarget { pub branch: S, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -45,7 +48,10 @@ pub struct LegacyTarget { } #[derive(Debug, Deserialize)] -#[serde(rename_all = "camelCase", bound(deserialize = "S: Deserialize<'de> + BosStr"))] +#[serde( + rename_all = "camelCase", + bound(deserialize = "S: Deserialize<'de> + BosStr") +)] pub struct LegacySource { pub branch: S, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/crates/types/src/search.rs b/crates/types/src/search.rs index e6d03b6..604d4a6 100644 --- a/crates/types/src/search.rs +++ b/crates/types/src/search.rs @@ -162,7 +162,14 @@ fn profile_doc(source: &AtUri, r: &Profile) -> SearchDoc if let Some(p) = &r.pronouns { parts.push(p.as_str().to_owned()); } - doc(source, "sh.tangled.actor.profile", &title, parts, None, None) + doc( + source, + "sh.tangled.actor.profile", + &title, + parts, + None, + None, + ) } fn repo_doc(source: &AtUri, r: &RepoRecord) -> SearchDoc { diff --git a/crates/xrpc/src/lib.rs b/crates/xrpc/src/lib.rs index f8f5b81..e5089ca 100644 --- a/crates/xrpc/src/lib.rs +++ b/crates/xrpc/src/lib.rs @@ -196,9 +196,15 @@ pub fn router(state: AppState) -> Router { .route("/xrpc/sh.tangled.knot.listKnots", get(list_knots)) .route("/xrpc/sh.tangled.knot.countKnots", get(count_knots)) .route("/xrpc/sh.tangled.spindle.listSpindles", get(list_spindles)) - .route("/xrpc/sh.tangled.spindle.countSpindles", get(count_spindles)) + .route( + "/xrpc/sh.tangled.spindle.countSpindles", + get(count_spindles), + ) .route("/xrpc/sh.tangled.publicKey.listKeys", get(list_public_keys)) - .route("/xrpc/sh.tangled.publicKey.countKeys", get(count_public_keys)) + .route( + "/xrpc/sh.tangled.publicKey.countKeys", + get(count_public_keys), + ) .route("/xrpc/sh.tangled.graph.listVouches", get(list_vouches)) .route("/xrpc/sh.tangled.graph.countVouches", get(count_vouches)) .route("/xrpc/sh.tangled.feed.listStarsBy", get(list_stars_by)) @@ -680,7 +686,6 @@ fn parse_uri(raw: &str) -> Result, XrpcError> { AtUri::::new_owned(raw).map_err(|e| XrpcError::InvalidParams(format!("uri: {e}"))) } - #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum SubjectShape { BareDid, @@ -811,8 +816,7 @@ impl MirrorOf for SpindleMemberBy { } impl HasSubject for StarRecord { - const SHAPE: SubjectShape = - SubjectShape::BareDidOrOneOfCollections(&["sh.tangled.string"]); + const SHAPE: SubjectShape = SubjectShape::BareDidOrOneOfCollections(&["sh.tangled.string"]); } impl HasSubject for FollowRecord { const SHAPE: SubjectShape = SubjectShape::BareDid; @@ -1128,7 +1132,9 @@ async fn get_repos( RawQuery(query): RawQuery, ) -> Result>>, XrpcError> { let uris = collect_repeated(query.as_deref(), BULK_REPOS_KEY); - bulk_fetch::>(&state, uris).await.map(Json) + bulk_fetch::>(&state, uris) + .await + .map(Json) } async fn get_profiles( @@ -1136,7 +1142,9 @@ async fn get_profiles( RawQuery(query): RawQuery, ) -> Result>>, XrpcError> { let uris = collect_repeated(query.as_deref(), BULK_PROFILES_KEY); - bulk_fetch::>(&state, uris).await.map(Json) + bulk_fetch::>(&state, uris) + .await + .map(Json) } async fn get_issues( @@ -1144,7 +1152,9 @@ async fn get_issues( RawQuery(query): RawQuery, ) -> Result>>, XrpcError> { let uris = collect_repeated(query.as_deref(), BULK_ISSUES_KEY); - bulk_fetch::>(&state, uris).await.map(Json) + bulk_fetch::>(&state, uris) + .await + .map(Json) } async fn get_pulls( @@ -1152,7 +1162,9 @@ async fn get_pulls( RawQuery(query): RawQuery, ) -> Result>>, XrpcError> { let uris = collect_repeated(query.as_deref(), BULK_PULLS_KEY); - bulk_fetch::>(&state, uris).await.map(Json) + bulk_fetch::>(&state, uris) + .await + .map(Json) } const BULK_REPOS_KEY: &str = "repos"; @@ -1190,22 +1202,24 @@ where let items: Vec> = stream::iter(parsed) .map(|uri| async move { match resolve(state, ExpectedNsid::new(R::NSID), uri).await { - Ok((body, _)) => match deserialize_or_upgrade::(state, R::NSID, &body.value).await { - Ok(value) => { - let Some(value) = value.normalize(&state.resolver).await else { - return Ok(None); - }; - Ok(Some(RecordView { - uri: body.uri.clone(), - cid: Some(body.cid.clone()), - value, - })) + Ok((body, _)) => { + match deserialize_or_upgrade::(state, R::NSID, &body.value).await { + Ok(value) => { + let Some(value) = value.normalize(&state.resolver).await else { + return Ok(None); + }; + Ok(Some(RecordView { + uri: body.uri.clone(), + cid: Some(body.cid.clone()), + value, + })) + } + Err(_) => Ok(None), } - Err(_) => Ok(None), - }, - Err(XrpcError::NotFound | XrpcError::UpstreamGone(_) | XrpcError::InvalidRecord(_)) => { - Ok(None) } + Err( + XrpcError::NotFound | XrpcError::UpstreamGone(_) | XrpcError::InvalidRecord(_), + ) => Ok(None), Err(other) => Err(other), } }) @@ -1301,10 +1315,7 @@ where }) } -fn count_mirror( - state: &AppState, - q: CountQuery, -) -> Result { +fn count_mirror(state: &AppState, q: CountQuery) -> Result { let subject = parse_subject(&q.subject, M::SHAPE)?; let key = EdgeKey::new(nsid_static(M::EDGE_KIND), subject); Ok(CountResponse { @@ -1419,7 +1430,9 @@ async fn list_ref_updates( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result>>, XrpcError> { - list_records::(&state, q).await.map(Json) + list_records::(&state, q) + .await + .map(Json) } async fn count_ref_updates( @@ -1523,7 +1536,9 @@ async fn list_public_keys( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result>>, XrpcError> { - list_records::(&state, q).await.map(Json) + list_records::(&state, q) + .await + .map(Json) } async fn count_public_keys( @@ -1655,7 +1670,9 @@ async fn list_pipeline_statuses_by( State(state): State, XrpcQuery(q): XrpcQuery, ) -> Result>>, XrpcError> { - list_mirror::(&state, q).await.map(Json) + list_mirror::(&state, q) + .await + .map(Json) } async fn count_pipeline_statuses_by( State(state): State, @@ -1995,9 +2012,7 @@ fn build_search_filters(q: &SearchQueryParams) -> Result u { - return Err(XrpcError::InvalidParams( - "since must be <= until".into(), - )); + return Err(XrpcError::InvalidParams("since must be <= until".into())); } Ok(SearchFilters { nsid, diff --git a/crates/xrpc/tests/aggregation.rs b/crates/xrpc/tests/aggregation.rs index abab05f..31f8e98 100644 --- a/crates/xrpc/tests/aggregation.rs +++ b/crates/xrpc/tests/aggregation.rs @@ -986,10 +986,7 @@ async fn star_endpoints_reject_unrelated_collection() { let (status, body) = json_response(resp).await; assert_eq!(status, StatusCode::BAD_REQUEST, "{endpoint}"); let msg = body["message"].as_str().unwrap_or_default(); - assert!( - msg.contains("sh.tangled.string"), - "{endpoint}: {msg}", - ); + assert!(msg.contains("sh.tangled.string"), "{endpoint}: {msg}",); } }) .await; diff --git a/crates/xrpc/tests/bulk.rs b/crates/xrpc/tests/bulk.rs index 837edd2..01025ab 100644 --- a/crates/xrpc/tests/bulk.rs +++ b/crates/xrpc/tests/bulk.rs @@ -72,9 +72,10 @@ impl Harness { .and(query_param("repo", did)) .and(query_param("collection", collection)) .and(query_param("rkey", rkey)) - .respond_with(ResponseTemplate::new(404).set_body_json( - json!({"error": "RecordNotFound", "message": "missing"}), - )) + .respond_with( + ResponseTemplate::new(404) + .set_body_json(json!({"error": "RecordNotFound", "message": "missing"})), + ) .mount(&self.server) .await; } @@ -145,10 +146,20 @@ fn profile_body(handle: &str) -> Value { #[tokio::test] async fn get_repos_returns_all_resolved_records() { let h = Harness::new().await; - h.mount("did:plc:nel", "sh.tangled.repo", "abalone", repo_body("abalone")) - .await; - h.mount("did:plc:teq", "sh.tangled.repo", "limpet", repo_body("limpet")) - .await; + h.mount( + "did:plc:nel", + "sh.tangled.repo", + "abalone", + repo_body("abalone"), + ) + .await; + h.mount( + "did:plc:teq", + "sh.tangled.repo", + "limpet", + repo_body("limpet"), + ) + .await; let app = router(h.state.clone()); let (status, body) = json_response( app.oneshot(bulk_request( @@ -278,8 +289,13 @@ async fn get_pulls_returns_all_resolved_pulls() { #[tokio::test] async fn missing_records_are_dropped_silently() { let h = Harness::new().await; - h.mount("did:plc:nel", "sh.tangled.repo", "abalone", repo_body("abalone")) - .await; + h.mount( + "did:plc:nel", + "sh.tangled.repo", + "abalone", + repo_body("abalone"), + ) + .await; h.mount_404("did:plc:teq", "sh.tangled.repo", "ghost").await; let app = router(h.state.clone()); let (status, body) = json_response( @@ -297,7 +313,11 @@ async fn missing_records_are_dropped_silently() { .await; assert_eq!(status, StatusCode::OK); let items = body["items"].as_array().unwrap(); - assert_eq!(items.len(), 1, "missing records must be dropped not fail the bulk call"); + assert_eq!( + items.len(), + 1, + "missing records must be dropped not fail the bulk call" + ); assert_eq!(items[0]["value"]["name"], json!("abalone")); } diff --git a/crates/xrpc/tests/coverage.rs b/crates/xrpc/tests/coverage.rs index a3c8cc3..ee64326 100644 --- a/crates/xrpc/tests/coverage.rs +++ b/crates/xrpc/tests/coverage.rs @@ -111,13 +111,7 @@ async fn promotion_flips_ready_field() { events_processed: 1, last_cursor: HydrantCursor::new(5), }); - let (_, before) = json_response( - app.clone() - .oneshot(coverage_request()) - .await - .unwrap(), - ) - .await; + let (_, before) = json_response(app.clone().oneshot(coverage_request()).await.unwrap()).await; assert_eq!(before["ready"], json!(false)); h.coverage.update(|_| Coverage::Ready { diff --git a/crates/xrpc/tests/search.rs b/crates/xrpc/tests/search.rs index 4c63ea2..2ae9a03 100644 --- a/crates/xrpc/tests/search.rs +++ b/crates/xrpc/tests/search.rs @@ -75,7 +75,8 @@ impl Harness { } async fn index_issue(&self, did: &str, rkey: &str, title: &str, body: &str) { - self.index_issue_at(did, rkey, title, body, None, None).await; + self.index_issue_at(did, rkey, title, body, None, None) + .await; } async fn index_issue_at( @@ -505,10 +506,8 @@ async fn second_query_short_circuits_via_lru_without_re_querying_slingshot() { #[tokio::test] async fn author_filter_narrows_to_matching_did() { let h = Harness::new().await; - h.index_issue("did:plc:nel", "i1", "kelp tide", "") - .await; - h.index_issue("did:plc:teq", "i2", "kelp wave", "") - .await; + h.index_issue("did:plc:nel", "i1", "kelp tide", "").await; + h.index_issue("did:plc:teq", "i2", "kelp wave", "").await; let app = router(h.state.clone()); let resp = app .oneshot(search_request(&[("q", "kelp"), ("author", "did:plc:nel")])) @@ -617,4 +616,3 @@ async fn since_after_until_returns_400() { .unwrap(); assert_eq!(resp.status(), StatusCode::BAD_REQUEST); } -