mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
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>
352 lines
12 KiB
Rust
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}");
|
|
}
|