From 126cb5e156879a041e2dfb16fce71a2daee8f54a Mon Sep 17 00:00:00 2001 From: Lewis Date: Sat, 09 May 2026 06:50:39 +0000 Subject: [PATCH] feat(runtime): SimClock, SeededEntropy, MemNetwork Lewis: May this revision serve well! --- crates/ingest/src/lib.rs | 81 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--- crates/runtime/src/clock.rs | 69 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ crates/runtime/src/entropy.rs | 66 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ crates/runtime/src/lib.rs | 9 +++++++-- crates/runtime/src/mem_network.rs | 352 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 5 file(s) changed, 572 insertion(s)(+), 5 deletion(s)(-) diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -190,8 +190,9 @@ backoff = RECONNECT_INITIAL_DELAY; } else { tokio::select! { - _ = runtime.clock.sleep(jittered(backoff, &*runtime.entropy)) => {} + biased; _ = runtime.cancel.cancelled() => return Ok(()), + _ = runtime.clock.sleep(jittered(backoff, &*runtime.entropy)) => {} } backoff = (backoff * 2).min(RECONNECT_MAX_DELAY); } @@ -518,9 +519,27 @@ } } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum Regime { + Replay, + Live, + NonRecord, +} + +impl Regime { + fn as_str(self) -> &'static str { + match self { + Self::Replay => "replay", + Self::Live => "live", + Self::NonRecord => "non_record", + } + } +} + struct Pending { cursor: HydrantCursor, signal: PromotionSignal, + regime: Regime, op: PendingOp, } @@ -610,6 +629,8 @@ ) { let nsid = pending_nsid(&staged.pending.op).cloned(); let edge_count = pending_edge_count(&staged.pending.op); + let regime = staged.pending.regime; + let cursor = staged.pending.cursor.raw(); let Resolved { pending, prepare_start, @@ -622,6 +643,8 @@ let commit_end = rt.clock.now_instant(); tracing::trace!( target: "bobbin_ingest::stage", + cursor, + regime = regime.as_str(), nsid = %nsid.as_ref().map(Nsid::as_str).unwrap_or(""), edge_count, prepare_us = prepare_end.duration_since(prepare_start).as_micros() as u64, @@ -642,6 +665,11 @@ ) -> Pending { let cursor = HydrantCursor::new(frame.id); let signal = promotion_signal(frame.record.as_ref(), now); + let regime = match frame.record.as_ref() { + Some(r) if r.live => Regime::Live, + Some(_) => Regime::Replay, + None => Regime::NonRecord, + }; let op = match frame.kind { FrameKind::Record => prepare_record(frame.record, resolver).await, FrameKind::Identity | FrameKind::Account => PendingOp::Noop, @@ -653,6 +681,7 @@ Pending { cursor, signal, + regime, op, } } @@ -725,7 +754,12 @@ } async fn resolve_pending(pending: Pending, resolver: &RepoIdResolver) -> Pending { - let Pending { cursor, signal, op } = pending; + let Pending { + cursor, + signal, + regime, + op, + } = pending; let op = match op { PendingOp::Upsert { source, @@ -750,6 +784,7 @@ Pending { cursor, signal, + regime, op, } } @@ -761,7 +796,12 @@ search: &S, records: &dyn RecordStore, ) { - let Pending { cursor, signal, op } = pending; + let Pending { + cursor, + signal, + regime: _, + op, + } = pending; match op { PendingOp::Noop => {} PendingOp::ClearCache { source } => records.remove(&source), @@ -1065,6 +1105,41 @@ bobbin_types::ids::EdgeKey::new(kind, AtUri::new_owned("at://did:plc:uni").unwrap()); assert_eq!(store.count(&old), 0); assert_eq!(store.count(&new), 1); + } + + #[tokio::test] + async fn prepare_frame_tags_regime_from_live_flag() { + let (_store, _cov, resolver) = fresh(); + let mk = |live: bool| -> HydrantFrame { + parse_frame(json!({ + "id": 1, + "type": "record", + "record": { + "live": live, + "did": "did:plc:olaren", + "rev": fresh_tid().as_str(), + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subjectDid": "did:plc:abalone" + } + } + })) + }; + let live_pending = prepare_frame(mk(true), &resolver, now()).await; + assert_eq!(live_pending.regime, Regime::Live); + let replay_pending = prepare_frame(mk(false), &resolver, now()).await; + assert_eq!(replay_pending.regime, Regime::Replay); + + let identity: HydrantFrame = parse_frame(json!({ + "id": 9, + "type": "identity", + })); + let id_pending = prepare_frame(identity, &resolver, now()).await; + assert_eq!(id_pending.regime, Regime::NonRecord); } #[tokio::test] diff --git a/crates/runtime/src/clock.rs b/crates/runtime/src/clock.rs --- a/crates/runtime/src/clock.rs +++ b/crates/runtime/src/clock.rs @@ -73,6 +73,43 @@ } } +#[derive(Clone, Debug)] +pub struct SimClock { + base_unix: UnixMicros, + base_instant: Instant, +} + +impl SimClock { + pub fn at(base_unix: UnixMicros) -> Self { + Self { + base_unix, + base_instant: Instant::now(), + } + } +} + +impl Clock for SimClock { + fn now_unix_micros(&self) -> UnixMicros { + let elapsed = self + .now_instant() + .saturating_duration_since(self.base_instant); + let micros = u64::try_from(elapsed.as_micros()).unwrap_or(u64::MAX); + UnixMicros::new(self.base_unix.raw().saturating_add(micros)) + } + + fn now_instant(&self) -> Instant { + Instant::now() + } + + fn sleep(&self, duration: Duration) -> SleepFuture { + Box::pin(tokio::time::sleep(duration)) + } + + fn sleep_until(&self, deadline: Instant) -> SleepFuture { + Box::pin(tokio::time::sleep_until(deadline)) + } +} + #[cfg(test)] mod tests { use super::*; @@ -139,6 +176,38 @@ assert_eq!( dyn_clock.now_instant(), instant_before + Duration::from_micros(500), + ); + } + + #[tokio::test(start_paused = true)] + async fn sim_clock_anchors_unix_at_explicit_base_under_paused_time() { + let base = UnixMicros::new(1_700_000_000_000_000); + let clock = SimClock::at(base); + assert_eq!( + clock.now_unix_micros(), + base, + "sim clock unix axis must reflect the explicit base, not wall clock", + ); + let bump = Duration::from_secs(10); + tokio::time::advance(bump).await; + assert_eq!( + clock.now_unix_micros().raw(), + base.raw() + u64::try_from(bump.as_micros()).unwrap(), + "sim clock unix axis must follow tokio virtual advance", + ); + } + + #[tokio::test(start_paused = true)] + async fn sim_clock_two_constructions_with_same_base_align() { + let base = UnixMicros::new(2_000_000_000_000_000); + let a = SimClock::at(base); + let b = SimClock::at(base); + assert_eq!(a.now_unix_micros(), b.now_unix_micros()); + tokio::time::advance(Duration::from_secs(5)).await; + assert_eq!( + a.now_unix_micros(), + b.now_unix_micros(), + "two sim clocks anchored at the same base must agree on virtual time", ); } diff --git a/crates/runtime/src/entropy.rs b/crates/runtime/src/entropy.rs --- a/crates/runtime/src/entropy.rs +++ b/crates/runtime/src/entropy.rs @@ -1,3 +1,5 @@ +use std::sync::atomic::{AtomicU64, Ordering}; + pub trait Entropy: Send + Sync + 'static { fn next_u64(&self) -> u64; } @@ -10,6 +12,31 @@ let mut buf = [0u8; 8]; getrandom::fill(&mut buf).expect("os entropy source unavailable"); u64::from_le_bytes(buf) + } +} + +#[derive(Debug)] +pub struct SeededEntropy { + state: AtomicU64, +} + +impl SeededEntropy { + const GOLDEN: u64 = 0x9E37_79B9_7F4A_7C15; + + pub fn new(seed: u64) -> Self { + Self { + state: AtomicU64::new(seed), + } + } +} + +impl Entropy for SeededEntropy { + fn next_u64(&self) -> u64 { + let prev = self.state.fetch_add(Self::GOLDEN, Ordering::Relaxed); + let z = prev.wrapping_add(Self::GOLDEN); + let z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); + let z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); + z ^ (z >> 31) } } @@ -32,6 +59,45 @@ let e: Arc = Arc::new(CountingEntropy(AtomicU64::new(100))); assert_eq!(e.next_u64(), 100); assert_eq!(e.next_u64(), 101); + } + + #[test] + fn seeded_entropy_two_constructions_with_same_seed_yield_same_stream() { + let a = SeededEntropy::new(0xCAFE_F00D); + let b = SeededEntropy::new(0xCAFE_F00D); + for _ in 0..32 { + assert_eq!( + a.next_u64(), + b.next_u64(), + "splitmix stream must reproduce exactly across constructions with same seed", + ); + } + } + + #[test] + fn seeded_entropy_different_seeds_diverge() { + let a = SeededEntropy::new(1); + let b = SeededEntropy::new(2); + let mut all_match = true; + for _ in 0..32 { + if a.next_u64() != b.next_u64() { + all_match = false; + break; + } + } + assert!( + !all_match, + "two seeded entropies with distinct seeds must diverge inside 32 draws", + ); + } + + #[test] + fn seeded_entropy_does_not_repeat_inside_short_window() { + let e = SeededEntropy::new(0); + let mut seen = std::collections::HashSet::new(); + for _ in 0..1024 { + assert!(seen.insert(e.next_u64()), "splitmix collided inside 1024 draws"); + } } #[test] diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs --- a/crates/runtime/src/lib.rs +++ b/crates/runtime/src/lib.rs @@ -1,11 +1,16 @@ mod clock; mod entropy; mod hasher; +mod mem_network; mod network; -pub use clock::{Clock, SleepFuture, SystemClock, UnixMicros}; -pub use entropy::{Entropy, OsEntropy}; +pub use clock::{Clock, SimClock, SleepFuture, SystemClock, UnixMicros}; +pub use entropy::{Entropy, OsEntropy, SeededEntropy}; pub use hasher::RuntimeHasher; +pub use mem_network::{ + DEFAULT_MEM_WS_CAPACITY, MemHttpBody, MemHttpResponder, MemHttpResponse, MemHttpTransport, + MemWsResponder, MemWsServerFuture, MemWsTransport, +}; pub use network::{ BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, WsMessage, WsMessageFuture, diff --git a/crates/runtime/src/mem_network.rs b/crates/runtime/src/mem_network.rs new file mode 100644 --- /dev/null +++ b/crates/runtime/src/mem_network.rs @@ -0,0 +1,352 @@ +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; +use std::time::Duration; + +use bytes::Bytes; +use futures::stream; +use http::{HeaderMap, StatusCode}; +use tokio::sync::mpsc; +use url::Url; + +use crate::clock::{Clock, SleepFuture}; +use crate::network::{ + BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, + NetworkError, WsConn, WsConnectFuture, WsMessage, WsMessageFuture, WsSendFuture, WsSink, + WsStream, WsTransport, +}; + +#[derive(Debug)] +pub struct MemHttpResponse { + pub latency: Duration, + pub result: Result, +} + +#[derive(Debug)] +pub struct MemHttpBody { + pub status: StatusCode, + pub headers: HeaderMap, + pub body: Bytes, +} + +impl MemHttpBody { + pub fn ok_json(body: Bytes) -> Self { + let mut headers = HeaderMap::new(); + headers.insert(http::header::CONTENT_TYPE, "application/json".parse().unwrap()); + Self { + status: StatusCode::OK, + headers, + body, + } + } + + pub fn status_only(status: StatusCode) -> Self { + Self { + status, + headers: HeaderMap::new(), + body: Bytes::new(), + } + } +} + +pub trait MemHttpResponder: Send + Sync + 'static { + fn respond(&self, request: &HttpRequest) -> MemHttpResponse; +} + +#[derive(Clone)] +pub struct MemHttpTransport { + responder: Arc, + clock: Arc, +} + +impl MemHttpTransport { + pub fn new(responder: Arc, clock: Arc) -> Self { + Self { responder, clock } + } + + pub fn shared(responder: Arc, clock: Arc) -> Arc { + Arc::new(Self::new(responder, clock)) + } +} + +impl HttpTransport for MemHttpTransport { + fn execute(&self, request: HttpRequest) -> HttpResponseFuture { + let response = self.responder.respond(&request); + let sleep: SleepFuture = self.clock.sleep(response.latency); + Box::pin(async move { + sleep.await; + into_response_head(response.result) + }) + } +} + +fn into_response_head(result: Result) -> HttpResult { + let body = result?; + let content_length = Some(body.body.len() as u64); + let body_stream: BodyStream = Box::pin(stream::once(async move { Ok(body.body) })); + Ok(HttpResponseHead { + status: body.status, + headers: body.headers, + content_length, + body: body_stream, + }) +} + +pub type MemWsServerFuture = Pin + Send + 'static>>; + +pub trait MemWsResponder: Send + Sync + 'static { + fn spawn_server( + &self, + url: Url, + recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture; +} + +pub const DEFAULT_MEM_WS_CAPACITY: usize = 4096; + +#[derive(Clone)] +pub struct MemWsTransport { + responder: Arc, + capacity: usize, +} + +impl MemWsTransport { + pub fn new(responder: Arc) -> Self { + Self::with_capacity(responder, DEFAULT_MEM_WS_CAPACITY) + } + + pub fn with_capacity(responder: Arc, capacity: usize) -> Self { + assert!(capacity > 0, "MemWsTransport capacity must be > 0"); + Self { + responder, + capacity, + } + } + + pub fn shared(responder: Arc) -> Arc { + Arc::new(Self::new(responder)) + } + + pub fn shared_with_capacity( + responder: Arc, + capacity: usize, + ) -> Arc { + Arc::new(Self::with_capacity(responder, capacity)) + } +} + +impl WsTransport for MemWsTransport { + fn connect(&self, url: Url) -> WsConnectFuture { + let responder = self.responder.clone(); + let capacity = self.capacity; + Box::pin(async move { + let (c2s_tx, c2s_rx) = mpsc::unbounded_channel(); + let (s2c_tx, s2c_rx) = mpsc::channel(capacity); + let server_future = responder.spawn_server(url, c2s_rx, s2c_tx); + tokio::spawn(server_future); + let sink: Box = Box::new(MemWsSink { sender: c2s_tx }); + let stream: Box = Box::new(MemWsStream { receiver: s2c_rx }); + Ok(WsConn { sink, stream }) + }) + } +} + +struct MemWsSink { + sender: mpsc::UnboundedSender, +} + +impl WsSink for MemWsSink { + fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a> { + let result = self + .sender + .send(message) + .map_err(|_| NetworkError::Transport("server side closed channel".into())); + Box::pin(async move { result }) + } +} + +struct MemWsStream { + receiver: mpsc::Receiver, +} + +impl WsStream for MemWsStream { + fn next<'a>(&'a mut self) -> WsMessageFuture<'a> { + Box::pin(async move { self.receiver.recv().await.map(Ok) }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::SimClock; + use crate::UnixMicros; + use std::sync::Mutex; + use std::sync::atomic::{AtomicUsize, Ordering}; + use tokio::time::Instant; + + struct ScriptedHttp { + responses: Mutex>, + cursor: AtomicUsize, + } + + impl ScriptedHttp { + fn new(responses: Vec) -> Self { + Self { + responses: Mutex::new(responses), + cursor: AtomicUsize::new(0), + } + } + } + + impl MemHttpResponder for ScriptedHttp { + fn respond(&self, _: &HttpRequest) -> MemHttpResponse { + let i = self.cursor.fetch_add(1, Ordering::Relaxed); + let mut guard = self.responses.lock().unwrap(); + let next = std::mem::replace( + &mut guard[i], + MemHttpResponse { + latency: Duration::ZERO, + result: Err(NetworkError::Transport("script consumed".into())), + }, + ); + next + } + } + + fn ok_response(body: &str, latency_ms: u64) -> MemHttpResponse { + MemHttpResponse { + latency: Duration::from_millis(latency_ms), + result: Ok(MemHttpBody::ok_json(Bytes::from(body.to_owned()))), + } + } + + fn err_response(latency_ms: u64) -> MemHttpResponse { + MemHttpResponse { + latency: Duration::from_millis(latency_ms), + result: Err(NetworkError::Transport("brownout".into())), + } + } + + #[tokio::test(start_paused = true)] + async fn http_returns_scripted_body_after_injected_latency() { + let clock: Arc = Arc::new(SimClock::at(UnixMicros::new(0))); + let responder: Arc = Arc::new(ScriptedHttp::new(vec![ + ok_response("{\"hello\":1}", 10), + ])); + let transport = MemHttpTransport::new(responder, clock); + let request = HttpRequest { + url: Url::parse("http://oyster.cafe/xrpc/x").unwrap(), + headers: HeaderMap::new(), + }; + + let before = Instant::now(); + let resp = transport.execute(request).await.unwrap(); + let elapsed = Instant::now().saturating_duration_since(before); + + assert_eq!(resp.status, StatusCode::OK); + assert_eq!(elapsed, Duration::from_millis(10)); + let chunks = collect_body(resp.body).await; + assert_eq!(chunks.as_ref(), b"{\"hello\":1}"); + } + + #[tokio::test(start_paused = true)] + async fn http_propagates_scripted_errors_with_latency() { + let clock: Arc = Arc::new(SimClock::at(UnixMicros::new(0))); + let responder: Arc = + Arc::new(ScriptedHttp::new(vec![err_response(50)])); + let transport = MemHttpTransport::new(responder, clock); + let request = HttpRequest { + url: Url::parse("http://oyster.cafe/xrpc/x").unwrap(), + headers: HeaderMap::new(), + }; + + let before = Instant::now(); + let resp = transport.execute(request).await; + let elapsed = Instant::now().saturating_duration_since(before); + + assert_eq!(elapsed, Duration::from_millis(50)); + assert!(matches!(resp, Err(NetworkError::Transport(_)))); + } + + async fn collect_body(mut body: BodyStream) -> Bytes { + use futures::StreamExt; + let mut acc = Vec::new(); + while let Some(chunk) = body.next().await { + acc.extend_from_slice(&chunk.unwrap()); + } + Bytes::from(acc) + } + + struct ScriptedWs { + frames: Mutex>, + } + + impl MemWsResponder for ScriptedWs { + fn spawn_server( + &self, + _: Url, + mut _recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + let frames: Vec = std::mem::take(&mut *self.frames.lock().unwrap()); + Box::pin(async move { + for frame in frames { + if send.send(frame).await.is_err() { + return; + } + } + }) + } + } + + #[tokio::test(start_paused = true)] + async fn ws_delivers_scripted_frames_to_client() { + let responder: Arc = Arc::new(ScriptedWs { + frames: Mutex::new(vec![ + WsMessage::Text("frame-a".into()), + WsMessage::Text("frame-b".into()), + ]), + }); + let transport = MemWsTransport::new(responder); + let mut conn = transport + .connect(Url::parse("ws://oyster.cafe/").unwrap()) + .await + .unwrap(); + let a = conn.stream.next().await.unwrap().unwrap(); + let b = conn.stream.next().await.unwrap().unwrap(); + assert!(matches!(a, WsMessage::Text(t) if t == "frame-a")); + assert!(matches!(b, WsMessage::Text(t) if t == "frame-b")); + } + + struct EchoWs; + + impl MemWsResponder for EchoWs { + fn spawn_server( + &self, + _: Url, + mut recv: mpsc::UnboundedReceiver, + send: mpsc::Sender, + ) -> MemWsServerFuture { + Box::pin(async move { + while let Some(msg) = recv.recv().await { + if send.send(msg).await.is_err() { + return; + } + } + }) + } + } + + #[tokio::test(start_paused = true)] + async fn ws_round_trip_via_server_echo() { + let transport = MemWsTransport::new(Arc::new(EchoWs)); + let mut conn = transport + .connect(Url::parse("ws://oyster.cafe/").unwrap()) + .await + .unwrap(); + conn.sink.send(WsMessage::Text("ping".into())).await.unwrap(); + let echoed = conn.stream.next().await.unwrap().unwrap(); + assert!(matches!(echoed, WsMessage::Text(t) if t == "ping")); + } +} -- tangled.sh