diff --git a/Cargo.lock b/Cargo.lock index 752796b37..1346c9796 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -831,6 +831,7 @@ dependencies = [ "chrono", "criterion", "futures", + "http", "jacquard-common", "scc", "serde", diff --git a/bobbin/crates/bobbin-sim/src/runtime.rs b/bobbin/crates/bobbin-sim/src/runtime.rs index 3fd599a1b..5c233a7a5 100644 --- a/bobbin/crates/bobbin-sim/src/runtime.rs +++ b/bobbin/crates/bobbin-sim/src/runtime.rs @@ -121,6 +121,7 @@ impl Sim { entropy: entropy.clone(), stream_health: Arc::new(StreamHealth::new()), ws: mem_ws, + hydrant: None, cancel: cancel.clone(), disconnects: Some(disconnects.clone()), warming_shadow: Some(warming_shadow.clone()), diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index 5381994d0..d9f717571 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -192,12 +192,13 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(lru_cap))); let slingshot = SlingshotClient::with_default_http(cfg.slingshot.url.clone())?; let actors = Arc::new(ActorIndex::new()); + let hydrant = HydrantClient::new( + cfg.hydrant.url.clone(), + ReqwestHttp::shared(default_http_client()?), + )?; let identity = Arc::new( IdentityResolver::with_slingshot(slingshot.clone(), clock.clone(), hasher.clone()) - .with_hydrant(HydrantClient::new( - cfg.hydrant.url.clone(), - ReqwestHttp::shared(default_http_client()?), - )?) + .with_hydrant(hydrant.clone()) .with_sink(actors.clone()), ); let mut resolver_opts = ResolverOptions::default(); @@ -375,6 +376,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { entropy, stream_health: stream_health.clone(), ws: ws.clone(), + hydrant: Some(hydrant), cancel: cancel.clone(), disconnects: None, warming_shadow: None, diff --git a/bobbin/crates/ingest/Cargo.toml b/bobbin/crates/ingest/Cargo.toml index 3506424ef..18d541103 100644 --- a/bobbin/crates/ingest/Cargo.toml +++ b/bobbin/crates/ingest/Cargo.toml @@ -29,6 +29,7 @@ tracing = { workspace = true } url = { workspace = true } [dev-dependencies] +http = { workspace = true } tracing-subscriber = { workspace = true } tokio-util = { workspace = true } tokio = { workspace = true, features = ["test-util"] } diff --git a/bobbin/crates/ingest/examples/smoke.rs b/bobbin/crates/ingest/examples/smoke.rs index 3c705fc5d..df1e3fd63 100644 --- a/bobbin/crates/ingest/examples/smoke.rs +++ b/bobbin/crates/ingest/examples/smoke.rs @@ -47,6 +47,7 @@ async fn main() { ws: TungsteniteWs::shared( WsTls::from_native_roots().expect("a system trust store for wss certificates"), ), + hydrant: None, cancel: cancel.clone(), disconnects: None, warming_shadow: None, diff --git a/bobbin/crates/ingest/src/lib.rs b/bobbin/crates/ingest/src/lib.rs index 189a67d7e..372d628ae 100644 --- a/bobbin/crates/ingest/src/lib.rs +++ b/bobbin/crates/ingest/src/lib.rs @@ -12,7 +12,7 @@ use bobbin_edge_index::{ use bobbin_knot_ingest::{CapabilityGate, KnotRegistry}; use bobbin_record_lru::RecordStore; use bobbin_resolver::{ - IdentityResolver, NormalizeRepoRefs, RepoClaim, decode_canon_or_upgrade_bytes, + HydrantClient, IdentityResolver, NormalizeRepoRefs, RepoClaim, decode_canon_or_upgrade_bytes, synthesize_created_at, }; use bobbin_runtime::{ @@ -57,6 +57,7 @@ const RECONNECT_MAX_DELAY: Duration = Duration::from_secs(30); const PING_INTERVAL: Duration = Duration::from_secs(20); const PONG_TIMEOUT: Duration = Duration::from_secs(15); const CONNECT_TIMEOUT: Duration = Duration::from_secs(15); +const HEAD_TIMEOUT: Duration = Duration::from_secs(5); const READY_SKEW: Duration = Duration::from_secs(60); const FRAME_CHANNEL_DEPTH: usize = 256; const CONTROL_CHANNEL_DEPTH: usize = 16; @@ -233,6 +234,7 @@ pub struct IngestRuntime { pub entropy: Arc, pub stream_health: Arc, pub ws: Arc, + pub hydrant: Option, pub cancel: CancellationToken, pub disconnects: Option>, pub warming_shadow: Option>, @@ -257,6 +259,7 @@ impl Clone for IngestRuntime { entropy: self.entropy.clone(), stream_health: self.stream_health.clone(), ws: self.ws.clone(), + hydrant: self.hydrant.clone(), cancel: self.cancel.clone(), disconnects: self.disconnects.clone(), warming_shadow: self.warming_shadow.clone(), @@ -329,13 +332,18 @@ async fn run_inner( frames: &FrameCounter, ) -> Result<(), IngestError> { let mut backoff = RECONNECT_INITIAL_DELAY; + let mut replay_head = None; loop { let cursor = next_connect_cursor(runtime.coverage.snapshot(), config.start_cursor); let opened = runtime.clock.now_instant(); runtime .stream_health .connecting(runtime.clock.now_unix_micros()); - let SessionEnd { outcome, error } = run_session(&config, cursor, runtime, frames).await; + if replay_head.is_none() && !runtime.coverage.snapshot().is_ready() { + replay_head = fetch_replay_head(runtime).await; + } + let SessionEnd { outcome, error } = + run_session(&config, cursor, replay_head, runtime, frames).await; let open_for = runtime.clock.now_instant().duration_since(opened); if runtime.cancel.is_cancelled() { info!( @@ -516,6 +524,59 @@ fn promote_if_caught_up( ); } +/// the newest event hydrant holds before we subscribe, so a warming replay can +/// promote the moment it gets there. `None` leaves promotion to the idle check +async fn fetch_replay_head( + runtime: &IngestRuntime, +) -> Option { + let client = runtime.hydrant.as_ref()?; + let head = tokio::select! { + biased; + _ = runtime.cancel.cancelled() => return None, + res = client.stream_head() => res, + _ = runtime.clock.sleep(HEAD_TIMEOUT) => { + warn!(timeout = ?HEAD_TIMEOUT, "hydrant didn't report its stream head, coverage waits for an idle stream"); + return None; + } + }; + match head { + Ok(Some(id)) => { + info!( + target: "bobbin_ingest::coverage", + head = id, + "coverage promotes once the replay reaches hydrant's stream head", + ); + Some(HydrantCursor::new(id)) + } + Ok(None) => { + debug!(target: "bobbin_ingest::coverage", "hydrant holds no events yet, coverage waits for an idle stream"); + None + } + Err(err) => { + warn!( + ?err, + "couldn't read hydrant's stream head, coverage waits for an idle stream" + ); + None + } + } +} + +fn promote_at_head(runtime: &IngestRuntime, head: HydrantCursor) { + let snap = runtime.coverage.snapshot(); + if snap.is_ready() || snap.last_cursor() < head { + return; + } + runtime.coverage.update(|c| c.force_ready()); + info!( + target: "bobbin_ingest::coverage", + events_processed = snap.events_processed(), + last_cursor = snap.last_cursor().raw(), + head = head.raw(), + "the replay reached hydrant's stream head, promoting coverage to ready", + ); +} + fn spawn_metrics_dumper( runtime: &IngestRuntime, ) -> tokio::task::JoinHandle<()> { @@ -563,6 +624,7 @@ fn jittered(base: Duration, entropy: &dyn Entropy) -> Duration { async fn run_session( config: &IngestConfig, cursor: HydrantCursor, + replay_head: Option, runtime: &IngestRuntime, frames: &FrameCounter, ) -> SessionEnd { @@ -619,7 +681,9 @@ async fn run_session( .then(move |staged| claim_stage(staged, claim_rt.clone())) .map(move |staged| resolve_stage(staged, resolve_rt.clone())) .buffered(parallelism) - .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)); + .for_each(move |staged| { + commit_stage(staged, commit_rt.clone(), parallelism, replay_head) + }); tokio::select! { biased; @@ -1087,6 +1151,7 @@ async fn commit_stage( staged: Resolved, rt: IngestRuntime, parallelism: usize, + replay_head: Option, ) { let nsid = pending_nsid(&staged.pending.op).cloned(); let edge_count = pending_edge_count(&staged.pending.op); @@ -1105,6 +1170,9 @@ async fn commit_stage( if let Some((watch, settled)) = rt.settlements.as_deref().zip(settled) { watch.settle(settled); } + if let Some(head) = replay_head { + promote_at_head(&rt, head); + } tracing::trace!( target: "bobbin_ingest::stage", cursor, @@ -3189,6 +3257,7 @@ mod tests { entropy: Arc::new(OsEntropy), stream_health: Arc::new(StreamHealth::new()), ws: ScriptedWs::undialed(), + hydrant: None, cancel, disconnects: None, warming_shadow: None, @@ -3292,6 +3361,7 @@ mod tests { run_session( &cfg, HydrantCursor::new(0), + None, &runtime, &FrameCounter::default(), ) @@ -3342,6 +3412,148 @@ mod tests { busy.session.await.expect("session panicked"); } + struct HeadResponder { + status: http::StatusCode, + head: serde_json::Value, + paths: std::sync::Mutex>, + } + + impl bobbin_runtime::MemHttpResponder for HeadResponder { + fn respond( + &self, + request: &bobbin_runtime::HttpRequest, + ) -> bobbin_runtime::MemHttpResponse { + self.paths + .lock() + .unwrap() + .push(request.url.path().to_owned()); + let body = if self.status == http::StatusCode::OK { + bobbin_runtime::MemHttpBody::ok_json(json!({ "id": self.head }).to_string().into()) + } else { + bobbin_runtime::MemHttpBody::status_only(self.status) + }; + bobbin_runtime::MemHttpResponse { + latency: Duration::ZERO, + result: Ok(body), + } + } + } + + struct HeadReplay { + coverage: Arc, + responder: Arc, + cancel: CancellationToken, + task: tokio::task::JoinHandle>, + } + + /// a cold ingest against a hydrant that answers `/stream/head` with `head` + /// (or `status`) and replays frames `1..=frames` before going quiet + fn head_replay(status: http::StatusCode, head: serde_json::Value, frames: u64) -> HeadReplay { + let (ws, ws_rx) = tokio::sync::mpsc::channel(16); + (1..=frames).for_each(|id| ws.try_send(other_frame(id)).expect("room for the replay")); + let cancel = CancellationToken::new(); + let mut runtime = fresh_runtime(cancel.clone()); + runtime.ws = ScriptedWs::dialing(WsConn { + sink: Box::new(PongingSink { + ws, + backlog: Backlog::Drained, + }), + stream: Box::new(ChannelWsStream { rx: ws_rx }), + }); + let responder = Arc::new(HeadResponder { + status, + head, + paths: std::sync::Mutex::new(Vec::new()), + }); + let http = + bobbin_runtime::MemHttpTransport::shared(responder.clone(), runtime.clock.clone()); + runtime.hydrant = Some( + HydrantClient::new(Url::parse("ws://127.0.0.1:1").expect("hydrant url"), http) + .expect("hydrant client"), + ); + let coverage = runtime.coverage.clone(); + let cfg = IngestConfig::new(Url::parse("ws://127.0.0.1:1").expect("hydrant url")); + let task = tokio::spawn(async move { run(cfg, runtime).await }); + HeadReplay { + coverage, + responder, + cancel, + task, + } + } + + async fn stop(replay: HeadReplay) { + replay.cancel.cancel(); + let outcome = replay.task.await.expect("ingest panicked"); + assert!(matches!(outcome, Ok(())), "got {outcome:?}"); + } + + #[tokio::test(start_paused = true)] + async fn a_replay_reaching_the_hydrant_head_promotes_without_waiting_for_idle() { + let replay = head_replay(http::StatusCode::OK, json!(3), 3); + let started = tokio::time::Instant::now(); + let mut coverage = replay.coverage.subscribe(); + coverage + .wait_for(|c| c.is_ready()) + .await + .expect("coverage watch outlives the ingest"); + assert!( + started.elapsed() < PING_INTERVAL, + "promotion came from the head, not the idle keepalive check, took {:?}", + started.elapsed(), + ); + assert_eq!( + replay.coverage.snapshot().last_cursor(), + HydrantCursor::new(3) + ); + assert_eq!( + *replay.responder.paths.lock().unwrap(), + vec!["/stream/head"] + ); + stop(replay).await; + } + + #[tokio::test(start_paused = true)] + async fn a_replay_short_of_the_hydrant_head_stays_warming() { + let replay = head_replay(http::StatusCode::OK, json!(4), 3); + assert!( + tokio::time::timeout( + PING_INTERVAL * 4, + replay.coverage.subscribe().wait_for(|c| c.is_ready()) + ) + .await + .is_err(), + "frame 4 never arrived, so the replay hasn't delivered what hydrant held", + ); + assert_eq!(replay.coverage.snapshot().events_processed(), 3); + stop(replay).await; + } + + #[tokio::test(start_paused = true)] + async fn a_hydrant_without_a_head_leaves_promotion_to_the_idle_check() { + for (status, head) in [ + (http::StatusCode::NOT_FOUND, json!(null)), + (http::StatusCode::OK, json!(null)), + ] { + let replay = head_replay(status, head, 3); + assert!( + tokio::time::timeout( + PING_INTERVAL * 4, + replay.coverage.subscribe().wait_for(|c| c.is_ready()) + ) + .await + .is_err(), + "no head and fewer than the idle minimum of events must stay warming ({status})", + ); + assert_eq!( + replay.coverage.snapshot().events_processed(), + 3, + "a missing head must not keep the stream from replaying ({status})", + ); + stop(replay).await; + } + } + #[tokio::test] async fn a_quiet_socket_leaves_a_cold_ingest_warming() { let runtime = fresh_runtime(CancellationToken::new()); @@ -4229,6 +4441,7 @@ mod tests { run_session( &cfg, HydrantCursor::new(0), + None, &runtime, &FrameCounter::default(), ), @@ -4357,6 +4570,7 @@ mod tests { entropy: Arc::new(OsEntropy), stream_health: Arc::new(StreamHealth::new()), ws: ScriptedWs::undialed(), + hydrant: None, cancel: CancellationToken::new(), disconnects: None, warming_shadow: None, @@ -4377,7 +4591,7 @@ mod tests { .then(move |frame| prep_stage(frame, prep_rt.clone())) .map(move |staged| resolve_stage(staged, resolve_rt.clone())) .buffered(parallelism) - .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism)) + .for_each(move |staged| commit_stage(staged, commit_rt.clone(), parallelism, None)) .await; }); @@ -4439,7 +4653,7 @@ mod tests { let staged = prep_stage(frame, rt.clone()).await; let staged = claim_stage(staged, rt.clone()).await; let staged = resolve_stage(staged, rt.clone()).await; - commit_stage(staged, rt.clone(), 1).await; + commit_stage(staged, rt.clone(), 1, None).await; } fn published(watch: &Arc, uri: &str) -> Option { diff --git a/bobbin/crates/resolver/src/hydrant.rs b/bobbin/crates/resolver/src/hydrant.rs index c13773d22..ac301b2cf 100644 --- a/bobbin/crates/resolver/src/hydrant.rs +++ b/bobbin/crates/resolver/src/hydrant.rs @@ -90,6 +90,11 @@ impl<'de> Deserialize<'de> for RepoStatus { } } +#[derive(Debug, Deserialize)] +struct StreamHead { + id: Option, +} + #[derive(Clone)] pub struct HydrantClient { http: Arc, @@ -119,15 +124,40 @@ impl HydrantClient { } // segment-wise so the colons in a did escape instead of parsing as a scheme - fn repo_url(&self, did: &Did) -> Result { + fn endpoint<'a>( + &self, + segments: impl IntoIterator, + ) -> Result { let mut url = self.base.clone(); url.path_segments_mut() .map_err(|()| HydrantError::BadScheme(self.base.scheme().to_owned()))? .pop_if_empty() - .extend(["repos", did.as_str()]); + .extend(segments); Ok(url) } + fn repo_url(&self, did: &Did) -> Result { + self.endpoint(["repos", did.as_str()]) + } + + /// the id of the newest event hydrant has stored, `None` while it stores none. + /// a `cursor=0` replay has delivered everything hydrant held at the time of + /// asking once it passes this id + pub async fn stream_head(&self) -> Result, HydrantError> { + let resp = self + .http + .execute(HttpRequest { + url: self.endpoint(["stream", "head"])?, + headers: HeaderMap::new(), + }) + .await?; + if resp.status != StatusCode::OK { + return Err(HydrantError::Upstream(resp.status)); + } + let bytes = read_bounded(resp).await?; + Ok(serde_json::from_slice::(&bytes)?.id) + } + /// `None` when hydrant does not track the repo or has no handle for it pub async fn repo_identity( &self, @@ -456,4 +486,51 @@ mod tests { "https://h/hydrant/repos/did:plc:dawn" ); } + + #[tokio::test] + async fn stream_head_reads_the_newest_event_id() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/hydrant/stream/head")) + .respond_with( + ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": 289455})), + ) + .mount(&server) + .await; + + let client = HydrantClient::new( + Url::parse(&format!("{}/hydrant/", server.uri())).unwrap(), + ReqwestHttp::shared(default_http_client().unwrap()), + ) + .unwrap(); + assert_eq!(client.stream_head().await.unwrap(), Some(289455)); + } + + #[tokio::test] + async fn an_empty_stream_has_no_head() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/stream/head")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": null}))) + .mount(&server) + .await; + + assert_eq!(client(&server).await.stream_head().await.unwrap(), None); + } + + #[tokio::test] + async fn a_hydrant_without_stream_head_is_an_upstream_error() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/stream/head")) + .respond_with(ResponseTemplate::new(404)) + .mount(&server) + .await; + + let err = client(&server).await.stream_head().await.unwrap_err(); + assert!( + matches!(err, HydrantError::Upstream(StatusCode::NOT_FOUND)), + "got {err:?}" + ); + } } diff --git a/bobbin/crates/xrpc/tests/await_record_e2e.rs b/bobbin/crates/xrpc/tests/await_record_e2e.rs index e732d34c0..a5d596046 100644 --- a/bobbin/crates/xrpc/tests/await_record_e2e.rs +++ b/bobbin/crates/xrpc/tests/await_record_e2e.rs @@ -103,6 +103,7 @@ async fn a_record_off_the_stream_answers_an_await() { ws: MemWsTransport::shared(Arc::new(Hydrant { frames: vec![star_frame()], })), + hydrant: None, cancel: cancel.clone(), disconnects: None, warming_shadow: None,