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::sync::Arc;
13use std::time::Duration;
14use tokio::sync::Semaphore;
15use tokio_util::sync::CancellationToken;
16
17#[derive(Clone)]
18pub struct ServerState {
19 inner: Arc<Inner>,
20}
21
22struct Inner {
23 auth: AuthMode,
24 prometheus: Option<PrometheusHandle>,
25 shutdown: CancellationToken,
26 registry: Registry,
27 semaphore: Arc<Semaphore>,
28 history: Arc<dyn RunHistory>,
29 log_hub: LogHub,
30 default_base: Option<Value>,
31 idempotency_retention: Duration,
32 probe_timeout: Duration,
33 cluster: crate::serve::cluster::ClusterHandle,
34 #[cfg(feature = "triggers")]
35 triggers: crate::serve::triggers::health::TriggersHandle,
36}
37
38impl ServerState {
39 #[allow(clippy::too_many_arguments)]
40 pub fn new(
41 config: &ServeConfig,
42 prometheus: Option<PrometheusHandle>,
43 shutdown: CancellationToken,
44 history: Arc<dyn RunHistory>,
45 log_hub: LogHub,
46 default_base: Option<Value>,
47 #[cfg(feature = "triggers")] triggers: crate::serve::triggers::health::TriggersHandle,
48 ) -> Self {
49 Self {
50 inner: Arc::new(Inner {
51 auth: config.auth.clone(),
52 prometheus,
53 shutdown,
54 registry: Registry::new(config.max_queued_runs),
55 semaphore: Arc::new(Semaphore::new(config.max_concurrent_runs)),
56 history,
57 log_hub,
58 default_base,
59 idempotency_retention: config.idempotency_retention,
60 probe_timeout: config.probe_timeout,
61 cluster: crate::serve::cluster::ClusterHandle::from_config(config),
62 #[cfg(feature = "triggers")]
63 triggers,
64 }),
65 }
66 }
67
68 pub fn auth_token(&self) -> Option<&str> {
69 match &self.inner.auth {
70 AuthMode::Token(t) => Some(t),
71 AuthMode::Rbac(_) | AuthMode::None => None,
72 }
73 }
74
75 pub fn auth_mode(&self) -> &AuthMode {
77 &self.inner.auth
78 }
79
80 pub fn render_metrics(&self) -> Option<String> {
81 self.inner.prometheus.as_ref().map(|h| h.render())
82 }
83
84 pub fn shutdown_token(&self) -> CancellationToken {
85 self.inner.shutdown.clone()
86 }
87
88 pub fn registry(&self) -> &Registry {
89 &self.inner.registry
90 }
91
92 pub fn semaphore(&self) -> Arc<Semaphore> {
93 Arc::clone(&self.inner.semaphore)
94 }
95
96 pub fn history(&self) -> Arc<dyn RunHistory> {
97 Arc::clone(&self.inner.history)
98 }
99
100 pub fn log_hub(&self) -> &LogHub {
102 &self.inner.log_hub
103 }
104
105 pub fn default_base(&self) -> Option<&Value> {
106 self.inner.default_base.as_ref()
107 }
108
109 pub fn idempotency_retention(&self) -> Duration {
110 self.inner.idempotency_retention
111 }
112
113 pub fn probe_timeout(&self) -> Duration {
114 self.inner.probe_timeout
115 }
116
117 pub fn cluster(&self) -> &crate::serve::cluster::ClusterHandle {
118 &self.inner.cluster
119 }
120
121 #[cfg(feature = "triggers")]
122 pub fn triggers(&self) -> &crate::serve::triggers::health::TriggersHandle {
123 &self.inner.triggers
124 }
125}
126
127#[cfg(test)]
128mod tests {
129 use super::*;
130 use crate::serve::config::HistoryBackendSpec;
131 use crate::serve::history::memory::MemoryHistory;
132
133 fn cfg(auth: AuthMode) -> ServeConfig {
134 ServeConfig {
135 listen: "127.0.0.1:0".parse().unwrap(),
136 auth,
137 max_concurrent_runs: 4,
138 max_queued_runs: 32,
139 default_config_path: None,
140 history: HistoryBackendSpec::Memory,
141 cors_origins: vec![],
142 body_limit_bytes: 1_048_576,
143 shutdown_grace: Duration::from_secs(60),
144 retain_terminal_runs: Duration::from_secs(60),
145 idempotency_retention: Duration::from_secs(60),
146 lease_ttl: Duration::from_secs(30),
147 probe_timeout: Duration::from_secs(10),
148 env_file: None,
149 no_env_file: false,
150 log_level: "info".into(),
151 ui_enabled: true,
152 cluster: crate::serve::cluster::ClusterConfig::disabled(),
153 triggers_path: None,
154 }
155 }
156
157 fn state(auth: AuthMode) -> ServerState {
158 use crate::serve::logs::LogHub;
159 let history = Arc::new(MemoryHistory::new(Duration::from_secs(60))) as Arc<dyn RunHistory>;
160 ServerState::new(
161 &cfg(auth),
162 None,
163 CancellationToken::new(),
164 history,
165 LogHub::new(),
166 None,
167 #[cfg(feature = "triggers")]
168 crate::serve::triggers::health::TriggersHandle::empty(),
169 )
170 }
171
172 #[test]
173 fn auth_token_reflects_mode() {
174 assert_eq!(state(AuthMode::Token("x".into())).auth_token(), Some("x"));
175 assert_eq!(state(AuthMode::None).auth_token(), None);
176 }
177
178 #[test]
179 fn render_metrics_none_without_handle() {
180 assert!(state(AuthMode::None).render_metrics().is_none());
181 }
182
183 #[test]
184 fn registry_starts_empty() {
185 let s = state(AuthMode::None);
186 assert_eq!(s.registry().queued(), 0);
187 assert_eq!(s.registry().in_flight(), 0);
188 }
189}