Skip to main content

max / makenotwork

4.2 KB · 117 lines History Blame Raw
1 //! Background route check task.
2
3 use tokio::task::JoinHandle;
4 use tracing::info;
5
6 use pom::alerts::Alerter;
7 use pom::checks::routes;
8 use pom::config::Config;
9 use pom::db;
10
11 use super::{CheckInterval, configured_targets};
12
13 pub(crate) fn spawn_route_tasks(
14 config: &Config,
15 pool: &sqlx::SqlitePool,
16 cancel: &tokio_util::sync::CancellationToken,
17 alerter: Option<&Alerter>,
18 ) -> Vec<JoinHandle<()>> {
19 let route_interval = config.serve.route_check_interval_secs;
20 let mut handles = Vec::new();
21
22 for (name, target_config) in configured_targets(config) {
23 if target_config.expected_routes.is_empty() {
24 continue;
25 }
26 let Some(ref health_config) = target_config.health else {
27 continue;
28 };
29 let Some(base_url) = routes::base_url_from_health_url(&health_config.url) else {
30 continue;
31 };
32 let route_paths = target_config.expected_routes.clone();
33 let timeout = std::time::Duration::from_secs(health_config.timeout_secs);
34 let label = target_config.label.clone();
35 let pool = pool.clone();
36 let alerter = alerter.cloned();
37 let n = route_paths.len();
38 let cancel = cancel.clone();
39
40 info!("{name}: route checks every {route_interval}s ({n} routes)");
41
42 handles.push(tokio::spawn(async move {
43 // Prune stale route data from DB (paths removed from config)
44 match db::prune_stale_routes(&pool, &name, &route_paths).await {
45 Ok(0) => {}
46 Ok(n) => info!("{name}: pruned {n} stale route check rows"),
47 Err(e) => tracing::error!("{name}: failed to prune stale routes: {e}"),
48 }
49
50 let mut ticks = CheckInterval::new(route_interval, cancel);
51 // reqwest 0.13's builder can fail and `unwrap_or_default()` panics on
52 // the same failure, silently killing this spawned route-check task. Log
53 // and exit visibly instead.
54 let client = match pom::tls::https_client_builder()
55 .redirect(reqwest::redirect::Policy::none())
56 .build()
57 {
58 Ok(c) => c,
59 Err(e) => {
60 tracing::error!("{name}: route-check client build failed: {e}");
61 return;
62 }
63 };
64 let mut prev_failed: std::collections::HashSet<String> =
65 std::collections::HashSet::new();
66
67 while ticks.next().await {
68 let results =
69 routes::check_routes(&client, &name, &base_url, &route_paths, timeout).await;
70
71 for result in &results {
72 if let Err(e) = db::insert_route_check(&pool, result).await {
73 tracing::error!(
74 "{}: failed to store route check for {}: {e}",
75 name,
76 result.path
77 );
78 }
79 }
80
81 let current_failed: std::collections::HashSet<String> = results
82 .iter()
83 .filter(|r| !r.ok)
84 .map(|r| r.path.clone())
85 .collect();
86
87 let ok_count = results.iter().filter(|r| r.ok).count();
88 info!("{name}: routes {ok_count}/{n} OK");
89
90 if let Some(ref alerter) = alerter {
91 // New failures
92 let new_failures: Vec<String> =
93 current_failed.difference(&prev_failed).cloned().collect();
94 if !new_failures.is_empty() {
95 alerter
96 .send_route_failure_alert(&name, &label, &new_failures)
97 .await;
98 }
99
100 // Recoveries
101 let recoveries: Vec<String> =
102 prev_failed.difference(&current_failed).cloned().collect();
103 if !recoveries.is_empty() {
104 alerter
105 .send_route_recovery_alert(&name, &label, &recoveries)
106 .await;
107 }
108 }
109
110 prev_failed = current_failed;
111 }
112 }));
113 }
114
115 handles
116 }
117