Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion crates/cellule-runtime/tests/protocol/client/typed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<CountComments>(&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() {
Expand Down
48 changes: 48 additions & 0 deletions crates/cellule-runtime/tests/runtime/lifecycle/ownership/lease.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Loading