mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-05 02:41:45 +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>
341 lines
10 KiB
Rust
341 lines
10 KiB
Rust
use std::path::PathBuf;
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use axum::body::{Body, to_bytes};
|
|
use axum::http::{Request, StatusCode};
|
|
use fabro_config::{RunEnvironmentLayer, RunLayer, ServerSettingsBuilder};
|
|
use fabro_server::server::{AppState, spawn_scheduler};
|
|
use fabro_server::test_support::{
|
|
TestAppStateBuilder, build_test_router, llm_overlay_with_provider_base_url,
|
|
test_app_state as server_test_app_state, test_app_state_with_runtime_settings_and_env_lookup,
|
|
test_app_state_with_runtime_settings_and_options_in_process,
|
|
};
|
|
use fabro_test::{
|
|
assert_axum_status, assert_reqwest_status, expect_axum_json, expect_axum_status,
|
|
expect_axum_status_in, expect_axum_text,
|
|
};
|
|
use fabro_types::ServerSettings;
|
|
use tokio::time::sleep;
|
|
use tower::ServiceExt;
|
|
|
|
pub(crate) const MINIMAL_DOT: &str = r#"digraph Test {
|
|
graph [goal="Test"]
|
|
start [shape=Mdiamond]
|
|
exit [shape=Msquare]
|
|
start -> exit
|
|
}"#;
|
|
|
|
pub(crate) const POLL_INTERVAL: Duration = Duration::from_millis(10);
|
|
pub(crate) const POLL_ATTEMPTS: usize = 500;
|
|
|
|
#[derive(Clone)]
|
|
pub(crate) struct TestAppSettings {
|
|
pub server_settings: ServerSettings,
|
|
pub manifest_run_defaults: RunLayer,
|
|
}
|
|
|
|
impl Default for TestAppSettings {
|
|
fn default() -> Self {
|
|
settings_from_toml("_version = 1\n")
|
|
}
|
|
}
|
|
|
|
fn ensure_test_auth_methods(document: &mut toml::Table) {
|
|
let server = document
|
|
.entry("server")
|
|
.or_insert_with(|| toml::Value::Table(toml::Table::new()))
|
|
.as_table_mut()
|
|
.expect("[server] should stay a table in test fixtures");
|
|
let auth = server
|
|
.entry("auth")
|
|
.or_insert_with(|| toml::Value::Table(toml::Table::new()))
|
|
.as_table_mut()
|
|
.expect("[server.auth] should stay a table in test fixtures");
|
|
auth.entry("methods")
|
|
.or_insert_with(|| toml::Value::Array(vec![toml::Value::String("dev-token".to_string())]));
|
|
}
|
|
|
|
pub(crate) fn settings_from_toml(source: &str) -> TestAppSettings {
|
|
let mut document: toml::Table = source.parse().expect("test fixture should parse as TOML");
|
|
ensure_test_auth_methods(&mut document);
|
|
let manifest_run_defaults = document
|
|
.remove("run")
|
|
.map(toml::Value::try_into::<RunLayer>)
|
|
.transpose()
|
|
.expect("test run settings should parse")
|
|
.unwrap_or_default();
|
|
let server_settings = ServerSettingsBuilder::from_toml(
|
|
&toml::to_string(&document).expect("test fixture should serialize"),
|
|
)
|
|
.expect("test server settings should resolve");
|
|
TestAppSettings {
|
|
server_settings,
|
|
manifest_run_defaults,
|
|
}
|
|
}
|
|
|
|
pub(crate) fn test_app_state() -> Arc<AppState> {
|
|
server_test_app_state()
|
|
}
|
|
|
|
pub(crate) fn test_app_state_with_options(
|
|
settings: TestAppSettings,
|
|
max_concurrent_runs: usize,
|
|
) -> Arc<AppState> {
|
|
test_app_state_with_runtime_settings_and_options_in_process(
|
|
settings.server_settings,
|
|
settings.manifest_run_defaults,
|
|
max_concurrent_runs,
|
|
)
|
|
}
|
|
|
|
pub(crate) fn test_settings() -> TestAppSettings {
|
|
TestAppSettings {
|
|
manifest_run_defaults: RunLayer {
|
|
environment: Some(RunEnvironmentLayer {
|
|
id: Some("local".to_string()),
|
|
..RunEnvironmentLayer::default()
|
|
}),
|
|
..RunLayer::default()
|
|
},
|
|
..TestAppSettings::default()
|
|
}
|
|
}
|
|
|
|
pub(crate) fn test_app_with_scheduler(state: Arc<AppState>) -> axum::Router {
|
|
spawn_scheduler(Arc::clone(&state));
|
|
build_test_router(state)
|
|
}
|
|
|
|
pub(crate) fn test_app_with_no_providers() -> axum::Router {
|
|
let settings = test_settings();
|
|
let state = test_app_state_with_runtime_settings_and_env_lookup(
|
|
settings.server_settings,
|
|
settings.manifest_run_defaults,
|
|
5,
|
|
|_| None,
|
|
);
|
|
build_test_router(state)
|
|
}
|
|
|
|
pub(crate) fn test_app_with_mock_anthropic(mock_base_url: &str) -> axum::Router {
|
|
let settings = test_settings();
|
|
let state = TestAppStateBuilder::new()
|
|
.runtime_settings(settings.server_settings, settings.manifest_run_defaults)
|
|
.max_concurrent_runs(5)
|
|
.llm_overlay(llm_overlay_with_provider_base_url(
|
|
"anthropic",
|
|
mock_base_url,
|
|
))
|
|
.vault_entries([("ANTHROPIC_API_KEY", "test-key")])
|
|
.build();
|
|
build_test_router(state)
|
|
}
|
|
|
|
pub(crate) fn api(path: &str) -> String {
|
|
format!("/api/v1{path}")
|
|
}
|
|
|
|
pub(crate) fn repo_root() -> PathBuf {
|
|
std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
|
|
.ancestors()
|
|
.nth(3)
|
|
.expect("fabro-server crate should be nested under lib/apps/fabro-server")
|
|
.to_path_buf()
|
|
}
|
|
|
|
#[expect(
|
|
clippy::disallowed_methods,
|
|
reason = "test fixture reads tracked files synchronously"
|
|
)]
|
|
pub(crate) fn read_repo_file(relative_path: &str) -> String {
|
|
let path = repo_root().join(relative_path);
|
|
std::fs::read_to_string(&path)
|
|
.unwrap_or_else(|err| panic!("failed to read {}: {err}", path.display()))
|
|
}
|
|
|
|
pub(crate) async fn body_json(body: Body) -> serde_json::Value {
|
|
let bytes = to_bytes(body, usize::MAX)
|
|
.await
|
|
.expect("response body should fit in memory");
|
|
serde_json::from_slice(&bytes).expect("response body should be valid JSON")
|
|
}
|
|
|
|
pub(crate) async fn response_status(
|
|
response: axum::response::Response,
|
|
expected: StatusCode,
|
|
context: impl std::fmt::Display,
|
|
) {
|
|
assert_axum_status(response, expected, context).await;
|
|
}
|
|
|
|
pub(crate) async fn response_json(
|
|
response: axum::response::Response,
|
|
expected: StatusCode,
|
|
context: impl std::fmt::Display,
|
|
) -> serde_json::Value {
|
|
expect_axum_json(response, expected, context).await
|
|
}
|
|
|
|
pub(crate) async fn response_text(
|
|
response: axum::response::Response,
|
|
expected: StatusCode,
|
|
context: impl std::fmt::Display,
|
|
) -> String {
|
|
expect_axum_text(response, expected, context).await
|
|
}
|
|
|
|
pub(crate) async fn checked_response(
|
|
response: axum::response::Response,
|
|
expected: StatusCode,
|
|
context: impl std::fmt::Display,
|
|
) -> axum::response::Response {
|
|
expect_axum_status(response, expected, context).await
|
|
}
|
|
|
|
pub(crate) async fn checked_response_in(
|
|
response: axum::response::Response,
|
|
expected: &[StatusCode],
|
|
context: impl std::fmt::Display,
|
|
) -> axum::response::Response {
|
|
expect_axum_status_in(response, expected, context).await
|
|
}
|
|
|
|
pub(crate) async fn reqwest_status(
|
|
response: fabro_http::Response,
|
|
expected: StatusCode,
|
|
context: impl std::fmt::Display,
|
|
) {
|
|
assert_reqwest_status(response, expected, context).await;
|
|
}
|
|
|
|
pub(crate) async fn create_and_start_run_from_intent(
|
|
app: &axum::Router,
|
|
intent: serde_json::Value,
|
|
) -> String {
|
|
let req = Request::builder()
|
|
.method("POST")
|
|
.uri(api("/runs"))
|
|
.header("content-type", "application/json")
|
|
.body(Body::from(
|
|
serde_json::to_string(&intent).expect("intent fixture should serialize"),
|
|
))
|
|
.expect("create-run request should build");
|
|
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()
|
|
.expect("create-run response should include an id")
|
|
.to_string();
|
|
|
|
let req = Request::builder()
|
|
.method("POST")
|
|
.uri(api(&format!("/runs/{run_id}/start")))
|
|
.body(Body::empty())
|
|
.expect("start-run request should build");
|
|
response_status(
|
|
app.clone().oneshot(req).await.unwrap(),
|
|
StatusCode::OK,
|
|
format!("POST /api/v1/runs/{run_id}/start"),
|
|
)
|
|
.await;
|
|
|
|
run_id
|
|
}
|
|
|
|
pub(crate) fn minimal_manifest_json(dot_source: &str) -> serde_json::Value {
|
|
serde_json::json!({
|
|
"version": 1,
|
|
"cwd": "/tmp",
|
|
"target": {
|
|
"path": "workflow.fabro"
|
|
},
|
|
"workflows": {
|
|
"workflow.fabro": {
|
|
"source": dot_source,
|
|
"files": {}
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
pub(crate) async fn minimal_intent_json(
|
|
app: &axum::Router,
|
|
source: &str,
|
|
workspace: &std::path::Path,
|
|
) -> serde_json::Value {
|
|
let path = fabro_types::WorkflowPath::new("workflow.fabro")
|
|
.expect("workflow fixture path should be valid");
|
|
let version = fabro_types::WorkflowVersion::new(
|
|
path.clone(),
|
|
std::collections::BTreeMap::from([(path, source.to_string())]),
|
|
std::collections::BTreeMap::new(),
|
|
)
|
|
.expect("workflow fixture version should be valid");
|
|
let id = fabro_server::test_support::test_register_workflow_version(app, &version, None).await;
|
|
serde_json::json!({"workflow_version_id": id, "target": {"kind": "folder", "path": workspace}, "environment_id": "local", "args": {}})
|
|
}
|
|
|
|
pub(crate) async fn minimal_intent_json_with_dry_run(
|
|
app: &axum::Router,
|
|
source: &str,
|
|
workspace: &std::path::Path,
|
|
) -> serde_json::Value {
|
|
let mut intent = minimal_intent_json(app, source, workspace).await;
|
|
intent["args"]["dry_run"] = serde_json::json!(true);
|
|
intent
|
|
}
|
|
|
|
pub(crate) async fn run_json(app: &axum::Router, run_id: &str) -> serde_json::Value {
|
|
let req = Request::builder()
|
|
.method("GET")
|
|
.uri(api(&format!("/runs/{run_id}")))
|
|
.body(Body::empty())
|
|
.expect("run lookup request should build");
|
|
let response = app.clone().oneshot(req).await.unwrap();
|
|
response_json(
|
|
response,
|
|
StatusCode::OK,
|
|
format!("GET /api/v1/runs/{run_id}"),
|
|
)
|
|
.await
|
|
}
|
|
|
|
pub(crate) async fn wait_for_run_status(
|
|
app: &axum::Router,
|
|
run_id: &str,
|
|
expected: &[&str],
|
|
) -> String {
|
|
for _ in 0..POLL_ATTEMPTS {
|
|
let body = run_json(app, run_id).await;
|
|
let status = body["lifecycle"]["status"]["kind"]
|
|
.as_str()
|
|
.expect("run response should include a tagged status kind")
|
|
.to_string();
|
|
if expected.iter().any(|candidate| *candidate == status) {
|
|
return status;
|
|
}
|
|
sleep(POLL_INTERVAL).await;
|
|
}
|
|
panic!("run {run_id} did not reach any of {expected:?}");
|
|
}
|
|
|
|
pub(crate) async fn wait_for_run_status_not_in(
|
|
app: &axum::Router,
|
|
run_id: &str,
|
|
unexpected: &[&str],
|
|
) -> String {
|
|
for _ in 0..POLL_ATTEMPTS {
|
|
let body = run_json(app, run_id).await;
|
|
let status = body["lifecycle"]["status"]["kind"]
|
|
.as_str()
|
|
.expect("run response should include a tagged status kind")
|
|
.to_string();
|
|
if unexpected.iter().all(|candidate| *candidate != status) {
|
|
return status;
|
|
}
|
|
sleep(POLL_INTERVAL).await;
|
|
}
|
|
panic!("run {run_id} stayed in {unexpected:?}");
|
|
}
|