This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

feat(runtime): SimClock, SeededEntropy, MemNetwork

Lewis: May this revision serve well! <lu5a@proton.me>

Lewis (May 9, 2026, 9:50 AM +0300) 126cb5e1 965f6cf8

+572 -5
+78 -3
crates/ingest/src/lib.rs
··· 190 190 backoff = RECONNECT_INITIAL_DELAY; 191 191 } else { 192 192 tokio::select! { 193 + biased; 194 + _ = runtime.cancel.cancelled() => return Ok(()), 193 195 _ = runtime.clock.sleep(jittered(backoff, &*runtime.entropy)) => {} 194 - _ = runtime.cancel.cancelled() => return Ok(()), 195 196 } 196 197 backoff = (backoff * 2).min(RECONNECT_MAX_DELAY); 197 198 } ··· 518 519 } 519 520 } 520 521 522 + #[derive(Clone, Copy, Debug, Eq, PartialEq)] 523 + enum Regime { 524 + Replay, 525 + Live, 526 + NonRecord, 527 + } 528 + 529 + impl Regime { 530 + fn as_str(self) -> &'static str { 531 + match self { 532 + Self::Replay => "replay", 533 + Self::Live => "live", 534 + Self::NonRecord => "non_record", 535 + } 536 + } 537 + } 538 + 521 539 struct Pending { 522 540 cursor: HydrantCursor, 523 541 signal: PromotionSignal, 542 + regime: Regime, 524 543 op: PendingOp, 525 544 } 526 545 ··· 610 629 ) { 611 630 let nsid = pending_nsid(&staged.pending.op).cloned(); 612 631 let edge_count = pending_edge_count(&staged.pending.op); 632 + let regime = staged.pending.regime; 633 + let cursor = staged.pending.cursor.raw(); 613 634 let Resolved { 614 635 pending, 615 636 prepare_start, ··· 622 643 let commit_end = rt.clock.now_instant(); 623 644 tracing::trace!( 624 645 target: "bobbin_ingest::stage", 646 + cursor, 647 + regime = regime.as_str(), 625 648 nsid = %nsid.as_ref().map(Nsid::as_str).unwrap_or(""), 626 649 edge_count, 627 650 prepare_us = prepare_end.duration_since(prepare_start).as_micros() as u64, ··· 642 665 ) -> Pending { 643 666 let cursor = HydrantCursor::new(frame.id); 644 667 let signal = promotion_signal(frame.record.as_ref(), now); 668 + let regime = match frame.record.as_ref() { 669 + Some(r) if r.live => Regime::Live, 670 + Some(_) => Regime::Replay, 671 + None => Regime::NonRecord, 672 + }; 645 673 let op = match frame.kind { 646 674 FrameKind::Record => prepare_record(frame.record, resolver).await, 647 675 FrameKind::Identity | FrameKind::Account => PendingOp::Noop, ··· 653 681 Pending { 654 682 cursor, 655 683 signal, 684 + regime, 656 685 op, 657 686 } 658 687 } ··· 725 754 } 726 755 727 756 async fn resolve_pending(pending: Pending, resolver: &RepoIdResolver) -> Pending { 728 - let Pending { cursor, signal, op } = pending; 757 + let Pending { 758 + cursor, 759 + signal, 760 + regime, 761 + op, 762 + } = pending; 729 763 let op = match op { 730 764 PendingOp::Upsert { 731 765 source, ··· 750 784 Pending { 751 785 cursor, 752 786 signal, 787 + regime, 753 788 op, 754 789 } 755 790 } ··· 761 796 search: &S, 762 797 records: &dyn RecordStore, 763 798 ) { 764 - let Pending { cursor, signal, op } = pending; 799 + let Pending { 800 + cursor, 801 + signal, 802 + regime: _, 803 + op, 804 + } = pending; 765 805 match op { 766 806 PendingOp::Noop => {} 767 807 PendingOp::ClearCache { source } => records.remove(&source), ··· 1065 1105 bobbin_types::ids::EdgeKey::new(kind, AtUri::new_owned("at://did:plc:uni").unwrap()); 1066 1106 assert_eq!(store.count(&old), 0); 1067 1107 assert_eq!(store.count(&new), 1); 1108 + } 1109 + 1110 + #[tokio::test] 1111 + async fn prepare_frame_tags_regime_from_live_flag() { 1112 + let (_store, _cov, resolver) = fresh(); 1113 + let mk = |live: bool| -> HydrantFrame { 1114 + parse_frame(json!({ 1115 + "id": 1, 1116 + "type": "record", 1117 + "record": { 1118 + "live": live, 1119 + "did": "did:plc:olaren", 1120 + "rev": fresh_tid().as_str(), 1121 + "collection": "sh.tangled.feed.star", 1122 + "rkey": "abcabcabcabcz", 1123 + "action": "create", 1124 + "record": { 1125 + "$type": "sh.tangled.feed.star", 1126 + "createdAt": "2026-05-01T00:00:00Z", 1127 + "subjectDid": "did:plc:abalone" 1128 + } 1129 + } 1130 + })) 1131 + }; 1132 + let live_pending = prepare_frame(mk(true), &resolver, now()).await; 1133 + assert_eq!(live_pending.regime, Regime::Live); 1134 + let replay_pending = prepare_frame(mk(false), &resolver, now()).await; 1135 + assert_eq!(replay_pending.regime, Regime::Replay); 1136 + 1137 + let identity: HydrantFrame = parse_frame(json!({ 1138 + "id": 9, 1139 + "type": "identity", 1140 + })); 1141 + let id_pending = prepare_frame(identity, &resolver, now()).await; 1142 + assert_eq!(id_pending.regime, Regime::NonRecord); 1068 1143 } 1069 1144 1070 1145 #[tokio::test]
+69
crates/runtime/src/clock.rs
··· 73 73 } 74 74 } 75 75 76 + #[derive(Clone, Debug)] 77 + pub struct SimClock { 78 + base_unix: UnixMicros, 79 + base_instant: Instant, 80 + } 81 + 82 + impl SimClock { 83 + pub fn at(base_unix: UnixMicros) -> Self { 84 + Self { 85 + base_unix, 86 + base_instant: Instant::now(), 87 + } 88 + } 89 + } 90 + 91 + impl Clock for SimClock { 92 + fn now_unix_micros(&self) -> UnixMicros { 93 + let elapsed = self 94 + .now_instant() 95 + .saturating_duration_since(self.base_instant); 96 + let micros = u64::try_from(elapsed.as_micros()).unwrap_or(u64::MAX); 97 + UnixMicros::new(self.base_unix.raw().saturating_add(micros)) 98 + } 99 + 100 + fn now_instant(&self) -> Instant { 101 + Instant::now() 102 + } 103 + 104 + fn sleep(&self, duration: Duration) -> SleepFuture { 105 + Box::pin(tokio::time::sleep(duration)) 106 + } 107 + 108 + fn sleep_until(&self, deadline: Instant) -> SleepFuture { 109 + Box::pin(tokio::time::sleep_until(deadline)) 110 + } 111 + } 112 + 76 113 #[cfg(test)] 77 114 mod tests { 78 115 use super::*; ··· 139 176 assert_eq!( 140 177 dyn_clock.now_instant(), 141 178 instant_before + Duration::from_micros(500), 179 + ); 180 + } 181 + 182 + #[tokio::test(start_paused = true)] 183 + async fn sim_clock_anchors_unix_at_explicit_base_under_paused_time() { 184 + let base = UnixMicros::new(1_700_000_000_000_000); 185 + let clock = SimClock::at(base); 186 + assert_eq!( 187 + clock.now_unix_micros(), 188 + base, 189 + "sim clock unix axis must reflect the explicit base, not wall clock", 190 + ); 191 + let bump = Duration::from_secs(10); 192 + tokio::time::advance(bump).await; 193 + assert_eq!( 194 + clock.now_unix_micros().raw(), 195 + base.raw() + u64::try_from(bump.as_micros()).unwrap(), 196 + "sim clock unix axis must follow tokio virtual advance", 197 + ); 198 + } 199 + 200 + #[tokio::test(start_paused = true)] 201 + async fn sim_clock_two_constructions_with_same_base_align() { 202 + let base = UnixMicros::new(2_000_000_000_000_000); 203 + let a = SimClock::at(base); 204 + let b = SimClock::at(base); 205 + assert_eq!(a.now_unix_micros(), b.now_unix_micros()); 206 + tokio::time::advance(Duration::from_secs(5)).await; 207 + assert_eq!( 208 + a.now_unix_micros(), 209 + b.now_unix_micros(), 210 + "two sim clocks anchored at the same base must agree on virtual time", 142 211 ); 143 212 } 144 213
+66
crates/runtime/src/entropy.rs
··· 1 + use std::sync::atomic::{AtomicU64, Ordering}; 2 + 1 3 pub trait Entropy: Send + Sync + 'static { 2 4 fn next_u64(&self) -> u64; 3 5 } ··· 13 15 } 14 16 } 15 17 18 + #[derive(Debug)] 19 + pub struct SeededEntropy { 20 + state: AtomicU64, 21 + } 22 + 23 + impl SeededEntropy { 24 + const GOLDEN: u64 = 0x9E37_79B9_7F4A_7C15; 25 + 26 + pub fn new(seed: u64) -> Self { 27 + Self { 28 + state: AtomicU64::new(seed), 29 + } 30 + } 31 + } 32 + 33 + impl Entropy for SeededEntropy { 34 + fn next_u64(&self) -> u64 { 35 + let prev = self.state.fetch_add(Self::GOLDEN, Ordering::Relaxed); 36 + let z = prev.wrapping_add(Self::GOLDEN); 37 + let z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9); 38 + let z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB); 39 + z ^ (z >> 31) 40 + } 41 + } 42 + 16 43 #[cfg(test)] 17 44 mod tests { 18 45 use super::*; ··· 32 59 let e: Arc<dyn Entropy> = Arc::new(CountingEntropy(AtomicU64::new(100))); 33 60 assert_eq!(e.next_u64(), 100); 34 61 assert_eq!(e.next_u64(), 101); 62 + } 63 + 64 + #[test] 65 + fn seeded_entropy_two_constructions_with_same_seed_yield_same_stream() { 66 + let a = SeededEntropy::new(0xCAFE_F00D); 67 + let b = SeededEntropy::new(0xCAFE_F00D); 68 + for _ in 0..32 { 69 + assert_eq!( 70 + a.next_u64(), 71 + b.next_u64(), 72 + "splitmix stream must reproduce exactly across constructions with same seed", 73 + ); 74 + } 75 + } 76 + 77 + #[test] 78 + fn seeded_entropy_different_seeds_diverge() { 79 + let a = SeededEntropy::new(1); 80 + let b = SeededEntropy::new(2); 81 + let mut all_match = true; 82 + for _ in 0..32 { 83 + if a.next_u64() != b.next_u64() { 84 + all_match = false; 85 + break; 86 + } 87 + } 88 + assert!( 89 + !all_match, 90 + "two seeded entropies with distinct seeds must diverge inside 32 draws", 91 + ); 92 + } 93 + 94 + #[test] 95 + fn seeded_entropy_does_not_repeat_inside_short_window() { 96 + let e = SeededEntropy::new(0); 97 + let mut seen = std::collections::HashSet::new(); 98 + for _ in 0..1024 { 99 + assert!(seen.insert(e.next_u64()), "splitmix collided inside 1024 draws"); 100 + } 35 101 } 36 102 37 103 #[test]
+7 -2
crates/runtime/src/lib.rs
··· 1 1 mod clock; 2 2 mod entropy; 3 3 mod hasher; 4 + mod mem_network; 4 5 mod network; 5 6 6 - pub use clock::{Clock, SleepFuture, SystemClock, UnixMicros}; 7 - pub use entropy::{Entropy, OsEntropy}; 7 + pub use clock::{Clock, SimClock, SleepFuture, SystemClock, UnixMicros}; 8 + pub use entropy::{Entropy, OsEntropy, SeededEntropy}; 8 9 pub use hasher::RuntimeHasher; 10 + pub use mem_network::{ 11 + DEFAULT_MEM_WS_CAPACITY, MemHttpBody, MemHttpResponder, MemHttpResponse, MemHttpTransport, 12 + MemWsResponder, MemWsServerFuture, MemWsTransport, 13 + }; 9 14 pub use network::{ 10 15 BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, 11 16 NetworkError, ReqwestHttp, TungsteniteWs, WsConn, WsConnectFuture, WsMessage, WsMessageFuture,
+352
crates/runtime/src/mem_network.rs
··· 1 + use std::future::Future; 2 + use std::pin::Pin; 3 + use std::sync::Arc; 4 + use std::time::Duration; 5 + 6 + use bytes::Bytes; 7 + use futures::stream; 8 + use http::{HeaderMap, StatusCode}; 9 + use tokio::sync::mpsc; 10 + use url::Url; 11 + 12 + use crate::clock::{Clock, SleepFuture}; 13 + use crate::network::{ 14 + BodyStream, HttpRequest, HttpResponseFuture, HttpResponseHead, HttpResult, HttpTransport, 15 + NetworkError, WsConn, WsConnectFuture, WsMessage, WsMessageFuture, WsSendFuture, WsSink, 16 + WsStream, WsTransport, 17 + }; 18 + 19 + #[derive(Debug)] 20 + pub struct MemHttpResponse { 21 + pub latency: Duration, 22 + pub result: Result<MemHttpBody, NetworkError>, 23 + } 24 + 25 + #[derive(Debug)] 26 + pub struct MemHttpBody { 27 + pub status: StatusCode, 28 + pub headers: HeaderMap, 29 + pub body: Bytes, 30 + } 31 + 32 + impl MemHttpBody { 33 + pub fn ok_json(body: Bytes) -> Self { 34 + let mut headers = HeaderMap::new(); 35 + headers.insert(http::header::CONTENT_TYPE, "application/json".parse().unwrap()); 36 + Self { 37 + status: StatusCode::OK, 38 + headers, 39 + body, 40 + } 41 + } 42 + 43 + pub fn status_only(status: StatusCode) -> Self { 44 + Self { 45 + status, 46 + headers: HeaderMap::new(), 47 + body: Bytes::new(), 48 + } 49 + } 50 + } 51 + 52 + pub trait MemHttpResponder: Send + Sync + 'static { 53 + fn respond(&self, request: &HttpRequest) -> MemHttpResponse; 54 + } 55 + 56 + #[derive(Clone)] 57 + pub struct MemHttpTransport { 58 + responder: Arc<dyn MemHttpResponder>, 59 + clock: Arc<dyn Clock>, 60 + } 61 + 62 + impl MemHttpTransport { 63 + pub fn new(responder: Arc<dyn MemHttpResponder>, clock: Arc<dyn Clock>) -> Self { 64 + Self { responder, clock } 65 + } 66 + 67 + pub fn shared(responder: Arc<dyn MemHttpResponder>, clock: Arc<dyn Clock>) -> Arc<dyn HttpTransport> { 68 + Arc::new(Self::new(responder, clock)) 69 + } 70 + } 71 + 72 + impl HttpTransport for MemHttpTransport { 73 + fn execute(&self, request: HttpRequest) -> HttpResponseFuture { 74 + let response = self.responder.respond(&request); 75 + let sleep: SleepFuture = self.clock.sleep(response.latency); 76 + Box::pin(async move { 77 + sleep.await; 78 + into_response_head(response.result) 79 + }) 80 + } 81 + } 82 + 83 + fn into_response_head(result: Result<MemHttpBody, NetworkError>) -> HttpResult { 84 + let body = result?; 85 + let content_length = Some(body.body.len() as u64); 86 + let body_stream: BodyStream = Box::pin(stream::once(async move { Ok(body.body) })); 87 + Ok(HttpResponseHead { 88 + status: body.status, 89 + headers: body.headers, 90 + content_length, 91 + body: body_stream, 92 + }) 93 + } 94 + 95 + pub type MemWsServerFuture = Pin<Box<dyn Future<Output = ()> + Send + 'static>>; 96 + 97 + pub trait MemWsResponder: Send + Sync + 'static { 98 + fn spawn_server( 99 + &self, 100 + url: Url, 101 + recv: mpsc::UnboundedReceiver<WsMessage>, 102 + send: mpsc::Sender<WsMessage>, 103 + ) -> MemWsServerFuture; 104 + } 105 + 106 + pub const DEFAULT_MEM_WS_CAPACITY: usize = 4096; 107 + 108 + #[derive(Clone)] 109 + pub struct MemWsTransport { 110 + responder: Arc<dyn MemWsResponder>, 111 + capacity: usize, 112 + } 113 + 114 + impl MemWsTransport { 115 + pub fn new(responder: Arc<dyn MemWsResponder>) -> Self { 116 + Self::with_capacity(responder, DEFAULT_MEM_WS_CAPACITY) 117 + } 118 + 119 + pub fn with_capacity(responder: Arc<dyn MemWsResponder>, capacity: usize) -> Self { 120 + assert!(capacity > 0, "MemWsTransport capacity must be > 0"); 121 + Self { 122 + responder, 123 + capacity, 124 + } 125 + } 126 + 127 + pub fn shared(responder: Arc<dyn MemWsResponder>) -> Arc<dyn WsTransport> { 128 + Arc::new(Self::new(responder)) 129 + } 130 + 131 + pub fn shared_with_capacity( 132 + responder: Arc<dyn MemWsResponder>, 133 + capacity: usize, 134 + ) -> Arc<dyn WsTransport> { 135 + Arc::new(Self::with_capacity(responder, capacity)) 136 + } 137 + } 138 + 139 + impl WsTransport for MemWsTransport { 140 + fn connect(&self, url: Url) -> WsConnectFuture { 141 + let responder = self.responder.clone(); 142 + let capacity = self.capacity; 143 + Box::pin(async move { 144 + let (c2s_tx, c2s_rx) = mpsc::unbounded_channel(); 145 + let (s2c_tx, s2c_rx) = mpsc::channel(capacity); 146 + let server_future = responder.spawn_server(url, c2s_rx, s2c_tx); 147 + tokio::spawn(server_future); 148 + let sink: Box<dyn WsSink> = Box::new(MemWsSink { sender: c2s_tx }); 149 + let stream: Box<dyn WsStream> = Box::new(MemWsStream { receiver: s2c_rx }); 150 + Ok(WsConn { sink, stream }) 151 + }) 152 + } 153 + } 154 + 155 + struct MemWsSink { 156 + sender: mpsc::UnboundedSender<WsMessage>, 157 + } 158 + 159 + impl WsSink for MemWsSink { 160 + fn send<'a>(&'a mut self, message: WsMessage) -> WsSendFuture<'a> { 161 + let result = self 162 + .sender 163 + .send(message) 164 + .map_err(|_| NetworkError::Transport("server side closed channel".into())); 165 + Box::pin(async move { result }) 166 + } 167 + } 168 + 169 + struct MemWsStream { 170 + receiver: mpsc::Receiver<WsMessage>, 171 + } 172 + 173 + impl WsStream for MemWsStream { 174 + fn next<'a>(&'a mut self) -> WsMessageFuture<'a> { 175 + Box::pin(async move { self.receiver.recv().await.map(Ok) }) 176 + } 177 + } 178 + 179 + #[cfg(test)] 180 + mod tests { 181 + use super::*; 182 + use crate::SimClock; 183 + use crate::UnixMicros; 184 + use std::sync::Mutex; 185 + use std::sync::atomic::{AtomicUsize, Ordering}; 186 + use tokio::time::Instant; 187 + 188 + struct ScriptedHttp { 189 + responses: Mutex<Vec<MemHttpResponse>>, 190 + cursor: AtomicUsize, 191 + } 192 + 193 + impl ScriptedHttp { 194 + fn new(responses: Vec<MemHttpResponse>) -> Self { 195 + Self { 196 + responses: Mutex::new(responses), 197 + cursor: AtomicUsize::new(0), 198 + } 199 + } 200 + } 201 + 202 + impl MemHttpResponder for ScriptedHttp { 203 + fn respond(&self, _: &HttpRequest) -> MemHttpResponse { 204 + let i = self.cursor.fetch_add(1, Ordering::Relaxed); 205 + let mut guard = self.responses.lock().unwrap(); 206 + let next = std::mem::replace( 207 + &mut guard[i], 208 + MemHttpResponse { 209 + latency: Duration::ZERO, 210 + result: Err(NetworkError::Transport("script consumed".into())), 211 + }, 212 + ); 213 + next 214 + } 215 + } 216 + 217 + fn ok_response(body: &str, latency_ms: u64) -> MemHttpResponse { 218 + MemHttpResponse { 219 + latency: Duration::from_millis(latency_ms), 220 + result: Ok(MemHttpBody::ok_json(Bytes::from(body.to_owned()))), 221 + } 222 + } 223 + 224 + fn err_response(latency_ms: u64) -> MemHttpResponse { 225 + MemHttpResponse { 226 + latency: Duration::from_millis(latency_ms), 227 + result: Err(NetworkError::Transport("brownout".into())), 228 + } 229 + } 230 + 231 + #[tokio::test(start_paused = true)] 232 + async fn http_returns_scripted_body_after_injected_latency() { 233 + let clock: Arc<dyn Clock> = Arc::new(SimClock::at(UnixMicros::new(0))); 234 + let responder: Arc<dyn MemHttpResponder> = Arc::new(ScriptedHttp::new(vec![ 235 + ok_response("{\"hello\":1}", 10), 236 + ])); 237 + let transport = MemHttpTransport::new(responder, clock); 238 + let request = HttpRequest { 239 + url: Url::parse("http://oyster.cafe/xrpc/x").unwrap(), 240 + headers: HeaderMap::new(), 241 + }; 242 + 243 + let before = Instant::now(); 244 + let resp = transport.execute(request).await.unwrap(); 245 + let elapsed = Instant::now().saturating_duration_since(before); 246 + 247 + assert_eq!(resp.status, StatusCode::OK); 248 + assert_eq!(elapsed, Duration::from_millis(10)); 249 + let chunks = collect_body(resp.body).await; 250 + assert_eq!(chunks.as_ref(), b"{\"hello\":1}"); 251 + } 252 + 253 + #[tokio::test(start_paused = true)] 254 + async fn http_propagates_scripted_errors_with_latency() { 255 + let clock: Arc<dyn Clock> = Arc::new(SimClock::at(UnixMicros::new(0))); 256 + let responder: Arc<dyn MemHttpResponder> = 257 + Arc::new(ScriptedHttp::new(vec![err_response(50)])); 258 + let transport = MemHttpTransport::new(responder, clock); 259 + let request = HttpRequest { 260 + url: Url::parse("http://oyster.cafe/xrpc/x").unwrap(), 261 + headers: HeaderMap::new(), 262 + }; 263 + 264 + let before = Instant::now(); 265 + let resp = transport.execute(request).await; 266 + let elapsed = Instant::now().saturating_duration_since(before); 267 + 268 + assert_eq!(elapsed, Duration::from_millis(50)); 269 + assert!(matches!(resp, Err(NetworkError::Transport(_)))); 270 + } 271 + 272 + async fn collect_body(mut body: BodyStream) -> Bytes { 273 + use futures::StreamExt; 274 + let mut acc = Vec::new(); 275 + while let Some(chunk) = body.next().await { 276 + acc.extend_from_slice(&chunk.unwrap()); 277 + } 278 + Bytes::from(acc) 279 + } 280 + 281 + struct ScriptedWs { 282 + frames: Mutex<Vec<WsMessage>>, 283 + } 284 + 285 + impl MemWsResponder for ScriptedWs { 286 + fn spawn_server( 287 + &self, 288 + _: Url, 289 + mut _recv: mpsc::UnboundedReceiver<WsMessage>, 290 + send: mpsc::Sender<WsMessage>, 291 + ) -> MemWsServerFuture { 292 + let frames: Vec<WsMessage> = std::mem::take(&mut *self.frames.lock().unwrap()); 293 + Box::pin(async move { 294 + for frame in frames { 295 + if send.send(frame).await.is_err() { 296 + return; 297 + } 298 + } 299 + }) 300 + } 301 + } 302 + 303 + #[tokio::test(start_paused = true)] 304 + async fn ws_delivers_scripted_frames_to_client() { 305 + let responder: Arc<dyn MemWsResponder> = Arc::new(ScriptedWs { 306 + frames: Mutex::new(vec![ 307 + WsMessage::Text("frame-a".into()), 308 + WsMessage::Text("frame-b".into()), 309 + ]), 310 + }); 311 + let transport = MemWsTransport::new(responder); 312 + let mut conn = transport 313 + .connect(Url::parse("ws://oyster.cafe/").unwrap()) 314 + .await 315 + .unwrap(); 316 + let a = conn.stream.next().await.unwrap().unwrap(); 317 + let b = conn.stream.next().await.unwrap().unwrap(); 318 + assert!(matches!(a, WsMessage::Text(t) if t == "frame-a")); 319 + assert!(matches!(b, WsMessage::Text(t) if t == "frame-b")); 320 + } 321 + 322 + struct EchoWs; 323 + 324 + impl MemWsResponder for EchoWs { 325 + fn spawn_server( 326 + &self, 327 + _: Url, 328 + mut recv: mpsc::UnboundedReceiver<WsMessage>, 329 + send: mpsc::Sender<WsMessage>, 330 + ) -> MemWsServerFuture { 331 + Box::pin(async move { 332 + while let Some(msg) = recv.recv().await { 333 + if send.send(msg).await.is_err() { 334 + return; 335 + } 336 + } 337 + }) 338 + } 339 + } 340 + 341 + #[tokio::test(start_paused = true)] 342 + async fn ws_round_trip_via_server_echo() { 343 + let transport = MemWsTransport::new(Arc::new(EchoWs)); 344 + let mut conn = transport 345 + .connect(Url::parse("ws://oyster.cafe/").unwrap()) 346 + .await 347 + .unwrap(); 348 + conn.sink.send(WsMessage::Text("ping".into())).await.unwrap(); 349 + let echoed = conn.stream.next().await.unwrap().unwrap(); 350 + assert!(matches!(echoed, WsMessage::Text(t) if t == "ping")); 351 + } 352 + }