diff --git a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs new file mode 100644 index 0000000..cf09c77 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs @@ -0,0 +1,248 @@ +use std::sync::Arc; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use bobbin_runtime::{ + MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage, +}; +use tokio::sync::mpsc; +use url::Url; + +use crate::report::{SimOutcome, SimReport}; +use crate::workload::{Workload, WorkloadCtx, WorkloadHooks}; +use crate::workloads::util::{AssertNoSlingshot, format_rkey, format_tid}; + +const NAME: &str = "hydrant-disconnect-barrage"; +const AUTHORITY_DID: &str = "did:plc:abalone"; + +#[derive(Clone, Debug)] +pub struct HydrantDisconnectBarrageConfig { + pub frames: usize, + pub disconnect_after_frames_per_session: Vec, +} + +impl Default for HydrantDisconnectBarrageConfig { + fn default() -> Self { + Self { + frames: 256, + disconnect_after_frames_per_session: vec![50, 80, 60], + } + } +} + +pub struct HydrantDisconnectBarrage { + config: HydrantDisconnectBarrageConfig, +} + +impl HydrantDisconnectBarrage { + pub fn new(config: HydrantDisconnectBarrageConfig) -> Self { + Self { config } + } +} + +impl Workload for HydrantDisconnectBarrage { + fn name(&self) -> &'static str { + NAME + } + + fn build(self: Box, ctx: WorkloadCtx) -> WorkloadHooks { + let cfg = self.config.clone(); + let WorkloadCtx { + seed, + clock, + coverage, + store, + cancel, + .. + } = ctx; + + let frames_total = cfg.frames as u64; + let frame_log = build_frame_log(cfg.frames); + let drop_schedule = Arc::new(Mutex::new(cfg.disconnect_after_frames_per_session.clone())); + let session_count = Arc::new(AtomicU64::new(0)); + let disconnect_count = Arc::new(AtomicU64::new(0)); + + let hydrant: Arc = Arc::new(BarrageHydrant { + frame_log: Arc::new(frame_log), + drop_schedule, + session_count: session_count.clone(), + disconnect_count: disconnect_count.clone(), + }); + let slingshot_probe = AssertNoSlingshot::new(); + let slingshot: Arc = Arc::new(slingshot_probe.clone()); + + let disconnects = disconnect_count.clone(); + let sessions = session_count.clone(); + let script = Box::pin(async move { + let started = clock.now_unix_micros(); + let mut rx = coverage.subscribe(); + let outcome = loop { + let snap = *rx.borrow_and_update(); + if snap.events_processed() == frames_total { + break SimOutcome::Passed; + } + if snap.events_processed() > frames_total { + break SimOutcome::Failed; + } + tokio::select! { + res = rx.changed() => match res { + Ok(()) => continue, + Err(_) => break SimOutcome::Failed, + }, + _ = cancel.cancelled() => break SimOutcome::Failed, + } + }; + let snap = coverage.snapshot(); + 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(); + let outcome = match outcome { + SimOutcome::Passed if dc == 0 => SimOutcome::Failed, + SimOutcome::Passed if stray_slingshot != 0 => SimOutcome::Failed, + other => other, + }; + let failure_reason = match outcome { + SimOutcome::Passed => None, + SimOutcome::Failed => Some(format!( + "events={} target={} disconnects={} sessions={} stray_slingshot={}", + snap.events_processed(), + frames_total, + dc, + ss, + stray_slingshot, + )), + SimOutcome::TimedOut => Some("timed out".into()), + }; + SimReport { + workload: NAME, + seed, + outcome, + virtual_runtime, + virtual_clock_end: clock.now_unix_micros(), + events_processed: snap.events_processed(), + last_cursor: snap.last_cursor().raw(), + edge_count: store.key_count() as u64, + resolver_hits: 0, + resolver_misses: 0, + consumer_too_slow_count: 0, + failure_reason, + } + }); + + WorkloadHooks { + slingshot, + hydrant, + script, + } + } +} + +fn build_frame_log(count: usize) -> Vec<(u64, String)> { + (0..count) + .map(|i| { + let id = (i + 1) as u64; + let body = serde_json::json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": AUTHORITY_DID, + "rev": format_tid(i), + "collection": "sh.tangled.repo", + "rkey": format_rkey(i), + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": format!("repo-{i}"), + "repoDid": format!("did:plc:abalone-{i}") + } + } + }) + .to_string(); + (id, body) + }) + .collect() +} + +struct BarrageHydrant { + frame_log: Arc>, + drop_schedule: Arc>>, + session_count: Arc, + disconnect_count: Arc, +} + +impl MemWsResponder for BarrageHydrant { + fn spawn_server( + &self, + url: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let cursor = parse_cursor(&url).unwrap_or(0); + let frame_log = self.frame_log.clone(); + let drop_after = { + let mut guard = self.drop_schedule.lock().unwrap(); + if guard.is_empty() { + usize::MAX + } else { + guard.remove(0) + } + }; + self.session_count.fetch_add(1, Ordering::Relaxed); + let disconnects = self.disconnect_count.clone(); + let pong_send = send.clone(); + Box::pin(async move { + let pong_loop = async move { + loop { + match recv.recv().await { + Some(WsMessage::Ping(payload)) => { + if pong_send.send(WsMessage::Pong(payload)).await.is_err() { + return; + } + } + Some(WsMessage::Close { .. }) | None => return, + Some(_) => {} + } + } + }; + let frame_emit = async move { + let mut emitted = 0usize; + for (id, body) in frame_log.iter() { + if *id < cursor { + continue; + } + if send.send(WsMessage::Text(body.clone())).await.is_err() { + return; + } + emitted += 1; + if emitted >= drop_after { + disconnects.fetch_add(1, Ordering::Relaxed); + let _ = send + .send(WsMessage::Close { + code: 1011, + reason: "scripted disconnect".into(), + }) + .await; + return; + } + } + }; + let _ = tokio::join!(pong_loop, frame_emit); + }) + } +} + +fn parse_cursor(url: &Url) -> Option { + for (k, v) in url.query_pairs() { + if k == "cursor" { + return v.parse().ok(); + } + } + None +} diff --git a/crates/bobbin-sim/tests/determinism_leak.rs b/crates/bobbin-sim/tests/determinism_leak.rs new file mode 100644 index 0000000..c342398 --- /dev/null +++ b/crates/bobbin-sim/tests/determinism_leak.rs @@ -0,0 +1,221 @@ +use std::num::NonZeroUsize; +use std::time::Duration; + +use bobbin_sim::workloads::{ + CancelMidHydration, CancelMidHydrationConfig, ColdStartUnderLiveLoad, + ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, + ConcurrentReadsDuringReplayConfig, FrameBurst, HydrantDisconnectBarrage, + HydrantDisconnectBarrageConfig, SlingshotFlap, SlingshotFlapConfig, +}; +use bobbin_sim::{LeakOutcome, LeakRunConfig, Workload, run_leak_check}; + +const PARALLELISM_SWEEP: &[usize] = &[1, 4, 16, 64]; +const SEEDS: &[u64] = &[1, 7, 42, 1337]; + +#[test] +fn frame_burst_is_byte_deterministic_across_parallelism_sweep() { + for &par in PARALLELISM_SWEEP { + for &seed in SEEDS { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(30), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { Box::new(FrameBurst::new(64)) }; + let result = run_leak_check(config.clone(), factory); + assert!( + result.passed(), + "frame-burst leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + match &result.outcome { + LeakOutcome::Match => {} + other => panic!("non-match outcome despite passed(): {other:?}"), + } + assert_eq!( + result.first_report.events_processed, 64, + "frame-burst should process all 64 frames at par={par} seed={seed}", + ); + } + } +} + +#[test] +fn cancel_mid_hydration_is_byte_deterministic() { + for &par in PARALLELISM_SWEEP { + for &seed in SEEDS { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(60), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { + Box::new(CancelMidHydration::new(CancelMidHydrationConfig { + frames: 200, + cancel_at_events: 50, + slingshot_latency_ms: 50, + frame_pace_us: 100, + })) + }; + let result = run_leak_check(config.clone(), factory); + assert!( + result.passed(), + "cancel-mid-hydration leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + match &result.outcome { + LeakOutcome::Match => {} + other => panic!("non-match outcome despite passed(): {other:?}"), + } + } + } +} + +#[test] +fn hydrant_disconnect_barrage_is_byte_deterministic() { + for &par in PARALLELISM_SWEEP { + for &seed in SEEDS { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(60), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { + 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!( + result.passed(), + "hydrant-disconnect-barrage leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + assert_eq!( + result.first_report.events_processed, 256, + "expected all 256 frames processed", + ); + } + } +} + +#[test] +fn cold_start_under_live_load_is_byte_deterministic() { + for &par in PARALLELISM_SWEEP { + for &seed in SEEDS { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(60), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { + Box::new(ColdStartUnderLiveLoad::new(ColdStartUnderLiveLoadConfig { + replay_frames: 200, + live_frames: 50, + live_pace_us: 1_000, + replay_pace_us: 0, + })) + }; + let result = run_leak_check(config.clone(), factory); + assert!( + result.passed(), + "cold-start-under-live-load leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + assert_eq!( + result.first_report.events_processed, 250, + "expected 250 (replay+live) frames", + ); + assert_eq!( + result.first_report.edge_count, 250, + "diversified owners should yield 250 distinct edge keys", + ); + } + } +} + +#[test] +fn concurrent_reads_during_replay_is_byte_deterministic() { + for &par in PARALLELISM_SWEEP { + for &seed in SEEDS { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(60), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { + Box::new(ConcurrentReadsDuringReplay::new( + ConcurrentReadsDuringReplayConfig { + follow_frames: 200, + frame_pace_us: 100, + read_pace_us: 1_000, + }, + )) + }; + let result = run_leak_check(config.clone(), factory); + assert!( + result.passed(), + "concurrent-reads-during-replay leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + assert_eq!( + result.first_report.events_processed, 200, + "expected 200 follow frames processed", + ); + } + } +} + +#[test] +fn slingshot_flap_short_brownout_is_byte_deterministic() { + for &par in &[4usize, 16, 64] { + for &seed in &[7u64, 99] { + let parallelism = NonZeroUsize::new(par).unwrap(); + let config = LeakRunConfig { + seed, + parallelism, + max_virtual_runtime: Duration::from_secs(30), + mem_ws_capacity: 4096, + }; + let factory = || -> Box { + Box::new(SlingshotFlap::new(SlingshotFlapConfig { + cross_did_stars: 50, + normal_latency_ms: 2, + brownout_start_ms: 0, + brownout_duration_ms: 100, + brownout_latency_ms: 50, + brownout_enabled: true, + hydrant_send_timeout_ms: 30_000, + hydrant_frame_pace_us: 0, + })) + }; + let result = run_leak_check(config.clone(), factory); + assert!( + result.passed(), + "slingshot-flap leak at par={par} seed={seed}: {:?}", + result.outcome, + ); + assert_eq!( + result.first_report.events_processed, 100, + "expected 100 (50 stars + 50 repos) frames processed", + ); + assert!( + result.first_report.resolver_misses > 0, + "brownout window should have produced resolver transient misses at par={par} seed={seed}, got hits={} misses={}", + result.first_report.resolver_hits, + result.first_report.resolver_misses, + ); + } + } +}