Skip to main content

lenso_kernel/
transition.rs

1//! A conservative atomic replacement at a proven single-lane quiescent point.
2
3use super::{
4    AppAdmission, AppReadyGate, BTreeMap, CancellationToken, Cell, DeactivationReason,
5    DriverControl, Duration, ExecutionAdapter, ExecutionAdapterCatalog, Future, FutureExt,
6    GenerationPreparationFailure, InvocationContext, ManagedResourceScope, ManagedTaskScope,
7    NativeApp, NativeAppRuntime, NativeEndpointSet, NativePluginGeneration, Rc, ResolvedAppPlan,
8    RuntimeFailure, await_with_context, ensure_context_active, oneshot,
9};
10use lenso_app_plan::PlanTransition;
11
12/// Physical retirement is distinct from snapshot commitment.
13#[derive(Clone, Debug, Eq, PartialEq)]
14pub enum TransitionRetirement {
15    /// Snapshot committed; the Driver owner still retains old generations.
16    Pending,
17    /// All replaced generations completed ordinary physical teardown.
18    Clean,
19    /// Commitment is final, but retirement failed; App admission is fenced.
20    Uncertain { error: RuntimeFailure },
21}
22
23/// Operator observation, including attempts whose caller has left.
24#[derive(Clone, Copy, Debug, Eq, PartialEq)]
25pub enum TransitionStatus {
26    /// No owned attempt or unresolved retirement remains.
27    Idle,
28    /// Candidate physical work or rollback is still owned; no rejection proof.
29    Preparing,
30    /// The successor is active and old generations are still retiring.
31    Retiring,
32    /// Shutdown or uncertain cleanup forbids another attempt or safe fallback.
33    Fenced,
34}
35
36/// Exact committed snapshot identities and affected generations.
37#[derive(Clone, Debug, Eq, PartialEq)]
38pub struct TransitionOutcome {
39    pub predecessor_digest: String,
40    pub successor_digest: String,
41    pub replaced_instances: Vec<String>,
42    pub retirement: TransitionRetirement,
43}
44
45pub(super) struct SnapshotCall(Rc<Cell<usize>>);
46
47impl SnapshotCall {
48    pub(super) fn new(runtime: &NativeAppRuntime) -> Self {
49        let count = runtime.transition_calls.clone();
50        count.set(
51            count
52                .get()
53                .checked_add(1)
54                .expect("bounded request call count"),
55        );
56        Self(count)
57    }
58}
59
60impl Drop for SnapshotCall {
61    fn drop(&mut self) {
62        self.0.set(self.0.get() - 1);
63    }
64}
65
66struct CancelCaller {
67    cancellation: CancellationToken,
68    cleanup: super::cleanup::StartupCleanupBudget,
69    driver: DriverControl,
70    deadline: Option<Duration>,
71}
72impl Drop for CancelCaller {
73    fn drop(&mut self) {
74        let now = (self.driver.now)();
75        self.cleanup.establish_at(
76            self.deadline
77                .filter(|deadline| now >= *deadline)
78                .unwrap_or(now),
79        );
80        self.cancellation.cancel();
81    }
82}
83
84struct Owner {
85    runtime: Rc<NativeAppRuntime>,
86    completed: bool,
87    started: bool,
88}
89impl Drop for Owner {
90    fn drop(&mut self) {
91        self.runtime.transition_pending.set(false);
92        if self.started && !self.completed {
93            self.runtime
94                .record_cleanup_failure(&invalid("transition owner was abandoned"));
95            self.runtime.begin_shutdown();
96        }
97    }
98}
99
100struct Staged {
101    key: String,
102    endpoints: NativeEndpointSet,
103    generation: NativePluginGeneration,
104    number: u64,
105    gate: AppReadyGate,
106    admission: AppAdmission,
107    adapter: Rc<dyn ExecutionAdapter>,
108}
109
110fn invalid(detail: impl Into<String>) -> RuntimeFailure {
111    RuntimeFailure::InvalidResolvedPlan {
112        detail: detail.into(),
113    }
114}
115fn busy() -> RuntimeFailure {
116    RuntimeFailure::ResourceExhausted {
117        capability: "lenso.plan-transition@1",
118        operation: "apply".into(),
119    }
120}
121
122impl NativeApp {
123    /// Observes transition ownership without guessing from a caller error.
124    pub fn transition_status(&self) -> TransitionStatus {
125        if self.runtime.shutdown_started.get() || self.runtime.cleanup_failure.borrow().is_some() {
126            TransitionStatus::Fenced
127        } else if !self.runtime.transition_pending.get() {
128            TransitionStatus::Idle
129        } else if self
130            .runtime
131            .last_transition
132            .borrow()
133            .as_ref()
134            .is_some_and(|outcome| outcome.retirement == TransitionRetirement::Pending)
135        {
136            TransitionStatus::Retiring
137        } else {
138            TransitionStatus::Preparing
139        }
140    }
141    /// Applies a fully resolved, interface-identical single-lane successor.
142    /// Busy lanes reject before commitment; no invocation observes mixed epochs.
143    /// Timeout/drop never abandons physical work. Inspect `last_transition` after
144    /// caller cancellation: a published commitment cannot be treated as rollback.
145    pub async fn apply_transition(
146        &self,
147        successor: ResolvedAppPlan,
148        transition: PlanTransition,
149        candidate_adapters: ExecutionAdapterCatalog,
150        context: InvocationContext,
151        cleanup_timeout: Duration,
152    ) -> Result<TransitionOutcome, RuntimeFailure> {
153        ensure_context_active(&self.runtime.driver, &context)?;
154        if self.runtime.shutdown_started.get() || self.runtime.cleanup_failure.borrow().is_some() {
155            return Err(RuntimeFailure::AdmissionClosed);
156        }
157        if self.runtime.transition_pending.replace(true) {
158            return Err(busy());
159        }
160        let mut owner = Owner {
161            runtime: self.runtime.clone(),
162            completed: false,
163            started: false,
164        };
165        let cancellation = CancellationToken::child(&context.cancellation());
166        let cleanup =
167            super::cleanup::StartupCleanupBudget::new(&self.runtime.driver, cleanup_timeout);
168        let caller = CancelCaller {
169            cancellation: cancellation.clone(),
170            cleanup: cleanup.clone(),
171            driver: self.runtime.driver.clone(),
172            deadline: context.deadline(),
173        };
174        let mut worker_context = context.clone();
175        worker_context.cancellation = cancellation;
176        let (publish, receive) = oneshot::channel();
177        let runtime = self.runtime.clone();
178        (self.runtime.driver.spawn_local)(Box::pin(async move {
179            owner.started = true;
180            let result = apply_owned(
181                &runtime,
182                successor,
183                transition,
184                candidate_adapters,
185                &worker_context,
186                &cleanup,
187            )
188            .await;
189            owner.completed = true;
190            drop(owner);
191            let _ = publish.send(result);
192        }))
193        .map_err(|error| invalid(format!("cannot schedule transition owner: {error}")))?;
194        let result = await_with_context(&self.runtime.driver, &context, receive)
195            .await?
196            .map_err(|_| invalid("transition owner ended without a result"))?;
197        drop(caller);
198        result
199    }
200}
201
202async fn retire(
203    runtime: &Rc<NativeAppRuntime>,
204    stages: Vec<(String, NativePluginGeneration, u64)>,
205    budget: super::cleanup::CleanupBudget,
206) -> Option<RuntimeFailure> {
207    let mut first = None;
208    for (key, generation, number) in stages.into_iter().rev() {
209        let error = super::supervision::cleanup_detached_generation_with_budget(
210            runtime,
211            &key,
212            generation,
213            DeactivationReason::SupervisionRestart,
214            number,
215            Some(budget.clone()),
216        )
217        .await;
218        if first.is_none() {
219            first = error;
220        }
221    }
222    if first.is_some() {
223        runtime.begin_shutdown();
224    }
225    first
226}
227
228#[allow(
229    clippy::too_many_lines,
230    reason = "one lane transaction stages, commits atomically and retires explicitly"
231)]
232async fn apply_owned(
233    runtime: &Rc<NativeAppRuntime>,
234    successor: ResolvedAppPlan,
235    transition: PlanTransition,
236    candidates: ExecutionAdapterCatalog,
237    context: &InvocationContext,
238    cleanup: &super::cleanup::StartupCleanupBudget,
239) -> Result<TransitionOutcome, RuntimeFailure> {
240    let previous = runtime.snapshot.borrow().clone();
241    let expected = PlanTransition::between(&previous, &successor)
242        .map_err(|error| invalid(error.to_string()))?;
243    if expected != transition {
244        return Err(invalid(
245            "transition does not match the active adjacent snapshots",
246        ));
247    }
248    // Validate every affected Adapter before recreating the first generation.
249    let mut selected = BTreeMap::new();
250    for key in transition.replaced_instances() {
251        let instance = successor
252            .plugin_instance(key)
253            .expect("validated replacement key");
254        let adapter = candidates
255            .adapter(instance.execution_class())
256            .ok_or_else(|| invalid(format!("missing candidate Execution Adapter for `{key}`")))?;
257        if !adapter
258            .supports_runtime_profile(instance.authoring_version(), instance.runtime_profile())
259            || !adapter.supports_plan_transition(instance)
260        {
261            return Err(invalid(format!(
262                "Adapter has no admitted stateless transition boundary for `{key}`"
263            )));
264        }
265        selected.insert(key.clone(), adapter);
266    }
267    let mut stages = Vec::new();
268    let preparation: Result<(), RuntimeFailure> = async {
269        for key in transition.replaced_instances() {
270            ensure_context_active(&runtime.driver, context)?;
271            if runtime.shutdown_started.get() {
272                return Err(RuntimeFailure::AdmissionClosed);
273            }
274            let instance = successor
275                .plugin_instance(key)
276                .expect("validated replacement key");
277            let adapter = selected[key].clone();
278            let prepared = adapter.recreate(&successor, key)?;
279            let (endpoints, lifecycle) = prepared.into_parts();
280            let number = runtime.supervision.borrow()[key]
281                .generation
282                .checked_add(1)
283                .ok_or_else(|| invalid("generation sequence exhausted"))?;
284            if let Err(error) = super::supervision::validate_native_endpoint_set(
285                key,
286                instance,
287                endpoints.request(),
288                endpoints.stream(),
289                endpoints.event(),
290            ) {
291                let generation = NativePluginGeneration {
292                    lifecycle,
293                    tasks: ManagedTaskScope::new_from_driver_control(&runtime.driver),
294                    resources: ManagedResourceScope::new(),
295                    stop_attempted: false,
296                    cleanup_timed_out: false,
297                    staged_admission: None,
298                };
299                let cleanup = super::supervision::cleanup_detached_generation_with_budget(
300                    runtime,
301                    key,
302                    generation,
303                    DeactivationReason::SupervisionRestart,
304                    number,
305                    Some(cleanup.establish()),
306                )
307                .await;
308                if cleanup.is_some() {
309                    runtime.begin_shutdown();
310                }
311                return Err(error);
312            }
313            let gate = AppReadyGate::new();
314            let admission = AppAdmission::new();
315            let tasks = ManagedTaskScope::new_from_driver_control(&runtime.driver);
316            let preparation = super::supervision::prepare_and_activate_generation_for_snapshot(
317                runtime,
318                &successor,
319                key,
320                lifecycle,
321                number,
322                gate.clone(),
323                admission.clone(),
324                true,
325                Some(tasks.clone()),
326                Some(cleanup.clone()),
327            );
328            futures::pin_mut!(preparation);
329            let mut cancelled = context.cancellation().cancelled();
330            let mut deadline = context.deadline().map_or_else(
331                || futures::future::pending().boxed_local(),
332                |deadline| (runtime.driver.sleep_until)(deadline),
333            );
334            // Signal cancellation but retain the actual lifecycle future until
335            // it returns. A timeout is never physical retirement evidence.
336            let generation = std::future::poll_fn(|cx| {
337                let _ = cancelled.as_mut().poll(cx);
338                let _ = deadline.as_mut().poll(cx);
339                if ensure_context_active(&runtime.driver, context).is_err()
340                    || runtime.shutdown_started.get()
341                {
342                    tasks.close();
343                }
344                preparation.as_mut().poll(cx)
345            })
346            .await
347            .map_err(|failure| match failure {
348                GenerationPreparationFailure::Lifecycle => {
349                    invalid("staged generation failed readiness")
350                }
351                GenerationPreparationFailure::Cleanup { primary } => {
352                    runtime.begin_shutdown();
353                    primary
354                }
355            })?;
356            stages.push(Staged {
357                key: key.clone(),
358                endpoints,
359                generation,
360                number,
361                gate,
362                admission,
363                adapter,
364            });
365        }
366        ensure_context_active(&runtime.driver, context)?;
367        if runtime.shutdown_started.get() {
368            return Err(RuntimeFailure::AdmissionClosed);
369        }
370        if runtime.transition_calls.get() != 0
371            || !runtime.executions.is_settled(None)
372            || runtime
373                .supervision
374                .borrow()
375                .values()
376                .any(|state| state.restarting)
377            || stages
378                .iter()
379                .any(|stage| runtime.plugins[&stage.key].generation.borrow().is_none())
380            || stages.iter().any(|stage| {
381                runtime.supervision.borrow()[&stage.key].generation != stage.number - 1
382            })
383            || stages
384                .iter()
385                .any(|stage| stage.generation.tasks.state.unreported_failure.get())
386        {
387            return Err(busy());
388        }
389        Ok(())
390    }
391    .await;
392    if let Err(error) = preparation {
393        let generations = stages
394            .into_iter()
395            .map(|stage| (stage.key, stage.generation, stage.number))
396            .collect();
397        let _ = retire(runtime, generations, cleanup.establish()).await;
398        return Err(error);
399    }
400    let mut outcome = TransitionOutcome {
401        predecessor_digest: transition.predecessor_digest().into(),
402        successor_digest: transition.successor_digest().into(),
403        replaced_instances: transition.replaced_instances().to_vec(),
404        retirement: TransitionRetirement::Pending,
405    };
406    let mut old = Vec::new();
407    // No await or user code: this complete binding/snapshot commit owns the lane.
408    let mut gates = Vec::new();
409    for stage in stages {
410        let Staged {
411            key,
412            endpoints,
413            generation,
414            number,
415            gate,
416            admission,
417            adapter,
418        } = stage;
419        let plugin = &runtime.plugins[&key];
420        let prior = plugin
421            .take_generation()
422            .expect("quiescent ready generation");
423        let previous_number = runtime.supervision.borrow()[&key].generation;
424        old.push((key.clone(), prior, previous_number));
425        super::kernel::attach_managed_task_failure_handler(runtime, &key, &generation.tasks);
426        plugin.install_generation(generation);
427        super::supervision::install_plugin_endpoints(
428            runtime,
429            &key,
430            endpoints.request().to_vec(),
431            endpoints.stream().to_vec(),
432            endpoints.event().to_vec(),
433            number,
434        );
435        runtime
436            .transition_adapters
437            .borrow_mut()
438            .insert(key.clone(), adapter);
439        let mut supervision = runtime.supervision.borrow_mut();
440        let state = supervision
441            .get_mut(&key)
442            .expect("validated Instance supervision state");
443        state.generation = number;
444        state.attempts.clear();
445        state.stable_since = Some((runtime.driver.now)());
446        gates.push((gate, admission));
447    }
448    runtime.snapshot.replace(Rc::new(successor));
449    runtime.last_transition.replace(Some(outcome.clone()));
450    for (gate, admission) in gates {
451        gate.open();
452        admission.open();
453    }
454    let error = retire(runtime, std::mem::take(&mut old), cleanup.establish()).await;
455    outcome.retirement = error.map_or(TransitionRetirement::Clean, |error| {
456        TransitionRetirement::Uncertain { error }
457    });
458    runtime.last_transition.replace(Some(outcome.clone()));
459    Ok(outcome)
460}