fix: retain Worker identity after delete persistence failure
This commit is contained in:
@@ -21,6 +21,12 @@ pub enum RuntimeError {
|
|||||||
message: String,
|
message: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
#[error("worker {worker_id} metadata deletion failed: {message}")]
|
||||||
|
WorkerDeletePersistenceFailed {
|
||||||
|
worker_id: WorkerId,
|
||||||
|
message: String,
|
||||||
|
},
|
||||||
|
|
||||||
#[error("worker creation has no execution backend: {message}")]
|
#[error("worker creation has no execution backend: {message}")]
|
||||||
ExecutionBackendUnavailable { message: String },
|
ExecutionBackendUnavailable { message: String },
|
||||||
|
|
||||||
|
|||||||
@@ -2482,6 +2482,7 @@ fn status_for_runtime_error(error: &RuntimeError) -> StatusCode {
|
|||||||
RuntimeError::StoreIo { .. }
|
RuntimeError::StoreIo { .. }
|
||||||
| RuntimeError::StoreMissing { .. }
|
| RuntimeError::StoreMissing { .. }
|
||||||
| RuntimeError::StoreCorrupt { .. }
|
| RuntimeError::StoreCorrupt { .. }
|
||||||
|
| RuntimeError::WorkerDeletePersistenceFailed { .. }
|
||||||
| RuntimeError::StatePoisoned => StatusCode::INTERNAL_SERVER_ERROR,
|
| RuntimeError::StatePoisoned => StatusCode::INTERNAL_SERVER_ERROR,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -2494,6 +2495,9 @@ fn code_for_runtime_error(error: &RuntimeError) -> String {
|
|||||||
RuntimeError::WorkerExecutionUnavailable { .. } => {
|
RuntimeError::WorkerExecutionUnavailable { .. } => {
|
||||||
"worker_execution_unavailable".to_string()
|
"worker_execution_unavailable".to_string()
|
||||||
}
|
}
|
||||||
|
RuntimeError::WorkerDeletePersistenceFailed { .. } => {
|
||||||
|
"worker_delete_persistence_failed".to_string()
|
||||||
|
}
|
||||||
RuntimeError::ExecutionBackendUnavailable { .. } => {
|
RuntimeError::ExecutionBackendUnavailable { .. } => {
|
||||||
"execution_backend_unavailable".to_string()
|
"execution_backend_unavailable".to_string()
|
||||||
}
|
}
|
||||||
@@ -3505,6 +3509,25 @@ mod tests {
|
|||||||
assert!(matches!(error, RuntimeHttpServerError::AuthRequired));
|
assert!(matches!(error, RuntimeHttpServerError::AuthRequired));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn worker_delete_persistence_error_is_bounded_and_typed() {
|
||||||
|
let worker_id = crate::identity::WorkerId::now_v7();
|
||||||
|
let error = RuntimeError::WorkerDeletePersistenceFailed {
|
||||||
|
worker_id,
|
||||||
|
message: "Worker metadata deletion failed; the persisted Worker identity was retained for retry".to_string(),
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
status_for_runtime_error(&error),
|
||||||
|
StatusCode::INTERNAL_SERVER_ERROR
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
code_for_runtime_error(&error),
|
||||||
|
"worker_delete_persistence_failed"
|
||||||
|
);
|
||||||
|
assert!(!error.to_string().contains('/'));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn workdir_runtime_errors_preserve_diagnostic_code() {
|
fn workdir_runtime_errors_preserve_diagnostic_code() {
|
||||||
let cases = [
|
let cases = [
|
||||||
|
|||||||
@@ -72,6 +72,10 @@ impl RuntimeWorkspaceScope {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const SUBSCRIPTION_QUEUE_CAPACITY: usize = 256;
|
const SUBSCRIPTION_QUEUE_CAPACITY: usize = 256;
|
||||||
|
const WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE: &str = "worker_delete_persistence_failed";
|
||||||
|
const WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE: &str =
|
||||||
|
"Worker metadata deletion failed; the persisted Worker identity was retained for retry";
|
||||||
|
const WORKER_DELETE_FAILURE_MESSAGE_MAX_BYTES: usize = 256;
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||||
enum WorkerRestoreMode {
|
enum WorkerRestoreMode {
|
||||||
@@ -2206,7 +2210,14 @@ impl Runtime {
|
|||||||
worker_ref.worker_id
|
worker_ref.worker_id
|
||||||
)));
|
)));
|
||||||
}
|
}
|
||||||
state.delete_worker_record(&worker_ref.worker_id)?;
|
if state.delete_worker_record(&worker_ref.worker_id).is_err() {
|
||||||
|
state.record_worker_delete_persistence_failure(worker_ref);
|
||||||
|
let _diagnostic_persistence = state.persist_runtime_snapshot();
|
||||||
|
return Err(RuntimeError::WorkerDeletePersistenceFailed {
|
||||||
|
worker_id: worker_ref.worker_id,
|
||||||
|
message: WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE.to_string(),
|
||||||
|
});
|
||||||
|
}
|
||||||
let removed = state.workers.remove(&worker_ref.worker_id).ok_or_else(|| {
|
let removed = state.workers.remove(&worker_ref.worker_id).ok_or_else(|| {
|
||||||
RuntimeError::WorkerNotFound {
|
RuntimeError::WorkerNotFound {
|
||||||
worker_id: worker_ref.worker_id,
|
worker_id: worker_ref.worker_id,
|
||||||
@@ -3281,6 +3292,40 @@ impl RuntimeState {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn record_worker_delete_persistence_failure(&mut self, worker_ref: &WorkerRef) {
|
||||||
|
debug_assert!(
|
||||||
|
WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE.len()
|
||||||
|
<= WORKER_DELETE_FAILURE_MESSAGE_MAX_BYTES
|
||||||
|
);
|
||||||
|
if self.diagnostics.iter().any(|diagnostic| {
|
||||||
|
diagnostic.code == WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE
|
||||||
|
&& diagnostic.worker_ref.as_ref() == Some(worker_ref)
|
||||||
|
}) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
#[cfg(feature = "fs-store")]
|
||||||
|
let diagnostic_id = {
|
||||||
|
let diagnostic_id = self.next_diagnostic_id;
|
||||||
|
self.next_diagnostic_id = self.next_diagnostic_id.saturating_add(1);
|
||||||
|
diagnostic_id
|
||||||
|
};
|
||||||
|
#[cfg(not(feature = "fs-store"))]
|
||||||
|
let diagnostic_id = self
|
||||||
|
.diagnostics
|
||||||
|
.iter()
|
||||||
|
.map(|diagnostic| diagnostic.id)
|
||||||
|
.max()
|
||||||
|
.unwrap_or(0)
|
||||||
|
.saturating_add(1);
|
||||||
|
self.diagnostics.push(RuntimeDiagnostic {
|
||||||
|
id: diagnostic_id,
|
||||||
|
code: WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE.to_string(),
|
||||||
|
severity: DiagnosticSeverity::Error,
|
||||||
|
message: WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE.to_string(),
|
||||||
|
worker_ref: Some(worker_ref.clone()),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "fs-store")]
|
#[cfg(feature = "fs-store")]
|
||||||
fn record_restore_failure(
|
fn record_restore_failure(
|
||||||
&mut self,
|
&mut self,
|
||||||
@@ -3893,6 +3938,9 @@ fn runtime_worker_create_failure_fields(
|
|||||||
RuntimeError::WorkerExecutionUnavailable { .. } => {
|
RuntimeError::WorkerExecutionUnavailable { .. } => {
|
||||||
("worker_execution_unavailable", None, None)
|
("worker_execution_unavailable", None, None)
|
||||||
}
|
}
|
||||||
|
RuntimeError::WorkerDeletePersistenceFailed { .. } => {
|
||||||
|
(WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE, None, None)
|
||||||
|
}
|
||||||
RuntimeError::ExecutionBackendUnavailable { .. } => {
|
RuntimeError::ExecutionBackendUnavailable { .. } => {
|
||||||
("execution_backend_unavailable", None, None)
|
("execution_backend_unavailable", None, None)
|
||||||
}
|
}
|
||||||
@@ -6679,6 +6727,105 @@ mod tests {
|
|||||||
assert_eq!(summary.worker_count, 0);
|
assert_eq!(summary.worker_count, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn delete_worker_persistence_failure_retains_identity_and_bounded_diagnostic() {
|
||||||
|
let root = tempfile::tempdir().unwrap();
|
||||||
|
let store_root = root.path().join("runtime");
|
||||||
|
let runtime = Runtime::with_fs_store_and_execution_backend(
|
||||||
|
crate::fs_store::FsRuntimeStoreOptions {
|
||||||
|
root: store_root.clone(),
|
||||||
|
runtime_id: "delete-persistence-failure".to_string(),
|
||||||
|
display_name: None,
|
||||||
|
},
|
||||||
|
Arc::new(TestExecutionBackend::default()),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
runtime.store_config_bundle(test_bundle()).unwrap();
|
||||||
|
let worker = runtime
|
||||||
|
.create_worker(task_request("retain identity after delete failure"))
|
||||||
|
.unwrap();
|
||||||
|
runtime
|
||||||
|
.stop_worker(&worker.worker_ref, Some("done".to_string()))
|
||||||
|
.unwrap();
|
||||||
|
let worker_dir = store_root
|
||||||
|
.join("workers")
|
||||||
|
.join(worker.worker_id.to_string());
|
||||||
|
let saved_worker_dir = store_root.join("saved-worker-record");
|
||||||
|
std::fs::rename(&worker_dir, &saved_worker_dir).unwrap();
|
||||||
|
std::fs::write(&worker_dir, b"not a directory").unwrap();
|
||||||
|
|
||||||
|
let error = runtime.delete_worker(&worker.worker_ref).unwrap_err();
|
||||||
|
|
||||||
|
assert!(matches!(
|
||||||
|
&error,
|
||||||
|
RuntimeError::WorkerDeletePersistenceFailed { worker_id, message }
|
||||||
|
if *worker_id == worker.worker_id
|
||||||
|
&& message == WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE
|
||||||
|
&& message.len() <= WORKER_DELETE_FAILURE_MESSAGE_MAX_BYTES
|
||||||
|
));
|
||||||
|
let error_text = error.to_string();
|
||||||
|
assert!(error_text.len() <= WORKER_DELETE_FAILURE_MESSAGE_MAX_BYTES + 96);
|
||||||
|
assert!(!error_text.contains(root.path().to_string_lossy().as_ref()));
|
||||||
|
assert_eq!(
|
||||||
|
runtime.worker_detail(&worker.worker_ref).unwrap().worker_id,
|
||||||
|
worker.worker_id
|
||||||
|
);
|
||||||
|
let diagnostic = runtime
|
||||||
|
.diagnostics()
|
||||||
|
.unwrap()
|
||||||
|
.into_iter()
|
||||||
|
.find(|diagnostic| {
|
||||||
|
diagnostic.code == WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE
|
||||||
|
&& diagnostic.worker_ref.as_ref() == Some(&worker.worker_ref)
|
||||||
|
})
|
||||||
|
.expect("retained delete failure diagnostic");
|
||||||
|
assert_eq!(diagnostic.message, WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE);
|
||||||
|
assert!(diagnostic.message.len() <= WORKER_DELETE_FAILURE_MESSAGE_MAX_BYTES);
|
||||||
|
assert!(
|
||||||
|
!diagnostic
|
||||||
|
.message
|
||||||
|
.contains(root.path().to_string_lossy().as_ref())
|
||||||
|
);
|
||||||
|
let persisted_runtime: serde_json::Value = serde_json::from_slice(
|
||||||
|
&std::fs::read(store_root.join("runtime.json")).expect("runtime snapshot"),
|
||||||
|
)
|
||||||
|
.expect("runtime snapshot json");
|
||||||
|
assert!(
|
||||||
|
persisted_runtime["diagnostics"]
|
||||||
|
.as_array()
|
||||||
|
.unwrap()
|
||||||
|
.iter()
|
||||||
|
.any(|diagnostic| {
|
||||||
|
diagnostic["code"] == WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE
|
||||||
|
&& diagnostic["message"] == WORKER_DELETE_FAILURE_DIAGNOSTIC_MESSAGE
|
||||||
|
})
|
||||||
|
);
|
||||||
|
assert!(matches!(
|
||||||
|
runtime.delete_worker(&worker.worker_ref),
|
||||||
|
Err(RuntimeError::WorkerDeletePersistenceFailed { .. })
|
||||||
|
));
|
||||||
|
assert_eq!(
|
||||||
|
runtime
|
||||||
|
.diagnostics()
|
||||||
|
.unwrap()
|
||||||
|
.iter()
|
||||||
|
.filter(|diagnostic| {
|
||||||
|
diagnostic.code == WORKER_DELETE_FAILURE_DIAGNOSTIC_CODE
|
||||||
|
&& diagnostic.worker_ref.as_ref() == Some(&worker.worker_ref)
|
||||||
|
})
|
||||||
|
.count(),
|
||||||
|
1
|
||||||
|
);
|
||||||
|
|
||||||
|
std::fs::remove_file(&worker_dir).unwrap();
|
||||||
|
std::fs::rename(&saved_worker_dir, &worker_dir).unwrap();
|
||||||
|
assert!(runtime.delete_worker(&worker.worker_ref).unwrap().deleted);
|
||||||
|
assert!(matches!(
|
||||||
|
runtime.worker_detail(&worker.worker_ref),
|
||||||
|
Err(RuntimeError::WorkerNotFound { .. })
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn delete_worker_waits_for_execution_shutdown_before_removing_snapshot() {
|
fn delete_worker_waits_for_execution_shutdown_before_removing_snapshot() {
|
||||||
let root = tempfile::tempdir().unwrap();
|
let root = tempfile::tempdir().unwrap();
|
||||||
|
|||||||
Reference in New Issue
Block a user