fabro/lib/apps/fabro-server/tests/it/scenario/lifecycle.rs
Bryan Helmkamp d90a5d9cbb
Delete fabro-core and the engine half of fabro-workflow
Every run executes on Petri, so the in-process legacy executor goes:
`fabro-core` and, in `fabro-workflow`, the handlers, lifecycle, pipeline
execution, routing, retry, conditions, node handlers, steering, agent
memory, artifacts, checkpoints, command log, and the `start`, `resume`,
`retry`, `fork`, `rewind` and `timeline` operations. The two are deleted
together because the engine half of `fabro-workflow` was the only user of
`fabro-core` and `fabro-core` the only runtime of that half; neither
compiles without the other.

Kept in `fabro-workflow`, narrowed: the parse/transform/validate/persist
pipeline and `create`, `archive`, `validate` (workflow definitions still
come from DOT and settings); the run tools (`run_tools`, moved from
`handler/llm/fabro_tools.rs`) for Ask Fabro, `fabro exec` and Petri's
host tools; the pull request pipeline (`pull_request`, moved from
`pipeline/`, for the step 0 port); Run Files' diff helpers in
`sandbox_git`; `git_identity`, `usage_rollup`, `run_status`,
`run_materialization`, `web_search` and `workflow_bundle`.

Server: `RegistryFactoryOverride` becomes `execute_in_process`;
`RunAnswerTransport::InProcess` carries only the interviewer; the
interrupt endpoint answers 501 `interrupt_unsupported` and every pair
endpoint 501 `pair_unsupported` (status lists none); rewind, fork, retry
and timeline handlers and routes are removed; the command log is served
from the stage output blob; usage rollups accumulate from the settled
projection after an in-process run as after a worker exit.

Ported while here:

- `materialize_admitted_run` materializes the goal and drops a disabled
  pull request block, as the legacy materializer did.
- A run whose admitted graph has an agent or prompt node is refused at
  create when no LLM provider is ready (`fabro.model.no_ready_provider`);
  a workflow of commands and gates needs no model and is admitted.
- The projection's question type falls back on the options, as the
  interview adapter does, so a gate with edge-label options answers as
  multiple choice.

Tests: the server scenarios (lifecycle, run completion, SSE, helpers)
run in process on Petri and assert Petri's stage labels and stream
names; the reconcile tests assert Petri's relaunch semantics; legacy
unit tests of the deleted executor are removed; three server unit tests
the removal took with it are restored; the pair fixtures go with the
pair feature. Petri test fixtures no longer name `[workflow] engine`.

Still red after this commit, all legacy consumers the next steps
delete or port: fabro-store's Slate/reducer fixtures and fabro-types
legacy JSON tests (step 4); server unit tests over legacy run events
(retry endpoints, list_run_events, artifacts, per-event pause/unpause,
run history activation, legacy sandbox fixtures) (steps 3-4); CLI tests
that parse legacy event envelopes, the legacy `events`/`attach`/`diff`/
`dump`/`inspect` snapshots, `run rewind`/`run fork`, the ACP and
git-identity workflow tests, and the runner tests that drive the legacy
worker by hand (steps 3-4); the web app's Petri fixtures still carry
`engine` (regenerate with `FABRO_CAPTURE_PETRI_FIXTURES` in step 4).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-18 10:44:41 -04:00

352 lines
12 KiB
Rust

