use super::{
fsm::{MetricsAggregatorContext, MetricsAggregatorEvent, MetricsAggregatorState},
inputs::MetricsInputs,
supervisor::MetricsAggregatorSupervisor,
};
use crate::supervised_base::{
BuilderError, ChannelBuilder, HandleBuilder, StandardHandle, SupervisorBuilder,
SupervisorTaskBuilder,
};
use obzenflow_core::{
event::SystemEvent,
journal::Journal,
metrics::{CompositeBoundary, MetricsSnapshotExporter, StageMetadata},
StageId,
};
use std::collections::HashMap;
use std::sync::Arc;
pub struct MetricsAggregatorBuilder {
inputs: MetricsInputs,
system_journal: Arc<dyn Journal<SystemEvent>>,
metrics_exporter: Arc<dyn MetricsSnapshotExporter>,
stage_metadata: HashMap<StageId, StageMetadata>,
composite_boundaries: Vec<CompositeBoundary>,
export_interval: std::time::Duration,
pipeline_writer: Option<obzenflow_core::event::WriterId>,
}
impl MetricsAggregatorBuilder {
pub fn new(
inputs: MetricsInputs,
system_journal: Arc<dyn Journal<SystemEvent>>,
metrics_exporter: Arc<dyn MetricsSnapshotExporter>,
) -> Self {
Self {
inputs,
system_journal,
metrics_exporter,
stage_metadata: HashMap::new(),
composite_boundaries: Vec::new(),
pipeline_writer: None,
export_interval: std::time::Duration::from_millis(
crate::runtime_config::schema::DEFAULT_OBSERVATION_EXPORT_INTERVAL_MS,
),
}
}
pub(crate) fn with_pipeline_writer(mut self, writer: obzenflow_core::event::WriterId) -> Self {
self.pipeline_writer = Some(writer);
self
}
pub(crate) fn with_observation_export_interval(
mut self,
interval: std::time::Duration,
) -> Self {
self.export_interval = interval;
self
}
pub fn with_stage_metadata(mut self, metadata: HashMap<StageId, StageMetadata>) -> Self {
self.stage_metadata = metadata;
self
}
#[doc(hidden)]
pub fn with_composite_boundaries(mut self, boundaries: Vec<CompositeBoundary>) -> Self {
self.composite_boundaries = boundaries;
self
}
}
#[async_trait::async_trait]
impl SupervisorBuilder for MetricsAggregatorBuilder {
type Handle = StandardHandle<MetricsAggregatorEvent, MetricsAggregatorState>;
type Error = BuilderError;
async fn build(self) -> Result<Self::Handle, Self::Error> {
self.prepare().await?.start()
}
}
pub(crate) struct PreparedMetricsAggregator {
context: MetricsAggregatorContext,
system_journal: Arc<dyn Journal<SystemEvent>>,
system_id: obzenflow_core::id::SystemId,
}
impl MetricsAggregatorBuilder {
pub(crate) async fn prepare(self) -> Result<PreparedMetricsAggregator, BuilderError> {
let system_id = obzenflow_core::id::SystemId::new();
let mut metrics_context = MetricsAggregatorContext::new(
self.inputs.clone(),
self.system_journal.clone(),
self.metrics_exporter,
self.export_interval,
system_id,
self.stage_metadata,
self.composite_boundaries,
)
.await
.map_err(BuilderError::Other)?;
metrics_context.pipeline_writer = self.pipeline_writer;
Ok(PreparedMetricsAggregator {
context: metrics_context,
system_journal: self.system_journal,
system_id,
})
}
}
impl PreparedMetricsAggregator {
pub(crate) fn writer_id(&self) -> obzenflow_core::event::WriterId {
self.system_id.into()
}
pub(crate) fn start(self) -> Result<super::MetricsHandle, BuilderError> {
let Self {
context: metrics_context,
system_journal,
system_id,
} = self;
let (event_sender, event_receiver, state_watcher) =
ChannelBuilder::<MetricsAggregatorEvent, MetricsAggregatorState>::new()
.with_event_buffer(10) .build(MetricsAggregatorState::Initializing);
let supervisor = MetricsAggregatorSupervisor {
name: "metrics_aggregator".to_string(),
system_journal,
system_id,
control: event_receiver,
readers: None,
final_refresh: None,
state_watcher: state_watcher.clone(),
last_state: Some(MetricsAggregatorState::Initializing),
};
let supervisor_task =
SupervisorTaskBuilder::<MetricsAggregatorSupervisor>::new("metrics_aggregator")
.spawn_self_supervised(
supervisor,
MetricsAggregatorState::Initializing,
metrics_context,
);
HandleBuilder::new()
.with_event_sender(event_sender)
.with_state_watcher(state_watcher)
.with_supervisor_task(supervisor_task)
.build_standard()
.map_err(|e| BuilderError::Other(e.to_string()))
}
}