1use 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#[derive(Clone, Debug, Default)]
25pub struct PackageOptions {
26 pub output_override: Option<PathBuf>,
37}
38
39#[derive(Clone, Debug, PartialEq)]
41pub struct ProjectReport {
42 pub packages: Vec<PackagedWorkflow>,
44 pub excluded: Vec<ExcludedModule>,
46}
47
48#[derive(Clone, Debug, PartialEq)]
50pub struct PackagedWorkflow {
51 pub workflow_type: String,
53 pub output_path: PathBuf,
55 pub package: Package,
57 pub version: WorkflowVersion,
59}
60
61#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
63pub struct ExcludedModule {
64 pub module: String,
66 pub package: String,
68 pub reason: ExcludedReason,
70}
71
72#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
74#[serde(rename_all = "snake_case")]
75pub enum ExcludedReason {
76 SdkTestOnly,
78 DevDependency,
80}
81
82pub 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 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 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}