diff --git a/crates/manifest/src/lib.rs b/crates/manifest/src/lib.rs index 5112fa6b..20d052ff 100644 --- a/crates/manifest/src/lib.rs +++ b/crates/manifest/src/lib.rs @@ -976,7 +976,8 @@ struct LegacyMemoryConfig { consolidation_threshold_bytes: Option, } -const RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION: u64 = 2; +const RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION: u64 = 3; +const PREVIOUS_RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION: u64 = 2; /// Serialize a resolved Worker Manifest for durable Worker-specific storage. pub fn write_persisted_worker_manifest_snapshot( @@ -1006,7 +1007,9 @@ pub fn read_persisted_worker_manifest_snapshot( "resolved Worker manifest snapshot schema_version must be an integer", )) })?; - if version != RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION { + if version != RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION + && version != PREVIOUS_RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION + { return Err(serde_json::Error::io(std::io::Error::new( std::io::ErrorKind::InvalidData, format!("unsupported resolved Worker manifest snapshot schema version {version}"), @@ -1018,7 +1021,7 @@ pub fn read_persisted_worker_manifest_snapshot( "resolved Worker manifest snapshot contains unknown fields", ))); } - let manifest = object.get("manifest").cloned().ok_or_else(|| { + let mut manifest = object.get("manifest").cloned().ok_or_else(|| { serde_json::Error::io(std::io::Error::new( std::io::ErrorKind::InvalidData, "resolved Worker manifest snapshot is missing manifest", @@ -1033,6 +1036,9 @@ pub fn read_persisted_worker_manifest_snapshot( "current resolved Worker manifest contains removed top-level memory authority", ))); } + if version == PREVIOUS_RESOLVED_MANIFEST_SNAPSHOT_SCHEMA_VERSION { + migrate_legacy_manifest_authority(&mut manifest)?; + } return validate_persisted_worker_manifest(serde_json::from_value(manifest)?); } @@ -1055,6 +1061,49 @@ fn validate_persisted_worker_manifest( Ok(manifest) } +fn migrate_legacy_manifest_authority( + manifest: &mut serde_json::Value, +) -> Result<(), serde_json::Error> { + let root = manifest.as_object_mut().ok_or_else(|| { + serde_json::Error::io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "resolved Worker manifest must be an object", + )) + })?; + root.remove("plugins"); + if let Some(feature) = root.get_mut("feature") { + let feature = feature.as_object_mut().ok_or_else(|| { + serde_json::Error::io(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "resolved Worker manifest feature must be an object", + )) + })?; + feature.remove("plugins"); + feature.remove("ticket_orchestration"); + if let Some(workers) = feature.remove("workers") { + feature + .entry("sub_worker".to_string()) + .or_insert_with(|| workers.clone()); + feature.entry("worker".to_string()).or_insert(workers); + } + if let Some(ticket) = feature + .get_mut("ticket") + .and_then(serde_json::Value::as_object_mut) + && let Some(access) = ticket.remove("access") + && ticket + .get("enabled") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false) + && access.as_str() == Some("lifecycle") + { + ticket.insert("authoring".to_string(), serde_json::Value::Bool(true)); + ticket.insert("thread".to_string(), serde_json::Value::Bool(true)); + ticket.insert("workflow".to_string(), serde_json::Value::Bool(true)); + } + } + Ok(()) +} + fn migrate_legacy_resolved_manifest_snapshot( mut snapshot: serde_json::Value, ) -> Result { @@ -1080,7 +1129,7 @@ fn migrate_legacy_resolved_manifest_snapshot( .remove("memory") .unwrap_or_else(|| serde_json::json!({})), )?; - let enabled = legacy_feature_memory.enabled; + let requested_enabled = legacy_feature_memory.enabled; let staging_tools = legacy_feature_memory.staging; let legacy_memory: LegacyMemoryConfig = @@ -1103,9 +1152,14 @@ fn migrate_legacy_resolved_manifest_snapshot( ))); } }; - if !enabled { + if !requested_enabled { workspace_settings = None; } + // Legacy standalone manifests could enable process-local Memory without a + // Workspace-owned settings snapshot. That authority no longer exists, so + // migration safely disables Memory instead of treating the whole Worker + // snapshot as corrupt. + let enabled = requested_enabled && workspace_settings.is_some(); let extraction_enabled = legacy_memory.extract_threshold.is_some(); if legacy_memory.consolidation_model.is_some() { return Err(serde_json::Error::io(std::io::Error::new( @@ -1151,6 +1205,7 @@ fn migrate_legacy_resolved_manifest_snapshot( .insert("workspace_settings".to_string(), workspace_settings); } feature.insert("memory".to_string(), resolved); + migrate_legacy_manifest_authority(&mut snapshot)?; validate_persisted_worker_manifest(serde_json::from_value(snapshot)?) } @@ -1601,7 +1656,7 @@ model_id = "claude-sonnet-4-20250514" "Français" ); let current = write_persisted_worker_manifest_snapshot(&migrated).unwrap(); - assert_eq!(current["schema_version"], 2); + assert_eq!(current["schema_version"], 3); assert!(current["manifest"].get("memory").is_none()); let mut disabled = @@ -1617,6 +1672,61 @@ model_id = "claude-sonnet-4-20250514" assert!(disabled.feature.memory.workspace_settings.is_none()); } + #[test] + fn persisted_manifest_adapter_drops_removed_plugin_authority() { + let manifest = WorkerManifest::from_toml(MINIMAL_REQUIRED).unwrap(); + let mut versioned = write_persisted_worker_manifest_snapshot(&manifest).unwrap(); + versioned["schema_version"] = serde_json::json!(2); + versioned["manifest"]["feature"]["plugins"] = serde_json::json!({ "enabled": true }); + versioned["manifest"]["feature"] + .as_object_mut() + .unwrap() + .remove("sub_worker"); + versioned["manifest"]["feature"] + .as_object_mut() + .unwrap() + .remove("worker"); + versioned["manifest"]["feature"]["workers"] = serde_json::json!({ "enabled": true }); + versioned["manifest"]["feature"]["ticket"] = + serde_json::json!({ "enabled": true, "access": "lifecycle" }); + versioned["manifest"]["feature"]["ticket_orchestration"] = + serde_json::json!({ "enabled": false }); + versioned["manifest"]["plugins"] = serde_json::json!({ + "enabled": ["legacy-plugin"], + "config": { "legacy-plugin": { "legacy": true } } + }); + + let restored = read_persisted_worker_manifest_snapshot(versioned).unwrap(); + let current = write_persisted_worker_manifest_snapshot(&restored).unwrap(); + assert_eq!(current["schema_version"], 3); + assert!(current["manifest"].get("plugins").is_none()); + assert!(current["manifest"]["feature"].get("plugins").is_none()); + assert!(current["manifest"]["feature"].get("workers").is_none()); + assert_eq!( + current["manifest"]["feature"]["sub_worker"]["enabled"], + true + ); + assert_eq!(current["manifest"]["feature"]["worker"]["enabled"], true); + assert_eq!(current["manifest"]["feature"]["ticket"]["authoring"], true); + assert_eq!(current["manifest"]["feature"]["ticket"]["thread"], true); + assert_eq!(current["manifest"]["feature"]["ticket"]["workflow"], true); + + let mut legacy = serde_json::to_value(manifest).unwrap(); + legacy.as_object_mut().unwrap().remove("memory"); + legacy["feature"]["memory"] = serde_json::json!({ + "enabled": true, + "staging": false + }); + legacy["feature"]["plugins"] = serde_json::json!({ "enabled": false }); + legacy["plugins"] = serde_json::json!({ "enabled": [] }); + let legacy = read_persisted_worker_manifest_snapshot(legacy).unwrap(); + let current = write_persisted_worker_manifest_snapshot(&legacy).unwrap(); + assert_eq!( + current["manifest"]["feature"]["memory"]["profile"]["enabled"], + false + ); + } + #[test] fn persisted_manifest_adapter_rejects_mixed_or_future_authority() { let manifest = @@ -1659,7 +1769,7 @@ model_id = "claude-sonnet-4-20250514" assert!( read_persisted_worker_manifest_snapshot(serde_json::json!({ - "schema_version": 3, + "schema_version": 4, "manifest": manifest, })) .is_err() diff --git a/crates/worker-runtime/src/fs_store.rs b/crates/worker-runtime/src/fs_store.rs index fdb33c1d..4c7c8301 100644 --- a/crates/worker-runtime/src/fs_store.rs +++ b/crates/worker-runtime/src/fs_store.rs @@ -15,7 +15,9 @@ use std::io::{BufReader, Write}; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicU64, Ordering}; -const SCHEMA_VERSION: u32 = 4; +const SCHEMA_VERSION: u32 = 5; +const PREVIOUS_SCHEMA_VERSION: u32 = 4; +const PRE_EXECUTION_SCHEMA_VERSION: u32 = 3; const RUNTIME_FILE: &str = "runtime.json"; const WORKERS_DIR: &str = "workers"; const WORKER_FILE: &str = "worker.json"; @@ -370,8 +372,8 @@ fn plan_runtime_store_migration( format!("Runtime store schema version {schema_version} is out of range"), ) })?; - let staging = migration_sibling(root, "schema-v4-staging")?; - let backup = migration_sibling(root, "pre-schema-v4-backup")?; + let staging = migration_sibling(root, "schema-v5-staging")?; + let backup = migration_sibling(root, "pre-schema-v5-backup")?; if staging.exists() || backup.exists() { return Err(runtime_store_corrupt( root, @@ -397,11 +399,14 @@ fn plan_runtime_store_migration( }; return Ok((plan, Vec::new())); } - if current_schema_version != 3 { + if !matches!( + current_schema_version, + PRE_EXECUTION_SCHEMA_VERSION | PREVIOUS_SCHEMA_VERSION + ) { return Err(runtime_store_corrupt( &runtime_path, format!( - "unsupported Runtime store schema version {schema_version}; expected 3 or {SCHEMA_VERSION}" + "unsupported Runtime store schema version {schema_version}; expected {PRE_EXECUTION_SCHEMA_VERSION}, {PREVIOUS_SCHEMA_VERSION}, or {SCHEMA_VERSION}" ), )); } @@ -428,6 +433,16 @@ fn plan_runtime_store_migration( runtime_store_corrupt(&source_dir, "Worker directory is not UTF-8".to_string()) })?; let snapshot_path = source_dir.join(WORKER_FILE); + if !snapshot_path + .try_exists() + .map_err(|source| RuntimeError::StoreIo { + operation: "inspect Worker snapshot", + path: snapshot_path.clone(), + source, + })? + { + continue; + } let snapshot: serde_json::Value = read_json(&snapshot_path, "read Worker snapshot")?; let (worker_id, workspace_id, legacy_mapping) = if current_schema_version == 1 { let legacy_worker_id = name.parse::().map_err(|_| { @@ -663,6 +678,42 @@ fn migrate_worker_document( object.insert("working_directory".to_string(), working_directory); } } + let legacy_materialization = object + .get("working_directory") + .and_then(|working_directory| working_directory.get("summary")) + .and_then(|summary| summary.get("materializer_kind")) + .and_then(serde_json::Value::as_str) + .is_some_and(|kind| matches!(kind, "runtime_git_cache" | "local_git_worktree")); + if legacy_materialization { + object.insert("working_directory".to_string(), serde_json::Value::Null); + } + if let Some(profile_source) = object + .get_mut("request") + .and_then(serde_json::Value::as_object_mut) + .and_then(|request| request.get_mut("profile_source")) + .and_then(serde_json::Value::as_object_mut) + && profile_source + .get("kind") + .and_then(serde_json::Value::as_str) + == Some("http") + { + let archive = profile_source + .get_mut("location") + .and_then(serde_json::Value::as_object_mut) + .and_then(|location| location.remove("archive")) + .ok_or_else(|| { + runtime_store_corrupt( + snapshot_path, + "legacy HTTP profile source is missing its archive".to_string(), + ) + })?; + profile_source.clear(); + profile_source.insert( + "kind".to_string(), + serde_json::Value::String("workspace_config".to_string()), + ); + profile_source.insert("archive".to_string(), archive); + } object.insert( "schema_version".to_string(), serde_json::Value::from(SCHEMA_VERSION), @@ -1048,8 +1099,8 @@ fn migrate_runtime_store( if !plan.migration_required { return Ok(plan); } - let staging = migration_sibling(root, "schema-v4-staging")?; - let backup = migration_sibling(root, "pre-schema-v4-backup")?; + let staging = migration_sibling(root, "schema-v5-staging")?; + let backup = migration_sibling(root, "pre-schema-v5-backup")?; if staging.exists() || backup.exists() { return Err(runtime_store_corrupt( root, @@ -1491,3 +1542,112 @@ fn sync_directory(path: &Path, operation: &'static str) -> Result<(), RuntimeErr source, }) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn schema_v4_migration_plan_ignores_orphan_worker_directories() { + let root = tempfile::tempdir().unwrap(); + fs::write( + root.path().join(RUNTIME_FILE), + serde_json::to_vec_pretty(&serde_json::json!({ + "schema_version": PREVIOUS_SCHEMA_VERSION, + "display_name": null, + "backend": "fs_store", + "status": "running", + "next_diagnostic_id": 1, + "config_bundles": {}, + "workspace_owners": {}, + "diagnostics": [] + })) + .unwrap(), + ) + .unwrap(); + fs::create_dir_all(root.path().join(WORKERS_DIR).join("orphan").join("session")).unwrap(); + fs::write( + root.path() + .join(WORKERS_DIR) + .join("orphan") + .join("session") + .join("history.json"), + b"[]", + ) + .unwrap(); + let (plan, _) = plan_runtime_store_migration(root.path(), "runtime-test").unwrap(); + + assert!(plan.migration_required); + assert_eq!(plan.current_schema_version, PREVIOUS_SCHEMA_VERSION); + assert_eq!(plan.target_schema_version, SCHEMA_VERSION); + assert_eq!(plan.worker_count, 0); + } + + #[test] + fn schema_v4_worker_migration_discards_unsupported_linked_worktree_binding() { + let source = serde_json::json!({ + "schema_version": 4, + "request": { + "profile_source": { + "kind": "http", + "location": { + "url": "https://workspace.example.test/archive", + "etag": "profile-source:test", + "archive": { + "id": "profiles-v1", + "digest": "sha256:test", + "size_bytes": 1, + "source_graph": { + "source_count": 1, + "total_source_bytes": 1, + "entrypoints": {}, + "import_count": 0 + } + } + } + } + }, + "working_directory": { + "summary": { + "materializer_kind": "runtime_git_cache" + } + } + }); + let path = Path::new("worker.json"); + + let migrated = + migrate_worker_document(source, PREVIOUS_SCHEMA_VERSION, None, path).unwrap(); + + assert_eq!(migrated["schema_version"], SCHEMA_VERSION); + assert_eq!(migrated["status"], "stopped"); + assert_eq!(migrated["working_directory"], serde_json::Value::Null); + assert_eq!( + migrated["request"]["profile_source"]["kind"], + "workspace_config" + ); + assert_eq!( + migrated["request"]["profile_source"]["archive"]["id"], + "profiles-v1" + ); + assert_eq!(migrated["execution"]["restore_intent"], "explicit"); + } + + #[test] + fn schema_v4_worker_migration_preserves_runtime_clone_observation() { + let source = serde_json::json!({ + "schema_version": 4, + "working_directory": { + "summary": { + "materializer_kind": "runtime_git_clone" + } + } + }); + let expected = source["working_directory"].clone(); + let path = Path::new("worker.json"); + + let migrated = + migrate_worker_document(source, PREVIOUS_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 7c05ecb7..f683870f 100644 --- a/crates/worker-runtime/src/main.rs +++ b/crates/worker-runtime/src/main.rs @@ -1273,7 +1273,7 @@ mod tests { } #[test] - fn migration_dry_run_accepts_previous_schema_document_without_workers_field() { + fn migration_dry_run_accepts_supported_schema_v3_without_workers_field() { let temp = tempfile::tempdir().unwrap(); let root = temp.path().join("runtime"); std::fs::create_dir_all(root.join("workers")).unwrap(); @@ -1312,7 +1312,7 @@ mod tests { } #[test] - fn migration_dry_run_rejects_previous_schema_document_that_cannot_decode_as_v4() { + fn migration_dry_run_rejects_schema_v3_document_that_cannot_decode_as_v5() { 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 c4a77b8a..06a70637 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -6079,7 +6079,7 @@ mod tests { assert!( error .to_string() - .contains("unsupported Runtime store schema version 2; expected 3 or 4") + .contains("unsupported Runtime store schema version 2; expected 3, 4, or 5") ); let _ = std::fs::remove_dir_all(root); @@ -6121,7 +6121,7 @@ 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!(4)); + assert_eq!(worker_snapshot["schema_version"], serde_json::json!(5)); assert_eq!(worker_snapshot["status"], serde_json::json!("stopped")); assert_eq!( worker_snapshot["execution"]["binding"]["run_generation"], @@ -6556,7 +6556,7 @@ 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!(4)); + assert_eq!(migrated_json["schema_version"], serde_json::json!(5)); assert_eq!(migrated_json["status"], serde_json::json!("stopped")); assert_eq!( migrated_json["execution"]["binding"]["run_generation"],