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: Some(workflow.timeout),
189 activities: workflow
190 .activities
191 .iter()
192 .map(|activity_type| DeclaredActivity {
193 activity_type: activity_type.clone(),
194 })
195 .collect(),
196 version: ManifestVersion::new("unstamped"),
197 format_version: CURRENT_FORMAT_VERSION,
198 additional_workflows: workflow
199 .additional_workflows
200 .iter()
201 .map(|entry| WorkflowEntry {
202 workflow_type: entry.workflow_type.clone(),
203 entry_module: entry.entry_module.clone(),
204 entry_function: entry.entry_function.clone(),
205 input_schema: entry.input_schema.clone(),
206 output_schema: entry.output_schema.clone(),
207 timeout: entry.timeout,
208 internal: entry.internal,
209 })
210 .collect(),
211 };
212
213 PackageBuilder::with_source(manifest, beams.clone(), source.clone())
214 .write_to_path(&workflow.output_path)
215 .map_err(|source| PackagingError::OutputWrite {
216 path: workflow.output_path.clone(),
217 source,
218 })?;
219 let package = Package::load_from_path(&workflow.output_path, ExtractionLimits::unbounded())?;
222
223 Ok(PackagedWorkflow {
224 workflow_type: workflow.entry_module.clone(),
225 output_path: workflow.output_path.clone(),
226 version: package.version_record(),
227 package,
228 })
229}
230
231#[cfg(test)]
232mod tests {
233 use std::{fs, path::PathBuf, time::Duration};
234
235 use serde_json::json;
236
237 use super::{ExcludedModule, ExcludedReason, PackageOptions, package_project};
238 use crate::{PackageError, project::error::PackagingError, project::fixture};
239
240 type TestResult = Result<(), Box<dyn std::error::Error>>;
241
242 const TWO_WORKFLOW_TOML: &str = r#"[[workflow]]
243entry_module = "demo"
244entry_function = "run"
245timeout_seconds = 30
246input_schema = "schemas/input.json"
247output_schema = "schemas/output.json"
248activities = ["greet"]
249
250[[workflow]]
251entry_module = "demo@nested"
252entry_function = "start"
253timeout_seconds = 60
254input_schema = "schemas/input.json"
255output_schema = "schemas/output.json"
256activities = []
257"#;
258
259 #[test]
260 fn packaged_workflow_round_trips_manifest_and_hash() -> TestResult {
261 let root = fixture::synthetic_built_project("assemble-happy")?;
262 let report = package_project(&root, &PackageOptions::default());
263 let reloaded = report
264 .as_ref()
265 .ok()
266 .map(|report| report.packages[0].output_path.clone())
267 .map(|path| crate::Package::load_from_path(path, crate::ExtractionLimits::unbounded()));
268 fs::remove_dir_all(&root)?;
269 let report = report?;
270
271 assert_eq!(report.packages.len(), 1);
272 let packaged = &report.packages[0];
273 assert_eq!(packaged.workflow_type, "demo");
274 assert!(packaged.output_path.is_absolute());
275 assert_eq!(
276 packaged
277 .output_path
278 .file_name()
279 .and_then(|name| name.to_str()),
280 Some("demo.aion")
281 );
282 let manifest = packaged.package.manifest();
283 assert_eq!(manifest.entry_module, "demo");
284 assert_eq!(manifest.entry_function, "run");
285 assert_eq!(manifest.timeout, Some(Duration::from_secs(30)));
286 assert_eq!(manifest.input_schema, json!({ "type": "object" }));
287 assert_eq!(manifest.output_schema, json!(true));
288 assert_eq!(manifest.activities.len(), 1);
289 assert_eq!(manifest.activities[0].activity_type, "greet");
290 assert_eq!(
291 manifest.version.as_str(),
292 packaged.package.content_hash().to_string()
293 );
294 assert_eq!(packaged.version, packaged.package.version_record());
295 let reloaded = reloaded.ok_or("report failed")??;
296 assert_eq!(&reloaded, &packaged.package);
297 Ok(())
298 }
299
300 #[test]
301 fn exclusions_and_sources_are_reported_and_shipped() -> TestResult {
302 let root = fixture::synthetic_built_project("assemble-exclusions")?;
303 let report = package_project(&root, &PackageOptions::default());
304 fs::remove_dir_all(&root)?;
305 let report = report?;
306
307 let expected_exclusions = vec![
308 ExcludedModule {
309 module: "dev_only".to_owned(),
310 package: "dev_only".to_owned(),
311 reason: ExcludedReason::DevDependency,
312 },
313 ExcludedModule {
314 module: "aion@testing".to_owned(),
315 package: "aion_flow".to_owned(),
316 reason: ExcludedReason::SdkTestOnly,
317 },
318 ExcludedModule {
319 module: "aion@testing@mock".to_owned(),
320 package: "aion_flow".to_owned(),
321 reason: ExcludedReason::SdkTestOnly,
322 },
323 ExcludedModule {
324 module: "aion_flow_ffi".to_owned(),
325 package: "aion_flow".to_owned(),
326 reason: ExcludedReason::SdkTestOnly,
327 },
328 ];
329 assert_eq!(report.excluded, expected_exclusions);
330
331 let package = &report.packages[0].package;
332 let source_names: Vec<&str> = package.source().keys().map(String::as_str).collect();
333 assert_eq!(source_names, vec!["demo", "demo/nested"]);
334 let beam_names: Vec<&str> = package
335 .beams()
336 .iter()
337 .map(crate::BeamModule::name)
338 .collect();
339 assert_eq!(
340 beam_names,
341 vec!["aion_flow", "demo", "demo@nested", "dep_a", "dep_b"]
342 );
343 Ok(())
344 }
345
346 #[test]
347 fn missing_entry_module_returns_entry_module_not_found() -> TestResult {
348 let root = fixture::synthetic_built_project("assemble-ghost-entry")?;
349 let descriptor = fixture::DEMO_WORKFLOW_TOML.replace("\"demo\"", "\"ghost\"");
350 fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
351 let result = package_project(&root, &PackageOptions::default());
352 fs::remove_dir_all(&root)?;
353
354 assert!(matches!(
355 result,
356 Err(PackagingError::EntryModuleNotFound { module, searched })
357 if module == "ghost" && searched.ends_with("build/dev/erlang")
358 ));
359 Ok(())
360 }
361
362 #[test]
363 fn explicit_output_field_is_respected() -> TestResult {
364 let root = fixture::synthetic_built_project("assemble-explicit-output")?;
365 let descriptor = format!(
366 "{}output = \"custom-name.aion\"\n",
367 fixture::DEMO_WORKFLOW_TOML
368 );
369 fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
370 let report = package_project(&root, &PackageOptions::default());
371 let written = root.join("custom-name.aion").is_file();
372 fs::remove_dir_all(&root)?;
373 let report = report?;
374
375 assert!(written);
376 assert_eq!(
377 report.packages[0]
378 .output_path
379 .file_name()
380 .and_then(|name| name.to_str()),
381 Some("custom-name.aion")
382 );
383 Ok(())
384 }
385
386 #[test]
387 fn output_override_applies_to_single_workflow_project() -> TestResult {
388 let root = fixture::synthetic_built_project("assemble-override")?;
389 let options = PackageOptions {
390 output_override: Some(PathBuf::from("override.aion")),
391 };
392 let report = package_project(&root, &options);
393 let written = root.join("override.aion").is_file();
394 let derived_absent = !root.join("demo.aion").exists();
395 fs::remove_dir_all(&root)?;
396 let report = report?;
397
398 assert!(written);
399 assert!(derived_absent);
400 assert_eq!(
401 report.packages[0]
402 .output_path
403 .file_name()
404 .and_then(|name| name.to_str()),
405 Some("override.aion")
406 );
407 Ok(())
408 }
409
410 #[test]
411 fn output_write_failure_names_the_output_path() -> TestResult {
412 let root = fixture::synthetic_built_project("assemble-missing-dir")?;
413 let descriptor = format!(
414 "{}output = \"missing-dir/demo.aion\"\n",
415 fixture::DEMO_WORKFLOW_TOML
416 );
417 fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
418 let result = package_project(&root, &PackageOptions::default());
419 fs::remove_dir_all(&root)?;
420
421 let expected = root.join("missing-dir/demo.aion");
422 let Err(error) = result else {
423 return Err("write into a missing directory unexpectedly succeeded".into());
424 };
425 assert!(
426 matches!(
427 &error,
428 PackagingError::OutputWrite { path, .. } if *path == expected
429 ),
430 "error does not carry the output path: {error:?}"
431 );
432 assert!(
433 error.to_string().contains(&expected.display().to_string()),
434 "message does not name the output path: {error}"
435 );
436 Ok(())
437 }
438
439 #[test]
440 fn output_override_may_point_outside_root_via_dotdot() -> TestResult {
441 let root = fixture::synthetic_built_project("assemble-override-outside")?;
444 let outside_name = format!("aion-override-outside-{}.aion", std::process::id());
445 let options = PackageOptions {
446 output_override: Some(PathBuf::from(format!("../{outside_name}"))),
447 };
448 let report = package_project(&root, &options);
449 let outside = std::env::temp_dir().join(&outside_name);
450 let written = outside.is_file();
451 fs::remove_dir_all(&root)?;
452 if written {
453 fs::remove_file(&outside)?;
454 }
455 let report = report?;
456
457 assert!(written, "override outside the root was not written");
458 assert_eq!(
459 report.packages[0]
460 .output_path
461 .file_name()
462 .and_then(|name| name.to_str()),
463 Some(outside_name.as_str())
464 );
465 Ok(())
466 }
467
468 #[test]
469 fn output_override_with_multiple_workflows_is_ambiguous() -> TestResult {
470 let root = fixture::synthetic_built_project("assemble-override-multi")?;
471 fixture::write_file(&root, "workflow.toml", TWO_WORKFLOW_TOML.as_bytes())?;
472 let options = PackageOptions {
473 output_override: Some(PathBuf::from("override.aion")),
474 };
475 let result = package_project(&root, &options);
476 fs::remove_dir_all(&root)?;
477
478 assert!(matches!(
479 result,
480 Err(PackagingError::OutputOverrideAmbiguous { count: 2 })
481 ));
482 Ok(())
483 }
484
485 #[test]
486 fn multi_workflow_project_shares_hash_with_distinct_deployed_entries() -> TestResult {
487 let root = fixture::synthetic_built_project("assemble-multi")?;
488 fixture::write_file(&root, "workflow.toml", TWO_WORKFLOW_TOML.as_bytes())?;
489 let report = package_project(&root, &PackageOptions::default());
490 fs::remove_dir_all(&root)?;
491 let report = report?;
492
493 assert_eq!(report.packages.len(), 2);
494 let first = &report.packages[0];
495 let second = &report.packages[1];
496 assert_eq!(first.workflow_type, "demo");
497 assert_eq!(second.workflow_type, "demo@nested");
498 assert_eq!(first.package.content_hash(), second.package.content_hash());
499 assert_ne!(
500 first.package.deployed_entry_module(),
501 second.package.deployed_entry_module()
502 );
503 assert_ne!(first.output_path, second.output_path);
504 Ok(())
505 }
506
507 #[test]
508 fn user_module_with_reserved_name_fails_typed() -> TestResult {
509 let root = fixture::synthetic_built_project("assemble-reserved")?;
510 fixture::write_file(
511 &root,
512 "build/dev/erlang/demo/ebin/aion_flow_ffi.beam",
513 b"user-owned-bytes",
514 )?;
515 let result = package_project(&root, &PackageOptions::default());
516 fs::remove_dir_all(&root)?;
517
518 assert!(matches!(
519 result,
520 Err(PackagingError::Package(PackageError::ReservedModuleName { module }))
521 if module == "aion_flow_ffi"
522 ));
523 Ok(())
524 }
525
526 #[test]
527 fn repackaging_produces_identical_archive_bytes() -> TestResult {
528 let root = fixture::synthetic_built_project("assemble-det-1")?;
529 let first_report = package_project(&root, &PackageOptions::default());
530 let first_bytes = fs::read(root.join("demo.aion"));
531 let second_report = package_project(&root, &PackageOptions::default());
532 let second_bytes = fs::read(root.join("demo.aion"));
533 fs::remove_dir_all(&root)?;
534 first_report?;
535 second_report?;
536
537 let first_bytes = first_bytes?;
538 assert!(!first_bytes.is_empty());
539 assert_eq!(first_bytes, second_bytes?);
540 Ok(())
541 }
542
543 #[test]
544 fn source_inclusion_changes_bytes_but_never_the_version() -> TestResult {
545 let root = fixture::synthetic_built_project("assemble-det-3")?;
546 let with_source = package_project(&root, &PackageOptions::default());
547 let with_source_bytes = fs::read(root.join("demo.aion"));
548 let descriptor = format!(
549 "[package]\ninclude_source = false\n\n{}",
550 fixture::DEMO_WORKFLOW_TOML
551 );
552 fixture::write_file(&root, "workflow.toml", descriptor.as_bytes())?;
553 let without_source = package_project(&root, &PackageOptions::default());
554 let without_source_bytes = fs::read(root.join("demo.aion"));
555 fs::remove_dir_all(&root)?;
556 let with_source = with_source?;
557 let without_source = without_source?;
558
559 assert!(!with_source.packages[0].package.source().is_empty());
560 assert!(without_source.packages[0].package.source().is_empty());
561 assert_ne!(with_source_bytes?, without_source_bytes?);
562 assert_eq!(
563 with_source.packages[0].package.content_hash(),
564 without_source.packages[0].package.content_hash()
565 );
566 assert_eq!(
567 with_source.packages[0].package.manifest().version,
568 without_source.packages[0].package.manifest().version
569 );
570 Ok(())
571 }
572}