From aafcd890e0c4e4953aa6d10286fc853db6bfb83e Mon Sep 17 00:00:00 2001 From: Nikolaos Dymitriadis Date: Mon, 17 Aug 2026 23:34:38 +0200 Subject: [PATCH] fix(grpc): don't panic on an unservable FollowTip/WatchTx intersect ChainStream::start ran ChainCrawler::start inside an async_stream::stream! block and unwrapped both its Result and its Option. A client resuming FollowTip or WatchTx from a point this node no longer holds (e.g. after a snapshot restore) hit the None case on every subscription, panicking a tokio worker while the client's stream just hung open with no status. Move the fallible crawler start ahead of the stream! block and have ChainStream::start return Result, DomainError>: None means the intersect wasn't found, Err means the crawler itself failed to start. Callers in the v1alpha/v1beta sync and watch services map None to Status::not_found (naming the requested points) and Err to Status::internal. Co-Authored-By: Claude Opus 5 --- src/serve/grpc/stream.rs | 45 ++++++++++++++++++++++++--------- src/serve/grpc/v1alpha/sync.rs | 8 +++++- src/serve/grpc/v1alpha/watch.rs | 8 +++++- src/serve/grpc/v1beta/sync.rs | 8 +++++- src/serve/grpc/v1beta/watch.rs | 8 +++++- 5 files changed, 61 insertions(+), 16 deletions(-) diff --git a/src/serve/grpc/stream.rs b/src/serve/grpc/stream.rs index f953acd58..fd6de8da8 100644 --- a/src/serve/grpc/stream.rs +++ b/src/serve/grpc/stream.rs @@ -9,17 +9,13 @@ impl ChainStream { domain: D, intersect: Vec, cancel: C, - ) -> impl Stream + 'static { - async_stream::stream! { - let result = ChainCrawler::::start( - &domain, - &intersect, - ); - - let start = result.expect("issue starting crawler"); - - let (mut crawler, intersected) = start.expect("crawler can't find start point"); + ) -> Result + 'static>, DomainError> { + let Some((mut crawler, intersected)) = ChainCrawler::::start(&domain, &intersect)? + else { + return Ok(None); + }; + Ok(Some(async_stream::stream! { yield TipEvent::Mark(intersected.clone()); while let Some((point, block)) = crawler.next_block() { @@ -36,7 +32,7 @@ impl ChainStream { } } } - } + })) } } @@ -81,7 +77,9 @@ mod tests { domain, vec![chain_point.clone()], CancelTokenImpl(CancellationToken::new()), - ); + ) + .unwrap() + .expect("intersect point should be found"); pin_mut!(s); @@ -105,4 +103,27 @@ mod tests { background.abort(); } + + #[tokio::test] + async fn test_stream_unknown_intersect() { + let domain = ToyDomain::new(None, None); + + for i in 0..=10 { + let (_, block) = make_conway_block(i * 10); + + use dolos_core::SyncExt; + domain.roll_forward(block).unwrap(); + } + + // this point was never rolled forward, so the domain can't intersect it. + let unknown_point = make_conway_block(9999).0; + + let result = ChainStream::start::( + domain, + vec![unknown_point], + CancelTokenImpl(CancellationToken::new()), + ); + + assert!(result.unwrap().is_none()); + } } diff --git a/src/serve/grpc/v1alpha/sync.rs b/src/serve/grpc/v1alpha/sync.rs index ab9d77f47..8dad14a67 100644 --- a/src/serve/grpc/v1alpha/sync.rs +++ b/src/serve/grpc/v1alpha/sync.rs @@ -286,7 +286,13 @@ where self.domain.clone(), intersect.clone(), self.cancel.clone(), - ); + ) + .map_err(|e| Status::internal(format!("failed to start chain stream: {e}")))? + .ok_or_else(|| { + Status::not_found(format!( + "none of the requested points intersect with local history: {intersect:?}" + )) + })?; let mapper = self.mapper.clone(); diff --git a/src/serve/grpc/v1alpha/watch.rs b/src/serve/grpc/v1alpha/watch.rs index 7808f36b5..f8af82e0d 100644 --- a/src/serve/grpc/v1alpha/watch.rs +++ b/src/serve/grpc/v1alpha/watch.rs @@ -486,7 +486,13 @@ where .collect::>(); let stream = - ChainStream::start::(self.domain.clone(), intersect, self.cancel.clone()); + ChainStream::start::(self.domain.clone(), intersect.clone(), self.cancel.clone()) + .map_err(|e| Status::internal(format!("failed to start chain stream: {e}")))? + .ok_or_else(|| { + Status::not_found(format!( + "none of the requested points intersect with local history: {intersect:?}" + )) + })?; let mapper = self.mapper.clone(); diff --git a/src/serve/grpc/v1beta/sync.rs b/src/serve/grpc/v1beta/sync.rs index 00a779d3e..8d58daf74 100644 --- a/src/serve/grpc/v1beta/sync.rs +++ b/src/serve/grpc/v1beta/sync.rs @@ -286,7 +286,13 @@ where self.domain.clone(), intersect.clone(), self.cancel.clone(), - ); + ) + .map_err(|e| Status::internal(format!("failed to start chain stream: {e}")))? + .ok_or_else(|| { + Status::not_found(format!( + "none of the requested points intersect with local history: {intersect:?}" + )) + })?; let mapper = self.mapper.clone(); diff --git a/src/serve/grpc/v1beta/watch.rs b/src/serve/grpc/v1beta/watch.rs index fce999dab..92661f207 100644 --- a/src/serve/grpc/v1beta/watch.rs +++ b/src/serve/grpc/v1beta/watch.rs @@ -486,7 +486,13 @@ where .collect::>(); let stream = - ChainStream::start::(self.domain.clone(), intersect, self.cancel.clone()); + ChainStream::start::(self.domain.clone(), intersect.clone(), self.cancel.clone()) + .map_err(|e| Status::internal(format!("failed to start chain stream: {e}")))? + .ok_or_else(|| { + Status::not_found(format!( + "none of the requested points intersect with local history: {intersect:?}" + )) + })?; let mapper = self.mapper.clone();