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.
This commit is contained in:
Bryan Helmkamp 2026-04-14 21:35:55 -04:00
parent bf9f4f9353
commit bdbfcd9d81
No known key found for this signature in database
4 changed files with 182 additions and 82 deletions

View file

@ -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!(

View file

@ -116,7 +116,7 @@ use crate::server_secrets::{
};
use crate::{demo, diagnostics, run_manifest, settings_view, static_files, web_auth};
type EnvLookup = Arc<dyn Fn(&str) -> Option<String> + Send + Sync>;
pub(crate) type EnvLookup = Arc<dyn Fn(&str) -> Option<String> + Send + Sync>;
pub fn default_page_limit() -> u32 {
20
@ -326,7 +326,8 @@ struct BillingAccumulator {
by_model: HashMap<String, ModelBillingTotals>,
}
type RegistryFactoryOverride = dyn Fn(Arc<dyn Interviewer>) -> HandlerRegistry + Send + Sync;
pub(crate) type RegistryFactoryOverride =
dyn Fn(Arc<dyn Interviewer>) -> 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<RwLock<SettingsLayer>>,
pub(crate) registry_factory_override: Option<Box<RegistryFactoryOverride>>,
pub(crate) max_concurrent_runs: usize,
pub(crate) store: Arc<Database>,
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<i64> {
(value != 0).then_some(value)
}
@ -2147,8 +2160,9 @@ pub fn create_app_state() -> Arc<AppState> {
pub fn create_app_state_with_registry_factory(
registry_factory_override: impl Fn(Arc<dyn Interviewer>) -> HandlerRegistry + Send + Sync + 'static,
) -> Arc<AppState> {
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<dyn Interviewer>) -> HandlerRegistry + Send + Sync + 'static,
) -> Arc<AppState> {
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<dyn Interviewer>) -> HandlerRegistry + Send + Sync + 'static,
) -> Arc<AppState> {
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<AppState> {
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<AppState> {
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<AppState> {
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<Database>, ArtifactStore) {
(store, artifact_store)
}
fn default_test_app_state_config(
settings: Arc<RwLock<SettingsLayer>>,
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<RwLock<SettingsLayer>>,
max_concurrent_runs: usize,
@ -2277,44 +2312,37 @@ fn create_app_state_with_store_and_env_lookup(
artifact_store: ArtifactStore,
env_lookup: &EnvLookup,
) -> Arc<AppState> {
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<RwLock<SettingsLayer>>,
registry_factory_override: Option<Box<RegistryFactoryOverride>>,
max_concurrent_runs: usize,
store: Arc<Database>,
artifact_store: ArtifactStore,
vault_path: &std::path::Path,
server_env_path: &std::path::Path,
local_daemon_mode: bool,
env_lookup: &EnvLookup,
) -> anyhow::Result<Arc<AppState>> {
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<Arc<AppState>> {
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);

View file

@ -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();

View file

@ -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<AppState> {
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 {