diff --git a/crates/bobbin-sim/src/determinism.rs b/crates/bobbin-sim/src/determinism.rs index b808771..34126b7 100644 --- a/crates/bobbin-sim/src/determinism.rs +++ b/crates/bobbin-sim/src/determinism.rs @@ -1,7 +1,7 @@ use std::num::NonZeroUsize; use std::time::Duration; -use bobbin_ingest::{DisconnectSnapshot, WarmingShadowSnapshot}; +use bobbin_ingest::{DisconnectSnapshot, WarmingBufferSnapshot, WarmingShadowSnapshot}; use tokio::runtime::Builder as TokioBuilder; use tracing_subscriber::Registry; use tracing_subscriber::layer::SubscriberExt; @@ -17,6 +17,24 @@ pub struct LeakRunConfig { pub parallelism: NonZeroUsize, pub max_virtual_runtime: Duration, pub mem_ws_capacity: usize, + pub warming_buffer_enabled: bool, +} + +impl LeakRunConfig { + pub fn new( + seed: u64, + parallelism: NonZeroUsize, + max_virtual_runtime: Duration, + mem_ws_capacity: usize, + ) -> Self { + Self { + seed, + parallelism, + max_virtual_runtime, + mem_ws_capacity, + warming_buffer_enabled: true, + } + } } #[derive(Clone, Debug, Eq, PartialEq)] @@ -62,6 +80,10 @@ pub enum LeakOutcome { first: WarmingShadowSnapshot, second: WarmingShadowSnapshot, }, + WarmingBufferMismatch { + first: WarmingBufferSnapshot, + second: WarmingBufferSnapshot, + }, TraceLengthMismatch { first: usize, second: usize, @@ -130,6 +152,7 @@ fn run_once( sim_config.max_virtual_runtime = config.max_virtual_runtime; sim_config.parallelism = config.parallelism; sim_config.mem_ws_capacity = config.mem_ws_capacity; + sim_config.warming_buffer_enabled = config.warming_buffer_enabled; let runtime = TokioBuilder::new_current_thread() .enable_all() @@ -212,6 +235,12 @@ fn compare_runs( second: b.warming_shadow, }; } + if a.warming_buffer != b.warming_buffer { + return LeakOutcome::WarmingBufferMismatch { + first: a.warming_buffer, + second: b.warming_buffer, + }; + } 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 38239df..a264249 100644 --- a/crates/bobbin-sim/src/main.rs +++ b/crates/bobbin-sim/src/main.rs @@ -173,6 +173,7 @@ fn run_leak(cli: &Cli, seed: u64) -> bobbin_sim::LeakRunResult { parallelism: NonZeroUsize::new(cli.parallelism).expect("parallelism > 0"), max_virtual_runtime: Duration::from_secs(cli.max_virtual_seconds), mem_ws_capacity: cli.mem_ws_capacity, + warming_buffer_enabled: true, }; let cli_snapshot = CliSnapshot::from(cli); let factory = move || -> Box { build_workload_from(&cli_snapshot) }; @@ -226,6 +227,8 @@ fn build_workload_from(snap: &CliSnapshot) -> Box { brownout_enabled: snap.brownout_enabled, hydrant_send_timeout_ms: snap.hydrant_send_timeout_ms, hydrant_frame_pace_us: snap.hydrant_frame_pace_us, + omit_target_repos: false, + emit_live_promotion_frame: false, })), WorkloadName::CancelMidHydration => Box::new(CancelMidHydration::new( CancelMidHydrationConfig::default(), @@ -269,6 +272,18 @@ fn report_as_json(r: &SimReport) -> serde_json::Value { "residual": r.warming_shadow.residual, "distinct_keys_seen": r.warming_shadow.distinct_keys_seen, }, + "warming_buffer": { + "enqueued_total": r.warming_buffer.enqueued_total, + "drained_observe_total": r.warming_buffer.drained_observe_total, + "drained_promote_total": r.warming_buffer.drained_promote_total, + "evicted_total": r.warming_buffer.evicted_total, + "rejected_after_seal": r.warming_buffer.rejected_after_seal, + "distinct_keys_seen": r.warming_buffer.distinct_keys_seen, + "current_entries": r.warming_buffer.current_entries, + "max_concurrent_entries": r.warming_buffer.max_concurrent_entries, + "dep_enqueued_total": r.warming_buffer.dep_enqueued_total, + "dep_drained_observe_total": r.warming_buffer.dep_drained_observe_total, + }, "failure_reason": r.failure_reason, }) } @@ -337,6 +352,11 @@ fn leak_outcome_as_json(outcome: &LeakOutcome) -> serde_json::Value { "first": format!("{first:?}"), "second": format!("{second:?}"), }), + LeakOutcome::WarmingBufferMismatch { first, second } => serde_json::json!({ + "kind": "warming_buffer_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 5d36b7f..3db2b01 100644 --- a/crates/bobbin-sim/src/report.rs +++ b/crates/bobbin-sim/src/report.rs @@ -1,6 +1,6 @@ use std::time::Duration; -use bobbin_ingest::{DisconnectSnapshot, WarmingShadowSnapshot}; +use bobbin_ingest::{DisconnectSnapshot, WarmingBufferSnapshot, WarmingShadowSnapshot}; use bobbin_runtime::UnixMicros; #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -26,6 +26,7 @@ pub struct SimReport { pub disconnect_count: u64, pub last_disconnect: Option, pub warming_shadow: WarmingShadowSnapshot, + pub warming_buffer: WarmingBufferSnapshot, pub failure_reason: Option, } diff --git a/crates/bobbin-sim/src/runtime.rs b/crates/bobbin-sim/src/runtime.rs index a347b37..037c23e 100644 --- a/crates/bobbin-sim/src/runtime.rs +++ b/crates/bobbin-sim/src/runtime.rs @@ -6,7 +6,7 @@ use std::time::Duration; use bobbin_edge_index::{CoverageWatch, EdgeStore}; use bobbin_ingest::{ DEFAULT_INGEST_PARALLELISM, DisconnectSink, IngestConfig, IngestRuntime, RepoIdResolver, - WarmingShadowBuffer, run as run_ingest, + WarmingBuffer, WarmingShadowBuffer, run as run_ingest, }; use bobbin_record_lru::NoopRecordStore; use bobbin_runtime::{ @@ -30,6 +30,7 @@ pub struct SimConfig { pub hydrant_base: Url, pub slingshot_base: Url, pub mem_ws_capacity: usize, + pub warming_buffer_enabled: bool, } impl SimConfig { @@ -42,6 +43,7 @@ impl SimConfig { hydrant_base: Url::parse("ws://hydrant.sim/").unwrap(), slingshot_base: Url::parse("http://slingshot.sim/").unwrap(), mem_ws_capacity: DEFAULT_MEM_WS_CAPACITY, + warming_buffer_enabled: true, } } } @@ -65,6 +67,7 @@ impl Sim { hydrant_base, slingshot_base, mem_ws_capacity, + warming_buffer_enabled, } = self.config; let entropy = Arc::new(SeededEntropy::new(seed)); @@ -78,6 +81,7 @@ impl Sim { 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 warming_buffer = Arc::new(WarmingBuffer::new(hasher.clone())); let workload_name = self.workload.name(); let ctx = WorkloadCtx { @@ -116,6 +120,7 @@ impl Sim { cancel: cancel.clone(), disconnects: Some(disconnects.clone()), warming_shadow: Some(warming_shadow.clone()), + warming_buffer: warming_buffer_enabled.then(|| warming_buffer.clone()), }; let ingest_config = IngestConfig { hydrant_base, @@ -127,7 +132,7 @@ impl Sim { }); let script = hooks.script; - let report = tokio::select! { + let initial_report = tokio::select! { biased; _ = clock.sleep(max_virtual_runtime) => SimReport { workload: workload_name, @@ -138,32 +143,22 @@ impl Sim { events_processed: coverage.snapshot().events_processed(), last_cursor: coverage.snapshot().last_cursor().raw(), edge_count: store.key_count() as u64, - resolver_hits: resolver.stats().hits, - resolver_misses: resolver.stats().miss_count(), + resolver_hits: 0, + resolver_misses: 0, consumer_too_slow_count: consumer_too_slow_count.load(Ordering::Relaxed), disconnect_count: disconnects.count(), last_disconnect: disconnects.snapshot(), warming_shadow: warming_shadow.snapshot(), + warming_buffer: warming_buffer.snapshot(), failure_reason: Some(format!( "max_virtual_runtime {max_virtual_runtime:?} exhausted", )), }, - r = script => { - let stats = resolver.stats(); - SimReport { - 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 - } - } + r = script => r, }; cancel.cancel(); - let drain_deadline = clock.sleep(Duration::from_secs(5)); + let drain_deadline = clock.sleep(Duration::from_secs(30)); tokio::pin!(drain_deadline); tokio::select! { _ = &mut drain_deadline => { @@ -174,6 +169,20 @@ impl Sim { let _ = res; } } - report + + let stats = resolver.stats(); + SimReport { + edge_count: store.key_count() as u64, + events_processed: coverage.snapshot().events_processed(), + last_cursor: coverage.snapshot().last_cursor().raw(), + 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(), + warming_buffer: warming_buffer.snapshot(), + ..initial_report + } } } diff --git a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs index 2fbfa16..87cb372 100644 --- a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs +++ b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs @@ -157,6 +157,7 @@ impl Workload for CancelMidHydration { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: 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 eea780c..ef6e2c2 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 @@ -123,6 +123,7 @@ impl Workload for ColdStartUnderLiveLoad { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: 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 a4ba397..fab4fd3 100644 --- a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs +++ b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs @@ -181,6 +181,7 @@ impl Workload for ConcurrentReadsDuringReplay { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: 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 16dde34..e314f6c 100644 --- a/crates/bobbin-sim/src/workloads/frame_burst.rs +++ b/crates/bobbin-sim/src/workloads/frame_burst.rs @@ -107,6 +107,7 @@ impl Workload for FrameBurst { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: 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 155dc1c..19050ae 100644 --- a/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs +++ b/crates/bobbin-sim/src/workloads/hydrant_disconnect_barrage.rs @@ -132,6 +132,7 @@ impl Workload for HydrantDisconnectBarrage { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: 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 a5a3aee..f0130de 100644 --- a/crates/bobbin-sim/src/workloads/slingshot_flap.rs +++ b/crates/bobbin-sim/src/workloads/slingshot_flap.rs @@ -8,6 +8,7 @@ use bobbin_runtime::{ }; use bytes::Bytes; use http::StatusCode; +use jacquard_common::types::tid::Tid; use tokio::sync::mpsc; use url::Url; @@ -28,6 +29,8 @@ pub struct SlingshotFlapConfig { pub brownout_enabled: bool, pub hydrant_send_timeout_ms: u64, pub hydrant_frame_pace_us: u64, + pub omit_target_repos: bool, + pub emit_live_promotion_frame: bool, } impl Default for SlingshotFlapConfig { @@ -41,6 +44,8 @@ impl Default for SlingshotFlapConfig { brownout_enabled: true, hydrant_send_timeout_ms: 30_000, hydrant_frame_pace_us: 0, + omit_target_repos: false, + emit_live_promotion_frame: false, } } } @@ -73,10 +78,13 @@ impl Workload for SlingshotFlap { } = ctx; let star_count = cfg.cross_did_stars; - let total_frames = (star_count as u64) * 2; + let repo_count = if cfg.omit_target_repos { 0 } else { star_count }; + let live_count = if cfg.emit_live_promotion_frame { 1 } else { 0 }; + let total_frames = (star_count + repo_count + live_count) as u64; let target_owner_count = star_count; - let frame_script = Arc::new(build_frame_script(&cfg)); + let started_unix_pre = clock.now_unix_micros(); + let frame_script = Arc::new(build_frame_script(&cfg, started_unix_pre)); let session_count = Arc::new(AtomicU64::new(0)); let hydrant: Arc = Arc::new(FlapHydrant { frames: frame_script, @@ -87,7 +95,7 @@ impl Workload for SlingshotFlap { clock: clock.clone(), }); - let started_unix = clock.now_unix_micros(); + let started_unix = started_unix_pre; let slingshot: Arc = Arc::new(FlapSlingshot { cfg: cfg.clone(), clock: clock.clone(), @@ -152,6 +160,7 @@ impl Workload for SlingshotFlap { disconnect_count: 0, last_disconnect: None, warming_shadow: Default::default(), + warming_buffer: Default::default(), failure_reason, } }); @@ -164,8 +173,8 @@ impl Workload for SlingshotFlap { } } -fn build_frame_script(cfg: &SlingshotFlapConfig) -> Vec<(u64, String)> { - let mut frames = Vec::with_capacity(cfg.cross_did_stars * 2); +fn build_frame_script(cfg: &SlingshotFlapConfig, started_unix: UnixMicros) -> Vec<(u64, String)> { + let mut frames = Vec::with_capacity(cfg.cross_did_stars * 2 + 1); for i in 0..cfg.cross_did_stars { let id = (i + 1) as u64; let star_did = format!("did:plc:starer-{i}"); @@ -191,26 +200,57 @@ fn build_frame_script(cfg: &SlingshotFlapConfig) -> Vec<(u64, 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}"); + if !cfg.omit_target_repos { + 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}"); + 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": { + "$type": REPO_COLLECTION, + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": format!("repo-{i}"), + "repoDid": repo_did, + } + } + }) + .to_string(); + frames.push((id, body)); + } + } + if cfg.emit_live_promotion_frame { + 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}"); + let promoter_rkey = format_rkey(prior + 1); + let live_rev = Tid::from_time(started_unix.raw(), 0); let body = serde_json::json!({ "id": id, "type": "record", "record": { - "live": false, - "did": owner, - "rev": format_tid(cfg.cross_did_stars + i), + "live": true, + "did": promoter_did, + "rev": live_rev.as_str(), "collection": REPO_COLLECTION, - "rkey": format_rkey(i), + "rkey": promoter_rkey, "action": "create", "record": { "$type": REPO_COLLECTION, "createdAt": "2026-05-01T00:00:00Z", "knot": "oyster.cafe", - "name": format!("repo-{i}"), - "repoDid": repo_did, + "name": promoter_name, + "repoDid": promoter_did, } } }) diff --git a/crates/bobbin-sim/tests/determinism_leak.rs b/crates/bobbin-sim/tests/determinism_leak.rs index c342398..73106ff 100644 --- a/crates/bobbin-sim/tests/determinism_leak.rs +++ b/crates/bobbin-sim/tests/determinism_leak.rs @@ -22,6 +22,7 @@ fn frame_burst_is_byte_deterministic_across_parallelism_sweep() { parallelism, max_virtual_runtime: Duration::from_secs(30), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(FrameBurst::new(64)) }; let result = run_leak_check(config.clone(), factory); @@ -52,6 +53,7 @@ fn cancel_mid_hydration_is_byte_deterministic() { parallelism, max_virtual_runtime: Duration::from_secs(60), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(CancelMidHydration::new(CancelMidHydrationConfig { @@ -85,6 +87,7 @@ fn hydrant_disconnect_barrage_is_byte_deterministic() { parallelism, max_virtual_runtime: Duration::from_secs(60), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(HydrantDisconnectBarrage::new(HydrantDisconnectBarrageConfig { @@ -116,6 +119,7 @@ fn cold_start_under_live_load_is_byte_deterministic() { parallelism, max_virtual_runtime: Duration::from_secs(60), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(ColdStartUnderLiveLoad::new(ColdStartUnderLiveLoadConfig { @@ -153,6 +157,7 @@ fn concurrent_reads_during_replay_is_byte_deterministic() { parallelism, max_virtual_runtime: Duration::from_secs(60), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(ConcurrentReadsDuringReplay::new( @@ -187,6 +192,7 @@ fn slingshot_flap_short_brownout_is_byte_deterministic() { parallelism, max_virtual_runtime: Duration::from_secs(30), mem_ws_capacity: 4096, + warming_buffer_enabled: true, }; let factory = || -> Box { Box::new(SlingshotFlap::new(SlingshotFlapConfig { @@ -198,6 +204,8 @@ fn slingshot_flap_short_brownout_is_byte_deterministic() { brownout_enabled: true, hydrant_send_timeout_ms: 30_000, hydrant_frame_pace_us: 0, + omit_target_repos: false, + emit_live_promotion_frame: false, })) }; let result = run_leak_check(config.clone(), factory); @@ -211,10 +219,15 @@ fn slingshot_flap_short_brownout_is_byte_deterministic() { "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, + result.first_report.warming_buffer.enqueued_total > 0, + "expected the buffer path to exercise during warming at par={par} seed={seed}, got snapshot={:?}", + result.first_report.warming_buffer, + ); + assert_eq!( + result.first_report.warming_buffer.enqueued_total, + result.first_report.warming_buffer.drained_observe_total, + "every parked star must drain via observe at par={par} seed={seed}, got snapshot={:?}", + result.first_report.warming_buffer, ); } } diff --git a/crates/bobbin-sim/tests/reconnect_replay.rs b/crates/bobbin-sim/tests/reconnect_replay.rs index f652ac3..df5e889 100644 --- a/crates/bobbin-sim/tests/reconnect_replay.rs +++ b/crates/bobbin-sim/tests/reconnect_replay.rs @@ -17,10 +17,13 @@ fn slingshot_flap_resumes_via_replay_after_pong_timeout() { brownout_enabled: true, hydrant_send_timeout_ms: 30_000, hydrant_frame_pace_us: 1_000, + omit_target_repos: false, + emit_live_promotion_frame: false, }; let mut sim_config = SimConfig::new(7); sim_config.parallelism = NonZeroUsize::new(4).unwrap(); sim_config.max_virtual_runtime = Duration::from_secs(600); + sim_config.warming_buffer_enabled = false; let workload = Box::new(SlingshotFlap::new(cfg.clone())); let runtime = TokioBuilder::new_current_thread() diff --git a/crates/bobbin-sim/tests/warming_buffer.rs b/crates/bobbin-sim/tests/warming_buffer.rs new file mode 100644 index 0000000..35339e0 --- /dev/null +++ b/crates/bobbin-sim/tests/warming_buffer.rs @@ -0,0 +1,231 @@ +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 lever_b_absorbs_sustained_brownout_without_disconnect() { + 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, + omit_target_repos: false, + emit_live_promotion_frame: false, + }; + 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, + "Lever B should keep slingshot off the hot path during cold replay, got {:?} reason={:?}", + report.outcome, + report.failure_reason, + ); + assert_eq!(report.events_processed, total_frames); + 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, + ); + + let b = report.warming_buffer; + assert_eq!( + b.enqueued_total, cfg.cross_did_stars as u64, + "every cross-DID star must park exactly once: {b:?}", + ); + assert_eq!( + b.drained_observe_total, cfg.cross_did_stars as u64, + "every observed repo must drain its parked star: {b:?}", + ); + assert_eq!(b.drained_promote_total, 0, "no residual at promote: {b:?}"); + assert_eq!(b.current_entries, 0, "buffer empty after drain: {b:?}"); + assert_eq!(b.evicted_total, 0, "no evictions in this workload: {b:?}"); + assert!( + b.max_concurrent_entries > 0 && b.max_concurrent_entries <= cfg.cross_did_stars as u64, + "peak depth bounded by cohort size, got {} for {} stars", + b.max_concurrent_entries, + cfg.cross_did_stars, + ); +} + +#[test] +fn shadow_and_buffer_observe_the_same_population() { + let cfg = SlingshotFlapConfig { + cross_did_stars: 2000, + 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, + omit_target_repos: false, + emit_live_promotion_frame: false, + }; + 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()); + + assert_eq!(report.outcome, SimOutcome::Passed, "{:?}", report.failure_reason); + + let s = report.warming_shadow; + let b = report.warming_buffer; + assert_eq!( + s.enqueued_total, b.dep_enqueued_total, + "shadow note_unresolved per dep must equal buffer dep_enqueued_total: shadow={s:?} buffer={b:?}", + ); + assert_eq!( + s.drained_via_observe_total, b.dep_drained_observe_total, + "shadow note_observed per dep must equal buffer dep_drained_observe_total: shadow={s:?} buffer={b:?}", + ); + assert_eq!( + s.distinct_keys_seen, b.distinct_keys_seen, + "shadow and buffer must see the same distinct (owner, rkey) keys: shadow={s:?} buffer={b:?}", + ); + assert_eq!( + s.residual, + b.dep_enqueued_total - b.dep_drained_observe_total, + "shadow residual must equal buffer un-observed deps: shadow={s:?} buffer={b:?}", + ); +} + +#[test] +fn warming_to_ready_promote_drains_residual_via_parallel_slingshot_wave() { + let cfg = SlingshotFlapConfig { + cross_did_stars: 200, + 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, + omit_target_repos: true, + emit_live_promotion_frame: true, + }; + 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()); + + assert_eq!( + report.outcome, + SimOutcome::Passed, + "promote-flush workload must complete, got {:?} reason={:?}", + report.outcome, + report.failure_reason, + ); + + let b = report.warming_buffer; + let stars = cfg.cross_did_stars as u64; + assert_eq!( + b.enqueued_total, stars, + "every cross-DID star must park (no repos arrive): {b:?}", + ); + assert_eq!( + b.drained_observe_total, 0, + "no repos arrive on the stream so observe-drain stays at zero: {b:?}", + ); + assert_eq!( + b.drained_promote_total, stars, + "Warming->Ready promote must drain every parked entry: {b:?}", + ); + assert_eq!( + b.current_entries, 0, + "buffer must be empty after promote drain: {b:?}", + ); + assert!( + report.resolver_hits + report.resolver_misses >= stars, + "promote wave must fan out one slingshot resolve per distinct dep, got hits={} misses={} for {stars} stars", + report.resolver_hits, + report.resolver_misses, + ); + assert_eq!( + report.edge_count, stars + 1, + "every drained star plus the live promoter repo should land in the index, got {} for {stars} stars", + report.edge_count, + ); +} + +#[test] +fn buffer_enabled_and_disabled_produce_identical_edge_index() { + let cfg = SlingshotFlapConfig { + cross_did_stars: 200, + 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, + omit_target_repos: false, + emit_live_promotion_frame: false, + }; + + let run_with_buffer = |enabled: bool| { + let mut sim_config = SimConfig::new(7); + sim_config.parallelism = NonZeroUsize::new(16).unwrap(); + sim_config.max_virtual_runtime = Duration::from_secs(120); + sim_config.warming_buffer_enabled = enabled; + 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"); + runtime.block_on(Sim::new(sim_config, workload).run()) + }; + + 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.events_processed, without_buffer.events_processed, + "events_processed must match across buffer modes", + ); + assert_eq!( + with_buffer.edge_count, without_buffer.edge_count, + "edge_count must match: lever B is correctness-preserving (with={} without={})", + with_buffer.edge_count, without_buffer.edge_count, + ); + assert_eq!( + with_buffer.last_cursor, without_buffer.last_cursor, + "final cursor must match across buffer modes", + ); +} diff --git a/crates/bobbin-sim/tests/warming_shadow.rs b/crates/bobbin-sim/tests/warming_shadow.rs index d8139a9..da06418 100644 --- a/crates/bobbin-sim/tests/warming_shadow.rs +++ b/crates/bobbin-sim/tests/warming_shadow.rs @@ -16,6 +16,8 @@ fn shadow_buffer_bounded_by_cross_did_attachment_cohort_and_drains_to_zero() { brownout_enabled: false, hydrant_send_timeout_ms: 30_000, hydrant_frame_pace_us: 1_000, + omit_target_repos: false, + emit_live_promotion_frame: false, }; let mut sim_config = SimConfig::new(7); sim_config.parallelism = NonZeroUsize::new(16).unwrap();