diff --git a/kernel/relayflowd/src/engine.rs b/kernel/relayflowd/src/engine.rs index d88b6a98a..91d6796d7 100644 --- a/kernel/relayflowd/src/engine.rs +++ b/kernel/relayflowd/src/engine.rs @@ -438,7 +438,10 @@ impl Engine { Registry::open(self.data_dir.join("relayflowd.sqlite3")).context("open run registry") } - pub(super) fn run_path(&self, run_id: &str) -> PathBuf { + /// The conventional location of a run's journal. Single source of truth: + /// the resume-repair path in `server.rs` opens by this too, so a change to + /// the layout cannot leave the two disagreeing. + pub(crate) fn run_path(&self, run_id: &str) -> PathBuf { self.data_dir.join("runs").join(format!("{run_id}.sqlite3")) } } diff --git a/kernel/relayflowd/src/server.rs b/kernel/relayflowd/src/server.rs index 4a4a08a9a..8ee07ed99 100644 --- a/kernel/relayflowd/src/server.rs +++ b/kernel/relayflowd/src/server.rs @@ -159,10 +159,47 @@ fn handle_request( .map_err(|error| internal_error(error.into()))? .is_none() { - return Err(( - "run_not_found", - format!("run {} does not exist", params.run_id), - )); + // A missing registry row does NOT mean the run does not exist. + // + // `Engine::start` creates the journal, appends RunSpawned, and + // registers LAST. A crash in that window leaves a complete, + // resumable journal on disk with no `runs` row -- and refusing + // here made that run unrecoverable forever, which is a + // durability hole, not a lookup miss. It is also how #174 + // presented: a resume that exited instantly with + // `run_not_found` while the test waited on a dispatch from a + // process that was already dead. + // + // The journal is the authority; the registry is an index over + // it. So repair the index from the authority. The sibling path + // that opens a journal already does exactly this + // (`unwrap_or_else(|| self.run_path(run_id))`), and the two + // disagreeing is what made this reachable. + let path = engine.run_path(¶ms.run_id); + // Adopt ONLY a journal that both opens and says it is this run. + // + // `SqliteJournal::open` does not check that: it reads whatever + // run id the file carries and hands it back. Opening alone + // would therefore register a well-formed journal for run A + // sitting at `runs/B.sqlite3` as B -- accepting a foreign file + // on the strength of its filename. DRIVE-LOG WP-12/F7 records + // filesystem-derived run existence being deliberately replaced + // by registry-owned lookup for exactly that reason, so the + // comparison below is what keeps this a repair of the index + // rather than a reopening of that hole. + let adopted = match relayflowd_journal::SqliteJournal::open(&path) { + Ok(journal) => journal.run_id() == params.run_id, + Err(_) => false, + }; + if !adopted { + return Err(( + "run_not_found", + format!("run {} does not exist", params.run_id), + )); + } + registry + .register(¶ms.run_id, &path) + .map_err(|error| internal_error(error.into()))?; } // Per-run serialization: the load-state -> next_actions -> append // sequence must be atomic, or two concurrent resumes both see a diff --git a/kernel/relayflowd/src/server/tests.rs b/kernel/relayflowd/src/server/tests.rs index ff67796a9..cf2809cf6 100644 --- a/kernel/relayflowd/src/server/tests.rs +++ b/kernel/relayflowd/src/server/tests.rs @@ -183,6 +183,143 @@ fn run_resume_asks_the_registry_instead_of_treating_an_orphan_file_as_a_run() { assert_eq!(error.code, "run_not_found"); } +/// #174: a run whose journal exists but whose registry row does not must be +/// RESUMABLE, not refused. +/// +/// `Engine::start` creates the journal, appends RunSpawned, then registers the +/// run last. A crash in that window leaves exactly this state, and refusing it +/// made the run unrecoverable forever -- the journal was on disk, complete, and +/// nothing could reach it. That is a durability hole, not a lookup miss. +/// +/// This is the counterpart to the orphan-file test above, and the pair states +/// the rule together: adopt a journal that is real, refuse a file that is not. +#[test] +fn run_resume_adopts_a_real_journal_whose_registry_row_is_missing() { + let directory = tempdir().unwrap(); + let data_dir = directory.path(); + + // A genuine run, created the way the engine creates one. + let outcome = Engine::new(data_dir) + .start( + serde_json::from_value(json!({ + "steps": [{ + "id": "x", + "type": "deterministic", + "command": ["/bin/sh", "-c", "printf x"], + "verification": {"output_contains": "x"} + }] + })) + .unwrap(), + "test", + None, + ) + .unwrap(); + let run_id = outcome.run_id.clone(); + assert!(data_dir.join("runs").join(format!("{run_id}.sqlite3")).exists()); + + // Reproduce the crash window: the journal survives, the index entry does + // not. Removing the registry outright is the same state a SIGKILL between + // the append and the register leaves behind, and strictly harsher -- the + // row is not merely stale, there is nothing to consult at all. + for suffix in ["", "-wal", "-shm"] { + let _ = std::fs::remove_file( + data_dir.join(format!("relayflowd.sqlite3{suffix}")), + ); + } + let registry = + relayflowd_journal::Registry::open(data_dir.join("relayflowd.sqlite3")).unwrap(); + assert!(registry.lookup(&run_id).unwrap().is_none()); + + let hub = Arc::new(ProtocolHub::default()); + let (writer, _peer) = shared_writer(); + let response = request( + data_dir, + &hub, + 1, + &writer, + &format!( + r#"{{"id":"resume","verb":"run.resume","params":{{"run_id":"{run_id}"}}}}"# + ), + ); + + assert!( + response.ok, + "a real journal with no registry row must be adopted, not refused: {:?}", + response.error + ); + // And the index is repaired, so the next lookup does not depend on this + // path running again. + assert!( + registry.lookup(&run_id).unwrap().is_some(), + "resume must re-register the run it adopted" + ); +} + +/// #174, third case: a journal that is structurally VALID but belongs to a +/// different run must still be refused. +/// +/// Opening a journal does not verify whose it is -- `SqliteJournal::open` +/// reports whatever run id the file carries. Without an explicit comparison, +/// a well-formed journal for run A dropped at `runs/B.sqlite3` would be +/// registered as B purely on the strength of its filename, which is the +/// filesystem-derived existence that WP-12/F7 removed on purpose. +#[test] +fn run_resume_refuses_a_valid_journal_that_belongs_to_another_run() { + let directory = tempdir().unwrap(); + let data_dir = directory.path(); + + let outcome = Engine::new(data_dir) + .start( + serde_json::from_value(json!({ + "steps": [{ + "id": "x", + "type": "deterministic", + "command": ["/bin/sh", "-c", "printf x"], + "verification": {"output_contains": "x"} + }] + })) + .unwrap(), + "test", + None, + ) + .unwrap(); + let real_id = outcome.run_id.clone(); + + // A complete, openable journal -- under someone else's name. + let impostor = "01IMPOSTORIMPOSTORIMPOSTOR"; + std::fs::copy( + data_dir.join("runs").join(format!("{real_id}.sqlite3")), + data_dir.join("runs").join(format!("{impostor}.sqlite3")), + ) + .unwrap(); + for suffix in ["", "-wal", "-shm"] { + let _ = std::fs::remove_file(data_dir.join(format!("relayflowd.sqlite3{suffix}"))); + } + + let hub = Arc::new(ProtocolHub::default()); + let (writer, _peer) = shared_writer(); + let response = request( + data_dir, + &hub, + 1, + &writer, + &format!(r#"{{"id":"resume","verb":"run.resume","params":{{"run_id":"{impostor}"}}}}"#), + ); + + let error = response + .error + .expect("a journal belonging to another run must be refused"); + assert_eq!(error.code, "run_not_found"); + + // And nothing was adopted under the impostor id. + let registry = + relayflowd_journal::Registry::open(data_dir.join("relayflowd.sqlite3")).unwrap(); + assert!( + registry.lookup(impostor).unwrap().is_none(), + "a refused journal must not leave a registry row behind" + ); +} + /// Finding 3: a hung worker that stops heartbeating past its lease deadline — /// socket still open, so no disconnect fires — must not leave the run in /// waiting_worker forever. The reconciler journals a `lease_expired`