use std::sync::Arc;
use std::time::Duration;
use beamr::atom::Atom;
use beamr::module::ModuleRegistry;
use beamr::namespace::NamespaceId;
use beamr::scheduler::{Scheduler, SchedulerConfig, SchedulerServices};
use crate::composition::{compose_scheduler, deferred_erlang_imports, resolve_atom};
use crate::registry::ComponentRegistry;
use crate::supervision::LifecycleConfig;
use super::error::RuntimeError;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct RuntimePolicy {
pub scheduler_threads: usize,
pub operation_timeout: Duration,
pub max_fragment_bytes: Option<std::num::NonZeroUsize>,
}
#[derive(Debug)]
#[must_use = "an unconsumed handle leaves its module loaded; unload it through the runtime"]
pub struct SupportModuleHandle {
module: Atom,
name: String,
}
impl SupportModuleHandle {
#[must_use]
pub fn swap_ref(&self) -> SupportModuleRef {
SupportModuleRef {
module: self.module,
name: self.name.clone(),
}
}
}
#[derive(Clone, Debug)]
pub struct SupportModuleRef {
module: Atom,
name: String,
}
pub struct ComponentRuntime {
scheduler: Arc<Scheduler>,
registry: Arc<ComponentRegistry>,
}
impl ComponentRuntime {
pub fn compose(policy: RuntimePolicy) -> Result<Self, RuntimeError> {
if policy.scheduler_threads == 0 {
return Err(RuntimeError::ZeroSchedulerThreads);
}
if policy.operation_timeout.is_zero() {
return Err(RuntimeError::ZeroOperationTimeout);
}
let scheduler = Arc::new(
compose_scheduler(
SchedulerConfig {
thread_count: Some(policy.scheduler_threads),
..SchedulerConfig::default()
},
SchedulerServices::minimal(),
Arc::new(ModuleRegistry::new()),
)
.map_err(|source| RuntimeError::Composition { source })?,
);
let registry = Arc::new(ComponentRegistry::new(
Arc::clone(&scheduler),
LifecycleConfig {
operation_timeout: policy.operation_timeout,
max_fragment_bytes: policy.max_fragment_bytes,
},
));
Ok(Self {
scheduler,
registry,
})
}
#[must_use]
pub fn registry(&self) -> &ComponentRegistry {
&self.registry
}
#[must_use]
pub fn registry_handle(&self) -> Arc<ComponentRegistry> {
Arc::clone(&self.registry)
}
pub fn load_support_module(
&self,
bytecode: &[u8],
) -> Result<SupportModuleHandle, RuntimeError> {
let loaded = self.scheduler.hot_load_module(bytecode).map_err(|error| {
RuntimeError::SupportModuleLoad {
detail: error.to_string(),
}
})?;
let module = loaded.module_name;
let name = resolve_atom(&self.scheduler, module);
let Some(deferred) = deferred_erlang_imports(&self.scheduler, module) else {
return Err(RuntimeError::SupportModuleAbsent { module: name });
};
if !deferred.is_empty() {
if !self.scheduler.delete_module(module) {
tracing::error!(
module = %name,
"failed to delete support module after deferred-BIF refusal"
);
}
return Err(RuntimeError::SupportModuleDeferredBifImports {
module: name,
imports: deferred.join(", "),
});
}
Ok(SupportModuleHandle { module, name })
}
pub fn unload_support_module(&self, handle: SupportModuleHandle) -> Result<(), RuntimeError> {
let SupportModuleHandle { module, name } = handle;
if self.scheduler.check_old_code(module) {
self.scheduler.purge_module(module).map_err(|error| {
RuntimeError::SupportModulePurge {
module: name.clone(),
detail: error.to_string(),
}
})?;
}
if !self.scheduler.delete_module(module) {
return Err(RuntimeError::SupportModuleDelete { module: name });
}
if self
.scheduler
.lookup_module_in(NamespaceId::DEFAULT, module)
.is_some()
{
return Err(RuntimeError::SupportModuleStillLoaded { module: name });
}
Ok(())
}
pub fn swap_support_module(
&self,
current: &SupportModuleRef,
bytecode: &[u8],
) -> Result<SupportModuleRef, RuntimeError> {
self.support_swapper().swap(current, bytecode)
}
#[must_use]
pub fn support_swapper(&self) -> SupportSwapper {
SupportSwapper {
scheduler: Arc::clone(&self.scheduler),
}
}
#[must_use]
pub fn live_process_count(&self) -> usize {
self.scheduler.process_count()
}
pub fn shutdown(self) -> Result<(), RuntimeError> {
let residue = self.scheduler.process_count();
if residue != 0 {
return Err(RuntimeError::ProcessResidue { count: residue });
}
self.scheduler.shutdown();
Ok(())
}
}
pub struct SupportSwapper {
scheduler: Arc<Scheduler>,
}
impl SupportSwapper {
pub fn swap(
&self,
current: &SupportModuleRef,
bytecode: &[u8],
) -> Result<SupportModuleRef, RuntimeError> {
if self.scheduler.check_old_code(current.module) {
self.scheduler
.purge_module(current.module)
.map_err(|error| RuntimeError::SupportModulePurge {
module: current.name.clone(),
detail: error.to_string(),
})?;
}
let loaded = self.scheduler.hot_load_module(bytecode).map_err(|error| {
RuntimeError::SupportModuleLoad {
detail: error.to_string(),
}
})?;
let module = loaded.module_name;
let name = resolve_atom(&self.scheduler, module);
let Some(deferred) = deferred_erlang_imports(&self.scheduler, module) else {
return Err(RuntimeError::SupportModuleAbsent { module: name });
};
if !deferred.is_empty() {
return Err(RuntimeError::SupportModuleDeferredBifImports {
module: name,
imports: deferred.join(", "),
});
}
Ok(SupportModuleRef { module, name })
}
}