Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion kernel/relayflowd/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -438,7 +438,10 @@ impl<C: Clock> Engine<C> {
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"))
}
}
Expand Down
45 changes: 41 additions & 4 deletions kernel/relayflowd/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(&params.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(&params.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
Expand Down
137 changes: 137 additions & 0 deletions kernel/relayflowd/src/server/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
Loading