use std::sync::Arc;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use fabro_server::server::spawn_scheduler;
use fabro_server::test_support::test_app_state_with_runtime_settings_in_process;
use tokio::time::sleep;
use tower::ServiceExt;
use crate::helpers::{
POLL_ATTEMPTS, POLL_INTERVAL, api, minimal_intent_json, response_json, response_status,
run_json, test_settings, wait_for_run_status,
};
async fn wait_for_question_id(app: &axum::Router, run_id: &str) -> String {
for _ in 0..POLL_ATTEMPTS {
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/questions")))
.body(Body::empty())
.expect("questions request should build");
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/questions"),
)
.await;
let arr = body["data"]
.as_array()
.expect("questions response should include a data array");
if let Some(question_id) = arr
.first()
.and_then(|item| item["id"].as_str())
.map(ToOwned::to_owned)
{
return question_id;
}
sleep(POLL_INTERVAL).await;
}
panic!("question should have appeared");
}
async fn wait_for_question(app: &axum::Router, run_id: &str) -> serde_json::Value {
for _ in 0..POLL_ATTEMPTS {
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/questions")))
.body(Body::empty())
.expect("questions request should build");
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/questions"),
)
.await;
let arr = body["data"]
.as_array()
.expect("questions response should include a data array");
if let Some(question) = arr.first() {
return question.clone();
}
sleep(POLL_INTERVAL).await;
}
panic!("question should have appeared");
}
async fn wait_for_run_state(
app: &axum::Router,
run_id: &str,
expected_status: &str,
expected_reason: &str,
) -> serde_json::Value {
for _ in 0..POLL_ATTEMPTS {
let body = run_json(app, run_id).await;
if body["lifecycle"]["status"]["kind"].as_str() == Some(expected_status)
&& body["lifecycle"]["status"]["reason"].as_str() == Some(expected_reason)
{
return body;
}
sleep(POLL_INTERVAL).await;
}
panic!("run {run_id} did not reach status={expected_status} reason={expected_reason}");
}
const GATE_DOT: &str = r#"digraph GateTest {
graph [goal="Test gate"]
start [shape=Mdiamond]
exit [shape=Msquare]
work [shape=box, prompt="Do work"]
gate [shape=hexagon, type="human", label="Approve?"]
done [shape=box, prompt="Finish"]
revise [shape=box, prompt="Revise"]
start -> work -> gate
gate -> done [label="[A] Approve"]
gate -> revise [label="[R] Revise"]
done -> exit
revise -> gate
}"#;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn full_http_lifecycle_approve_and_complete() {
let workspace = tempfile::tempdir().unwrap();
let settings = test_settings();
let state = test_app_state_with_runtime_settings_in_process(
settings.server_settings,
settings.manifest_run_defaults,
);
spawn_scheduler(Arc::clone(&state));
let app = fabro_server::test_support::build_test_router(Arc::clone(&state));
// 1. Create run
let req = Request::builder()
.method("POST")
.uri(api("/runs"))
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_string(&minimal_intent_json(&app, GATE_DOT, workspace.path()).await)
.unwrap(),
))
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(response, StatusCode::CREATED, "POST /api/v1/runs").await;
let run_id = body["id"].as_str().unwrap().to_string();
// 1b. Start the run
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/start")))
.body(Body::empty())
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
response_status(
response,
StatusCode::OK,
format!("POST /api/v1/runs/{run_id}/start"),
)
.await;
// 2. Poll for question to appear (run goes start -> work -> gate, then blocks)
let question = wait_for_question(&app, &run_id).await;
let question_id = question["id"].as_str().unwrap().to_string();
assert_eq!(question["stage"], "gate@1");
assert!(question["timeout_seconds"].is_null());
assert!(question["context_display"].is_null() || question["context_display"].is_string());
// 3. Submit answer selecting first option (Approve). Petri's id
// (`gate#3`) travels as one percent-encoded path segment.
let encoded_id =
percent_encoding::utf8_percent_encode(&question_id, percent_encoding::NON_ALPHANUMERIC)
.to_string();
let req = Request::builder()
.method("POST")
.uri(api(&format!(
"/runs/{run_id}/questions/{encoded_id}/answer"
)))
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_string(&serde_json::json!({
"kind": "selected",
"option_key": "A",
}))
.unwrap(),
))
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
response_status(
response,
StatusCode::NO_CONTENT,
format!("POST /api/v1/runs/{run_id}/questions/{question_id}/answer"),
)
.await;
// 4. Poll until the run reaches a terminal success or failure state.
let final_status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
assert_eq!(final_status, "succeeded");
// 5. Verify no pending questions
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}/questions")))
.body(Body::empty())
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}/questions"),
)
.await;
assert!(
body["data"].as_array().unwrap().is_empty(),
"no pending questions after completion"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn full_http_lifecycle_cancel() {
let workspace = tempfile::tempdir().unwrap();
let settings = test_settings();
let state = test_app_state_with_runtime_settings_in_process(
settings.server_settings,
settings.manifest_run_defaults,
);
spawn_scheduler(Arc::clone(&state));
let app = fabro_server::test_support::build_test_router(Arc::clone(&state));
// Create and start a run that will block at the human gate
let req = Request::builder()
.method("POST")
.uri(api("/runs"))
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_string(&minimal_intent_json(&app, GATE_DOT, workspace.path()).await)
.unwrap(),
))
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(response, StatusCode::CREATED, "POST /api/v1/runs").await;
let run_id = body["id"].as_str().unwrap().to_string();
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/start")))
.body(Body::empty())
.unwrap();
response_status(
app.clone().oneshot(req).await.unwrap(),
StatusCode::OK,
format!("POST /api/v1/runs/{run_id}/start"),
)
.await;
// Wait until the worker has reached the human gate so cancel exercises the
// live-running path rather than racing the in-memory queue transition.
let _question_id = wait_for_question_id(&app, &run_id).await;
// Cancel it
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/cancel")))
.body(Body::empty())
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(
response,
StatusCode::ACCEPTED,
format!("POST /api/v1/runs/{run_id}/cancel"),
)
.await;
// Both fields below are read from the same post-signal projection, so both
// race the worker the same way. `status.kind` is "blocked" while the worker
// still sits at the gate and "running" once it has been notified and
// resumed to process the cancel. What matters here is that cancel reached a
// live run rather than racing the in-memory queue transition, so this
// asserts "not queued" via the two live states. Durable convergence is
// asserted below.
let status_kind = &body["lifecycle"]["status"]["kind"];
assert!(
status_kind == "blocked" || status_kind == "running",
"expected status.kind to be \"blocked\" or \"running\", got {status_kind}"
);
// `pending_control` is computed from the store projection after the cancel
// event is appended AND the worker is signaled. The worker is sitting at a
// human gate; once notified it can emit a clearing event before this
// handler re-reads the projection, so the response can legitimately
// observe either the still-pending "cancel" or a null where the worker
// already consumed it. Durable convergence is asserted below.
let pending_control = &body["lifecycle"]["pending_control"];
assert!(
pending_control == "cancel" || pending_control.is_null(),
"expected pending_control to be \"cancel\" or null, got {pending_control}"
);
// Verify the durable store view converges to cancelled failure.
let body = wait_for_run_state(&app, &run_id, "failed", "cancelled").await;
assert_eq!(body["lifecycle"]["status"]["reason"], "cancelled");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn cancel_at_human_gate_persists_cancelled_terminal_event() {
let workspace = tempfile::tempdir().unwrap();
let settings = test_settings();
let state = test_app_state_with_runtime_settings_in_process(
settings.server_settings,
settings.manifest_run_defaults,
);
spawn_scheduler(Arc::clone(&state));
let app = fabro_server::test_support::build_test_router(Arc::clone(&state));
let req = Request::builder()
.method("POST")
.uri(api("/runs"))
.header("content-type", "application/json")
.body(Body::from(
serde_json::to_string(&minimal_intent_json(&app, GATE_DOT, workspace.path()).await)
.unwrap(),
))
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
let body = response_json(response, StatusCode::CREATED, "POST /api/v1/runs").await;
let run_id = body["id"].as_str().unwrap().to_string();
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/start")))
.body(Body::empty())
.unwrap();
response_status(
app.clone().oneshot(req).await.unwrap(),
StatusCode::OK,
format!("POST /api/v1/runs/{run_id}/start"),
)
.await;
let _question_id = wait_for_question_id(&app, &run_id).await;
let req = Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/cancel")))
.body(Body::empty())
.unwrap();
let response = app.clone().oneshot(req).await.unwrap();
response_status(
response,
StatusCode::ACCEPTED,
format!("POST /api/v1/runs/{run_id}/cancel"),
)
.await;
let status = wait_for_run_status(&app, &run_id, &["failed"]).await;
assert_eq!(status, "failed");
// The run's record says it was cancelled: Petri's finish, and the
// terminal lifecycle record Fabro wrote after it, both name the reason.
let req = Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}")))
.body(Body::empty())
.unwrap();
let response = app.oneshot(req).await.unwrap();
let body = response_json(
response,
StatusCode::OK,
format!("GET /api/v1/runs/{run_id}"),
)
.await;
assert_eq!(body["lifecycle"]["status"]["reason"], "cancelled", "{body}");
}