fix: preserve worker run generations across restore

This commit is contained in:
2026-09-15 01:17:30 +09:00
parent 7210d3c202
commit 016dbd7cb1
3 changed files with 138 additions and 27 deletions
+123 -16
View File
@@ -15,8 +15,9 @@ use std::io::{BufReader, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
const SCHEMA_VERSION: u32 = 5;
const PREVIOUS_SCHEMA_VERSION: u32 = 4;
const SCHEMA_VERSION: u32 = 6;
const PREVIOUS_SCHEMA_VERSION: u32 = 5;
const EXECUTION_SCHEMA_VERSION: u32 = 4;
const PRE_EXECUTION_SCHEMA_VERSION: u32 = 3;
const RUNTIME_FILE: &str = "runtime.json";
const WORKERS_DIR: &str = "workers";
@@ -285,6 +286,7 @@ pub(crate) struct PersistedWorkerExecutionBinding {
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct PersistedWorkerExecution {
pub(crate) last_run_generation: u64,
pub(crate) binding: Option<PersistedWorkerExecutionBinding>,
pub(crate) restore_intent: WorkerRestoreIntent,
}
@@ -372,8 +374,8 @@ fn plan_runtime_store_migration(
format!("Runtime store schema version {schema_version} is out of range"),
)
})?;
let staging = migration_sibling(root, "schema-v5-staging")?;
let backup = migration_sibling(root, "pre-schema-v5-backup")?;
let staging = migration_sibling(root, "schema-v6-staging")?;
let backup = migration_sibling(root, "pre-schema-v6-backup")?;
if staging.exists() || backup.exists() {
return Err(runtime_store_corrupt(
root,
@@ -401,12 +403,12 @@ fn plan_runtime_store_migration(
}
if !matches!(
current_schema_version,
PRE_EXECUTION_SCHEMA_VERSION | PREVIOUS_SCHEMA_VERSION
PRE_EXECUTION_SCHEMA_VERSION | EXECUTION_SCHEMA_VERSION | PREVIOUS_SCHEMA_VERSION
) {
return Err(runtime_store_corrupt(
&runtime_path,
format!(
"unsupported Runtime store schema version {schema_version}; expected {PRE_EXECUTION_SCHEMA_VERSION}, {PREVIOUS_SCHEMA_VERSION}, or {SCHEMA_VERSION}"
"unsupported Runtime store schema version {schema_version}; expected {PRE_EXECUTION_SCHEMA_VERSION}, {EXECUTION_SCHEMA_VERSION}, {PREVIOUS_SCHEMA_VERSION}, or {SCHEMA_VERSION}"
),
));
}
@@ -631,6 +633,38 @@ fn migrate_v1_worker_document(
Ok(snapshot)
}
fn max_persisted_run_generation(snapshot_path: &Path) -> Result<u64, RuntimeError> {
let worker_dir = snapshot_path.parent().ok_or_else(|| {
runtime_store_corrupt(
snapshot_path,
"Worker snapshot path is missing its aggregate directory".to_string(),
)
})?;
let runs_dir = worker_dir.join("runs");
if !runs_dir
.try_exists()
.map_err(|source| runtime_io_error("inspect Worker runs", &runs_dir, source))?
{
return Ok(0);
}
let entries = fs::read_dir(&runs_dir)
.map_err(|source| runtime_io_error("read Worker runs", &runs_dir, source))?;
let mut max_generation = 0;
for entry in entries {
let entry =
entry.map_err(|source| runtime_io_error("read Worker runs", &runs_dir, source))?;
let Some(generation) = entry
.file_name()
.to_str()
.and_then(|name| name.parse::<u64>().ok())
else {
continue;
};
max_generation = max_generation.max(generation);
}
Ok(max_generation)
}
fn migrate_worker_document(
mut document: serde_json::Value,
source_schema_version: u32,
@@ -655,7 +689,7 @@ fn migrate_worker_document(
"Worker snapshot must be an object".to_string(),
)
})?;
let run_generation = object
let declared_run_generation = object
.remove("run_generation")
.map(|value| {
value.as_u64().ok_or_else(|| {
@@ -665,9 +699,45 @@ fn migrate_worker_document(
)
})
})
.transpose()?
.filter(|generation| *generation > 0);
.transpose()?;
let legacy_execution = object.remove("execution");
let execution = legacy_execution
.as_ref()
.and_then(serde_json::Value::as_object);
let persisted_last_run_generation = execution
.and_then(|execution| execution.get("last_run_generation"))
.map(|value| {
value.as_u64().ok_or_else(|| {
runtime_store_corrupt(
snapshot_path,
"Worker execution last_run_generation must be an unsigned integer".to_string(),
)
})
})
.transpose()?;
let binding_run_generation = execution
.and_then(|execution| execution.get("binding"))
.and_then(serde_json::Value::as_object)
.and_then(|binding| binding.get("run_generation"))
.map(|value| {
value.as_u64().ok_or_else(|| {
runtime_store_corrupt(
snapshot_path,
"Worker execution binding run_generation must be an unsigned integer"
.to_string(),
)
})
})
.transpose()?;
let run_generation = declared_run_generation
.into_iter()
.chain(persisted_last_run_generation)
.chain(binding_run_generation)
.chain(std::iter::once(max_persisted_run_generation(
snapshot_path,
)?))
.max()
.unwrap_or(0);
if !object.contains_key("working_directory") {
if let Some(working_directory) = legacy_execution
.as_ref()
@@ -725,9 +795,8 @@ fn migrate_worker_document(
object.insert(
"execution".to_string(),
serde_json::json!({
"binding": run_generation.map(|run_generation| {
serde_json::json!({ "run_generation": run_generation })
}),
"last_run_generation": run_generation,
"binding": null,
"restore_intent": "explicit",
}),
);
@@ -1099,8 +1168,8 @@ fn migrate_runtime_store(
if !plan.migration_required {
return Ok(plan);
}
let staging = migration_sibling(root, "schema-v5-staging")?;
let backup = migration_sibling(root, "pre-schema-v5-backup")?;
let staging = migration_sibling(root, "schema-v6-staging")?;
let backup = migration_sibling(root, "pre-schema-v6-backup")?;
if staging.exists() || backup.exists() {
return Err(runtime_store_corrupt(
root,
@@ -1373,6 +1442,18 @@ impl WorkerSnapshot {
),
});
}
if let Some(binding) = self.execution.binding.as_ref()
&& binding.run_generation != self.execution.last_run_generation
{
return Err(RuntimeError::StoreCorrupt {
operation: "read worker snapshot",
path: path.to_path_buf(),
message: format!(
"execution binding run_generation {} does not match last_run_generation {}",
binding.run_generation, self.execution.last_run_generation
),
});
}
match (self.status, self.execution.restore_intent) {
(status, WorkerRestoreIntent::Automatic) if status.is_active() => {
let Some(binding) = self.execution.binding.as_ref() else {
@@ -1583,6 +1664,32 @@ mod tests {
assert_eq!(plan.worker_count, 0);
}
#[test]
fn schema_v5_worker_migration_recovers_last_generation_from_run_aggregates() {
let root = tempfile::tempdir().unwrap();
let worker_dir = root.path().join("worker-a");
fs::create_dir_all(worker_dir.join("runs/1")).unwrap();
fs::create_dir_all(worker_dir.join("runs/7")).unwrap();
fs::create_dir_all(worker_dir.join("runs/incomplete")).unwrap();
let path = worker_dir.join(WORKER_FILE);
let source = serde_json::json!({
"schema_version": 5,
"execution": {
"binding": null,
"restore_intent": "explicit"
}
});
let migrated =
migrate_worker_document(source, PREVIOUS_SCHEMA_VERSION, None, &path).unwrap();
assert_eq!(
migrated["execution"]["last_run_generation"],
serde_json::json!(7)
);
assert_eq!(migrated["execution"]["binding"], serde_json::Value::Null);
}
#[test]
fn schema_v4_worker_migration_discards_unsupported_linked_worktree_binding() {
let source = serde_json::json!({
@@ -1616,7 +1723,7 @@ mod tests {
let path = Path::new("worker.json");
let migrated =
migrate_worker_document(source, PREVIOUS_SCHEMA_VERSION, None, path).unwrap();
migrate_worker_document(source, EXECUTION_SCHEMA_VERSION, None, path).unwrap();
assert_eq!(migrated["schema_version"], SCHEMA_VERSION);
assert_eq!(migrated["status"], "stopped");
@@ -1646,7 +1753,7 @@ mod tests {
let path = Path::new("worker.json");
let migrated =
migrate_worker_document(source, PREVIOUS_SCHEMA_VERSION, None, path).unwrap();
migrate_worker_document(source, EXECUTION_SCHEMA_VERSION, None, path).unwrap();
assert_eq!(migrated["working_directory"], expected);
}
+1 -1
View File
@@ -1312,7 +1312,7 @@ mod tests {
}
#[test]
fn migration_dry_run_rejects_schema_v3_document_that_cannot_decode_as_v5() {
fn migration_dry_run_rejects_schema_v3_document_that_cannot_decode_as_v6() {
let temp = tempfile::tempdir().unwrap();
let root = temp.path().join("runtime");
std::fs::create_dir_all(root.join("workers")).unwrap();
+14 -10
View File
@@ -2606,12 +2606,7 @@ impl RuntimeState {
let diagnostics = persisted.diagnostics;
let next_diagnostic_id = persisted.next_diagnostic_id;
for (worker_id, worker) in persisted.workers {
let run_generation = worker
.execution
.binding
.as_ref()
.map(|binding| binding.run_generation)
.unwrap_or(0);
let run_generation = worker.execution.last_run_generation;
workers.insert(
worker_id,
WorkerRecord {
@@ -3397,6 +3392,7 @@ impl WorkerRecord {
request: self.request.clone(),
status: self.status,
execution: PersistedWorkerExecution {
last_run_generation: self.run_generation,
binding: self
.execution_bound
.then_some(PersistedWorkerExecutionBinding {
@@ -6079,7 +6075,7 @@ mod tests {
assert!(
error
.to_string()
.contains("unsupported Runtime store schema version 2; expected 3, 4, or 5")
.contains("unsupported Runtime store schema version 2; expected 3, 4, 5, or 6")
);
let _ = std::fs::remove_dir_all(root);
@@ -6121,8 +6117,12 @@ mod tests {
let worker_snapshot: serde_json::Value =
serde_json::from_slice(&std::fs::read(worker_store_dir.join("worker.json")).unwrap())
.unwrap();
assert_eq!(worker_snapshot["schema_version"], serde_json::json!(5));
assert_eq!(worker_snapshot["schema_version"], serde_json::json!(6));
assert_eq!(worker_snapshot["status"], serde_json::json!("stopped"));
assert_eq!(
worker_snapshot["execution"]["last_run_generation"],
serde_json::json!(1)
);
assert_eq!(
worker_snapshot["execution"]["binding"]["run_generation"],
serde_json::json!(1)
@@ -6556,12 +6556,16 @@ mod tests {
);
let migrated_json: serde_json::Value =
serde_json::from_slice(&std::fs::read(&worker_path).unwrap()).unwrap();
assert_eq!(migrated_json["schema_version"], serde_json::json!(5));
assert_eq!(migrated_json["schema_version"], serde_json::json!(6));
assert_eq!(migrated_json["status"], serde_json::json!("stopped"));
assert_eq!(
migrated_json["execution"]["binding"]["run_generation"],
migrated_json["execution"]["last_run_generation"],
serde_json::json!(7)
);
assert_eq!(
migrated_json["execution"]["binding"],
serde_json::Value::Null
);
assert_eq!(
migrated_json["execution"]["restore_intent"],
serde_json::json!("explicit")