From 5ee287497511a081790065a5d00c48f19ded381f Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 23:29:13 -0700 Subject: [PATCH] Verify forwarded reads remain fenced after owner loss --- .../tests/protocol/client/typed.rs | 13 ++++- .../runtime/lifecycle/ownership/lease.rs | 48 +++++++++++++++++++ 2 files changed, 60 insertions(+), 1 deletion(-) diff --git a/crates/cellule-runtime/tests/protocol/client/typed.rs b/crates/cellule-runtime/tests/protocol/client/typed.rs index 5966484..e25ca19 100644 --- a/crates/cellule-runtime/tests/protocol/client/typed.rs +++ b/crates/cellule-runtime/tests/protocol/client/typed.rs @@ -669,8 +669,19 @@ async fn remote_runtime_route_skips_local_metadata_without_skipping_local_owner_ 0 ); assert_eq!(counted.counts().body_requests(), 3); - caller.shutdown().await.unwrap(); fixture.handle().drain().await.unwrap(); + + // A peer still holding the old handle must refuse after the owner drains, + // even though the forwarding node no longer reads catalog or control. + counted.reset(); + assert!(matches!( + remote + .query::(&fixture.target, None, ()) + .await, + Err(InvocationError::NotStarted(_)) + )); + assert_eq!(counted.counts().body_requests(), 0); + caller.shutdown().await.unwrap(); } #[tokio::test] async fn typed_client_rejects_conflicting_identity_receipt_and_module_before_execution() { diff --git a/crates/cellule-runtime/tests/runtime/lifecycle/ownership/lease.rs b/crates/cellule-runtime/tests/runtime/lifecycle/ownership/lease.rs index 6be56eb..dfd74e5 100644 --- a/crates/cellule-runtime/tests/runtime/lifecycle/ownership/lease.rs +++ b/crates/cellule-runtime/tests/runtime/lifecycle/ownership/lease.rs @@ -90,6 +90,54 @@ async fn node_lease_expiry_hides_an_inflight_committed_command() { Err(cellule_runtime::Error::Fenced) )); } + +#[tokio::test(flavor = "multi_thread")] +async fn node_lease_expiry_hides_an_inflight_query_result() { + let fixture = fixture_for(b"node-lease-query-output-gate"); + let session = SessionId::from_bytes([45; 16]); + let runtime = CellRuntime::new_with_replica_host_requiring_node_lease( + SqlWorkerPool::new(1, 1).unwrap(), + 2 * 1024 * 1024, + session, + ReplicaHost::default(), + ) + .unwrap(); + let lease = NodeLeaseGuard::new(0, 60_000).unwrap(); + runtime.install_node_lease(lease.clone()).unwrap(); + let handle = bootstrap_on(&runtime, &fixture, session).await; + let (entered_tx, entered_rx) = mpsc::channel(); + let (resume_tx, resume_rx) = mpsc::channel(); + + let query = tokio::spawn(async move { + handle + .query(1, 16, move |_connection| { + entered_tx.send(()).unwrap(); + resume_rx.recv().unwrap(); + Ok(b"stale-result".to_vec()) + }) + .await + }); + tokio::task::spawn_blocking(move || { + entered_rx + .recv_timeout(std::time::Duration::from_secs(5)) + .unwrap() + }) + .await + .unwrap(); + lease.fence(); + resume_tx.send(()).unwrap(); + + assert!(matches!( + query.await.unwrap(), + Err(cellule_runtime::Error::Fenced) + )); + let shutdown = runtime.shutdown().await; + assert!(matches!( + shutdown, + Ok(()) | Err(cellule_runtime::Error::Fenced) + )); + assert_eq!(runtime.stats().active_cells(), 0); +} #[tokio::test] async fn activation_rejects_control_owned_by_another_node_session() { let fixture = fixture();