use serde::de::DeserializeOwned;
use crate::{
ActPackage, ActPlugin, ChannelOptions, Signal,
builder::EngineBuilder,
config::{Config, ConfigResolver},
export::{Channel, Executor, Extender},
package::{self, ActPackageRegister},
scheduler::Runtime,
store::KvStore,
};
use std::sync::Arc;
use tracing::info;
#[derive(Clone)]
pub struct Engine {
config: Arc<Config>,
plugins: Vec<Arc<dyn ActPlugin>>,
packages: Vec<ActPackageRegister>,
resolvers: Vec<(String, Arc<dyn ConfigResolver>)>,
store: Option<Arc<dyn KvStore>>,
runtime: Option<Arc<Runtime>>,
}
impl Default for Engine {
fn default() -> Self {
Self::new()
}
}
impl Engine {
pub fn new() -> Self {
Self {
config: Arc::new(Config::default()),
plugins: Vec::new(),
packages: Vec::new(),
resolvers: Vec::new(),
store: None,
runtime: None,
}
}
pub fn config(&self) -> Arc<Config> {
self.config.clone()
}
pub fn with_config(mut self, config: &Config) -> Self {
self.config = Arc::new(config.clone());
self
}
pub fn add_plugin<T>(mut self, plugin: &T) -> Self
where
T: ActPlugin + Clone + 'static,
{
self.plugins.push(Arc::new(plugin.clone()));
self
}
pub fn set_plugins(mut self, plugins: Vec<Arc<dyn ActPlugin>>) -> Self {
self.plugins = plugins;
self
}
pub fn add_package<T>(mut self) -> Self
where
T: ActPackage + Clone + DeserializeOwned + 'static,
{
let package_register = ActPackageRegister::new::<T>();
self.packages.push(package_register);
self
}
pub fn add_resolver(&self, name: &str, resolver: Arc<dyn ConfigResolver>) {
self.runtime().register_resolver(name, resolver);
}
pub fn set_packages(mut self, packages: Vec<ActPackageRegister>) -> Self {
self.packages = packages;
self
}
pub fn set_resolvers(mut self, resolvers: Vec<(String, Arc<dyn ConfigResolver>)>) -> Self {
self.resolvers = resolvers;
self
}
pub fn set_store(mut self, store: Option<Arc<dyn KvStore>>) -> Self {
self.store = store;
self
}
pub fn executor(&self) -> Arc<Executor> {
Arc::new(Executor::new(&self.runtime()))
}
pub fn channel(&self) -> Arc<Channel> {
Arc::new(Channel::new(&self.runtime()))
}
pub fn channel_with_options(&self, matcher: &ChannelOptions) -> Arc<Channel> {
Arc::new(Channel::channel(&self.runtime(), matcher))
}
pub fn extender(&self) -> Arc<Extender> {
Arc::new(Extender::new(&self.runtime()))
}
pub fn builder() -> EngineBuilder {
EngineBuilder::new()
}
pub(crate) fn runtime(&self) -> Arc<Runtime> {
self.runtime.clone().expect("runtime not initialized")
}
pub async fn close(&self) {
self.runtime().close().await;
}
pub fn signal<T: Clone>(&self, init: T) -> Signal<T> {
Signal::new(init)
}
pub async fn start(mut self) -> crate::Result<Self> {
self.runtime = Some(Runtime::new(&self.config(), self.store.clone())?);
let rt = self.runtime();
let init = (|| async {
self.prepare().await?;
rt.event_loop();
rt.recover_actions().await?;
rt.resume().await?;
rt.init_retry_timer()?;
rt.init_trigger_timer();
Ok::<_, crate::ActError>(())
})()
.await;
if let Err(err) = init {
rt.close().await;
self.runtime = None;
return Err(err);
}
info!("engine started");
Ok(self)
}
async fn prepare(&self) -> crate::Result<()> {
for (name, resolver) in self.resolvers.iter() {
self.runtime().register_resolver(name, resolver.clone());
}
for plugin in self.plugins.iter() {
plugin.on_init(self)?;
}
package::init(self).await?;
for package_register in self.packages.iter() {
let meta = (package_register.meta)();
self.extender().register_package(&meta).await?;
if meta.run_as == crate::ActRunAs::Func {
self.runtime().package().register(meta.id, package_register);
}
}
Ok(())
}
}