saluki_core/runtime/tree/
worker.rs1use async_trait::async_trait;
2use saluki_api::{DynamicRoute, EndpointType};
3use saluki_common::sync::shutdown::ShutdownHandle;
4use saluki_error::generic_error;
5
6use super::SupervisionTreeHandle;
7use crate::{
8 diagnostic::DiagnosticsEmitter,
9 runtime::{state::DataspaceRegistry, InitializationError, Supervisable, SupervisorFuture},
10 support::SubsystemIdentifier,
11};
12
13pub struct SupervisionTreeWorker {
24 tree: SupervisionTreeHandle,
25}
26
27impl SupervisionTreeWorker {
28 pub(super) fn new(tree: SupervisionTreeHandle) -> Self {
29 Self { tree }
30 }
31}
32
33#[async_trait]
34impl Supervisable for SupervisionTreeWorker {
35 fn name(&self) -> &str {
36 "supervision-tree"
37 }
38
39 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
40 let tree_routes = DynamicRoute::http(EndpointType::Privileged, self.tree.api_handler());
41
42 let tree = self.tree.clone();
43
44 Ok(Box::pin(async move {
45 let dataspace =
46 DataspaceRegistry::try_current().ok_or_else(|| generic_error!("Dataspace not available."))?;
47
48 dataspace.assert(tree_routes, "supervision-tree-api");
50
51 let diagnostics =
53 DiagnosticsEmitter::from_dataspace(SubsystemIdentifier::from_segments(["supervision-tree"]), dataspace);
54 diagnostics.register_collector("supervision_tree.json", move || tree.snapshot_json());
55
56 process_shutdown.await;
57
58 Ok(())
59 }))
60 }
61}