saluki_core/diagnostic/
emitter.rs1use snafu::{OptionExt as _, Snafu};
8use stringtheory::MetaString;
9
10use super::{DiagnosticCollector, DiagnosticEvent};
11use crate::{
12 runtime::state::{DataspaceRegistry, Identifier, IdentifierFilter, Subscription},
13 support::SubsystemIdentifier,
14};
15
16#[derive(Debug, Snafu)]
18#[snafu(context(suffix(false)))]
19pub enum DiagnosticsEmitterError {
20 #[snafu(display("no dataspace available in the current context (not running inside a supervision tree)"))]
25 NoDataspace,
26}
27
28#[derive(Clone)]
66pub struct DiagnosticsEmitter {
67 base_id: MetaString,
68 dataspace: DataspaceRegistry,
69}
70
71impl DiagnosticsEmitter {
72 pub fn from_current(id: SubsystemIdentifier) -> Result<Self, DiagnosticsEmitterError> {
78 let dataspace = DataspaceRegistry::try_current().context(NoDataspace)?;
79 Ok(Self::from_dataspace(id, dataspace))
80 }
81
82 pub fn from_dataspace(id: SubsystemIdentifier, dataspace: DataspaceRegistry) -> Self {
84 let base_id = id.to_string();
85 Self {
86 base_id: base_id.into(),
87 dataspace,
88 }
89 }
90
91 pub fn register_collector<F, T>(&self, artifact_name: impl Into<String>, collect_fn: F)
107 where
108 F: Fn() -> T + Send + Sync + 'static,
109 T: Into<Vec<u8>>,
110 {
111 let collector = DiagnosticCollector::new(artifact_name, collect_fn);
112 let id = self.build_collector_identifier(collector.artifact_name());
113 self.dataspace.assert(collector, id);
114 }
115
116 pub fn unregister_collector(&self, artifact_name: impl AsRef<str>) {
120 let id = self.build_collector_identifier(artifact_name.as_ref());
121 self.dataspace.retract::<DiagnosticCollector>(id);
122 }
123
124 pub fn emit(&self, event: DiagnosticEvent) {
128 self.dataspace.send(event, self.base_id.clone());
129 }
130
131 fn build_collector_identifier(&self, artifact_name: &str) -> Identifier {
132 Identifier::named(format!("{}-{}", self.base_id, artifact_name))
133 }
134}
135
136pub fn subscribe_events(filter: IdentifierFilter) -> Result<Subscription<DiagnosticEvent>, DiagnosticsEmitterError> {
145 let dataspace = DataspaceRegistry::try_current().context(NoDataspace)?;
146 Ok(dataspace.subscribe::<DiagnosticEvent>(filter))
147}
148
149#[cfg(test)]
150mod tests {
151 use tokio_test::{assert_pending, assert_ready, assert_ready_eq, task::spawn as test_spawn};
152
153 use super::*;
154 use crate::{
155 diagnostic::DiagnosticDetails,
156 runtime::{
157 state::{DataspaceUpdate, CURRENT_DATASPACE},
158 ProcessId,
159 },
160 };
161
162 fn emitter(dataspace: DataspaceRegistry) -> DiagnosticsEmitter {
163 DiagnosticsEmitter::from_dataspace(SubsystemIdentifier::from_segments(["sub"]), dataspace)
164 }
165
166 #[test]
167 fn from_current_without_dataspace_errors() {
168 let result = DiagnosticsEmitter::from_current(SubsystemIdentifier::from_segments(["sub"]));
169 assert!(matches!(result, Err(DiagnosticsEmitterError::NoDataspace)));
170 }
171
172 #[test]
173 fn from_current_inside_dataspace_succeeds() {
174 let registry = DataspaceRegistry::new();
175 CURRENT_DATASPACE.sync_scope(registry, || {
176 assert!(DiagnosticsEmitter::from_current(SubsystemIdentifier::from_segments(["sub"])).is_ok());
177 });
178 }
179
180 #[test]
181 fn register_collector_is_discoverable() {
182 let registry = DataspaceRegistry::new();
183 emitter(registry.clone()).register_collector("state.json", || b"hello");
184
185 let collectors = registry.current_values::<DiagnosticCollector>(IdentifierFilter::all());
186 assert_eq!(collectors.len(), 1);
187 assert_eq!(collectors[0].artifact_name(), "state.json");
188 assert_eq!(collectors[0].collect(), b"hello".to_vec());
189 }
190
191 #[test]
192 fn multiple_collectors_coexist() {
193 let registry = DataspaceRegistry::new();
194 let emitter = emitter(registry.clone());
195 emitter.register_collector("tags.json", || vec![1]);
196 emitter.register_collector("eds.json", || vec![2]);
197
198 let mut names: Vec<String> = registry
199 .current_values::<DiagnosticCollector>(IdentifierFilter::all())
200 .iter()
201 .map(|c| c.artifact_name().to_string())
202 .collect();
203 names.sort();
204 assert_eq!(names, vec!["eds.json".to_string(), "tags.json".to_string()]);
205 }
206
207 #[test]
208 fn reregister_same_artifact_updates() {
209 let registry = DataspaceRegistry::new();
210 let emitter = emitter(registry.clone());
211 emitter.register_collector("state.json", || b"v1");
212 emitter.register_collector("state.json", || b"v2");
213
214 let collectors = registry.current_values::<DiagnosticCollector>(IdentifierFilter::all());
215 assert_eq!(collectors.len(), 1);
216 assert_eq!(collectors[0].collect(), b"v2".to_vec());
217 }
218
219 #[test]
220 fn unregister_collector_retracts() {
221 let registry = DataspaceRegistry::new();
222 let emitter = emitter(registry.clone());
223
224 let mut sub = registry.subscribe::<DiagnosticCollector>(IdentifierFilter::all());
225 emitter.register_collector("state.json", || vec![0]);
226 emitter.unregister_collector("state.json");
227
228 let mut recv = test_spawn(sub.recv());
230 match assert_ready!(recv.poll()) {
231 Some(DataspaceUpdate::Asserted(id, collector)) => {
232 assert_eq!(id, Identifier::named("sub-state.json"));
233 assert_eq!(collector.artifact_name(), "state.json");
234 }
235 _ => panic!("expected an assertion first"),
236 }
237 drop(recv);
238
239 let mut recv = test_spawn(sub.recv());
241 match assert_ready!(recv.poll()) {
242 Some(DataspaceUpdate::Retracted(id)) => assert_eq!(id, Identifier::named("sub-state.json")),
243 _ => panic!("expected a retraction second"),
244 }
245 }
246
247 #[test]
248 fn register_collector_is_tagged_to_current_process() {
249 let registry = DataspaceRegistry::new();
250 let emitter = emitter(registry.clone());
251
252 let pid = ProcessId::current();
254 emitter.register_collector("state.json", || vec![0]);
255 assert_eq!(
256 registry
257 .current_values::<DiagnosticCollector>(IdentifierFilter::all())
258 .len(),
259 1
260 );
261
262 registry.retract_all_for_process(pid);
264 assert!(registry
265 .current_values::<DiagnosticCollector>(IdentifierFilter::all())
266 .is_empty());
267 }
268
269 #[test]
270 fn emit_delivers_event_to_subscriber() {
271 let registry = DataspaceRegistry::new();
272 let emitter = emitter(registry.clone());
273
274 let mut sub = registry.subscribe::<DiagnosticEvent>(IdentifierFilter::all());
275 emitter.emit(DiagnosticEvent::new(
276 "credentials rejected",
277 DiagnosticDetails::InvalidApiKey,
278 ));
279
280 let mut recv = test_spawn(sub.recv());
281 assert_ready_eq!(
282 recv.poll(),
283 Some(DataspaceUpdate::Message(
284 Identifier::named("sub"),
285 DiagnosticEvent::new("credentials rejected", DiagnosticDetails::InvalidApiKey)
286 ))
287 );
288 }
289
290 #[test]
291 fn subscribe_events_receives_emitted_event() {
292 let registry = DataspaceRegistry::new();
293 let emitter = emitter(registry.clone());
294
295 CURRENT_DATASPACE.sync_scope(registry, || {
296 let mut sub = subscribe_events(IdentifierFilter::all()).expect("dataspace should be available");
297 emitter.emit(DiagnosticEvent::new("boom", DiagnosticDetails::InvalidApiKey));
298
299 let mut recv = test_spawn(sub.recv());
300 assert_ready_eq!(
301 recv.poll(),
302 Some(DataspaceUpdate::Message(
303 Identifier::named("sub"),
304 DiagnosticEvent::new("boom", DiagnosticDetails::InvalidApiKey)
305 ))
306 );
307 });
308 }
309
310 #[test]
311 fn emit_is_transient() {
312 let registry = DataspaceRegistry::new();
313 let emitter = emitter(registry.clone());
314
315 emitter.emit(DiagnosticEvent::new("boom", DiagnosticDetails::InvalidApiKey));
317
318 assert!(registry
320 .current_values::<DiagnosticEvent>(IdentifierFilter::all())
321 .is_empty());
322
323 let mut sub = registry.subscribe::<DiagnosticEvent>(IdentifierFilter::all());
325 let mut recv = test_spawn(sub.recv());
326 assert_pending!(recv.poll());
327 }
328}