mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-03 02:24:33 +00:00
Drive the Petri projector from the server
The server holds one projector over its database and signals it after
each committed worker append, after each committed platform record
(through the run summary store's hook), at worker exit, and over every
Petri run at startup after the restart reconcile. A run executing in
the server process under the test override appends through the
projector's observing store, so it is signalled the same way. The
scenario tests read GET /runs/{id}/state after the view settles: the
hello prompt stage with its response, the command stage with its
output, and a two-branch parallel bundle whose branches are grouped
under the fork with the fork's results.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
e2bf05c0f0
commit
7b466ad9f3
5 changed files with 162 additions and 4 deletions
|
|
@ -838,6 +838,11 @@ where
|
|||
"Reconciled stale in-flight runs on startup"
|
||||
);
|
||||
}
|
||||
state
|
||||
.petri_projector
|
||||
.startup_pass()
|
||||
.await
|
||||
.context("catching Petri projections up at startup")?;
|
||||
spawn_scheduler(Arc::clone(&state));
|
||||
spawn_automation_scheduler(Arc::clone(&state));
|
||||
let pull_request_creation_supervisor =
|
||||
|
|
|
|||
|
|
@ -61,6 +61,7 @@ use fabro_llm::credentials::CredentialProvider;
|
|||
use fabro_llm::lithos_catalog::Catalog;
|
||||
use fabro_llm::{ClientOptions, FabroClient};
|
||||
use fabro_mcp_store::McpServerStore;
|
||||
use fabro_petri::projector::Projector;
|
||||
use fabro_redact::redact_jsonl_line;
|
||||
use fabro_sandbox::details::sandbox_details;
|
||||
use fabro_sandbox::driver::{DaytonaCredentials, ProviderAccess, ProviderConnectOptions};
|
||||
|
|
@ -1116,6 +1117,8 @@ pub struct AppState {
|
|||
pub(crate) worker_runtime: Arc<dyn WorkerRuntime>,
|
||||
/// The Petri runs held open for workers over the API.
|
||||
pub(crate) petri_runs: PetriRuns,
|
||||
/// The projector of Petri runs: signalled after each committed record.
|
||||
pub(crate) petri_projector: Arc<Projector>,
|
||||
scheduler_notify: Notify,
|
||||
automation_scheduler_notify: Notify,
|
||||
pull_request_scheduler_notify: Notify,
|
||||
|
|
@ -1187,6 +1190,13 @@ impl AppState {
|
|||
self.petri_runs.store()
|
||||
}
|
||||
|
||||
/// The projector of Petri runs, so a test can wait for a run's view to
|
||||
/// settle before it reads it.
|
||||
#[cfg(any(test, feature = "test-support"))]
|
||||
pub fn test_petri_projector(&self) -> &Arc<Projector> {
|
||||
&self.petri_projector
|
||||
}
|
||||
|
||||
/// A worker token for `run_id` with the plain `run:worker` scope, as the
|
||||
/// server mints for the worker it launches.
|
||||
pub fn test_issue_worker_token(&self, run_id: &RunId) -> String {
|
||||
|
|
@ -2490,6 +2500,11 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
);
|
||||
let variables = Arc::new(VariableStore::new(db_pool.clone()));
|
||||
let petri_runs = PetriRuns::new(db_pool.clone());
|
||||
let petri_projector = Projector::new(db_pool.clone());
|
||||
{
|
||||
let projector = Arc::clone(&petri_projector);
|
||||
store.set_platform_record_hook(Arc::new(move |run_id| projector.signal(run_id)));
|
||||
}
|
||||
let session_records = Arc::new(RunSessionRecordStore::new(db_pool.clone()));
|
||||
let secret_store = Arc::new(SecretStore::new(db_pool));
|
||||
let vault = preloaded_vault;
|
||||
|
|
@ -2607,6 +2622,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
|
|||
worker_control_bus,
|
||||
worker_runtime,
|
||||
petri_runs,
|
||||
petri_projector,
|
||||
scheduler_notify: Notify::new(),
|
||||
automation_scheduler_notify: Notify::new(),
|
||||
pull_request_scheduler_notify: Notify::new(),
|
||||
|
|
@ -4538,6 +4554,7 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
|
|||
// The worker is gone: whatever Petri run handles it held open over the
|
||||
// API drop here, so its lease never outlives it.
|
||||
state.petri_runs.worker_exited(run_id);
|
||||
state.petri_projector.signal(run_id);
|
||||
append_worker_exit_failure(&run_store, run_id, &worker_exit).await;
|
||||
|
||||
let final_state = match run_store.state().await {
|
||||
|
|
|
|||
|
|
@ -145,7 +145,11 @@ async fn append_records(
|
|||
Err(err) => return store_error_response(id, &err),
|
||||
};
|
||||
match writer.append(&log, &records).await {
|
||||
Ok(()) => StatusCode::NO_CONTENT.into_response(),
|
||||
Ok(()) => {
|
||||
// The records are durable; the projection trails them from here.
|
||||
state.petri_projector.signal(id);
|
||||
StatusCode::NO_CONTENT.into_response()
|
||||
}
|
||||
Err(err) => store_error_response(id, &err),
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -341,7 +341,9 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
|
|||
run_id: run_id.to_string(),
|
||||
run_dir: run_dir.join("petri"),
|
||||
execution,
|
||||
store: Arc::new(SqliteRunStore::new(state.db_pool.clone())),
|
||||
store: state
|
||||
.petri_projector
|
||||
.observe_store(Arc::new(SqliteRunStore::new(state.db_pool.clone()))),
|
||||
runtime: runtime_spec(&state, &eligible, dry_run),
|
||||
provider: run_state.spec.settings.run.environment.provider.clone(),
|
||||
cancel,
|
||||
|
|
|
|||
|
|
@ -26,8 +26,8 @@ use std::sync::Arc;
|
|||
|
||||
use axum::body::Body;
|
||||
use axum::http::{Request, StatusCode};
|
||||
use fabro_petri::SqliteRunStore;
|
||||
use fabro_petri::engine::{self, RunStatus};
|
||||
use fabro_petri::{SqliteRunStore, projector};
|
||||
use fabro_server::server::AppState;
|
||||
use fabro_server::test_support::{
|
||||
TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool,
|
||||
|
|
@ -35,7 +35,7 @@ use fabro_server::test_support::{
|
|||
};
|
||||
use fabro_static::EnvVars;
|
||||
use fabro_test::{TwinScenario, TwinScenarios, twin_openai};
|
||||
use fabro_types::{WorkflowPath, WorkflowVersion};
|
||||
use fabro_types::{RunId, WorkflowPath, WorkflowVersion};
|
||||
use tower::ServiceExt;
|
||||
|
||||
use crate::helpers::{
|
||||
|
|
@ -87,6 +87,23 @@ const UNKNOWN_MODEL_DOT: &str = r#"digraph Bad {
|
|||
start -> work -> exit
|
||||
}"#;
|
||||
|
||||
/// Two command branches joined by a fan-in.
|
||||
const PARALLEL_DOT: &str = r#"digraph Parallel {
|
||||
graph [goal="Run two branches"]
|
||||
start [shape=Mdiamond]
|
||||
exit [shape=Msquare]
|
||||
fork [shape=component]
|
||||
a [shape=parallelogram, script="echo a"]
|
||||
b [shape=parallelogram, script="echo b"]
|
||||
merge [shape=tripleoctagon]
|
||||
start -> fork
|
||||
fork -> a
|
||||
fork -> b
|
||||
a -> merge
|
||||
b -> merge
|
||||
merge -> exit
|
||||
}"#;
|
||||
|
||||
const PLAIN_SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
|
||||
const PETRI_SETTINGS: &str =
|
||||
"_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\nengine = \"petri\"\n";
|
||||
|
|
@ -168,6 +185,37 @@ async fn petri_outcome(state: &AppState, run_id: &str) -> engine::RunOutcome {
|
|||
.expect("the run's Petri record inspects")
|
||||
}
|
||||
|
||||
/// The run's projected state once its projector settled.
|
||||
async fn settled_state(state: &AppState, app: &axum::Router, run_id: &str) -> serde_json::Value {
|
||||
let id: RunId = run_id.parse().expect("the run id parses");
|
||||
state.test_petri_projector().settle(id).await;
|
||||
let req = Request::builder()
|
||||
.method("GET")
|
||||
.uri(api(&format!("/runs/{run_id}/state")))
|
||||
.body(Body::empty())
|
||||
.expect("state request should build");
|
||||
let response = app
|
||||
.clone()
|
||||
.oneshot(req)
|
||||
.await
|
||||
.expect("state request routes");
|
||||
response_json(
|
||||
response,
|
||||
StatusCode::OK,
|
||||
format!("GET /api/v1/runs/{run_id}/state"),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// How many items the run's projected stream holds.
|
||||
async fn petri_stream_len(state: &AppState, run_id: &str) -> usize {
|
||||
let id: RunId = run_id.parse().expect("the run id parses");
|
||||
projector::stored_stream(&test_app_db_pool(state), id)
|
||||
.await
|
||||
.expect("the stream reads")
|
||||
.len()
|
||||
}
|
||||
|
||||
async fn run_engine(app: &axum::Router, run_id: &str) -> serde_json::Value {
|
||||
let req = Request::builder()
|
||||
.method("GET")
|
||||
|
|
@ -260,6 +308,31 @@ async fn the_hello_bundle_runs_on_petri_when_the_version_names_the_engine() {
|
|||
let outcome = petri_outcome(&state, &run_id).await;
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
let projection = settled_state(&state, &app, &run_id).await;
|
||||
assert_eq!(projection["status"]["kind"], "succeeded", "{projection}");
|
||||
assert_eq!(
|
||||
projection["conclusion"]["status"], "succeeded",
|
||||
"{projection}"
|
||||
);
|
||||
let stages = projection["stages"]
|
||||
.as_object()
|
||||
.expect("the state carries its stages");
|
||||
let prompt = stages
|
||||
.values()
|
||||
.find(|stage| stage["handler"] == "prompt")
|
||||
.unwrap_or_else(|| panic!("the hello prompt stage is projected: {projection}"));
|
||||
assert_eq!(prompt["state"], "succeeded", "{prompt}");
|
||||
assert!(
|
||||
prompt["response"]
|
||||
.as_str()
|
||||
.is_some_and(|response| response.contains("A haiku, added.")),
|
||||
"the prompt's response is projected: {prompt}"
|
||||
);
|
||||
assert_eq!(
|
||||
run["usage"]["tokens"]["input"].as_u64().is_some(),
|
||||
true,
|
||||
"{run}"
|
||||
);
|
||||
let logs = twin.request_logs(&namespace).await;
|
||||
let requests = logs["requests"]
|
||||
.as_array()
|
||||
|
|
@ -302,6 +375,63 @@ async fn a_command_bundle_runs_on_petri_under_the_server_setting() {
|
|||
let outcome = petri_outcome(&state, &run_id).await;
|
||||
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
|
||||
assert!(outcome.complete, "{:?}", outcome.incomplete);
|
||||
let projection = settled_state(&state, &app, &run_id).await;
|
||||
let say = &projection["stages"]["say@1"];
|
||||
assert_eq!(say["state"], "succeeded", "{projection}");
|
||||
assert_eq!(say["handler"], "command", "{say}");
|
||||
assert!(
|
||||
say["output"]
|
||||
.as_str()
|
||||
.is_some_and(|output| output.contains("hello from petri")),
|
||||
"{say}"
|
||||
);
|
||||
let stream = petri_stream_len(&state, &run_id).await;
|
||||
assert!(stream > 0, "the run's stream holds its events");
|
||||
}
|
||||
|
||||
/// A parallel bundle with two command branches runs on Petri through the
|
||||
/// server: each branch is a child execution, projected as a stage grouped
|
||||
/// under the fork, and the fork carries the branch results.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn a_parallel_bundle_projects_its_branches_through_the_server() {
|
||||
if host_plugin().is_none() {
|
||||
return;
|
||||
}
|
||||
let workspace = tempfile::tempdir().expect("workspace tempdir");
|
||||
let settings = settings_from_toml(
|
||||
"_version = 1\n\n[run.environment]\nid = \"local\"\n\n[server.execution]\nengine = \
|
||||
\"petri\"\n",
|
||||
);
|
||||
let state = test_app_state_with_options(settings, 5);
|
||||
let app = test_app_with_scheduler(Arc::clone(&state));
|
||||
|
||||
let version_id = register_version(&app, &[
|
||||
("workflow.fabro", PARALLEL_DOT),
|
||||
("workflow.toml", PLAIN_SETTINGS),
|
||||
])
|
||||
.await;
|
||||
let run_id =
|
||||
create_and_start_run_from_intent(&app, intent(&version_id, workspace.path())).await;
|
||||
|
||||
let status = wait_for_run_status(&app, &run_id, &["succeeded", "failed"]).await;
|
||||
let run = run_json(&app, &run_id).await;
|
||||
assert_eq!(status, "succeeded", "run: {run}");
|
||||
let projection = settled_state(&state, &app, &run_id).await;
|
||||
for branch in ["a@1", "b@1"] {
|
||||
let stage = &projection["stages"][branch];
|
||||
assert_eq!(stage["state"], "succeeded", "{branch}: {projection}");
|
||||
assert_eq!(stage["parallel_branch_id"]["group"], "fork@1", "{stage}");
|
||||
}
|
||||
let fork = &projection["stages"]["fork@1"];
|
||||
assert_eq!(
|
||||
fork["parallel_results"].as_array().map(Vec::len),
|
||||
Some(2),
|
||||
"{fork}"
|
||||
);
|
||||
assert_eq!(
|
||||
projection["conclusion"]["status"], "succeeded",
|
||||
"{projection}"
|
||||
);
|
||||
}
|
||||
|
||||
/// A version that names no engine on a server whose setting is the default
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue