worker: retry attachment cleanup stages

This commit is contained in:
2026-08-12 17:31:29 +09:00
parent 1e10cdecc8
commit 4e7eaac7d5
2 changed files with 95 additions and 1 deletions
+39 -1
View File
@@ -106,12 +106,15 @@ pub struct WorkerRemovalPlan {
pub reason: String, pub reason: String,
pub created_at: String, pub created_at: String,
pub updated_at: String, pub updated_at: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub failure_category: Option<String>,
} }
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct PreparedWorkerRemoval { pub struct PreparedWorkerRemoval {
pub plan: WorkerRemovalPlan, pub plan: WorkerRemovalPlan,
pub runtime_request: WorkerRetentionExecutionRequest, pub runtime_request: WorkerRetentionExecutionRequest,
pub prior_failure_category: Option<String>,
} }
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
@@ -369,6 +372,7 @@ impl SqliteWorkspaceStore {
) )
})?; })?;
let removed_at = plan.created_at.clone(); let removed_at = plan.created_at.clone();
let prior_failure_category = plan.failure_category.clone();
Ok(PreparedWorkerRemoval { Ok(PreparedWorkerRemoval {
runtime_request: WorkerRetentionExecutionRequest { runtime_request: WorkerRetentionExecutionRequest {
operation_id: plan.operation_id.clone(), operation_id: plan.operation_id.clone(),
@@ -389,6 +393,7 @@ impl SqliteWorkspaceStore {
diagnostics_disposition: plan.diagnostics_disposition, diagnostics_disposition: plan.diagnostics_disposition,
}, },
plan, plan,
prior_failure_category,
}) })
} }
@@ -429,6 +434,7 @@ impl SqliteWorkspaceStore {
let Some(plan) = plan else { let Some(plan) = plan else {
return Ok(None); return Ok(None);
}; };
let prior_failure_category = plan.failure_category.clone();
let worker_number = plan.worker.worker_id.parse::<u64>().map_err(|_| { let worker_number = plan.worker.worker_id.parse::<u64>().map_err(|_| {
WorkerRetentionError::Invalid( WorkerRetentionError::Invalid(
"Runtime Worker id is not a canonical unsigned integer".to_string(), "Runtime Worker id is not a canonical unsigned integer".to_string(),
@@ -468,6 +474,7 @@ impl SqliteWorkspaceStore {
diagnostics_disposition: plan.diagnostics_disposition, diagnostics_disposition: plan.diagnostics_disposition,
}, },
plan, plan,
prior_failure_category,
})) }))
} }
@@ -758,7 +765,7 @@ fn load_plan_q(c: &Connection, key: &str, id: &str) -> crate::Result<Option<Work
worker_revision,run_generation,policy_id,policy_revision,session_disposition, worker_revision,run_generation,policy_id,policy_revision,session_disposition,
metadata_disposition,archive_retention_kind,archive_retention_seconds, metadata_disposition,archive_retention_kind,archive_retention_seconds,
diagnostics_disposition,diagnostics_retention_seconds,archive_id,blockers_json, diagnostics_disposition,diagnostics_retention_seconds,archive_id,blockers_json,
state,reason,created_at,updated_at state,reason,created_at,updated_at,failure_category
FROM worker_removal_operations WHERE {key}=?1" FROM worker_removal_operations WHERE {key}=?1"
); );
c.query_row(&query, params![id], |row| { c.query_row(&query, params![id], |row| {
@@ -799,6 +806,7 @@ fn load_plan_q(c: &Connection, key: &str, id: &str) -> crate::Result<Option<Work
reason: row.get(19)?, reason: row.get(19)?,
created_at: row.get(20)?, created_at: row.get(20)?,
updated_at: row.get(21)?, updated_at: row.get(21)?,
failure_category: row.get(22)?,
}) })
}) })
.optional() .optional()
@@ -1532,6 +1540,36 @@ mod tests {
); );
} }
#[test]
fn recovery_preserves_attachment_failure_stage_for_retry_ordering() {
let s = setup();
let request = req();
let plan = s.plan_worker_removal(&request, &inv()).unwrap();
s.prepare_worker_removal_execution("w", &plan.plan_id, &plan.input_fingerprint)
.unwrap();
s.fail_worker_removal(
"w",
&plan.operation_id,
&plan.input_fingerprint,
"workdir_attachment_release_failed",
)
.unwrap();
let recovered = s
.recover_worker_removal_execution(
"w",
&request.worker,
&request.expected_worker_revision,
&request.reason,
)
.unwrap()
.unwrap();
assert_eq!(
recovered.prior_failure_category.as_deref(),
Some("workdir_attachment_release_failed")
);
assert_eq!(recovered.plan.state, WorkerRemovalPlanState::Failed);
}
#[test] #[test]
fn old_schema_upgrade_seeds_existing_workspace() { fn old_schema_upgrade_seeds_existing_workspace() {
let temp = tempfile::tempdir().unwrap(); let temp = tempfile::tempdir().unwrap();
+56
View File
@@ -396,6 +396,11 @@ impl WorkspaceWorkerRemoveExecutor {
if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Succeeded { if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Succeeded {
return Ok(worker_remove_success_response(&target)); return Ok(worker_remove_success_response(&target));
} }
let must_close_session =
prepared.prior_failure_category.as_deref() == Some("workdir_session_close_failed");
let must_release_attachment = must_close_session
|| prepared.prior_failure_category.as_deref()
== Some("workdir_attachment_release_failed");
let prepared = let prepared =
if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Failed { if prepared.plan.state == crate::retention::WorkerRemovalPlanState::Failed {
match self.store.prepare_worker_removal_execution( match self.store.prepare_worker_removal_execution(
@@ -409,6 +414,57 @@ impl WorkspaceWorkerRemoveExecutor {
} else { } else {
prepared prepared
}; };
if must_close_session {
let session = {
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.get(&target)
.cloned()
};
if let Some(session) = session {
if session.close().await.is_err() {
let _ = self.store.fail_worker_removal(
&self.workspace_id,
&prepared.plan.operation_id,
&prepared.plan.input_fingerprint,
"workdir_session_close_failed",
);
return Ok(worker_remove_error_response(
StatusCode::SERVICE_UNAVAILABLE,
"attachment_close_failed",
"Worker Workdir session could not be closed; removal can be retried",
));
}
self.workdir_sessions
.lock()
.map_err(|_| "Workdir session registry was poisoned".to_string())?
.remove(&target);
}
}
if must_release_attachment
&& self
.store
.detach_worker_workdir(
&self.workspace_id,
&target,
None,
&Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true),
)
.is_err()
{
let _ = self.store.fail_worker_removal(
&self.workspace_id,
&prepared.plan.operation_id,
&prepared.plan.input_fingerprint,
"workdir_attachment_release_failed",
);
return Ok(worker_remove_error_response(
StatusCode::SERVICE_UNAVAILABLE,
"attachment_release_failed",
"Worker Workdir attachment could not be released; removal can be retried",
));
}
return self return self
.resume_worker_retention(&runtime, &target, prepared) .resume_worker_retention(&runtime, &target, prepared)
.await; .await;