Skip to main content

aion_package/project/
assemble.rs

1//! `package_project` pipeline: config → discovery → build → verify-after-write.
2
3use std::{
4    collections::BTreeMap,
5    path::{Path, PathBuf},
6};
7
8use serde::Serialize;
9
10use super::{
11    config::{self, WorkflowConfig},
12    discover,
13    error::PackagingError,
14};
15use crate::{
16    BeamSet, CURRENT_FORMAT_VERSION, DeclaredActivity, ExtractionLimits, Manifest, ManifestVersion,
17    Package, PackageBuilder, WorkflowEntry, WorkflowVersion,
18};
19
20/// Options for packaging an already-built Gleam workflow project.
21///
22/// Construct via [`Default`] and assign fields, so call sites keep compiling
23/// when options are added.
24#[derive(Clone, Debug, Default)]
25pub struct PackageOptions {
26    /// Overrides the single workflow's output path, resolved against the
27    /// project root when relative. Packaging fails with
28    /// [`PackagingError::OutputOverrideAmbiguous`] when the project declares
29    /// more than one workflow.
30    ///
31    /// This is the caller's own path and is intentionally exempt from the
32    /// root confinement applied to `workflow.toml`-declared paths: it may
33    /// point anywhere, including outside the project root (the CLI resolves
34    /// `--out` against the invoker's working directory before passing it
35    /// here).
36    pub output_override: Option<PathBuf>,
37}
38
39/// Result of packaging every workflow a project declares.
40#[derive(Clone, Debug, PartialEq)]
41pub struct ProjectReport {
42    /// One built package per `[[workflow]]` entry, in declaration order.
43    pub packages: Vec<PackagedWorkflow>,
44    /// Modules excluded by the SDK test filter or the dependency-closure filter.
45    pub excluded: Vec<ExcludedModule>,
46}
47
48/// One workflow archive written and verified by [`package_project`].
49#[derive(Clone, Debug, PartialEq)]
50pub struct PackagedWorkflow {
51    /// Workflow type, identical to the manifest entry module.
52    pub workflow_type: String,
53    /// Absolute path of the written `.aion` archive.
54    pub output_path: PathBuf,
55    /// The archive re-loaded from disk after writing, proving integrity.
56    pub package: Package,
57    /// Canonical version record of the verified package.
58    pub version: WorkflowVersion,
59}
60
61/// A compiled module excluded from packaging, with provenance.
62#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
63pub struct ExcludedModule {
64    /// Logical module name that was excluded.
65    pub module: String,
66    /// Gleam package whose ebin provided the module.
67    pub package: String,
68    /// Why the module was excluded.
69    pub reason: ExcludedReason,
70}
71
72/// Reason a discovered compiled module was excluded from packaging.
73#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
74#[serde(rename_all = "snake_case")]
75pub enum ExcludedReason {
76    /// SDK test machinery from the `aion_flow` package's ebin.
77    SdkTestOnly,
78    /// Module of a package outside the production dependency closure.
79    DevDependency,
80}
81
82/// Packages every workflow declared by `<root>/workflow.toml`.
83///
84/// The project must already be built (`gleam build`); this function never
85/// spawns processes. The pipeline parses and validates the descriptor,
86/// discovers the production-closure compiled modules once, then writes one
87/// deterministic `.aion` archive per `[[workflow]]` entry. Every written
88/// archive is re-loaded through [`Package::load_from_path`] before this
89/// function returns, so the full read-path validation (integrity hash, format
90/// version, entry module) gates success.
91///
92/// All archives from one project share a single content hash (it covers beams
93/// only), while deployed entry names remain distinct per entry module. First
94/// party sources ship by default and never affect the hash.
95///
96/// Pure with respect to the environment: reads only under `root` (which is
97/// made absolute against the current directory once, up front), writes only
98/// the declared outputs, reads no environment variables, never prints, and
99/// blocks on synchronous filesystem I/O — async callers should wrap it in a
100/// blocking task.
101///
102/// The confinement is enforced, not assumed: every `workflow.toml`-declared
103/// path (`output`, `input_schema`, `output_schema`) is lexically normalized
104/// and must resolve inside `root` — absolute paths and `..` traversal that
105/// escapes the root fail with [`PackagingError::PathEscapesRoot`] before any
106/// file is touched. The sole exception is
107/// [`PackageOptions::output_override`]: that path belongs to the caller and
108/// may point anywhere, including outside the root.
109///
110/// # Errors
111///
112/// Returns [`PackagingError`] variants for missing or invalid `workflow.toml`
113/// descriptors, descriptor paths that are absolute or escape the project
114/// root, unreadable or non-JSON schema files, unbuilt projects, broken Gleam
115/// metadata, unresolved dependencies, duplicate or unreadable compiled
116/// modules, missing entry modules, output conflicts, ambiguous output
117/// overrides, and archive write (path-carrying) or verify-after-write
118/// failures.
119pub fn package_project(
120    root: &Path,
121    options: &PackageOptions,
122) -> Result<ProjectReport, PackagingError> {
123    let root = std::path::absolute(root).map_err(|source| PackagingError::ConfigRead {
124        path: root.to_path_buf(),
125        source,
126    })?;
127
128    let mut config = config::load_config(&root)?;
129    apply_output_override(&root, options, &mut config.workflows)?;
130
131    let discovered = discover::discover_modules(&root)?;
132    let beams = BeamSet::new(discovered.modules)?;
133    let source = if config.include_source {
134        discover::discover_sources(&root)?
135    } else {
136        BTreeMap::new()
137    };
138
139    let mut packages = Vec::with_capacity(config.workflows.len());
140    for workflow in &config.workflows {
141        if beams.get(&workflow.entry_module).is_none() {
142            return Err(PackagingError::EntryModuleNotFound {
143                module: workflow.entry_module.clone(),
144                searched: discovered.searched.clone(),
145            });
146        }
147        packages.push(build_workflow_package(workflow, &beams, &source)?);
148    }
149
150    Ok(ProjectReport {
151        packages,
152        excluded: discovered.excluded,
153    })
154}
155
156fn apply_output_override(
157    root: &Path,
158    options: &PackageOptions,
159    workflows: &mut [WorkflowConfig],
160) -> Result<(), PackagingError> {
161    let Some(output_override) = &options.output_override else {
162        return Ok(());
163    };
164    match workflows {
165        [workflow] => {
166            workflow.output_path = root.join(output_override);
167            Ok(())
168        }
169        _ => Err(PackagingError::OutputOverrideAmbiguous {
170            count: workflows.len(),
171        }),
172    }
173}
174
175fn build_workflow_package(
176    workflow: &WorkflowConfig,
177    beams: &BeamSet,
178    source: &BTreeMap<String, Vec<u8>>,
179) -> Result<PackagedWorkflow, PackagingError> {
180    let manifest = Manifest {
181        entry_module: workflow.entry_module.clone(),
182        entry_function: workflow.entry_function.clone(),
183        input_schema: workflow.input_schema.clone(),
184        output_schema: workflow.output_schema.clone(),
185        timeout: workflow.timeout,
186        activities: workflow
187            .activities
188            .iter()
189            .map(|activity_type| DeclaredActivity {
190                activity_type: activity_type.clone(),
191            })
192            .collect(),
193        version: ManifestVersion::new("unstamped"),
194        format_version: CURRENT_FORMAT_VERSION,
195        additional_workflows: workflow
196            .additional_workflows
197            .iter()
198            .map(|entry| WorkflowEntry {
199                workflow_type: entry.workflow_type.clone(),
200                entry_module: entry.entry_module.clone(),
201                entry_function: entry.entry_function.clone(),
202                input_schema: entry.input_schema.clone(),
203                output_schema: entry.output_schema.clone(),
204                timeout: entry.timeout,
205                internal: entry.internal,
206            })
207            .collect(),
208    };
209
210    PackageBuilder::with_source(manifest, beams.clone(), source.clone())
211        .write_to_path(&workflow.output_path)
212        .map_err(|source| PackagingError::OutputWrite {
213            path: workflow.output_path.clone(),
214            source,
215        })?;
216    // Trusted local input: the archive was written by this process moments
217    // ago, so extraction runs unbounded.
218    let package = Package::load_from_path(&workflow.output_path, ExtractionLimits::unbounded())?;
219
220    Ok(PackagedWorkflow {
221        workflow_type: workflow.entry_module.clone(),
222        output_path: workflow.output_path.clone(),
223        version: package.version_record(),
224        package,
225    })
226}
227
228#[cfg(test)]
229mod tests {
230    use std::{fs, path::PathBuf, time::Duration};
231
232    use serde_json::json;
233
234    use super::{ExcludedModule, ExcludedReason, PackageOptions, package_project};
235    use crate::{PackageError, project::error::PackagingError, project::fixture};
236
237    type TestResult = Result<(), Box<dyn std::error::Error>>;
238
239    const TWO_WORKFLOW_TOML: &str = r#"[[workflow]]
240entry_module = "demo"
241entry_function = "run"
242timeout_seconds = 30
243input_schema = "schemas/input.json"
244output_schema = "schemas/output.json"
245activities = ["greet"]
246
247[[workflow]]
248entry_module = "demo@nested"
249entry_function = "start"
250timeout_seconds = 60
251input_schema = "schemas/input.json"
252output_schema = "schemas/output.json"
253activities = []
254"#;
255
256    #[test]
257    fn packaged_workflow_round_trips_manifest_and_hash() -> TestResult {
258        let root = fixture::synthetic_built_project("assemble-happy")?;
259        let report = package_project(&root, &PackageOptions::default());
260        let reloaded = report
261            .as_ref()
262            .ok()
263            .map(|report| report.packages[0].output_path.clone())
264            .map(|path| crate::Package::load_from_path(path, crate::ExtractionLimits::unbounded()));
265        fs::remove_dir_all(&root)?;
266        let report = report?;
267
268        assert_eq!(report.packages.len(), 1);
269        let packaged = &report.packages[0];
270        assert_eq!(packaged.workflow_type, "demo");
271        assert!(packaged.output_path.is_absolute());
272        assert_eq!(
273            packaged
274                .output_path
275                .file_name()
276                .and_then(|name| name.to_str()),
277            Some("demo.aion")
278        );
279        let manifest = packaged.package.manifest();
280        assert_eq!(manifest.entry_module, "demo");
281        assert_eq!(manifest.entry_function, "run");
282        assert_eq!(manifest.timeout, Duration::from_secs(30));
283        assert_eq!(manifest.input_schema, json!({ "type": "object" }));
284        assert_eq!(manifest.output_schema, json!(true));
285        assert_eq!(manifest.activities.len(), 1);
286        assert_eq!(manifest.activities[0].activity_type, "greet");
287        assert_eq!(
288            manifest.version.as_str(),
289            packaged.package.content_hash().to_string()
290        );
291        assert_eq!(packaged.version, packaged.package.version_record());
292        let reloaded = reloaded.ok_or("report failed")??;
293        assert_eq!(&reloaded, &packaged.package);
294        Ok(())
295    }
296
297    #[test]
298    fn exclusions_and_sources_are_reported_and_shipped() -> TestResult {
299        let root = fixture::synthetic_built_project("assemble-exclusions")?;
300        let report = package_project(&root, &PackageOptions::default());
301        fs::remove_dir_all(&root)?;
302        let report = report?;
303
304        let expected_exclusions = vec![
305            ExcludedModule {
306                module: "dev_only".to_owned(),
307                package: "dev_only".to_owned(),
308                reason: ExcludedReason::DevDependency,
309            },
310            ExcludedModule {
311                module: "aion@testing".to_owned(),
312                package: "aion_flow".to_owned(),
313                reason: ExcludedReason::SdkTestOnly,
314            },
315            ExcludedModule {
316                module: "aion@testing@mock".to_owned(),
317                package: "aion_flow".to_owned(),
318                reason: ExcludedReason::SdkTestOnly,
319            },
320            ExcludedModule {
321                module: "aion_flow_ffi".to_owned(),
322                package: "aion_flow".to_owned(),
323                reason: ExcludedReason::SdkTestOnly,
324            },
325        ];
326        assert_eq!(report.excluded, expected_exclusions);
327
328        let package = &report.packages[0].package;
329        let source_names: Vec<&str> = package.source().keys().map(String::as_str).collect();
330        assert_eq!(source_names, vec!["demo", "demo/nested"]);
331        let beam_names: Vec<&str> = package
332            .beams()
333            .iter()
334            .map(crate::BeamModule::name)
335            .collect();
336        assert_eq!(
337            beam_names,
338            vec!["aion_flow", "demo", "demo@nested", "dep_a", "dep_b"]
339        );
340        Ok(())
341    }
342
343    #[test]
344    fn missing_entry_module_returns_entry_module_not_found() -> TestResult {
345        let root = fixture::synthetic_built_project("assemble-ghost-entry")?;
346        let descriptor = fixture::DEMO_WORKFLOW_TOML.replace("\"demo\"", "\"ghost\"");
347        fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
348        let result = package_project(&root, &PackageOptions::default());
349        fs::remove_dir_all(&root)?;
350
351        assert!(matches!(
352            result,
353            Err(PackagingError::EntryModuleNotFound { module, searched })
354                if module == "ghost" && searched.ends_with("build/dev/erlang")
355        ));
356        Ok(())
357    }
358
359    #[test]
360    fn explicit_output_field_is_respected() -> TestResult {
361        let root = fixture::synthetic_built_project("assemble-explicit-output")?;
362        let descriptor = format!(
363            "{}output = \"custom-name.aion\"\n",
364            fixture::DEMO_WORKFLOW_TOML
365        );
366        fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
367        let report = package_project(&root, &PackageOptions::default());
368        let written = root.join("custom-name.aion").is_file();
369        fs::remove_dir_all(&root)?;
370        let report = report?;
371
372        assert!(written);
373        assert_eq!(
374            report.packages[0]
375                .output_path
376                .file_name()
377                .and_then(|name| name.to_str()),
378            Some("custom-name.aion")
379        );
380        Ok(())
381    }
382
383    #[test]
384    fn output_override_applies_to_single_workflow_project() -> TestResult {
385        let root = fixture::synthetic_built_project("assemble-override")?;
386        let options = PackageOptions {
387            output_override: Some(PathBuf::from("override.aion")),
388        };
389        let report = package_project(&root, &options);
390        let written = root.join("override.aion").is_file();
391        let derived_absent = !root.join("demo.aion").exists();
392        fs::remove_dir_all(&root)?;
393        let report = report?;
394
395        assert!(written);
396        assert!(derived_absent);
397        assert_eq!(
398            report.packages[0]
399                .output_path
400                .file_name()
401                .and_then(|name| name.to_str()),
402            Some("override.aion")
403        );
404        Ok(())
405    }
406
407    #[test]
408    fn output_write_failure_names_the_output_path() -> TestResult {
409        let root = fixture::synthetic_built_project("assemble-missing-dir")?;
410        let descriptor = format!(
411            "{}output = \"missing-dir/demo.aion\"\n",
412            fixture::DEMO_WORKFLOW_TOML
413        );
414        fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
415        let result = package_project(&root, &PackageOptions::default());
416        fs::remove_dir_all(&root)?;
417
418        let expected = root.join("missing-dir/demo.aion");
419        let Err(error) = result else {
420            return Err("write into a missing directory unexpectedly succeeded".into());
421        };
422        assert!(
423            matches!(
424                &error,
425                PackagingError::OutputWrite { path, .. } if *path == expected
426            ),
427            "error does not carry the output path: {error:?}"
428        );
429        assert!(
430            error.to_string().contains(&expected.display().to_string()),
431            "message does not name the output path: {error}"
432        );
433        Ok(())
434    }
435
436    #[test]
437    fn output_override_may_point_outside_root_via_dotdot() -> TestResult {
438        // The exemption under test: workflow.toml paths are confined to the
439        // root, but the caller's `output_override` may point anywhere.
440        let root = fixture::synthetic_built_project("assemble-override-outside")?;
441        let outside_name = format!("aion-override-outside-{}.aion", std::process::id());
442        let options = PackageOptions {
443            output_override: Some(PathBuf::from(format!("../{outside_name}"))),
444        };
445        let report = package_project(&root, &options);
446        let outside = std::env::temp_dir().join(&outside_name);
447        let written = outside.is_file();
448        fs::remove_dir_all(&root)?;
449        if written {
450            fs::remove_file(&outside)?;
451        }
452        let report = report?;
453
454        assert!(written, "override outside the root was not written");
455        assert_eq!(
456            report.packages[0]
457                .output_path
458                .file_name()
459                .and_then(|name| name.to_str()),
460            Some(outside_name.as_str())
461        );
462        Ok(())
463    }
464
465    #[test]
466    fn output_override_with_multiple_workflows_is_ambiguous() -> TestResult {
467        let root = fixture::synthetic_built_project("assemble-override-multi")?;
468        fixture::write_file(&root, "workflow.toml", TWO_WORKFLOW_TOML.as_bytes())?;
469        let options = PackageOptions {
470            output_override: Some(PathBuf::from("override.aion")),
471        };
472        let result = package_project(&root, &options);
473        fs::remove_dir_all(&root)?;
474
475        assert!(matches!(
476            result,
477            Err(PackagingError::OutputOverrideAmbiguous { count: 2 })
478        ));
479        Ok(())
480    }
481
482    #[test]
483    fn multi_workflow_project_shares_hash_with_distinct_deployed_entries() -> TestResult {
484        let root = fixture::synthetic_built_project("assemble-multi")?;
485        fixture::write_file(&root, "workflow.toml", TWO_WORKFLOW_TOML.as_bytes())?;
486        let report = package_project(&root, &PackageOptions::default());
487        fs::remove_dir_all(&root)?;
488        let report = report?;
489
490        assert_eq!(report.packages.len(), 2);
491        let first = &report.packages[0];
492        let second = &report.packages[1];
493        assert_eq!(first.workflow_type, "demo");
494        assert_eq!(second.workflow_type, "demo@nested");
495        assert_eq!(first.package.content_hash(), second.package.content_hash());
496        assert_ne!(
497            first.package.deployed_entry_module(),
498            second.package.deployed_entry_module()
499        );
500        assert_ne!(first.output_path, second.output_path);
501        Ok(())
502    }
503
504    #[test]
505    fn user_module_with_reserved_name_fails_typed() -> TestResult {
506        let root = fixture::synthetic_built_project("assemble-reserved")?;
507        fixture::write_file(
508            &root,
509            "build/dev/erlang/demo/ebin/aion_flow_ffi.beam",
510            b"user-owned-bytes",
511        )?;
512        let result = package_project(&root, &PackageOptions::default());
513        fs::remove_dir_all(&root)?;
514
515        assert!(matches!(
516            result,
517            Err(PackagingError::Package(PackageError::ReservedModuleName { module }))
518                if module == "aion_flow_ffi"
519        ));
520        Ok(())
521    }
522
523    #[test]
524    fn repackaging_produces_identical_archive_bytes() -> TestResult {
525        let root = fixture::synthetic_built_project("assemble-det-1")?;
526        let first_report = package_project(&root, &PackageOptions::default());
527        let first_bytes = fs::read(root.join("demo.aion"));
528        let second_report = package_project(&root, &PackageOptions::default());
529        let second_bytes = fs::read(root.join("demo.aion"));
530        fs::remove_dir_all(&root)?;
531        first_report?;
532        second_report?;
533
534        let first_bytes = first_bytes?;
535        assert!(!first_bytes.is_empty());
536        assert_eq!(first_bytes, second_bytes?);
537        Ok(())
538    }
539
540    #[test]
541    fn source_inclusion_changes_bytes_but_never_the_version() -> TestResult {
542        let root = fixture::synthetic_built_project("assemble-det-3")?;
543        let with_source = package_project(&root, &PackageOptions::default());
544        let with_source_bytes = fs::read(root.join("demo.aion"));
545        let descriptor = format!(
546            "[package]\ninclude_source = false\n\n{}",
547            fixture::DEMO_WORKFLOW_TOML
548        );
549        fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
550        let without_source = package_project(&root, &PackageOptions::default());
551        let without_source_bytes = fs::read(root.join("demo.aion"));
552        fs::remove_dir_all(&root)?;
553        let with_source = with_source?;
554        let without_source = without_source?;
555
556        assert!(!with_source.packages[0].package.source().is_empty());
557        assert!(without_source.packages[0].package.source().is_empty());
558        assert_ne!(with_source_bytes?, without_source_bytes?);
559        assert_eq!(
560            with_source.packages[0].package.content_hash(),
561            without_source.packages[0].package.content_hash()
562        );
563        assert_eq!(
564            with_source.packages[0].package.manifest().version,
565            without_source.packages[0].package.manifest().version
566        );
567        Ok(())
568    }
569}