From bdbfcd9d81927121e1eee956e103d6a6176eeabf Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Tue, 14 Apr 2026 21:35:55 -0400 Subject: [PATCH] refactor(server): simplify AppState construction Replace the internal positional AppState builder with an AppStateConfig and route both production and test setup through the new config-backed path. Preserve the in-process test helper behavior while fixing the ignored max_concurrent_runs argument with a regression test. --- lib/crates/fabro-server/src/serve.rs | 18 +- lib/crates/fabro-server/src/server.rs | 162 ++++++++++-------- .../fabro-server/tests/it/api/system.rs | 73 +++++++- lib/crates/fabro-server/tests/it/helpers.rs | 11 +- 4 files changed, 182 insertions(+), 82 deletions(-) diff --git a/lib/crates/fabro-server/src/serve.rs b/lib/crates/fabro-server/src/serve.rs index 904e97057..ffb90d604 100644 --- a/lib/crates/fabro-server/src/serve.rs +++ b/lib/crates/fabro-server/src/serve.rs @@ -27,7 +27,7 @@ use crate::bind::{self, Bind, BindRequest}; use crate::github_webhooks::WebhookManager; use crate::jwt_auth::resolve_auth_mode_with_lookup; use crate::server::{ - RouterOptions, build_app_state_with_path, build_router_with_options, + AppStateConfig, RouterOptions, build_app_state, build_router_with_options, reconcile_incomplete_runs_on_startup, shutdown_active_workers, spawn_scheduler, }; use crate::server_secrets::ServerSecrets; @@ -326,17 +326,17 @@ where build_artifact_object_store(&resolved_server_settings)?; let artifact_store = fabro_store::ArtifactStore::new(artifact_object_store, artifact_prefix); let env_lookup: EnvLookup = Arc::new(|name| std::env::var(name).ok()); - let state = build_app_state_with_path( - Arc::clone(&shared_settings), - None, + let state = build_app_state(AppStateConfig { + settings: Arc::clone(&shared_settings), + registry_factory_override: None, max_concurrent_runs, store, artifact_store, - &vault_path, - &server_env_path, - true, - &env_lookup, - )?; + vault_path, + server_env_path, + local_daemon_mode: true, + env_lookup, + })?; let reconciled = reconcile_incomplete_runs_on_startup(&state).await?; if reconciled > 0 { info!( diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 6a79492b9..324ecd68f 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -116,7 +116,7 @@ use crate::server_secrets::{ }; use crate::{demo, diagnostics, run_manifest, settings_view, static_files, web_auth}; -type EnvLookup = Arc Option + Send + Sync>; +pub(crate) type EnvLookup = Arc Option + Send + Sync>; pub fn default_page_limit() -> u32 { 20 @@ -326,7 +326,8 @@ struct BillingAccumulator { by_model: HashMap, } -type RegistryFactoryOverride = dyn Fn(Arc) -> HandlerRegistry + Send + Sync; +pub(crate) type RegistryFactoryOverride = + dyn Fn(Arc) -> HandlerRegistry + Send + Sync; #[derive(Clone)] enum RunAnswerTransport { @@ -548,6 +549,18 @@ pub struct AppState { slack_started: AtomicBool, } +pub(crate) struct AppStateConfig { + pub(crate) settings: Arc>, + pub(crate) registry_factory_override: Option>, + pub(crate) max_concurrent_runs: usize, + pub(crate) store: Arc, + pub(crate) artifact_store: ArtifactStore, + pub(crate) vault_path: PathBuf, + pub(crate) server_env_path: PathBuf, + pub(crate) local_daemon_mode: bool, + pub(crate) env_lookup: EnvLookup, +} + fn nonzero_i64(value: i64) -> Option { (value != 0).then_some(value) } @@ -2147,8 +2160,9 @@ pub fn create_app_state() -> Arc { pub fn create_app_state_with_registry_factory( registry_factory_override: impl Fn(Arc) -> HandlerRegistry + Send + Sync + 'static, ) -> Arc { - create_app_state_with_settings_and_registry_factory( + create_app_state_with_options_and_registry_factory( SettingsLayer::default(), + 5, registry_factory_override, ) } @@ -2158,22 +2172,23 @@ pub fn create_app_state_with_settings_and_registry_factory( settings: SettingsLayer, registry_factory_override: impl Fn(Arc) -> HandlerRegistry + Send + Sync + 'static, ) -> Arc { - let (store, artifact_store) = test_store_bundle(); + create_app_state_with_options_and_registry_factory(settings, 5, registry_factory_override) +} + +#[doc(hidden)] +pub fn create_app_state_with_options_and_registry_factory( + settings: SettingsLayer, + max_concurrent_runs: usize, + registry_factory_override: impl Fn(Arc) -> HandlerRegistry + Send + Sync + 'static, +) -> Arc { let env_lookup = default_env_lookup(); - let secrets_path = test_secret_store_path(); - let server_env_path = secrets_path.with_file_name("server.env"); - build_app_state_with_path( + let mut config = default_test_app_state_config( Arc::new(RwLock::new(settings)), - Some(Box::new(registry_factory_override)), - 5, - store, - artifact_store, - &secrets_path, - &server_env_path, - false, - &env_lookup, - ) - .expect("test app state should build") + max_concurrent_runs, + env_lookup, + ); + config.registry_factory_override = Some(Box::new(registry_factory_override)); + build_app_state(config).expect("test app state should build") } /// Create an `AppState` with the given settings and concurrency limit. @@ -2181,15 +2196,13 @@ pub fn create_app_state_with_options( settings: SettingsLayer, max_concurrent_runs: usize, ) -> Arc { - let (store, artifact_store) = test_store_bundle(); let env_lookup = default_env_lookup(); - create_app_state_with_store_and_env_lookup( + build_app_state(default_test_app_state_config( Arc::new(RwLock::new(settings)), max_concurrent_runs, - store, - artifact_store, - &env_lookup, - ) + env_lookup, + )) + .expect("test app state should build") } #[doc(hidden)] @@ -2200,13 +2213,14 @@ pub fn create_app_state_with_env_lookup( ) -> Arc { let (store, artifact_store) = test_store_bundle(); let env_lookup: EnvLookup = Arc::new(env_lookup); - create_app_state_with_store_and_env_lookup( + let mut config = default_test_app_state_config( Arc::new(RwLock::new(settings)), max_concurrent_runs, - store, - artifact_store, - &env_lookup, - ) + env_lookup, + ); + config.store = store; + config.artifact_store = artifact_store; + build_app_state(config).expect("test app state should build") } #[cfg(test)] @@ -2215,8 +2229,8 @@ pub(crate) fn create_test_app_state_with_session_key( session_secret: Option<&str>, local_daemon_mode: bool, ) -> Arc { - let secrets_path = test_secret_store_path(); - let server_env_path = secrets_path + let vault_path = test_secret_store_path(); + let server_env_path = vault_path .parent() .expect("test secrets path should have parent") .join("server.env"); @@ -2229,17 +2243,17 @@ pub(crate) fn create_test_app_state_with_session_key( } let (store, artifact_store) = test_store_bundle(); let env_lookup = default_env_lookup(); - build_app_state_with_path( - Arc::new(RwLock::new(settings)), - None, - 5, + build_app_state(AppStateConfig { + settings: Arc::new(RwLock::new(settings)), + registry_factory_override: None, + max_concurrent_runs: 5, store, artifact_store, - &secrets_path, - &server_env_path, + vault_path, + server_env_path, local_daemon_mode, - &env_lookup, - ) + env_lookup, + }) .expect("test app state should build") } @@ -2254,6 +2268,27 @@ fn test_store_bundle() -> (Arc, ArtifactStore) { (store, artifact_store) } +fn default_test_app_state_config( + settings: Arc>, + max_concurrent_runs: usize, + env_lookup: EnvLookup, +) -> AppStateConfig { + let (store, artifact_store) = test_store_bundle(); + let vault_path = test_secret_store_path(); + let server_env_path = vault_path.with_file_name("server.env"); + AppStateConfig { + settings, + registry_factory_override: None, + max_concurrent_runs, + store, + artifact_store, + vault_path, + server_env_path, + local_daemon_mode: false, + env_lookup, + } +} + pub fn create_app_state_with_store( settings: Arc>, max_concurrent_runs: usize, @@ -2277,44 +2312,37 @@ fn create_app_state_with_store_and_env_lookup( artifact_store: ArtifactStore, env_lookup: &EnvLookup, ) -> Arc { - let secrets_path = test_secret_store_path(); - let server_env_path = secrets_path.with_file_name("server.env"); - build_app_state_with_path( - settings, - None, - max_concurrent_runs, - store, - artifact_store, - &secrets_path, - &server_env_path, - false, - env_lookup, - ) - .expect("test app state should build") + let mut config = + default_test_app_state_config(settings, max_concurrent_runs, Arc::clone(env_lookup)); + config.store = store; + config.artifact_store = artifact_store; + build_app_state(config).expect("test app state should build") } fn default_env_lookup() -> EnvLookup { Arc::new(|name| std::env::var(name).ok()) } -pub(crate) fn build_app_state_with_path( - settings: Arc>, - registry_factory_override: Option>, - max_concurrent_runs: usize, - store: Arc, - artifact_store: ArtifactStore, - vault_path: &std::path::Path, - server_env_path: &std::path::Path, - local_daemon_mode: bool, - env_lookup: &EnvLookup, -) -> anyhow::Result> { - let vault = Arc::new(AsyncRwLock::new(Vault::load(vault_path.to_path_buf())?)); - let server_secrets = ServerSecrets::with_env_lookup(server_env_path.to_path_buf(), { - let env_lookup = Arc::clone(env_lookup); +pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result> { + let AppStateConfig { + settings, + registry_factory_override, + max_concurrent_runs, + store, + artifact_store, + vault_path, + server_env_path, + local_daemon_mode, + env_lookup, + } = config; + + let vault = Arc::new(AsyncRwLock::new(Vault::load(vault_path)?)); + let server_secrets = ServerSecrets::with_env_lookup(server_env_path, { + let env_lookup = Arc::clone(&env_lookup); move |name| env_lookup(name) })?; let provider_credentials = ProviderCredentials::with_env_lookup(Arc::clone(&vault), { - let env_lookup = Arc::clone(env_lookup); + let env_lookup = Arc::clone(&env_lookup); move |name| env_lookup(name) }); let (global_event_tx, _) = broadcast::channel(4096); diff --git a/lib/crates/fabro-server/tests/it/api/system.rs b/lib/crates/fabro-server/tests/it/api/system.rs index 72e1ebf0e..8f315ddfa 100644 --- a/lib/crates/fabro-server/tests/it/api/system.rs +++ b/lib/crates/fabro-server/tests/it/api/system.rs @@ -14,10 +14,27 @@ use tokio::time::timeout; use tower::ServiceExt; use crate::helpers::{ - MINIMAL_DOT, api, body_json, minimal_manifest_json_with_dry_run, test_app_state_with_options, - test_app_with_scheduler, test_settings, wait_for_run_status, + MINIMAL_DOT, POLL_ATTEMPTS, POLL_INTERVAL, api, body_json, minimal_manifest_json, + minimal_manifest_json_with_dry_run, test_app_state_with_options, test_app_with_scheduler, + test_settings, wait_for_run_status, }; +const HUMAN_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 +}"#; + fn temp_storage_settings() -> (tempfile::TempDir, SettingsLayer, PathBuf) { let temp = tempdir().expect("tempdir should create"); let mut settings = test_settings(); @@ -51,6 +68,34 @@ async fn start_run(app: &axum::Router, run_id: &str) { assert_eq!(response.status(), StatusCode::OK); } +async fn wait_for_question(app: &axum::Router, run_id: &str) -> serde_json::Value { + for _ in 0..POLL_ATTEMPTS { + let request = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/questions"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(request).await.unwrap(); + let body = body_json(response.into_body()).await; + if let Some(question) = body["data"].as_array().and_then(|items| items.first()) { + return question.clone(); + } + tokio::time::sleep(POLL_INTERVAL).await; + } + panic!("question should have appeared for {run_id}"); +} + +async fn load_questions(app: &axum::Router, run_id: &str) -> serde_json::Value { + let request = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}/questions"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(request).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + body_json(response.into_body()).await +} + #[tokio::test] async fn get_system_info_returns_runtime_fields() { let (_temp, settings, expected_storage_dir) = temp_storage_settings(); @@ -79,6 +124,30 @@ async fn get_system_info_returns_runtime_fields() { assert!(body["uptime_secs"].as_i64().is_some()); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn test_app_state_with_options_respects_max_concurrent_runs() { + let app = test_app_with_scheduler(test_app_state_with_options(test_settings(), 1)); + + let first_run = create_run(&app, minimal_manifest_json(HUMAN_GATE_DOT)).await; + let second_run = create_run(&app, minimal_manifest_json(HUMAN_GATE_DOT)).await; + + start_run(&app, &first_run).await; + start_run(&app, &second_run).await; + + let question = wait_for_question(&app, &first_run).await; + assert_eq!(question["stage"], "gate"); + + tokio::time::sleep(POLL_INTERVAL * 5).await; + + let second_questions = load_questions(&app, &second_run).await; + assert!( + second_questions["data"] + .as_array() + .is_some_and(|items| items.is_empty()), + "second run should still be queued while the first waits at the human gate: {second_questions}" + ); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn get_system_disk_usage_returns_summary_and_verbose_rows() { let (_temp, settings, storage_dir) = temp_storage_settings(); diff --git a/lib/crates/fabro-server/tests/it/helpers.rs b/lib/crates/fabro-server/tests/it/helpers.rs index 6ca6c0653..0f26d20ff 100644 --- a/lib/crates/fabro-server/tests/it/helpers.rs +++ b/lib/crates/fabro-server/tests/it/helpers.rs @@ -6,7 +6,7 @@ use axum::http::{Request, StatusCode}; use fabro_server::jwt_auth::AuthMode; use fabro_server::server::{ AppState, build_router, create_app_state, create_app_state_with_env_lookup, - create_app_state_with_settings_and_registry_factory, spawn_scheduler, + create_app_state_with_options_and_registry_factory, spawn_scheduler, }; use fabro_types::settings::SettingsLayer; use fabro_types::settings::run::{LocalSandboxLayer, RunLayer, RunSandboxLayer, WorktreeMode}; @@ -31,10 +31,13 @@ pub(crate) fn test_app_state_with_options( settings: SettingsLayer, max_concurrent_runs: usize, ) -> Arc { - let _ = max_concurrent_runs; - create_app_state_with_settings_and_registry_factory(settings, |interviewer| { + create_app_state_with_options_and_registry_factory( + settings, + max_concurrent_runs, + |interviewer| { fabro_workflow::handler::default_registry(interviewer, || None) - }) + }, + ) } pub(crate) fn test_settings() -> SettingsLayer {