From 1d0648b3b5ae0cca040c614b505b4d72bb769ab1 Mon Sep 17 00:00:00 2001 From: forhappy Date: Fri, 2 Oct 2026 19:42:57 -0700 Subject: [PATCH] fix(store): keep observed streams finished after EOF --- crates/cellule-store/docs/transport.md | 5 ++ crates/cellule-store/src/observation/mod.rs | 12 +++-- crates/cellule-store/src/observation/tests.rs | 52 +++++++++++++++++++ 3 files changed, 65 insertions(+), 4 deletions(-) diff --git a/crates/cellule-store/docs/transport.md b/crates/cellule-store/docs/transport.md index f3ed1453..2922196a 100644 --- a/crates/cellule-store/docs/transport.md +++ b/crates/cellule-store/docs/transport.md @@ -20,6 +20,11 @@ sequenceDiagram | Read admission | Reserve before backend work and release on completion/cancellation. | | Observation | Count requests and bytes without changing ownership semantics. | +Observed GET, listing, and deletion streams remain finished after EOF, including +empty bodies. Polling them again returns `None` and records no second terminal +observation. Completion, provider errors, and cancellation each release the +active operation exactly once. + `StorageError` keeps its source where available. `RetryClass` separates transient, throttled, state-dependent, and fatal failures. A caller controls its own deadline and whether an ambiguous write can be retried. diff --git a/crates/cellule-store/src/observation/mod.rs b/crates/cellule-store/src/observation/mod.rs index 636411df..b85f301a 100644 --- a/crates/cellule-store/src/observation/mod.rs +++ b/crates/cellule-store/src/observation/mod.rs @@ -506,7 +506,7 @@ fn observe_get_stream( stream: BoxStream<'static, object_store::Result>, observation: ActiveObservation, ) -> BoxStream<'static, object_store::Result> { - Box::pin(futures_util::stream::unfold( + futures_util::stream::unfold( (stream, Some(observation)), |(mut stream, mut observation)| async move { let active = observation.as_mut()?; @@ -525,14 +525,16 @@ fn observe_get_stream( } } }, - )) + ) + .fuse() + .boxed() } fn observe_stream( stream: BoxStream<'static, object_store::Result>, observation: ActiveObservation, ) -> BoxStream<'static, object_store::Result> { - Box::pin(futures_util::stream::unfold( + futures_util::stream::unfold( (stream, Some(observation)), |(mut stream, mut observation)| async move { let active = observation.as_mut()?; @@ -548,7 +550,9 @@ fn observe_stream( } } }, - )) + ) + .fuse() + .boxed() } fn classify_error(error: &object_store::Error) -> StorageOutcome { diff --git a/crates/cellule-store/src/observation/tests.rs b/crates/cellule-store/src/observation/tests.rs index ab8b0b17..972107e1 100644 --- a/crates/cellule-store/src/observation/tests.rs +++ b/crates/cellule-store/src/observation/tests.rs @@ -363,3 +363,55 @@ fn provider_failures_use_bounded_outcomes() { assert_eq!(classify_error(&forbidden), StorageOutcome::Auth); assert_eq!(classify_error(&cancelled), StorageOutcome::Cancelled); } + +#[tokio::test] +async fn empty_observed_body_can_be_collected_through_get_result() { + let observer = Arc::new(RecordingObserver::default()); + let erased: Arc = observer.clone(); + let provider = InMemory::new(); + let path = Path::from("repository/empty"); + provider + .put(&path, Bytes::new().into()) + .await + .expect("put empty"); + let mut result = provider.get(&path).await.expect("get metadata"); + // S3 can expose a zero-chunk body. GetResult::bytes calls collect_bytes, + // which polls twice even when the first poll reaches EOF. + result.payload = GetResultPayload::Stream(observe_get_stream( + futures_util::stream::empty().boxed(), + ActiveObservation::new(StorageOperation::Get, &erased), + )); + assert!(result.bytes().await.expect("collect empty body").is_empty()); + assert_eq!(observer.active(StorageOperation::Get), 0); + let observations = observer.observations(); + assert_eq!(observations.len(), 1); + assert_eq!(observations[0].outcome, StorageOutcome::Success); + assert_eq!(observations[0].bytes_read, 0); +} + +#[tokio::test] +async fn observed_streams_remain_finished_after_repeated_eof_polls() { + let observer = Arc::new(RecordingObserver::default()); + let erased: Arc = observer.clone(); + let mut body = observe_get_stream( + futures_util::stream::empty().boxed(), + ActiveObservation::new(StorageOperation::Get, &erased), + ); + let mut listing = observe_stream::( + futures_util::stream::empty().boxed(), + ActiveObservation::new(StorageOperation::List, &erased), + ); + for _ in 0..3 { + assert!(body.next().await.is_none()); + assert!(listing.next().await.is_none()); + } + assert_eq!(observer.active(StorageOperation::Get), 0); + assert_eq!(observer.active(StorageOperation::List), 0); + let observations = observer.observations(); + assert_eq!(observations.len(), 2); + assert!( + observations + .iter() + .all(|entry| entry.outcome == StorageOutcome::Success) + ); +}