use std::collections::HashSet;
use std::num::NonZeroUsize;
use std::thread::JoinHandle;
use frame_core::component::ComponentId;
use frame_core::event::{
EventReceiveError, LifecycleEventKind, LifecycleState, LifecycleSubscription,
};
use frame_core::registry::ComponentRegistry;
use frame_core::runtime::{ComponentRuntime, RuntimePolicy, SupportModuleHandle};
use frame_core::status::ComponentStatus;
use crate::error::HostError;
use crate::spec::ComponentInstall;
const LIFECYCLE_EVENT_BUFFER: usize = 256;
const LOGGER_WAIT_QUANTUM: std::time::Duration = std::time::Duration::from_secs(3600);
pub struct HostRuntime {
runtime: ComponentRuntime,
components: Vec<ComponentId>,
support: Vec<SupportModuleHandle>,
}
pub struct EventLoggerHandle {
join: JoinHandle<()>,
}
impl HostRuntime {
pub fn new(policy: RuntimePolicy) -> Result<Self, HostError> {
let runtime = ComponentRuntime::compose(policy)?;
Ok(Self {
runtime,
components: Vec::new(),
support: Vec::new(),
})
}
#[must_use]
pub fn registry(&self) -> &ComponentRegistry {
self.runtime.registry()
}
#[must_use]
pub fn registry_handle(&self) -> std::sync::Arc<ComponentRegistry> {
self.runtime.registry_handle()
}
#[must_use]
pub fn support_swapper(&self) -> frame_core::runtime::SupportSwapper {
self.runtime.support_swapper()
}
#[must_use]
pub fn first_support_ref(&self) -> Option<frame_core::runtime::SupportModuleRef> {
self.support.first().map(SupportModuleHandle::swap_ref)
}
pub fn spawn_event_logger(
&self,
expected: HashSet<ComponentId>,
) -> Result<EventLoggerHandle, HostError> {
let capacity =
NonZeroUsize::new(LIFECYCLE_EVENT_BUFFER).ok_or(HostError::SynchronizationPoisoned)?;
let subscription = self.registry().subscribe(capacity)?;
let join = std::thread::Builder::new()
.name("frame-host-lifecycle-log".to_owned())
.spawn(move || log_lifecycle_stream(&subscription, &expected))
.map_err(|source| HostError::EventLoggerSpawn { source })?;
Ok(EventLoggerHandle { join })
}
pub fn install(
&mut self,
installs: Vec<ComponentInstall>,
) -> Result<Vec<ComponentStatus>, HostError> {
let mut statuses = Vec::with_capacity(installs.len());
for install in installs {
for bytecode in &install.support_modules {
let handle = self.runtime.load_support_module(bytecode)?;
self.support.push(handle);
}
let id = install.meta.id;
tracing::info!(
component = %id,
name = %install.meta.name,
"registering application component"
);
self.registry().register(install.meta, install.bytecode)?;
self.components.push(id);
let grants = self.registry().host_capabilities().grants_for(id)?;
tracing::info!(
component = %id,
grants = grants.len(),
"capability table registered component"
);
self.registry().start(id)?;
let status = self
.registry()
.status(id)?
.ok_or(HostError::StatusMissing { id })?;
if status.state != LifecycleState::Running {
return Err(HostError::NotRunning {
id,
state: status.state,
});
}
for child in &status.children {
tracing::info!(
component = %id,
child = %child.name,
pid = child.pid,
"supervised child is live"
);
}
statuses.push(status);
}
Ok(statuses)
}
pub fn abandon_after_boot_failure(&self) {
for id in self.components.iter().rev() {
match self.registry().remove(*id) {
Ok(outcome) => {
tracing::warn!(component = %id, outcome = ?outcome, "removed component after failed boot");
}
Err(cleanup) => {
tracing::error!(component = %id, %cleanup, "cleanup after failed boot also failed");
}
}
}
}
pub fn shutdown(mut self, logger: EventLoggerHandle) -> Result<(), HostError> {
for id in self.components.iter().rev() {
let stop_outcome = self.registry().stop(*id)?;
tracing::info!(component = %id, outcome = ?stop_outcome, "ordered stop completed");
let remove_outcome = self.registry().remove(*id)?;
tracing::info!(component = %id, outcome = ?remove_outcome, "component removed");
}
while let Some(handle) = self.support.pop() {
self.runtime.unload_support_module(handle)?;
}
let residue = self.runtime.live_process_count();
if residue != 0 {
return Err(HostError::ProcessResidue { count: residue });
}
self.runtime.shutdown()?;
logger
.join
.join()
.map_err(|_| HostError::EventLoggerPanicked)
}
}
fn log_lifecycle_stream(subscription: &LifecycleSubscription, expected: &HashSet<ComponentId>) {
let mut reported_lag = 0;
let mut removed: HashSet<ComponentId> = HashSet::new();
loop {
let event = match subscription.recv_timeout(LOGGER_WAIT_QUANTUM) {
Ok(event) => event,
Err(EventReceiveError::Timeout) => continue,
Err(EventReceiveError::Closed) => {
tracing::info!("lifecycle event stream closed; console logger exiting");
return;
}
Err(EventReceiveError::Poisoned) => {
tracing::error!(
"lifecycle event stream synchronization poisoned; console logger exiting"
);
return;
}
};
let lagged = subscription.lagged_events();
if lagged > reported_lag {
tracing::warn!(
dropped = lagged - reported_lag,
total_dropped = lagged,
"lifecycle logger overflowed its bounded buffer; oldest events were dropped"
);
reported_lag = lagged;
}
match &event.kind {
LifecycleEventKind::Transition { from, to } => {
tracing::info!(
sequence = event.sequence,
component = %event.component_id,
from = ?from,
to = ?to,
"component lifecycle transition"
);
if *to == LifecycleState::Removed {
removed.insert(event.component_id);
if removed.is_superset(expected) {
tracing::info!(
"every application component removed; console logger exiting"
);
return;
}
}
}
LifecycleEventKind::CapabilityDenied(denial) => {
tracing::warn!(
sequence = event.sequence,
component = %denial.component_id,
kind = ?denial.kind,
scope = ?denial.scope,
declared = denial.declared,
"capability check denied"
);
}
LifecycleEventKind::FragmentContentUpdated { key } => {
tracing::info!(
sequence = event.sequence,
component = %event.component_id,
fragment = key.fragment_id.as_str(),
"fragment content updated"
);
}
}
}
}
#[cfg(test)]
mod tests {
use crate::manifest;
#[test]
fn manifest_declares_the_contract_and_presence_child() {
let meta = manifest::component_meta();
assert_eq!(meta.provides.len(), 1);
assert_eq!(meta.provides[0].id, manifest::CONTRACT_ID);
assert_eq!(meta.children.len(), 1);
assert_eq!(meta.children[0].name, manifest::PRESENCE_CHILD);
assert!(
meta.needs.is_empty(),
"class-(b) manifest must declare no host needs"
);
assert!(meta.supervision.is_some());
}
}