From 818c8f0cafc3d897e5e9ee9505f1e44374690505 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sat, 9 May 2026 16:42:45 +0300 Subject: [PATCH] feat(bobbin-sim): cancel_mid_hydration, concurrent_reads_during_replay Lewis: May this revision serve well! --- .../src/workloads/cancel_mid_hydration.rs | 280 ++++++++++++++++++ .../concurrent_reads_during_replay.rs | 264 +++++++++++++++++ 2 files changed, 544 insertions(+) create mode 100644 crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs create mode 100644 crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs diff --git a/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs new file mode 100644 index 0000000..312504e --- /dev/null +++ b/crates/bobbin-sim/src/workloads/cancel_mid_hydration.rs @@ -0,0 +1,280 @@ +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, 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 = "cancel-mid-hydration"; +const REPO_COLLECTION: &str = "sh.tangled.repo"; + +#[derive(Clone, Debug)] +pub struct CancelMidHydrationConfig { + pub frames: usize, + pub cancel_at_events: u64, + pub slingshot_latency_ms: u64, + pub frame_pace_us: u64, +} + +impl Default for CancelMidHydrationConfig { + fn default() -> Self { + Self { + frames: 200, + cancel_at_events: 50, + slingshot_latency_ms: 50, + frame_pace_us: 100, + } + } +} + +pub struct CancelMidHydration { + config: CancelMidHydrationConfig, +} + +impl CancelMidHydration { + pub fn new(config: CancelMidHydrationConfig) -> Self { + Self { config } + } +} + +impl Workload for CancelMidHydration { + 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 = build_frames(cfg.frames); + let slingshot_calls = Arc::new(AtomicU64::new(0)); + + let hydrant: Arc = Arc::new(PacedHydrant { + frames: Mutex::new(Some(frames)), + frame_pace: Duration::from_micros(cfg.frame_pace_us), + clock: clock.clone(), + }); + let slingshot: Arc = Arc::new(LatentSlingshot { + latency: Duration::from_millis(cfg.slingshot_latency_ms), + slingshot_calls: slingshot_calls.clone(), + }); + + let cancel_at = cfg.cancel_at_events; + let frames_total = cfg.frames as u64; + let slingshot_calls_probe = slingshot_calls.clone(); + let script = Box::pin(async move { + let started = clock.now_unix_micros(); + let mut rx = coverage.subscribe(); + let cancel_outcome = loop { + let snap = *rx.borrow_and_update(); + if snap.events_processed() >= cancel_at { + break SimOutcome::Passed; + } + tokio::select! { + res = rx.changed() => match res { + Ok(()) => continue, + Err(_) => break SimOutcome::Failed, + }, + _ = cancel.cancelled() => break SimOutcome::Failed, + } + }; + + cancel.cancel(); + clock.sleep(Duration::from_millis(100)).await; + let events_after_observe = coverage.snapshot().events_processed(); + let slingshot_after_observe = slingshot_calls_probe.load(Ordering::Relaxed); + + clock.sleep(Duration::from_secs(2)).await; + let snap = coverage.snapshot(); + let virtual_runtime = Duration::from_micros( + clock.now_unix_micros().raw().saturating_sub(started.raw()), + ); + let slingshot_after_drain = slingshot_calls_probe.load(Ordering::Relaxed); + let post_observe_events = + snap.events_processed().saturating_sub(events_after_observe); + let post_observe_slingshot = + slingshot_after_drain.saturating_sub(slingshot_after_observe); + let outcome = if !matches!(cancel_outcome, SimOutcome::Passed) { + SimOutcome::Failed + } else if post_observe_slingshot != 0 { + SimOutcome::Failed + } else if post_observe_events != 0 { + SimOutcome::Failed + } else if snap.events_processed() > frames_total { + SimOutcome::Failed + } else if snap.events_processed() < cancel_at { + SimOutcome::Failed + } else { + SimOutcome::Passed + }; + let failure_reason = match outcome { + SimOutcome::Passed => None, + SimOutcome::Failed => Some(format!( + "cancel-mid-hydration anomaly: events={} cancel_at={} \ + events_after_observe={} post_observe_events={} \ + slingshot_after_observe={} slingshot_after_drain={} \ + post_observe_slingshot={}", + snap.events_processed(), + cancel_at, + events_after_observe, + post_observe_events, + slingshot_after_observe, + slingshot_after_drain, + post_observe_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_frames(count: usize) -> Vec { + (0..count) + .map(|i| { + let star_did = format!("did:plc:starer-{i}"); + let target_owner = format!("did:plc:owner-{i}"); + let target_rkey = format_rkey(i); + 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() + }) + .collect() +} + +struct PacedHydrant { + frames: Mutex>>, + frame_pace: Duration, + clock: Arc, +} + +impl MemWsResponder for PacedHydrant { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + let frame_pace = self.frame_pace; + 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 text in frames { + if !frame_pace.is_zero() { + clock.sleep(frame_pace).await; + } + if send.send(WsMessage::Text(text)).await.is_err() { + return; + } + } + }; + let _ = tokio::join!(pong_loop, frame_emit); + }) + } +} + +struct LatentSlingshot { + latency: Duration, + slingshot_calls: Arc, +} + +impl MemHttpResponder for LatentSlingshot { + fn respond(&self, request: &HttpRequest) -> MemHttpResponse { + self.slingshot_calls.fetch_add(1, Ordering::Relaxed); + 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: self.latency, + result: Ok(MemHttpBody::status_only(StatusCode::NOT_FOUND)), + }; + } + }; + MemHttpResponse { + latency: self.latency, + result: Ok(MemHttpBody::ok_json(Bytes::from(body))), + } + } +} diff --git a/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs new file mode 100644 index 0000000..6c558e0 --- /dev/null +++ b/crates/bobbin-sim/src/workloads/concurrent_reads_during_replay.rs @@ -0,0 +1,264 @@ +use std::sync::Arc; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use bobbin_edge_index::{PageCursor, PageLimit}; +use bobbin_runtime::{ + Clock, MemHttpResponder, MemWsResponder, MemWsServerFuture, WsMessage, +}; +use bobbin_types::ids::EdgeKey; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::string::AtUri; +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 = "concurrent-reads-during-replay"; +const TARGET_DID: &str = "did:plc:lyna"; + +#[derive(Clone, Debug)] +pub struct ConcurrentReadsDuringReplayConfig { + pub follow_frames: usize, + pub frame_pace_us: u64, + pub read_pace_us: u64, +} + +impl Default for ConcurrentReadsDuringReplayConfig { + fn default() -> Self { + Self { + follow_frames: 200, + frame_pace_us: 100, + read_pace_us: 50, + } + } +} + +pub struct ConcurrentReadsDuringReplay { + config: ConcurrentReadsDuringReplayConfig, +} + +impl ConcurrentReadsDuringReplay { + pub fn new(config: ConcurrentReadsDuringReplayConfig) -> Self { + Self { config } + } +} + +impl Workload for ConcurrentReadsDuringReplay { + 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 = cfg.follow_frames as u64; + let frame_log = build_frame_log(cfg.follow_frames); + + let hydrant: Arc = Arc::new(FollowHydrant { + frames: Mutex::new(Some(frame_log)), + frame_pace: Duration::from_micros(cfg.frame_pace_us), + clock: clock.clone(), + }); + let slingshot_probe = AssertNoSlingshot::new(); + let slingshot: Arc = Arc::new(slingshot_probe.clone()); + + let read_clock = clock.clone(); + let read_store = store.clone(); + let read_cancel = cancel.clone(); + let read_pace = Duration::from_micros(cfg.read_pace_us); + let monotonicity_violations = Arc::new(AtomicU64::new(0)); + let max_observed = Arc::new(AtomicU64::new(0)); + let read_iterations = Arc::new(AtomicU64::new(0)); + let monotonicity_violations_w = monotonicity_violations.clone(); + let max_observed_w = max_observed.clone(); + let read_iterations_w = read_iterations.clone(); + + let read_key = EdgeKey::new( + Nsid::new_static("sh.tangled.graph.follow").unwrap(), + AtUri::new_owned(format!("at://{TARGET_DID}")).unwrap(), + ); + let final_key = read_key.clone(); + + let reader_task = tokio::spawn(async move { + let mut last_count: u64 = 0; + let mut last_list_len: usize = 0; + loop { + tokio::select! { + biased; + _ = read_cancel.cancelled() => return, + _ = read_clock.sleep(read_pace) => {} + } + let count = read_store.count(&read_key); + let listed = + read_store.list(&read_key, PageCursor::Start, PageLimit::new(50).unwrap()); + if count < last_count || listed.items.len() < last_list_len { + monotonicity_violations_w.fetch_add(1, Ordering::Relaxed); + } + last_count = count.max(last_count); + last_list_len = listed.items.len().max(last_list_len); + max_observed_w.store(last_count, Ordering::Relaxed); + read_iterations_w.fetch_add(1, Ordering::Relaxed); + } + }); + + 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 { + break SimOutcome::Passed; + } + tokio::select! { + res = rx.changed() => match res { + Ok(()) => continue, + Err(_) => break SimOutcome::Failed, + }, + _ = cancel.cancelled() => break SimOutcome::Failed, + } + }; + + cancel.cancel(); + let _ = reader_task.await; + + let snap = coverage.snapshot(); + let virtual_runtime = Duration::from_micros( + clock.now_unix_micros().raw().saturating_sub(started.raw()), + ); + let violations = monotonicity_violations.load(Ordering::Relaxed); + let observed = max_observed.load(Ordering::Relaxed); + let iters = read_iterations.load(Ordering::Relaxed); + let final_count = store.count(&final_key); + let stray_slingshot = slingshot_probe.calls(); + let outcome = match outcome { + SimOutcome::Passed + if violations == 0 && final_count == total && stray_slingshot == 0 => + { + SimOutcome::Passed + } + SimOutcome::Passed => SimOutcome::Failed, + other => other, + }; + let failure_reason = match outcome { + SimOutcome::Passed => None, + SimOutcome::Failed => Some(format!( + "events={} target={} reads={} max_observed_count={} final_count={} \ + monotonicity_violations={} stray_slingshot={}", + snap.events_processed(), + total, + iters, + observed, + final_count, + violations, + 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 { + (0..count) + .map(|i| { + let id = (i + 1) as u64; + let actor = format!("did:plc:actor-{i}"); + serde_json::json!({ + "id": id, + "type": "record", + "record": { + "live": false, + "did": actor, + "rev": format_tid(i), + "collection": "sh.tangled.graph.follow", + "rkey": format_rkey(i), + "action": "create", + "record": { + "$type": "sh.tangled.graph.follow", + "createdAt": "2026-05-01T00:00:00Z", + "subject": TARGET_DID, + } + } + }) + .to_string() + }) + .collect() +} + +struct FollowHydrant { + frames: Mutex>>, + frame_pace: Duration, + clock: Arc, +} + +impl MemWsResponder for FollowHydrant { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames = self.frames.lock().unwrap().take().unwrap_or_default(); + let frame_pace = self.frame_pace; + 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 text in frames { + if !frame_pace.is_zero() { + clock.sleep(frame_pace).await; + } + if send.send(WsMessage::Text(text)).await.is_err() { + return; + } + } + }; + let _ = tokio::join!(pong_loop, frame_emit); + }) + } +} -- 2.51.2