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