diff --git a/src/ingest/firehose.rs b/src/ingest/firehose.rs index af7976d..ac6b198 100644 --- a/src/ingest/firehose.rs +++ b/src/ingest/firehose.rs @@ -1,5 +1,7 @@ use crate::filter::{FilterHandle, FilterMode}; -use crate::ingest::stream::{FirehoseError, FirehoseStream, SubscribeReposMessage, decode_frame}; +use crate::ingest::stream::{ + FirehoseError, FirehoseStream, InfoName, SubscribeReposMessage, decode_frame, +}; use crate::ingest::{BufferTx, IngestMessage}; use crate::pds_meta::HostStatus; use crate::state::AppState; @@ -71,7 +73,12 @@ fn classify_firehose_error(err: &FirehoseError) -> FirehoseFailure { FirehoseFailure::new("stream_closed", format!("close code {code}: {reason}")) } FirehoseError::TcpDropped => FirehoseFailure::new("tcp_dropped", "tcp layer dropped"), - FirehoseError::FutureCursor => FirehoseFailure::new("cursor", "future cursor"), + FirehoseError::FutureCursor { message } => FirehoseFailure::new( + "future_cursor", + message + .as_deref() + .map_or_else(|| "future cursor".to_owned(), |message| message.to_owned()), + ), } } @@ -268,6 +275,7 @@ impl FirehoseIngestor { let active_sleep_secs = if cfg!(debug_assertions) { 1 } else { 60 }; let mut active_sleep = std::pin::pin!(tokio::time::sleep(Duration::from_secs(active_sleep_secs))); + let mut reported_outdated_cursor = false; let res = loop { tokio::select! { @@ -281,6 +289,22 @@ impl FirehoseIngestor { Ok(msg) => { let (kind, seq) = message_stats(&msg); self.stats.record_decoded(kind, seq); + + if let SubscribeReposMessage::Info(info) = &msg { + if info.name == InfoName::OutdatedCursor + && !reported_outdated_cursor + { + reported_outdated_cursor = true; + warn!( + relay = %self.relay_host, + requested_cursor = ?start_cursor, + message = ?info.message, + "firehose cursor predates relay retention; consuming the oldest retained events" + ); + } + continue; + } + if self.is_pds { let tier = { let meta = self.state.pds_meta.load(); @@ -384,18 +408,57 @@ impl FirehoseIngestor { debug!(reason = %reason, "host gone away"); tokio::time::sleep(Duration::from_secs(1)).await; } - Err(FirehoseError::FutureCursor) => { - self.stats.record_stream_error("future_cursor"); - if self.is_pds - && let Err(e) = self.set_host_status(HostStatus::Idle) - { - error!(err = %e, "failed to update host status to idle"); - } - if let Err(e) = self.clear_stale_cursor() { - error!(err = %e, "failed to clear outdated cursor"); + Err(FirehoseError::FutureCursor { message }) => { + let failure = FirehoseFailure::new( + "future_cursor", + message + .as_deref() + .map_or_else(|| "future cursor".to_owned(), str::to_owned), + ); + self.stats.record_stream_error(failure.kind); + if start_cursor.is_some() { + if self.is_pds + && let Err(e) = self.set_host_status(HostStatus::Idle) + { + error!(err = %e, "failed to update host status to idle"); + } + if let Err(e) = self.clear_stale_cursor() { + error!(err = %e, "failed to clear future cursor"); + } + warn!( + relay = %self.relay_host, + requested_cursor = ?start_cursor, + message = ?message, + "relay rejected a future cursor; cleared stored cursor and retrying without it" + ); + tokio::time::sleep(Duration::from_secs(1)).await; + continue; } - warn!("outdated cursor, cleared stored cursor and retrying from live tail"); - tokio::time::sleep(Duration::from_secs(1)).await; + + let secs = match self.on_failure(&failure).await { + Some(secs) => secs, + None => { + error!( + failure_kind = failure.kind, + failure = %failure.detail, + failures = self.throttle.consecutive_failures(), + max_failures = self.max_failures, + "relay reported a future cursor without a requested cursor, giving up" + ); + break Ok(()); + } + }; + let timeout = rng.add_jitter(Duration::from_secs(secs).min(MAX_BACKOFF)); + let fmt = humantime::format_duration(timeout); + error!( + failure_kind = failure.kind, + failure = %failure.detail, + failures = self.throttle.consecutive_failures(), + max_failures = self.max_failures, + in = %fmt, + "relay reported a future cursor without a requested cursor, reconnecting later" + ); + tokio::time::sleep(timeout).await; } Err(FirehoseError::RelayError { error, message }) => { let message = message @@ -623,7 +686,45 @@ mod tests { use super::*; use crate::config::Config; use crate::control::Hydrant; + use crate::filter::FilterConfig; + use crate::ingest::stream::{Datetime, Identity, Info}; + use futures::{SinkExt as _, StreamExt as _}; + use serde::Serialize; use std::sync::atomic::AtomicI64; + use tokio_websockets::{Message as WsMessage, ServerBuilder}; + + #[derive(Serialize)] + struct TestFrameHeader<'a> { + op: i64, + #[serde(skip_serializing_if = "Option::is_none")] + t: Option<&'a str>, + } + + fn test_frame(op: i64, ty: Option<&str>, body: &T) -> Vec { + let mut bytes = serde_ipld_dagcbor::to_vec(&TestFrameHeader { op, t: ty }).unwrap(); + bytes.extend_from_slice(&serde_ipld_dagcbor::to_vec(body).unwrap()); + bytes + } + + fn test_identity(seq: i64) -> Identity { + Identity { + did: Did::new_static("did:plc:ewvi7nxzyoun6zhxrhs64oiz").unwrap(), + handle: None, + seq, + time: serde_json::from_str::(r#""2026-08-10T12:34:56Z""#).unwrap(), + } + } + + fn outdated_cursor_frame() -> Vec { + test_frame( + 1, + Some("#info"), + &Info { + message: Some("cursor is outside the retained window".into()), + name: InfoName::OutdatedCursor, + }, + ) + } #[tokio::test] async fn operational_transition_does_not_clear_an_operator_ban() -> Result<()> { @@ -731,4 +832,120 @@ mod tests { assert_eq!(select_start_cursor(Some(-5), true), None); assert_eq!(select_start_cursor(Some(-5), false), None); } + + #[tokio::test] + async fn in_window_resume_requests_cursor_and_receives_replay() -> Result<()> { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .into_diagnostic()?; + let addr = listener.local_addr().into_diagnostic()?; + let server = tokio::spawn(async move { + let (tcp, _) = listener.accept().await.unwrap(); + let (request, mut ws) = ServerBuilder::new().accept(tcp).await.unwrap(); + assert_eq!( + request.uri().path(), + "/xrpc/com.atproto.sync.subscribeRepos" + ); + assert_eq!(request.uri().query(), Some("cursor=42")); + ws.send(WsMessage::binary(test_frame( + 1, + Some("#identity"), + &test_identity(42), + ))) + .await + .unwrap(); + }); + + let relay = Url::parse(&format!("http://{addr}/ignored")).into_diagnostic()?; + let mut stream = FirehoseStream::connect(relay, Some(42)).await?; + let bytes = tokio::time::timeout(Duration::from_secs(2), stream.next()) + .await + .into_diagnostic()??; + let SubscribeReposMessage::Identity(identity) = decode_frame(&bytes)? else { + panic!("expected identity replay"); + }; + assert_eq!(identity.seq, 42); + server.await.into_diagnostic()?; + Ok(()) + } + + #[tokio::test] + async fn outdated_cursor_info_keeps_stream_open_and_forwards_replay() -> Result<()> { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .into_diagnostic()?; + let addr = listener.local_addr().into_diagnostic()?; + let server = tokio::spawn(async move { + let (tcp, _) = listener.accept().await.unwrap(); + let (request, mut ws) = ServerBuilder::new().accept(tcp).await.unwrap(); + assert_eq!(request.uri().query(), Some("cursor=42")); + ws.send(WsMessage::binary(outdated_cursor_frame())) + .await + .unwrap(); + ws.send(WsMessage::binary(test_frame( + 1, + Some("#identity"), + &test_identity(100), + ))) + .await + .unwrap(); + let _ = ws.next().await; + }); + + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = Config { + database_path: tmp.path().to_path_buf(), + ..Config::default() + }; + let state = Arc::new(AppState::new(&cfg)?); + state + .filter + .store(Arc::new(FilterConfig::new(FilterMode::Full))); + let relay = Url::parse(&format!("http://{addr}/")).into_diagnostic()?; + let _ = state + .firehose_cursors + .insert_async(relay.clone(), AtomicI64::new(42)) + .await; + let (buffer_tx, mut rxs) = BufferTx::channel(1); + let ingestor = FirehoseIngestor::new( + state.clone(), + buffer_tx, + relay.clone(), + false, + state.filter.clone(), + state.firehose_enabled.subscribe(), + false, + 2, + ) + .await; + let run = tokio::spawn(ingestor.run()); + + let replay = tokio::time::timeout(Duration::from_secs(2), rxs[0].recv()) + .await + .into_diagnostic()? + .ok_or_else(|| miette::miette!("firehose buffer closed before replay arrived"))?; + let IngestMessage::Firehose { + msg: SubscribeReposMessage::Identity(identity), + .. + } = replay + else { + panic!("expected retained identity event after OutdatedCursor info"); + }; + assert_eq!(identity.seq, 100); + assert_eq!( + state + .firehose_cursors + .peek_with(&relay, |_, cursor| cursor.load(Ordering::SeqCst)), + Some(42), + "receiving a frame must not bypass the downstream durable checkpoint" + ); + + run.abort(); + let _ = run.await; + tokio::time::timeout(Duration::from_secs(2), server) + .await + .into_diagnostic()? + .into_diagnostic()?; + Ok(()) + } } diff --git a/src/ingest/stream.rs b/src/ingest/stream.rs index 482ee58..f1be735 100644 --- a/src/ingest/stream.rs +++ b/src/ingest/stream.rs @@ -54,8 +54,8 @@ pub enum FirehoseError { StreamClosed { code: u16, reason: String }, #[error("tcp layer dropped")] TcpDropped, - #[error("future cursor")] - FutureCursor, + #[error("future cursor: {message:?}")] + FutureCursor { message: Option }, } impl From> for FirehoseError { diff --git a/src/ingest/stream/codec.rs b/src/ingest/stream/codec.rs index 9f41ca1..a7c6a36 100644 --- a/src/ingest/stream/codec.rs +++ b/src/ingest/stream/codec.rs @@ -1,5 +1,5 @@ use super::FirehoseError; -use super::types::{Info, InfoName, SubscribeReposMessage}; +use super::types::{Info, SubscribeReposMessage}; use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] @@ -21,10 +21,15 @@ pub fn decode_frame<'i>(bytes: &'i [u8]) -> Result, Fi match header.op { -1 => { let err = ErrorFrame::deserialize(&mut de)?; - return Err(FirehoseError::RelayError { - error: err.error, - message: err.message, - }); + return match err.error.as_str() { + "FutureCursor" => Err(FirehoseError::FutureCursor { + message: err.message, + }), + _ => Err(FirehoseError::RelayError { + error: err.error, + message: err.message, + }), + }; } 1 => {} op => return Err(FirehoseError::UnknownOp(op)), @@ -41,9 +46,6 @@ pub fn decode_frame<'i>(bytes: &'i [u8]) -> Result, Fi "#sync" => SubscribeReposMessage::Sync(Box::new(Deserialize::deserialize(&mut de)?)), "#info" => { let info: Info<'i> = Deserialize::deserialize(&mut de)?; - if info.name == InfoName::OutdatedCursor { - return Err(FirehoseError::FutureCursor); - } SubscribeReposMessage::Info(Box::new(info)) } other => return Err(FirehoseError::UnknownType(other.to_string())), @@ -97,7 +99,7 @@ pub fn encode_error_frame(error: &str, message: Option<&str>) -> miette::Result< #[cfg(test)] mod tests { use super::*; - use crate::ingest::stream::{Datetime, Identity, Sync}; + use crate::ingest::stream::{Datetime, Identity, InfoName, Sync}; use bytes::Bytes; use jacquard_common::{CowStr, types::string::Did}; @@ -211,20 +213,36 @@ mod tests { } #[test] - fn firehose_frame_maps_outdated_cursor_info() { + fn firehose_frame_preserves_outdated_cursor_info() { let bytes = frame( EventHeader { op: 1, t: Some("#info".to_owned()), }, &Info { - message: Some("cursor is in the future".into()), + message: Some("cursor is outside the retained window".into()), name: InfoName::OutdatedCursor, }, ); + let SubscribeReposMessage::Info(info) = decode_frame(&bytes).unwrap() else { + panic!("expected info"); + }; + assert_eq!(info.name, InfoName::OutdatedCursor); + } + + #[test] + fn firehose_frame_maps_future_cursor_error() { + let bytes = frame( + EventHeader { op: -1, t: None }, + &serde_json::json!({ + "error": "FutureCursor", + "message": "requested cursor is ahead of the relay", + }), + ); assert!(matches!( decode_frame(&bytes), - Err(FirehoseError::FutureCursor) + Err(FirehoseError::FutureCursor { message }) + if message.as_deref() == Some("requested cursor is ahead of the relay") )); } }