diff --git a/crates/nemoclaw-e2e/tests/deployment.rs b/crates/nemoclaw-e2e/tests/deployment.rs index ffd5f4995c8..344c5d8af5f 100644 --- a/crates/nemoclaw-e2e/tests/deployment.rs +++ b/crates/nemoclaw-e2e/tests/deployment.rs @@ -939,7 +939,19 @@ async fn destroy_does_not_require_the_inference_credential_or_rewrite_its_refere .apply(&document, &cancel) .await .unwrap(); - let deployment = Deployment::new(directory.path(), &bundle); + let initializations = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let counted = initializations.clone(); + let deployment = Deployment::new(directory.path(), &bundle).with_progress(std::sync::Arc::new( + move |event| { + if let nemoclaw_sdk::Progress::Completed { + operation: "tofu.init", + .. + } = event + { + counted.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + } + }, + )); let effects = fixture.state.lock().unwrap().effects; assert_eq!( deployment @@ -951,10 +963,17 @@ async fn destroy_does_not_require_the_inference_credential_or_rewrite_its_refere 4 ); assert_eq!(fixture.state.lock().unwrap().effects, effects); + // Teardown initializes the stage once to read its bindings and once for + // its teardown graph; it does not read the same state again. + assert_eq!( + initializations.swap(0, std::sync::atomic::Ordering::SeqCst), + 2 + ); assert_eq!( deployment.destroy(&cancel).await.unwrap().outcome, Outcome::Destroyed ); + assert_eq!(initializations.load(std::sync::atomic::Ordering::SeqCst), 2); let record: serde_json::Value = serde_json::from_slice(&fs::read(directory.path().join("intent.json")).unwrap()).unwrap(); assert_eq!(record["document"], serde_json::to_value(&document).unwrap()); diff --git a/crates/nemoclaw-sdk/src/bundle/mod.rs b/crates/nemoclaw-sdk/src/bundle/mod.rs index 98abea72ab4..26e58742c3d 100644 --- a/crates/nemoclaw-sdk/src/bundle/mod.rs +++ b/crates/nemoclaw-sdk/src/bundle/mod.rs @@ -43,24 +43,13 @@ impl Bundle { return Err(Error::Bundle("bundle is incomplete")); } } - for (name, digest) in &manifest.files { - if name.is_empty() - || Path::new(name) - .components() - .any(|c| !matches!(c, Component::Normal(_))) - || name.contains('\\') - { - return Err(Error::Bundle("invalid bundle path")); - } - let path = directory.join(name); - if !fs::symlink_metadata(&path) - .map_err(|_| Error::Bundle("bundle file is unavailable"))? - .is_file() - || hash_file(&path)? != *digest - { - return Err(Error::Bundle("bundle file integrity check failed")); - } - } + let files: Vec<_> = manifest.files.iter().collect(); + // Hash files concurrently; the first failure in manifest order decides + // the error, as it would when checking them one at a time. + let checks = verify_concurrently(&files, |(name, digest)| { + verify_file(&directory, name, digest) + }); + checks.into_iter().collect::>()?; Ok(Self { directory, manifest, @@ -70,6 +59,57 @@ impl Bundle { self.directory.join("libexec").join(executable("tofu")) } } +fn verify_file(directory: &Path, name: &str, digest: &str) -> Result<(), Error> { + if name.is_empty() + || Path::new(name) + .components() + .any(|c| !matches!(c, Component::Normal(_))) + || name.contains('\\') + { + return Err(Error::Bundle("invalid bundle path")); + } + let path = directory.join(name); + if !fs::symlink_metadata(&path) + .map_err(|_| Error::Bundle("bundle file is unavailable"))? + .is_file() + || hash_file(&path)? != digest + { + return Err(Error::Bundle("bundle file integrity check failed")); + } + Ok(()) +} +/// Run `check` on every item across a few threads, returning results in item order. +fn verify_concurrently(items: &[T], check: impl Fn(&T) -> R + Sync) -> Vec { + let workers = std::thread::available_parallelism() + .map_or(1, std::num::NonZeroUsize::get) + .min(items.len()); + if workers <= 1 { + return items.iter().map(check).collect(); + } + let next = std::sync::atomic::AtomicUsize::new(0); + let mut results: Vec<(usize, R)> = std::thread::scope(|scope| { + let handles: Vec<_> = (0..workers) + .map(|_| { + scope.spawn(|| { + let mut done = Vec::new(); + loop { + let index = next.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let Some(item) = items.get(index) else { + return done; + }; + done.push((index, check(item))); + } + }) + }) + .collect(); + handles + .into_iter() + .flat_map(|handle| handle.join().expect("bundle verification thread")) + .collect() + }); + results.sort_by_key(|(index, _)| *index); + results.into_iter().map(|(_, result)| result).collect() +} /// Metadata that changes when a file is written or replaced. #[derive(Clone, Debug, PartialEq, Eq)] struct Stamp { diff --git a/crates/nemoclaw-sdk/src/bundle/tests.rs b/crates/nemoclaw-sdk/src/bundle/tests.rs index e0f76abeb81..e9348995eb9 100644 --- a/crates/nemoclaw-sdk/src/bundle/tests.rs +++ b/crates/nemoclaw-sdk/src/bundle/tests.rs @@ -162,6 +162,61 @@ fn fixture_bundle(directory: &Path) { .unwrap(); } +#[test] +fn every_listed_file_is_verified_and_the_first_listed_failure_is_reported() { + let names: Vec = { + let directory = tempfile::tempdir().unwrap(); + fixture_bundle(directory.path()); + Bundle::open(directory.path()) + .unwrap() + .manifest + .files + .into_keys() + .collect() + }; + assert!(names.len() > 2); + let failure = |directory: &Path| match Bundle::open(directory) { + Err(Error::Bundle(message)) => message, + other => panic!("expected a bundle error, got {other:?}"), + }; + for name in &names { + let directory = tempfile::tempdir().unwrap(); + fixture_bundle(directory.path()); + fs::write(directory.path().join(name), b"tampered binary").unwrap(); + assert_eq!( + failure(directory.path()), + "bundle file integrity check failed", + "{name}" + ); + fs::remove_file(directory.path().join(name)).unwrap(); + assert_eq!( + failure(directory.path()), + "bundle file is unavailable", + "{name}" + ); + fs::create_dir(directory.path().join(name)).unwrap(); + assert_eq!( + failure(directory.path()), + "bundle file integrity check failed", + "{name}" + ); + } + // With several failures, the first file in manifest order decides the error. + let directory = tempfile::tempdir().unwrap(); + fixture_bundle(directory.path()); + fs::remove_file(directory.path().join(&names[0])).unwrap(); + fs::write(directory.path().join(names.last().unwrap()), b"tampered").unwrap(); + assert_eq!(failure(directory.path()), "bundle file is unavailable"); + let directory = tempfile::tempdir().unwrap(); + fixture_bundle(directory.path()); + fs::write(directory.path().join(&names[0]), b"tampered").unwrap(); + fs::remove_file(directory.path().join(names.last().unwrap())).unwrap(); + assert_eq!( + failure(directory.path()), + "bundle file integrity check failed" + ); +} + #[test] fn an_unchanged_bundle_is_hashed_once_and_any_written_file_is_hashed_again() { let directory = tempfile::tempdir().unwrap(); diff --git a/crates/nemoclaw-sdk/src/deployment/runtime/teardown.rs b/crates/nemoclaw-sdk/src/deployment/runtime/teardown.rs index 1f4bd083fc2..e84289fb637 100644 --- a/crates/nemoclaw-sdk/src/deployment/runtime/teardown.rs +++ b/crates/nemoclaw-sdk/src/deployment/runtime/teardown.rs @@ -106,14 +106,28 @@ impl Deployment { let mut stages = Vec::new(); if !record.root_destroyed() { let (changes, planned) = operation - .plan_teardown_stage(&bundle, &store, &record, false, &root_graph, cancel) + .plan_teardown_stage( + &bundle, + &store, + &record, + (false, &bindings), + &root_graph, + cancel, + ) .await?; result.changes.extend(changes); stages.push((&store, false, planned)); } if let Some((stage, graph)) = runtime.as_ref().zip(runtime_graph.as_ref()) { let (changes, planned) = operation - .plan_teardown_stage(&bundle, stage, &record, true, graph, cancel) + .plan_teardown_stage( + &bundle, + stage, + &record, + (true, &runtime_bindings), + graph, + cancel, + ) .await?; result.changes.extend(changes); stages.push((stage, true, planned)); @@ -170,25 +184,17 @@ impl Deployment { result.outcome = Outcome::Destroyed; Ok(result) } + /// Plan one stage's teardown from the bindings already read from its state; + /// nothing has changed that state since. async fn plan_teardown_stage( &self, bundle: &Bundle, store: &Store, record: &Record, - runtime: bool, + (runtime, bindings): (bool, &BTreeMap), compiled: &compile::CompiledTeardown, cancel: &CancellationToken, ) -> Result<(Vec, bool), Error> { - let bindings = self - .state_bindings( - bundle, - store, - &record.document, - &record.generations, - runtime, - cancel, - ) - .await?; if bindings.is_empty() { if record.succeeded() || record.destroying() { return Err(Error::Conflict( @@ -198,7 +204,7 @@ impl Deployment { return Ok((Vec::new(), false)); } let expected = with_observations( - &teardown_expected(record, &bindings, runtime)?, + &teardown_expected(record, bindings, runtime)?, &compiled.observations, ); self.initialize(bundle, store, &compiled.graph, cancel) @@ -213,7 +219,7 @@ impl Deployment { ) .await?; Ok(( - check_destroy_plan(&plan, &expected, &bindings, &compiled.retained)?, + check_destroy_plan(&plan, &expected, bindings, &compiled.retained)?, true, )) } @@ -405,7 +411,7 @@ mod tests { &bundle, &store, &record, - false, + (false, &BTreeMap::new()), &compiled, &CancellationToken::new(), )