frame-core 0.3.0

Component model, lifecycle, process isolation — hosts components as supervised BEAM process trees
Documentation
//! The composed component runtime and its opaque support-module handles.

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;

/// Stated-by-the-application runtime policy. No field has a default: frame
/// supplies no lifecycle values, so the embedding application declares each
/// one explicitly.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct RuntimePolicy {
    /// Scheduler worker threads for the component tree. Must be positive —
    /// the wrapper refuses zero rather than letting the runtime silently
    /// substitute the machine's parallelism.
    pub scheduler_threads: usize,
    /// Deadline bounding each spawn, mailbox round-trip, or tombstone
    /// observation. Must be nonzero.
    pub operation_timeout: Duration,
    /// Caller-supplied per-fragment content byte limit (F-5a R1 — no hidden
    /// default). `None` states loudly that this host configured no limit;
    /// registering a fragment-carrying component then refuses typed.
    pub max_fragment_bytes: Option<std::num::NonZeroUsize>,
}

/// Opaque, single-use handle to a loaded support module.
///
/// Returned by [`ComponentRuntime::load_support_module`] and consumed by
/// [`ComponentRuntime::unload_support_module`]; because it is neither `Clone`
/// nor `Copy`, a double unload is unrepresentable.
#[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 {
    /// A swap reference to the same module NAME: what the F-7b dev reload
    /// swaps against. Deliberately separate from the handle itself — a
    /// ref cannot unload, so the double-unload law the handle encodes
    /// stays unrepresentable while swaps stay cheap to thread around.
    #[must_use]
    pub fn swap_ref(&self) -> SupportModuleRef {
        SupportModuleRef {
            module: self.module,
            name: self.name.clone(),
        }
    }
}

/// A cloneable reference to a loaded support module's NAME, for hot
/// swaps. Carries no unload authority.
#[derive(Clone, Debug)]
pub struct SupportModuleRef {
    module: Atom,
    name: String,
}

/// The composed component runtime: scheduler + registry, correctly assembled.
///
/// Composition happens once, in [`ComponentRuntime::compose`], with the
/// built-in-function registry populated before any bytecode can load (see the
/// [module docs](crate::runtime) for the full composition story). The
/// scheduler itself stays private: the host's authority surface is
/// [`ComponentRuntime::registry`], support modules travel as opaque handles,
/// and teardown is the residue-checked [`ComponentRuntime::shutdown`].
pub struct ComponentRuntime {
    scheduler: Arc<Scheduler>,
    registry: Arc<ComponentRegistry>,
}

impl ComponentRuntime {
    /// Composes the runtime with its built-in-function registry populated
    /// before any bytecode can load (the composition
    /// [`crate::composition::compose_scheduler`] performs, absorbed).
    ///
    /// # Errors
    ///
    /// Refuses a zero thread count or zero operation timeout with typed
    /// policy errors, and propagates every typed scheduler-composition
    /// failure.
    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,
        })
    }

    /// Borrows the host's authority surface (register, grant, start, probe,
    /// stop — the unchanged [`ComponentRegistry`] API).
    #[must_use]
    pub fn registry(&self) -> &ComponentRegistry {
        &self.registry
    }

    /// A shared handle to the same authority surface, for a management
    /// adapter running on its own thread (the F-7b dev door). The handle
    /// may outlive `shutdown`; operations against a stopped scheduler
    /// fail typed, never silently.
    #[must_use]
    pub fn registry_handle(&self) -> Arc<ComponentRegistry> {
        Arc::clone(&self.registry)
    }

    /// Loads a support module (e.g. a component's FFI bytecode) and returns
    /// an opaque handle for ordered unload.
    ///
    /// The freshly committed module passes the same deferred-`erlang:*` wall
    /// the registry applies to component bytecode: a module whose import
    /// table still defers a built-in would die at first dispatch, so it is
    /// refused — and removed again — up front.
    ///
    /// # Errors
    ///
    /// Returns typed failures for loader refusal, a module absent immediately
    /// after load, or deferred `erlang:*` imports.
    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() {
            // The just-loaded module is unusable; take it back out before
            // refusing so the refusal does not leak a loaded generation. A
            // cleanup failure is logged loudly but never masks the refusal
            // itself.
            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 })
    }

    /// Ordered unload with verification: purge any retained old generation,
    /// delete the current one, then verify the module is actually gone.
    ///
    /// # Errors
    ///
    /// Returns typed failures for an unsafe purge, a failed delete (a stale
    /// handle to an already-unloaded module lands here), or post-delete
    /// lookup residue.
    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(())
    }

    /// Hot-swaps a support module to new bytes of the SAME module name —
    /// see [`SupportSwapper::swap`], which this delegates to.
    ///
    /// # Errors
    ///
    /// Returns typed failures for an unsafe purge, loader refusal, a
    /// module absent immediately after load, or deferred `erlang:*`
    /// imports.
    pub fn swap_support_module(
        &self,
        current: &SupportModuleRef,
        bytecode: &[u8],
    ) -> Result<SupportModuleRef, RuntimeError> {
        self.support_swapper().swap(current, bytecode)
    }

    /// A thread-shareable handle for support-module swaps, for a
    /// management adapter running on its own thread (the F-7b dev door).
    #[must_use]
    pub fn support_swapper(&self) -> SupportSwapper {
        SupportSwapper {
            scheduler: Arc::clone(&self.scheduler),
        }
    }

    /// Live process count — the zero-residue check after ordered stop.
    #[must_use]
    pub fn live_process_count(&self) -> usize {
        self.scheduler.process_count()
    }

    /// Final teardown; refuses if the process tree still has residue.
    ///
    /// `Scheduler::shutdown` is never cleanup: every component must have
    /// completed its ordered stop first. A refusal deliberately does NOT stop
    /// the scheduler — live processes are the embedder's unfinished ordered
    /// stop, and reaping them here would hide it.
    ///
    /// # Errors
    ///
    /// Returns the live process count as a typed residue refusal.
    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(())
    }
}

/// A `Send + Sync` handle for hot-swapping support modules from the
/// management adapter's thread (the F-7b dev reload's FFI leg).
pub struct SupportSwapper {
    scheduler: Arc<Scheduler>,
}

impl SupportSwapper {
    /// Hot-swaps a support module to new bytes of the SAME module name:
    /// purges the generation the previous swap retained, then hot-loads
    /// the new bytes — the module-generation wall never bites because
    /// every swap purges before it loads. Only legal while every
    /// component using the module is stopped (the barrier's order
    /// guarantees this); a purge refusal is the typed drain-failure
    /// signal.
    ///
    /// A deferred-`erlang:*` refusal here does NOT delete the module the
    /// way `load_support_module`'s initial-load cleanup does: at this
    /// point the bad bytes are current and the previous bytes are
    /// retained old — the caller's rollback swap purges and reloads the
    /// previous bytes, which a delete would sabotage.
    ///
    /// # Errors
    ///
    /// Returns typed failures for an unsafe purge, loader refusal, a
    /// module absent immediately after load, or deferred `erlang:*`
    /// imports.
    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 })
    }
}