1use std::collections::BTreeMap;
2use std::collections::HashSet;
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use saluki_error::GenericError;
8use serde::Deserialize;
9use stringtheory::MetaString;
10use tokio::fs;
11use tokio::sync::mpsc::{self, Receiver, Sender};
12use tokio::sync::{Mutex, OnceCell};
13use tokio::time::{interval, Duration};
14use tracing::{debug, info, warn};
15
16use crate::autodiscovery::{AutodiscoveryEvent, AutodiscoveryProvider, CheckConfig, Config, Data};
17
18const BG_MONITOR_INTERVAL: u64 = 30;
19
20type AutodiscoverySubscribers = Arc<Mutex<Vec<Sender<AutodiscoveryEvent>>>>;
21
22pub struct LocalAutodiscoveryProvider {
24 search_paths: Vec<PathBuf>,
25 subscribers: AutodiscoverySubscribers,
26 listener_init: OnceCell<()>,
27}
28
29impl LocalAutodiscoveryProvider {
30 pub fn new<P: AsRef<Path>>(paths: Vec<P>) -> Self {
32 let search_paths: Vec<PathBuf> = paths
33 .iter()
34 .filter_map(|p| {
35 if !p.as_ref().exists() {
36 warn!("Skipping path '{}' as it does not exist", p.as_ref().display());
37 return None;
38 }
39 if !p.as_ref().is_dir() {
40 warn!("Skipping path '{}', it is not a directory.", p.as_ref().display());
41 return None;
42 }
43 Some(p.as_ref().to_path_buf())
44 })
45 .collect();
46
47 Self {
48 search_paths,
49 subscribers: Arc::new(Mutex::new(Vec::new())),
50 listener_init: OnceCell::new(),
51 }
52 }
53
54 async fn start_background_monitor(&self, interval_sec: u64) {
56 let mut interval = interval(Duration::from_secs(interval_sec));
57 let subscribers = self.subscribers.clone();
58 let search_paths = self.search_paths.clone();
59
60 info!(
61 "Scanning for local autodiscovery events every {} seconds.",
62 interval_sec
63 );
64
65 tokio::spawn(async move {
66 let mut known_configs = HashSet::new();
67 let mut configs = BTreeMap::new();
68 loop {
69 interval.tick().await;
70
71 if let Err(e) =
73 scan_and_emit_events(&search_paths, &mut known_configs, &subscribers, &mut configs).await
74 {
75 warn!("Error scanning for configurations: {}", e);
76 }
77 }
78 });
79 }
80}
81
82#[derive(Debug, Deserialize)]
83struct LocalCheckConfig {
84 #[serde(default)]
85 init_config: BTreeMap<String, serde_yaml::Value>,
86 instances: Vec<BTreeMap<String, serde_yaml::Value>>,
87}
88
89async fn parse_config_file(path: &PathBuf, check_name: &str) -> Result<(String, CheckConfig), GenericError> {
91 let content = fs::read_to_string(path).await?;
92
93 let check_config: LocalCheckConfig = match serde_yaml::from_str(&content) {
94 Ok(read) => read,
95 Err(e) => {
96 return Err(GenericError::from(e).context("Failed to decode yaml as check configuration."));
97 }
98 };
99
100 let canonicalized_path = fs::canonicalize(&path).await?;
101
102 let config_id = canonicalized_path.to_string_lossy().replace(['/', '\\'], "_");
104
105 let instances: Vec<Data> = check_config
106 .instances
107 .into_iter()
108 .map(|instance| {
109 let mut result = BTreeMap::new();
110 for (key, value) in instance {
111 result.insert(key.into(), value);
112 }
113 Data { value: result }
114 })
115 .collect();
116
117 let init_config = {
118 let mut result = BTreeMap::new();
119 for (key, value) in check_config.init_config {
120 result.insert(key.into(), value);
121 }
122 Data { value: result }
123 };
124
125 let config = Config {
127 name: MetaString::from(check_name),
128 init_config,
129 instances,
130 metric_config: Data::default(),
131 logs_config: Data::default(),
132 ad_identifiers: Vec::new(),
133 provider: MetaString::empty(),
134 service_id: MetaString::empty(),
135 tagger_entity: MetaString::empty(),
136 cluster_check: false,
137 node_name: MetaString::empty(),
138 source: MetaString::from_static("local"),
139 ignore_autodiscovery_tags: false,
140 metrics_excluded: false,
141 logs_excluded: false,
142 advanced_ad_identifiers: Vec::new(),
143 };
144
145 let check_config = CheckConfig::from(config);
146
147 Ok((config_id, check_config))
148}
149
150async fn process_yaml_file(
153 path: PathBuf, check_name: &str, found_configs: &mut HashSet<String>, known_configs: &mut HashSet<String>,
154 subscribers: &AutodiscoverySubscribers, configs: &mut BTreeMap<String, CheckConfig>,
155) {
156 match parse_config_file(&path, check_name).await {
157 Ok((config_id, config)) => {
158 found_configs.insert(config_id.clone());
159
160 if !known_configs.contains(&config_id) {
161 debug!("New configuration found: {}", config_id);
162 let event = AutodiscoveryEvent::CheckSchedule { config: config.clone() };
163 send_to_subscribers(subscribers, event).await;
164 known_configs.insert(config_id.clone());
165 configs.insert(config_id, config);
166 } else {
167 let existing_config = configs.get(&config_id).unwrap();
169 if *existing_config != config {
170 configs.insert(config_id.clone(), config.clone());
171 debug!("Configuration updated: {}", config_id);
172 let event = AutodiscoveryEvent::CheckSchedule { config };
173 send_to_subscribers(subscribers, event).await;
174 }
175 }
176 }
177 Err(e) => {
178 warn!("Failed to parse config file {}: {}", path.display(), e);
179 }
180 }
181}
182
183fn is_yaml_file(path: &Path) -> bool {
185 matches!(path.extension().and_then(|e| e.to_str()), Some("yaml") | Some("yml"))
186}
187
188async fn scan_and_emit_events(
195 paths: &[PathBuf], known_configs: &mut HashSet<String>, subscribers: &AutodiscoverySubscribers,
196 configs: &mut BTreeMap<String, CheckConfig>,
197) -> Result<(), GenericError> {
198 let mut found_configs = HashSet::new();
199
200 for path in paths {
201 let mut entries = fs::read_dir(path).await?;
202 while let Ok(Some(entry)) = entries.next_entry().await {
203 let path = entry.path();
204
205 if is_yaml_file(&path) {
206 let check_name = path.file_stem().unwrap().to_string_lossy().into_owned();
208 process_yaml_file(
209 path,
210 &check_name,
211 &mut found_configs,
212 known_configs,
213 subscribers,
214 configs,
215 )
216 .await;
217 } else if path.is_dir() {
218 let dir_name = path.file_name().unwrap_or_default().to_string_lossy();
220 if let Some(check_name) = dir_name.strip_suffix(".d") {
221 let Ok(mut sub_entries) = fs::read_dir(&path).await else {
222 warn!("Failed to read directory {}", path.display());
223 continue;
224 };
225 while let Ok(Some(sub_entry)) = sub_entries.next_entry().await {
226 let sub_path = sub_entry.path();
227 if is_yaml_file(&sub_path) {
228 process_yaml_file(
229 sub_path,
230 check_name,
231 &mut found_configs,
232 known_configs,
233 subscribers,
234 configs,
235 )
236 .await;
237 }
238 }
239 }
240 }
241 }
242 }
243
244 let to_remove: Vec<String> = known_configs
246 .iter()
247 .filter(|config_id| !found_configs.contains(*config_id))
248 .cloned()
249 .collect();
250
251 for config_id in to_remove {
252 debug!("Configuration removed: {}", config_id);
253 known_configs.remove(&config_id);
254
255 let config = configs.remove(&config_id).unwrap();
256
257 let event = AutodiscoveryEvent::CheckUnscheduled { config };
259 send_to_subscribers(subscribers, event).await;
260 }
261
262 Ok(())
263}
264
265async fn send_to_subscribers(subscribers: &AutodiscoverySubscribers, event: AutodiscoveryEvent) {
266 let mut subscribers = subscribers.lock().await;
267 let mut active_subscribers = Vec::with_capacity(subscribers.len());
268
269 for sender in subscribers.drain(..) {
270 if sender.send(event.clone()).await.is_ok() {
271 active_subscribers.push(sender);
272 }
273 }
274
275 *subscribers = active_subscribers;
276}
277
278#[async_trait]
279impl AutodiscoveryProvider for LocalAutodiscoveryProvider {
280 async fn subscribe(&self) -> Option<Receiver<AutodiscoveryEvent>> {
281 self.listener_init
282 .get_or_init(|| async {
283 self.start_background_monitor(BG_MONITOR_INTERVAL).await;
284 })
285 .await;
286
287 let (sender, receiver) = mpsc::channel::<AutodiscoveryEvent>(16);
288 self.subscribers.lock().await.push(sender);
289 Some(receiver)
290 }
291}
292
293#[cfg(test)]
294mod tests {
295 use std::path::Path;
296
297 use tempfile::tempdir;
298 use tokio::io::AsyncWriteExt;
299
300 use super::*;
301
302 fn test_data_path() -> PathBuf {
304 let manifest_dir = std::env::var("CARGO_MANIFEST_DIR").unwrap_or_else(|_| ".".to_string());
305 PathBuf::from(manifest_dir)
306 .join("src")
307 .join("autodiscovery")
308 .join("providers")
309 .join("test_data")
310 }
311
312 fn check_config(name: &'static str) -> CheckConfig {
313 CheckConfig {
314 name: MetaString::from_static(name),
315 init_config: Data::default(),
316 instances: Vec::new(),
317 metric_config: Data::default(),
318 logs_config: Data::default(),
319 ad_identifiers: Vec::new(),
320 advanced_ad_identifiers: Vec::new(),
321 provider: MetaString::default(),
322 service_id: MetaString::default(),
323 tagger_entity: MetaString::default(),
324 cluster_check: false,
325 node_name: MetaString::default(),
326 source: MetaString::from_static("local"),
327 ignore_autodiscovery_tags: false,
328 metrics_excluded: false,
329 logs_excluded: false,
330 }
331 }
332
333 async fn copy_test_check_dir(check_dir_name: &str, temp_dir: &Path) {
335 let source_dir = test_data_path().join(check_dir_name);
336 let target_dir = temp_dir.join(check_dir_name);
337 fs::create_dir_all(&target_dir).await.unwrap();
338
339 let mut entries = fs::read_dir(&source_dir).await.unwrap();
340 while let Ok(Some(entry)) = entries.next_entry().await {
341 fs::copy(entry.path(), target_dir.join(entry.file_name()))
342 .await
343 .unwrap();
344 }
345 }
346
347 async fn copy_test_file(source_name: &str, temp_dir: &Path) -> PathBuf {
349 copy_test_file_as(source_name, source_name, temp_dir).await
350 }
351
352 async fn copy_test_file_as(source_name: &str, target_name: &str, temp_dir: &Path) -> PathBuf {
354 let source_path = test_data_path().join(source_name);
355 let target_path = temp_dir.join(target_name);
356
357 let content = fs::read_to_string(&source_path)
358 .await
359 .unwrap_or_else(|_| panic!("Failed to read test file: {:?}", source_path));
360
361 let mut file = fs::File::create(&target_path).await.unwrap();
362 file.write_all(content.as_bytes()).await.unwrap();
363
364 target_path
365 }
366
367 #[tokio::test]
368 async fn test_parse_config_file() {
369 let test_file = test_data_path().join("test-config.yaml");
370
371 let (id, config) = parse_config_file(&test_file, "test-config").await.unwrap();
372
373 assert!(id.contains("saluki-env_src_autodiscovery_providers_test_data_test-config.yaml"));
374 assert_eq!(config.name, "test-config");
375 assert_eq!(
376 config.init_config.value.get("service"),
377 Some(&serde_yaml::Value::String("test-service".to_string()))
378 );
379 assert_eq!(config.source, "local");
380 }
381
382 #[tokio::test]
383 async fn test_parse_minimal_config_file() {
384 let test_file = test_data_path().join("test-minimal-config.yaml");
385
386 let (_, config) = parse_config_file(&test_file, "test-minimal-config").await.unwrap();
387
388 assert!(config.init_config.value.is_empty());
390 }
391
392 #[tokio::test]
393 async fn test_scan_and_emit_events_new_config() {
394 let dir = tempdir().unwrap();
395 let _test_file = copy_test_file("config1.yaml", dir.path()).await;
396
397 let mut known_configs = HashSet::new();
398 let mut configs = BTreeMap::new();
399 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
400 let subscribers = Arc::new(Mutex::new(vec![sender]));
401
402 scan_and_emit_events(
403 &[dir.path().to_path_buf()],
404 &mut known_configs,
405 &subscribers,
406 &mut configs,
407 )
408 .await
409 .unwrap();
410
411 assert_eq!(known_configs.len(), 1);
412
413 let event = receiver.try_recv().unwrap();
414 assert!(matches!(event, AutodiscoveryEvent::CheckSchedule { .. }));
415
416 if let AutodiscoveryEvent::Schedule { config } = event {
417 assert_eq!(config.name, "config1.yaml");
418 assert_eq!(config.instances.len(), 1);
419 assert_eq!(
420 config.instances[0].value.get("server"),
421 Some(&serde_yaml::Value::String("localhost".to_string()))
422 );
423 assert_eq!(
424 config.instances[0].value.get("port"),
425 Some(&serde_yaml::Value::Number(8080.into()))
426 );
427 assert_eq!(
428 config.instances[0].value.get("tags"),
429 Some(&serde_yaml::Value::Sequence(vec![
430 serde_yaml::Value::String("test:true".to_string()),
431 serde_yaml::Value::String("env:test".to_string())
432 ]))
433 );
434 }
435 assert!(receiver.try_recv().is_err());
436 }
437
438 #[tokio::test]
439 async fn test_scan_and_emit_events_preserves_bursts() {
440 let dir = tempdir().unwrap();
441 let _test_file1 = copy_test_file_as("config1.yaml", "config1.yaml", dir.path()).await;
442 let _test_file2 = copy_test_file_as("config1.yaml", "config2.yaml", dir.path()).await;
443
444 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(1);
445 let subscribers = Arc::new(Mutex::new(vec![sender]));
446 let search_path = dir.path().to_path_buf();
447
448 let scan = tokio::spawn(async move {
449 let mut known_configs = HashSet::new();
450 let mut configs = BTreeMap::new();
451
452 scan_and_emit_events(&[search_path], &mut known_configs, &subscribers, &mut configs).await
453 });
454
455 let event1 = tokio::time::timeout(Duration::from_secs(1), receiver.recv())
456 .await
457 .unwrap()
458 .unwrap();
459 let event2 = tokio::time::timeout(Duration::from_secs(1), receiver.recv())
460 .await
461 .unwrap()
462 .unwrap();
463
464 scan.await.unwrap().unwrap();
465
466 assert!(matches!(event1, AutodiscoveryEvent::CheckSchedule { .. }));
467 assert!(matches!(event2, AutodiscoveryEvent::CheckSchedule { .. }));
468 assert!(receiver.try_recv().is_err());
469 }
470
471 #[tokio::test]
472 async fn test_send_to_subscribers_fans_out_and_prunes_closed_receivers() {
473 let (sender1, mut receiver1) = mpsc::channel::<AutodiscoveryEvent>(10);
474 let (sender2, mut receiver2) = mpsc::channel::<AutodiscoveryEvent>(10);
475 let (closed_sender, closed_receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
476 drop(closed_receiver);
477
478 let subscribers = Arc::new(Mutex::new(vec![sender1, sender2, closed_sender]));
479 let event = AutodiscoveryEvent::CheckSchedule {
480 config: check_config("test-check"),
481 };
482
483 send_to_subscribers(&subscribers, event).await;
484
485 assert_eq!(subscribers.lock().await.len(), 2);
486 assert!(matches!(
487 receiver1.try_recv().unwrap(),
488 AutodiscoveryEvent::CheckSchedule { .. }
489 ));
490 assert!(matches!(
491 receiver2.try_recv().unwrap(),
492 AutodiscoveryEvent::CheckSchedule { .. }
493 ));
494 }
495
496 #[tokio::test]
497 async fn test_scan_and_emit_events_removed_config() {
498 let dir = tempdir().unwrap();
499
500 let mut known_configs = HashSet::new();
501 known_configs.insert("removed-config".to_string());
502 let mut configs = BTreeMap::new();
503 configs.insert("removed-config".to_string(), check_config("removed-config"));
504
505 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
506 let subscribers = Arc::new(Mutex::new(vec![sender]));
507
508 scan_and_emit_events(
509 &[dir.path().to_path_buf()],
510 &mut known_configs,
511 &subscribers,
512 &mut configs,
513 )
514 .await
515 .unwrap();
516
517 assert_eq!(known_configs.len(), 0);
518
519 let event = receiver.try_recv().unwrap();
520 assert!(matches!(event, AutodiscoveryEvent::CheckUnscheduled { config } if config.name == "removed-config"));
521
522 assert!(receiver.try_recv().is_err());
523 }
524
525 #[tokio::test]
526 async fn test_scan_and_emit_events_check_dir() {
527 let dir = tempdir().unwrap();
528 copy_test_check_dir("test-check.d", dir.path()).await;
529
530 let mut known_configs = HashSet::new();
531 let mut configs = BTreeMap::new();
532 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
533 let subscribers = Arc::new(Mutex::new(vec![sender]));
534
535 scan_and_emit_events(
536 &[dir.path().to_path_buf()],
537 &mut known_configs,
538 &subscribers,
539 &mut configs,
540 )
541 .await
542 .unwrap();
543
544 assert_eq!(known_configs.len(), 1);
546
547 let event = receiver.try_recv().unwrap();
548 assert!(
549 matches!(&event, AutodiscoveryEvent::CheckSchedule { config } if config.name == "test-check"),
550 "expected check name 'test-check', got: {:?}",
551 event
552 );
553
554 assert!(receiver.try_recv().is_err());
555 }
556}