runtime: harden retention reconciliation and retry

This commit is contained in:
2026-08-12 00:46:39 +09:00
parent da8313fa1a
commit 5e5ce73fd0
3 changed files with 619 additions and 77 deletions
+400 -50
View File
@@ -46,6 +46,69 @@ pub struct WorkerRetentionInventory {
pub diagnostics_bytes: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
pub struct RuntimeWorkerAggregateDiagnostic {
worker_id: String,
category: String,
detail: String,
}
impl RuntimeWorkerAggregateDiagnostic {
pub fn worker_id(&self) -> &str {
&self.worker_id
}
pub fn category(&self) -> &str {
&self.category
}
pub fn detail(&self) -> &str {
&self.detail
}
}
/// Opaque host-derived inventory snapshot. Callers can inspect but cannot
/// construct or alter the Runtime/Workspace scope used by reconciliation.
#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
pub struct WorkerRetentionInventorySnapshot {
workspace_id: String,
runtime_id: String,
workers: Vec<WorkerRetentionInventory>,
diagnostics: Vec<RuntimeWorkerAggregateDiagnostic>,
}
impl WorkerRetentionInventorySnapshot {
pub(crate) fn new(
workspace_id: String,
runtime_id: String,
workers: Vec<WorkerRetentionInventory>,
diagnostics: Vec<RuntimeWorkerAggregateDiagnostic>,
) -> Self {
Self {
workspace_id,
runtime_id,
workers,
diagnostics,
}
}
pub fn workspace_id(&self) -> &str {
&self.workspace_id
}
pub fn runtime_id(&self) -> &str {
&self.runtime_id
}
pub fn workers(&self) -> &[WorkerRetentionInventory] {
&self.workers
}
pub fn diagnostics(&self) -> &[RuntimeWorkerAggregateDiagnostic] {
&self.diagnostics
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkerRetentionExecutionRequest {
pub operation_id: String,
@@ -113,12 +176,6 @@ pub(crate) trait WorkerRetentionProvider: Send + Sync {
&self,
request: &WorkerRetentionExecutionRequest,
) -> Result<WorkerRetentionExecutionResult, RuntimeError>;
fn completed(
&self,
operation_id: &str,
input_fingerprint: &str,
) -> Result<Option<WorkerRetentionExecutionResult>, RuntimeError>;
}
/// Filesystem provider for the canonical Runtime Worker aggregate.
@@ -134,6 +191,137 @@ impl FsWorkerRetentionProvider {
}
}
pub(crate) fn completed_for(
&self,
request: &WorkerRetentionExecutionRequest,
) -> Result<Option<WorkerRetentionExecutionResult>, RuntimeError> {
let path = self.operation_path(&request.operation_id)?;
let Some(receipt) = read_operation_receipt(&path)? else {
return Ok(None);
};
validate_receipt_request(&receipt, request)?;
Ok(receipt.result.source_removed.then_some(receipt.result))
}
pub(crate) fn snapshot(
&self,
workspace_id: &str,
runtime_id: &str,
) -> Result<WorkerRetentionInventorySnapshot, RuntimeError> {
fs::create_dir_all(&self.runtime_root).map_err(|source| RuntimeError::StoreIo {
operation: "prepare Worker retention inventory",
path: self.runtime_root.clone(),
source,
})?;
let lock_path = self.runtime_root.join(RETENTION_LOCK);
let lock = OpenOptions::new()
.create(true)
.read(true)
.write(true)
.open(&lock_path)
.map_err(|source| RuntimeError::StoreIo {
operation: "lock Worker retention inventory",
path: lock_path.clone(),
source,
})?;
lock.lock().map_err(|source| RuntimeError::StoreIo {
operation: "lock Worker retention inventory",
path: lock_path,
source,
})?;
let workers_root = self.runtime_root.join("workers");
if !workers_root.is_dir() {
return Ok(WorkerRetentionInventorySnapshot::new(
workspace_id.to_string(),
runtime_id.to_string(),
Vec::new(),
Vec::new(),
));
}
let mut entries = fs::read_dir(&workers_root)
.map_err(|source| RuntimeError::StoreIo {
operation: "scan Worker retention inventory",
path: workers_root.clone(),
source,
})?
.collect::<Result<Vec<_>, _>>()
.map_err(|source| RuntimeError::StoreIo {
operation: "scan Worker retention inventory",
path: workers_root.clone(),
source,
})?;
entries.sort_by_key(|entry| entry.file_name());
let mut workers = Vec::new();
let mut diagnostics = Vec::new();
for entry in entries {
let raw_id = entry.file_name().to_string_lossy().to_string();
let bounded_id = bounded_diagnostic_id(&raw_id);
let file_type = match entry.file_type() {
Ok(file_type) => file_type,
Err(_) => {
diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_metadata_unreadable",
));
continue;
}
};
if !file_type.is_dir() || file_type.is_symlink() {
diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_entry_unsupported",
));
continue;
}
let Ok(worker_number) = raw_id.parse::<u64>() else {
diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_worker_id_invalid",
));
continue;
};
let worker_id = WorkerId::new(worker_number);
let worker_dir = self.worker_dir(worker_id);
let snapshot: WorkerGenerationSnapshot = match read_json(
&worker_dir.join("worker.json"),
"scan Worker retention inventory",
) {
Ok(snapshot) => snapshot,
Err(_) => {
diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_worker_record_corrupt",
));
continue;
}
};
if snapshot.workspace_id.as_deref() != Some(workspace_id) {
diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_workspace_mismatch",
));
continue;
}
match self.inventory(workspace_id, runtime_id, worker_id, snapshot.run_generation) {
Ok(item) => workers.push(item),
Err(_) => diagnostics.push(runtime_aggregate_diagnostic(
&bounded_id,
"aggregate_session_inventory_failed",
)),
}
}
workers.sort_by_key(|item| item.worker_id);
diagnostics.sort_by(|left, right| {
(&left.worker_id, &left.category).cmp(&(&right.worker_id, &right.category))
});
Ok(WorkerRetentionInventorySnapshot::new(
workspace_id.to_string(),
runtime_id.to_string(),
workers,
diagnostics,
))
}
pub(crate) fn recover_after_source_removal(
&self,
request: &WorkerRetentionExecutionRequest,
@@ -192,6 +380,19 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
if !worker_dir.is_dir() {
return Err(RuntimeError::WorkerNotFound { worker_id });
}
let worker: WorkerGenerationSnapshot = read_json(
&worker_dir.join("worker.json"),
"inventory Worker retention",
)?;
if worker.workspace_id.as_deref() != Some(workspace_id) {
return Err(RuntimeError::WorkerNotFound { worker_id });
}
if worker.run_generation != run_generation {
return Err(RuntimeError::InvalidRequest(format!(
"Worker retention inventory expected generation {run_generation}, current generation is {}",
worker.run_generation
)));
}
let session_dir = worker_dir.join("session");
let (session_id, segment_ids, session_bytes) = if session_dir.is_dir() {
let manifest: CanonicalSessionManifest = read_json(
@@ -264,11 +465,9 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
source,
})?;
let pending = read_operation_receipt(
&self.operation_path(&request.operation_id)?,
&request.input_fingerprint,
)?;
let pending = read_operation_receipt(&self.operation_path(&request.operation_id)?)?;
if let Some(receipt) = &pending {
validate_receipt_request(receipt, request)?;
if receipt.result.source_removed {
return Ok(receipt.result.clone());
}
@@ -293,6 +492,11 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
}
let snapshot: WorkerGenerationSnapshot =
read_json(&worker_dir.join("worker.json"), "execute Worker retention")?;
if snapshot.workspace_id.as_deref() != Some(request.workspace_id.as_str()) {
return Err(RuntimeError::WorkerNotFound {
worker_id: request.worker_id,
});
}
if snapshot.run_generation != request.expected_run_generation {
return Err(RuntimeError::InvalidRequest(format!(
"Worker retention plan expected generation {}, current generation is {}",
@@ -326,6 +530,7 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
};
let mut receipt = RetentionOperationReceipt {
schema_version: OPERATION_SCHEMA_VERSION,
request: request.clone(),
result: result.clone(),
};
// Pending receipt makes the delete/final-receipt crash window
@@ -355,22 +560,12 @@ impl WorkerRetentionProvider for FsWorkerRetentionProvider {
)?;
Ok(result)
}
fn completed(
&self,
operation_id: &str,
input_fingerprint: &str,
) -> Result<Option<WorkerRetentionExecutionResult>, RuntimeError> {
let path = self.operation_path(operation_id)?;
let Some(receipt) = read_operation_receipt(&path, input_fingerprint)? else {
return Ok(None);
};
Ok(receipt.result.source_removed.then_some(receipt.result))
}
}
#[derive(Deserialize)]
struct WorkerGenerationSnapshot {
#[serde(default)]
workspace_id: Option<String>,
#[serde(default)]
run_generation: u64,
}
@@ -383,13 +578,11 @@ struct CanonicalSessionManifest {
#[derive(Serialize, Deserialize)]
struct RetentionOperationReceipt {
schema_version: u32,
request: WorkerRetentionExecutionRequest,
result: WorkerRetentionExecutionResult,
}
fn read_operation_receipt(
path: &Path,
input_fingerprint: &str,
) -> Result<Option<RetentionOperationReceipt>, RuntimeError> {
fn read_operation_receipt(path: &Path) -> Result<Option<RetentionOperationReceipt>, RuntimeError> {
if !path.is_file() {
return Ok(None);
}
@@ -404,16 +597,26 @@ fn read_operation_receipt(
),
});
}
if receipt.result.input_fingerprint != input_fingerprint {
return Err(RuntimeError::InvalidRequest(format!(
"retention operation {} was already used with different input",
receipt.result.operation_id
)));
}
Ok(Some(receipt))
}
#[derive(Serialize, Deserialize)]
fn validate_receipt_request(
receipt: &RetentionOperationReceipt,
request: &WorkerRetentionExecutionRequest,
) -> Result<(), RuntimeError> {
if &receipt.request != request
|| receipt.result.operation_id != request.operation_id
|| receipt.result.input_fingerprint != request.input_fingerprint
{
return Err(RuntimeError::InvalidRequest(format!(
"retention operation {} was already used with different input",
request.operation_id
)));
}
Ok(())
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
struct DiagnosticsArchiveManifest {
schema_version: u32,
operation_id: String,
@@ -426,6 +629,36 @@ struct DiagnosticsArchiveManifest {
content_file_count: u64,
}
fn runtime_aggregate_diagnostic(
worker_id: &str,
category: &str,
) -> RuntimeWorkerAggregateDiagnostic {
RuntimeWorkerAggregateDiagnostic {
worker_id: worker_id.to_string(),
category: category.to_string(),
detail: "Canonical Runtime Worker aggregate requires diagnostic reconciliation".to_string(),
}
}
fn bounded_diagnostic_id(value: &str) -> String {
let sanitized = value
.chars()
.take(64)
.map(|character| {
if character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.') {
character
} else {
'_'
}
})
.collect::<String>();
if sanitized.is_empty() {
"unknown".to_string()
} else {
sanitized
}
}
fn validate_request(request: &WorkerRetentionExecutionRequest) -> Result<(), RuntimeError> {
validate_id("operation_id", &request.operation_id)?;
validate_id("workspace_id", &request.workspace_id)?;
@@ -611,16 +844,7 @@ fn commit_diagnostics_archive(
let files = diagnostics_files(worker_dir, "archive Worker diagnostics")?;
let target = provider.diagnostics_dir(&request.operation_id)?;
if target.exists() {
let manifest: DiagnosticsArchiveManifest = read_json(
&target.join("manifest.json"),
"verify Worker diagnostics archive",
)?;
if manifest.input_fingerprint != request.input_fingerprint {
return Err(RuntimeError::InvalidRequest(format!(
"diagnostics archive operation {} was reused with different input",
request.operation_id
)));
}
validate_existing_diagnostics_archive(&target, request)?;
return Ok(());
}
let parent = target.parent().ok_or_else(|| RuntimeError::StoreCorrupt {
@@ -677,7 +901,8 @@ fn commit_diagnostics_archive(
path: target.clone(),
source,
})?;
sync_directory(parent, "archive Worker diagnostics")
sync_directory(parent, "archive Worker diagnostics")?;
validate_existing_diagnostics_archive(&target, request)
})();
if result.is_err() {
let _ = fs::remove_dir_all(&staging);
@@ -685,6 +910,42 @@ fn commit_diagnostics_archive(
result
}
fn validate_existing_diagnostics_archive(
target: &Path,
request: &WorkerRetentionExecutionRequest,
) -> Result<(), RuntimeError> {
let manifest: DiagnosticsArchiveManifest = read_json(
&target.join("manifest.json"),
"verify Worker diagnostics archive",
)?;
if manifest.schema_version != ARCHIVE_SCHEMA_VERSION
|| manifest.operation_id != request.operation_id
|| manifest.workspace_id != request.workspace_id
|| manifest.source_runtime_id != request.source_runtime_id
|| manifest.source_worker_id != request.worker_id
|| manifest.input_fingerprint != request.input_fingerprint
{
return Err(RuntimeError::InvalidRequest(format!(
"diagnostics archive operation {} does not match the retention request",
request.operation_id
)));
}
let mut files = collect_files(target, "verify Worker diagnostics archive")?;
files.retain(|(relative, _)| relative != Path::new("manifest.json"));
let (checksum, bytes, count) = checksum_files(&files, "verify Worker diagnostics archive")?;
if checksum != manifest.content_checksum_sha256
|| bytes != manifest.content_bytes
|| count != manifest.content_file_count
{
return Err(RuntimeError::StoreCorrupt {
operation: "verify Worker diagnostics archive",
path: target.to_path_buf(),
message: "diagnostics archive checksum or content summary mismatch".to_string(),
});
}
Ok(())
}
fn diagnostics_files(
worker_dir: &Path,
operation: &'static str,
@@ -1001,7 +1262,7 @@ mod tests {
let worker = root.join("workers").join(worker_id.to_string());
write_json(
&worker.join("worker.json"),
&serde_json::json!({"run_generation": generation}),
&serde_json::json!({"workspace_id": "workspace-a", "run_generation": generation}),
);
write_json(
&worker.join("session/session.json"),
@@ -1102,6 +1363,36 @@ mod tests {
);
}
#[test]
fn target_inventory_and_execute_reject_cross_workspace_aggregate() {
let temp = tempfile::tempdir().unwrap();
let worker_id = WorkerId::new(16);
source(temp.path(), worker_id, 3);
let provider = FsWorkerRetentionProvider::new(temp.path());
assert!(matches!(
provider.inventory("other-workspace", "runtime-a", worker_id, 3),
Err(RuntimeError::WorkerNotFound { .. })
));
let mut request = request(worker_id, 3, SessionDisposition::Purge);
request.workspace_id = "other-workspace".to_string();
assert!(matches!(
provider.execute(&request),
Err(RuntimeError::WorkerNotFound { .. })
));
assert!(temp.path().join("workers/16/session").is_dir());
assert!(
!temp
.path()
.join("retention/operations/operation-a.json")
.exists()
);
request.workspace_id = "workspace-a".to_string();
provider.execute(&request).unwrap();
request.workspace_id = "other-workspace".to_string();
assert!(provider.execute(&request).is_err());
}
#[test]
fn purge_removes_aggregate_and_rejects_stale_generation() {
let temp = tempfile::tempdir().unwrap();
@@ -1141,12 +1432,71 @@ mod tests {
let recovered = provider.execute(&request).unwrap();
assert_eq!(recovered, completed);
assert!(
provider
.completed("operation-a", "fingerprint-a")
.unwrap()
.is_some()
assert!(provider.completed_for(&request).unwrap().is_some());
}
#[test]
fn provider_snapshot_scans_aggregate_storage_independent_of_runtime_catalog() {
let temp = tempfile::tempdir().unwrap();
source(temp.path(), WorkerId::new(13), 2);
source(temp.path(), WorkerId::new(14), 1);
write_json(
&temp.path().join("workers/14/worker.json"),
&serde_json::json!({"workspace_id": "other-workspace", "run_generation": 1}),
);
fs::create_dir_all(temp.path().join("workers/not-a-worker")).unwrap();
fs::write(
temp.path().join("workers/not-a-worker/worker.json"),
b"not-json",
)
.unwrap();
fs::create_dir_all(temp.path().join("workers/15")).unwrap();
fs::write(temp.path().join("workers/15/worker.json"), b"not-json").unwrap();
let provider = FsWorkerRetentionProvider::new(temp.path());
let snapshot = provider.snapshot("workspace-a", "runtime-a").unwrap();
assert_eq!(snapshot.workers().len(), 1);
assert_eq!(snapshot.workers()[0].worker_id, WorkerId::new(13));
assert!(snapshot.diagnostics().iter().any(|diagnostic| {
diagnostic.worker_id() == "14"
&& diagnostic.category() == "aggregate_workspace_mismatch"
}));
assert!(snapshot.diagnostics().iter().any(|diagnostic| {
diagnostic.worker_id() == "not-a-worker"
&& diagnostic.category() == "aggregate_worker_id_invalid"
}));
assert!(snapshot.diagnostics().iter().any(|diagnostic| {
diagnostic.worker_id() == "15"
&& diagnostic.category() == "aggregate_worker_record_corrupt"
}));
}
#[test]
fn diagnostics_retry_rejects_corrupt_existing_archive_before_source_delete() {
let temp = tempfile::tempdir().unwrap();
let worker_id = WorkerId::new(12);
source(temp.path(), worker_id, 1);
let provider = FsWorkerRetentionProvider::new(temp.path());
let mut request = request(worker_id, 1, SessionDisposition::Archive);
request.diagnostics_disposition = DiagnosticsDisposition::Retain;
provider.execute(&request).unwrap();
let receipt_path = temp.path().join("retention/operations/operation-a.json");
let mut receipt: RetentionOperationReceipt =
serde_json::from_slice(&fs::read(&receipt_path).unwrap()).unwrap();
receipt.result.source_removed = false;
fs::write(&receipt_path, serde_json::to_vec_pretty(&receipt).unwrap()).unwrap();
source(temp.path(), worker_id, 1);
fs::write(
temp.path()
.join("archives/diagnostics/operation-a/runs/1/worker.out.log"),
b"corrupt\n",
)
.unwrap();
assert!(provider.execute(&request).is_err());
assert!(temp.path().join("workers/12/session").is_dir());
assert!(provider.completed_for(&request).unwrap().is_none());
}
#[test]
+25 -4
View File
@@ -29,7 +29,7 @@ use crate::observation::{WorkerObservationCursor, WorkerObservationEvent};
#[cfg(feature = "fs-store")]
use crate::retention::{
FsWorkerRetentionProvider, WorkerRetentionExecutionRequest, WorkerRetentionExecutionResult,
WorkerRetentionInventory, WorkerRetentionProvider,
WorkerRetentionInventory, WorkerRetentionInventorySnapshot, WorkerRetentionProvider,
};
use protocol::subscription::{
EventSubscriptionSelector, SubscriptionEventPayload, SubscriptionSnapshot,
@@ -1706,6 +1706,29 @@ impl Runtime {
)
}
/// Enumerate host-authoritative Runtime inventory for Backend orphan
/// reconciliation. Runtime identity and Workspace scope are derived here,
/// not accepted in a diagnostic payload.
#[cfg(feature = "fs-store")]
pub fn list_worker_retention_inventory(
&self,
workspace_id: &str,
) -> Result<WorkerRetentionInventorySnapshot, RuntimeError> {
let state = self.lock()?;
let runtime_id = state.runtime_identity.as_deref().ok_or_else(|| {
RuntimeError::InvalidRequest(
"Runtime identity is not bound for Worker retention".to_string(),
)
})?;
let store = state.fs_store().ok_or_else(|| {
RuntimeError::InvalidRequest(
"Worker retention inventory requires an fs-backed Runtime".to_string(),
)
})?;
let provider = FsWorkerRetentionProvider::new(store.runtime_dir());
provider.snapshot(workspace_id, runtime_id)
}
/// Execute a Backend-resolved retention plan. Only stopped Workers are
/// eligible. Provider receipt lookup happens before live lookup so exact
/// retries converge after aggregate removal.
@@ -1731,9 +1754,7 @@ impl Runtime {
)
})?;
let provider = FsWorkerRetentionProvider::new(store.runtime_dir());
if let Some(completed) =
provider.completed(&request.operation_id, &request.input_fingerprint)?
{
if let Some(completed) = provider.completed_for(request)? {
state.workers.remove(&request.worker_id);
state.persist_runtime_snapshot()?;
return Ok(completed);