use crate::error::{LoaderError, Result};
use crate::lock;
use crate::registry::{PluginRegistry, WithInject};
use cordis::{Config, Context, EffectHandle, EventOptions, Fiber, FiberState, PluginHandle};
use cordis_include::{Entry, EntryOptions, EntryTree, LoaderFile, Node, PluginResolver, TreeDiff};
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::{Arc, Mutex, Weak};
#[derive(Clone, Default)]
pub struct LoaderConfig {
pub filename: PathBuf,
pub initial: Option<cordis_include::Document>,
pub registry: Option<PluginRegistry>,
}
impl LoaderConfig {
pub fn new(filename: impl Into<PathBuf>) -> Self {
Self {
filename: filename.into(),
initial: None,
registry: None,
}
}
pub fn with_initial(mut self, initial: cordis_include::Document) -> Self {
self.initial = Some(initial);
self
}
pub fn with_registry(mut self, registry: PluginRegistry) -> Self {
self.registry = Some(registry);
self
}
}
struct LoaderState {
entries: HashMap<u64, Entry>,
operating: u16,
last_error: Option<String>,
_keep_alive: Vec<EffectHandle>,
}
#[derive(Clone)]
pub struct Loader {
pub(crate) inner: Arc<LoaderInner>,
}
pub(crate) struct LoaderInner {
root: Context,
file: LoaderFile,
tree: EntryTree,
registry: Mutex<PluginRegistry>,
state: Mutex<LoaderState>,
}
pub struct LoaderHandle {
inner: Weak<LoaderInner>,
}
impl LoaderHandle {
pub fn upgrade(&self) -> Option<Loader> {
self.inner.upgrade().map(|inner| Loader { inner })
}
}
impl Loader {
pub fn open(root: &Context, config: LoaderConfig) -> Result<Loader> {
let file = LoaderFile::open(&config.filename)?;
if !file.path().exists() {
if let Some(initial) = &config.initial {
file.write(initial)?;
}
}
let tree = EntryTree::new();
tree.update(file.read()?.entries)?;
let inner = Arc::new(LoaderInner {
root: root.clone(),
file,
tree,
registry: Mutex::new(config.registry.unwrap_or_default()),
state: Mutex::new(LoaderState {
entries: HashMap::new(),
operating: 0,
last_error: None,
_keep_alive: Vec::new(),
}),
});
let weak = Arc::downgrade(&inner);
let status = root.events().on(
"internal/status",
move |event| {
if let Some(inner) = weak.upgrade() {
handle_status(&inner, &event)?;
}
Ok(None)
},
EventOptions {
global: true,
..EventOptions::default()
},
)?;
let service = root.provide_arc(
"loader",
Arc::new(LoaderHandle {
inner: Arc::downgrade(&inner),
}),
)?;
lock(&inner.state)._keep_alive = vec![status, service];
let loader = Loader { inner };
loader.start_all();
Ok(loader)
}
pub fn context(&self) -> &Context {
&self.inner.root
}
pub fn tree(&self) -> &EntryTree {
&self.inner.tree
}
pub fn file(&self) -> &LoaderFile {
&self.inner.file
}
pub fn registry(&self) -> PluginRegistry {
lock(&self.inner.registry).clone()
}
pub fn register_plugin<P: cordis::Plugin>(&self, plugin: P) {
lock(&self.inner.registry).register_plugin(plugin);
}
pub fn register<F>(&self, name: impl Into<String>, factory: F)
where
F: Fn() -> PluginHandle + Send + Sync + 'static,
{
lock(&self.inner.registry).register(name, factory);
}
pub fn last_error(&self) -> Option<String> {
lock(&self.inner.state).last_error.clone()
}
pub fn locate(&self, fiber: &Fiber) -> Option<Entry> {
let state = lock(&self.inner.state);
if let Some(uid) = fiber.uid() {
return state.entries.get(&uid).cloned();
}
state
.entries
.values()
.find(|entry| entry.fiber().is_some_and(|started| started.ptr_eq(fiber)))
.cloned()
}
fn start_all(&self) {
for entry in self.inner.tree.entries() {
if let Err(error) = start_entry(&self.inner, &entry) {
self.record_error(error);
}
}
}
pub fn reload(&self) -> Result<TreeDiff> {
let inner = &self.inner;
let _suspend = inner.file.suspend();
let document = inner.file.read()?;
let diff = inner.tree.update(document.entries)?;
for entry in &diff.removed {
if let Err(error) = stop_entry(inner, entry) {
self.record_error(error);
}
}
for entry in &diff.moved {
if let Err(error) = stop_entry(inner, entry) {
self.record_error(error);
}
}
for entry in &diff.updated {
if let Err(error) = patch_entry(inner, entry) {
self.record_error(error);
}
}
for entry in &diff.created {
if let Err(error) = start_entry(inner, entry) {
self.record_error(error);
}
}
drop(_suspend);
if !diff.created.is_empty() {
write_back(inner)?;
}
Ok(diff)
}
pub fn update_config(&self, id: &str, config: Node) -> Result<()> {
let inner = &self.inner;
let entry = inner.tree.resolve(id).ok_or_else(|| {
LoaderError::Include(cordis_include::IncludeError::EntryNotFound { id: id.to_owned() })
})?;
if let Some(fiber) = entry.fiber() {
fiber.update_value(Config::new(config.clone()))?;
}
let mut options = entry_options_with_children(&entry);
options.config = Some(config);
inner
.tree
.update_entry(&entry.path(), options, None, None)?;
write_back(inner)
}
pub fn dispose(&self) -> Result<()> {
let inner = &self.inner;
for entry in inner.tree.top_level() {
if let Err(error) = stop_entry(inner, &entry) {
self.record_error(error);
}
}
Ok(())
}
#[cfg(feature = "watch")]
pub fn watch(&self) -> Result<cordis_include::FileWatcher> {
let loader = self.clone();
self.inner
.file
.watch(move || {
if let Err(error) = loader.reload() {
loader.record_error(error);
}
})
.map_err(LoaderError::Include)
}
fn record_error(&self, error: LoaderError) {
lock(&self.inner.state).last_error = Some(error.to_string());
}
}
struct OperatingGuard<'a> {
state: &'a Mutex<LoaderState>,
}
impl<'a> OperatingGuard<'a> {
fn new(state: &'a Mutex<LoaderState>) -> Self {
lock(state).operating += 1;
Self { state }
}
}
impl Drop for OperatingGuard<'_> {
fn drop(&mut self) {
let mut state = lock(self.state);
state.operating = state.operating.saturating_sub(1);
}
}
fn start_entry(inner: &LoaderInner, entry: &Entry) -> Result<()> {
if !entry.enabled() || entry.fiber().is_some() {
return Ok(());
}
let name = entry.name();
let handle: PluginHandle = lock(&inner.registry)
.resolve(&name)
.map_err(LoaderError::Cordis)?;
let inject = entry.options().inject;
let handle = WithInject::wrap(handle, inject);
let config = entry.resolved_config()?.unwrap_or(Node::Null);
let parent_ctx = entry
.parent()
.and_then(|parent| parent.fiber())
.and_then(|fiber| fiber.context())
.unwrap_or_else(|| inner.root.clone());
let fiber = parent_ctx.plugin(handle, config);
entry.set_fiber(Some(fiber.clone()));
if let Some(uid) = fiber.uid() {
lock(&inner.state).entries.insert(uid, entry.clone());
}
Ok(())
}
fn stop_entry(inner: &LoaderInner, entry: &Entry) -> Result<()> {
for child in entry.children() {
stop_entry(inner, &child)?;
}
let Some(fiber) = entry.fiber() else {
return Ok(());
};
entry.set_fiber(None);
if let Some(uid) = fiber.uid() {
lock(&inner.state).entries.remove(&uid);
}
let _guard = OperatingGuard::new(&inner.state);
fiber.dispose().map_err(LoaderError::Cordis)
}
fn patch_entry(inner: &LoaderInner, entry: &Entry) -> Result<()> {
if !entry.enabled() {
return stop_entry(inner, entry);
}
let Some(fiber) = entry.fiber() else {
return start_entry(inner, entry);
};
let options = entry.options();
let inject_changed =
fiber.inject().names().collect::<Vec<_>>() != options.inject.iter().collect::<Vec<_>>();
if fiber.name() != options.name || inject_changed {
stop_entry(inner, entry)?;
return start_entry(inner, entry);
}
let new_config = entry.resolved_config()?.unwrap_or(Node::Null);
let current = fiber
.config()
.downcast::<Node>()
.ok()
.map(|node| (*node).clone());
if current.as_ref() != Some(&new_config) {
fiber.update_value(Config::new(new_config))?;
}
Ok(())
}
fn entry_options_with_children(entry: &Entry) -> EntryOptions {
let mut options = entry.options();
options.group = entry
.children()
.iter()
.map(entry_options_with_children)
.collect();
options
}
fn write_back(inner: &LoaderInner) -> Result<()> {
let mut document = inner.file.read()?;
document.entries = inner.tree.serialize();
inner.file.write(&document)?;
Ok(())
}
fn handle_status(inner: &LoaderInner, event: &cordis::Event) -> cordis::EventResult {
let Some(fiber) = event.arg::<Fiber>(0).ok().flatten() else {
return Ok(None);
};
if fiber.state() != FiberState::Disposed {
return Ok(None);
}
if lock(&inner.state).operating > 0 {
return Ok(None);
}
let Some(entry) = lock(&inner.state)
.entries
.values()
.find(|entry| entry.fiber().is_some_and(|started| started.ptr_eq(&fiber)))
.cloned()
else {
return Ok(None);
};
if let Err(error) = persist_self_dispose(inner, &entry) {
lock(&inner.state).last_error = Some(error.to_string());
}
Ok(None)
}
fn persist_self_dispose(inner: &LoaderInner, entry: &Entry) -> Result<()> {
{
let mut state = lock(&inner.state);
let key = state
.entries
.iter()
.find(|(_, mapped)| Entry::ptr_eq(mapped, entry))
.map(|(uid, _)| *uid);
if let Some(uid) = key {
state.entries.remove(&uid);
}
}
entry.set_fiber(None);
let mut options = entry_options_with_children(entry);
options.disabled = true;
inner
.tree
.update_entry(&entry.path(), options, None, None)?;
write_back(inner)
}