use super::{
AppAdmission, AppReadyGate, BTreeMap, CancellationToken, Cell, DeactivationReason,
DriverControl, Duration, ExecutionAdapter, ExecutionAdapterCatalog, Future, FutureExt,
GenerationPreparationFailure, InvocationContext, ManagedResourceScope, ManagedTaskScope,
NativeApp, NativeAppRuntime, NativeEndpointSet, NativePluginGeneration, Rc, ResolvedAppPlan,
RuntimeFailure, await_with_context, ensure_context_active, oneshot,
};
use lenso_app_plan::PlanTransition;
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum TransitionRetirement {
Pending,
Clean,
Uncertain { error: RuntimeFailure },
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum TransitionStatus {
Idle,
Preparing,
Retiring,
Fenced,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TransitionOutcome {
pub predecessor_digest: String,
pub successor_digest: String,
pub replaced_instances: Vec<String>,
pub retirement: TransitionRetirement,
}
pub(super) struct SnapshotCall(Rc<Cell<usize>>);
impl SnapshotCall {
pub(super) fn new(runtime: &NativeAppRuntime) -> Self {
let count = runtime.transition_calls.clone();
count.set(
count
.get()
.checked_add(1)
.expect("bounded request call count"),
);
Self(count)
}
}
impl Drop for SnapshotCall {
fn drop(&mut self) {
self.0.set(self.0.get() - 1);
}
}
struct CancelCaller {
cancellation: CancellationToken,
cleanup: super::cleanup::StartupCleanupBudget,
driver: DriverControl,
deadline: Option<Duration>,
}
impl Drop for CancelCaller {
fn drop(&mut self) {
let now = (self.driver.now)();
self.cleanup.establish_at(
self.deadline
.filter(|deadline| now >= *deadline)
.unwrap_or(now),
);
self.cancellation.cancel();
}
}
struct Owner {
runtime: Rc<NativeAppRuntime>,
completed: bool,
started: bool,
}
impl Drop for Owner {
fn drop(&mut self) {
self.runtime.transition_pending.set(false);
if self.started && !self.completed {
self.runtime
.record_cleanup_failure(&invalid("transition owner was abandoned"));
self.runtime.begin_shutdown();
}
}
}
struct Staged {
key: String,
endpoints: NativeEndpointSet,
generation: NativePluginGeneration,
number: u64,
gate: AppReadyGate,
admission: AppAdmission,
adapter: Rc<dyn ExecutionAdapter>,
}
fn invalid(detail: impl Into<String>) -> RuntimeFailure {
RuntimeFailure::InvalidResolvedPlan {
detail: detail.into(),
}
}
fn busy() -> RuntimeFailure {
RuntimeFailure::ResourceExhausted {
capability: "lenso.plan-transition@1",
operation: "apply".into(),
}
}
impl NativeApp {
pub fn transition_status(&self) -> TransitionStatus {
if self.runtime.shutdown_started.get() || self.runtime.cleanup_failure.borrow().is_some() {
TransitionStatus::Fenced
} else if !self.runtime.transition_pending.get() {
TransitionStatus::Idle
} else if self
.runtime
.last_transition
.borrow()
.as_ref()
.is_some_and(|outcome| outcome.retirement == TransitionRetirement::Pending)
{
TransitionStatus::Retiring
} else {
TransitionStatus::Preparing
}
}
pub async fn apply_transition(
&self,
successor: ResolvedAppPlan,
transition: PlanTransition,
candidate_adapters: ExecutionAdapterCatalog,
context: InvocationContext,
cleanup_timeout: Duration,
) -> Result<TransitionOutcome, RuntimeFailure> {
ensure_context_active(&self.runtime.driver, &context)?;
if self.runtime.shutdown_started.get() || self.runtime.cleanup_failure.borrow().is_some() {
return Err(RuntimeFailure::AdmissionClosed);
}
if self.runtime.transition_pending.replace(true) {
return Err(busy());
}
let mut owner = Owner {
runtime: self.runtime.clone(),
completed: false,
started: false,
};
let cancellation = CancellationToken::child(&context.cancellation());
let cleanup =
super::cleanup::StartupCleanupBudget::new(&self.runtime.driver, cleanup_timeout);
let caller = CancelCaller {
cancellation: cancellation.clone(),
cleanup: cleanup.clone(),
driver: self.runtime.driver.clone(),
deadline: context.deadline(),
};
let mut worker_context = context.clone();
worker_context.cancellation = cancellation;
let (publish, receive) = oneshot::channel();
let runtime = self.runtime.clone();
(self.runtime.driver.spawn_local)(Box::pin(async move {
owner.started = true;
let result = apply_owned(
&runtime,
successor,
transition,
candidate_adapters,
&worker_context,
&cleanup,
)
.await;
owner.completed = true;
drop(owner);
let _ = publish.send(result);
}))
.map_err(|error| invalid(format!("cannot schedule transition owner: {error}")))?;
let result = await_with_context(&self.runtime.driver, &context, receive)
.await?
.map_err(|_| invalid("transition owner ended without a result"))?;
drop(caller);
result
}
}
async fn retire(
runtime: &Rc<NativeAppRuntime>,
stages: Vec<(String, NativePluginGeneration, u64)>,
budget: super::cleanup::CleanupBudget,
) -> Option<RuntimeFailure> {
let mut first = None;
for (key, generation, number) in stages.into_iter().rev() {
let error = super::supervision::cleanup_detached_generation_with_budget(
runtime,
&key,
generation,
DeactivationReason::SupervisionRestart,
number,
Some(budget.clone()),
)
.await;
if first.is_none() {
first = error;
}
}
if first.is_some() {
runtime.begin_shutdown();
}
first
}
#[allow(
clippy::too_many_lines,
reason = "one lane transaction stages, commits atomically and retires explicitly"
)]
async fn apply_owned(
runtime: &Rc<NativeAppRuntime>,
successor: ResolvedAppPlan,
transition: PlanTransition,
candidates: ExecutionAdapterCatalog,
context: &InvocationContext,
cleanup: &super::cleanup::StartupCleanupBudget,
) -> Result<TransitionOutcome, RuntimeFailure> {
let previous = runtime.snapshot.borrow().clone();
let expected = PlanTransition::between(&previous, &successor)
.map_err(|error| invalid(error.to_string()))?;
if expected != transition {
return Err(invalid(
"transition does not match the active adjacent snapshots",
));
}
let mut selected = BTreeMap::new();
for key in transition.replaced_instances() {
let instance = successor
.plugin_instance(key)
.expect("validated replacement key");
let adapter = candidates
.adapter(instance.execution_class())
.ok_or_else(|| invalid(format!("missing candidate Execution Adapter for `{key}`")))?;
if !adapter
.supports_runtime_profile(instance.authoring_version(), instance.runtime_profile())
|| !adapter.supports_plan_transition(instance)
{
return Err(invalid(format!(
"Adapter has no admitted stateless transition boundary for `{key}`"
)));
}
selected.insert(key.clone(), adapter);
}
let mut stages = Vec::new();
let preparation: Result<(), RuntimeFailure> = async {
for key in transition.replaced_instances() {
ensure_context_active(&runtime.driver, context)?;
if runtime.shutdown_started.get() {
return Err(RuntimeFailure::AdmissionClosed);
}
let instance = successor
.plugin_instance(key)
.expect("validated replacement key");
let adapter = selected[key].clone();
let prepared = adapter.recreate(&successor, key)?;
let (endpoints, lifecycle) = prepared.into_parts();
let number = runtime.supervision.borrow()[key]
.generation
.checked_add(1)
.ok_or_else(|| invalid("generation sequence exhausted"))?;
if let Err(error) = super::supervision::validate_native_endpoint_set(
key,
instance,
endpoints.request(),
endpoints.stream(),
endpoints.event(),
) {
let generation = NativePluginGeneration {
lifecycle,
tasks: ManagedTaskScope::new_from_driver_control(&runtime.driver),
resources: ManagedResourceScope::new(),
stop_attempted: false,
cleanup_timed_out: false,
staged_admission: None,
};
let cleanup = super::supervision::cleanup_detached_generation_with_budget(
runtime,
key,
generation,
DeactivationReason::SupervisionRestart,
number,
Some(cleanup.establish()),
)
.await;
if cleanup.is_some() {
runtime.begin_shutdown();
}
return Err(error);
}
let gate = AppReadyGate::new();
let admission = AppAdmission::new();
let tasks = ManagedTaskScope::new_from_driver_control(&runtime.driver);
let preparation = super::supervision::prepare_and_activate_generation_for_snapshot(
runtime,
&successor,
key,
lifecycle,
number,
gate.clone(),
admission.clone(),
true,
Some(tasks.clone()),
Some(cleanup.clone()),
);
futures::pin_mut!(preparation);
let mut cancelled = context.cancellation().cancelled();
let mut deadline = context.deadline().map_or_else(
|| futures::future::pending().boxed_local(),
|deadline| (runtime.driver.sleep_until)(deadline),
);
let generation = std::future::poll_fn(|cx| {
let _ = cancelled.as_mut().poll(cx);
let _ = deadline.as_mut().poll(cx);
if ensure_context_active(&runtime.driver, context).is_err()
|| runtime.shutdown_started.get()
{
tasks.close();
}
preparation.as_mut().poll(cx)
})
.await
.map_err(|failure| match failure {
GenerationPreparationFailure::Lifecycle => {
invalid("staged generation failed readiness")
}
GenerationPreparationFailure::Cleanup { primary } => {
runtime.begin_shutdown();
primary
}
})?;
stages.push(Staged {
key: key.clone(),
endpoints,
generation,
number,
gate,
admission,
adapter,
});
}
ensure_context_active(&runtime.driver, context)?;
if runtime.shutdown_started.get() {
return Err(RuntimeFailure::AdmissionClosed);
}
if runtime.transition_calls.get() != 0
|| !runtime.executions.is_settled(None)
|| runtime
.supervision
.borrow()
.values()
.any(|state| state.restarting)
|| stages
.iter()
.any(|stage| runtime.plugins[&stage.key].generation.borrow().is_none())
|| stages.iter().any(|stage| {
runtime.supervision.borrow()[&stage.key].generation != stage.number - 1
})
|| stages
.iter()
.any(|stage| stage.generation.tasks.state.unreported_failure.get())
{
return Err(busy());
}
Ok(())
}
.await;
if let Err(error) = preparation {
let generations = stages
.into_iter()
.map(|stage| (stage.key, stage.generation, stage.number))
.collect();
let _ = retire(runtime, generations, cleanup.establish()).await;
return Err(error);
}
let mut outcome = TransitionOutcome {
predecessor_digest: transition.predecessor_digest().into(),
successor_digest: transition.successor_digest().into(),
replaced_instances: transition.replaced_instances().to_vec(),
retirement: TransitionRetirement::Pending,
};
let mut old = Vec::new();
let mut gates = Vec::new();
for stage in stages {
let Staged {
key,
endpoints,
generation,
number,
gate,
admission,
adapter,
} = stage;
let plugin = &runtime.plugins[&key];
let prior = plugin
.take_generation()
.expect("quiescent ready generation");
let previous_number = runtime.supervision.borrow()[&key].generation;
old.push((key.clone(), prior, previous_number));
super::kernel::attach_managed_task_failure_handler(runtime, &key, &generation.tasks);
plugin.install_generation(generation);
super::supervision::install_plugin_endpoints(
runtime,
&key,
endpoints.request().to_vec(),
endpoints.stream().to_vec(),
endpoints.event().to_vec(),
number,
);
runtime
.transition_adapters
.borrow_mut()
.insert(key.clone(), adapter);
let mut supervision = runtime.supervision.borrow_mut();
let state = supervision
.get_mut(&key)
.expect("validated Instance supervision state");
state.generation = number;
state.attempts.clear();
state.stable_since = Some((runtime.driver.now)());
gates.push((gate, admission));
}
runtime.snapshot.replace(Rc::new(successor));
runtime.last_transition.replace(Some(outcome.clone()));
for (gate, admission) in gates {
gate.open();
admission.open();
}
let error = retire(runtime, std::mem::take(&mut old), cleanup.establish()).await;
outcome.retirement = error.map_or(TransitionRetirement::Clean, |error| {
TransitionRetirement::Uncertain { error }
});
runtime.last_transition.replace(Some(outcome.clone()));
Ok(outcome)
}