saluki_components/sources/heartbeat/
mod.rs1use 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#[derive(Clone, Debug)]
22pub struct HeartbeatConfiguration {
23 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 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 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 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 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}