Skip to main content

acts/package/
mod.rs

1pub mod core;
2pub mod transform;
3
4#[cfg(test)]
5mod tests;
6
7use crate::{
8    ActError, Config, Engine, Result, Vars, data,
9    scheduler::{Context, Runtime},
10    store::DbCollectionIden,
11};
12use dashmap::DashMap;
13use parking_lot::Mutex;
14use serde::{Deserialize, Serialize};
15use std::{fmt, sync::Arc};
16use tracing::debug;
17
18#[cfg(test)]
19pub use core::RunningMode;
20
21struct PackageEntry {
22    register: ActPackageRegister,
23    instance: Mutex<Option<Arc<dyn ActPackage>>>,
24}
25
26type SharedPackageEntry = Arc<PackageEntry>;
27
28impl fmt::Debug for PackageEntry {
29    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
30        f.debug_struct("PackageEntry")
31            .field("register", &self.register)
32            .finish()
33    }
34}
35
36#[derive(Clone)]
37pub struct Package {
38    packages: Arc<DashMap<String, SharedPackageEntry>>,
39}
40
41impl fmt::Debug for Package {
42    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
43        f.debug_struct("Package")
44            .field("packages", &self.packages)
45            .finish()
46    }
47}
48
49#[async_trait::async_trait]
50
51pub trait ActPackage: Send + Sync {
52    /// create package instance with config
53    fn new(config: &Config) -> Result<Self>
54    where
55        Self: Sized;
56    /// get package meta definition
57    fn definition() -> ActPackageDefinition
58    where
59        Self: Sized;
60    /// executing with task context
61    async fn execute(&self, _ctx: &Context, _params: &serde_json::Value) -> Result<Option<Vars>> {
62        Ok(None)
63    }
64    /// start with non-context, such as workflow event
65    async fn start(
66        &self,
67        _rt: &Arc<Runtime>,
68        _params: &serde_json::Value,
69        _options: &Vars,
70    ) -> Result<Option<Vars>> {
71        Ok(None)
72    }
73}
74
75#[derive(
76    Serialize,
77    Deserialize,
78    Debug,
79    Clone,
80    Copy,
81    Default,
82    PartialEq,
83    strum::AsRefStr,
84    strum::EnumString,
85)]
86#[serde(rename_all = "snake_case")]
87#[strum(serialize_all = "snake_case")]
88pub enum ActRunAs {
89    /// only used internally
90    Func,
91    /// interrupt request, need to response
92    #[default]
93    Irq,
94    /// message without response
95    Msg,
96}
97
98#[derive(
99    Serialize,
100    Deserialize,
101    Debug,
102    Clone,
103    Copy,
104    Default,
105    PartialEq,
106    strum::AsRefStr,
107    strum::EnumString,
108)]
109#[serde(rename_all = "snake_case")]
110#[strum(serialize_all = "snake_case")]
111pub enum ActPackageCatalog {
112    /// acts core packages
113    Core,
114
115    /// workflow event
116    Event,
117
118    /// workflow trace package
119    Output,
120
121    /// data transform
122    Transform,
123
124    /// form submition
125    Form,
126
127    /// AI related for LLMs
128    Ai,
129
130    /// the other applications to integrate into acts
131    /// such as Store, State, Observability, Pubsub
132    #[default]
133    App,
134}
135
136#[derive(Debug, Clone, Deserialize, Serialize)]
137pub struct ActPackageDefinition {
138    /// package id, used to identify the package
139    pub id: &'static str,
140
141    /// package simple name
142    pub name: &'static str,
143
144    /// package description
145    pub desc: &'static str,
146
147    /// icon name to display in the editor ui
148    pub icon: &'static str,
149
150    /// releated doc url to show the help
151    pub doc: &'static str,
152
153    /// package version
154    pub version: &'static str,
155
156    /// json schema for package params
157    pub schema: serde_json::Value,
158
159    /// extra options
160    #[serde(default)]
161    pub options: Option<serde_json::Value>,
162
163    /// package run as Irq, Msg or Func
164    /// Func is only used internally
165    pub run_as: ActRunAs,
166
167    /// package resources to the orgnize multiple resources
168    /// it is used for the editor ui to search and select the resources
169    /// each resource value can fill the special value into the UI
170    pub resources: Vec<ActResource>,
171
172    /// package catalog
173    pub catalog: ActPackageCatalog,
174}
175
176#[derive(Debug, Clone, Deserialize, Serialize)]
177pub struct ActResource {
178    pub name: String,
179    pub desc: String,
180    pub value: serde_json::Value,
181}
182
183#[derive(Debug, Clone)]
184pub struct ActPackageRegister {
185    pub meta: fn() -> ActPackageDefinition,
186    pub create: fn(config: &Config) -> Result<Arc<dyn ActPackage>>,
187}
188
189impl ActPackageRegister {
190    pub(crate) const fn new<T>() -> Self
191    where
192        T: ActPackage + 'static,
193    {
194        Self {
195            meta: T::definition,
196            create: (|config: &Config| {
197                // let meta = T::definition();
198                // jsonschema::validate(&meta.schema, params).map_err(|err| {
199                //     ActError::Package(format!(
200                //         "package({}) schema validation error: {}",
201                //         meta.id, err
202                //     ))
203                // })?;
204
205                let ret = T::new(config)?;
206                Ok(Arc::new(ret) as Arc<dyn ActPackage>)
207            }),
208        }
209    }
210}
211
212impl Default for Package {
213    fn default() -> Self {
214        Self::new()
215    }
216}
217
218impl Package {
219    pub fn new() -> Self {
220        Self {
221            packages: Arc::new(DashMap::new()),
222        }
223    }
224
225    pub fn register(&self, id: &str, register: &ActPackageRegister) {
226        // Replacing the entry also replaces its instance slot, so an old
227        // registration can never initialize a cache entry for a new one.
228        self.packages.insert(
229            id.to_string(),
230            Arc::new(PackageEntry {
231                register: register.clone(),
232                instance: Mutex::new(None),
233            }),
234        );
235    }
236
237    pub fn get(&self, id: &str) -> Option<ActPackageRegister> {
238        self.packages.get(id).map(|entry| entry.register.clone())
239    }
240
241    /// Return the cached instance for a registration, creating it on first use.
242    ///
243    /// Package constructors are intended to initialize reusable resources (or
244    /// leave them to async execution), so this avoids rebuilding a package for
245    /// every act/event and lets async packages keep one connection alive.
246    pub(crate) fn create(&self, id: &str, config: &Config) -> Result<Arc<dyn ActPackage>> {
247        let entry = self
248            .packages
249            .get(id)
250            .ok_or_else(|| ActError::Runtime(format!("cannot find package '{id}'")))?
251            .clone();
252        let mut instance = entry.instance.lock();
253
254        if let Some(package) = &*instance {
255            return Ok(package.clone());
256        }
257
258        let package = (entry.register.create)(config)?;
259        *instance = Some(package.clone());
260        Ok(package)
261    }
262}
263
264impl ActPackageDefinition {
265    pub fn into_data(&self) -> Result<data::Package> {
266        let pack = self.clone();
267        Ok(data::Package {
268            id: pack.id.to_string(),
269            name: pack.name.to_string(),
270            desc: pack.desc.to_string(),
271            icon: pack.icon.to_string(),
272            doc: pack.doc.to_string(),
273            version: pack.version.to_string(),
274            schema: pack.schema.to_string(),
275            options: pack.options.map(|v| v.to_string()),
276            run_as: pack.run_as,
277            resources: serde_json::to_string(&pack.resources)
278                .expect("cannot convert ActPackageMeta.resources to json"),
279            catalog: pack.catalog,
280            create_time: 0,
281            update_time: 0,
282            timestamp: 0,
283            built_in: false,
284            v: data::Package::version(),
285        })
286    }
287}
288
289inventory::collect!(ActPackageRegister);
290
291pub async fn init(engine: &Engine) -> Result<()> {
292    // The built-in packages are the engine's own registrations, not requests,
293    // so they are published as the unrestricted `system` principal — the
294    // deployment's policy governs its callers, not the engine's startup.
295    let executor = engine.executor(&crate::Principal::unrestricted());
296    for register in inventory::iter::<ActPackageRegister> {
297        let meta = (register.meta)();
298        debug!("package: {}", meta.name);
299
300        let mut pack = meta.into_data()?;
301        pack.built_in = true;
302        executor.pack().publish(&pack).await?;
303        engine.runtime().package().register(meta.id, register);
304    }
305    Ok(())
306}
307
308#[cfg(test)]
309mod cached_instance_tests {
310    use super::*;
311    use serde_json::json;
312    #[derive(Clone, Debug)]
313    struct CachedPackage;
314
315    #[async_trait::async_trait]
316    impl ActPackage for CachedPackage {
317        fn new(_config: &Config) -> Result<Self> {
318            Ok(Self)
319        }
320
321        fn definition() -> ActPackageDefinition {
322            ActPackageDefinition {
323                id: "test.cached",
324                name: "Cached",
325                desc: "",
326                icon: "",
327                doc: "",
328                version: "0.1.0",
329                schema: json!({}),
330                options: None,
331                run_as: ActRunAs::Func,
332                resources: vec![],
333                catalog: ActPackageCatalog::Core,
334            }
335        }
336    }
337
338    #[test]
339    fn create_reuses_instance_until_registration_is_replaced() {
340        let package = Package::new();
341        package.register(
342            CachedPackage::definition().id,
343            &ActPackageRegister::new::<CachedPackage>(),
344        );
345        let config = Config::default();
346
347        let first = package.create("test.cached", &config).unwrap();
348        let second = package.create("test.cached", &config).unwrap();
349        assert!(Arc::ptr_eq(&first, &second));
350
351        package.register(
352            CachedPackage::definition().id,
353            &ActPackageRegister::new::<CachedPackage>(),
354        );
355        let third = package.create("test.cached", &config).unwrap();
356        assert!(!Arc::ptr_eq(&first, &third));
357    }
358}