use crate::serve::config::{AuthMode, ServeConfig};
use crate::serve::history::RunHistory;
use crate::serve::logs::LogHub;
use crate::serve::registry::Registry;
use metrics_exporter_prometheus::PrometheusHandle;
use serde_json::Value;
use std::path::PathBuf;
use std::sync::{Arc, RwLock};
use std::time::Duration;
use tokio::sync::Semaphore;
use tokio_util::sync::CancellationToken;
#[derive(Clone)]
pub struct ServerState {
inner: Arc<Inner>,
}
struct Inner {
auth: AuthMode,
prometheus: Option<PrometheusHandle>,
shutdown: CancellationToken,
registry: Registry,
semaphore: Arc<Semaphore>,
history: Arc<dyn RunHistory>,
log_hub: LogHub,
default_base: RwLock<Option<Value>>,
default_config_path: Option<PathBuf>,
idempotency_retention: Duration,
probe_timeout: Duration,
local_output_retention_days: u32,
local_output_in_flight_grace: Duration,
preview: crate::serve::preview::PreviewConfig,
callback_allow_hosts: Vec<String>,
cluster: crate::serve::cluster::ClusterHandle,
#[cfg(feature = "triggers")]
triggers: crate::serve::triggers::health::TriggersHandle,
#[cfg(feature = "templates-sync")]
templates_sync: RwLock<Option<Arc<crate::templates::sync::SyncFile>>>,
#[cfg(feature = "policy")]
policy: RwLock<Option<Arc<faucet_core::PolicySpec>>>,
#[cfg(feature = "tenants")]
tenants: RwLock<Arc<crate::serve::tenants::TenantsRuntime>>,
require_approval: Vec<crate::serve::changes::ChangeKind>,
approval_expiry: Duration,
}
impl ServerState {
#[allow(clippy::too_many_arguments)]
pub fn new(
config: &ServeConfig,
prometheus: Option<PrometheusHandle>,
shutdown: CancellationToken,
history: Arc<dyn RunHistory>,
log_hub: LogHub,
default_base: Option<Value>,
#[cfg(feature = "triggers")] triggers: crate::serve::triggers::health::TriggersHandle,
) -> Self {
Self {
inner: Arc::new(Inner {
auth: config.auth.clone(),
prometheus,
shutdown,
registry: Registry::new(config.max_queued_runs),
semaphore: Arc::new(Semaphore::new(config.max_concurrent_runs)),
history,
log_hub,
default_base: RwLock::new(default_base),
default_config_path: config.default_config_path.clone(),
idempotency_retention: config.idempotency_retention,
probe_timeout: config.probe_timeout,
local_output_retention_days: config.local_output_retention_days,
local_output_in_flight_grace: config.local_output_in_flight_grace,
preview: config.preview,
callback_allow_hosts: config.callback_allow_hosts.clone(),
cluster: crate::serve::cluster::ClusterHandle::from_config(config),
#[cfg(feature = "triggers")]
triggers,
#[cfg(feature = "templates-sync")]
templates_sync: RwLock::new(None),
#[cfg(feature = "policy")]
policy: RwLock::new(None),
#[cfg(feature = "tenants")]
tenants: RwLock::new(Arc::new(Default::default())),
require_approval: config.require_approval.clone(),
approval_expiry: config.approval_expiry,
}),
}
}
pub fn requires_approval(&self, kind: crate::serve::changes::ChangeKind) -> bool {
self.inner.require_approval.contains(&kind)
}
pub fn require_approval(&self) -> &[crate::serve::changes::ChangeKind] {
&self.inner.require_approval
}
pub fn approvals(&self) -> crate::serve::changes::ApprovalPolicy {
match &self.inner.auth {
AuthMode::Rbac(cfg) => cfg.approvals().clone(),
_ => crate::serve::changes::ApprovalPolicy::permissive(),
}
}
pub fn approval_expiry(&self) -> Duration {
match self.approvals().expire_secs {
Some(secs) => Duration::from_secs(secs),
None => self.inner.approval_expiry,
}
}
#[cfg(feature = "tenants")]
pub fn set_tenants(&self, rt: crate::serve::tenants::TenantsRuntime) {
*self
.inner
.tenants
.write()
.unwrap_or_else(|e| e.into_inner()) = Arc::new(rt);
}
#[cfg(feature = "tenants")]
pub fn tenants(&self) -> Arc<crate::serve::tenants::TenantsRuntime> {
self.inner
.tenants
.read()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
#[cfg(feature = "policy")]
pub fn set_policy(&self, policy: Arc<faucet_core::PolicySpec>) {
*self.inner.policy.write().unwrap_or_else(|e| e.into_inner()) = Some(policy);
}
#[cfg(feature = "policy")]
pub fn policy(&self) -> Option<Arc<faucet_core::PolicySpec>> {
self.inner
.policy
.read()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
#[cfg(feature = "templates-sync")]
pub fn set_templates_sync(&self, file: Arc<crate::templates::sync::SyncFile>) {
*self
.inner
.templates_sync
.write()
.unwrap_or_else(|e| e.into_inner()) = Some(file);
}
#[cfg(feature = "templates-sync")]
pub fn templates_sync(&self) -> Option<Arc<crate::templates::sync::SyncFile>> {
self.inner
.templates_sync
.read()
.unwrap_or_else(|e| e.into_inner())
.clone()
}
pub fn auth_token(&self) -> Option<&str> {
match &self.inner.auth {
AuthMode::Token(t) => Some(t),
AuthMode::Rbac(_) | AuthMode::None => None,
}
}
pub fn auth_mode(&self) -> &AuthMode {
&self.inner.auth
}
pub fn render_metrics(&self) -> Option<String> {
self.inner.prometheus.as_ref().map(|h| h.render())
}
pub fn shutdown_token(&self) -> CancellationToken {
self.inner.shutdown.clone()
}
pub fn registry(&self) -> &Registry {
&self.inner.registry
}
pub fn semaphore(&self) -> Arc<Semaphore> {
Arc::clone(&self.inner.semaphore)
}
pub fn local_output_retention_days(&self) -> u32 {
self.inner.local_output_retention_days
}
pub fn local_output_in_flight_grace(&self) -> Duration {
self.inner.local_output_in_flight_grace
}
pub fn preview(&self) -> &crate::serve::preview::PreviewConfig {
&self.inner.preview
}
pub fn history(&self) -> Arc<dyn RunHistory> {
Arc::clone(&self.inner.history)
}
pub fn log_hub(&self) -> &LogHub {
&self.inner.log_hub
}
pub fn default_base(&self) -> Option<Value> {
self.inner.default_base.read().unwrap().clone()
}
pub fn default_config_path(&self) -> Option<&PathBuf> {
self.inner.default_config_path.as_ref()
}
pub fn set_default_base(&self, base: Option<Value>) {
*self.inner.default_base.write().unwrap() = base;
}
pub fn idempotency_retention(&self) -> Duration {
self.inner.idempotency_retention
}
pub fn probe_timeout(&self) -> Duration {
self.inner.probe_timeout
}
pub fn callback_allow_hosts(&self) -> &[String] {
&self.inner.callback_allow_hosts
}
pub fn cluster(&self) -> &crate::serve::cluster::ClusterHandle {
&self.inner.cluster
}
#[cfg(feature = "triggers")]
pub fn triggers(&self) -> &crate::serve::triggers::health::TriggersHandle {
&self.inner.triggers
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::serve::config::HistoryBackendSpec;
use crate::serve::history::memory::MemoryHistory;
fn cfg(auth: AuthMode) -> ServeConfig {
ServeConfig {
listen: "127.0.0.1:0".parse().unwrap(),
log_format: crate::cli::LogFormat::Text,
auth,
max_concurrent_runs: 4,
max_queued_runs: 32,
default_config_path: None,
history: HistoryBackendSpec::Memory,
cors_origins: vec![],
body_limit_bytes: 1_048_576,
shutdown_grace: Duration::from_secs(60),
retain_terminal_runs: Duration::from_secs(60),
idempotency_retention: Duration::from_secs(60),
log_retention: Duration::from_secs(0),
log_max_lines_per_run: 100_000,
local_output_retention_days: 7,
local_output_in_flight_grace: Duration::from_secs(60),
preview: crate::serve::preview::PreviewConfig::default(),
lease_ttl: Duration::from_secs(30),
probe_timeout: Duration::from_secs(10),
env_file: None,
no_env_file: false,
log_level: "info".into(),
ui_enabled: true,
cluster: crate::serve::cluster::ClusterConfig::disabled(),
triggers_path: None,
templates_sync_path: None,
policy_path: None,
callback_allow_hosts: Vec::new(),
require_approval: Vec::new(),
approval_expiry: std::time::Duration::from_secs(86_400),
vault: None,
connect_providers_path: None,
}
}
fn state(auth: AuthMode) -> ServerState {
use crate::serve::logs::LogHub;
let history = Arc::new(MemoryHistory::new(Duration::from_secs(60))) as Arc<dyn RunHistory>;
ServerState::new(
&cfg(auth),
None,
CancellationToken::new(),
history,
LogHub::new(),
None,
#[cfg(feature = "triggers")]
crate::serve::triggers::health::TriggersHandle::empty(),
)
}
#[test]
fn auth_token_reflects_mode() {
assert_eq!(state(AuthMode::Token("x".into())).auth_token(), Some("x"));
assert_eq!(state(AuthMode::None).auth_token(), None);
}
#[test]
fn render_metrics_none_without_handle() {
assert!(state(AuthMode::None).render_metrics().is_none());
}
#[test]
fn default_base_swaps_atomically() {
let s = state(AuthMode::None);
assert!(s.default_base().is_none());
assert!(s.default_config_path().is_none());
s.set_default_base(Some(serde_json::json!({"version": 1})));
assert_eq!(s.default_base(), Some(serde_json::json!({"version": 1})));
s.set_default_base(None);
assert!(s.default_base().is_none());
}
#[tokio::test]
async fn reload_handler_noop_without_default_config() {
let s = state(AuthMode::None);
let actor = crate::serve::rbac::AuthContext {
principal: "admin".into(),
role: crate::serve::rbac::Role::Admin,
source_ip: None,
tenant: None,
};
let axum::Json(body) = crate::serve::handlers::reload::reload(
axum::extract::State(s),
axum::extract::Extension(actor),
)
.await
.expect("reload ok");
assert_eq!(body["reloaded"], serde_json::json!(false));
}
#[test]
fn registry_starts_empty() {
let s = state(AuthMode::None);
assert_eq!(s.registry().queued(), 0);
assert_eq!(s.registry().in_flight(), 0);
}
#[test]
fn approval_settings_come_from_the_config() {
let mut cfg = crate::serve::test_support::test_config();
cfg.require_approval = vec![crate::serve::changes::ChangeKind::Run];
cfg.approval_expiry = Duration::from_secs(90);
let st = crate::serve::test_support::state_from(&cfg);
assert_eq!(
st.require_approval(),
&[crate::serve::changes::ChangeKind::Run]
);
assert!(st.requires_approval(crate::serve::changes::ChangeKind::Run));
assert_eq!(st.approval_expiry(), Duration::from_secs(90));
}
}