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
50 changes: 40 additions & 10 deletions docs/SURFACE.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,19 +108,49 @@ The herdr model: first-party helpers are just plugins that ship in the box; the

`flows build` seals a flow into a content-addressed, immutable bundle: canonical spec JSON, compiled TS with pinned deps, helper/plugin lockfile, assets, preflight declaration, identity signature — `flow@sha256:…`, pushed to a bucket/registry. `flows deploy` points a trigger at a digest; `flows run flow@sha256:…` executes from the bucket on any cell, no checkout. Preflight runs at build time for everything build-provable and again at deploy time for environment facts (credentials, workers, MCP servers). The working tree is for authoring; **production only ever runs digests.**

## 5. Invocation: three ways in
## 5. Invocation: the gate-1 CLI

A deployed flow is reachable three ways, all landing on the same digest and the same journal:
Gate 1 ships three CLI verbs over the journal protocol:

1. **Events** — `on(...)` triggers: relayfile webhooks, mentions, file changes.
2. **Schedules** — RelayCron: durable alarms + sweep.
3. **Direct call** — a flow is a *function*:
- CLI: `flows run release-note --input '{"branch":"main"}'` (or `flows run flow@sha256:…`)
- HTTP: every deployed flow is an endpoint — `POST /flows/release-note` returns the result for short flows, or a run handle (`202 + run id`) to poll/stream for long ones
- SDK: `await flows.call("release-note", input)` from any app (this is how sage and consumer apps invoke pipelines)
- Flow-to-flow: `f.dispatch("garden/implement", plan)` — same mechanism, child run with its own journal
```text
flows check [--json] <flow.yaml|spec.json>
flows run [--json] [--data-dir <dir>] <flow.yaml|spec.json>
flows resume [--json] [--data-dir <dir>] <run-id>
```

The caller always gets the same contract back: a typed result on completion, or a durable run handle it can await, stream, or abandon — the run finishes either way, journaled.
`check` compiles and preflights without starting a run. `run` performs that
same preflight before contacting `relayflowd`, then submits the compiled spec
to `<data-dir>/relayflowd.sock`; `resume` asks that daemon to continue an
existing run from its journal. The data directory defaults to `.relayflowd`.
Neither verb starts the daemon implicitly. `--json` writes one report-shaped
object to stdout while diagnostics remain on stderr.

The exit codes are part of the surface contract:

| Exit | Outcome |
|---:|---|
| `0` | The run completed with `completionReason: success`. |
| `1` | The run failed with a declared `completionReason`, or a transport, runtime, or daemon protocol error left the outcome unknown. |
| `2` | The command was refused before a journal write: invalid input, failed preflight, unreachable daemon, or a `run_not_found` resume target. |
| `3` | The run parked. `PARKED [run_parked]` names the step and its `llm` or `agent` type, and distinguishes an unavailable worker from a `needs_human` recovery wait. |

At gate 1 no `llm` or `agent` worker is attached by the CLI. Reaching either
step therefore returns the durable parked outcome instead of hanging or
reporting success. Event, schedule, deployed-digest, HTTP, SDK-call, and
flow-to-flow invocation remain later-gate surface work; they are not shipped
by this CLI.

When a worker is attached, the CLI follows the typed snapshot while its lease
is live and prints `WAITING [worker_lease]` with the step and lease deadline.
If the lease expires without a completion, the command fails closed instead of
polling forever. A manual-recovery agent whose worker dies parks in
`needs_human`; the same exit-3 report says it is waiting for human recovery.

`flows resume` reports `run_unavailable` only when relayflowd returns the
typed `run_not_found` refusal. A dropped connection, request failure, or
`journal_write_failed` response exits 1 as `protocol_error`, because the
journal may already have changed and the CLI cannot honestly claim the resume
was refused before a write.

## 6. Open surface questions (for gate-1 SDK work)

Expand Down
14 changes: 14 additions & 0 deletions kernel/relayflowd-core/tests/spec_parity.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,20 @@ fn the_kernel_parses_the_sdk_compiled_spec_and_stamps_the_same_hash() {
assert_parity(LADDER_CANONICAL, LADDER_SHA256);
}

#[test]
fn the_kernel_parses_the_deterministic_rung_and_stamps_the_same_hash() {
assert_parity(
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../testdata/hello-deterministic.spec.canonical.json"
)),
include_str!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/../../testdata/hello-deterministic.spec.sha256"
)),
);
}

