Skip to main content

faucet_cli/serve/
state.rs

1//! Shared, cheaply-cloneable server state handed to every handler via
2//! `axum::extract::State`. Holds auth, the Prometheus render handle, the
3//! server-wide shutdown token, the run registry, the execution semaphore, the
4//! run-history backend, and the `--default-config` merge base.
5
6use crate::serve::config::{AuthMode, ServeConfig};
7use crate::serve::history::RunHistory;
8use crate::serve::logs::LogHub;
9use crate::serve::registry::Registry;
10use metrics_exporter_prometheus::PrometheusHandle;
11use serde_json::Value;
12use std::path::PathBuf;
13use std::sync::{Arc, RwLock};
14use std::time::Duration;
15use tokio::sync::Semaphore;
16use tokio_util::sync::CancellationToken;
17
18#[derive(Clone)]
19pub struct ServerState {
20    inner: Arc<Inner>,
21}
22
23struct Inner {
24    auth: AuthMode,
25    prometheus: Option<PrometheusHandle>,
26    shutdown: CancellationToken,
27    registry: Registry,
28    semaphore: Arc<Semaphore>,
29    history: Arc<dyn RunHistory>,
30    log_hub: LogHub,
31    /// The `--default-config` merge base, hot-reloadable via `POST /v1/reload`
32    /// (#198). An `RwLock` so a reload can swap it while runs read it.
33    default_base: RwLock<Option<Value>>,
34    /// Path the default-config was loaded from, so a reload can re-read it.
35    default_config_path: Option<PathBuf>,
36    idempotency_retention: Duration,
37    probe_timeout: Duration,
38    /// Allowlist of hosts a per-run completion callback may target (#481).
39    callback_allow_hosts: Vec<String>,
40    cluster: crate::serve::cluster::ClusterHandle,
41    #[cfg(feature = "triggers")]
42    triggers: crate::serve::triggers::health::TriggersHandle,
43}
44
45impl ServerState {
46    #[allow(clippy::too_many_arguments)]
47    pub fn new(
48        config: &ServeConfig,
49        prometheus: Option<PrometheusHandle>,
50        shutdown: CancellationToken,
51        history: Arc<dyn RunHistory>,
52        log_hub: LogHub,
53        default_base: Option<Value>,
54        #[cfg(feature = "triggers")] triggers: crate::serve::triggers::health::TriggersHandle,
55    ) -> Self {
56        Self {
57            inner: Arc::new(Inner {
58                auth: config.auth.clone(),
59                prometheus,
60                shutdown,
61                registry: Registry::new(config.max_queued_runs),
62                semaphore: Arc::new(Semaphore::new(config.max_concurrent_runs)),
63                history,
64                log_hub,
65                default_base: RwLock::new(default_base),
66                default_config_path: config.default_config_path.clone(),
67                idempotency_retention: config.idempotency_retention,
68                probe_timeout: config.probe_timeout,
69                callback_allow_hosts: config.callback_allow_hosts.clone(),
70                cluster: crate::serve::cluster::ClusterHandle::from_config(config),
71                #[cfg(feature = "triggers")]
72                triggers,
73            }),
74        }
75    }
76
77    pub fn auth_token(&self) -> Option<&str> {
78        match &self.inner.auth {
79            AuthMode::Token(t) => Some(t),
80            AuthMode::Rbac(_) | AuthMode::None => None,
81        }
82    }
83
84    /// The configured authentication mode (bearer resolution + RBAC).
85    pub fn auth_mode(&self) -> &AuthMode {
86        &self.inner.auth
87    }
88
89    pub fn render_metrics(&self) -> Option<String> {
90        self.inner.prometheus.as_ref().map(|h| h.render())
91    }
92
93    pub fn shutdown_token(&self) -> CancellationToken {
94        self.inner.shutdown.clone()
95    }
96
97    pub fn registry(&self) -> &Registry {
98        &self.inner.registry
99    }
100
101    pub fn semaphore(&self) -> Arc<Semaphore> {
102        Arc::clone(&self.inner.semaphore)
103    }
104
105    pub fn history(&self) -> Arc<dyn RunHistory> {
106        Arc::clone(&self.inner.history)
107    }
108
109    /// The per-run log buffer registry shared with the tracing [`LogHub`] layer.
110    pub fn log_hub(&self) -> &LogHub {
111        &self.inner.log_hub
112    }
113
114    /// A snapshot of the `--default-config` merge base (cloned under the read
115    /// lock, so a concurrent hot reload can't tear it).
116    pub fn default_base(&self) -> Option<Value> {
117        self.inner.default_base.read().unwrap().clone()
118    }
119
120    /// Path the default-config was loaded from (`None` when `--default-config`
121    /// was not passed).
122    pub fn default_config_path(&self) -> Option<&PathBuf> {
123        self.inner.default_config_path.as_ref()
124    }
125
126    /// Atomically swap the `--default-config` merge base (hot reload, #198).
127    pub fn set_default_base(&self, base: Option<Value>) {
128        *self.inner.default_base.write().unwrap() = base;
129    }
130
131    pub fn idempotency_retention(&self) -> Duration {
132        self.inner.idempotency_retention
133    }
134
135    pub fn probe_timeout(&self) -> Duration {
136        self.inner.probe_timeout
137    }
138
139    /// Hosts a per-run completion callback may target. Empty = any host except
140    /// link-local / cloud-metadata addresses (#481).
141    pub fn callback_allow_hosts(&self) -> &[String] {
142        &self.inner.callback_allow_hosts
143    }
144
145    pub fn cluster(&self) -> &crate::serve::cluster::ClusterHandle {
146        &self.inner.cluster
147    }
148
149    #[cfg(feature = "triggers")]
150    pub fn triggers(&self) -> &crate::serve::triggers::health::TriggersHandle {
151        &self.inner.triggers
152    }
153}
154
155#[cfg(test)]
156mod tests {
157    use super::*;
158    use crate::serve::config::HistoryBackendSpec;
159    use crate::serve::history::memory::MemoryHistory;
160
161    fn cfg(auth: AuthMode) -> ServeConfig {
162        ServeConfig {
163            listen: "127.0.0.1:0".parse().unwrap(),
164            auth,
165            max_concurrent_runs: 4,
166            max_queued_runs: 32,
167            default_config_path: None,
168            history: HistoryBackendSpec::Memory,
169            cors_origins: vec![],
170            body_limit_bytes: 1_048_576,
171            shutdown_grace: Duration::from_secs(60),
172            retain_terminal_runs: Duration::from_secs(60),
173            idempotency_retention: Duration::from_secs(60),
174            log_retention: Duration::from_secs(0),
175            log_max_lines_per_run: 100_000,
176            lease_ttl: Duration::from_secs(30),
177            probe_timeout: Duration::from_secs(10),
178            env_file: None,
179            no_env_file: false,
180            log_level: "info".into(),
181            ui_enabled: true,
182            cluster: crate::serve::cluster::ClusterConfig::disabled(),
183            triggers_path: None,
184            callback_allow_hosts: Vec::new(),
185        }
186    }
187
188    fn state(auth: AuthMode) -> ServerState {
189        use crate::serve::logs::LogHub;
190        let history = Arc::new(MemoryHistory::new(Duration::from_secs(60))) as Arc<dyn RunHistory>;
191        ServerState::new(
192            &cfg(auth),
193            None,
194            CancellationToken::new(),
195            history,
196            LogHub::new(),
197            None,
198            #[cfg(feature = "triggers")]
199            crate::serve::triggers::health::TriggersHandle::empty(),
200        )
201    }
202
203    #[test]
204    fn auth_token_reflects_mode() {
205        assert_eq!(state(AuthMode::Token("x".into())).auth_token(), Some("x"));
206        assert_eq!(state(AuthMode::None).auth_token(), None);
207    }
208
209    #[test]
210    fn render_metrics_none_without_handle() {
211        assert!(state(AuthMode::None).render_metrics().is_none());
212    }
213
214    #[test]
215    fn default_base_swaps_atomically() {
216        let s = state(AuthMode::None);
217        assert!(s.default_base().is_none());
218        assert!(s.default_config_path().is_none());
219        s.set_default_base(Some(serde_json::json!({"version": 1})));
220        assert_eq!(s.default_base(), Some(serde_json::json!({"version": 1})));
221        s.set_default_base(None);
222        assert!(s.default_base().is_none());
223    }
224
225    // The reload handler is a no-op (200 `reloaded:false`) when the server was
226    // started without `--default-config`.
227    #[tokio::test]
228    async fn reload_handler_noop_without_default_config() {
229        let s = state(AuthMode::None);
230        let actor = crate::serve::rbac::AuthContext {
231            principal: "admin".into(),
232            role: crate::serve::rbac::Role::Admin,
233            source_ip: None,
234        };
235        let axum::Json(body) = crate::serve::handlers::reload::reload(
236            axum::extract::State(s),
237            axum::extract::Extension(actor),
238        )
239        .await
240        .expect("reload ok");
241        assert_eq!(body["reloaded"], serde_json::json!(false));
242    }
243
244    #[test]
245    fn registry_starts_empty() {
246        let s = state(AuthMode::None);
247        assert_eq!(s.registry().queued(), 0);
248        assert_eq!(s.registry().in_flight(), 0);
249    }
250}