saluki_components/sources/heartbeat/
mod.rs

1use std::time::Duration;
2
3use async_trait::async_trait;
4use saluki_core::{
5    accounting::{MemoryBounds, MemoryBoundsBuilder},
6    components::{sources::*, BuildContext},
7    data_model::event::{
8        metric::{context::Context, Metric},
9        Event, EventType,
10    },
11    topology::OutputDefinition,
12};
13use saluki_error::GenericError;
14use tokio::pin;
15use tokio::{select, time::interval};
16use tracing::{debug, error};
17
18/// Heartbeat source.
19///
20/// Emits a "heartbeat" metric on a configurable interval.
21#[derive(Clone, Debug)]
22pub struct HeartbeatConfiguration {
23    /// Interval for heartbeat metrics in seconds
24    pub heartbeat_interval_secs: u64,
25}
26
27impl Default for HeartbeatConfiguration {
28    fn default() -> Self {
29        Self {
30            heartbeat_interval_secs: 10,
31        }
32    }
33}
34
35struct Heartbeat {
36    heartbeat_interval_secs: u64,
37}
38
39#[async_trait]
40impl Source for Heartbeat {
41    async fn run(self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
42        let global_shutdown = context.take_shutdown_handle();
43        pin!(global_shutdown);
44
45        let mut health = context.take_health_handle();
46        let mut tick_interval = interval(Duration::from_secs(self.heartbeat_interval_secs));
47
48        health.mark_ready();
49        debug!("Heartbeat source started.");
50
51        loop {
52            select! {
53                _ = &mut global_shutdown => {
54                    debug!("Received shutdown signal.");
55                    break;
56                },
57                _ = health.live() => continue,
58                _ = tick_interval.tick() => {
59                    // Create a simple heartbeat metric
60                    let metric_context = Context::from_static_name("heartbeat");
61                    let metric = Metric::gauge(metric_context, 1.0);
62                    let mut buffered_dispatcher = context.dispatcher().buffered().expect("default output must always exist");
63
64                    if let Err(e) = buffered_dispatcher.push(Event::Metric(metric)).await {
65                        error!(error = %e, "Failed to dispatch event.");
66                    } else if let Err(e) = buffered_dispatcher.flush().await {
67                        error!(error = %e, "Failed to dispatch events.");
68                    } else {
69                        debug!("Emitted heartbeat metric.");
70                    }
71                }
72            }
73        }
74
75        debug!("Heartbeat source stopped.");
76        Ok(())
77    }
78}
79
80#[async_trait]
81impl SourceBuilder for HeartbeatConfiguration {
82    async fn build(&self, _context: BuildContext) -> Result<Box<dyn Source + Send>, GenericError> {
83        Ok(Box::new(Heartbeat {
84            heartbeat_interval_secs: self.heartbeat_interval_secs,
85        }))
86    }
87
88    fn outputs(&self) -> &[OutputDefinition<EventType>] {
89        static OUTPUTS: &[OutputDefinition<EventType>] = &[OutputDefinition::default_output(EventType::Metric)];
90        OUTPUTS
91    }
92}
93
94impl MemoryBounds for HeartbeatConfiguration {
95    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
96        // Minimal memory footprint when emitting heartbeat metrics
97        builder.minimum().with_single_value::<Heartbeat>("component struct");
98    }
99}
100
101#[cfg(test)]
102mod tests {
103    use std::mem::size_of;
104
105    use saluki_core::{accounting::ComponentRegistry, support::SubsystemIdentifier};
106
107    use super::*;
108
109    #[test]
110    fn default_configuration_uses_ten_second_interval() {
111        let config = HeartbeatConfiguration::default();
112        assert_eq!(config.heartbeat_interval_secs, 10);
113    }
114
115    #[test]
116    fn declares_single_default_output_for_metrics() {
117        // The source's documented behavior is to emit a heartbeat metric, so it must declare exactly one output --
118        // the default (unnamed) one -- carrying metric events.
119        let config = HeartbeatConfiguration::default();
120        let outputs = config.outputs();
121
122        assert_eq!(outputs.len(), 1);
123        assert_eq!(
124            outputs[0].output_name(),
125            None,
126            "heartbeat emits on the default (unnamed) output"
127        );
128        assert_eq!(outputs[0].data_ty(), EventType::Metric);
129    }
130
131    #[test]
132    fn specify_bounds_accounts_only_for_the_boxed_component_struct() {
133        // The bounds are documented as a "minimal memory footprint": a single boxed `Heartbeat` value and nothing
134        // else. The firm limit includes the minimum, so with no additional firm usage both totals equal the struct
135        // size.
136        let config = HeartbeatConfiguration::default();
137
138        let registry = ComponentRegistry::default();
139        config.specify_bounds(&mut registry.bounds_builder(&SubsystemIdentifier::from_dotted("test")));
140        let bounds = registry.as_bounds();
141
142        assert_eq!(bounds.total_minimum_required_bytes(), size_of::<Heartbeat>());
143        assert_eq!(bounds.total_firm_limit_bytes(), size_of::<Heartbeat>());
144    }
145}