diff --git a/crates/cellule-runtime/docs/runtime.md b/crates/cellule-runtime/docs/runtime.md index ca928a3..4c626cb 100644 --- a/crates/cellule-runtime/docs/runtime.md +++ b/crates/cellule-runtime/docs/runtime.md @@ -431,7 +431,7 @@ activation, never correctness. ## 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`. diff --git a/crates/cellule-runtime/src/cell/actor/lifecycle/scheduling.rs b/crates/cellule-runtime/src/cell/actor/lifecycle/scheduling.rs index 3cbb164..f268f55 100644 --- a/crates/cellule-runtime/src/cell/actor/lifecycle/scheduling.rs +++ b/crates/cellule-runtime/src/cell/actor/lifecycle/scheduling.rs @@ -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, - tasks: &mut JoinSet, - 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) { + 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, + tasks: &mut JoinSet, + 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, + } + }); + } } } @@ -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()); + } +} diff --git a/crates/cellule-runtime/src/cell/actor/task.rs b/crates/cellule-runtime/src/cell/actor/task.rs index 7de1229..2174cb4 100644 --- a/crates/cellule-runtime/src/cell/actor/task.rs +++ b/crates/cellule-runtime/src/cell/actor/task.rs @@ -17,6 +17,7 @@ pub(super) async fn run( let mut cells = HashMap::::new(); let mut transitioning = HashSet::::new(); let mut tasks = JoinSet::::new(); + let mut renewals = RenewalScheduler::new(); let mut next_generation = 0_u64; let mut shutdown = None::; let mut pressure = match PressureClassifier::new(800, 600, 1_000) { @@ -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, @@ -67,6 +69,9 @@ pub(super) async fn run( &mut movement, &mut movement_permits, ); + if renewed { + renewals.finished(); + } continue; } if tasks.is_empty() { @@ -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( @@ -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; }; @@ -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(