Skip to main content

max / makenotwork

9.2 KB · 270 lines History Blame Raw
1 //! `pom serve`: starts the API/dashboard server and every background check task.
2
3 use tokio::task::JoinHandle;
4 use tracing::info;
5
6 use pom::alerts::Alerter;
7 use pom::config::Config;
8 use pom::db;
9 use pom::error::Result;
10 use pom::peer;
11
12 use super::tasks;
13 use super::tasks::CheckInterval;
14
15 pub(crate) async fn cmd_serve(pool: &sqlx::SqlitePool, config: &Config) -> Result<()> {
16 let default_interval = config.serve.interval_secs;
17 let prune_days = config.serve.prune_days;
18 let listen_addr = config.serve.listen.clone();
19
20 // Cancellation token for graceful shutdown
21 let token = tokio_util::sync::CancellationToken::new();
22
23 let instance_id = peer::load_or_create_instance_id(config.instance.id.as_deref())?;
24 let instance_name = config.instance_name();
25 let instance_info = peer::InstanceInfo {
26 id: instance_id.clone(),
27 name: instance_name.clone(),
28 version: env!("CARGO_PKG_VERSION").to_string(),
29 targets: config.target_names(),
30 started_at: chrono::Utc::now().to_rfc3339(),
31 };
32
33 let alerter = match config.alerts.as_ref() {
34 Some(alert_config) => {
35 info!("Alerts enabled (to: {})", alert_config.to);
36 match Alerter::new(alert_config.clone(), pool.clone(), instance_name.clone()) {
37 Ok(a) => Some(a),
38 Err(e) => {
39 tracing::error!("alerts: client build failed, alerting disabled: {e}");
40 None
41 }
42 }
43 }
44 None => None,
45 };
46
47 info!("Instance: {instance_name} (id={instance_id})");
48 info!("Starting serve mode (default interval: {default_interval}s, prune: {prune_days}d)");
49
50 let mesh = peer::new_mesh_state(instance_info, &config.peers);
51
52 // Load known peer identities from DB
53 {
54 let mut state = mesh.write().await;
55 for (name, peer) in &mut state.peers {
56 if let Ok(Some(known_id)) = db::get_peer_identity(pool, name).await {
57 peer.known_id = Some(known_id);
58 }
59 }
60 }
61
62 let mut handles: Vec<JoinHandle<()>> = Vec::new();
63
64 handles.extend(tasks::spawn_health_tasks(
65 config,
66 pool,
67 &token,
68 alerter.as_ref(),
69 ));
70 handles.extend(tasks::spawn_tls_tasks(
71 config,
72 pool,
73 &token,
74 alerter.as_ref(),
75 ));
76 handles.extend(tasks::spawn_route_tasks(
77 config,
78 pool,
79 &token,
80 alerter.as_ref(),
81 ));
82 handles.extend(tasks::spawn_dns_tasks(
83 config,
84 pool,
85 &token,
86 alerter.as_ref(),
87 ));
88 handles.extend(tasks::spawn_whois_tasks(
89 config,
90 pool,
91 &token,
92 alerter.as_ref(),
93 ));
94 handles.extend(tasks::spawn_cors_tasks(
95 config,
96 pool,
97 &token,
98 alerter.as_ref(),
99 ));
100 handles.extend(tasks::spawn_backup_tasks(
101 config,
102 pool,
103 &token,
104 alerter.as_ref(),
105 ));
106 handles.extend(tasks::spawn_scan_pipeline_tasks(
107 config,
108 pool,
109 &token,
110 alerter.as_ref(),
111 ));
112 handles.extend(tasks::spawn_systemd_tasks(
113 config,
114 pool,
115 &token,
116 alerter.as_ref(),
117 ));
118 // No alerter: the fleet readout has nothing to alert on. See the task module.
119 handles.extend(tasks::spawn_synckit_fleet_tasks(config, pool, &token));
120 handles.push(tasks::spawn_prune_task(pool, prune_days, &token));
121
122 // Spawn peer heartbeat tasks
123 if !config.peers.is_empty() {
124 let heartbeat_secs = config.serve.peer_heartbeat_secs;
125 info!(
126 "Peer mesh: {} peers, heartbeat every {heartbeat_secs}s",
127 config.peers.len()
128 );
129 let hb_handles = peer::spawn_heartbeat_tasks(
130 mesh.clone(),
131 pool.clone(),
132 heartbeat_secs,
133 alerter.clone(),
134 token.clone(),
135 )
136 .await;
137 handles.extend(hb_handles);
138 }
139
140 if let Some(handle) =
141 tasks::spawn_meta_alert_task(config, pool, default_interval, &token, alerter.as_ref())
142 {
143 handles.push(handle);
144 }
145
146 // Spawn the pending-alert retry drainer. When an alert's original send failed
147 // it was queued (the status transition fires only once, so without this the
148 // outage would go silent); re-deliver it here until it lands. This is the
149 // durable half of the fix for the transition-once alert-loss CRITICAL.
150 if let Some(retry_alerter) = alerter.clone() {
151 let pool = pool.clone();
152 let cancel = token.clone();
153 handles.push(tokio::spawn(async move {
154 const MAX_ATTEMPTS: i64 = 8; // ~ several hours of backoff before dead-letter
155 let mut ticks = CheckInterval::new(30, cancel);
156 while ticks.next().await {
157 let due = db::due_pending_alerts(&pool, 50).await.unwrap_or_default();
158 for p in due {
159 if retry_alerter.retry_pending(&p).await {
160 if let Err(e) = db::delete_pending_alert(&pool, p.id).await {
161 tracing::warn!(
162 alert_id = p.id,
163 alert_key = %p.alert_key,
164 error = %e,
165 "pending alert delete failed, alert will re-send"
166 );
167 }
168 info!("delivered previously-queued alert to {}", p.alert_key);
169 } else if p.attempts + 1 >= MAX_ATTEMPTS {
170 // Give up so the queue can't grow without bound; losing the
171 // row is logged loudly rather than silently.
172 tracing::error!(
173 "alert to {} undeliverable after {} attempts; dropping",
174 p.alert_key,
175 p.attempts + 1
176 );
177 if let Err(e) = db::delete_pending_alert(&pool, p.id).await {
178 tracing::warn!(
179 alert_id = p.id,
180 alert_key = %p.alert_key,
181 error = %e,
182 "pending alert delete failed, alert will re-send"
183 );
184 }
185 } else {
186 // Exponential backoff, capped at 1h.
187 let backoff =
188 std::cmp::min(30 * 2i64.pow((p.attempts.min(7)) as u32), 3600);
189 let next =
190 (chrono::Utc::now() + chrono::Duration::seconds(backoff)).to_rfc3339();
191 if let Err(e) = db::bump_pending_alert(&pool, p.id, &next).await {
192 tracing::warn!(
193 alert_id = p.id,
194 alert_key = %p.alert_key,
195 error = %e,
196 "pending alert backoff bump failed, row retries next tick"
197 );
198 }
199 }
200 }
201 }
202 }));
203 }
204
205 // Start HTTP API server
206 let api_app = pom::api::router(pool.clone(), config.clone(), Some(mesh.clone()));
207 let api_listener = tokio::net::TcpListener::bind(&listen_addr).await?;
208 info!("API server listening on {listen_addr}");
209 let api_cancel = token.clone();
210 handles.push(tokio::spawn(async move {
211 if let Err(e) = axum::serve(
212 api_listener,
213 api_app.into_make_service_with_connect_info::<std::net::SocketAddr>(),
214 )
215 .with_graceful_shutdown(async move { api_cancel.cancelled().await })
216 .await
217 {
218 tracing::error!("API server error: {e}");
219 }
220 }));
221
222 // Wait for shutdown signal or unexpected task exit
223 let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
224
225 let mut watchdog_interval = tokio::time::interval(std::time::Duration::from_mins(1));
226 watchdog_interval.tick().await; // consume immediate first tick
227
228 loop {
229 tokio::select! {
230 _ = tokio::signal::ctrl_c() => {
231 info!("Received SIGINT, shutting down");
232 break;
233 }
234 _ = sigterm.recv() => {
235 info!("Received SIGTERM, shutting down");
236 break;
237 }
238 _ = watchdog_interval.tick() => {
239 for (i, handle) in handles.iter().enumerate() {
240 if handle.is_finished() {
241 tracing::error!("Background task {i} exited unexpectedly (possible panic)");
242 }
243 }
244 }
245 }
246 }
247
248 // Graceful shutdown: cancel all tasks, then wait with grace period
249 token.cancel();
250 info!("Waiting for tasks to finish (5s grace period)...");
251
252 let shutdown = async {
253 for handle in handles {
254 if let Err(e) = handle.await {
255 tracing::error!("Task shutdown error: {e}");
256 }
257 }
258 };
259
260 if tokio::time::timeout(std::time::Duration::from_secs(5), shutdown)
261 .await
262 .is_err()
263 {
264 tracing::warn!("Grace period elapsed, some tasks may not have finished cleanly");
265 }
266
267 info!("Shutdown complete");
268 Ok(())
269 }
270