use std::{any::Any, cell::RefCell, collections::BTreeMap, rc::Rc, time::Duration};
use futures::{FutureExt, channel::oneshot, future::LocalBoxFuture};
use lenso_app_plan::{
CapabilityBinding, CapabilityEndpointPlan, CapabilityRequirementPlan, PlanTransition,
PluginInstancePlan, ResolvedAppPlan,
};
use lenso_kernel::{
CancellationToken, DeactivateContext, DeterministicDriver, ExecutionAdapterCatalog,
ExecutionLease, InvocationContext, Kernel, NativeExecutionAdapter, NativeRequestEndpoint,
PluginFuture, PluginLifecycle, PrepareContext, PreparedBinding, PreparedNativeApp,
PreparedNativePlugin, RequestCapability, RuntimeDriver, RuntimeFailure, TransitionRetirement,
};
#[derive(Debug)]
struct Echo;
impl RequestCapability for Echo {
type Request = String;
type Response = String;
type DomainError = String;
const ID: &'static str = "test.transition@1";
const DESCRIPTOR_VERSION: &'static str = "1.0.0";
}
#[derive(Debug, Default)]
struct Probe {
recreated: RefCell<Vec<String>>,
prepared: RefCell<Vec<(String, String)>>,
stopped: RefCell<Vec<(String, String)>>,
finish: RefCell<Vec<oneshot::Sender<()>>>,
lease: RefCell<Option<ExecutionLease>>,
}
#[derive(Debug)]
struct Endpoint {
configuration: String,
probe: Rc<Probe>,
}
impl NativeRequestEndpoint for Endpoint {
fn capability_id(&self) -> &'static str {
Echo::ID
}
fn descriptor_version(&self) -> &'static str {
Echo::DESCRIPTOR_VERSION
}
fn operations(&self) -> &'static [&'static str] {
&["read"]
}
fn invoke(
&self,
_: &str,
request: Box<dyn Any>,
context: InvocationContext,
) -> LocalBoxFuture<'static, Result<Result<Box<dyn Any>, Box<dyn Any>>, RuntimeFailure>> {
let request = *request.downcast::<String>().unwrap();
let configuration = self.configuration.clone();
let probe = self.probe.clone();
Box::pin(async move {
if request == "wait" {
let (send, receive) = oneshot::channel();
probe.finish.borrow_mut().push(send);
let _ = receive.await;
}
if request == "detached" {
probe.lease.replace(Some(context.retain_execution()?));
}
Ok(Ok(Box::new(configuration) as Box<dyn Any>))
})
}
}
#[derive(Debug)]
struct Lifecycle {
key: String,
config: String,
probe: Rc<Probe>,
}
impl PluginLifecycle for Lifecycle {
fn prepare(&self, context: PrepareContext) -> PluginFuture {
self.probe
.prepared
.borrow_mut()
.push((self.key.clone(), context.configuration().to_owned()));
let config = self.config.clone();
let probe = self.probe.clone();
Box::pin(async move {
if config == "slow" {
let (send, receive) = oneshot::channel();
probe.finish.borrow_mut().push(send);
let _ = receive.await;
}
if config == "fail" {
Err(RuntimeFailure::Internal {
detail: "fixture prepare failure".into(),
})
} else {
Ok(())
}
})
}
fn deactivate(&self, _: DeactivateContext) -> PluginFuture {
self.probe
.stopped
.borrow_mut()
.push((self.key.clone(), self.config.clone()));
let fail = self.config == "stop-error";
Box::pin(async move {
if fail {
Err(RuntimeFailure::Internal {
detail: "fixture stop failure".into(),
})
} else {
Ok(())
}
})
}
}
#[derive(Debug)]
struct Adapter(Rc<Probe>);
impl Adapter {
fn generation(&self, instance: &PluginInstancePlan) -> PreparedNativePlugin {
let endpoints: Vec<Rc<dyn NativeRequestEndpoint>> =
if instance.provided_capabilities().is_empty() {
vec![]
} else {
vec![Rc::new(Endpoint {
configuration: instance.configuration().into(),
probe: self.0.clone(),
})]
};
PreparedNativePlugin::new(
endpoints,
Lifecycle {
key: instance.instance_key().into(),
config: instance.configuration().into(),
probe: self.0.clone(),
},
)
}
}
impl NativeExecutionAdapter for Adapter {
fn supports_runtime_profile(&self, version: u32, profile: &str) -> bool {
version == 2 && profile == "lenso.native-authoring@2"
}
fn prepare(&self, plan: &ResolvedAppPlan) -> Result<PreparedNativeApp, RuntimeFailure> {
let generations: BTreeMap<_, _> = plan
.plugin_instances()
.iter()
.map(|instance| (instance.instance_key().into(), self.generation(instance)))
.collect();
let bindings = plan
.capability_bindings()
.iter()
.map(|binding| {
PreparedBinding::new(
binding.consumer_instance(),
binding.provider_instance(),
generations[binding.provider_instance()].endpoints()[0].clone(),
)
.with_requirement_id(binding.requirement_id())
})
.collect();
Ok(PreparedNativeApp::new(bindings, generations))
}
fn recreate(
&self,
plan: &ResolvedAppPlan,
key: &str,
) -> Result<PreparedNativePlugin, RuntimeFailure> {
self.0.recreated.borrow_mut().push(key.into());
Ok(self.generation(plan.plugin_instance(key).unwrap()))
}
}
fn plan(configuration: &str) -> ResolvedAppPlan {
let instance = |key: &str| {
PluginInstancePlan::new(key, "test.plugin").with_authoring(2, "lenso.native-authoring@2")
};
ResolvedAppPlan::new(
vec![
instance("provider-a")
.with_configuration(configuration)
.with_capability(CapabilityEndpointPlan::new(
Echo::ID,
Echo::DESCRIPTOR_VERSION,
["read"],
)),
instance("provider-b")
.with_configuration("untouched")
.with_capability(CapabilityEndpointPlan::new(
Echo::ID,
Echo::DESCRIPTOR_VERSION,
["read"],
)),
instance("client-a").with_requirement(
CapabilityRequirementPlan::one(Echo::ID, Echo::DESCRIPTOR_VERSION)
.with_requirement_id("echo"),
),
instance("client-b").with_requirement(
CapabilityRequirementPlan::one(Echo::ID, Echo::DESCRIPTOR_VERSION)
.with_requirement_id("echo"),
),
],
vec![
CapabilityBinding::new("client-a", Echo::ID, Echo::DESCRIPTOR_VERSION, "provider-a")
.with_requirement_id("echo"),
CapabilityBinding::new("client-b", Echo::ID, Echo::DESCRIPTOR_VERSION, "provider-b")
.with_requirement_id("echo"),
],
)
}
fn setup(config: &str) -> (DeterministicDriver, lenso_kernel::NativeApp, Rc<Probe>) {
let driver = DeterministicDriver::new();
let probe = Rc::new(Probe::default());
let app = driver
.run(Kernel::start(
plan(config),
driver.clone(),
ExecutionAdapterCatalog::single(Adapter(probe.clone())),
))
.unwrap();
(driver, app, probe)
}
fn candidate(next: &ResolvedAppPlan, probe: &Rc<Probe>) -> ExecutionAdapterCatalog {
ExecutionAdapterCatalog::single(Adapter(probe.clone()))
.with_stateless_transitions(next, &["provider-a".into()])
.unwrap()
}
fn context() -> InvocationContext {
InvocationContext::new(1, None, CancellationToken::new())
}
#[test]
fn replacement_updates_stable_handle_and_only_affected_generation() {
let (driver, app, probe) = setup("old");
let handle = app.handle::<Echo>("client-a").unwrap();
let next = plan("new");
let transition = PlanTransition::between(&app.plan_snapshot(), &next).unwrap();
let outcome = driver
.run(app.apply_transition(
next.clone(),
transition,
candidate(&next, &probe),
context(),
Duration::from_secs(1),
))
.unwrap();
assert_eq!(outcome.retirement, TransitionRetirement::Clean);
assert_eq!(app.plan_snapshot(), next);
assert_eq!(app.last_transition(), Some(outcome));
assert_eq!(*probe.recreated.borrow(), ["provider-a"]);
assert!(
probe
.prepared
.borrow()
.contains(&("provider-a".into(), "new".into()))
);
assert_eq!(
driver
.run(handle.invoke("read", String::new()))
.unwrap()
.unwrap(),
"new"
);
assert_eq!(
driver
.run(app.invoke::<Echo>("client-b", "read", String::new()))
.unwrap()
.unwrap(),
"untouched"
);
assert_eq!(app.plugin_generation("provider-a"), Some(2));
assert_eq!(app.plugin_generation("provider-b"), Some(1));
assert_eq!(
driver.run(app.shutdown(Duration::from_secs(1))),
lenso_kernel::ShutdownOutcome::Clean
);
}
#[test]
fn failed_readiness_keeps_predecessor_and_retires_candidate() {
let (driver, app, probe) = setup("old");
let next = plan("fail");
assert!(
driver
.run(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1)
))
.is_err()
);
assert_eq!(app.plan_snapshot(), plan("old"));
assert!(app.last_transition().is_none());
assert!(
probe
.stopped
.borrow()
.contains(&("provider-a".into(), "fail".into()))
);
assert_eq!(
driver
.run(app.invoke::<Echo>("client-a", "read", String::new()))
.unwrap()
.unwrap(),
"old"
);
}
#[test]
fn inflight_request_rejects_commit_and_finishes_against_old_snapshot() {
let (driver, app, probe) = setup("old");
let next = plan("new");
driver.run(async {
let call = app.invoke::<Echo>("client-a", "read", "wait".into());
futures::pin_mut!(call);
assert!(call.as_mut().now_or_never().is_none());
let error = app
.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1),
)
.await
.unwrap_err();
assert!(matches!(error, RuntimeFailure::ResourceExhausted { .. }));
assert_eq!(app.plan_snapshot(), plan("old"));
probe.finish.borrow_mut().pop().unwrap().send(()).unwrap();
assert_eq!(call.await.unwrap().unwrap(), "old");
});
}
#[test]
fn retained_physical_execution_blocks_commit_after_caller_reply() {
let (driver, app, probe) = setup("old");
driver
.run(app.invoke::<Echo>("client-a", "read", "detached".into()))
.unwrap()
.unwrap();
let next = plan("new");
assert!(matches!(
driver.run(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1)
)),
Err(RuntimeFailure::ResourceExhausted { .. })
));
assert_eq!(app.plan_snapshot(), plan("old"));
probe.lease.borrow_mut().take().unwrap().settle();
driver
.run(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1),
))
.unwrap();
}
#[test]
fn stale_or_forged_receipt_and_missing_opt_in_fail_before_recreation() {
let (driver, app, probe) = setup("old");
let next = plan("new");
let stale = PlanTransition::between(&plan("another"), &next).unwrap();
assert!(
driver
.run(app.apply_transition(
next.clone(),
stale,
candidate(&next, &probe),
context(),
Duration::from_secs(1)
))
.is_err()
);
let mut forged =
serde_json::to_value(PlanTransition::between(&plan("old"), &next).unwrap()).unwrap();
forged["replaced_instances"] = serde_json::json!(["provider-b"]);
assert!(
driver
.run(app.apply_transition(
next.clone(),
serde_json::from_value(forged).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1)
))
.is_err()
);
assert!(
driver
.run(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
ExecutionAdapterCatalog::single(Adapter(probe.clone())),
context(),
Duration::from_secs(1)
))
.is_err()
);
assert!(probe.recreated.borrow().is_empty());
}
#[test]
fn dropped_caller_keeps_driver_owned_rollback_and_shutdown_waits() {
let (driver, app, probe) = setup("old");
let next = plan("slow");
driver.run(async {
let mut transition = Box::pin(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1),
));
assert!(transition.as_mut().now_or_never().is_none());
driver.yield_now().await;
drop(transition);
assert_eq!(
app.transition_status(),
lenso_kernel::TransitionStatus::Preparing
);
assert_eq!(app.plan_snapshot(), plan("old"));
probe.finish.borrow_mut().pop().unwrap().send(()).unwrap();
driver.yield_now().await;
driver.yield_now().await;
assert!(app.last_transition().is_none());
assert_eq!(
app.transition_status(),
lenso_kernel::TransitionStatus::Idle
);
assert_eq!(
app.shutdown(Duration::from_secs(1)).await,
lenso_kernel::ShutdownOutcome::Clean
);
});
}
#[test]
fn retirement_failure_preserves_committed_snapshot_and_fences_admission() {
let (driver, app, probe) = setup("stop-error");
let next = plan("new");
let result = driver
.run(app.apply_transition(
next.clone(),
PlanTransition::between(&plan("stop-error"), &next).unwrap(),
candidate(&next, &probe),
context(),
Duration::from_secs(1),
))
.unwrap();
assert!(matches!(
result.retirement,
TransitionRetirement::Uncertain { .. }
));
assert_eq!(app.plan_snapshot(), next);
assert!(!app.is_accepting());
assert!(matches!(
driver.run(app.shutdown(Duration::from_secs(1))),
lenso_kernel::ShutdownOutcome::RuntimeFailure { .. }
));
}
#[test]
fn late_constructor_cannot_reset_cancelled_attempts_cleanup_budget() {
let (driver, app, probe) = setup("old");
let next = plan("slow");
let cancellation = CancellationToken::new();
driver.run(async {
let attempt = app.apply_transition(
next.clone(),
PlanTransition::between(&plan("old"), &next).unwrap(),
candidate(&next, &probe),
InvocationContext::new(9, None, cancellation.clone()),
Duration::from_secs(1),
);
futures::pin_mut!(attempt);
assert!(attempt.as_mut().now_or_never().is_none());
driver.yield_now().await;
cancellation.cancel();
assert!(matches!(
attempt.await,
Err(RuntimeFailure::Cancelled { request_id: 9 })
));
assert_eq!(
app.transition_status(),
lenso_kernel::TransitionStatus::Preparing
);
driver.advance(Duration::from_secs(2));
probe.finish.borrow_mut().pop().unwrap().send(()).unwrap();
driver.yield_now().await;
driver.yield_now().await;
assert_eq!(app.plan_snapshot(), plan("old"));
assert!(app.last_transition().is_none());
assert_eq!(
app.transition_status(),
lenso_kernel::TransitionStatus::Fenced
);
assert!(
!probe
.stopped
.borrow()
.contains(&("provider-a".into(), "slow".into()))
);
assert!(matches!(
app.shutdown(Duration::from_secs(1)).await,
lenso_kernel::ShutdownOutcome::RuntimeFailure { .. }
));
});
}