Skip to content
2 changes: 1 addition & 1 deletion crates/cellule-app/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
7 changes: 7 additions & 0 deletions crates/cellule-app/PERFORMANCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
74 changes: 72 additions & 2 deletions crates/cellule-app/tests/fleet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -181,10 +181,26 @@ pub(super) async fn send_tcp(
request: Vec<u8>,
remaining_ms: u32,
) -> Result<Vec<u8>> {
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<Output = std::io::Result<TcpStream>>,
request: Vec<u8>,
remaining_ms: u32,
) -> Result<Vec<u8>> {
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),
Expand Down Expand Up @@ -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:?}");
}
2 changes: 2 additions & 0 deletions crates/cellule-app/tests/host.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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,
Expand Down
50 changes: 40 additions & 10 deletions crates/cellule-app/tests/host/replicas.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand All @@ -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]);
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(&notified)
.filter(|session| Some(**session) != held_session)
.collect::<Vec<_>>(),
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.
Expand All @@ -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);
}
Expand Down
9 changes: 9 additions & 0 deletions crates/cellule-runtime/docs/runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -450,6 +450,15 @@ activation, never correctness.
<a id="ownership-renewal"></a>
## 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.
Expand Down
16 changes: 16 additions & 0 deletions crates/cellule-runtime/docs/rust-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
3 changes: 2 additions & 1 deletion crates/cellule-runtime/src/cell/actor/admission.rs
Original file line number Diff line number Diff line change
Expand Up @@ -98,8 +98,9 @@ pub(super) fn fence_admission(admission: &CellAdmission) {
admission.bytes.close();
}

pub(super) fn new_cell_admission() -> Arc<CellAdmission> {
pub(super) fn new_cell_admission(owner_fence: crate::control::OwnerFence) -> Arc<CellAdmission> {
Arc::new(CellAdmission {
owner_fence,
requests: Arc::new(Semaphore::new(CELL_REQUESTS)),
bytes: Arc::new(Semaphore::new(CELL_BYTES)),
draining: AtomicBool::new(false),
Expand Down
12 changes: 11 additions & 1 deletion crates/cellule-runtime/src/cell/actor/handle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ pub struct CellHandle {
}

pub(super) struct CellAdmission {
pub(super) owner_fence: crate::control::OwnerFence,
pub(super) requests: Arc<Semaphore>,
pub(super) bytes: Arc<Semaphore>,
pub(super) draining: AtomicBool,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions crates/cellule-runtime/src/cell/actor/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
4 changes: 3 additions & 1 deletion crates/cellule-runtime/src/cell/actor/tasks/activation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
2 changes: 2 additions & 0 deletions crates/cellule-runtime/src/client/local.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
Expand Down
Loading
Loading