1use 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#[derive(Clone, Debug, Eq, PartialEq)]
14pub enum TransitionRetirement {
15 Pending,
17 Clean,
19 Uncertain { error: RuntimeFailure },
21}
22
23#[derive(Clone, Copy, Debug, Eq, PartialEq)]
25pub enum TransitionStatus {
26 Idle,
28 Preparing,
30 Retiring,
32 Fenced,
34}
35
36#[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 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 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 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 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 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}