diff --git a/crates/bobbin-sim/src/main.rs b/crates/bobbin-sim/src/main.rs new file mode 100644 index 0000000..3610017 --- /dev/null +++ b/crates/bobbin-sim/src/main.rs @@ -0,0 +1,334 @@ +use std::num::NonZeroUsize; +use std::path::PathBuf; +use std::process::ExitCode; +use std::time::Duration; + +use bobbin_sim::workloads::{ + CancelMidHydration, CancelMidHydrationConfig, ColdStartUnderLiveLoad, + ColdStartUnderLiveLoadConfig, ConcurrentReadsDuringReplay, + ConcurrentReadsDuringReplayConfig, FrameBurst, HydrantDisconnectBarrage, + HydrantDisconnectBarrageConfig, SlingshotFlap, SlingshotFlapConfig, +}; +use bobbin_sim::{ + LeakOutcome, LeakRunConfig, Sim, SimConfig, SimOutcome, SimReport, Workload, + run_leak_check, +}; +use clap::{Parser, ValueEnum}; + +#[derive(Parser, Debug)] +#[command(name = "bobbin-sim", version)] +struct Cli { + #[arg(long, default_value_t = 0)] + seed: u64, + #[arg(long, value_enum, default_value_t = WorkloadName::FrameBurst)] + workload: WorkloadName, + #[arg(long, default_value_t = 64)] + frames: usize, + #[arg(long, default_value_t = 60)] + max_virtual_seconds: u64, + #[arg(long)] + no_brownout: bool, + #[arg(long, default_value_t = 200)] + cross_did_stars: usize, + #[arg(long, default_value_t = 1_000)] + brownout_start_ms: u64, + #[arg(long, default_value_t = 5_000)] + brownout_duration_ms: u64, + #[arg(long, default_value_t = 200)] + brownout_latency_ms: u64, + #[arg(long, default_value_t = 2)] + normal_latency_ms: u64, + #[arg(long, default_value_t = 16)] + parallelism: usize, + #[arg(long, default_value_t = 4096)] + mem_ws_capacity: usize, + #[arg(long, default_value_t = 30_000)] + hydrant_send_timeout_ms: u64, + #[arg(long, default_value_t = 0)] + hydrant_frame_pace_us: u64, + #[arg(long, default_value_t = 1)] + seeds: u64, + #[arg(long)] + quarantine_out: Option, + #[arg(long)] + leak_check: bool, +} + +#[derive(Clone, Copy, Debug, ValueEnum, Eq, PartialEq)] +enum WorkloadName { + FrameBurst, + SlingshotFlap, + CancelMidHydration, + HydrantDisconnectBarrage, + ColdStartUnderLiveLoad, + ConcurrentReadsDuringReplay, +} + +fn main() -> ExitCode { + let cli = Cli::parse(); + tracing_subscriber::fmt() + .with_env_filter( + tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("warn,bobbin_sim=info")), + ) + .with_writer(std::io::stderr) + .init(); + + if cli.seeds <= 1 && !cli.leak_check { + let report = run_single(&cli, cli.seed); + println!( + "{}", + serde_json::to_string_pretty(&report_as_json(&report)).unwrap() + ); + exit_code_from(report.outcome) + } else if cli.leak_check && cli.seeds <= 1 { + let result = run_leak(&cli, cli.seed); + let json = leak_result_as_json(&result); + println!("{}", serde_json::to_string_pretty(&json).unwrap()); + if matches!(result.outcome, LeakOutcome::Match) + && matches!(result.first_report.outcome, SimOutcome::Passed) + { + ExitCode::from(0) + } else { + ExitCode::from(1) + } + } else { + run_sweep(&cli) + } +} + +fn run_sweep(cli: &Cli) -> ExitCode { + let mut quarantined: Vec = Vec::new(); + let mut passed = 0u64; + let total = cli.seeds; + for offset in 0..total { + let seed = cli.seed.wrapping_add(offset); + let outcome_record = if cli.leak_check { + let result = run_leak(cli, seed); + let leak_match = matches!(result.outcome, LeakOutcome::Match); + let workload_passed = matches!(result.first_report.outcome, SimOutcome::Passed); + let ok = leak_match && workload_passed; + if ok { + passed += 1; + None + } else { + Some(leak_result_as_json(&result)) + } + } else { + let report = run_single(cli, seed); + if matches!(report.outcome, SimOutcome::Passed) { + passed += 1; + None + } else { + Some(report_as_json(&report)) + } + }; + if let Some(rec) = outcome_record { + quarantined.push(rec); + } + } + let summary = serde_json::json!({ + "workload": format!("{:?}", cli.workload), + "seeds_run": total, + "passed": passed, + "failed": total - passed, + "leak_check": cli.leak_check, + "parallelism": cli.parallelism, + "quarantine_count": quarantined.len(), + }); + println!("{}", serde_json::to_string_pretty(&summary).unwrap()); + if let Some(path) = &cli.quarantine_out { + let payload = serde_json::json!({ + "summary": summary, + "failures": quarantined, + }); + std::fs::write(path, serde_json::to_vec_pretty(&payload).unwrap()) + .expect("write quarantine output"); + eprintln!("quarantine written to {}", path.display()); + } + if quarantined.is_empty() { + ExitCode::from(0) + } else { + ExitCode::from(1) + } +} + +fn run_single(cli: &Cli, seed: u64) -> SimReport { + let workload = build_workload(cli); + let mut config = SimConfig::new(seed); + config.max_virtual_runtime = Duration::from_secs(cli.max_virtual_seconds); + config.parallelism = NonZeroUsize::new(cli.parallelism).expect("parallelism > 0"); + config.mem_ws_capacity = cli.mem_ws_capacity; + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .start_paused(true) + .build() + .expect("build current_thread runtime with paused time"); + runtime.block_on(Sim::new(config, workload).run()) +} + +fn run_leak(cli: &Cli, seed: u64) -> bobbin_sim::LeakRunResult { + let leak_config = LeakRunConfig { + seed, + 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, + }; + let cli_snapshot = CliSnapshot::from(cli); + let factory = move || -> Box { build_workload_from(&cli_snapshot) }; + run_leak_check(leak_config, factory) +} + +#[derive(Clone)] +struct CliSnapshot { + workload: WorkloadName, + frames: usize, + cross_did_stars: usize, + normal_latency_ms: u64, + brownout_start_ms: u64, + brownout_duration_ms: u64, + brownout_latency_ms: u64, + brownout_enabled: bool, + hydrant_send_timeout_ms: u64, + hydrant_frame_pace_us: u64, +} + +impl From<&Cli> for CliSnapshot { + fn from(cli: &Cli) -> Self { + Self { + workload: cli.workload, + frames: cli.frames, + cross_did_stars: cli.cross_did_stars, + normal_latency_ms: cli.normal_latency_ms, + brownout_start_ms: cli.brownout_start_ms, + brownout_duration_ms: cli.brownout_duration_ms, + brownout_latency_ms: cli.brownout_latency_ms, + brownout_enabled: !cli.no_brownout, + hydrant_send_timeout_ms: cli.hydrant_send_timeout_ms, + hydrant_frame_pace_us: cli.hydrant_frame_pace_us, + } + } +} + +fn build_workload(cli: &Cli) -> Box { + build_workload_from(&CliSnapshot::from(cli)) +} + +fn build_workload_from(snap: &CliSnapshot) -> Box { + match snap.workload { + WorkloadName::FrameBurst => Box::new(FrameBurst::new(snap.frames)), + WorkloadName::SlingshotFlap => Box::new(SlingshotFlap::new(SlingshotFlapConfig { + cross_did_stars: snap.cross_did_stars, + normal_latency_ms: snap.normal_latency_ms, + brownout_start_ms: snap.brownout_start_ms, + brownout_duration_ms: snap.brownout_duration_ms, + brownout_latency_ms: snap.brownout_latency_ms, + brownout_enabled: snap.brownout_enabled, + hydrant_send_timeout_ms: snap.hydrant_send_timeout_ms, + hydrant_frame_pace_us: snap.hydrant_frame_pace_us, + })), + WorkloadName::CancelMidHydration => Box::new(CancelMidHydration::new( + CancelMidHydrationConfig::default(), + )), + WorkloadName::HydrantDisconnectBarrage => Box::new(HydrantDisconnectBarrage::new( + HydrantDisconnectBarrageConfig::default(), + )), + WorkloadName::ColdStartUnderLiveLoad => Box::new(ColdStartUnderLiveLoad::new( + ColdStartUnderLiveLoadConfig::default(), + )), + WorkloadName::ConcurrentReadsDuringReplay => Box::new(ConcurrentReadsDuringReplay::new( + ConcurrentReadsDuringReplayConfig::default(), + )), + } +} + +fn report_as_json(r: &SimReport) -> serde_json::Value { + serde_json::json!({ + "workload": r.workload, + "seed": r.seed, + "outcome": format!("{:?}", r.outcome), + "virtual_runtime_micros": r.virtual_runtime.as_micros() as u64, + "virtual_clock_end_unix_micros": r.virtual_clock_end.raw(), + "events_processed": r.events_processed, + "last_cursor": r.last_cursor, + "edge_count": r.edge_count, + "resolver_hits": r.resolver_hits, + "resolver_misses": r.resolver_misses, + "consumer_too_slow_count": r.consumer_too_slow_count, + "failure_reason": r.failure_reason, + }) +} + +fn leak_result_as_json(r: &bobbin_sim::LeakRunResult) -> serde_json::Value { + serde_json::json!({ + "workload": r.workload, + "seed": r.config.seed, + "parallelism": r.config.parallelism.get(), + "leak_outcome": leak_outcome_as_json(&r.outcome), + "first": report_as_json(&r.first_report), + "second": report_as_json(&r.second_report), + }) +} + +fn leak_outcome_as_json(outcome: &LeakOutcome) -> serde_json::Value { + match outcome { + LeakOutcome::Match => serde_json::json!({ "kind": "match" }), + LeakOutcome::OutcomeMismatch { first, second } => serde_json::json!({ + "kind": "outcome_mismatch", + "first": format!("{first:?}"), + "second": format!("{second:?}"), + }), + LeakOutcome::EventCountMismatch { first, second } => serde_json::json!({ + "kind": "event_count_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::EdgeCountMismatch { first, second } => serde_json::json!({ + "kind": "edge_count_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::LastCursorMismatch { first, second } => serde_json::json!({ + "kind": "last_cursor_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::ResolverHitsMismatch { first, second } => serde_json::json!({ + "kind": "resolver_hits_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::ResolverMissesMismatch { first, second } => serde_json::json!({ + "kind": "resolver_misses_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::ConsumerTooSlowMismatch { first, second } => serde_json::json!({ + "kind": "consumer_too_slow_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::TraceLengthMismatch { first, second } => serde_json::json!({ + "kind": "trace_length_mismatch", + "first": first, + "second": second, + }), + LeakOutcome::TraceLineMismatch { + index, + first, + second, + } => serde_json::json!({ + "kind": "trace_line_mismatch", + "index": index, + "first": first, + "second": second, + }), + } +} + +fn exit_code_from(outcome: SimOutcome) -> ExitCode { + match outcome { + SimOutcome::Passed => ExitCode::from(0), + SimOutcome::Failed | SimOutcome::TimedOut => ExitCode::from(1), + } +} diff --git a/crates/bobbin-sim/src/workloads/frame_burst.rs b/crates/bobbin-sim/src/workloads/frame_burst.rs new file mode 100644 index 0000000..3d31fe5 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/frame_burst.rs @@ -0,0 +1,179 @@ +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 = "frame-burst"; +const AUTHORITY_DID: &str = "did:plc:abalone"; + +pub struct FrameBurst { + frames: usize, +} + +impl FrameBurst { + pub fn new(frames: usize) -> Self { + Self { frames } + } +} + +impl Workload for FrameBurst { + fn name(&self) -> &'static str { + NAME + } + + fn build(self: Box, ctx: WorkloadCtx) -> WorkloadHooks { + let count = self.frames; + let frames = build_frame_script(count); + let session_count = Arc::new(AtomicU64::new(0)); + + let hydrant: Arc = Arc::new(BurstHydrant { + frames: Mutex::new(Some(frames)), + session_count: session_count.clone(), + }); + let slingshot_probe = AssertNoSlingshot::new(); + let slingshot: Arc = Arc::new(slingshot_probe.clone()); + + let WorkloadCtx { + seed, + clock, + coverage, + store, + cancel, + .. + } = ctx; + let target = count as u64; + let sessions = session_count.clone(); + + let script = Box::pin(async move { + let mut rx = coverage.subscribe(); + let started = clock.now_unix_micros(); + let outcome = loop { + let snap = *rx.borrow_and_update(); + if snap.events_processed() >= target { + break SimOutcome::Passed; + } + 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 stray_slingshot = slingshot_probe.calls(); + let session_total = sessions.load(Ordering::Relaxed); + let outcome = match outcome { + SimOutcome::Passed if stray_slingshot == 0 && session_total == 1 => { + SimOutcome::Passed + } + SimOutcome::Passed => SimOutcome::Failed, + other => other, + }; + let failure_reason = match outcome { + SimOutcome::Passed => None, + SimOutcome::Failed => Some(format!( + "events={}/{} stray_slingshot={} sessions={}", + snap.events_processed(), + target, + stray_slingshot, + session_total, + )), + 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_script(count: usize) -> Vec { + (0..count) + .map(|i| { + serde_json::json!({ + "id": (i + 1) as u64, + "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() + }) + .collect() +} + +struct BurstHydrant { + frames: Mutex>>, + session_count: Arc, +} + +impl MemWsResponder for BurstHydrant { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + self.session_count.fetch_add(1, Ordering::Relaxed); + Box::pin(async move { + for text in frames { + if send.send(WsMessage::Text(text)).await.is_err() { + return; + } + } + loop { + match recv.recv().await { + Some(WsMessage::Ping(payload)) => { + if send.send(WsMessage::Pong(payload)).await.is_err() { + return; + } + } + Some(WsMessage::Close { .. }) | None => return, + Some(_) => {} + } + } + }) + } +} diff --git a/crates/bobbin-sim/src/workloads/mod.rs b/crates/bobbin-sim/src/workloads/mod.rs new file mode 100644 index 0000000..4f52264 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/mod.rs @@ -0,0 +1,16 @@ +pub mod cancel_mid_hydration; +pub mod cold_start_under_live_load; +pub mod concurrent_reads_during_replay; +pub mod frame_burst; +pub mod hydrant_disconnect_barrage; +pub mod slingshot_flap; +pub mod util; + +pub use cancel_mid_hydration::{CancelMidHydration, CancelMidHydrationConfig}; +pub use cold_start_under_live_load::{ColdStartUnderLiveLoad, ColdStartUnderLiveLoadConfig}; +pub use concurrent_reads_during_replay::{ + ConcurrentReadsDuringReplay, ConcurrentReadsDuringReplayConfig, +}; +pub use frame_burst::FrameBurst; +pub use hydrant_disconnect_barrage::{HydrantDisconnectBarrage, HydrantDisconnectBarrageConfig}; +pub use slingshot_flap::{SlingshotFlap, SlingshotFlapConfig}; diff --git a/crates/bobbin-sim/src/workloads/util.rs b/crates/bobbin-sim/src/workloads/util.rs new file mode 100644 index 0000000..8022b17 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/util.rs @@ -0,0 +1,73 @@ +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use bobbin_runtime::{HttpRequest, MemHttpBody, MemHttpResponder, MemHttpResponse}; +use http::StatusCode; +use url::Url; + +const ALPHABET_RKEY: &[u8] = b"abcdefghijklmnopqrstuvwxyz234567"; +const ALPHABET_TID: &[u8] = b"234567abcdefghijklmnopqrstuvwxyz"; +const ENCODED_LEN: usize = 13; + +pub fn format_rkey(i: usize) -> String { + encode_padded(i.saturating_add(1), ALPHABET_RKEY, b'a') +} + +pub fn format_tid(i: usize) -> String { + encode_padded(i.saturating_add(1), ALPHABET_TID, b'2') +} + +fn encode_padded(mut idx: usize, alphabet: &[u8], pad: u8) -> String { + let mut buf = [pad; ENCODED_LEN]; + let mut pos = buf.len(); + while idx > 0 && pos > 0 { + pos -= 1; + buf[pos] = alphabet[idx % alphabet.len()]; + idx /= alphabet.len(); + } + String::from_utf8(buf.to_vec()).unwrap() +} + +pub fn parse_repo_lookup(url: &Url, expected_collection: &str) -> Option<(String, String)> { + let mut repo: Option = None; + let mut collection: Option = None; + let mut rkey: Option = None; + for (k, v) in url.query_pairs() { + match k.as_ref() { + "repo" => repo = Some(v.into_owned()), + "collection" => collection = Some(v.into_owned()), + "rkey" => rkey = Some(v.into_owned()), + _ => {} + } + } + if collection.as_deref() != Some(expected_collection) { + return None; + } + Some((repo?, rkey?)) +} + +#[derive(Clone, Default)] +pub struct AssertNoSlingshot { + calls: Arc, +} + +impl AssertNoSlingshot { + pub fn new() -> Self { + Self::default() + } + + pub fn calls(&self) -> u64 { + self.calls.load(Ordering::Relaxed) + } +} + +impl MemHttpResponder for AssertNoSlingshot { + fn respond(&self, _: &HttpRequest) -> MemHttpResponse { + self.calls.fetch_add(1, Ordering::Relaxed); + MemHttpResponse { + latency: Duration::ZERO, + result: Ok(MemHttpBody::status_only(StatusCode::NOT_FOUND)), + } + } +}