diff --git a/Cargo.lock b/Cargo.lock index 54ff9d1..da518e3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -290,6 +290,7 @@ dependencies = [ "bobbin-ingest", "bobbin-knot-proxy", "bobbin-record-lru", + "bobbin-runtime", "bobbin-search", "bobbin-slingshot-client", "bobbin-xrpc", @@ -310,6 +311,7 @@ dependencies = [ name = "bobbin-edge-index" version = "0.0.1" dependencies = [ + "bobbin-runtime", "bobbin-types", "jacquard-common", "lasso", @@ -325,6 +327,7 @@ version = "0.0.1" dependencies = [ "bobbin-edge-index", "bobbin-record-lru", + "bobbin-runtime", "bobbin-slingshot-client", "bobbin-types", "bytes", @@ -335,7 +338,6 @@ dependencies = [ "serde_json", "thiserror 2.0.18", "tokio", - "tokio-tungstenite 0.29.0", "tokio-util", "tracing", "tracing-subscriber", @@ -347,6 +349,7 @@ dependencies = [ name = "bobbin-knot-proxy" version = "0.0.1" dependencies = [ + "bobbin-runtime", "bytes", "futures", "http", @@ -368,10 +371,27 @@ dependencies = [ "quick_cache", ] +[[package]] +name = "bobbin-runtime" +version = "0.0.1" +dependencies = [ + "ahash", + "bytes", + "futures", + "getrandom 0.3.4", + "http", + "reqwest", + "thiserror 2.0.18", + "tokio", + "tokio-tungstenite 0.29.0", + "url", +] + [[package]] name = "bobbin-search" version = "0.0.1" dependencies = [ + "bobbin-runtime", "bobbin-types", "jacquard-common", "tantivy", @@ -384,10 +404,12 @@ dependencies = [ name = "bobbin-slingshot-client" version = "0.0.1" dependencies = [ + "bobbin-runtime", "bobbin-types", "bytes", "cid", "futures", + "http", "jacquard-common", "reqwest", "serde", @@ -424,6 +446,7 @@ dependencies = [ "bobbin-edge-index", "bobbin-knot-proxy", "bobbin-record-lru", + "bobbin-runtime", "bobbin-search", "bobbin-slingshot-client", "bobbin-types", diff --git a/Cargo.toml b/Cargo.toml index 18339b1..aee74c9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -8,6 +8,7 @@ members = [ "crates/slingshot-client", "crates/record-lru", "crates/knot-proxy", + "crates/runtime", "crates/search", "crates/xrpc", ] @@ -25,6 +26,7 @@ bobbin-ingest = { path = "crates/ingest" } bobbin-slingshot-client = { path = "crates/slingshot-client" } bobbin-record-lru = { path = "crates/record-lru" } bobbin-knot-proxy = { path = "crates/knot-proxy" } +bobbin-runtime = { path = "crates/runtime" } bobbin-search = { path = "crates/search" } bobbin-xrpc = { path = "crates/xrpc" } @@ -54,6 +56,8 @@ scc = "3" roaring = "0.11" lasso = { version = "0.7", features = ["multi-threaded"] } quick_cache = "0.6" +getrandom = "0.3" +ahash = { version = "0.8", default-features = false, features = ["std"] } reqwest = { version = "0.12", default-features = false, features = ["rustls-tls-webpki-roots", "http2", "json", "gzip", "stream"] } axum = "0.8" diff --git a/crates/runtime/Cargo.toml b/crates/runtime/Cargo.toml new file mode 100644 index 0000000..4d78689 --- /dev/null +++ b/crates/runtime/Cargo.toml @@ -0,0 +1,21 @@ +[package] +name = "bobbin-runtime" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +ahash = { workspace = true } +bytes = { workspace = true } +futures = { workspace = true } +getrandom = { workspace = true } +http = { workspace = true } +reqwest = { workspace = true } +thiserror = { workspace = true } +tokio = { workspace = true } +tokio-tungstenite = { workspace = true } +url = { workspace = true } + +[dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } diff --git a/crates/runtime/src/clock.rs b/crates/runtime/src/clock.rs new file mode 100644 index 0000000..6e59d97 --- /dev/null +++ b/crates/runtime/src/clock.rs @@ -0,0 +1,165 @@ +use std::future::Future; +use std::pin::Pin; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; + +use tokio::time::Instant; + +#[derive(Clone, Copy, Debug, Default, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct UnixMicros(u64); + +impl UnixMicros { + pub const fn new(value: u64) -> Self { + Self(value) + } + + pub const fn raw(self) -> u64 { + self.0 + } +} + +pub type SleepFuture = Pin + Send + 'static>>; + +pub trait Clock: Send + Sync + 'static { + fn now_unix_micros(&self) -> UnixMicros; + fn now_instant(&self) -> Instant; + fn sleep(&self, duration: Duration) -> SleepFuture; + fn sleep_until(&self, deadline: Instant) -> SleepFuture; +} + +#[derive(Clone, Copy, Debug)] +pub struct SystemClock { + base_unix: UnixMicros, + base_instant: Instant, +} + +impl SystemClock { + pub fn new() -> Self { + let raw = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("system clock before unix epoch") + .as_micros(); + Self { + base_unix: UnixMicros::new(u64::try_from(raw).unwrap_or(u64::MAX)), + base_instant: Instant::now(), + } + } +} + +impl Default for SystemClock { + fn default() -> Self { + Self::new() + } +} + +impl Clock for SystemClock { + 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::*; + use std::sync::Arc; + use std::sync::atomic::{AtomicU64, Ordering}; + + struct ManualClock { + base: Instant, + unix_micros: AtomicU64, + } + + impl ManualClock { + fn new(unix_micros: u64) -> Self { + Self { + base: Instant::now(), + unix_micros: AtomicU64::new(unix_micros), + } + } + + fn advance(&self, micros: u64) { + self.unix_micros.fetch_add(micros, Ordering::SeqCst); + } + } + + impl Clock for ManualClock { + fn now_unix_micros(&self) -> UnixMicros { + UnixMicros::new(self.unix_micros.load(Ordering::SeqCst)) + } + + fn now_instant(&self) -> Instant { + self.base + Duration::from_micros(self.unix_micros.load(Ordering::SeqCst)) + } + + fn sleep(&self, _: Duration) -> SleepFuture { + Box::pin(std::future::ready(())) + } + + fn sleep_until(&self, _: Instant) -> SleepFuture { + Box::pin(std::future::ready(())) + } + } + + #[tokio::test] + async fn external_clock_works_behind_dyn_arc() { + let clock: Arc = Arc::new(ManualClock::new(42)); + assert_eq!(clock.now_unix_micros().raw(), 42); + let before = clock.now_instant(); + clock.sleep(Duration::from_secs(60)).await; + assert_eq!( + clock.now_instant(), + before, + "manual clock's sleep must not advance virtual time on its own", + ); + } + + #[tokio::test] + async fn manual_clock_advance_moves_both_unix_and_instant() { + let clock = Arc::new(ManualClock::new(1_000)); + let dyn_clock: Arc = clock.clone(); + let unix_before = dyn_clock.now_unix_micros().raw(); + let instant_before = dyn_clock.now_instant(); + clock.advance(500); + assert_eq!(dyn_clock.now_unix_micros().raw(), unix_before + 500); + assert_eq!( + dyn_clock.now_instant(), + instant_before + Duration::from_micros(500), + ); + } + + #[tokio::test(start_paused = true)] + async fn system_clock_axes_advance_together_under_paused_time() { + let clock = SystemClock::new(); + let unix_before = clock.now_unix_micros().raw(); + let instant_before = clock.now_instant(); + let bump = Duration::from_secs(60); + tokio::time::advance(bump).await; + let unix_after = clock.now_unix_micros().raw(); + let instant_after = clock.now_instant(); + assert_eq!( + instant_after.saturating_duration_since(instant_before), + bump, + "instant axis must reflect virtual advance", + ); + assert_eq!( + unix_after - unix_before, + u64::try_from(bump.as_micros()).unwrap(), + "unix axis must follow the same virtual advance, not wall clock", + ); + } +} diff --git a/crates/runtime/src/entropy.rs b/crates/runtime/src/entropy.rs new file mode 100644 index 0000000..f1e258b --- /dev/null +++ b/crates/runtime/src/entropy.rs @@ -0,0 +1,47 @@ +pub trait Entropy: Send + Sync + 'static { + fn next_u64(&self) -> u64; +} + +#[derive(Clone, Copy, Debug, Default)] +pub struct OsEntropy; + +impl Entropy for OsEntropy { + fn next_u64(&self) -> u64 { + let mut buf = [0u8; 8]; + getrandom::fill(&mut buf).expect("os entropy source unavailable"); + u64::from_le_bytes(buf) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::Arc; + use std::sync::atomic::{AtomicU64, Ordering}; + + struct CountingEntropy(AtomicU64); + + impl Entropy for CountingEntropy { + fn next_u64(&self) -> u64 { + self.0.fetch_add(1, Ordering::SeqCst) + } + } + + #[test] + fn external_impl_works_behind_dyn_arc() { + let e: Arc = Arc::new(CountingEntropy(AtomicU64::new(100))); + assert_eq!(e.next_u64(), 100); + assert_eq!(e.next_u64(), 101); + } + + #[test] + fn os_entropy_does_not_collide_on_back_to_back_calls() { + let e = OsEntropy; + let a = e.next_u64(); + let b = e.next_u64(); + assert_ne!( + a, b, + "back-to-back os entropy collided; getrandom likely not actually wired up", + ); + } +} diff --git a/crates/runtime/src/hasher.rs b/crates/runtime/src/hasher.rs new file mode 100644 index 0000000..5c21d63 --- /dev/null +++ b/crates/runtime/src/hasher.rs @@ -0,0 +1,101 @@ +use std::hash::BuildHasher; + +use ahash::{AHasher, RandomState}; + +use crate::{Entropy, OsEntropy}; + +#[derive(Clone, Debug)] +pub struct RuntimeHasher { + inner: RandomState, +} + +impl RuntimeHasher { + pub fn from_entropy(entropy: &dyn Entropy) -> Self { + Self { + inner: RandomState::with_seeds( + entropy.next_u64(), + entropy.next_u64(), + entropy.next_u64(), + entropy.next_u64(), + ), + } + } + + pub const fn from_seeds(k0: u64, k1: u64, k2: u64, k3: u64) -> Self { + Self { + inner: RandomState::with_seeds(k0, k1, k2, k3), + } + } +} + +impl Default for RuntimeHasher { + fn default() -> Self { + Self::from_entropy(&OsEntropy) + } +} + +impl BuildHasher for RuntimeHasher { + type Hasher = AHasher; + + fn build_hasher(&self) -> Self::Hasher { + self.inner.build_hasher() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::OsEntropy; + use std::collections::HashMap; + use std::sync::Arc; + use std::sync::atomic::{AtomicU64, Ordering}; + + struct CountingEntropy(AtomicU64); + + impl Entropy for CountingEntropy { + fn next_u64(&self) -> u64 { + self.0.fetch_add(1, Ordering::SeqCst) + } + } + + #[test] + fn same_seed_inputs_produce_identical_hashes() { + let a = RuntimeHasher::from_seeds(1, 2, 3, 4); + let b = RuntimeHasher::from_seeds(1, 2, 3, 4); + assert_eq!(a.hash_one("limpet"), b.hash_one("limpet")); + assert_eq!(a.hash_one(42_u64), b.hash_one(42_u64)); + } + + #[test] + fn different_seeds_produce_different_hashes() { + let a = RuntimeHasher::from_seeds(1, 2, 3, 4); + let b = RuntimeHasher::from_seeds(5, 6, 7, 8); + assert_ne!( + a.hash_one("conch"), + b.hash_one("conch"), + "two seedings must not collapse to the same key state", + ); + } + + #[test] + fn from_entropy_consumes_four_words_in_call_order() { + let entropy: Arc = Arc::new(CountingEntropy(AtomicU64::new(100))); + let hasher = RuntimeHasher::from_entropy(&*entropy); + let direct = RuntimeHasher::from_seeds(100, 101, 102, 103); + assert_eq!( + hasher.hash_one("nautilus"), + direct.hash_one("nautilus"), + "from_entropy must seed by draining four u64s in order", + ); + } + + #[test] + fn os_entropy_seeds_a_usable_hashmap() { + let hasher = RuntimeHasher::from_entropy(&OsEntropy); + let mut map: HashMap<&str, u32, RuntimeHasher> = HashMap::with_hasher(hasher); + map.insert("kelp", 1); + map.insert("uni", 2); + assert_eq!(map.get("kelp"), Some(&1)); + assert_eq!(map.get("uni"), Some(&2)); + } +} diff --git a/crates/runtime/src/lib.rs b/crates/runtime/src/lib.rs new file mode 100644 index 0000000..635fba3 --- /dev/null +++ b/crates/runtime/src/lib.rs @@ -0,0 +1,13 @@ +mod clock; +mod entropy; +mod hasher; +mod network; + +pub use clock::{Clock, SleepFuture, SystemClock, UnixMicros}; +pub use entropy::{Entropy, OsEntropy}; +pub use hasher::RuntimeHasher; +pub use network::{ + BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, + NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, WsMessage, WsMessageFuture, + WsSendFuture, WsSink, WsStream, WsTransport, +}; diff --git a/crates/runtime/src/network.rs b/crates/runtime/src/network.rs new file mode 100644 index 0000000..65c2efe --- /dev/null +++ b/crates/runtime/src/network.rs @@ -0,0 +1,230 @@ +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; + +use bytes::Bytes; +use futures::stream::{Stream, StreamExt}; +use http::{HeaderMap, StatusCode}; +use thiserror::Error; +use tokio_tungstenite::tungstenite::{ + Bytes as WsBytes, Message as TungsteniteMessage, protocol::CloseFrame as TungsteniteClose, + protocol::frame::coding::CloseCode as TungsteniteCloseCode, +}; +use url::Url; + +#[derive(Debug, Error)] +pub enum NetworkError { + #[error("connect: {0}")] + Connect(String), + #[error("timeout: {0}")] + Timeout(String), + #[error("redirect: {0}")] + Redirect(String), + #[error("transport: {0}")] + Transport(String), + #[error("body: {0}")] + Body(String), + #[error("protocol: {0}")] + Protocol(String), +} + +pub struct HttpRequest { + pub url: Url, + pub headers: HeaderMap, +} + +pub type BodyStream = Pin> + Send + 'static>>; + +pub struct HttpResponseHead { + pub status: StatusCode, + pub headers: HeaderMap, + pub content_length: Option, + pub body: BodyStream, +} + +pub type HttpResult = Result; +pub type HttpResponseFuture = Pin + Send + 'static>>; + +pub trait HttpTransport: Send + Sync + 'static { + fn execute(&self, request: HttpRequest) -> HttpResponseFuture; +} + +#[derive(Clone, Debug)] +pub struct ReqwestHttp { + client: reqwest::Client, +} + +impl ReqwestHttp { + pub fn new(client: reqwest::Client) -> Self { + Self { client } + } + + pub fn shared(client: reqwest::Client) -> Arc { + Arc::new(Self::new(client)) + } +} + +impl HttpTransport for ReqwestHttp { + fn execute(&self, request: HttpRequest) -> HttpResponseFuture { + let client = self.client.clone(); + Box::pin(async move { + let resp = client + .get(request.url) + .headers(request.headers) + .send() + .await + .map_err(map_reqwest)?; + let status = resp.status(); + let headers = resp.headers().clone(); + let content_length = resp.content_length(); + let body: BodyStream = Box::pin( + resp.bytes_stream() + .map(|chunk| chunk.map_err(|e| NetworkError::Body(e.to_string()))), + ); + Ok(HttpResponseHead { + status, + headers, + content_length, + body, + }) + }) + } +} + +fn map_reqwest(err: reqwest::Error) -> NetworkError { + let msg = err.to_string(); + if err.is_timeout() { + NetworkError::Timeout(msg) + } else if err.is_connect() { + NetworkError::Connect(msg) + } else if err.is_redirect() { + NetworkError::Redirect(msg) + } else { + NetworkError::Transport(msg) + } +} + +#[derive(Clone, Debug)] +pub enum WsMessage { + Text(String), + Binary(Bytes), + Ping(Bytes), + Pong(Bytes), + Close { code: u16, reason: String }, +} + +pub type WsSendFuture<'a> = Pin> + Send + 'a>>; +pub type WsMessageFuture<'a> = + Pin>> + Send + 'a>>; + +pub trait WsSink: Send + 'static { + fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a>; +} + +pub trait WsStream: Send + 'static { + fn next<'a>(&'a mut self) -> WsMessageFuture<'a>; +} + +pub struct WsConn { + pub sink: Box, + pub stream: Box, +} + +pub type WsConnectFuture = + Pin> + Send + 'static>>; + +pub trait WsTransport: Send + Sync + 'static { + fn connect(&self, url: Url) -> WsConnectFuture; +} + +#[derive(Clone, Copy, Debug, Default)] +pub struct TungsteniteWs; + +impl TungsteniteWs { + pub fn shared() -> Arc { + Arc::new(Self) + } +} + +impl WsTransport for TungsteniteWs { + fn connect(&self, url: Url) -> WsConnectFuture { + Box::pin(async move { + let url_str = url.as_str().to_owned(); + let (ws, _resp) = tokio_tungstenite::connect_async(&url_str) + .await + .map_err(|e| NetworkError::Connect(e.to_string()))?; + let (sink_inner, stream_inner) = futures::StreamExt::split(ws); + let sink: Box = Box::new(TungsteniteSink { inner: sink_inner }); + let stream: Box = Box::new(TungsteniteStream { + inner: stream_inner, + }); + Ok(WsConn { sink, stream }) + }) + } +} + +type TungsteniteWsStream = + tokio_tungstenite::WebSocketStream>; + +struct TungsteniteSink { + inner: futures::stream::SplitSink, +} + +impl WsSink for TungsteniteSink { + fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a> { + Box::pin(async move { + use futures::SinkExt; + self.inner + .send(message_to_tungstenite(message)) + .await + .map_err(|e| NetworkError::Transport(e.to_string())) + }) + } +} + +struct TungsteniteStream { + inner: futures::stream::SplitStream, +} + +impl WsStream for TungsteniteStream { + fn next<'a>(&'a mut self) -> WsMessageFuture<'a> { + Box::pin(async move { + let item = StreamExt::next(&mut self.inner).await?; + Some( + item.map_err(|e| NetworkError::Transport(e.to_string())) + .and_then(message_from_tungstenite), + ) + }) + } +} + +fn message_to_tungstenite(message: WsMessage) -> TungsteniteMessage { + match message { + WsMessage::Text(text) => TungsteniteMessage::Text(text.into()), + WsMessage::Binary(bytes) => TungsteniteMessage::Binary(WsBytes::copy_from_slice(&bytes)), + WsMessage::Ping(bytes) => TungsteniteMessage::Ping(WsBytes::copy_from_slice(&bytes)), + WsMessage::Pong(bytes) => TungsteniteMessage::Pong(WsBytes::copy_from_slice(&bytes)), + WsMessage::Close { code, reason } => TungsteniteMessage::Close(Some(TungsteniteClose { + code: TungsteniteCloseCode::from(code), + reason: reason.into(), + })), + } +} + +fn message_from_tungstenite(message: TungsteniteMessage) -> Result { + match message { + TungsteniteMessage::Text(t) => Ok(WsMessage::Text(t.to_string())), + TungsteniteMessage::Binary(b) => Ok(WsMessage::Binary(Bytes::copy_from_slice(&b))), + TungsteniteMessage::Ping(b) => Ok(WsMessage::Ping(Bytes::copy_from_slice(&b))), + TungsteniteMessage::Pong(b) => Ok(WsMessage::Pong(Bytes::copy_from_slice(&b))), + TungsteniteMessage::Close(close) => { + let (code, reason) = close + .map(|c| (u16::from(c.code), c.reason.to_string())) + .unwrap_or((1000, String::new())); + Ok(WsMessage::Close { code, reason }) + } + TungsteniteMessage::Frame(_) => Err(NetworkError::Protocol( + "tungstenite raw frame surfaced unexpectedly".to_owned(), + )), + } +}