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 lease_ttl: Duration::from_secs(30),
175 probe_timeout: Duration::from_secs(10),
176 env_file: None,
177 no_env_file: false,
178 log_level: "info".into(),
179 ui_enabled: true,
180 cluster: crate::serve::cluster::ClusterConfig::disabled(),
181 triggers_path: None,
182 callback_allow_hosts: Vec::new(),
183 }
184 }
185
186 fn state(auth: AuthMode) -> ServerState {
187 use crate::serve::logs::LogHub;
188 let history = Arc::new(MemoryHistory::new(Duration::from_secs(60))) as Arc<dyn RunHistory>;
189 ServerState::new(
190 &cfg(auth),
191 None,
192 CancellationToken::new(),
193 history,
194 LogHub::new(),
195 None,
196 #[cfg(feature = "triggers")]
197 crate::serve::triggers::health::TriggersHandle::empty(),
198 )
199 }
200
201 #[test]
202 fn auth_token_reflects_mode() {
203 assert_eq!(state(AuthMode::Token("x".into())).auth_token(), Some("x"));
204 assert_eq!(state(AuthMode::None).auth_token(), None);
205 }
206
207 #[test]
208 fn render_metrics_none_without_handle() {
209 assert!(state(AuthMode::None).render_metrics().is_none());
210 }
211
212 #[test]
213 fn default_base_swaps_atomically() {
214 let s = state(AuthMode::None);
215 assert!(s.default_base().is_none());
216 assert!(s.default_config_path().is_none());
217 s.set_default_base(Some(serde_json::json!({"version": 1})));
218 assert_eq!(s.default_base(), Some(serde_json::json!({"version": 1})));
219 s.set_default_base(None);
220 assert!(s.default_base().is_none());
221 }
222
223 #[tokio::test]
226 async fn reload_handler_noop_without_default_config() {
227 let s = state(AuthMode::None);
228 let actor = crate::serve::rbac::AuthContext {
229 principal: "admin".into(),
230 role: crate::serve::rbac::Role::Admin,
231 source_ip: None,
232 };
233 let axum::Json(body) = crate::serve::handlers::reload::reload(
234 axum::extract::State(s),
235 axum::extract::Extension(actor),
236 )
237 .await
238 .expect("reload ok");
239 assert_eq!(body["reloaded"], serde_json::json!(false));
240 }
241
242 #[test]
243 fn registry_starts_empty() {
244 let s = state(AuthMode::None);
245 assert_eq!(s.registry().queued(), 0);
246 assert_eq!(s.registry().in_flight(), 0);
247 }
248}