diff --git a/crates/cellule-app/Cargo.toml b/crates/cellule-app/Cargo.toml index a0016e10..4dd2feb1 100644 --- a/crates/cellule-app/Cargo.toml +++ b/crates/cellule-app/Cargo.toml @@ -26,7 +26,7 @@ futures-util.workspace = true object_store.workspace = true prost.workspace = true tempfile.workspace = true -tokio = { workspace = true, features = ["io-util", "macros", "net", "rt-multi-thread"] } +tokio = { workspace = true, features = ["io-util", "macros", "net", "rt-multi-thread", "test-util"] } tokio-util.workspace = true tracing-subscriber = "0.3" diff --git a/crates/cellule-app/PERFORMANCE.md b/crates/cellule-app/PERFORMANCE.md index ccd736a6..0992cbd3 100644 --- a/crates/cellule-app/PERFORMANCE.md +++ b/crates/cellule-app/PERFORMANCE.md @@ -32,6 +32,13 @@ Set `CELLULE_TEST_ENDPOINT`, `CELLULE_TEST_BUCKET`, and a unique correctness from provider, fleet, and production evidence. Local and Compose runs do not establish production capacity. +The reference fleet's TCP transport bounds connection setup to 250 ms within +its overall request deadline. A setup timeout is a known pre-dispatch failure; +a deadline after connection remains an unknown outcome. This lets reader +selection move past a killed node while preserving the full reply budget for +connected peers. Reader-loss evidence still requires every load lane to make +progress before, during, and after replacement, with receipt minima enforced. + ## Dated evidence | Question | Reports | diff --git a/crates/cellule-app/tests/fleet.rs b/crates/cellule-app/tests/fleet.rs index 2c18bf28..4e45614c 100644 --- a/crates/cellule-app/tests/fleet.rs +++ b/crates/cellule-app/tests/fleet.rs @@ -181,10 +181,26 @@ pub(super) async fn send_tcp( request: Vec, remaining_ms: u32, ) -> Result> { - tokio::time::timeout(Duration::from_millis(u64::from(remaining_ms)), async { + send_tcp_connecting(TcpStream::connect(address), request, remaining_ms).await +} + +async fn send_tcp_connecting( + connection: impl Future>, + request: Vec, + remaining_ms: u32, +) -> Result> { + let remaining = Duration::from_millis(u64::from(remaining_ms)); + tokio::time::timeout(remaining, async { + // A killed private-network peer can blackhole SYNs. No request has + // been dispatched yet, so bound setup separately to preserve failover + // time without shortening the reply budget of a connected peer. let mut socket = - TcpStream::connect(address) + tokio::time::timeout(remaining.min(Duration::from_millis(250)), connection) .await + .map_err(|source| Error::PeerTransport { + context: "fleet peer connect deadline", + source: Box::new(source), + })? .map_err(|source| Error::PeerTransport { context: "fleet peer connect", source: Box::new(source), @@ -410,3 +426,57 @@ async fn serve_peer( socket.write_all(&reply).await.map_err(peer_io)?; Ok(()) } + +#[tokio::test(start_paused = true)] +async fn stalled_connect_preserves_failover_budget() { + for (remaining_ms, connect_ms) in [(5_000, 250), (100, 100)] { + let started = tokio::time::Instant::now(); + let error = send_tcp_connecting(std::future::pending(), vec![1], remaining_ms) + .await + .unwrap_err(); + assert!( + matches!( + error, + Error::PeerTransport { + context: "fleet peer connect deadline", + .. + } + ), + "{error:?}" + ); + assert_eq!(started.elapsed(), Duration::from_millis(connect_ms)); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn dispatched_request_keeps_full_reply_budget() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let (received, dispatched) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.unwrap(); + let mut request = [0; 5]; + socket.read_exact(&mut request).await.unwrap(); + assert_eq!(request, [0, 0, 0, 1, 42]); + received.send(()).unwrap(); + std::future::pending::<()>().await; + drop(socket); + }); + let started = Instant::now(); + let result = send_tcp(address, vec![42], 1_000).await; + let elapsed = started.elapsed(); + server.abort(); + let _ = server.await; + dispatched.await.unwrap(); + assert!( + matches!( + result, + Err(Error::PeerTransportUnknown { + context: "fleet peer deadline", + .. + }) + ), + "{result:?}" + ); + assert!(elapsed >= Duration::from_millis(1_000), "{elapsed:?}"); +} diff --git a/crates/cellule-app/tests/host.rs b/crates/cellule-app/tests/host.rs index e07c7593..b7ac2a4d 100644 --- a/crates/cellule-app/tests/host.rs +++ b/crates/cellule-app/tests/host.rs @@ -94,6 +94,7 @@ async fn due_workflow_activity_survives_a_clock_rollback() { .unwrap(); let input = encoder.finish(); let registry = Arc::clone(&fixture.registry); + let owner_fence = handle.owner_fence(); // Publish at a later clock sample, then let the ordinary application // supervisor use the earlier wall clock without changing global time. let future = now_ms() + 10_000; @@ -113,6 +114,7 @@ async fn due_workflow_activity_survives_a_clock_rollback() { codec_version: 1, schema: 1, target, + owner_fence, sequence: 1, now_ms: future, input: &input, diff --git a/crates/cellule-app/tests/host/replicas.rs b/crates/cellule-app/tests/host/replicas.rs index 8b9028b5..007ac133 100644 --- a/crates/cellule-app/tests/host/replicas.rs +++ b/crates/cellule-app/tests/host/replicas.rs @@ -78,22 +78,22 @@ impl PeerRoundTrip for ReaderHints { } } -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(start_paused = true)] async fn published_command_wakes_reader_recruitment_before_periodic_scan() { publication_hints(1, false).await; } -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(start_paused = true)] async fn stalled_reader_does_not_delay_healthy_reader_publication_hints() { publication_hints(2, true).await; } -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(start_paused = true)] async fn publication_hints_reach_readers_beyond_the_activation_concurrency() { publication_hints(20, true).await; } -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(start_paused = true)] async fn pending_reader_activation_retains_a_new_publication_hint() { publication_hints_with_gate(20, true, true).await; } @@ -109,6 +109,14 @@ async fn publication_hints_with_gate(readers: usize, stalled_reader: bool, hold_ }; use std::time::Duration; + // A periodic pass may coalesce queued hints with the next publication. Keep + // its clock stationary so every observed refresh proves a publication wake-up. + // SQL workers run outside Tokio's blocking pool; this guard inhibits automatic + // clock advancement while they reply, with a real-time bound for broken tests. + let (clock_guard, clock_release) = std::sync::mpsc::channel::<()>(); + let clock_task = tokio::task::spawn_blocking(move || { + let _ = clock_release.recv_timeout(Duration::from_secs(10)); + }); let application = Arc::new(compiled()); let registry = application.registry(); let tenant = TenantId::from_bytes([81; 16]); @@ -232,11 +240,13 @@ async fn publication_hints_with_gate(readers: usize, stalled_reader: bool, hold_ ) .unwrap(); node.start().unwrap(); + let recruitment_clock = tokio::time::Instant::now(); + let recruitment_started = std::time::Instant::now(); // Always drain the host before propagating a failed assertion. A regression // must not leave the intentionally stalled transport alive in the suite. let observed = std::panic::AssertUnwindSafe(async { // Consume the immediate periodic pass before publishing. The next tick is - // five seconds away, so only a publication wake-up can satisfy this bound. + // five virtual seconds away; the two-second bound excludes its repair. for _ in 1..count { tokio::time::timeout(Duration::from_secs(2), hints.recv()) .await @@ -268,21 +278,32 @@ async fn publication_hints_with_gate(readers: usize, stalled_reader: bool, hold_ .as_ref() .filter(|_| occurrence == 1) .map(|held| held.session); - let notified = tokio::time::timeout(Duration::from_secs(2), async { - let mut received = std::collections::HashSet::new(); + let started = std::time::Instant::now(); + let mut notified = std::collections::HashSet::new(); + let completed = tokio::time::timeout(Duration::from_secs(2), async { for _ in 0..healthy.len() - usize::from(held_session.is_some()) { let (cell, session) = hints.recv().await.unwrap(); assert_eq!(cell, target.cell_id()); assert_ne!(Some(session), held_session); assert!( - received.insert(session), + notified.insert(session), "duplicate activation in one publication pass" ); } - received }) .await; - let mut notified = notified.unwrap(); + assert!( + completed.is_ok(), + "publication {occurrence} timed out after {:?}; recruitment elapsed: {:?}; directory age: {}ms; missing readers: {:?}; pending stalled hints: {}", + started.elapsed(), + recruitment_started.elapsed(), + now_ms() - now, + healthy + .difference(¬ified) + .filter(|session| Some(**session) != held_session) + .collect::>(), + pending.load(std::sync::atomic::Ordering::Relaxed), + ); if let Some(held) = held.as_ref().filter(|_| occurrence == 1) { // Healthy peers prove that the new publication pass prepared // while this old activation remained in flight. @@ -303,15 +324,24 @@ async fn publication_hints_with_gate(readers: usize, stalled_reader: bool, hold_ pending.load(std::sync::atomic::Ordering::Relaxed), usize::from(stalled_reader) ); + assert_eq!( + tokio::time::Instant::now(), + recruitment_clock, + "periodic repair must not supply a publication hint" + ); } }) .catch_unwind() .await; + // Drain against real time even if the watchdog released the paused clock. + tokio::time::resume(); tokio::time::timeout(Duration::from_secs(2), node.shutdown()) .await .unwrap() .unwrap(); assert_eq!(pending.load(std::sync::atomic::Ordering::Relaxed), 0); + drop(clock_guard); + clock_task.await.unwrap(); if let Err(failure) = observed { std::panic::resume_unwind(failure); } diff --git a/crates/cellule-runtime/docs/runtime.md b/crates/cellule-runtime/docs/runtime.md index c07d7450..0366675c 100644 --- a/crates/cellule-runtime/docs/runtime.md +++ b/crates/cellule-runtime/docs/runtime.md @@ -450,6 +450,15 @@ activation, never correctness. ## Renew and self-fence ownership +`CommandContext::owner_fence()` exposes the incarnation and epoch stamped on the +admission that accepted the execution. Local typed commands and authenticated +inbox delivery use the same stamp; no authority fetch is added to their command +path. Operation tokens can bind this fence and reject delayed preparation from +a predecessor. Admission replacement within the same ownership epoch retains +the stamp. Recorded outcomes and follower-restored SQL bytes remain unchanged; +recovery does not rerun application handlers with a guessed current fence. + + - **Scanner.** Each owner normally becomes due every three seconds. One node-level scanner finds due owners every 100 ms, orders them by their original deadline, and starts at most 32 renewals concurrently. A completed renewal immediately frees a slot for the next due owner; each scan discards stale candidates. This removes the former 320-starts/s tick ceiling without reducing the one-control-update-per-active-Cell cost. A mutation publication also advances owner progress. - **Renewal budget.** The runtime gives a control-record renewal up to thirty seconds under object-store pressure. - **Session guard.** A separate node-session guard closes admission at its signed expiry, even while renewal I/O is pending. diff --git a/crates/cellule-runtime/docs/rust-api.md b/crates/cellule-runtime/docs/rust-api.md index cbe654b6..af3b3dcb 100644 --- a/crates/cellule-runtime/docs/rust-api.md +++ b/crates/cellule-runtime/docs/rust-api.md @@ -254,11 +254,27 @@ Add exact byte fixtures for every new input and output version. | --- | --- | | `cell_id()` | Read the verified target Cell ID | | `target()` | Derive deterministic effect targets | +| `owner_fence()` | Compare a stored operation token with the admitted incarnation/epoch | | `sequence()` | Allocate stable transition-local identities | | `now_ms()` | Use the runtime-sampled logical timestamp | | `sql()` | Execute an authorized bounded SQL batch | | `emit_effect(&EffectCommandIntent)` | Append one typed cross-Cell intention to the command ledger | +`OwnerFence` is available from `control` and `registry`. The runtime stamps it +on the activation's admission capability and supplies it to typed commands and +inbox effect handlers. Renewal, ordinary publication and schema migration keep +the same incarnation/epoch; an ownership claim advances the epoch. A stale +handle keeps its original value but cannot admit new work. The value alone is +not authority and does not replace current application policy checks or the +runtime's durable response gate. Exact stored-outcome replay skips the handler, +so it preserves the original result rather than substituting the new owner's +fence. + +Bind operation tokens to their Cell ID or target as well as `OwnerFence`. +Embedders that invoke `Registry::execute_command` directly must supply the +admitting handle's fence in `CommandInvocation`; the registry does not verify +ownership independently. + **`QueryContext` exposes** the Cell ID, commit sequence, logical timestamp, and bounded read-only SQL. diff --git a/crates/cellule-runtime/src/cell/actor/admission.rs b/crates/cellule-runtime/src/cell/actor/admission.rs index 742c302c..fab46b15 100644 --- a/crates/cellule-runtime/src/cell/actor/admission.rs +++ b/crates/cellule-runtime/src/cell/actor/admission.rs @@ -98,8 +98,9 @@ pub(super) fn fence_admission(admission: &CellAdmission) { admission.bytes.close(); } -pub(super) fn new_cell_admission() -> Arc { +pub(super) fn new_cell_admission(owner_fence: crate::control::OwnerFence) -> Arc { Arc::new(CellAdmission { + owner_fence, requests: Arc::new(Semaphore::new(CELL_REQUESTS)), bytes: Arc::new(Semaphore::new(CELL_BYTES)), draining: AtomicBool::new(false), diff --git a/crates/cellule-runtime/src/cell/actor/handle.rs b/crates/cellule-runtime/src/cell/actor/handle.rs index ac9bf6bf..9ed69594 100644 --- a/crates/cellule-runtime/src/cell/actor/handle.rs +++ b/crates/cellule-runtime/src/cell/actor/handle.rs @@ -40,6 +40,7 @@ pub struct CellHandle { } pub(super) struct CellAdmission { + pub(super) owner_fence: crate::control::OwnerFence, pub(super) requests: Arc, pub(super) bytes: Arc, pub(super) draining: AtomicBool, @@ -115,6 +116,15 @@ impl CellHandle { self.incarnation } + /// Returns the fence stamped on this activation's admission capability. + /// + /// A stale handle retains its original value but cannot admit new work. + /// Command handlers receive this same value through `CommandContext`. + #[must_use] + pub fn owner_fence(&self) -> crate::control::OwnerFence { + self.admission.owner_fence + } + /// Returns the application code digest the Cell serves. #[must_use] pub const fn code(&self) -> Digest { @@ -375,7 +385,7 @@ impl CellHandle { } self.admission.requests.close(); self.admission.bytes.close(); - let successor_admission = new_cell_admission(); + let successor_admission = new_cell_admission(self.owner_fence()); let (reply, response) = oneshot::channel(); self.inner .sender diff --git a/crates/cellule-runtime/src/cell/actor/task.rs b/crates/cellule-runtime/src/cell/actor/task.rs index 81261401..38dc6276 100644 --- a/crates/cellule-runtime/src/cell/actor/task.rs +++ b/crates/cellule-runtime/src/cell/actor/task.rs @@ -352,7 +352,7 @@ pub(super) fn handle_message( } *next_generation = next_generation.wrapping_add(1).max(1); let generation = *next_generation; - let admission = new_cell_admission(); + let admission = new_cell_admission(publisher.control().value().owner_fence()); let pool = pool.clone(); tasks.spawn(async move { let mut publisher = publisher; @@ -761,7 +761,7 @@ pub(super) fn handle_message( active.admission.draining.store(true, Ordering::Release); active.admission.requests.close(); active.admission.bytes.close(); - active.admission = new_cell_admission(); + active.admission = new_cell_admission(active.admission.owner_fence); active.transfer = Some(TransferPreflight { reply }); }; movement_permits.insert(cell, permit); diff --git a/crates/cellule-runtime/src/cell/actor/tasks/activation.rs b/crates/cellule-runtime/src/cell/actor/tasks/activation.rs index 36e1127a..52b7ded1 100644 --- a/crates/cellule-runtime/src/cell/actor/tasks/activation.rs +++ b/crates/cellule-runtime/src/cell/actor/tasks/activation.rs @@ -34,7 +34,9 @@ pub(super) fn handle_activated( } = context; match result { Ok((interrupt, hydration)) => { - if node_lease.check().is_err() { + if node_lease.check().is_err() + || publisher.control().value().owner_fence() != admission.owner_fence + { fence_admission(&admission); let _ = reply.send(Err(Error::Fenced)); start_orphan_deactivate( diff --git a/crates/cellule-runtime/src/client/local.rs b/crates/cellule-runtime/src/client/local.rs index b9ec5631..140c16db 100644 --- a/crates/cellule-runtime/src/client/local.rs +++ b/crates/cellule-runtime/src/client/local.rs @@ -50,6 +50,7 @@ impl CellTransport for LocalCellTransport { return Err(Error::Command("encoded command input exceeds limit")); } let schema = handle.schema(); + let owner_fence = handle.owner_fence(); let input_bytes = command.input.len(); let output_limit = command.output_limit as usize; handle @@ -70,6 +71,7 @@ impl CellTransport for LocalCellTransport { codec_version: command.codec_version, schema, target: command.target.clone(), + owner_fence, sequence, now_ms, input: &command.input, diff --git a/crates/cellule-runtime/src/control/mod.rs b/crates/cellule-runtime/src/control/mod.rs index 387619ef..798a6bea 100644 --- a/crates/cellule-runtime/src/control/mod.rs +++ b/crates/cellule-runtime/src/control/mod.rs @@ -66,6 +66,21 @@ impl RootRef { } } +/// Incarnation and ownership epoch bound to one admitted Cell activation. +/// +/// Applications compare this value with their stored operation token. A value +/// alone is not authority: runtime admission and durable publication still +/// fence the command. Lease renewal and root publication preserve this fence. +/// This identity is scoped to one Cell; application tokens must also bind their +/// Cell ID or target. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct OwnerFence { + /// Incarnation whose storage and command outcomes this activation serves. + pub incarnation: IncarnationId, + /// Authority epoch that admitted the activation. + pub epoch: u64, +} + /// Enrolled process currently responsible for one Cell. #[derive(Clone, Debug, PartialEq, Eq)] pub struct Owner { @@ -187,6 +202,15 @@ pub enum Transition { } impl Control { + /// Returns this record's incarnation/epoch identity, without granting authority. + #[must_use] + pub const fn owner_fence(&self) -> OwnerFence { + OwnerFence { + incarnation: self.incarnation, + epoch: self.epoch, + } + } + /// Creates and validates the only legal first control record. pub fn initial( cell: CellId, diff --git a/crates/cellule-runtime/src/peer/dispatch/mod.rs b/crates/cellule-runtime/src/peer/dispatch/mod.rs index f27897ec..26359a57 100644 --- a/crates/cellule-runtime/src/peer/dispatch/mod.rs +++ b/crates/cellule-runtime/src/peer/dispatch/mod.rs @@ -567,6 +567,7 @@ impl PeerDispatcher { let registry = Arc::clone(&self.registry); let telemetry = self.telemetry.clone(); let schema = transport.handle.schema(); + let owner_fence = transport.handle.owner_fence(); let input = command.input.clone(); let operation_id = command.command_id; let codec_version = command.codec_version; @@ -588,6 +589,7 @@ impl PeerDispatcher { codec_version, schema, target: target.clone(), + owner_fence, sequence, now_ms: logical_time_ms, input: &input, diff --git a/crates/cellule-runtime/src/registry/builder/mod.rs b/crates/cellule-runtime/src/registry/builder/mod.rs index 666fb197..6e1003ae 100644 --- a/crates/cellule-runtime/src/registry/builder/mod.rs +++ b/crates/cellule-runtime/src/registry/builder/mod.rs @@ -941,7 +941,7 @@ impl Registry { invocation: CommandInvocation<'_>, issued_at_ms: i64, ) -> Result { - if invocation.sequence == 0 || invocation.now_ms < 0 { + if invocation.sequence == 0 || invocation.now_ms < 0 || invocation.owner_fence.epoch == 0 { return Err(Error::Command("invalid registered command context")); } let key = BindingKey::new( @@ -969,6 +969,7 @@ impl Registry { let mut context = CommandContext { transaction, target: invocation.target.clone(), + owner_fence: invocation.owner_fence, effect_targets: self .namespace_modules .get(&invocation.target.namespace()) diff --git a/crates/cellule-runtime/src/registry/handlers.rs b/crates/cellule-runtime/src/registry/handlers.rs index 34107dd3..73851b64 100644 --- a/crates/cellule-runtime/src/registry/handlers.rs +++ b/crates/cellule-runtime/src/registry/handlers.rs @@ -11,6 +11,7 @@ use super::*; pub struct CommandContext<'borrow, 'connection> { pub(super) transaction: &'borrow Transaction<'connection>, pub(super) target: CellTarget, + pub(super) owner_fence: OwnerFence, pub(super) effect_targets: &'static [NamespaceId], pub(super) sequence: u64, pub(super) now_ms: i64, @@ -33,6 +34,16 @@ impl CommandContext<'_, '_> { &self.target } + /// Returns the incarnation/epoch stamped on the admitted activation. + /// + /// Compare it with application operation tokens before staging or publishing. + /// Client inputs cannot select this value. A recorded result is replayed + /// without running the handler again; replay does not substitute a new fence. + #[must_use] + pub const fn owner_fence(&self) -> OwnerFence { + self.owner_fence + } + /// Returns the actor-ordered sequence the command was admitted at. #[must_use] pub const fn sequence(&self) -> u64 { @@ -325,6 +336,8 @@ pub struct CommandInvocation<'a> { pub schema: u32, /// Validated target the command must own. pub target: CellTarget, + /// Fence supplied by the runtime-validated admission, never client input. + pub owner_fence: OwnerFence, /// Actor-ordered sequence of the command. pub sequence: u64, /// Logical runtime time for the command. diff --git a/crates/cellule-runtime/src/registry/mod.rs b/crates/cellule-runtime/src/registry/mod.rs index f1533a9b..75404d54 100644 --- a/crates/cellule-runtime/src/registry/mod.rs +++ b/crates/cellule-runtime/src/registry/mod.rs @@ -1,4 +1,5 @@ //! Compiled primitive registry: descriptors, schemas, handlers, and builder. +pub use crate::control::OwnerFence; mod descriptor; pub use builder::{Registry, RegistryBuilder}; pub use handlers::{ diff --git a/crates/cellule-runtime/tests/contracts/registry.rs b/crates/cellule-runtime/tests/contracts/registry.rs index 1550ced4..7df4efb3 100644 --- a/crates/cellule-runtime/tests/contracts/registry.rs +++ b/crates/cellule-runtime/tests/contracts/registry.rs @@ -342,6 +342,10 @@ fn compiled_registry_is_canonical_and_executes_only_declared_bindings() { codec_version: 1, schema: 1, target, + owner_fence: cellule_runtime::control::OwnerFence { + incarnation: cellule_runtime::identity::IncarnationId::from_bytes([7; 16]), + epoch: 3, + }, sequence: 1, now_ms: 10, input: &input, @@ -414,6 +418,10 @@ fn command_execution_rejects_a_module_targeting_another_namespace_owner() { codec_version: 1, schema: 1, target, + owner_fence: cellule_runtime::control::OwnerFence { + incarnation: cellule_runtime::identity::IncarnationId::from_bytes([7; 16]), + epoch: 3 + }, sequence: 1, now_ms: 10, input: &input, @@ -667,3 +675,45 @@ fn dead_letter_cycles_are_rejected() { Err(Error::Registry("dead-letter cycle")) )); } + +#[test] +fn registered_commands_reject_an_unadmitted_zero_epoch_before_running_the_handler() { + let registry = build_registry(false); + let mut connection = cellule_ltx::rusqlite::Connection::open_in_memory().unwrap(); + connection.execute_batch(MIGRATION).unwrap(); + let tx = connection.transaction().unwrap(); + let input = wire(&b"unadmitted".to_vec(), 16); + let target = CellTarget::new( + TenantId::from_bytes([1; 16]), + ApplicationId::from_bytes([2; 16]), + NamespaceId::from_bytes([1; 16]), + b"unadmitted", + ) + .unwrap(); + let outcome = registry.execute_command( + &tx, + CommandInvocation { + module: "first", + operation_id: 1, + codec_version: 1, + schema: 1, + target, + sequence: 1, + now_ms: 10, + input: &input, + owner_fence: cellule_runtime::control::OwnerFence { + incarnation: cellule_runtime::identity::IncarnationId::from_bytes([7; 16]), + epoch: 0, + }, + }, + ); + assert!(matches!( + outcome, + Err(Error::Command("invalid registered command context")) + )); + assert_eq!( + tx.query_row("SELECT count(*) FROM items", [], |r| r.get::<_, u64>(0)) + .unwrap(), + 0 + ); +} diff --git a/crates/cellule-runtime/tests/protocol/client.rs b/crates/cellule-runtime/tests/protocol/client.rs index 4d30cf38..1d770e86 100644 --- a/crates/cellule-runtime/tests/protocol/client.rs +++ b/crates/cellule-runtime/tests/protocol/client.rs @@ -117,6 +117,14 @@ const COMMANDS: &[OperationDescriptor] = &[ input_limit: 8, output_limit: 1, }, + OperationDescriptor { + id: 9, + codec_version: 1, + schema_min: 1, + schema_max: 1, + input_limit: 64, + output_limit: 64, + }, ]; const QUERIES: &[OperationDescriptor] = &[ OperationDescriptor { @@ -161,6 +169,7 @@ impl CellModule for RepositoryModule { registry.bind_command::()?; registry.bind_command::()?; registry.bind_command::()?; + registry.bind_command::()?; registry.bind_query::()?; register_effect_delivery::(registry)?; Ok(()) @@ -175,6 +184,24 @@ impl EffectModule for RepositoryModule { const STATUS_QUERY_ID: u32 = 3; } +struct ObserveOwnerFence; +impl Command for ObserveOwnerFence { + const MODULE: &'static str = MODULE; + const ID: u32 = 9; + const CODEC_VERSION: u32 = 1; + type Input = Vec; + type Output = Vec; + fn execute( + context: &mut CommandContext<'_, '_>, + _untrusted_input: Vec, + ) -> cellule_runtime::Result>> { + let fence = context.owner_fence(); + let mut bytes = fence.incarnation.as_bytes().to_vec(); + bytes.extend_from_slice(&fence.epoch.to_be_bytes()); + Ok(CommandResult::Success(bytes)) + } +} + struct CreateComment; impl Command for CreateComment { diff --git a/crates/cellule-runtime/tests/protocol/client/effects.rs b/crates/cellule-runtime/tests/protocol/client/effects.rs index 739aaa68..982a4eb3 100644 --- a/crates/cellule-runtime/tests/protocol/client/effects.rs +++ b/crates/cellule-runtime/tests/protocol/client/effects.rs @@ -2,9 +2,10 @@ use super::*; -#[tokio::test] -async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() { - let fixture = fixture().await; +fn effect_client( + fixture: &Fixture, + handle: cellule_runtime::cell::actor::CellHandle, +) -> EffectPeerClient { let signer = PeerSigner::new( SessionId::from_bytes([12; 16]), fixture.registry.release_digest(), @@ -19,7 +20,7 @@ async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() Arc::clone(&fixture.registry), Arc::new(LocalResolver { target: fixture.target.clone(), - handle: fixture.handle().clone(), + handle, }), Arc::new(RepositoryAuthorizer), )); @@ -32,6 +33,13 @@ async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() subject: "source-session".into(), actions: vec!["repository.issue.create".into()], }; + EffectPeerClient::new(Arc::new(signer), principal, round_trip) +} + +#[tokio::test] +async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() { + let mut fixture = fixture().await; + let client = effect_client(&fixture, fixture.handle().clone()); let now_ms = i64::try_from( std::time::SystemTime::now() .duration_since(UNIX_EPOCH) @@ -87,7 +95,6 @@ async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() expires_at_ms: identity.expires_at_ms, created_sequence: source_sequence, }; - let client = EffectPeerClient::new(Arc::new(signer), principal, round_trip); let delivered = client.deliver(&claim, now_ms).await.unwrap(); assert_eq!(delivered.commit_sequence(), 1); assert_eq!(client.deliver(&claim, now_ms + 1).await.unwrap(), delivered); @@ -104,7 +111,106 @@ async fn authenticated_effect_delivery_publishes_once_and_resolves_from_inbox() .output, 1 ); - fixture.handle().drain().await.unwrap(); + // Inbox dispatch supplies the destination's admitted fence; request bytes + // cannot replace it. The exact stored result also survives effect replay. + let fence = fixture.handle().owner_fence(); + let mut expected_fence = fence.incarnation.as_bytes().to_vec(); + expected_fence.extend_from_slice(&fence.epoch.to_be_bytes()); + let mut fence_request = request; + let mut fence_identity = identity; + fence_identity.source_sequence = 10; + fence_identity.ordinal = 4; + let fence_id = + cellule_runtime::primitives::effects::effect_id(source_cell, source_incarnation, 10, 4); + fence_identity.effect_id = fence_id.to_vec(); + fence_request.identity = Some(fence_identity); + let mut fence_input = BoundedEncoder::new(64).unwrap(); + vec![99; 24].encode(&mut fence_input).unwrap(); + fence_request.operation = Some(wire::effect_request::Operation::CellCommand( + wire::CellCommand { + command_id: ObserveOwnerFence::ID, + codec_version: 1, + input: fence_input.finish(), + }, + )); + let mut fence_claim = claim.clone(); + fence_claim.effect_id = fence_id; + fence_claim.operation = prost::Message::encode_to_vec(&fence_request); + fence_claim.operation_digest = + effect_operation_digest(fixture.target.cell_id(), fence_id, &fence_claim.operation); + fence_claim.created_sequence = 10; + let observed = client.deliver(&fence_claim, now_ms + 3).await.unwrap(); + let mut decoder = BoundedDecoder::new(observed.result(), 64).unwrap(); + assert_eq!(Vec::::decode(&mut decoder).unwrap(), expected_fence); + decoder.finish().unwrap(); + assert_eq!( + client.deliver(&fence_claim, now_ms + 4).await.unwrap(), + observed + ); + assert_eq!( + client.resolve(&fence_claim, now_ms + 5).await.unwrap(), + Resolution::Committed(observed.clone()) + ); + let old_handle = fixture.take_handle(); + old_handle.drain().await.unwrap(); + fixture.runtime.take().unwrap().shutdown().await.unwrap(); + let session = SessionId::from_bytes([19; 16]); + let runtime = CellRuntime::new(SqlWorkerPool::new(1, 4).unwrap(), 4 << 20, session).unwrap(); + let idle = fixture + .authority + .load(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(); + let successor = runtime + .acquire_idle_restored( + fixture.proof.clone(), + fixture.replica.clone(), + fixture.authority.clone(), + idle, + fixture + ._directory + .path() + .join("inbox-fence-successor.sqlite"), + Owner { + session, + endpoint: "https://inbox-fence-successor.invalid".into(), + }, + ) + .await + .unwrap(); + let new_fence = successor.owner_fence(); + assert_eq!(new_fence.incarnation, fence.incarnation); + assert!(new_fence.epoch > fence.epoch); + let client = effect_client(&fixture, successor); + // A new destination owner must replay the original inbox output, while a + // new effect invokes the same handler with the successor's admission fence. + assert_eq!( + client.deliver(&fence_claim, now_ms + 6).await.unwrap(), + observed + ); + assert_eq!( + client.resolve(&fence_claim, now_ms + 7).await.unwrap(), + Resolution::Committed(observed) + ); + let fresh_id = + cellule_runtime::primitives::effects::effect_id(source_cell, source_incarnation, 11, 5); + let fresh_identity = fence_request.identity.as_mut().unwrap(); + fresh_identity.effect_id = fresh_id.to_vec(); + fresh_identity.source_sequence = 11; + fresh_identity.ordinal = 5; + fence_claim.effect_id = fresh_id; + fence_claim.created_sequence = 11; + fence_claim.operation = prost::Message::encode_to_vec(&fence_request); + fence_claim.operation_digest = + effect_operation_digest(fixture.target.cell_id(), fresh_id, &fence_claim.operation); + let fresh = client.deliver(&fence_claim, now_ms + 8).await.unwrap(); + let mut expected = new_fence.incarnation.as_bytes().to_vec(); + expected.extend_from_slice(&new_fence.epoch.to_be_bytes()); + let mut decoder = BoundedDecoder::new(fresh.result(), 64).unwrap(); + assert_eq!(Vec::::decode(&mut decoder).unwrap(), expected); + decoder.finish().unwrap(); + runtime.shutdown().await.unwrap(); } #[tokio::test] async fn typed_effect_source_publishes_claim_validation_ack_and_lost_lease() { diff --git a/crates/cellule-runtime/tests/runtime.rs b/crates/cellule-runtime/tests/runtime.rs index 8c2e0600..cbdcb9b0 100644 --- a/crates/cellule-runtime/tests/runtime.rs +++ b/crates/cellule-runtime/tests/runtime.rs @@ -12,6 +12,7 @@ mod runtime { pub mod fault_fs; pub mod lifecycle; pub mod migration; + pub mod owner_fence; pub mod publication; pub mod release_progress; pub mod scheduler; diff --git a/crates/cellule-runtime/tests/runtime/lifecycle.rs b/crates/cellule-runtime/tests/runtime/lifecycle.rs index e8f4c391..7c68045a 100644 --- a/crates/cellule-runtime/tests/runtime/lifecycle.rs +++ b/crates/cellule-runtime/tests/runtime/lifecycle.rs @@ -57,7 +57,7 @@ pub mod read_replica; pub mod residency; #[derive(Debug)] -struct PausingStore { +pub(super) struct PausingStore { inner: Arc, armed: AtomicBool, update_armed: AtomicBool, @@ -81,7 +81,7 @@ struct PausingStore { } impl PausingStore { - fn new(inner: Arc) -> Self { + pub(super) fn new(inner: Arc) -> Self { Self { inner, armed: AtomicBool::new(false), @@ -114,7 +114,7 @@ impl PausingStore { self.armed.store(true, Ordering::Release); } - fn arm_next_update(&self) { + pub(super) fn arm_next_update(&self) { self.update_armed.store(true, Ordering::Release); } @@ -144,13 +144,13 @@ impl PausingStore { } } - async fn wait_until_blocked(&self) { + pub(super) async fn wait_until_blocked(&self) { while !self.blocked.load(Ordering::Acquire) { self.entered.notified().await; } } - fn release(&self) { + pub(super) fn release(&self) { self.released.store(true, Ordering::Release); self.release.notify_waiters(); } @@ -319,7 +319,7 @@ impl ObjectStore for PausingStore { } #[derive(Default)] -struct TestNodeAuthority { +pub(super) struct TestNodeAuthority { activations: Mutex>, coverage: Mutex>, closes: Mutex>, @@ -481,7 +481,7 @@ impl NodeLogTransport for LostAckFollowerTransport { } } -async fn fence_log_session( +pub(super) async fn fence_log_session( layout: &CellStorageLayout, session: SessionId, claimant: SessionId, diff --git a/crates/cellule-runtime/tests/runtime/lifecycle/idle/release.rs b/crates/cellule-runtime/tests/runtime/lifecycle/idle/release.rs index b7dfb1a8..c2f7f28d 100644 --- a/crates/cellule-runtime/tests/runtime/lifecycle/idle/release.rs +++ b/crates/cellule-runtime/tests/runtime/lifecycle/idle/release.rs @@ -175,6 +175,12 @@ async fn exact_idle_release_refuses_persisted_work() { .is_err() ); assert_eq!(runtime.stats().active_cells(), 1); + let refreshed = runtime + .resident_handle(&fixture.target, CatalogRole::Application) + .await + .unwrap() + .unwrap(); + assert_eq!(refreshed.owner_fence(), handle.owner_fence()); assert!(matches!( handle.drain().await, Err(cellule_runtime::Error::CellDraining) diff --git a/crates/cellule-runtime/tests/runtime/lifecycle/residency.rs b/crates/cellule-runtime/tests/runtime/lifecycle/residency.rs index dd73e4d1..1336d5d5 100644 --- a/crates/cellule-runtime/tests/runtime/lifecycle/residency.rs +++ b/crates/cellule-runtime/tests/runtime/lifecycle/residency.rs @@ -1918,6 +1918,7 @@ async fn slow_control_renewal_keeps_a_live_owner() { Store::new(store.clone()), ); let (runtime, handle, _) = activate_runtime(&fixture, 16 * 1024 * 1024).await; + let owner_fence = handle.owner_fence(); let authority = CellAuthority::new(fixture.layout.clone()); let initial = authority .load(fixture.target.cell_id()) @@ -1955,6 +1956,8 @@ async fn slow_control_renewal_keeps_a_live_owner() { renewed.value().owner.as_ref().unwrap().session, SessionId::from_bytes([4; 16]) ); + assert_eq!(renewed.value().owner_fence(), owner_fence); + assert_eq!(handle.owner_fence(), owner_fence); handle.drain().await.unwrap(); runtime.shutdown().await.unwrap(); } diff --git a/crates/cellule-runtime/tests/runtime/migration/schema.rs b/crates/cellule-runtime/tests/runtime/migration/schema.rs index cea96385..6f468d74 100644 --- a/crates/cellule-runtime/tests/runtime/migration/schema.rs +++ b/crates/cellule-runtime/tests/runtime/migration/schema.rs @@ -111,6 +111,7 @@ async fn migration_replaces_capability_publishes_schema_and_restores_exact_root( .await .unwrap() .unwrap(); + assert_eq!(migrated.handle.owner_fence(), old_handle.owner_fence()); object_store.started.wait().await; let queued_handle = migrated.handle.clone(); let mut queued = tokio::spawn(async move { @@ -325,6 +326,7 @@ async fn code_only_migration_publishes_new_code_without_schema_ledger_entry() { assert_eq!(plan.sql(), None); let migrated = handle.migrate(plan, 10).await.unwrap(); + assert_eq!(migrated.handle.owner_fence(), old_handle.owner_fence()); assert_eq!(migrated.outcome.code, code); assert_eq!(migrated.outcome.schema, 2); assert_eq!(migrated.outcome.commit_sequence, 1); diff --git a/crates/cellule-runtime/tests/runtime/owner_fence.rs b/crates/cellule-runtime/tests/runtime/owner_fence.rs new file mode 100644 index 00000000..95b0c749 --- /dev/null +++ b/crates/cellule-runtime/tests/runtime/owner_fence.rs @@ -0,0 +1,515 @@ +//! Admission fence exposure, stale operation rejection and exact outcome replay. +use crate::support::fixtures::mutation_identity; +use cellule_ltx::{CellReplica, CellStorageLayout, Limits}; +use cellule_runtime::cell::actor::CellHandle; +use cellule_runtime::cell::catalog::{CatalogEntry, CatalogRole, CellCatalog}; +use cellule_runtime::control::authority::CellAuthority; +use cellule_runtime::control::{Owner, OwnerFence}; +use cellule_runtime::identity::{ + ApplicationId, CellTarget, Digest, IncarnationId, NamespaceId, SessionId, TenantId, +}; +use cellule_runtime::primitives::sql::{SqlBatch, SqlStatement, SqlValue}; +use cellule_runtime::registry::{ + BuildDescriptor, CellModule, Command, CommandContext, CommandResult, MigrationDescriptor, + ModuleDescriptor, NamespaceDescriptor, OperationDescriptor, +}; +use cellule_runtime::{ + CellClient, CellRuntime, InvocationError, Registry, RegistryBuilder, SqlWorkerPool, +}; +use cellule_store::Store; +use object_store::{memory::InMemory, path::Path}; +use std::sync::{Arc, OnceLock}; + +const MODULE: &str = "owner-fence-test"; +const NAMESPACE: NamespaceId = NamespaceId::from_bytes([12; 16]); +const SCHEMA: &str = "CREATE TABLE fence_executions(fence BLOB NOT NULL)"; +const COMMANDS: &[OperationDescriptor] = &[OperationDescriptor { + id: 1, + codec_version: 1, + schema_min: 1, + schema_max: 1, + input_limit: 64, + output_limit: 64, +}]; +struct FenceModule; +impl CellModule for FenceModule { + const NAME: &'static str = MODULE; + fn descriptor(&self) -> &'static ModuleDescriptor { + static DESCRIPTOR: OnceLock = OnceLock::new(); + DESCRIPTOR.get_or_init(|| ModuleDescriptor { + name: MODULE, + source_digest: Digest::from_bytes([13; 32]), + retained_codes: &[], + schema_min: 1, + schema_max: 1, + migrations: Box::leak(Box::new([MigrationDescriptor { + version: 1, + sql: SCHEMA, + digest: Digest::from_bytes(*blake3::hash(SCHEMA.as_bytes()).as_bytes()), + }])), + commands: COMMANDS, + queries: &[], + workflow_definitions: &[], + activity_types: &[], + namespaces: Box::leak(Box::new([NamespaceDescriptor { + id: NAMESPACE, + name: MODULE, + role: CatalogRole::Application, + shards: 1, + effect_targets: &[], + dead_letter: None, + }])), + }) + } + fn register(self, registry: &mut RegistryBuilder) -> cellule_runtime::Result<()> { + registry.bind_command::() + } +} +fn encoded(fence: OwnerFence) -> Vec { + let mut bytes = fence.incarnation.as_bytes().to_vec(); + bytes.extend_from_slice(&fence.epoch.to_be_bytes()); + bytes +} +struct FencedCommand; +impl Command for FencedCommand { + const MODULE: &'static str = MODULE; + const ID: u32 = 1; + const CODEC_VERSION: u32 = 1; + type Input = Vec; + type Output = Vec; + fn execute( + context: &mut CommandContext<'_, '_>, + expected: Vec, + ) -> cellule_runtime::Result>> { + let fence = encoded(context.owner_fence()); + if !expected.is_empty() && expected != fence { + return Ok(CommandResult::Rejected(fence)); + } + context.sql(&SqlBatch { + statements: vec![SqlStatement { + sql: "INSERT INTO fence_executions VALUES(?1)".into(), + parameters: vec![SqlValue::Blob(fence.clone())], + }], + })?; + Ok(CommandResult::Success(fence)) + } +} +struct Fixture { + root: tempfile::TempDir, + target: CellTarget, + layout: CellStorageLayout, + replica: CellReplica, + registry: Arc, + runtime: CellRuntime, + handle: CellHandle, +} +impl Fixture { + async fn new() -> Self { + let session = SessionId::from_bytes([18; 16]); + let runtime = + CellRuntime::new(SqlWorkerPool::new(1, 4).unwrap(), 64 << 20, session).unwrap(); + Self::with_runtime(session, Store::new(Arc::new(InMemory::new())), runtime).await + } + + async fn with_runtime(session: SessionId, store: Store, runtime: CellRuntime) -> Self { + let mut builder = RegistryBuilder::new(BuildDescriptor { + source_revision: "owner-fence-test".into(), + cargo_lock_digest: Digest::from_bytes([14; 32]), + }); + builder.register(FenceModule).unwrap(); + let registry = Arc::new(builder.finish().unwrap()); + let target = CellTarget::new( + TenantId::from_bytes([15; 16]), + ApplicationId::from_bytes([16; 16]), + NAMESPACE, + b"operation", + ) + .unwrap(); + let layout = CellStorageLayout::new(store, Path::from("owner-fence"), [16; 16]); + let incarnation = IncarnationId::from_bytes([17; 16]); + let replica = CellReplica::new( + layout.clone(), + *target.cell_id().as_bytes(), + *incarnation.as_bytes(), + Limits::default(), + ) + .unwrap(); + let root = tempfile::TempDir::new().unwrap(); + let proof = CellCatalog::new(layout.clone(), target.tenant()) + .provision( + CatalogEntry::new( + &target, + CatalogRole::Application, + registry.module_code(MODULE).unwrap(), + 1, + ) + .unwrap(), + ) + .await + .unwrap(); + let authority = CellAuthority::new(layout.clone()); + let control = authority + .create_initial( + &proof, + incarnation, + Owner { + session, + endpoint: "https://owner-a.invalid".into(), + }, + ) + .await + .unwrap(); + let handle = runtime + .bootstrap( + proof, + replica.clone(), + authority, + control, + root.path().join("a.sqlite"), + |tx| { + tx.execute_batch(SCHEMA)?; + Ok(()) + }, + ) + .await + .unwrap(); + Self { + root, + target, + layout, + replica, + registry, + runtime, + handle, + } + } + fn client(&self) -> CellClient { + CellClient::local(self.registry.clone(), self.handle.clone()) + } +} +async fn rows(handle: &CellHandle) -> Vec> { + let bytes = handle + .query(0, 1024, |connection| { + let mut statement = + connection.prepare("SELECT fence FROM fence_executions ORDER BY rowid")?; + let rows: Vec> = statement + .query_map([], |row| row.get(0))? + .collect::>()?; + Ok(rows.concat()) + }) + .await + .unwrap(); + bytes + .as_chunks::<24>() + .0 + .iter() + .map(|row| row.to_vec()) + .collect() +} + +#[tokio::test] +async fn typed_handlers_observe_the_admitted_fence_and_replay_does_not_reexecute() { + let fixture = Fixture::new().await; + let client = fixture.client(); + let fence = fixture.handle.owner_fence(); + let observed = CellAuthority::new(fixture.layout.clone()) + .load(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(); + assert_eq!(fence, observed.value().owner_fence()); + let identity = mutation_identity(21); + let first = client + .command::(&fixture.target, identity, Vec::new()) + .await + .unwrap(); + assert_eq!(first.output, encoded(fence)); + let replay = client + .command::(&fixture.target, identity, Vec::new()) + .await + .unwrap(); + assert_eq!(replay.output, first.output); + assert_eq!(replay.receipt, first.receipt); + assert_eq!(rows(&fixture.handle).await, vec![encoded(fence)]); + // An unrelated root publication keeps the epoch of the same admission. + let second = client + .command::(&fixture.target, mutation_identity(22), encoded(fence)) + .await + .unwrap(); + assert_eq!(second.output, encoded(fence)); + let mut wrong_incarnation = fence; + wrong_incarnation.incarnation = IncarnationId::from_bytes([99; 16]); + let mismatch = client + .command::( + &fixture.target, + mutation_identity(27), + encoded(wrong_incarnation), + ) + .await; + assert!( + matches!(mismatch,Err(InvocationError::Rejected(ref outcome)) if outcome.output==encoded(fence)) + ); + assert_eq!( + rows(&fixture.handle).await, + vec![encoded(fence), encoded(fence)] + ); + let resident = fixture + .runtime + .resident_handle(&fixture.target, CatalogRole::Application) + .await + .unwrap() + .unwrap(); + assert_eq!(resident.owner_fence(), fence); + fixture.runtime.shutdown().await.unwrap(); +} + +#[tokio::test] +async fn successor_owner_rejects_the_old_token_but_replays_its_recorded_outcome() { + let fixture = Fixture::new().await; + let old = fixture.handle.owner_fence(); + let identity = mutation_identity(23); + let committed = fixture + .client() + .command::(&fixture.target, identity, Vec::new()) + .await + .unwrap(); + fixture.handle.drain().await.unwrap(); + fixture.runtime.shutdown().await.unwrap(); + let session = SessionId::from_bytes([19; 16]); + let runtime = CellRuntime::new(SqlWorkerPool::new(1, 4).unwrap(), 64 << 20, session).unwrap(); + let authority = CellAuthority::new(fixture.layout.clone()); + let idle = authority + .load(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(); + let proof = CellCatalog::new(fixture.layout.clone(), fixture.target.tenant()) + .lookup(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(); + let successor = runtime + .acquire_idle_restored( + proof, + fixture.replica.clone(), + authority, + idle, + fixture.root.path().join("b.sqlite"), + Owner { + session, + endpoint: "https://owner-b.invalid".into(), + }, + ) + .await + .unwrap(); + let new = successor.owner_fence(); + assert_eq!(new.incarnation, old.incarnation); + assert!(new.epoch > old.epoch); + let client = CellClient::local(fixture.registry.clone(), successor.clone()); + let replay = client + .command::(&fixture.target, identity, Vec::new()) + .await + .unwrap(); + assert_eq!(replay.output, committed.output); // Original fence A, not B. + assert_eq!(replay.receipt, committed.receipt); + assert_eq!(rows(&successor).await, vec![encoded(old)]); + let stale = client + .command::(&fixture.target, mutation_identity(24), encoded(old)) + .await; + assert!( + matches!(stale,Err(InvocationError::Rejected(ref result)) if result.output==encoded(new)) + ); + assert_eq!(rows(&successor).await, vec![encoded(old)]); + let result = client + .command::(&fixture.target, mutation_identity(25), encoded(new)) + .await + .unwrap(); + assert_eq!(result.output, encoded(new)); + assert_eq!(rows(&successor).await, vec![encoded(old), encoded(new)]); + assert_eq!(fixture.handle.owner_fence(), old); + assert!( + fixture + .client() + .command::(&fixture.target, mutation_identity(26), Vec::new()) + .await + .is_err() + ); + runtime.shutdown().await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread")] +async fn follower_recovered_typed_outcome_keeps_its_fence_after_owner_loss() { + use super::lifecycle::{PausingStore, TestNodeAuthority, fence_log_session}; + use cellule_runtime::identity::NodeId; + use cellule_runtime::ltx::{DiskBudget, Host as ReplicaHost}; + use cellule_runtime::node::durability::NodeDurability; + use cellule_runtime::node::lease::NodeLeaseGuard; + use cellule_runtime::node::log::DurabilityGate; + use cellule_runtime::node::log_recovery::{ + NodeLogRecovery, RecoveryCoordinator, recoverable_cells, + }; + use cellule_runtime::node::log_shipper::NodeLogShipper; + use cellule_runtime::node::log_transport::{LocalFollowerTransport, NodeLogTransport}; + use cellule_runtime::recovery::manifest::RecoveryManifestStore; + use std::time::Duration; + + let session = SessionId::from_bytes([18; 16]); + let follower_session = SessionId::from_bytes([30; 16]); + let follower = NodeId::from_bytes(*follower_session.as_bytes()); + let successor_session = SessionId::from_bytes([31; 16]); + let runtime = CellRuntime::new_with_replica_host_requiring_node_lease( + SqlWorkerPool::new(1, 4).unwrap(), + 64 << 20, + session, + ReplicaHost::default(), + ) + .unwrap(); + let lease = NodeLeaseGuard::new(0, 60_000).unwrap(); + runtime.install_node_lease(lease.clone()).unwrap(); + let follower_directory = tempfile::TempDir::new().unwrap(); + let follower_store = cellule_runtime::FollowerStore::open( + follower_directory.path().to_owned(), + Limits::default(), + DiskBudget::new(1 << 30), + ) + .unwrap(); + let transport: Arc = Arc::new(LocalFollowerTransport::new( + follower, + follower_store.clone(), + )); + let gate = DurabilityGate::new( + session, + NodeId::from_bytes(*session.as_bytes()), + 1, + [follower], + ) + .unwrap(); + let shipper = NodeLogShipper::new(gate.clone(), transport.clone(), Limits::default()).unwrap(); + runtime + .install_node_durability( + ApplicationId::from_bytes([16; 16]), + Arc::new(NodeDurability::new( + gate, + shipper, + Arc::new(TestNodeAuthority::default()), + transport.clone(), + lease.clone(), + )), + ) + .unwrap(); + let store = Arc::new(PausingStore::new(Arc::new(InMemory::new()))); + let fixture = Fixture::with_runtime(session, Store::new(store.clone()), runtime).await; + let old = fixture.handle.owner_fence(); + let authority = CellAuthority::new(fixture.layout.clone()); + let catalog = CellCatalog::new(fixture.layout.clone(), fixture.target.tenant()); + let identity = mutation_identity(28); + store.arm_next_update(); + let committed = tokio::time::timeout( + Duration::from_secs(5), + fixture + .client() + .command::(&fixture.target, identity, Vec::new()), + ) + .await + .unwrap() + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), store.wait_until_blocked()) + .await + .unwrap(); + let pending = authority + .load(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(); + assert_eq!(pending.value().root.as_ref().unwrap().commit_sequence, 0); + assert_eq!(committed.receipt.commit_sequence, 1); + assert_eq!(committed.output, encoded(old)); + assert!(follower_store.retained_bytes() > 0); + + // Fence the owner with its object CAS still paused. Recovery must obtain + // the command and its original output from the actual fsynced follower log. + lease.fence(); + let fenced = fence_log_session( + &fixture.layout, + session, + successor_session, + follower_session, + 0, + ) + .await; + let recovery = NodeLogRecovery::from_fenced(transport, &fenced, Limits::default()).unwrap(); + let manifests = RecoveryManifestStore::new(fixture.layout.clone(), Limits::default()); + let coordinator = RecoveryCoordinator::new(recovery, manifests.clone()); + let inventory = recoverable_cells(&catalog, &authority, session, 10) + .await + .unwrap(); + let directory = cellule_runtime::node::NodeDirectory::new( + fixture.layout.clone(), + Digest::from_bytes([90; 32]), + Digest::from_bytes([91; 32]), + Digest::from_bytes([92; 32]), + ); + let completed = coordinator + .recover_and_seal(&directory, fenced, inventory, 10_002) + .await + .unwrap(); + // The recovery attachment supersedes the paused CAS before it resumes. + store.release(); + assert!( + tokio::time::timeout(Duration::from_secs(5), fixture.runtime.shutdown()) + .await + .unwrap() + .is_err() + ); + let runtime = CellRuntime::new( + SqlWorkerPool::new(1, 4).unwrap(), + 64 << 20, + successor_session, + ) + .unwrap(); + let successor = runtime + .takeover_restored( + catalog + .lookup(fixture.target.cell_id()) + .await + .unwrap() + .unwrap(), + fixture.replica.clone(), + authority, + completed.controls.into_iter().next().unwrap(), + completed.takeover, + manifests, + fixture.root.path().join("follower-successor.sqlite"), + Owner { + session: successor_session, + endpoint: "https://follower-successor.invalid".into(), + }, + ) + .await + .unwrap(); + let new = successor.owner_fence(); + assert_eq!(new.incarnation, old.incarnation); + assert!(new.epoch > old.epoch); + let client = CellClient::local(fixture.registry.clone(), successor.clone()); + let replay = client + .command::(&fixture.target, identity, Vec::new()) + .await + .unwrap(); + assert_eq!(replay.output, committed.output); + assert_eq!(replay.receipt, committed.receipt); + assert_eq!(rows(&successor).await, vec![encoded(old)]); + let stale = client + .command::(&fixture.target, mutation_identity(29), encoded(old)) + .await; + assert!( + matches!(stale, Err(InvocationError::Rejected(ref outcome)) if outcome.output == encoded(new)) + ); + assert_eq!(rows(&successor).await, vec![encoded(old)]); + let current = client + .command::(&fixture.target, mutation_identity(32), encoded(new)) + .await + .unwrap(); + assert_eq!(current.output, encoded(new)); + assert_eq!(rows(&successor).await, vec![encoded(old), encoded(new)]); + runtime.shutdown().await.unwrap(); +}