faucet_cli/serve/
cluster.rs1use crate::serve::config::ServeConfig;
7use crate::serve::state::ServerState;
8use std::sync::Arc;
9use std::sync::atomic::{AtomicUsize, Ordering};
10use std::time::Duration;
11use tokio::sync::Notify;
12use tokio_util::sync::CancellationToken;
13
14#[derive(Debug, Clone)]
16pub struct ClusterConfig {
17 pub enabled: bool,
18 pub poll: Duration,
20 pub max_attempts: u32,
22}
23
24impl ClusterConfig {
25 pub fn disabled() -> Self {
27 Self {
28 enabled: false,
29 poll: Duration::from_secs(2),
30 max_attempts: 3,
31 }
32 }
33}
34
35#[derive(Clone)]
40pub struct ClusterHandle {
41 inner: Arc<ClusterInner>,
42}
43
44struct ClusterInner {
45 cfg: ClusterConfig,
46 listen: String,
47 max_concurrent: u32,
48 started_at: chrono::DateTime<chrono::Utc>,
49 kick: Notify,
50 members: AtomicUsize,
51}
52
53impl ClusterHandle {
54 pub fn from_config(config: &ServeConfig) -> Self {
57 Self {
58 inner: Arc::new(ClusterInner {
59 cfg: config.cluster.clone(),
60 listen: config.listen.to_string(),
61 max_concurrent: config.max_concurrent_runs as u32,
62 started_at: chrono::Utc::now(),
63 kick: Notify::new(),
64 members: AtomicUsize::new(0),
65 }),
66 }
67 }
68
69 pub fn enabled(&self) -> bool {
70 self.inner.cfg.enabled
71 }
72 pub fn poll(&self) -> Duration {
73 self.inner.cfg.poll
74 }
75 pub fn max_attempts(&self) -> u32 {
76 self.inner.cfg.max_attempts
77 }
78 pub fn listen(&self) -> &str {
79 &self.inner.listen
80 }
81 pub fn max_concurrent(&self) -> u32 {
82 self.inner.max_concurrent
83 }
84 pub fn started_at(&self) -> chrono::DateTime<chrono::Utc> {
85 self.inner.started_at
86 }
87
88 pub fn kick(&self) {
90 self.inner.kick.notify_one();
91 }
92 pub async fn kicked(&self) {
94 self.inner.kick.notified().await;
95 }
96
97 pub fn members(&self) -> usize {
99 self.inner.members.load(Ordering::Acquire)
100 }
101 pub fn set_members(&self, n: usize) {
102 self.inner.members.store(n, Ordering::Release);
103 }
104}
105
106pub async fn claim_loop(state: ServerState, shutdown: CancellationToken) {
120 let handle = state.cluster().clone();
121 let mut tick = tokio::time::interval(handle.poll());
122 tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
123 loop {
124 tokio::select! {
125 biased;
126 _ = shutdown.cancelled() => break,
127 _ = tick.tick() => {}
128 _ = handle.kicked() => {}
129 }
130
131 match state.history().pending_cancellations().await {
133 Ok(ids) => {
134 for id in ids {
135 state.registry().cancel(&id);
136 }
137 }
138 Err(e) => tracing::warn!(error = %e, "cluster: pending_cancellations failed"),
139 }
140
141 let free = state.semaphore().available_permits();
143 if free == 0 {
144 continue;
145 }
146 match state.history().claim_pending(free).await {
147 Ok(claimed) => {
148 if !claimed.is_empty() {
149 crate::serve::metrics::record_runs_claimed(claimed.len());
150 for rec in claimed {
151 crate::serve::runner::resume_claimed_run(state.clone(), rec);
152 }
153 }
154 }
155 Err(e) => tracing::warn!(error = %e, "cluster: claim_pending failed"),
156 }
157 }
158}
159
160#[cfg(test)]
161mod tests {
162 use super::*;
163
164 #[test]
165 fn disabled_handle_reports_disabled() {
166 let cfg = ClusterConfig::disabled();
167 assert!(!cfg.enabled);
168 assert_eq!(cfg.max_attempts, 3);
169 }
170
171 #[tokio::test]
172 async fn kick_wakes_a_waiter() {
173 use crate::serve::config::{AuthMode, HistoryBackendSpec, ServeConfig};
176 let cfg = ServeConfig {
177 listen: "127.0.0.1:0".parse().unwrap(),
178 auth: AuthMode::None,
179 max_concurrent_runs: 4,
180 max_queued_runs: 4,
181 default_config_path: None,
182 history: HistoryBackendSpec::Memory,
183 cors_origins: vec![],
184 body_limit_bytes: 1_048_576,
185 shutdown_grace: std::time::Duration::from_secs(60),
186 retain_terminal_runs: std::time::Duration::from_secs(60),
187 idempotency_retention: std::time::Duration::from_secs(60),
188 lease_ttl: std::time::Duration::from_secs(30),
189 probe_timeout: std::time::Duration::from_secs(10),
190 env_file: None,
191 no_env_file: false,
192 log_level: "info".into(),
193 ui_enabled: true,
194 cluster: ClusterConfig::disabled(),
195 triggers_path: None,
196 };
197 let h = ClusterHandle::from_config(&cfg);
198 h.kick();
199 tokio::time::timeout(std::time::Duration::from_secs(1), h.kicked())
200 .await
201 .expect("kick must wake the waiter");
202 h.set_members(2);
203 assert_eq!(h.members(), 2);
204 }
205}