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 async fn copy_test_check_dir(check_dir_name: &str, temp_dir: &Path) {
314 let source_dir = test_data_path().join(check_dir_name);
315 let target_dir = temp_dir.join(check_dir_name);
316 fs::create_dir_all(&target_dir).await.unwrap();
317
318 let mut entries = fs::read_dir(&source_dir).await.unwrap();
319 while let Ok(Some(entry)) = entries.next_entry().await {
320 fs::copy(entry.path(), target_dir.join(entry.file_name()))
321 .await
322 .unwrap();
323 }
324 }
325
326 async fn copy_test_file(source_name: &str, temp_dir: &Path) -> PathBuf {
328 copy_test_file_as(source_name, source_name, temp_dir).await
329 }
330
331 async fn copy_test_file_as(source_name: &str, target_name: &str, temp_dir: &Path) -> PathBuf {
333 let source_path = test_data_path().join(source_name);
334 let target_path = temp_dir.join(target_name);
335
336 let content = fs::read_to_string(&source_path)
337 .await
338 .unwrap_or_else(|_| panic!("Failed to read test file: {:?}", source_path));
339
340 let mut file = fs::File::create(&target_path).await.unwrap();
341 file.write_all(content.as_bytes()).await.unwrap();
342
343 target_path
344 }
345
346 #[tokio::test]
347 async fn test_parse_config_file() {
348 let test_file = test_data_path().join("test-config.yaml");
349
350 let (id, config) = parse_config_file(&test_file, "test-config").await.unwrap();
351
352 assert!(id.contains("saluki-env_src_autodiscovery_providers_test_data_test-config.yaml"));
353 assert_eq!(config.name, "test-config");
354 assert_eq!(
355 config.init_config.value.get("service"),
356 Some(&serde_yaml::Value::String("test-service".to_string()))
357 );
358 assert_eq!(config.source, "local");
359 }
360
361 #[tokio::test]
362 async fn test_parse_minimal_config_file() {
363 let test_file = test_data_path().join("test-minimal-config.yaml");
364
365 let (_, config) = parse_config_file(&test_file, "test-minimal-config").await.unwrap();
366
367 assert!(config.init_config.value.is_empty());
369 }
370
371 #[tokio::test]
372 async fn test_scan_and_emit_events_new_config() {
373 let dir = tempdir().unwrap();
374 let _test_file = copy_test_file("config1.yaml", dir.path()).await;
375
376 let mut known_configs = HashSet::new();
377 let mut configs = BTreeMap::new();
378 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
379 let subscribers = Arc::new(Mutex::new(vec![sender]));
380
381 scan_and_emit_events(
382 &[dir.path().to_path_buf()],
383 &mut known_configs,
384 &subscribers,
385 &mut configs,
386 )
387 .await
388 .unwrap();
389
390 assert_eq!(known_configs.len(), 1);
391
392 let event = receiver.try_recv().unwrap();
393 assert!(matches!(event, AutodiscoveryEvent::CheckSchedule { .. }));
394
395 if let AutodiscoveryEvent::Schedule { config } = event {
396 assert_eq!(config.name, "config1.yaml");
397 assert_eq!(config.instances.len(), 1);
398 assert_eq!(
399 config.instances[0].value.get("server"),
400 Some(&serde_yaml::Value::String("localhost".to_string()))
401 );
402 assert_eq!(
403 config.instances[0].value.get("port"),
404 Some(&serde_yaml::Value::Number(8080.into()))
405 );
406 assert_eq!(
407 config.instances[0].value.get("tags"),
408 Some(&serde_yaml::Value::Sequence(vec![
409 serde_yaml::Value::String("test:true".to_string()),
410 serde_yaml::Value::String("env:test".to_string())
411 ]))
412 );
413 }
414 assert!(receiver.try_recv().is_err());
415 }
416
417 #[tokio::test]
418 async fn test_scan_and_emit_events_preserves_bursts() {
419 let dir = tempdir().unwrap();
420 let _test_file1 = copy_test_file_as("config1.yaml", "config1.yaml", dir.path()).await;
421 let _test_file2 = copy_test_file_as("config1.yaml", "config2.yaml", dir.path()).await;
422
423 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(1);
424 let subscribers = Arc::new(Mutex::new(vec![sender]));
425 let search_path = dir.path().to_path_buf();
426
427 let scan = tokio::spawn(async move {
428 let mut known_configs = HashSet::new();
429 let mut configs = BTreeMap::new();
430
431 scan_and_emit_events(&[search_path], &mut known_configs, &subscribers, &mut configs).await
432 });
433
434 let event1 = tokio::time::timeout(Duration::from_secs(1), receiver.recv())
435 .await
436 .unwrap()
437 .unwrap();
438 let event2 = tokio::time::timeout(Duration::from_secs(1), receiver.recv())
439 .await
440 .unwrap()
441 .unwrap();
442
443 scan.await.unwrap().unwrap();
444
445 assert!(matches!(event1, AutodiscoveryEvent::CheckSchedule { .. }));
446 assert!(matches!(event2, AutodiscoveryEvent::CheckSchedule { .. }));
447 assert!(receiver.try_recv().is_err());
448 }
449
450 #[tokio::test]
451 async fn test_send_to_subscribers_fans_out_and_prunes_closed_receivers() {
452 let (sender1, mut receiver1) = mpsc::channel::<AutodiscoveryEvent>(10);
453 let (sender2, mut receiver2) = mpsc::channel::<AutodiscoveryEvent>(10);
454 let (closed_sender, closed_receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
455 drop(closed_receiver);
456
457 let subscribers = Arc::new(Mutex::new(vec![sender1, sender2, closed_sender]));
458 let event = AutodiscoveryEvent::CheckSchedule {
459 config: CheckConfig {
460 name: MetaString::from("test-check"),
461 init_config: Data::default(),
462 instances: Vec::new(),
463 source: MetaString::from_static("local"),
464 },
465 };
466
467 send_to_subscribers(&subscribers, event).await;
468
469 assert_eq!(subscribers.lock().await.len(), 2);
470 assert!(matches!(
471 receiver1.try_recv().unwrap(),
472 AutodiscoveryEvent::CheckSchedule { .. }
473 ));
474 assert!(matches!(
475 receiver2.try_recv().unwrap(),
476 AutodiscoveryEvent::CheckSchedule { .. }
477 ));
478 }
479
480 #[tokio::test]
481 async fn test_scan_and_emit_events_removed_config() {
482 let dir = tempdir().unwrap();
483
484 let mut known_configs = HashSet::new();
485 known_configs.insert("removed-config".to_string());
486 let mut configs = BTreeMap::new();
487 configs.insert(
488 "removed-config".to_string(),
489 CheckConfig {
490 name: MetaString::from("removed-config"),
491 init_config: Data::default(),
492 instances: Vec::new(),
493 source: MetaString::from_static("local"),
494 },
495 );
496
497 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
498 let subscribers = Arc::new(Mutex::new(vec![sender]));
499
500 scan_and_emit_events(
501 &[dir.path().to_path_buf()],
502 &mut known_configs,
503 &subscribers,
504 &mut configs,
505 )
506 .await
507 .unwrap();
508
509 assert_eq!(known_configs.len(), 0);
510
511 let event = receiver.try_recv().unwrap();
512 assert!(matches!(event, AutodiscoveryEvent::CheckUnscheduled { config } if config.name == "removed-config"));
513
514 assert!(receiver.try_recv().is_err());
515 }
516
517 #[tokio::test]
518 async fn test_scan_and_emit_events_check_dir() {
519 let dir = tempdir().unwrap();
520 copy_test_check_dir("test-check.d", dir.path()).await;
521
522 let mut known_configs = HashSet::new();
523 let mut configs = BTreeMap::new();
524 let (sender, mut receiver) = mpsc::channel::<AutodiscoveryEvent>(10);
525 let subscribers = Arc::new(Mutex::new(vec![sender]));
526
527 scan_and_emit_events(
528 &[dir.path().to_path_buf()],
529 &mut known_configs,
530 &subscribers,
531 &mut configs,
532 )
533 .await
534 .unwrap();
535
536 assert_eq!(known_configs.len(), 1);
538
539 let event = receiver.try_recv().unwrap();
540 assert!(
541 matches!(&event, AutodiscoveryEvent::CheckSchedule { config } if config.name == "test-check"),
542 "expected check name 'test-check', got: {:?}",
543 event
544 );
545
546 assert!(receiver.try_recv().is_err());
547 }
548}