faucet_cli/serve/
state.rs1use 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 default_base: RwLock<Option<Value>>,
34 default_config_path: Option<PathBuf>,
36 idempotency_retention: Duration,
37 probe_timeout: Duration,
38 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 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 pub fn log_hub(&self) -> &LogHub {
111 &self.inner.log_hub
112 }
113
114 pub fn default_base(&self) -> Option<Value> {
117 self.inner.default_base.read().unwrap().clone()
118 }
119
120 pub fn default_config_path(&self) -> Option<&PathBuf> {
123 self.inner.default_config_path.as_ref()
124 }
125
126 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 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 #[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}