From 6e7ea8b9d6679b30e9e5dfa8f1252ec7cc03649b Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 02:05:52 -0400 Subject: [PATCH] Add managed runtime adapter and require operation-owned write drain through update --- core/archipelago/src/container/mod.rs | 1 + .../src/container/supervised_runtime.rs | 427 ++++++++++++++++++ .../src/container/supervised_update.rs | 117 ++++- 3 files changed, 539 insertions(+), 6 deletions(-) create mode 100644 core/archipelago/src/container/supervised_runtime.rs diff --git a/core/archipelago/src/container/mod.rs b/core/archipelago/src/container/mod.rs index 98db75e1..465dc48d 100644 --- a/core/archipelago/src/container/mod.rs +++ b/core/archipelago/src/container/mod.rs @@ -34,4 +34,5 @@ mod staged_update; pub(crate) mod update_transaction; +pub(crate) mod supervised_runtime; pub(crate) mod supervised_update; diff --git a/core/archipelago/src/container/supervised_runtime.rs b/core/archipelago/src/container/supervised_runtime.rs new file mode 100644 index 00000000..21f51705 --- /dev/null +++ b/core/archipelago/src/container/supervised_runtime.rs @@ -0,0 +1,427 @@ +//! Production systemd/Podman adapter. Application write admission/drain is an +//! explicit dependency: neither process pause nor a filesystem receipt is drain. +use super::{ + supervised_update::{self, PreparedTarget, RecoveryImage, Supervisor, Unit}, + update_transaction::{Observed, Podman, Runtime, Target}, +}; +use anyhow::{Context, Result}; +use sha2::{Digest, Sha256}; +use std::{ + collections::HashMap, + future::Future, + io::Write, + os::unix::fs::{MetadataExt, OpenOptionsExt}, + path::{Path, PathBuf}, + time::Duration, +}; + +pub(crate) trait DrainBarrier: Sync { + /// Persist ownership before blocking admissions. Return only after API, + /// direct uploads and workers have drained and the coherent backup finished. + fn acquire( + &self, + operation: &str, + originals: &[Unit], + ) -> impl Future> + Send; + /// Must inspect the live admission fence and same-operation ownership, not + /// merely trust an old 'backup complete' file. Keep it through snapshot/stop. + fn verify(&self, operation: &str) -> impl Future> + Send; + /// Idempotent; an old recovery must never clear a newer operation's fence. + fn release(&self, operation: &str) -> impl Future> + Send; +} + +#[derive(Clone, serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +pub(crate) struct MigrationPlan { + original_sha256: String, + prepared: PreparedTarget, +} +impl MigrationPlan { + /// The old render must be reproducible from an image-bound retained original + /// manifest. A current target catalog is not an installed-version receipt. + pub(crate) fn from_installed_render( + original: &str, + old_render: &str, + target_render: &str, + manifest: archipelago_container::AppManifest, + ) -> Result { + manifest.validate()?; + let body = supervised_update::merge_reviewed_unit(original, old_render, target_render)?; + Ok(Self { + original_sha256: hex::encode(Sha256::digest(original.as_bytes())), + prepared: PreparedTarget { body, manifest }, + }) + } + /// Legacy overrides whose provenance cannot be reconstructed need a reviewed + /// exact-original-hash-bound migration, prepared by node administration code. + /// No browser RPC accepts unit bodies or this type. + pub(crate) fn from_reviewed_legacy( + original_sha256: &str, + body: String, + manifest: archipelago_container::AppManifest, + ) -> Result { + anyhow::ensure!( + original_sha256.len() == 64 && original_sha256.bytes().all(|v| v.is_ascii_hexdigit()), + "Invalid original unit commitment" + ); + manifest.validate()?; + Ok(Self { + original_sha256: original_sha256.into(), + prepared: PreparedTarget { body, manifest }, + }) + } +} + +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct PlanFile { + schema: u8, + package: String, + plans: HashMap, +} +/// Administrative preparation is separate from the Update RPC. The RPC never +/// accepts unit bodies, hook commands, or a browser-chosen preparation path. +pub(crate) fn load_reviewed_plans( + data_dir: &Path, + package: &str, +) -> Result> { + anyhow::ensure!( + !package.is_empty() + && package.len() <= 128 + && package + .bytes() + .all(|v| v.is_ascii_alphanumeric() || matches!(v, b'-' | b'_')), + "Invalid managed package" + ); + let dir = data_dir.join("managed-update-plans"); + let directory = std::fs::symlink_metadata(&dir) + .context("No original-bound managed update plan; existing app remains unchanged")?; + anyhow::ensure!( + directory.is_dir() + && !directory.file_type().is_symlink() + && directory.uid() == unsafe { libc::geteuid() } + && directory.mode() & 0o077 == 0, + "Managed update plan directory must be private and node-owned" + ); + let path = dir.join(format!("{package}.json")); + let metadata = std::fs::symlink_metadata(&path)?; + anyhow::ensure!( + metadata.is_file() + && !metadata.file_type().is_symlink() + && metadata.uid() == unsafe { libc::geteuid() } + && metadata.mode() & 0o077 == 0 + && metadata.len() <= 4 * 1024 * 1024, + "Managed update plan must be private and node-owned" + ); + let record: PlanFile = serde_json::from_slice(&std::fs::read(path)?)?; + anyhow::ensure!( + record.schema == 1 + && record.package == package + && !record.plans.is_empty() + && record.plans.len() <= 32, + "Managed update plan does not match this package" + ); + for plan in record.plans.values() { + anyhow::ensure!( + plan.original_sha256.len() == 64 + && plan.original_sha256.bytes().all(|v| v.is_ascii_hexdigit()), + "Invalid original unit commitment" + ); + plan.prepared.manifest.validate()?; + } + Ok(record.plans) +} + +pub(crate) struct SystemdSupervisor { + data_dir: PathBuf, + unit_dir: PathBuf, + plans: HashMap, + barrier: B, +} +impl SystemdSupervisor { + pub(crate) async fn new( + data_dir: PathBuf, + plans: HashMap, + barrier: B, + ) -> Result { + let unit_dir = super::quadlet::unit_dir().await?; + let meta = std::fs::symlink_metadata(&unit_dir)?; + anyhow::ensure!( + meta.is_dir() + && !meta.file_type().is_symlink() + && meta.uid() == unsafe { libc::geteuid() } + && meta.mode() & 0o022 == 0, + "Quadlet directory ownership changed" + ); + Ok(Self { + data_dir, + unit_dir, + plans, + barrier, + }) + } + fn path(&self, name: &str) -> Result { + anyhow::ensure!( + !name.is_empty() + && name.len() <= 128 + && name + .bytes() + .all(|v| v.is_ascii_alphanumeric() || matches!(v, b'-' | b'_')), + "Invalid supervised unit name" + ); + Ok(self.unit_dir.join(format!("{name}.container"))) + } + async fn manager(args: &[&str]) -> Result { + let output = tokio::time::timeout( + Duration::from_secs(120), + tokio::process::Command::new("systemctl") + .arg("--user") + .args(args) + .kill_on_drop(true) + .output(), + ) + .await + .context("User service manager timed out")??; + anyhow::ensure!( + output.status.success(), + "User service manager rejected operation" + ); + Ok(String::from_utf8(output.stdout)?) + } + async fn owned_file(&self, name: &str) -> Result<(PathBuf, std::fs::Metadata)> { + let expected = self.path(name)?; + let service = format!("{name}.service"); + let output = Self::manager(&[ + "show", + &service, + "--property=SourcePath", + "--property=DropInPaths", + "--property=LoadState", + ]) + .await?; + let values: HashMap<_, _> = output + .lines() + .filter_map(|line| line.split_once('=')) + .collect(); + anyhow::ensure!( + values.get("LoadState") == Some(&"loaded") + && values.get("DropInPaths") == Some(&"") + && values + .get("SourcePath") + .is_some_and(|path| Path::new(path) == expected), + "Service is not owned by the exact original source Quadlet or has external overrides" + ); + let meta = std::fs::symlink_metadata(&expected)?; + anyhow::ensure!( + meta.is_file() + && !meta.file_type().is_symlink() + && meta.uid() == unsafe { libc::geteuid() } + && meta.mode() & 0o022 == 0 + && meta.len() <= 1024 * 1024, + "Original Quadlet ownership changed" + ); + Ok((expected, meta)) + } +} +impl Supervisor for SystemdSupervisor { + async fn begin_barrier(&self, operation: &str, originals: &[Unit]) -> Result<()> { + self.barrier.acquire(operation, originals).await + } + async fn verify_barrier(&self, operation: &str) -> Result<()> { + self.barrier.verify(operation).await + } + async fn release_barrier(&self, operation: &str) -> Result<()> { + self.barrier.release(operation).await + } + async fn prepare_target(&self, target: &Target, original: &Unit) -> Result { + let plan = self + .plans + .get(&target.name) + .context("No reviewed original-bound managed migration; app unchanged")?; + anyhow::ensure!( + hex::encode(Sha256::digest(original.body.as_bytes())) == plan.original_sha256 + && plan.prepared.manifest.app.container.image.as_deref() + == Some(target.reference.as_str()) + && super::prod_orchestrator::compute_container_name(&plan.prepared.manifest) + == target.name, + "Original unit or reviewed immutable target changed before update" + ); + Ok(plan.prepared.clone()) + } + async fn target_hooks( + &self, + name: &str, + manifest: &archipelago_container::AppManifest, + ) -> Result<()> { + super::hooks::run_post_install_strict(manifest, name, &self.data_dir).await + } + async fn capture(&self, name: &str) -> Result { + let (path, meta) = self.owned_file(name).await?; + let body = std::fs::read_to_string(path)?; + let observed = Podman + .inspect(name) + .await? + .context("Original supervised runtime missing")?; + Ok(Unit { + name: name.into(), + body, + image: observed.image, + container_id: observed.id, + file_mode: meta.mode() & 0o777, + running: observed.running, + config_sha256: observed.config_sha256, + }) + } + async fn validate_original_file(&self, original: &Unit) -> Result<()> { + let (_, meta) = self.owned_file(&original.name).await?; + anyhow::ensure!( + meta.mode() & 0o777 == original.file_mode, + "Original Quadlet mode changed" + ); + Ok(()) + } + async fn read(&self, name: &str) -> Result { + let (path, _) = self.owned_file(name).await?; + Ok(std::fs::read_to_string(path)?) + } + async fn write(&self, original: &Unit, expected: &[String], body: &str) -> Result<()> { + let (path, meta) = self.owned_file(&original.name).await?; + anyhow::ensure!( + meta.mode() & 0o777 == original.file_mode + && expected.contains(&std::fs::read_to_string(&path)?), + "Original Quadlet changed; refusing to replace a foreign edit" + ); + let temporary = self + .unit_dir + .join(format!(".archy-update-{}.tmp", uuid::Uuid::new_v4())); + // Commit point is synchronous while the transaction guard is held; + // cancelled futures cannot rename over a later recovery after unlock. + let result = (|| -> Result<()> { + let mut file = std::fs::OpenOptions::new() + .create_new(true) + .write(true) + .mode(original.file_mode) + .open(&temporary)?; + file.write_all(body.as_bytes())?; + file.sync_all()?; + std::fs::rename(&temporary, &path)?; + std::fs::File::open(&self.unit_dir)?.sync_all()?; + Ok(()) + })(); + if result.is_err() { + let _ = std::fs::remove_file(&temporary); + } + result + } + async fn snapshot(&self, original: &Unit, operation: &str, tag: &str) -> Result { + self.barrier.verify(operation).await?; + supervised_update::capture_local_recovery_image(original, operation, tag).await + } + async fn pin(&self, image: &str, tag: &str) -> Result<()> { + anyhow::ensure!( + image.len() == 64 + && image.bytes().all(|v| v.is_ascii_hexdigit()) + && tag.starts_with("localhost/archy-update-recovery:"), + "Invalid recovery image pin" + ); + let output = tokio::time::timeout( + Duration::from_secs(30), + tokio::process::Command::new("podman") + .args(["tag", &format!("sha256:{image}"), tag]) + .kill_on_drop(true) + .output(), + ) + .await??; + anyhow::ensure!( + output.status.success(), + "Original recovery image unavailable" + ); + Ok(()) + } + async fn stop(&self, name: &str) -> Result<()> { + self.owned_file(name).await?; + super::quadlet::stop_service_with_timeout( + &format!("{name}.service"), + Duration::from_secs(archipelago_container::runtime::stop_grace_secs_for(name) + 30), + ) + .await + } + async fn reload(&self) -> Result<()> { + super::quadlet::daemon_reload_user().await + } + async fn start(&self, name: &str) -> Result<()> { + self.owned_file(name).await?; + super::quadlet::enable_now(&format!("{name}.service")).await + } + async fn observed(&self, name: &str) -> Result> { + Podman.inspect(name).await + } + async fn healthy(&self, name: &str) -> Result { + Podman.healthy(name).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + struct UnusedBarrier; + impl DrainBarrier for UnusedBarrier { + async fn acquire(&self, _: &str, _: &[Unit]) -> Result<()> { + anyhow::bail!("No application barrier installed") + } + async fn verify(&self, _: &str) -> Result<()> { + anyhow::bail!("No application barrier installed") + } + async fn release(&self, _: &str) -> Result<()> { + Ok(()) + } + } + #[tokio::test] + async fn original_bound_plan_rejects_changed_unit_or_image_without_service_calls() { + let image = format!("localhost/movie@sha256:{}", "b".repeat(64)); + let manifest = archipelago_container::AppManifest::parse(&format!( + "app:\n id: movie\n name: Movie\n version: 2.0.0\n container:\n image: {image}\n")).unwrap(); + let original_body = "[Container]\nContainerName=movie\nImage=old\n"; + let next_body = format!("[Container]\nContainerName=movie\nImage={image}\n"); + let plan = MigrationPlan::from_reviewed_legacy( + &hex::encode(Sha256::digest(original_body)), + next_body.clone(), + manifest, + ) + .unwrap(); + let adapter = SystemdSupervisor { + data_dir: PathBuf::from("/unused"), + unit_dir: PathBuf::from("/unused"), + plans: HashMap::from([("movie".into(), plan)]), + barrier: UnusedBarrier, + }; + let mut original = Unit { + name: "movie".into(), + body: original_body.into(), + image: "a".repeat(64), + container_id: "c".repeat(64), + file_mode: 0o600, + running: true, + config_sha256: "d".repeat(64), + }; + let mut target = Target { + name: "movie".into(), + reference: image, + image: "b".repeat(64), + }; + assert_eq!( + adapter + .prepare_target(&target, &original) + .await + .unwrap() + .body, + next_body + ); + original.body.push_str("Environment=OPERATOR=changed\n"); + assert!(adapter.prepare_target(&target, &original).await.is_err()); + original.body = original_body.into(); + target.reference = format!("localhost/movie@sha256:{}", "e".repeat(64)); + assert!(adapter.prepare_target(&target, &original).await.is_err()); + assert!(adapter.begin_barrier("unused", &[original]).await.is_err()); + } +} diff --git a/core/archipelago/src/container/supervised_update.rs b/core/archipelago/src/container/supervised_update.rs index 0f4d736b..3d96d3fe 100644 --- a/core/archipelago/src/container/supervised_update.rs +++ b/core/archipelago/src/container/supervised_update.rs @@ -15,6 +15,7 @@ pub(crate) struct Unit { pub body: String, pub image: String, pub container_id: String, + pub file_mode: u32, pub running: bool, pub config_sha256: String, } @@ -59,6 +60,16 @@ struct Journal { } pub(crate) trait Supervisor: Sync { + /// Admission must be fenced and in-flight application work drained under + /// this durable operation before snapshots. A paused process is not proof. + /// Acquisition/release are idempotent and operation-owned across restart. + fn begin_barrier( + &self, + operation: &str, + originals: &[Unit], + ) -> impl Future> + Send; + fn verify_barrier(&self, operation: &str) -> impl Future> + Send; + fn release_barrier(&self, operation: &str) -> impl Future> + Send; /// Only an internal reviewed signed-manifest planner may supply this value; /// browser parameters must never become a unit body or hook recipe. fn prepare_target( @@ -83,10 +94,11 @@ pub(crate) trait Supervisor: Sync { tag: &str, ) -> impl Future> + Send; fn capture(&self, name: &str) -> impl Future> + Send; + fn validate_original_file(&self, original: &Unit) -> impl Future> + Send; fn read(&self, name: &str) -> impl Future> + Send; fn write( &self, - name: &str, + original: &Unit, expected: &[String], body: &str, ) -> impl Future> + Send; @@ -406,6 +418,8 @@ fn validate(record: &Journal) -> Result<()> { && digest(&member.target.image) && digest(&member.original.image) && digest(&member.original.container_id) + && member.original.file_mode & !0o777 == 0 + && member.original.file_mode & 0o022 == 0 && digest(&member.original.config_sha256), "Invalid original supervised identity" ); @@ -567,6 +581,10 @@ pub(crate) async fn execute( } let result = apply(guard, &mut record, supervisor).await; if let Err(error) = result { + if record.phase == Phase::Committed { + return Err(error + .context("Update committed; completion cleanup must resume without rolling back")); + } return match restore(guard,&mut record,supervisor).await { Ok(())=>Err(error.context("Original supervised image/configuration and running intent restored; container IDs may change and data was not rolled back")), Err(recovery)=>Err(error.context(format!("Supervised runtime recovery remains unresolved: {recovery:#}"))), @@ -575,12 +593,20 @@ pub(crate) async fn execute( Ok(()) } async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor) -> Result<()> { + let originals: Vec<_> = record + .members + .iter() + .map(|member| member.original.clone()) + .collect(); + supervisor.begin_barrier(&record.id, &originals).await?; for index in 0..record.members.len() { let member = &record.members[index]; + supervisor.validate_original_file(&member.original).await?; anyhow::ensure!( supervisor.read(&member.original.name).await? == member.original.body, "Unit edited before update; originals retained" ); + supervisor.verify_barrier(&record.id).await?; let image = supervisor .snapshot(&member.original, &record.id, &member.original_tag) .await?; @@ -600,15 +626,17 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor // The image acknowledgement becomes durable before any original stop. save(guard, record)?; } + supervisor.verify_barrier(&record.id).await?; record.phase = Phase::Editing; save(guard, record)?; for member in record.members.iter().rev() { + supervisor.verify_barrier(&record.id).await?; supervisor.stop(&member.original.name).await?; } for member in &record.members { supervisor .write( - &member.original.name, + &member.original, &[member.original.body.clone()], &member.target_body, ) @@ -635,6 +663,7 @@ async fn apply(guard: &Guard, record: &mut Journal, supervisor: &impl Supervisor } record.phase = Phase::Committed; save(guard, record)?; + supervisor.release_barrier(&record.id).await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } @@ -660,13 +689,16 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis } record.phase = Phase::Aborted; save(guard, record)?; + supervisor.release_barrier(&record.id).await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } return Ok(()); } + supervisor.verify_barrier(&record.id).await?; // Refuse to overwrite a foreign edit before stopping any surviving member. for member in &record.members { + supervisor.validate_original_file(&member.original).await?; let body = supervisor.read(&member.original.name).await?; anyhow::ensure!( [ @@ -694,7 +726,7 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis for member in &record.members { supervisor .write( - &member.original.name, + &member.original, &[ member.original.body.clone(), member.target_body.clone(), @@ -734,17 +766,21 @@ async fn restore(guard: &Guard, record: &mut Journal, supervisor: &impl Supervis guard.hold(&member.original.name, &record.id)?; } record.phase = Phase::Restored; - save(guard, record) + save(guard, record)?; + supervisor.release_barrier(&record.id).await } pub(crate) async fn recover(guard: &Guard, supervisor: &impl Supervisor) -> Result<()> { for mut record in records(guard)? { match record.phase { Phase::Committed | Phase::Aborted => { + supervisor.release_barrier(&record.id).await?; for member in &record.members { guard.release_hold(&member.original.name, &record.id)?; } } - Phase::Restored => {} // A later retry may own the current hold. + Phase::Restored => { + supervisor.release_barrier(&record.id).await?; + } // Never release a newer owner. _ => restore(guard, &mut record, supervisor).await?, } } @@ -765,6 +801,8 @@ mod tests { calls: Mutex>, fail_new_hooks: AtomicBool, fail_snapshot: AtomicBool, + fail_barrier_release: AtomicBool, + barrier: Mutex>, generation: AtomicUsize, } impl Mock { @@ -776,6 +814,7 @@ mod tests { body: body.clone(), image: "a".repeat(64), container_id: format!("{:064x}", 1), + file_mode: 0o600, running: true, config_sha256: "c".repeat(64), }, @@ -784,6 +823,8 @@ mod tests { calls: Default::default(), fail_new_hooks: AtomicBool::new(false), fail_snapshot: AtomicBool::new(false), + fail_barrier_release: AtomicBool::new(false), + barrier: Mutex::new(None), generation: AtomicUsize::new(1), } } @@ -796,6 +837,33 @@ mod tests { } } impl Supervisor for Mock { + async fn begin_barrier(&self, operation: &str, _originals: &[Unit]) -> Result<()> { + let mut held = self.barrier.lock().unwrap(); + anyhow::ensure!( + held.as_deref().is_none_or(|id| id == operation), + "Another drain barrier owns admission" + ); + *held = Some(operation.into()); + Ok(()) + } + async fn verify_barrier(&self, operation: &str) -> Result<()> { + anyhow::ensure!( + self.barrier.lock().unwrap().as_deref() == Some(operation), + "Write barrier not held" + ); + Ok(()) + } + async fn release_barrier(&self, operation: &str) -> Result<()> { + anyhow::ensure!( + !self.fail_barrier_release.load(Ordering::SeqCst), + "Barrier release unavailable" + ); + let mut held = self.barrier.lock().unwrap(); + if held.as_deref() == Some(operation) { + *held = None; + } + Ok(()) + } async fn prepare_target( &self, target: &Target, @@ -843,10 +911,13 @@ mod tests { async fn capture(&self, _name: &str) -> Result { Ok(self.original.clone()) } + async fn validate_original_file(&self, _original: &Unit) -> Result<()> { + Ok(()) + } async fn read(&self, _name: &str) -> Result { Ok(self.body.lock().unwrap().clone()) } - async fn write(&self, _name: &str, expected: &[String], body: &str) -> Result<()> { + async fn write(&self, _original: &Unit, expected: &[String], body: &str) -> Result<()> { let mut current = self.body.lock().unwrap(); anyhow::ensure!(expected.contains(&*current), "Foreign unit edit"); *current = body.into(); @@ -933,6 +1004,38 @@ mod tests { assert!(merge_reviewed_unit("[Container]\nEnvironment=A=1 B=2\n", old, old).is_err()); } #[tokio::test] + async fn lost_barrier_release_after_commit_never_rolls_back_a_verified_update() { + let root = tempfile::tempdir().unwrap(); + let guard = Guard::acquire(root.path()).unwrap(); + let runtime = Mock::new(); + runtime.fail_barrier_release.store(true, Ordering::SeqCst); + let error = execute(&guard, "movie", &[Mock::target()], &runtime) + .await + .unwrap_err(); + assert!(error.to_string().contains("Update committed")); + assert_eq!(records(&guard).unwrap()[0].phase, Phase::Committed); + assert_eq!( + runtime.observed("movie").await.unwrap().unwrap().image, + "b".repeat(64) + ); + assert_eq!( + runtime + .calls + .lock() + .unwrap() + .iter() + .filter(|call| *call == "stop-unit") + .count(), + 1 + ); + runtime.calls.lock().unwrap().clear(); + runtime.fail_barrier_release.store(false, Ordering::SeqCst); + recover(&guard, &runtime).await.unwrap(); + assert!(runtime.calls.lock().unwrap().is_empty()); + assert!(runtime.barrier.lock().unwrap().is_none()); + assert!(!super::super::update_transaction::is_held(root.path(), "movie").unwrap()); + } + #[tokio::test] async fn snapshot_failure_never_stops_or_recreates_original_runtime() { let root = tempfile::tempdir().unwrap(); let guard = Guard::acquire(root.path()).unwrap(); @@ -1022,6 +1125,7 @@ mod tests { .unwrap(); let mut record = records(&guard).unwrap().pop().unwrap(); record.phase = Phase::Starting; + *runtime.barrier.lock().unwrap() = Some(record.id.clone()); save(&guard, &record).unwrap(); runtime.calls.lock().unwrap().clear(); recover(&guard, &runtime).await.unwrap(); @@ -1046,6 +1150,7 @@ mod tests { .unwrap(); let mut record = records(&guard).unwrap().pop().unwrap(); record.phase = Phase::Starting; + *runtime.barrier.lock().unwrap() = Some(record.id.clone()); save(&guard, &record).unwrap(); *runtime.body.lock().unwrap() = "operator replaced unit".into(); runtime.calls.lock().unwrap().clear();