From e8b555ba5e856d3076b9449eb35f58406cd506a6 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sun, 10 May 2026 00:07:16 +0300 Subject: [PATCH] fix(bobbin-sim): added the actual tests I wanted Lewis: May this revision serve well! --- crates/bobbin-sim/src/determinism.rs | 31 +++ crates/bobbin-sim/src/main.rs | 29 +++ crates/bobbin-sim/src/report.rs | 4 + crates/bobbin-sim/src/runtime.rs | 13 +- .../src/workloads/cancel_mid_hydration.rs | 3 + .../workloads/cold_start_under_live_load.rs | 3 + .../concurrent_reads_during_replay.rs | 3 + .../bobbin-sim/src/workloads/frame_burst.rs | 3 + .../workloads/hydrant_disconnect_barrage.rs | 3 + .../src/workloads/slingshot_flap.rs | 120 +++++++----- crates/bobbin-sim/tests/reconnect_replay.rs | 68 +++++++ crates/bobbin-sim/tests/warming_shadow.rs | 66 +++++++ crates/bobbin/src/main.rs | 2 + crates/ingest/examples/smoke.rs | 2 + crates/ingest/src/lib.rs | 152 +++++++++++++-- crates/ingest/src/resolver.rs | 12 ++ crates/ingest/src/shadow.rs | 181 ++++++++++++++++++ 17 files changed, 633 insertions(+), 62 deletions(-) create mode 100644 crates/bobbin-sim/tests/reconnect_replay.rs create mode 100644 crates/bobbin-sim/tests/warming_shadow.rs create mode 100644 crates/ingest/src/shadow.rs diff --git a/crates/bobbin-sim/src/determinism.rs b/crates/bobbin-sim/src/determinism.rs index 2dfec0f..b808771 100644 --- a/crates/bobbin-sim/src/determinism.rs +++ b/crates/bobbin-sim/src/determinism.rs @@ -1,6 +1,7 @@ use std::num::NonZeroUsize; use std::time::Duration; +use bobbin_ingest::{DisconnectSnapshot, WarmingShadowSnapshot}; use tokio::runtime::Builder as TokioBuilder; use tracing_subscriber::Registry; use tracing_subscriber::layer::SubscriberExt; @@ -49,6 +50,18 @@ pub enum LeakOutcome { first: u64, second: u64, }, + DisconnectCountMismatch { + first: u64, + second: u64, + }, + LastDisconnectMismatch { + first: Option, + second: Option, + }, + WarmingShadowMismatch { + first: WarmingShadowSnapshot, + second: WarmingShadowSnapshot, + }, TraceLengthMismatch { first: usize, second: usize, @@ -181,6 +194,24 @@ fn compare_runs( second: b.consumer_too_slow_count, }; } + if a.disconnect_count != b.disconnect_count { + return LeakOutcome::DisconnectCountMismatch { + first: a.disconnect_count, + second: b.disconnect_count, + }; + } + if a.last_disconnect != b.last_disconnect { + return LeakOutcome::LastDisconnectMismatch { + first: a.last_disconnect.clone(), + second: b.last_disconnect.clone(), + }; + } + if a.warming_shadow != b.warming_shadow { + return LeakOutcome::WarmingShadowMismatch { + first: a.warming_shadow, + second: b.warming_shadow, + }; + } if a_lines.len() != b_lines.len() { return LeakOutcome::TraceLengthMismatch { first: a_lines.len(), diff --git a/crates/bobbin-sim/src/main.rs b/crates/bobbin-sim/src/main.rs index 3610017..38239df 100644 --- a/crates/bobbin-sim/src/main.rs +++ b/crates/bobbin-sim/src/main.rs @@ -255,6 +255,20 @@ fn report_as_json(r: &SimReport) -> serde_json::Value { "resolver_hits": r.resolver_hits, "resolver_misses": r.resolver_misses, "consumer_too_slow_count": r.consumer_too_slow_count, + "disconnect_count": r.disconnect_count, + "last_disconnect": r.last_disconnect.as_ref().map(|d| serde_json::json!({ + "kind": format!("{:?}", d.kind), + "message": d.message, + "at_unix_micros": d.at_unix_micros.raw(), + "last_cursor": d.last_cursor.raw(), + })), + "warming_shadow": { + "enqueued_total": r.warming_shadow.enqueued_total, + "drained_via_observe_total": r.warming_shadow.drained_via_observe_total, + "max_concurrent": r.warming_shadow.max_concurrent, + "residual": r.warming_shadow.residual, + "distinct_keys_seen": r.warming_shadow.distinct_keys_seen, + }, "failure_reason": r.failure_reason, }) } @@ -308,6 +322,21 @@ fn leak_outcome_as_json(outcome: &LeakOutcome) -> serde_json::Value { "first": first, "second": second, }), + LeakOutcome::DisconnectCountMismatch { first, second } => serde_json::json!({ + "kind": "disconnect_count_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::LastDisconnectMismatch { first, second } => serde_json::json!({ + "kind": "last_disconnect_mismatch", + "first": first.as_ref().map(|d| format!("{:?}: {}", d.kind, d.message)), + "second": second.as_ref().map(|d| format!("{:?}: {}", d.kind, d.message)), + }), + LeakOutcome::WarmingShadowMismatch { first, second } => serde_json::json!({ + "kind": "warming_shadow_mismatch", + "first": format!("{first:?}"), + "second": format!("{second:?}"), + }), LeakOutcome::TraceLengthMismatch { first, second } => serde_json::json!({ "kind": "trace_length_mismatch", "first": first, diff --git a/crates/bobbin-sim/src/report.rs b/crates/bobbin-sim/src/report.rs index 21053c0..5d36b7f 100644 --- a/crates/bobbin-sim/src/report.rs +++ b/crates/bobbin-sim/src/report.rs @@ -1,5 +1,6 @@ use std::time::Duration; +use bobbin_ingest::{DisconnectSnapshot, WarmingShadowSnapshot}; use bobbin_runtime::UnixMicros; #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -22,6 +23,9 @@ pub struct SimReport { pub resolver_hits: u64, pub resolver_misses: u64, pub consumer_too_slow_count: u64, + pub disconnect_count: u64, + pub last_disconnect: Option, + pub warming_shadow: WarmingShadowSnapshot, pub failure_reason: Option, } diff --git a/crates/bobbin-sim/src/runtime.rs b/crates/bobbin-sim/src/runtime.rs index 15b422e..a347b37 100644 --- a/crates/bobbin-sim/src/runtime.rs +++ b/crates/bobbin-sim/src/runtime.rs @@ -5,7 +5,8 @@ use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; use bobbin_ingest::{ - DEFAULT_INGEST_PARALLELISM, IngestConfig, IngestRuntime, RepoIdResolver, run as run_ingest, + DEFAULT_INGEST_PARALLELISM, DisconnectSink, IngestConfig, IngestRuntime, RepoIdResolver, + WarmingShadowBuffer, run as run_ingest, }; use bobbin_record_lru::NoopRecordStore; use bobbin_runtime::{ @@ -75,6 +76,8 @@ impl Sim { let records = Arc::new(NoopRecordStore); let cancel = CancellationToken::new(); let consumer_too_slow_count = Arc::new(AtomicU64::new(0)); + let disconnects = Arc::new(DisconnectSink::new()); + let warming_shadow = Arc::new(WarmingShadowBuffer::new(hasher.clone())); let workload_name = self.workload.name(); let ctx = WorkloadCtx { @@ -111,6 +114,8 @@ impl Sim { entropy: entropy.clone(), ws: mem_ws, cancel: cancel.clone(), + disconnects: Some(disconnects.clone()), + warming_shadow: Some(warming_shadow.clone()), }; let ingest_config = IngestConfig { hydrant_base, @@ -136,6 +141,9 @@ impl Sim { resolver_hits: resolver.stats().hits, resolver_misses: resolver.stats().miss_count(), consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), + disconnect_count: disconnects.count(), + last_disconnect: disconnects.snapshot(), + warming_shadow: warming_shadow.snapshot(), failure_reason: Some(format!( "max_virtual_runtime {max_virtual_runtime:?} exhausted", )), @@ -146,6 +154,9 @@ impl Sim { resolver_hits: stats.hits, resolver_misses: stats.miss_count(), consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), + disconnect_count: disconnects.count(), + last_disconnect: disconnects.snapshot(), + warming_shadow: warming_shadow.snapshot(), ..r } } diff --git a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs index 312504e..2fbfa16 100644 --- a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs +++ b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs @@ -154,6 +154,9 @@ impl Workload for CancelMidHydration { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: 0, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); 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 1ab7a84..eea780c 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 @@ -120,6 +120,9 @@ impl Workload for ColdStartUnderLiveLoad { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: 0, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); 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 6c558e0..a4ba397 100644 --- a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs +++ b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs @@ -178,6 +178,9 @@ impl Workload for ConcurrentReadsDuringReplay { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: 0, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); diff --git a/crates/bobbin-sim/src/workloads/frame_burst.rs b/crates/bobbin-sim/src/workloads/frame_burst.rs index 3d31fe5..16dde34 100644 --- a/crates/bobbin-sim/src/workloads/frame_burst.rs +++ b/crates/bobbin-sim/src/workloads/frame_burst.rs @@ -104,6 +104,9 @@ impl Workload for FrameBurst { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: 0, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); diff --git a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs index cf09c77..155dc1c 100644 --- a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs +++ b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs @@ -129,6 +129,9 @@ impl Workload for HydrantDisconnectBarrage { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: 0, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); diff --git a/crates/bobbin-sim/src/workloads/slingshot_flap.rs b/crates/bobbin-sim/src/workloads/slingshot_flap.rs index 6019921..a5a3aee 100644 --- a/crates/bobbin-sim/src/workloads/slingshot_flap.rs +++ b/crates/bobbin-sim/src/workloads/slingshot_flap.rs @@ -1,5 +1,4 @@ use std::sync::Arc; -use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; @@ -77,12 +76,14 @@ impl Workload for SlingshotFlap { let total_frames = (star_count as u64) * 2; let target_owner_count = star_count; - let frame_script = build_frame_script(&cfg); + let frame_script = Arc::new(build_frame_script(&cfg)); + let session_count = Arc::new(AtomicU64::new(0)); let hydrant: Arc = Arc::new(FlapHydrant { - frames: Mutex::new(Some(frame_script)), + frames: frame_script, send_timeout: Duration::from_millis(cfg.hydrant_send_timeout_ms), frame_pace: Duration::from_micros(cfg.hydrant_frame_pace_us), consumer_too_slow: consumer_too_slow_count.clone(), + session_count: session_count.clone(), clock: clock.clone(), }); @@ -94,6 +95,7 @@ impl Workload for SlingshotFlap { }); let cts_counter = consumer_too_slow_count.clone(); + let sessions_counter = session_count.clone(); let script = Box::pin(async move { let started = started_unix; @@ -119,15 +121,16 @@ impl Workload for SlingshotFlap { 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 { SimOutcome::Passed => None, SimOutcome::Failed if cts > 0 => Some(format!( - "ConsumerTooSlow fired {cts} times; processed {}/{} events before disconnect", + "ConsumerTooSlow fired {cts} times across {sessions} sessions; processed {}/{} events before disconnect", snap.events_processed(), total_frames, )), SimOutcome::Failed => Some(format!( - "only {}/{} events processed (expected {} owners staged)", + "only {}/{} events processed across {sessions} sessions (expected {} owners staged)", snap.events_processed(), total_frames, target_owner_count, @@ -146,6 +149,9 @@ impl Workload for SlingshotFlap { resolver_hits: 0, resolver_misses: 0, consumer_too_slow_count: cts, + disconnect_count: 0, + last_disconnect: None, + warming_shadow: Default::default(), failure_reason, } }); @@ -158,82 +164,91 @@ impl Workload for SlingshotFlap { } } -fn build_frame_script(cfg: &SlingshotFlapConfig) -> Vec { +fn build_frame_script(cfg: &SlingshotFlapConfig) -> Vec<(u64, String)> { let mut frames = Vec::with_capacity(cfg.cross_did_stars * 2); for i in 0..cfg.cross_did_stars { + let id = (i + 1) as u64; let star_did = format!("did:plc:starer-{i}"); let target_owner = format!("did:plc:owner-{i}"); let target_rkey = format_rkey(i); - frames.push( - serde_json::json!({ - "id": (i + 1) as u64, - "type": "record", + let body = serde_json::json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": star_did, + "rev": format_tid(i), + "collection": "sh.tangled.feed.star", + "rkey": format_rkey(i + 1_000_000), + "action": "create", "record": { - "live": false, - "did": star_did, - "rev": format_tid(i), - "collection": "sh.tangled.feed.star", - "rkey": format_rkey(i + 1_000_000), - "action": "create", - "record": { - "$type": "sh.tangled.feed.star", - "createdAt": "2026-05-01T00:00:00Z", - "subject": format!("at://{target_owner}/{REPO_COLLECTION}/{target_rkey}") - } + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": format!("at://{target_owner}/{REPO_COLLECTION}/{target_rkey}") } - }) - .to_string(), - ); + } + }) + .to_string(); + frames.push((id, body)); } for i in 0..cfg.cross_did_stars { + let id = (cfg.cross_did_stars + i + 1) as u64; let owner = format!("did:plc:owner-{i}"); let repo_did = format!("did:plc:repo-{i}"); - frames.push( - serde_json::json!({ - "id": (cfg.cross_did_stars + i + 1) as u64, - "type": "record", + let body = serde_json::json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": owner, + "rev": format_tid(cfg.cross_did_stars + i), + "collection": REPO_COLLECTION, + "rkey": format_rkey(i), + "action": "create", "record": { - "live": false, - "did": owner, - "rev": format_tid(cfg.cross_did_stars + i), - "collection": REPO_COLLECTION, - "rkey": format_rkey(i), - "action": "create", - "record": { - "$type": REPO_COLLECTION, - "createdAt": "2026-05-01T00:00:00Z", - "knot": "oyster.cafe", - "name": format!("repo-{i}"), - "repoDid": repo_did, - } + "$type": REPO_COLLECTION, + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": format!("repo-{i}"), + "repoDid": repo_did, } - }) - .to_string(), - ); + } + }) + .to_string(); + frames.push((id, body)); } frames } +fn parse_cursor(url: &Url) -> u64 { + url.query_pairs() + .find_map(|(k, v)| (k == "cursor").then(|| v.parse().ok()).flatten()) + .unwrap_or(0) +} + struct FlapHydrant { - frames: Mutex>>, + frames: Arc>, send_timeout: Duration, frame_pace: Duration, consumer_too_slow: Arc, + session_count: Arc, clock: Arc, } impl MemWsResponder for FlapHydrant { fn spawn_server( &self, - _: Url, + url: Url, mut recv: mpsc::UnboundedReceiver, send: mpsc::Sender, ) -> MemWsServerFuture { - let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + let cursor = parse_cursor(&url); + let frames = self.frames.clone(); let send_timeout = self.send_timeout; let frame_pace = self.frame_pace; let consumer_too_slow = self.consumer_too_slow.clone(); let clock = self.clock.clone(); + self.session_count.fetch_add(1, Ordering::Relaxed); let pong_send = send.clone(); Box::pin(async move { let pong_loop = async move { @@ -251,11 +266,18 @@ impl MemWsResponder for FlapHydrant { }; let clock_for_emit = clock.clone(); let frame_emit = async move { - for text in frames { + for (id, text) in frames.iter() { + if *id < cursor { + continue; + } if !frame_pace.is_zero() { clock_for_emit.sleep(frame_pace).await; } - let res = tokio::time::timeout(send_timeout, send.send(WsMessage::Text(text))).await; + let res = tokio::time::timeout( + send_timeout, + send.send(WsMessage::Text(text.clone())), + ) + .await; match res { Ok(Ok(())) => continue, Ok(Err(_)) => return, diff --git a/crates/bobbin-sim/tests/reconnect_replay.rs b/crates/bobbin-sim/tests/reconnect_replay.rs new file mode 100644 index 0000000..f652ac3 --- /dev/null +++ b/crates/bobbin-sim/tests/reconnect_replay.rs @@ -0,0 +1,68 @@ +use std::num::NonZeroUsize; +use std::time::Duration; + +use bobbin_ingest::DisconnectKind; +use bobbin_sim::workloads::{SlingshotFlap, SlingshotFlapConfig}; +use bobbin_sim::{Sim, SimConfig, SimOutcome}; +use tokio::runtime::Builder as TokioBuilder; + +#[test] +fn slingshot_flap_resumes_via_replay_after_pong_timeout() { + let cfg = SlingshotFlapConfig { + cross_did_stars: 5000, + normal_latency_ms: 2, + brownout_start_ms: 0, + brownout_duration_ms: 60_000, + brownout_latency_ms: 200, + brownout_enabled: true, + hydrant_send_timeout_ms: 30_000, + hydrant_frame_pace_us: 1_000, + }; + let mut sim_config = SimConfig::new(7); + sim_config.parallelism = NonZeroUsize::new(4).unwrap(); + sim_config.max_virtual_runtime = Duration::from_secs(600); + + let workload = Box::new(SlingshotFlap::new(cfg.clone())); + let runtime = TokioBuilder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("build current_thread runtime"); + let report = runtime.block_on(Sim::new(sim_config, workload).run()); + + let total_frames = (cfg.cross_did_stars as u64) * 2; + + assert_eq!( + report.outcome, + SimOutcome::Passed, + "expected eventual recovery via replay after disconnect, got {:?} reason={:?}", + report.outcome, + report.failure_reason, + ); + assert_eq!( + report.events_processed, total_frames, + "expected all {total_frames} frames processed after reconnect, got {}", + report.events_processed, + ); + assert!( + report.disconnect_count >= 1, + "expected at least one disconnect during sustained brownout, got {}", + report.disconnect_count, + ); + let last = report + .last_disconnect + .as_ref() + .expect("last_disconnect populated when disconnect_count > 0"); + assert_eq!( + last.kind, + DisconnectKind::PongTimeout, + "expected PongTimeout under sustained brownout (paced upstream + N=4 + 200ms slingshot RTT), got {:?}: {}", + last.kind, + last.message, + ); + assert!( + last.last_cursor.raw() > 0 && last.last_cursor.raw() < total_frames, + "expected disconnect mid-stream, got last_cursor={}", + last.last_cursor.raw(), + ); +} diff --git a/crates/bobbin-sim/tests/warming_shadow.rs b/crates/bobbin-sim/tests/warming_shadow.rs new file mode 100644 index 0000000..d8139a9 --- /dev/null +++ b/crates/bobbin-sim/tests/warming_shadow.rs @@ -0,0 +1,66 @@ +use std::num::NonZeroUsize; +use std::time::Duration; + +use bobbin_sim::workloads::{SlingshotFlap, SlingshotFlapConfig}; +use bobbin_sim::{Sim, SimConfig, SimOutcome}; +use tokio::runtime::Builder as TokioBuilder; + +#[test] +fn shadow_buffer_bounded_by_cross_did_attachment_cohort_and_drains_to_zero() { + let cfg = SlingshotFlapConfig { + cross_did_stars: 5000, + normal_latency_ms: 2, + brownout_start_ms: 0, + brownout_duration_ms: 0, + brownout_latency_ms: 0, + brownout_enabled: false, + hydrant_send_timeout_ms: 30_000, + hydrant_frame_pace_us: 1_000, + }; + let mut sim_config = SimConfig::new(7); + sim_config.parallelism = NonZeroUsize::new(16).unwrap(); + sim_config.max_virtual_runtime = Duration::from_secs(120); + + let workload = Box::new(SlingshotFlap::new(cfg.clone())); + let runtime = TokioBuilder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("build current_thread runtime"); + let report = runtime.block_on(Sim::new(sim_config, workload).run()); + + let stars = cfg.cross_did_stars as u64; + let total = stars * 2; + + assert_eq!( + report.outcome, + SimOutcome::Passed, + "expected clean drain with brownout disabled, got {:?} reason={:?}", + report.outcome, + report.failure_reason, + ); + assert_eq!(report.events_processed, total); + + let s = report.warming_shadow; + assert_eq!( + s.enqueued_total, stars, + "every cross-DID star must enqueue exactly once: {s:?}", + ); + assert_eq!( + s.drained_via_observe_total, stars, + "every observe must drain its corresponding enqueue (stars send their repos in this workload): {s:?}", + ); + assert_eq!( + s.distinct_keys_seen, stars, + "one (owner, rkey) pair per repo: {s:?}", + ); + assert_eq!( + s.residual, 0, + "no entry must remain pending after observe arrives: {s:?}", + ); + assert!( + s.max_concurrent > 0 && s.max_concurrent <= stars, + "peak depth must be bounded by the cross-DID attachment cohort, got max_concurrent={} for {stars} stars: {s:?}", + s.max_concurrent, + ); +} diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs index ae56be3..34d129a 100644 --- a/crates/bobbin/src/main.rs +++ b/crates/bobbin/src/main.rs @@ -169,6 +169,8 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { entropy, ws, cancel: cancel.clone(), + disconnects: None, + warming_shadow: None, }; 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 c2ebcd6..e2ac8d9 100644 --- a/crates/ingest/examples/smoke.rs +++ b/crates/ingest/examples/smoke.rs @@ -41,6 +41,8 @@ async fn main() { entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), cancel: cancel.clone(), + disconnects: None, + warming_shadow: 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 35e917d..14d3d3a 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -28,9 +28,11 @@ use url::Url; mod frame; mod resolver; +mod shadow; pub use frame::{FrameKind, HydrantFrame, RecordAction, RecordFrame}; use frame::HydrantStreamErrorFrame; pub use resolver::{RepoIdResolver, Resolution}; +pub use shadow::{WarmingShadowBuffer, WarmingShadowSnapshot}; const TANGLED_PREFIX: &str = "sh.tangled."; const RECONNECT_INITIAL_DELAY: Duration = Duration::from_millis(500); @@ -109,6 +111,77 @@ pub enum IngestError { HydrantStream { code: String, message: String }, } +#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash)] +pub enum DisconnectKind { + Url, + UnknownScheme, + Network, + Decode, + InvalidAtUri, + Extract, + PongTimeout, + SendTimeout, + ConsumerTooSlow, + HydrantStream, +} + +impl DisconnectKind { + pub fn from_error(err: &IngestError) -> Self { + match err { + IngestError::Url(_) => Self::Url, + IngestError::UnknownScheme(_) => Self::UnknownScheme, + IngestError::Network(_) => Self::Network, + IngestError::Decode(_) => Self::Decode, + IngestError::InvalidAtUri(_) => Self::InvalidAtUri, + IngestError::Extract(_) => Self::Extract, + IngestError::PongTimeout(_) => Self::PongTimeout, + IngestError::SendTimeout(_) => Self::SendTimeout, + IngestError::ConsumerTooSlow { .. } => Self::ConsumerTooSlow, + IngestError::HydrantStream { .. } => Self::HydrantStream, + } + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DisconnectSnapshot { + pub kind: DisconnectKind, + pub message: String, + pub at_unix_micros: UnixMicros, + pub last_cursor: HydrantCursor, +} + +#[derive(Default)] +pub struct DisconnectSink { + last: std::sync::Mutex>, + count: std::sync::atomic::AtomicU64, +} + +impl DisconnectSink { + pub fn new() -> Self { + Self::default() + } + + pub fn record(&self, snap: DisconnectSnapshot) { + *self + .last + .lock() + .expect("disconnect sink mutex poisoned") = Some(snap); + self.count + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } + + pub fn snapshot(&self) -> Option { + self.last + .lock() + .expect("disconnect sink mutex poisoned") + .clone() + } + + pub fn count(&self) -> u64 { + self.count.load(std::sync::atomic::Ordering::Relaxed) + } +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum SessionOutcome { Progressed, @@ -131,6 +204,8 @@ pub struct IngestRuntime { pub entropy: Arc, pub ws: Arc, pub cancel: CancellationToken, + pub disconnects: Option>, + pub warming_shadow: Option>, } impl Clone for IngestRuntime { @@ -145,6 +220,8 @@ impl Clone for IngestRuntime { entropy: self.entropy.clone(), ws: self.ws.clone(), cancel: self.cancel.clone(), + disconnects: self.disconnects.clone(), + warming_shadow: self.warming_shadow.clone(), } } } @@ -185,6 +262,14 @@ async fn run_inner( } (_, Some(err)) => warn!(?err, "hydrant stream errored"), } + if let (Some(sink), Some(err)) = (runtime.disconnects.as_ref(), error.as_ref()) { + sink.record(DisconnectSnapshot { + kind: DisconnectKind::from_error(err), + message: err.to_string(), + at_unix_micros: runtime.clock.now_unix_micros(), + last_cursor: runtime.coverage.snapshot().last_cursor(), + }); + } let made_progress = matches!(outcome, SessionOutcome::Progressed); if made_progress { backoff = RECONNECT_INITIAL_DELAY; @@ -597,7 +682,13 @@ 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, now).await; + let pending = prepare_frame( + frame, + &rt.resolver, + rt.warming_shadow.as_deref(), + now, + ) + .await; let prepare_end = rt.clock.now_instant(); Prepared { pending, @@ -611,7 +702,13 @@ async fn resolve_stage( rt: IngestRuntime, ) -> Resolved { let resolve_start = rt.clock.now_instant(); - let pending = resolve_pending(staged.pending, &rt.resolver).await; + let pending = resolve_pending( + staged.pending, + &rt.resolver, + &rt.coverage, + rt.warming_shadow.as_deref(), + ) + .await; let resolve_end = rt.clock.now_instant(); Resolved { pending, @@ -661,6 +758,7 @@ async fn commit_stage( async fn prepare_frame( frame: HydrantFrame, resolver: &RepoIdResolver, + shadow: Option<&WarmingShadowBuffer>, now: UnixMicros, ) -> Pending { let cursor = HydrantCursor::new(frame.id); @@ -671,7 +769,7 @@ async fn prepare_frame( None => Regime::NonRecord, }; let op = match frame.kind { - FrameKind::Record => prepare_record(frame.record, resolver).await, + FrameKind::Record => prepare_record(frame.record, resolver, shadow).await, FrameKind::Identity | FrameKind::Account => PendingOp::Noop, FrameKind::Other => { debug!(id = frame.id, "ignoring unknown hydrant frame kind"); @@ -686,7 +784,11 @@ async fn prepare_frame( } } -async fn prepare_record(record: Option, resolver: &RepoIdResolver) -> PendingOp { +async fn prepare_record( + record: Option, + resolver: &RepoIdResolver, + shadow: Option<&WarmingShadowBuffer>, +) -> PendingOp { let Some(record) = record else { debug!("record-typed frame missing payload, skipping"); return PendingOp::Noop; @@ -721,6 +823,9 @@ async fn prepare_record(record: Option, resolver: &RepoIdResolver) } }; if let Record::Repo(repo) = &parsed { + if let Some(shadow) = shadow { + shadow.note_observed(&record.did, &record.rkey).await; + } resolver .observe( record.did.clone(), @@ -753,7 +858,12 @@ async fn prepare_record(record: Option, resolver: &RepoIdResolver) } } -async fn resolve_pending(pending: Pending, resolver: &RepoIdResolver) -> Pending { +async fn resolve_pending( + pending: Pending, + resolver: &RepoIdResolver, + coverage: &CoverageWatch, + shadow: Option<&WarmingShadowBuffer>, +) -> Pending { let Pending { cursor, signal, @@ -769,7 +879,7 @@ async fn resolve_pending(pending: Pending, resolver: &RepoIdResolver) -> Pending cid, edges, } => { - let edges = normalize_subjects(edges, resolver).await; + let edges = normalize_subjects(edges, resolver, coverage, shadow).await; PendingOp::Upsert { source, nsid, @@ -840,8 +950,8 @@ async fn handle_frame( now: UnixMicros, ) { let _ = clock; - let pending = prepare_frame(frame, resolver, now).await; - let pending = resolve_pending(pending, resolver).await; + let pending = prepare_frame(frame, resolver, None, now).await; + let pending = resolve_pending(pending, resolver, coverage, None).await; commit_pending(pending, store, coverage, search, records).await; } @@ -873,12 +983,26 @@ fn cache_body( } } -async fn normalize_subjects(edges: Vec, resolver: &RepoIdResolver) -> Vec { +async fn normalize_subjects( + edges: Vec, + resolver: &RepoIdResolver, + coverage: &CoverageWatch, + shadow: Option<&WarmingShadowBuffer>, +) -> Vec { + let warming = shadow.is_some() && !coverage.snapshot().is_ready(); futures::stream::iter(edges) .then(|edge| async move { let Some((owner, rkey)) = parse_repo_subject(&edge.subject) else { return edge; }; + if warming + && let Some(shadow) = shadow + && resolver.cached_resolution(&owner, &rkey).await.is_none() + { + shadow + .note_unresolved(owner.clone(), rkey.clone()) + .await; + } match resolver.resolve(&owner, &rkey).await { Resolution::Mapped(repo_did) => Edge { subject: owned_did_aturi(&repo_did), @@ -1129,16 +1253,16 @@ mod tests { } })) }; - let live_pending = prepare_frame(mk(true), &resolver, now()).await; + let live_pending = prepare_frame(mk(true), &resolver, None, now()).await; assert_eq!(live_pending.regime, Regime::Live); - let replay_pending = prepare_frame(mk(false), &resolver, now()).await; + let replay_pending = prepare_frame(mk(false), &resolver, None, 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, now()).await; + let id_pending = prepare_frame(identity, &resolver, None, now()).await; assert_eq!(id_pending.regime, Regime::NonRecord); } @@ -1718,6 +1842,8 @@ mod tests { entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), cancel, + disconnects: None, + warming_shadow: None, } } @@ -2226,6 +2352,8 @@ mod tests { entropy: Arc::new(OsEntropy), ws: TungsteniteWs::shared(), cancel: CancellationToken::new(), + disconnects: None, + warming_shadow: None, }; let parallelism = 4usize; diff --git a/crates/ingest/src/resolver.rs b/crates/ingest/src/resolver.rs index 8b5ac7e..39aa10e 100644 --- a/crates/ingest/src/resolver.rs +++ b/crates/ingest/src/resolver.rs @@ -190,6 +190,18 @@ impl RepoIdResolver { self.stats.snapshot() } + pub async fn cached_resolution( + &self, + owner: &Did, + rkey: &Rkey, + ) -> Option { + let key = RepoRef::new(owner.clone(), rkey.clone()); + self.cache + .get_async(&key) + .await + .map(|entry| entry.get().clone().into_resolution()) + } + pub async fn observe( &self, owner: Did, diff --git a/crates/ingest/src/shadow.rs b/crates/ingest/src/shadow.rs new file mode 100644 index 0000000..bfc5da1 --- /dev/null +++ b/crates/ingest/src/shadow.rs @@ -0,0 +1,181 @@ +use std::sync::atomic::{AtomicU64, Ordering}; + +use bobbin_runtime::RuntimeHasher; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::recordkey::Rkey; +use scc::HashMap as SccMap; + +#[derive(Clone, Debug, Eq, Hash, PartialEq)] +struct ShadowKey { + owner: Did, + rkey: Rkey, +} + +pub struct WarmingShadowBuffer { + pending: SccMap, + enqueued_total: AtomicU64, + drained_via_observe_total: AtomicU64, + max_concurrent: AtomicU64, + current_concurrent: AtomicU64, + distinct_keys_seen: AtomicU64, +} + +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub struct WarmingShadowSnapshot { + pub enqueued_total: u64, + pub drained_via_observe_total: u64, + pub max_concurrent: u64, + pub residual: u64, + pub distinct_keys_seen: u64, +} + +impl WarmingShadowBuffer { + pub fn new(hasher: RuntimeHasher) -> Self { + Self { + pending: SccMap::with_hasher(hasher), + enqueued_total: AtomicU64::new(0), + drained_via_observe_total: AtomicU64::new(0), + max_concurrent: AtomicU64::new(0), + current_concurrent: AtomicU64::new(0), + distinct_keys_seen: AtomicU64::new(0), + } + } + + pub async fn note_unresolved(&self, owner: Did, rkey: Rkey) { + let key = ShadowKey { owner, rkey }; + let mut entry = self.pending.entry_async(key).await.or_insert_with(|| { + self.distinct_keys_seen.fetch_add(1, Ordering::Relaxed); + 0 + }); + *entry.get_mut() += 1; + self.enqueued_total.fetch_add(1, Ordering::Relaxed); + let cur = self.current_concurrent.fetch_add(1, Ordering::Relaxed) + 1; + self.max_concurrent.fetch_max(cur, Ordering::Relaxed); + } + + pub async fn note_observed(&self, owner: &Did, rkey: &Rkey) { + let key = ShadowKey { + owner: owner.clone(), + rkey: rkey.clone(), + }; + let Some(mut entry) = self.pending.get_async(&key).await else { + return; + }; + let count = std::mem::replace(entry.get_mut(), 0); + if count > 0 { + self.drained_via_observe_total + .fetch_add(count, Ordering::Relaxed); + self.current_concurrent + .fetch_sub(count, Ordering::Relaxed); + } + } + + pub fn snapshot(&self) -> WarmingShadowSnapshot { + WarmingShadowSnapshot { + enqueued_total: self.enqueued_total.load(Ordering::Relaxed), + drained_via_observe_total: self.drained_via_observe_total.load(Ordering::Relaxed), + max_concurrent: self.max_concurrent.load(Ordering::Relaxed), + residual: self.current_concurrent.load(Ordering::Relaxed), + distinct_keys_seen: self.distinct_keys_seen.load(Ordering::Relaxed), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use jacquard_common::types::did::Did; + use jacquard_common::types::recordkey::Rkey; + + fn d(s: &str) -> Did { + Did::new_owned(s).unwrap() + } + + fn r(s: &str) -> Rkey { + Rkey::new_owned(s).unwrap() + } + + #[tokio::test] + async fn enqueue_then_observe_drains_to_zero() { + let shadow = WarmingShadowBuffer::new(RuntimeHasher::default()); + shadow + .note_unresolved(d("did:plc:nel"), r("abcabcabcabcz")) + .await; + shadow + .note_unresolved(d("did:plc:nel"), r("abcabcabcabcz")) + .await; + shadow + .note_observed(&d("did:plc:nel"), &r("abcabcabcabcz")) + .await; + let s = shadow.snapshot(); + assert_eq!(s.enqueued_total, 2); + assert_eq!(s.drained_via_observe_total, 2); + assert_eq!(s.max_concurrent, 2); + assert_eq!(s.residual, 0); + assert_eq!(s.distinct_keys_seen, 1); + } + + #[tokio::test] + async fn distinct_keys_tracked_separately() { + let shadow = WarmingShadowBuffer::new(RuntimeHasher::default()); + shadow + .note_unresolved(d("did:plc:nel"), r("abcabcabcabcz")) + .await; + shadow + .note_unresolved(d("did:plc:olaren"), r("abcabcabcabd1")) + .await; + let s = shadow.snapshot(); + assert_eq!(s.enqueued_total, 2); + assert_eq!(s.distinct_keys_seen, 2); + assert_eq!(s.max_concurrent, 2); + assert_eq!(s.residual, 2); + } + + #[tokio::test] + async fn observe_without_prior_enqueue_is_a_noop() { + let shadow = WarmingShadowBuffer::new(RuntimeHasher::default()); + shadow + .note_observed(&d("did:plc:nel"), &r("abcabcabcabcz")) + .await; + let s = shadow.snapshot(); + assert_eq!(s.enqueued_total, 0); + assert_eq!(s.drained_via_observe_total, 0); + assert_eq!(s.distinct_keys_seen, 0); + } + + #[tokio::test] + async fn second_observe_after_drain_does_not_double_count() { + let shadow = WarmingShadowBuffer::new(RuntimeHasher::default()); + shadow + .note_unresolved(d("did:plc:nel"), r("abcabcabcabcz")) + .await; + shadow + .note_observed(&d("did:plc:nel"), &r("abcabcabcabcz")) + .await; + shadow + .note_observed(&d("did:plc:nel"), &r("abcabcabcabcz")) + .await; + let s = shadow.snapshot(); + assert_eq!(s.drained_via_observe_total, 1); + assert_eq!(s.residual, 0); + } + + #[tokio::test] + async fn residual_reflects_unobserved_keys() { + let shadow = WarmingShadowBuffer::new(RuntimeHasher::default()); + shadow + .note_unresolved(d("did:plc:nel"), r("abcabcabcabcz")) + .await; + shadow + .note_unresolved(d("did:plc:olaren"), r("abcabcabcabd1")) + .await; + shadow + .note_observed(&d("did:plc:nel"), &r("abcabcabcabcz")) + .await; + let s = shadow.snapshot(); + assert_eq!(s.enqueued_total, 2); + assert_eq!(s.drained_via_observe_total, 1); + assert_eq!(s.residual, 1); + } +} -- 2.51.2