diff --git a/crates/worker-runtime/src/fs_store.rs b/crates/worker-runtime/src/fs_store.rs index 4c7c8301..f6f295af 100644 --- a/crates/worker-runtime/src/fs_store.rs +++ b/crates/worker-runtime/src/fs_store.rs @@ -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, 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 { + 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::().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); } diff --git a/crates/worker-runtime/src/main.rs b/crates/worker-runtime/src/main.rs index f683870f..95cd115f 100644 --- a/crates/worker-runtime/src/main.rs +++ b/crates/worker-runtime/src/main.rs @@ -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(); diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 06a70637..05cc1258 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -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")