From 19ca6a2b89b26d5a939af8860b6e602ca3f934e5 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sat, 9 May 2026 15:10:58 +0300 Subject: [PATCH] feat(bobbin-sim): slingshot_flap + cold_start_under_live_load Lewis: May this revision serve well! --- .../workloads/cold_start_under_live_load.rs | 219 +++++++++++ .../src/workloads/slingshot_flap.rs | 346 ++++++++++++++++++ 2 files changed, 565 insertions(+) create mode 100644 crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs create mode 100644 crates/bobbin-sim/src/workloads/slingshot_flap.rs 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 new file mode 100644 index 0000000..1ab7a84 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/cold_start_under_live_load.rs @@ -0,0 +1,219 @@ +use std::sync::Arc; +use std::sync::Mutex; +use std::time::Duration; + +use bobbin_runtime::{ + Clock, 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 = "cold-start-under-live-load"; + +#[derive(Clone, Debug)] +pub struct ColdStartUnderLiveLoadConfig { + pub replay_frames: usize, + pub live_frames: usize, + pub live_pace_us: u64, + pub replay_pace_us: u64, +} + +impl Default for ColdStartUnderLiveLoadConfig { + fn default() -> Self { + Self { + replay_frames: 200, + live_frames: 50, + live_pace_us: 1_000, + replay_pace_us: 0, + } + } +} + +pub struct ColdStartUnderLiveLoad { + config: ColdStartUnderLiveLoadConfig, +} + +impl ColdStartUnderLiveLoad { + pub fn new(config: ColdStartUnderLiveLoadConfig) -> Self { + Self { config } + } +} + +impl Workload for ColdStartUnderLiveLoad { + 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 total_frames = (cfg.replay_frames + cfg.live_frames) as u64; + let frame_log = build_frame_log(cfg.replay_frames, cfg.live_frames); + + let hydrant: Arc = Arc::new(LiveLoadHydrant { + frames: Mutex::new(Some(frame_log)), + replay_pace: Duration::from_micros(cfg.replay_pace_us), + live_pace: Duration::from_micros(cfg.live_pace_us), + replay_count: cfg.replay_frames, + clock: clock.clone(), + }); + let slingshot_probe = AssertNoSlingshot::new(); + let slingshot: Arc = Arc::new(slingshot_probe.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() >= total_frames { + 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 outcome = match outcome { + SimOutcome::Passed if stray_slingshot != 0 => SimOutcome::Failed, + other => other, + }; + let failure_reason = match outcome { + SimOutcome::Passed => None, + SimOutcome::Failed => Some(format!( + "events={} target={} last_cursor={} stray_slingshot={}", + snap.events_processed(), + total_frames, + snap.last_cursor().raw(), + 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(replay: usize, live: usize) -> Vec<(bool, String)> { + let total = replay + live; + (0..total) + .map(|i| { + let id = (i + 1) as u64; + let is_live = i >= replay; + let owner = format!("did:plc:owner-{i}"); + let body = serde_json::json!({ + "id": id, + "type": "record", + "record": { + "live": is_live, + "did": owner, + "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:repo-{i}") + } + } + }) + .to_string(); + (is_live, body) + }) + .collect() +} + +struct LiveLoadHydrant { + frames: Mutex>>, + replay_pace: Duration, + live_pace: Duration, + replay_count: usize, + clock: Arc, +} + +impl MemWsResponder for LiveLoadHydrant { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + let replay_pace = self.replay_pace; + let live_pace = self.live_pace; + let replay_count = self.replay_count; + let clock = self.clock.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 { + for (i, (is_live, text)) in frames.into_iter().enumerate() { + let pace = if is_live { live_pace } else { replay_pace }; + if !pace.is_zero() { + clock.sleep(pace).await; + } + if send.send(WsMessage::Text(text)).await.is_err() { + return; + } + if i + 1 == replay_count { + clock.sleep(Duration::from_millis(1)).await; + } + } + }; + let _ = tokio::join!(pong_loop, frame_emit); + }) + } +} diff --git a/crates/bobbin-sim/src/workloads/slingshot_flap.rs b/crates/bobbin-sim/src/workloads/slingshot_flap.rs new file mode 100644 index 0000000..6019921 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/slingshot_flap.rs @@ -0,0 +1,346 @@ +use std::sync::Arc; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use bobbin_runtime::{ + Clock, HttpRequest, MemHttpBody, MemHttpResponder, MemHttpResponse, MemWsResponder, + MemWsServerFuture, NetworkError, UnixMicros, WsMessage, +}; +use bytes::Bytes; +use http::StatusCode; +use tokio::sync::mpsc; +use url::Url; + +use crate::report::{SimOutcome, SimReport}; +use crate::workload::{Workload, WorkloadCtx, WorkloadHooks}; +use crate::workloads::util::{format_rkey, format_tid, parse_repo_lookup}; + +const NAME: &str = "slingshot-flap"; +const REPO_COLLECTION: &str = "sh.tangled.repo"; + +#[derive(Clone, Debug)] +pub struct SlingshotFlapConfig { + pub cross_did_stars: usize, + pub normal_latency_ms: u64, + pub brownout_start_ms: u64, + pub brownout_duration_ms: u64, + pub brownout_latency_ms: u64, + pub brownout_enabled: bool, + pub hydrant_send_timeout_ms: u64, + pub hydrant_frame_pace_us: u64, +} + +impl Default for SlingshotFlapConfig { + fn default() -> Self { + Self { + cross_did_stars: 200, + normal_latency_ms: 2, + brownout_start_ms: 1_000, + brownout_duration_ms: 5_000, + brownout_latency_ms: 200, + brownout_enabled: true, + hydrant_send_timeout_ms: 30_000, + hydrant_frame_pace_us: 0, + } + } +} + +pub struct SlingshotFlap { + config: SlingshotFlapConfig, +} + +impl SlingshotFlap { + pub fn new(config: SlingshotFlapConfig) -> Self { + Self { config } + } +} + +impl Workload for SlingshotFlap { + 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, + consumer_too_slow_count, + .. + } = ctx; + + let star_count = cfg.cross_did_stars; + let total_frames = (star_count as u64) * 2; + let target_owner_count = star_count; + + let frame_script = build_frame_script(&cfg); + let hydrant: Arc = Arc::new(FlapHydrant { + frames: Mutex::new(Some(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(), + clock: clock.clone(), + }); + + let started_unix = clock.now_unix_micros(); + let slingshot: Arc = Arc::new(FlapSlingshot { + cfg: cfg.clone(), + clock: clock.clone(), + started_unix, + }); + + let cts_counter = consumer_too_slow_count.clone(); + let script = Box::pin(async move { + let started = started_unix; + + let mut rx = coverage.subscribe(); + let outcome = loop { + let snap = *rx.borrow_and_update(); + if snap.events_processed() >= total_frames { + break SimOutcome::Passed; + } + if cts_counter.load(Ordering::Relaxed) > 0 { + 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 cts = cts_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", + snap.events_processed(), + total_frames, + )), + SimOutcome::Failed => Some(format!( + "only {}/{} events processed (expected {} owners staged)", + snap.events_processed(), + total_frames, + target_owner_count, + )), + 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: cts, + failure_reason, + } + }); + + WorkloadHooks { + slingshot, + hydrant, + script, + } + } +} + +fn build_frame_script(cfg: &SlingshotFlapConfig) -> Vec { + let mut frames = Vec::with_capacity(cfg.cross_did_stars * 2); + for i in 0..cfg.cross_did_stars { + 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", + "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}") + } + } + }) + .to_string(), + ); + } + for i in 0..cfg.cross_did_stars { + 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", + "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 +} + +struct FlapHydrant { + frames: Mutex>>, + send_timeout: Duration, + frame_pace: Duration, + consumer_too_slow: Arc, + clock: Arc, +} + +impl MemWsResponder for FlapHydrant { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + 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(); + 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 clock_for_emit = clock.clone(); + let frame_emit = async move { + for text in frames { + 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; + match res { + Ok(Ok(())) => continue, + Ok(Err(_)) => return, + Err(_) => { + consumer_too_slow.fetch_add(1, Ordering::Relaxed); + let err_frame = serde_json::json!({ + "type": "error", + "error": "ConsumerTooSlow", + "message": format!( + "sim hydrant send_timeout {send_timeout:?} exhausted", + ), + }) + .to_string(); + let _ = tokio::time::timeout( + Duration::from_secs(1), + send.send(WsMessage::Text(err_frame)), + ) + .await; + let _ = tokio::time::timeout( + Duration::from_secs(1), + send.send(WsMessage::Close { + code: 1011, + reason: "ConsumerTooSlow".into(), + }), + ) + .await; + return; + } + } + } + }; + let _ = tokio::join!(pong_loop, frame_emit); + }) + } +} + +struct FlapSlingshot { + cfg: SlingshotFlapConfig, + clock: Arc, + started_unix: UnixMicros, +} + +impl MemHttpResponder for FlapSlingshot { + fn respond(&self, request: &HttpRequest) -> MemHttpResponse { + let now = self.clock.now_unix_micros().raw(); + let elapsed_ms = now.saturating_sub(self.started_unix.raw()) / 1000; + let in_brownout = self.cfg.brownout_enabled + && elapsed_ms >= self.cfg.brownout_start_ms + && elapsed_ms < self.cfg.brownout_start_ms + self.cfg.brownout_duration_ms; + + if in_brownout { + return MemHttpResponse { + latency: Duration::from_millis(self.cfg.brownout_latency_ms), + result: Err(NetworkError::Transport("slingshot brownout".into())), + }; + } + + let owner_rkey = parse_repo_lookup(&request.url, REPO_COLLECTION); + let body = match owner_rkey { + Some((owner, rkey)) => { + let repo_did = owner.replace("owner-", "repo-"); + serde_json::json!({ + "uri": format!("at://{owner}/{REPO_COLLECTION}/{rkey}"), + "cid": "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i", + "value": { + "$type": REPO_COLLECTION, + "createdAt": "2026-05-01T00:00:00Z", + "knot": "oyster.cafe", + "name": "abalone", + "repoDid": repo_did, + } + }) + .to_string() + } + None => { + return MemHttpResponse { + latency: Duration::from_millis(self.cfg.normal_latency_ms), + result: Ok(MemHttpBody::status_only(StatusCode::NOT_FOUND)), + }; + } + }; + + MemHttpResponse { + latency: Duration::from_millis(self.cfg.normal_latency_ms), + result: Ok(MemHttpBody::ok_json(Bytes::from(body))), + } + } +} -- 2.51.2