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
2 changes: 1 addition & 1 deletion crates/cellule-runtime/docs/runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -431,7 +431,7 @@ activation, never correctness.
<a id="ownership-renewal"></a>
## Renew and self-fence ownership

- **Scanner.** One node-level scanner renews owned Cells every three seconds. A mutation publication also advances owner progress.
- **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.
- **Scope.** Renewal changes owner liveness fields only. It preserves root, code, schema, and durable `next_due_ms`.
Expand Down
182 changes: 129 additions & 53 deletions crates/cellule-runtime/src/cell/actor/lifecycle/scheduling.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,65 +148,111 @@ pub(in crate::cell::actor) fn start_transfer_inspection(
});
}

pub(in crate::cell::actor) fn start_due_renewals(
pool: &SqlWorkerPool,
cells: &mut HashMap<CellId, ActiveCell>,
tasks: &mut JoinSet<TaskResult>,
node_lease: &RuntimeNodeLease,
) {
let active_renewals = cells.values().filter(|active| active.renewing()).count();
let mut available = MAX_RENEWALS_IN_FLIGHT.saturating_sub(active_renewals);
if available == 0 {
return;
}
let now = std::time::Instant::now();
for (cell, active) in cells.iter_mut() {
if available == 0 {
break;
}
if active
.publisher
.as_ref()
.is_none_or(|publisher| !publisher.renewal_due(now))
{
continue;
/// Due owners are scanned once per tick, then drained in deadline order as
/// renewal I/O completes. Rebuilding the pending set bounds stale entries
/// when Cells depart or change generation while all I/O slots are occupied.
pub(in crate::cell::actor) struct RenewalScheduler {
due: Vec<(std::time::Instant, [u8; 32], u64)>,
in_flight: usize,
}

impl RenewalScheduler {
pub(in crate::cell::actor) fn new() -> Self {
Self {
due: Vec::new(),
in_flight: 0,
}
match active.coordination.step(CoordinationInput::BeginRenewal {
queue_empty: active.queue.is_empty(),
publication_idle: active.coordination.publication_count() == 0,
lease_live: node_lease.check().is_ok(),
}) {
CoordinationDecision::Started => {}
CoordinationDecision::Fence => {
fence_active(active);
}

pub(in crate::cell::actor) fn scan_due(&mut self, cells: &HashMap<CellId, ActiveCell>) {
self.due.clear();
let now = std::time::Instant::now();
for (cell, active) in cells {
let Some(publisher) = active.publisher.as_ref() else {
continue;
};
if publisher.renewal_due(now) {
self.due
.push((publisher.renewal_at(), *cell.as_bytes(), active.generation));
}
_ => continue,
}
let Some(mut publisher) = active.publisher.take() else {
active
.coordination
.step(CoordinationInput::FinishRenewal { fenced: true });
continue;
};
let generation = active.generation;
let effect_id = active.begin_task(CoordinationEffect::Renewal);
available -= 1;
let cell = *cell;
let pool = pool.clone();
tasks.spawn(async move {
let result = publisher.renew().await;
if result.is_err() {
let _ = pool.fence(cell).await;
// Pop the oldest deadline first. Identity breaks ties independently of
// HashMap iteration order, so newly due Cells cannot overtake old debt.
self.order_due();
}

fn order_due(&mut self) {
self.due.sort_unstable_by(|left, right| right.cmp(left));
}

fn next_due(&mut self) -> Option<(CellId, u64)> {
self.due
.pop()
.map(|(_, bytes, generation)| (CellId::from_bytes(bytes), generation))
}

pub(in crate::cell::actor) fn finished(&mut self) {
self.in_flight = self.in_flight.saturating_sub(1);
}

pub(in crate::cell::actor) fn dispatch(
&mut self,
pool: &SqlWorkerPool,
cells: &mut HashMap<CellId, ActiveCell>,
tasks: &mut JoinSet<TaskResult>,
node_lease: &RuntimeNodeLease,
) {
let now = std::time::Instant::now();
while self.in_flight < MAX_RENEWALS_IN_FLIGHT {
let Some((cell, generation)) = self.next_due() else {
break;
};
let Some(active) = cells.get_mut(&cell) else {
continue;
};
if active.generation != generation
|| active
.publisher
.as_ref()
.is_none_or(|publisher| !publisher.renewal_due(now))
{
continue;
}
TaskResult::Renewed {
cell,
generation,
effect_id,
publisher: Box::new(publisher),
result,
match active.coordination.step(CoordinationInput::BeginRenewal {
queue_empty: active.queue.is_empty(),
publication_idle: active.coordination.publication_count() == 0,
lease_live: node_lease.check().is_ok(),
}) {
CoordinationDecision::Started => {}
CoordinationDecision::Fence => {
fence_active(active);
continue;
}
_ => continue,
}
});
let Some(mut publisher) = active.publisher.take() else {
active
.coordination
.step(CoordinationInput::FinishRenewal { fenced: true });
continue;
};
let effect_id = active.begin_task(CoordinationEffect::Renewal);
self.in_flight += 1;
let pool = pool.clone();
tasks.spawn(async move {
let result = publisher.renew().await;
if result.is_err() {
let _ = pool.fence(cell).await;
}
TaskResult::Renewed {
cell,
generation,
effect_id,
publisher: Box::new(publisher),
result,
}
});
}
}
}

Expand Down Expand Up @@ -329,3 +375,33 @@ pub(in crate::cell::actor) fn start_orphan_deactivate(
});
transitioning.insert(cell);
}

#[cfg(test)]
mod renewal_tests {
use super::*;

#[test]
fn due_owners_are_ordered_and_stale_entries_are_discarded_on_scan() {
let mut scheduler = RenewalScheduler::new();
let now = std::time::Instant::now();
for index in (0_u64..1_000).rev() {
let mut bytes = [0; 32];
bytes[..8].copy_from_slice(&index.to_be_bytes());
scheduler
.due
.push((now + std::time::Duration::from_millis(index), bytes, index));
}
scheduler.order_due();
for index in 0_u64..1_000 {
let mut bytes = [0; 32];
bytes[..8].copy_from_slice(&index.to_be_bytes());
assert_eq!(
scheduler.next_due(),
Some((CellId::from_bytes(bytes), index))
);
}
scheduler.due.push((now, [7; 32], 1));
scheduler.scan_due(&HashMap::new());
assert!(scheduler.due.is_empty());
}
}
22 changes: 20 additions & 2 deletions crates/cellule-runtime/src/cell/actor/task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ pub(super) async fn run(
let mut cells = HashMap::<CellId, ActiveCell>::new();
let mut transitioning = HashSet::<CellId>::new();
let mut tasks = JoinSet::<TaskResult>::new();
let mut renewals = RenewalScheduler::new();
let mut next_generation = 0_u64;
let mut shutdown = None::<ShutdownState>;
let mut pressure = match PressureClassifier::new(800, 600, 1_000) {
Expand Down Expand Up @@ -54,6 +55,7 @@ pub(super) async fn run(
finish_shutdown(&mut shutdown);
return;
};
let renewed = matches!(result, TaskResult::Renewed { .. });
super::tasks::handle_task(
result,
&pool,
Expand All @@ -67,6 +69,9 @@ pub(super) async fn run(
&mut movement,
&mut movement_permits,
);
if renewed {
renewals.finished();
}
continue;
}
if tasks.is_empty() {
Expand All @@ -89,7 +94,8 @@ pub(super) async fn run(
handle_message(message, &mut receiver, &pool, &mut cells, &mut transitioning, &mut tasks, &mut shutdown, &node_lease, &telemetry, &mut pressure, &mut movement, &mut movement_permits, &mut next_generation);
}
_ = renewal_tick.tick() => {
start_due_renewals(&pool, &mut cells, &mut tasks, &node_lease);
renewals.scan_due(&cells);
renewals.dispatch(&pool, &mut cells, &mut tasks, &node_lease);
}
_ = hydration_tick.tick() => {
start_background_hydration(
Expand Down Expand Up @@ -132,7 +138,11 @@ pub(super) async fn run(
}
while let Some(result) = tasks.join_next().await {
let Ok(result) = result else { return; };
let renewed = matches!(result, TaskResult::Renewed { .. });
super::tasks::handle_task(result, &pool, &mut cells, &mut transitioning, &mut tasks, &mut shutdown, &node_lease, &unpublished_node_log_bytes, &publications, &mut movement, &mut movement_permits);
if renewed {
renewals.finished();
}
}
break;
};
Expand All @@ -142,10 +152,18 @@ pub(super) async fn run(
let Some(Ok(result)) = result else {
return;
};
let renewed = matches!(result, TaskResult::Renewed { .. });
super::tasks::handle_task(result, &pool, &mut cells, &mut transitioning, &mut tasks, &mut shutdown, &node_lease, &unpublished_node_log_bytes, &publications, &mut movement, &mut movement_permits);
if renewed {
renewals.finished();
if !shutdown.as_ref().is_some_and(|state| state.draining) {
renewals.dispatch(&pool, &mut cells, &mut tasks, &node_lease);
}
}
}
_ = renewal_tick.tick() => {
start_due_renewals(&pool, &mut cells, &mut tasks, &node_lease);
renewals.scan_due(&cells);
renewals.dispatch(&pool, &mut cells, &mut tasks, &node_lease);
}
_ = hydration_tick.tick() => {
start_background_hydration(
Expand Down
Loading