Add managed runtime adapter and require operation-owned write drain through update
This commit is contained in:
@@ -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<Output = Result<()>> + Send;
|
||||
fn verify_barrier(&self, operation: &str) -> impl Future<Output = Result<()>> + Send;
|
||||
fn release_barrier(&self, operation: &str) -> impl Future<Output = Result<()>> + 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<Output = Result<RecoveryImage>> + Send;
|
||||
fn capture(&self, name: &str) -> impl Future<Output = Result<Unit>> + Send;
|
||||
fn validate_original_file(&self, original: &Unit) -> impl Future<Output = Result<()>> + Send;
|
||||
fn read(&self, name: &str) -> impl Future<Output = Result<String>> + Send;
|
||||
fn write(
|
||||
&self,
|
||||
name: &str,
|
||||
original: &Unit,
|
||||
expected: &[String],
|
||||
body: &str,
|
||||
) -> impl Future<Output = Result<()>> + 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<Vec<String>>,
|
||||
fail_new_hooks: AtomicBool,
|
||||
fail_snapshot: AtomicBool,
|
||||
fail_barrier_release: AtomicBool,
|
||||
barrier: Mutex<Option<String>>,
|
||||
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<Unit> {
|
||||
Ok(self.original.clone())
|
||||
}
|
||||
async fn validate_original_file(&self, _original: &Unit) -> Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
async fn read(&self, _name: &str) -> Result<String> {
|
||||
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();
|
||||
|
||||
Reference in New Issue
Block a user