#![allow(clippy::expect_used, clippy::unwrap_used)]
mod common;
use std::collections::HashMap;
use std::num::NonZeroUsize;
use std::sync::mpsc::{self, Sender};
use std::sync::{Arc, Barrier, Condvar, Mutex, OnceLock};
use std::thread;
use std::time::Duration;
use beamr::atom::{Atom, AtomTable};
use beamr::module::ModuleRegistry;
use beamr::native::{BifRegistryImpl, Capability, ProcessContext};
use beamr::scheduler::{Scheduler, SchedulerConfig, SchedulerServices};
use beamr::term::Term;
use frame_core::component::{
ChildArgument, ChildSpec, ComponentId, ComponentMeta, ServiceCapability, SupervisionPolicy,
};
use frame_core::error::{RegistryError, TombstoneKind};
use frame_core::event::{LifecycleEventKind, LifecycleState};
use frame_core::registry::{ComponentRegistry, RemoveOutcome, StopOutcome};
use frame_core::supervision::LifecycleConfig;
use common::ExitWitness;
const WORKER: &[u8] = include_bytes!("fixtures/spike_worker_v1.beam");
const WORKER_B: &[u8] = include_bytes!("fixtures/spike_worker_b.beam");
const FFI: &[u8] = include_bytes!("fixtures/spike_worker_ffi.beam");
const STOP_ORDER_WORKER: &[u8] = include_bytes!("fixtures/stop_order_worker.beam");
const OPERATION_DEADLINE: Duration = Duration::from_secs(3);
const CONCURRENT_STOP_OBSERVATION: Duration = Duration::from_millis(250);
const STOP_RECORDER_TOKEN: i64 = 7_041;
const FIRST_STOP_TAG: i64 = 10;
#[derive(Clone)]
struct StopRecorder {
sender: Sender<i64>,
first_child_release: Arc<(Mutex<bool>, Condvar)>,
}
static STOP_RECORDERS: OnceLock<Mutex<HashMap<i64, StopRecorder>>> = OnceLock::new();
fn scheduler() -> Arc<Scheduler> {
Arc::new(
frame_core::composition::compose_scheduler(
SchedulerConfig {
thread_count: Some(2),
..SchedulerConfig::default()
},
SchedulerServices::minimal(),
Arc::new(ModuleRegistry::new()),
)
.expect("published beamr scheduler starts"),
)
}
fn scheduler_with_stop_recorder() -> Arc<Scheduler> {
let atom_table = Arc::new(AtomTable::with_common_atoms());
let bif_registry = Arc::new(BifRegistryImpl::new());
beamr::native::bifs::register_gate1_bifs(&bif_registry, &atom_table)
.expect("gate-1 BIF registration succeeds");
bif_registry
.register(
atom_table.intern("frame_test"),
atom_table.intern("record_stop"),
2,
record_stop,
Capability::Pure,
)
.expect("test stop recorder BIF registration succeeds");
Arc::new(
Scheduler::with_services_and_code_server(
SchedulerConfig {
thread_count: Some(4),
..SchedulerConfig::default()
},
SchedulerServices::minimal(),
Arc::new(ModuleRegistry::new()),
atom_table,
bif_registry,
)
.expect("published beamr scheduler starts with the test BIF"),
)
}
fn stop_recorders() -> &'static Mutex<HashMap<i64, StopRecorder>> {
STOP_RECORDERS.get_or_init(|| Mutex::new(HashMap::new()))
}
fn registry(scheduler: &Arc<Scheduler>) -> Arc<ComponentRegistry> {
Arc::new(ComponentRegistry::new(
Arc::clone(scheduler),
LifecycleConfig {
operation_timeout: OPERATION_DEADLINE,
max_fragment_bytes: None,
},
))
}
fn id(name: &str) -> ComponentId {
ComponentId::derive("frame.tests", name)
}
fn policy(max_restarts: usize) -> SupervisionPolicy {
SupervisionPolicy {
max_restarts,
window: Duration::from_secs(30),
}
}
fn ffi_meta() -> ComponentMeta {
ComponentMeta {
id: id("ffi"),
name: "fixture ffi".to_owned(),
version: "1.0.0".to_owned(),
requires: Vec::new(),
provides: vec![ServiceCapability {
id: "fixture.ffi".to_owned(),
description: "fixture mailbox implementation".to_owned(),
}],
needs: Vec::new(),
actions: Vec::new(),
fragments: Vec::new(),
children: Vec::new(),
supervision: Some(policy(1)),
}
}
fn worker_meta(name: &str, max_restarts: usize) -> ComponentMeta {
ComponentMeta {
id: id(name),
name: name.to_owned(),
version: "1.0.0".to_owned(),
requires: vec![id("ffi")],
provides: vec![ServiceCapability {
id: format!("fixture.{name}"),
description: "fixture worker".to_owned(),
}],
needs: Vec::new(),
actions: Vec::new(),
fragments: Vec::new(),
children: vec![child("a"), child("b")],
supervision: Some(policy(max_restarts)),
}
}
fn registration_meta(name: &str) -> ComponentMeta {
ComponentMeta {
id: id(name),
name: name.to_owned(),
version: "1.0.0".to_owned(),
requires: Vec::new(),
provides: Vec::new(),
needs: Vec::new(),
actions: Vec::new(),
fragments: Vec::new(),
children: Vec::new(),
supervision: Some(policy(1)),
}
}
fn stop_order_meta() -> ComponentMeta {
let mut meta = registration_meta("drain-order");
meta.children = [("a", 10), ("b", 20), ("c", 30)]
.into_iter()
.map(|(name, tag)| ChildSpec {
name: name.to_owned(),
module: "stop_order_worker".to_owned(),
function: "run".to_owned(),
arguments: vec![
ChildArgument::SupervisorPid,
ChildArgument::Integer(STOP_RECORDER_TOKEN),
ChildArgument::Integer(tag),
],
liveness_message: 0,
stop_message: 2,
})
.collect();
meta
}
fn child(name: &str) -> ChildSpec {
ChildSpec {
name: name.to_owned(),
module: "spike_worker".to_owned(),
function: "run".to_owned(),
arguments: vec![ChildArgument::SupervisorPid],
liveness_message: 0,
stop_message: 2,
}
}
fn install_fixture(registry: &ComponentRegistry, worker: ComponentMeta) {
registry.register(ffi_meta(), FFI.to_vec()).unwrap();
registry.register(worker, WORKER.to_vec()).unwrap();
registry.start(id("ffi")).unwrap();
}
fn cleanup(registry: &ComponentRegistry, scheduler: &Scheduler, worker: ComponentId) {
if registry.status(worker).unwrap().is_some() {
let _ = registry.remove(worker).unwrap();
}
let _ = registry.remove(id("ffi")).unwrap();
assert_eq!(scheduler.process_count(), 0);
scheduler.shutdown();
}
fn record_stop(args: &[Term], _context: &mut ProcessContext) -> Result<Term, Term> {
let [token, tag] = args else {
return Err(Term::atom(Atom::BADARG));
};
let (Some(token), Some(tag)) = (token.as_small_int(), tag.as_small_int()) else {
return Err(Term::atom(Atom::BADARG));
};
let Some(recorder) = stop_recorders()
.lock()
.ok()
.and_then(|recorders| recorders.get(&token).cloned())
else {
return Err(Term::atom(Atom::BADARG));
};
if recorder.sender.send(tag).is_err() {
return Err(Term::atom(Atom::BADARG));
}
if tag == FIRST_STOP_TAG {
let (released, wake) = &*recorder.first_child_release;
let Ok(released) = released.lock() else {
return Err(Term::atom(Atom::BADARG));
};
let Ok((released, timeout)) =
wake.wait_timeout_while(released, OPERATION_DEADLINE, |released| !*released)
else {
return Err(Term::atom(Atom::BADARG));
};
if timeout.timed_out() && !*released {
return Err(Term::atom(Atom::BADARG));
}
}
Ok(Term::atom(Atom::OK))
}
#[test]
fn registration_failures_are_atomic_and_policy_has_no_default() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let mut missing_policy = ffi_meta();
missing_policy.supervision = None;
assert!(matches!(
registry.register(missing_policy, FFI.to_vec()),
Err(RegistryError::UndeclaredSupervision { .. })
));
assert_eq!(registry.len().unwrap(), 0);
registry.register(ffi_meta(), FFI.to_vec()).unwrap();
let before = registry.len().unwrap();
let mut missing = worker_meta("atomic", 1);
missing.requires.push(id("absent"));
assert!(matches!(
registry.register(missing, WORKER.to_vec()),
Err(RegistryError::MissingDependency { .. })
));
assert_eq!(registry.len().unwrap(), before);
let mut collision = worker_meta("collision", 1);
collision.provides[0].id = "fixture.ffi".to_owned();
assert!(matches!(
registry.register(collision, WORKER.to_vec()),
Err(RegistryError::CapabilityConflict { .. })
));
assert_eq!(registry.len().unwrap(), before);
registry.remove(id("ffi")).unwrap();
scheduler.shutdown();
}
#[test]
fn register_refuses_duplicate_component_identity() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let meta = registration_meta("duplicate-component");
let component_id = meta.id;
registry.register(meta.clone(), Vec::new()).unwrap();
assert!(matches!(
registry.register(meta, Vec::new()),
Err(RegistryError::DuplicateComponent { id }) if id == component_id
));
assert_eq!(registry.len().unwrap(), 1);
registry.remove(component_id).unwrap();
scheduler.shutdown();
}
#[test]
fn register_refuses_zero_length_supervision_window() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let mut meta = registration_meta("zero-supervision-window");
let component_id = meta.id;
meta.supervision = Some(SupervisionPolicy {
max_restarts: 1,
window: Duration::ZERO,
});
assert!(matches!(
registry.register(meta, Vec::new()),
Err(RegistryError::InvalidSupervisionWindow { id }) if id == component_id
));
assert_eq!(registry.len().unwrap(), 0);
scheduler.shutdown();
}
#[test]
fn register_refuses_duplicate_child_name() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let mut meta = registration_meta("duplicate-child");
let component_id = meta.id;
meta.children = vec![child("same-name"), child("same-name")];
assert!(matches!(
registry.register(meta, Vec::new()),
Err(RegistryError::DuplicateChild { id, child })
if id == component_id && child == "same-name"
));
assert_eq!(registry.len().unwrap(), 0);
scheduler.shutdown();
}
#[test]
fn register_refuses_dependency_cycle() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let component_a = id("cycle-a");
let component_b = id("cycle-b");
registry
.register(registration_meta("cycle-a"), Vec::new())
.unwrap();
let mut b = registration_meta("cycle-b");
b.requires.push(component_a);
registry.register(b, Vec::new()).unwrap();
registry.remove(component_a).unwrap();
let mut replacement_a = registration_meta("cycle-a");
replacement_a.requires.push(component_b);
assert!(matches!(
registry.register(replacement_a, Vec::new()),
Err(RegistryError::CyclicDependency {
component,
required,
}) if component == component_a && required == component_b
));
assert_eq!(registry.len().unwrap(), 1);
registry.remove(component_b).unwrap();
scheduler.shutdown();
}
#[test]
fn stop_drains_children_in_declaration_order() {
let scheduler = scheduler_with_stop_recorder();
let registry = registry(&scheduler);
let component_id = id("drain-order");
let (stop_tx, stop_rx) = mpsc::channel();
let first_child_release = Arc::new((Mutex::new(false), Condvar::new()));
stop_recorders()
.lock()
.expect("stop recorder table lock is available")
.insert(
STOP_RECORDER_TOKEN,
StopRecorder {
sender: stop_tx,
first_child_release: Arc::clone(&first_child_release),
},
);
registry
.register(stop_order_meta(), STOP_ORDER_WORKER.to_vec())
.expect("stop-order fixture registers");
registry
.start(component_id)
.expect("stop-order fixture starts");
let stop_registry = Arc::clone(®istry);
let stop = thread::spawn(move || stop_registry.stop(component_id));
assert_eq!(
stop_rx
.recv_timeout(OPERATION_DEADLINE)
.expect("first declared child records its stop"),
FIRST_STOP_TAG
);
assert_eq!(
stop_rx.recv_timeout(CONCURRENT_STOP_OBSERVATION),
Err(mpsc::RecvTimeoutError::Timeout),
"later children must not receive stop while the first child is still draining"
);
let (released, wake) = &*first_child_release;
*released
.lock()
.expect("first-child release gate lock is available") = true;
wake.notify_all();
assert_eq!(
stop.join()
.expect("stop worker thread does not panic")
.expect("ordered component stop succeeds"),
StopOutcome::Stopped
);
assert_eq!(
[
stop_rx
.recv_timeout(OPERATION_DEADLINE)
.expect("second declared child records its stop"),
stop_rx
.recv_timeout(OPERATION_DEADLINE)
.expect("third declared child records its stop"),
],
[20, 30]
);
assert!(stop_rx.try_recv().is_err(), "no undeclared child may stop");
stop_recorders()
.lock()
.expect("stop recorder table lock is available")
.remove(&STOP_RECORDER_TOKEN);
registry
.remove(component_id)
.expect("stopped fixture is removable");
assert_eq!(scheduler.process_count(), 0);
scheduler.shutdown();
}
#[test]
fn full_cycle_restarts_one_for_one_and_emits_every_transition() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let events = registry.subscribe(NonZeroUsize::new(64).unwrap()).unwrap();
let worker_id = id("cycle");
install_fixture(®istry, worker_meta("cycle", 3));
registry.start(worker_id).unwrap();
let started = registry.status(worker_id).unwrap().unwrap();
assert_eq!(started.state, LifecycleState::Running);
assert_eq!(started.children.len(), 2);
let old_a = started.children.iter().find(|c| c.name == "a").unwrap().pid;
let old_b = started.children.iter().find(|c| c.name == "b").unwrap().pid;
let exited = ExitWitness::register(&scheduler, old_a, OPERATION_DEADLINE).unwrap();
registry.send_child(worker_id, "a", 1).unwrap();
exited.wait().unwrap();
assert_eq!(registry.probe_child(worker_id, "a").unwrap(), 0);
assert_eq!(registry.probe_child(worker_id, "b").unwrap(), 1);
let after = registry.status(worker_id).unwrap().unwrap();
assert_eq!(
after.children.iter().find(|c| c.name == "b").unwrap().pid,
old_b
);
assert_eq!(registry.stop(worker_id).unwrap(), StopOutcome::Stopped);
assert_eq!(
registry.stop(worker_id).unwrap(),
StopOutcome::AlreadyStopped
);
assert_eq!(registry.remove(worker_id).unwrap(), RemoveOutcome::Removed);
assert_eq!(
registry.remove(worker_id).unwrap(),
RemoveOutcome::AlreadyRemoved
);
assert!(registry.status(worker_id).unwrap().is_none());
assert!(
scheduler
.lookup_module_in(
beamr::namespace::NamespaceId::DEFAULT,
scheduler.atom_table().intern("spike_worker")
)
.is_none()
);
registry
.register(worker_meta("cycle", 3), WORKER.to_vec())
.unwrap();
registry.start(worker_id).unwrap();
registry.stop(worker_id).unwrap();
registry.remove(worker_id).unwrap();
let mut worker_states = Vec::new();
while let Ok(event) = events.try_recv() {
if event.component_id == worker_id
&& let LifecycleEventKind::Transition { to, .. } = event.kind
{
worker_states.push(to);
}
}
assert_eq!(
worker_states,
vec![
LifecycleState::Registered,
LifecycleState::Starting,
LifecycleState::Running,
LifecycleState::Stopping,
LifecycleState::Stopped,
LifecycleState::Removed,
LifecycleState::Registered,
LifecycleState::Starting,
LifecycleState::Running,
LifecycleState::Stopping,
LifecycleState::Stopped,
LifecycleState::Removed,
]
);
cleanup(®istry, &scheduler, worker_id);
}
#[test]
fn restart_intensity_escalates_with_tombstone_and_drains_sibling() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let worker_id = id("intensity");
install_fixture(®istry, worker_meta("intensity", 1));
registry.start(worker_id).unwrap();
let first = registry.status(worker_id).unwrap().unwrap().children[0].pid;
let first_exit = ExitWitness::register(&scheduler, first, OPERATION_DEADLINE).unwrap();
registry.send_child(worker_id, "a", 1).unwrap();
first_exit.wait().unwrap();
assert_eq!(registry.probe_child(worker_id, "a").unwrap(), 0);
let before_escalation = registry.status(worker_id).unwrap().unwrap();
let replacement = before_escalation
.children
.iter()
.find(|child| child.name == "a")
.unwrap()
.pid;
let sibling = before_escalation
.children
.iter()
.find(|child| child.name == "b")
.unwrap()
.pid;
let replacement_exit =
ExitWitness::register(&scheduler, replacement, OPERATION_DEADLINE).unwrap();
let sibling_exit = ExitWitness::register(&scheduler, sibling, OPERATION_DEADLINE).unwrap();
registry.send_child(worker_id, "a", 1).unwrap();
replacement_exit.wait().unwrap();
sibling_exit.wait().unwrap();
let failed = registry.status(worker_id).unwrap().unwrap();
assert_eq!(failed.state, LifecycleState::Failed);
assert_eq!(
failed.failure.unwrap().tombstone,
Some(TombstoneKind::Error)
);
let outcome = registry.remove(worker_id).unwrap();
assert!(matches!(
outcome,
RemoveOutcome::RemovedAfterForcedCleanup(_)
));
let RemoveOutcome::RemovedAfterForcedCleanup(report) = outcome else {
return;
};
assert!(report.children.is_empty());
registry.remove(id("ffi")).unwrap();
assert_eq!(scheduler.process_count(), 0);
scheduler.shutdown();
}
#[test]
fn slow_subscriber_lags_without_blocking_live_subscriber() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let stalled = registry.subscribe(NonZeroUsize::new(1).unwrap()).unwrap();
let live = registry.subscribe(NonZeroUsize::new(32).unwrap()).unwrap();
let worker_id = id("events");
install_fixture(®istry, worker_meta("events", 1));
registry.start(worker_id).unwrap();
registry.stop(worker_id).unwrap();
registry.remove(worker_id).unwrap();
assert!(stalled.lagged_events() > 0);
let mut states = Vec::new();
while let Ok(event) = live.try_recv() {
if event.component_id == worker_id
&& let LifecycleEventKind::Transition { to, .. } = event.kind
{
states.push(to);
}
}
assert_eq!(
states,
vec![
LifecycleState::Registered,
LifecycleState::Starting,
LifecycleState::Running,
LifecycleState::Stopping,
LifecycleState::Stopped,
LifecycleState::Removed,
]
);
cleanup(®istry, &scheduler, worker_id);
}
#[test]
fn failed_start_is_removable_and_corrected_incarnation_starts() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let component_id = id("correctable");
let mut broken = worker_meta("correctable", 1);
broken.children[0].liveness_message = 1;
install_fixture(®istry, broken);
assert!(matches!(
registry.start(component_id),
Err(RegistryError::StartFailed { .. })
));
let failed = registry.status(component_id).unwrap().unwrap();
assert_eq!(failed.state, LifecycleState::Failed);
assert_eq!(
failed.failure.unwrap().tombstone,
Some(TombstoneKind::Error)
);
registry.remove(component_id).unwrap();
registry
.register(worker_meta("correctable", 1), WORKER.to_vec())
.unwrap();
registry.start(component_id).unwrap();
assert_eq!(registry.probe_child(component_id, "a").unwrap(), 1);
cleanup(®istry, &scheduler, component_id);
}
#[test]
fn concurrent_child_crashes_are_both_observed_and_restarted() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let component_id = id("concurrent-crashes");
install_fixture(®istry, worker_meta("concurrent-crashes", 4));
registry.start(component_id).unwrap();
let before = registry.status(component_id).unwrap().unwrap();
let old_a = before
.children
.iter()
.find(|child| child.name == "a")
.unwrap()
.pid;
let old_b = before
.children
.iter()
.find(|child| child.name == "b")
.unwrap()
.pid;
let exit_a = ExitWitness::register(&scheduler, old_a, OPERATION_DEADLINE).unwrap();
let exit_b = ExitWitness::register(&scheduler, old_b, OPERATION_DEADLINE).unwrap();
registry.send_child(component_id, "a", 1).unwrap();
registry.send_child(component_id, "b", 1).unwrap();
exit_a.wait().unwrap();
exit_b.wait().unwrap();
assert_eq!(registry.probe_child(component_id, "a").unwrap(), 0);
assert_eq!(registry.probe_child(component_id, "b").unwrap(), 0);
cleanup(®istry, &scheduler, component_id);
}
#[test]
fn failed_component_is_invisible_to_independent_running_component() {
let scheduler = scheduler();
let registry = registry(&scheduler);
registry.register(ffi_meta(), FFI.to_vec()).unwrap();
let failed_id = id("failure-a");
let healthy_id = id("healthy-b");
registry
.register(worker_meta("failure-a", 0), WORKER.to_vec())
.unwrap();
let mut healthy = worker_meta("healthy-b", 1);
for child in &mut healthy.children {
child.module = "spike_worker_b".to_owned();
}
registry.register(healthy, WORKER_B.to_vec()).unwrap();
registry.start(id("ffi")).unwrap();
registry.start(failed_id).unwrap();
registry.start(healthy_id).unwrap();
let failing = registry.status(failed_id).unwrap().unwrap();
let failed_a = failing
.children
.iter()
.find(|child| child.name == "a")
.unwrap()
.pid;
let failed_b = failing
.children
.iter()
.find(|child| child.name == "b")
.unwrap()
.pid;
let exit_a = ExitWitness::register(&scheduler, failed_a, OPERATION_DEADLINE).unwrap();
let exit_b = ExitWitness::register(&scheduler, failed_b, OPERATION_DEADLINE).unwrap();
registry.send_child(failed_id, "a", 1).unwrap();
exit_a.wait().unwrap();
exit_b.wait().unwrap();
let failed = registry.status(failed_id).unwrap().unwrap();
assert_eq!(failed.state, LifecycleState::Failed);
let healthy = registry.status(healthy_id).unwrap().unwrap();
assert_eq!(healthy.state, LifecycleState::Running);
assert_eq!(registry.probe_child(healthy_id, "a").unwrap(), 1);
let outcome = registry.remove(failed_id).unwrap();
assert!(matches!(
outcome,
RemoveOutcome::RemovedAfterForcedCleanup(_)
));
let RemoveOutcome::RemovedAfterForcedCleanup(report) = outcome else {
return;
};
assert!(report.children.is_empty());
registry.remove(healthy_id).unwrap();
registry.remove(id("ffi")).unwrap();
assert_eq!(scheduler.process_count(), 0);
scheduler.shutdown();
}
#[test]
fn running_dependent_gates_dependency_removal() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let component_id = id("dependent");
install_fixture(®istry, worker_meta("dependent", 1));
registry.start(component_id).unwrap();
assert!(matches!(
registry.remove(id("ffi")),
Err(RegistryError::RunningDependent { dependent, .. }) if dependent == component_id
));
cleanup(®istry, &scheduler, component_id);
}
#[test]
fn concurrent_starts_have_one_winner_and_one_inflight_refusal() {
let scheduler = scheduler();
let registry = registry(&scheduler);
let worker_id = id("concurrent");
install_fixture(®istry, worker_meta("concurrent", 1));
let events = registry.subscribe(NonZeroUsize::new(8).unwrap()).unwrap();
let barrier = Arc::new(Barrier::new(3));
let mut joins = Vec::new();
for _ in 0..2 {
let registry = Arc::clone(®istry);
let barrier = Arc::clone(&barrier);
joins.push(thread::spawn(move || {
barrier.wait();
registry.start(worker_id)
}));
}
barrier.wait();
let results: Vec<_> = joins.into_iter().map(|join| join.join().unwrap()).collect();
assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
assert_eq!(
results
.iter()
.filter(|result| matches!(result, Err(RegistryError::AlreadyStarting { .. })))
.count(),
1
);
let mut states = Vec::new();
while let Ok(event) = events.try_recv() {
if event.component_id == worker_id
&& let LifecycleEventKind::Transition { to, .. } = event.kind
{
states.push(to);
}
}
assert_eq!(
states,
vec![LifecycleState::Starting, LifecycleState::Running]
);
cleanup(®istry, &scheduler, worker_id);
}