#[test]
fn the_kernel_parses_the_rung_b_spec_and_stamps_the_same_hash() {
assert_parity(
Expand Down
14 changes: 12 additions & 2 deletions kernel/relayflowd/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ mod drive;
mod effects;
mod model;
mod remote;
pub use model::{RunOutcome, RunSnapshot, RunStatus};
pub use model::{RunOutcome, RunSnapshot, RunStatus, StepSnapshot, StepStatus};
use model::{outcome_from_state, snapshot_from_state};
pub use remote::OutOfBandCompletion;

Expand Down Expand Up @@ -171,7 +171,17 @@ impl<C: Clock> Engine<C> {
let journal = self.open_run(run_id)?;
let spec = journal.run_spec().context("read run spec")?;
let state = self.load_state(&journal, spec)?;
Ok(snapshot_from_state(&state))
let mut snapshot = snapshot_from_state(&state);
let registry_record = self.registry()?.lookup(run_id)?;
if let Some(record) = registry_record.filter(|record| record.status == "waiting_worker") {
for step in snapshot.steps.values_mut().filter(|step| {
step.step_type != relayflowd_core::StepType::Deterministic
&& step.state == model::StepStatus::Running
}) {
step.lease_deadline_ms = record.next_wake_at_ms;
}
}
Ok(snapshot)
}

pub fn journal_entries(
Expand Down
49 changes: 46 additions & 3 deletions kernel/relayflowd/src/engine/model.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
use std::collections::BTreeMap;

use relayflowd_core::{Budget, RunCompletionReason, RunState};
use relayflowd_core::{Budget, RunCompletionReason, RunState, StepState, StepType};
use serde::{Deserialize, Serialize};

#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
Expand All @@ -25,10 +25,31 @@ pub struct RunOutcome {
pub struct RunSnapshot {
pub run_id: String,
pub status: RunStatus,
pub steps: BTreeMap<String, String>,
pub steps: BTreeMap<String, StepSnapshot>,
pub budget: Budget,
}

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct StepSnapshot {
#[serde(rename = "type")]
pub step_type: StepType,
pub state: StepStatus,
#[serde(skip_serializing_if = "Option::is_none")]
pub lease_deadline_ms: Option<i64>,
}

#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum StepStatus {
Pending,
Runnable,
Running,
Backoff,
Waiting,
NeedsHuman,
Done,
}

pub(super) fn outcome_from_state(state: &RunState, reason: RunCompletionReason) -> RunOutcome {
RunOutcome {
run_id: state.run_id.clone(),
Expand Down Expand Up @@ -57,9 +78,31 @@ pub(super) fn snapshot_from_state(state: &RunState) -> RunSnapshot {
None => RunStatus::Running,
},
steps: state
.spec
.steps
.iter()
.map(|(id, runtime)| (id.clone(), format!("{:?}", runtime.state)))
.map(|step| {
let runtime = &state.steps[&step.id];
let (step_state, lease_deadline_ms) = match runtime.state {
StepState::Pending => (StepStatus::Pending, None),
StepState::Runnable => (StepStatus::Runnable, None),
StepState::Running {
lease_deadline_ms, ..
} => (StepStatus::Running, Some(lease_deadline_ms)),
StepState::Backoff { .. } => (StepStatus::Backoff, None),
StepState::Waiting { .. } => (StepStatus::Waiting, None),
StepState::NeedsHuman { .. } => (StepStatus::NeedsHuman, None),
StepState::Done { .. } => (StepStatus::Done, None),
};
(
step.id.clone(),
StepSnapshot {
step_type: step.step_type(),
state: step_state,
lease_deadline_ms,
},
)
})
.collect(),
budget: state.budget.clone(),
}
Expand Down
5 changes: 4 additions & 1 deletion kernel/relayflowd/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,4 +4,7 @@ pub mod exec_det;
pub mod server;
pub mod worker;

pub use engine::{DriveOptions, Engine, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus};
pub use engine::{
DriveOptions, Engine, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus, StepSnapshot,
StepStatus,
};
12 changes: 12 additions & 0 deletions kernel/relayflowd/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,18 @@ fn handle_request(
}
"run.resume" => {
let params: RunIdParams = decode_params(request.params)?;
let registry = relayflowd_journal::Registry::open(data_dir.join("relayflowd.sqlite3"))
.map_err(|error| internal_error(error.into()))?;
if registry
.lookup(&params.run_id)
.map_err(|error| internal_error(error.into()))?
.is_none()
Comment on lines +170 to +173

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Preserve resume when the rebuildable registry is absent

When relayflowd.sqlite3 is lost or rebuilt while a valid per-run journal remains, this check returns run_not_found before Engine::resume_filtered can execute its existing “repair missing run registry entry” path. That makes an authoritative journal impossible to resume solely because its documented non-authoritative index is missing; validate the run journal and rebuild the registry row instead of treating an absent index entry as proof that the run does not exist.

AGENTS.md reference: AGENTS.md:L3-L5

Useful? React with 👍 / 👎.

{
return Err((
"run_not_found",
format!("run {} does not exist", params.run_id),
));
}
// Per-run serialization: the load-state -> next_actions -> append
// sequence must be atomic, or two concurrent resumes both see a
// step Runnable and double-dispatch the same attempt.
Expand Down
22 changes: 22 additions & 0 deletions kernel/relayflowd/src/server/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,28 @@ fn run_start_fails_closed_on_an_unknown_verification_key() {
assert_eq!(response.error.unwrap().code, "invalid_spec");
}

#[test]
fn run_resume_asks_the_registry_instead_of_treating_an_orphan_file_as_a_run() {
let directory = tempdir().unwrap();
let data_dir = directory.path();
let runs = data_dir.join("runs");
std::fs::create_dir_all(&runs).unwrap();
std::fs::File::create(runs.join("orphan.sqlite3")).unwrap();
let hub = Arc::new(ProtocolHub::default());
let (writer, _peer) = shared_writer();

let response = request(
data_dir,
&hub,
1,
&writer,
r#"{"id":"resume","verb":"run.resume","params":{"run_id":"orphan"}}"#,
);

let error = response.error.expect("orphan file must be refused");
assert_eq!(error.code, "run_not_found");
}

/// 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