use obzenflow_core::event::chain_event::ChainEventContent;
use obzenflow_core::event::context::StageType;
use obzenflow_core::event::observability::{HttpPullTelemetry, HttpSurfaceRouteMetricsSnapshot};
use obzenflow_core::event::payloads::observability_payload::{
CircuitBreakerEvent, CircuitBreakerOpenTrigger, MetricsLifecycle, MiddlewareLifecycle,
ObservabilityPayload,
};
use obzenflow_core::event::status::processing_status::{ErrorKind, ProcessingStatus};
use obzenflow_core::event::{JournalEvent, SinkOperationFailed, SinkOperationPhase, WriterId};
use obzenflow_core::id::{FlowId, StageId, SystemId};
use obzenflow_core::ingress::IngressKey;
use obzenflow_core::metrics::{
BoundaryMetricsView, CompositeDurationAccumulator, ContractMetricEdgeKey,
ContractMetricResultKey, ContractMetricViolationKey, ContractMetricsSnapshot,
ContractViolationCauseLabel, Percentile, StageMetadata,
};
use obzenflow_core::time::MetricsDuration;
use obzenflow_core::web::HttpMethod;
use obzenflow_core::{ChainEvent, EventId, EventType, Journal, TypedPayload};
use obzenflow_fsm::{
fsm, EventVariant, FsmAction, FsmContext, StateMachine, StateVariant, Transition,
};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::str::FromStr;
use std::sync::Arc;
use crate::metrics::tail_read;
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub enum MetricsAggregatorState {
Initializing,
Running,
Draining,
Drained { last_event_id: Option<EventId> },
Failed { error: String },
}
impl StateVariant for MetricsAggregatorState {
fn variant_name(&self) -> &str {
match self {
MetricsAggregatorState::Initializing => "Initializing",
MetricsAggregatorState::Running => "Running",
MetricsAggregatorState::Draining => "Draining",
MetricsAggregatorState::Drained { .. } => "Drained",
MetricsAggregatorState::Failed { .. } => "Failed",
}
}
}
#[derive(Clone, Debug)]
pub enum MetricsAggregatorEvent {
StartRunning,
ProcessBatch {
events: Vec<obzenflow_core::EventEnvelope<obzenflow_core::ChainEvent>>,
journal_kind: MetricsJournalKind,
journal_stage: StageId,
},
ProcessSystemEvent {
envelope: Box<obzenflow_core::EventEnvelope<obzenflow_core::event::SystemEvent>>,
},
ExportMetrics,
StartDraining,
FlowTerminal,
Error(String),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum MetricsJournalKind {
Data,
Error,
}
impl EventVariant for MetricsAggregatorEvent {
fn variant_name(&self) -> &str {
match self {
MetricsAggregatorEvent::StartRunning => "StartRunning",
MetricsAggregatorEvent::ProcessBatch { .. } => "ProcessBatch",
MetricsAggregatorEvent::ProcessSystemEvent { .. } => "ProcessSystemEvent",
MetricsAggregatorEvent::ExportMetrics => "ExportMetrics",
MetricsAggregatorEvent::StartDraining => "StartDraining",
MetricsAggregatorEvent::FlowTerminal => "FlowTerminal",
MetricsAggregatorEvent::Error(_) => "Error",
}
}
}
#[derive(Clone, Debug)]
pub enum MetricsAggregatorAction {
Initialize,
UpdateMetrics {
envelope: Box<obzenflow_core::EventEnvelope<obzenflow_core::ChainEvent>>,
journal_kind: MetricsJournalKind,
journal_stage: StageId,
},
ProcessSystemEvent {
envelope: Box<obzenflow_core::EventEnvelope<obzenflow_core::event::SystemEvent>>,
},
ExportMetrics,
PublishDrainComplete { last_event_id: Option<EventId> },
}
pub struct MetricsAggregatorContext {
pub system_journal: Arc<dyn Journal<obzenflow_core::event::SystemEvent>>,
pub stage_data_journals: HashMap<StageId, Arc<dyn Journal<ChainEvent>>>,
pub stage_error_journals: HashMap<StageId, Arc<dyn Journal<ChainEvent>>>,
pub backpressure_registry: Option<Arc<crate::backpressure::BackpressureRegistry>>,
pub include_error_journals: bool,
pub exporter: Option<Arc<dyn obzenflow_core::metrics::MetricsExporter>>,
pub metrics_store: MetricsStore,
pub export_interval_secs: u64,
pub system_id: SystemId,
pub stage_metadata: HashMap<StageId, StageMetadata>,
#[doc(hidden)]
pub composite_boundaries: Vec<obzenflow_core::metrics::CompositeBoundary>,
#[doc(hidden)]
pub composite_durations: CompositeDurationAccumulator,
}
pub(crate) struct MetricsAggregatorIo {
pub(crate) data_subscription:
crate::messaging::upstream_subscription::UpstreamSubscription<ChainEvent>,
pub(crate) error_subscription:
Option<crate::messaging::upstream_subscription::UpstreamSubscription<ChainEvent>>,
pub(crate) system_subscription: crate::messaging::system_subscription::SystemSubscription<
obzenflow_core::event::SystemEvent,
>,
}
#[derive(Default)]
#[doc(hidden)]
pub struct MetricsStore {
pub stage_metrics: std::collections::HashMap<StageId, StageMetrics>,
pub last_event_id: Option<EventId>,
pub flow_start_time: Option<std::time::Instant>,
pub first_event_time: Option<std::time::Instant>,
pub last_event_time: Option<std::time::Instant>,
pub total_events_processed: u64,
pub sink_operation_failures: HashMap<(StageId, SinkOperationPhase, ErrorKind), u64>,
pub stage_vector_clocks: HashMap<StageId, u64>,
pub system_vector_clocks: HashMap<SystemId, u64>,
pub circuit_breaker_state: HashMap<StageId, f64>,
pub circuit_breaker_rejection_rate: HashMap<StageId, f64>,
pub circuit_breaker_consecutive_failures: HashMap<StageId, f64>,
pub circuit_breaker_requests_total: HashMap<StageId, u64>,
pub circuit_breaker_rejections_total: HashMap<StageId, u64>,
pub circuit_breaker_opened_total: HashMap<StageId, u64>,
pub circuit_breaker_successes_total: HashMap<StageId, u64>,
pub circuit_breaker_failures_total: HashMap<StageId, u64>,
pub circuit_breaker_slow_total: HashMap<StageId, u64>,
pub circuit_breaker_time_in_state_seconds_total: HashMap<(StageId, String), f64>,
pub circuit_breaker_state_transitions_total: HashMap<(StageId, String, String), u64>,
circuit_breaker_last_state: HashMap<StageId, String>,
pub rate_limiter_events_total: HashMap<StageId, u64>,
pub rate_limiter_delayed_total: HashMap<StageId, u64>,
pub rate_limiter_tokens_consumed_total: HashMap<StageId, f64>,
pub rate_limiter_delay_seconds_total: HashMap<StageId, f64>,
pub rate_limiter_bucket_tokens: HashMap<StageId, f64>,
pub rate_limiter_bucket_capacity: HashMap<StageId, f64>,
pub contract_metrics: ContractMetricsSnapshot,
pub edge_liveness_state: HashMap<(StageId, StageId), f64>,
pub http_surface_metrics:
HashMap<(String, HttpMethod, String, String), HttpSurfaceRouteMetricsSnapshot>,
pub ingestion_refusals_total: HashMap<(IngressKey, String), u64>,
pub http_pull_metrics: HashMap<StageId, HttpPullTelemetry>,
pub ai_chunking_metrics: HashMap<StageId, obzenflow_core::metrics::AiChunkingMetricsSnapshot>,
pub stage_lifecycle_states: HashMap<(StageId, String), bool>,
pub pipeline_state: String,
}
#[derive(Clone, Default)]
pub struct StageMetrics {
pub errors_by_kind: HashMap<ErrorKind, u64>,
pub last_in_flight: Option<u32>,
pub last_failures_total: Option<u64>,
pub join_reference_since_last_stream: Option<u64>,
pub latest_events_processed_total: Option<u64>,
pub latest_events_accumulated_total: Option<u64>,
pub latest_events_emitted_total: Option<u64>,
pub latest_data_outputs_by_event_type: HashMap<EventType, u64>,
pub latest_data_inputs_by_upstream_event_type: HashMap<(StageId, EventType), u64>,
pub latest_errors_total: Option<u64>,
pub event_loops_total: u64,
pub event_loops_with_work_total: u64,
pub snapshot_p50_ms: Option<u64>,
pub snapshot_p90_ms: Option<u64>,
pub snapshot_p95_ms: Option<u64>,
pub snapshot_p99_ms: Option<u64>,
pub snapshot_p999_ms: Option<u64>,
pub processing_time_sum_nanos: Option<u64>,
pub first_event_time: Option<std::time::Instant>,
pub last_event_time: Option<std::time::Instant>,
}
impl StageMetrics {
fn merge_runtime_context(
&mut self,
runtime_ctx: &obzenflow_core::runtime_context::RuntimeContext,
) {
self.last_in_flight = Some(runtime_ctx.in_flight);
self.last_failures_total = Some(runtime_ctx.failures_total);
self.join_reference_since_last_stream = Some(runtime_ctx.join_reference_since_last_stream);
self.latest_events_processed_total = Some(
self.latest_events_processed_total
.unwrap_or(0)
.max(runtime_ctx.events_processed_total),
);
self.latest_events_accumulated_total = Some(
self.latest_events_accumulated_total
.unwrap_or(0)
.max(runtime_ctx.events_accumulated_total),
);
self.latest_events_emitted_total = Some(
self.latest_events_emitted_total
.unwrap_or(0)
.max(runtime_ctx.events_emitted_total),
);
self.latest_errors_total = Some(
self.latest_errors_total
.unwrap_or(0)
.max(runtime_ctx.errors_total),
);
self.event_loops_total = self.event_loops_total.max(runtime_ctx.event_loops_total);
self.event_loops_with_work_total = self
.event_loops_with_work_total
.max(runtime_ctx.event_loops_with_work_total);
self.snapshot_p50_ms = Some(runtime_ctx.recent_p50_ms);
self.snapshot_p90_ms = Some(runtime_ctx.recent_p90_ms);
self.snapshot_p95_ms = Some(runtime_ctx.recent_p95_ms);
self.snapshot_p99_ms = Some(runtime_ctx.recent_p99_ms);
self.snapshot_p999_ms = Some(runtime_ctx.recent_p999_ms);
self.processing_time_sum_nanos = Some(runtime_ctx.processing_time_sum_nanos);
for count in &runtime_ctx.data_outputs_by_event_type {
let current = self
.latest_data_outputs_by_event_type
.entry(count.event_type.clone())
.or_insert(0);
*current = (*current).max(count.total);
}
for count in &runtime_ctx.data_inputs_by_upstream_event_type {
let current = self
.latest_data_inputs_by_upstream_event_type
.entry((count.upstream, count.event_type.clone()))
.or_insert(0);
*current = (*current).max(count.total);
}
}
}
impl BoundaryMetricsView for MetricsStore {
fn data_inputs(&self, member: StageId, upstream: StageId, event_type: &EventType) -> u64 {
self.stage_metrics
.get(&member)
.and_then(|metrics| {
metrics
.latest_data_inputs_by_upstream_event_type
.get(&(upstream, event_type.clone()))
})
.copied()
.unwrap_or(0)
}
fn data_outputs(&self, member: StageId, event_type: &EventType) -> u64 {
self.stage_metrics
.get(&member)
.and_then(|metrics| metrics.latest_data_outputs_by_event_type.get(event_type))
.copied()
.unwrap_or(0)
}
fn errors(&self, member: StageId) -> u64 {
self.stage_metrics
.get(&member)
.and_then(|metrics| metrics.latest_errors_total)
.unwrap_or(0)
}
}
async fn fold_composite_duration_prefix(
reader: &mut dyn obzenflow_core::journal::journal_reader::JournalReader<ChainEvent>,
journal_stage: StageId,
boundaries: &[obzenflow_core::metrics::CompositeBoundary],
accumulator: &mut CompositeDurationAccumulator,
) -> Result<u64, obzenflow_core::journal::journal_error::JournalError> {
let mut position = 0;
while let Some(envelope) = reader.next().await? {
accumulator.observe_event(boundaries, journal_stage, &envelope.event);
position += 1;
}
Ok(position)
}
fn observe_live_composite_duration(
journal_kind: MetricsJournalKind,
accumulator: &mut CompositeDurationAccumulator,
boundaries: &[obzenflow_core::metrics::CompositeBoundary],
journal_stage: StageId,
event: &ChainEvent,
) {
if journal_kind == MetricsJournalKind::Data {
accumulator.observe_event(boundaries, journal_stage, event);
}
}
impl MetricsAggregatorContext {
pub(crate) async fn new(
inputs: crate::metrics::inputs::MetricsInputs,
system_journal: Arc<dyn Journal<obzenflow_core::event::SystemEvent>>,
exporter: Option<Arc<dyn obzenflow_core::metrics::MetricsExporter>>,
export_interval_secs: u64,
system_id: SystemId,
stage_metadata: HashMap<StageId, StageMetadata>,
composite_boundaries: Vec<obzenflow_core::metrics::CompositeBoundary>,
) -> Result<(Self, MetricsAggregatorIo), String> {
let mut metrics_store = MetricsStore::default();
let mut composite_durations = CompositeDurationAccumulator::default();
let with_names = |journals: &[(StageId, Arc<dyn Journal<ChainEvent>>)]| {
journals
.iter()
.map(|(id, journal)| {
let name = stage_metadata
.get(id)
.map(|m| m.name.clone())
.unwrap_or_else(|| format!("{id:?}"));
(*id, name, journal.clone())
})
.collect::<Vec<_>>()
};
let data_with_names = with_names(&inputs.stage_data_journals);
let mut data_start_positions = Vec::with_capacity(data_with_names.len());
for (stage_id, stage_name, journal) in &data_with_names {
let tail_snapshot = match journal.read_last_n(1).await {
Ok(mut events) => events
.drain(..)
.find(|env| env.event.runtime_context.is_some()),
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed to tail-read data journal for snapshot; seeding skipped for this stage"
);
None
}
};
if let Some(envelope) = &tail_snapshot {
if let Some(runtime_ctx) = &envelope.event.runtime_context {
{
let metrics = metrics_store
.stage_metrics
.entry(*stage_id)
.or_insert_with(StageMetrics::default);
metrics.merge_runtime_context(runtime_ctx);
}
metrics_store
.update_control_metrics_from_runtime_context(*stage_id, runtime_ctx);
}
let writer_id = *envelope.event.writer_id();
let writer_key = writer_id.to_string();
let seq = envelope.vector_clock.get(&writer_key);
let entry = metrics_store
.stage_vector_clocks
.entry(*stage_id)
.or_insert(0);
*entry = (*entry).max(seq);
}
let start_position = match journal.reader().await {
Ok(mut reader) => {
match fold_composite_duration_prefix(
reader.as_mut(),
*stage_id,
&composite_boundaries,
&mut composite_durations,
)
.await
{
Ok(position) => position,
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed while streaming data journal to determine tail position; starting from 0"
);
0
}
}
}
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed to create reader for data journal; starting from 0"
);
0
}
};
data_start_positions.push(start_position);
}
tracing::info!(
upstream_count = inputs.stage_data_journals.len(),
upstream_stages = tracing::field::debug(
&inputs
.stage_data_journals
.iter()
.map(|(id, _)| *id)
.collect::<Vec<_>>(),
),
"MetricsAggregator creating data subscription (tail-start)"
);
let data_subscription =
crate::messaging::upstream_subscription::UpstreamSubscription::new_at_tail(
"metrics_aggregator",
&data_with_names,
&data_start_positions,
)
.await
.map_err(|e| format!("Failed to create data subscription: {e}"))?;
let error_with_names = with_names(&inputs.error_journals);
for (stage_id, stage_name, journal) in &error_with_names {
match journal.read_last_n(1).await {
Ok(mut events) => {
if let Some(envelope) = events
.drain(..)
.find(|env| env.event.runtime_context.is_some())
{
if let Some(runtime_ctx) = &envelope.event.runtime_context {
{
let metrics = metrics_store
.stage_metrics
.entry(*stage_id)
.or_insert_with(StageMetrics::default);
metrics.merge_runtime_context(runtime_ctx);
}
metrics_store.update_control_metrics_from_runtime_context(
*stage_id,
runtime_ctx,
);
}
let writer_id = *envelope.event.writer_id();
let writer_key = writer_id.to_string();
let seq = envelope.vector_clock.get(&writer_key);
let entry = metrics_store
.stage_vector_clocks
.entry(*stage_id)
.or_insert(0);
*entry = (*entry).max(seq);
}
}
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed to tail-read error journal for snapshot; seeding skipped for this stage"
);
}
}
}
let mut error_start_positions = Vec::with_capacity(error_with_names.len());
for (stage_id, stage_name, journal) in &error_with_names {
let start_position = match journal.reader().await {
Ok(mut reader) => {
let mut pos: u64 = 0;
loop {
match reader.next().await {
Ok(Some(_)) => {
pos += 1;
}
Ok(None) => break,
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed while streaming error journal to determine tail position; starting from 0"
);
pos = 0;
break;
}
}
}
pos
}
Err(e) => {
tracing::warn!(
target: "flowip-059",
owner = "metrics_aggregator",
stage_id = ?stage_id,
stage_name = stage_name,
error = ?e,
"Failed to create reader for error journal; starting from 0"
);
0
}
};
error_start_positions.push(start_position);
}
if !inputs.error_journals.is_empty() {
tracing::info!(
upstream_count = inputs.error_journals.len(),
upstream_stages = tracing::field::debug(
&inputs
.error_journals
.iter()
.map(|(id, _)| *id)
.collect::<Vec<_>>(),
),
"MetricsAggregator creating error subscription (tail-start)"
);
}
let error_subscription = if !inputs.error_journals.is_empty() {
Some(
crate::messaging::upstream_subscription::UpstreamSubscription::new_at_tail(
"metrics_aggregator",
&error_with_names,
&error_start_positions,
)
.await
.map_err(|e| format!("Failed to create error subscription: {e}"))?,
)
} else {
None
};
let system_reader = system_journal
.reader()
.await
.map_err(|e| format!("Failed to create system journal reader: {e:?}"))?;
let system_subscription = crate::messaging::system_subscription::SystemSubscription::new(
system_reader,
"metrics_aggregator".to_string(),
);
let stage_data_journals: HashMap<StageId, Arc<dyn Journal<ChainEvent>>> = inputs
.stage_data_journals
.iter()
.map(|(id, journal)| (*id, journal.clone()))
.collect();
let stage_error_journals: HashMap<StageId, Arc<dyn Journal<ChainEvent>>> = inputs
.error_journals
.iter()
.map(|(id, journal)| (*id, journal.clone()))
.collect();
let context = Self {
system_journal,
stage_data_journals,
stage_error_journals,
backpressure_registry: inputs.backpressure_registry.clone(),
include_error_journals: true, exporter,
metrics_store,
export_interval_secs,
system_id,
stage_metadata,
composite_boundaries,
composite_durations,
};
let io = MetricsAggregatorIo {
data_subscription,
error_subscription,
system_subscription,
};
Ok((context, io))
}
}
impl FsmContext for MetricsAggregatorContext {}
fn infer_join_reference_mode_from_fsm_state(fsm_state: &str) -> Option<&'static str> {
match fsm_state {
"Live" => Some("live"),
"Hydrating" | "Enriching" => Some("finite_eof"),
_ => None,
}
}
impl MetricsAggregatorContext {
fn build_app_metrics_snapshot(&self) -> obzenflow_core::metrics::AppMetricsSnapshot {
let store = &self.metrics_store;
let mut snapshot = obzenflow_core::metrics::AppMetricsSnapshot::default();
tracing::debug!(
"Exporting metrics: {} stage entries",
store.stage_metrics.len()
);
let mut flow_events_in_total: u64 = 0;
let mut flow_events_out_total: u64 = 0;
let mut flow_errors_total_snapshot: u64 = 0;
let mut total_events_processed_snapshot: u64 = 0;
let mut total_event_loops: u64 = 0;
let mut total_event_loops_with_work: u64 = 0;
for (stage_id, metrics) in &store.stage_metrics {
let events_count = metrics.latest_events_processed_total.unwrap_or(0);
snapshot.event_counts.insert(*stage_id, events_count);
let accumulated_count = metrics.latest_events_accumulated_total.unwrap_or(0);
snapshot
.events_accumulated_total
.insert(*stage_id, accumulated_count);
let emitted_count = metrics.latest_events_emitted_total.unwrap_or(0);
snapshot
.events_emitted_total
.insert(*stage_id, emitted_count);
if let Some(value) = metrics.join_reference_since_last_stream {
if let Some(metadata) = self.stage_metadata.get(stage_id) {
if metadata.stage_type == StageType::Join {
snapshot
.join_reference_since_last_stream
.insert(*stage_id, value);
}
}
}
let stage_errors_total = metrics.latest_errors_total.unwrap_or(0);
snapshot.error_counts.insert(*stage_id, stage_errors_total);
if !metrics.errors_by_kind.is_empty() && stage_errors_total > 0 {
snapshot
.error_counts_by_kind
.insert(*stage_id, metrics.errors_by_kind.clone());
}
if events_count > 0 && metrics.snapshot_p50_ms.is_some() {
let mut percentiles = std::collections::HashMap::new();
if let Some(p50) = metrics.snapshot_p50_ms {
percentiles.insert(Percentile::P50, (p50 * 1_000_000) as f64);
}
if let Some(p90) = metrics.snapshot_p90_ms {
percentiles.insert(Percentile::P90, (p90 * 1_000_000) as f64);
}
if let Some(p95) = metrics.snapshot_p95_ms {
percentiles.insert(Percentile::P95, (p95 * 1_000_000) as f64);
}
if let Some(p99) = metrics.snapshot_p99_ms {
percentiles.insert(Percentile::P99, (p99 * 1_000_000) as f64);
}
if let Some(p999) = metrics.snapshot_p999_ms {
percentiles.insert(Percentile::P999, (p999 * 1_000_000) as f64);
}
let sum_nanos = metrics.processing_time_sum_nanos.unwrap_or(0);
let hist_snapshot = obzenflow_core::metrics::HistogramSnapshot {
count: events_count,
sum: sum_nanos as f64,
min: (metrics.snapshot_p50_ms.unwrap_or(0) * 1_000_000) as f64,
max: (metrics.snapshot_p999_ms.unwrap_or(0) * 1_000_000) as f64,
percentiles,
};
snapshot.processing_times.insert(*stage_id, hist_snapshot);
}
if let Some(in_flight) = metrics.last_in_flight {
snapshot.in_flight.insert(*stage_id, in_flight as f64);
}
if let Some(failures_total) = metrics.last_failures_total {
snapshot.failures_total.insert(*stage_id, failures_total);
}
snapshot
.event_loops_total
.insert(*stage_id, metrics.event_loops_total);
snapshot
.event_loops_with_work_total
.insert(*stage_id, metrics.event_loops_with_work_total);
total_events_processed_snapshot =
total_events_processed_snapshot.saturating_add(events_count);
flow_errors_total_snapshot =
flow_errors_total_snapshot.saturating_add(stage_errors_total);
total_event_loops = total_event_loops.saturating_add(metrics.event_loops_total);
total_event_loops_with_work =
total_event_loops_with_work.saturating_add(metrics.event_loops_with_work_total);
if let Some(metadata) = self.stage_metadata.get(stage_id) {
match metadata.stage_type {
obzenflow_core::event::context::StageType::FiniteSource
| obzenflow_core::event::context::StageType::InfiniteSource => {
flow_events_in_total = flow_events_in_total.saturating_add(events_count);
}
obzenflow_core::event::context::StageType::Sink => {
flow_events_out_total = flow_events_out_total.saturating_add(events_count);
}
_ => {}
}
}
tracing::debug!(
"Exported metrics for {:?}: events={}, errors_total_snapshot={}",
stage_id,
events_count,
stage_errors_total
);
}
if let (Some(first_time), Some(last_time)) = (store.first_event_time, store.last_event_time)
{
let flow_duration = last_time.duration_since(first_time);
let flow_metrics = obzenflow_core::metrics::FlowMetricsSnapshot {
flow_duration: MetricsDuration::from(flow_duration),
total_events_processed: total_events_processed_snapshot,
events_in: flow_events_in_total,
events_out: flow_events_out_total,
errors_total: flow_errors_total_snapshot,
event_loops_total: total_event_loops,
event_loops_with_work_total: total_event_loops_with_work,
};
snapshot.flow_metrics = Some(flow_metrics);
}
snapshot.stage_metadata = self.stage_metadata.clone();
snapshot.sink_operation_failures = store
.sink_operation_failures
.iter()
.map(|((stage_id, phase, error_kind), count)| {
obzenflow_core::metrics::SinkOperationFailureMetric {
stage_id: *stage_id,
phase: *phase,
error_kind: error_kind.clone(),
count: *count,
}
})
.collect();
snapshot.circuit_breaker_state = store.circuit_breaker_state.clone();
snapshot.circuit_breaker_rejection_rate = store.circuit_breaker_rejection_rate.clone();
snapshot.circuit_breaker_consecutive_failures =
store.circuit_breaker_consecutive_failures.clone();
snapshot.circuit_breaker_requests_total = store.circuit_breaker_requests_total.clone();
snapshot.circuit_breaker_rejections_total = store.circuit_breaker_rejections_total.clone();
snapshot.circuit_breaker_opened_total = store.circuit_breaker_opened_total.clone();
snapshot.circuit_breaker_successes_total = store.circuit_breaker_successes_total.clone();
snapshot.circuit_breaker_failures_total = store.circuit_breaker_failures_total.clone();
snapshot.circuit_breaker_slow_total = store.circuit_breaker_slow_total.clone();
snapshot.circuit_breaker_time_in_state_seconds_total =
store.circuit_breaker_time_in_state_seconds_total.clone();
snapshot.circuit_breaker_state_transitions_total =
store.circuit_breaker_state_transitions_total.clone();
snapshot.rate_limiter_utilization = store
.rate_limiter_bucket_capacity
.iter()
.filter_map(|(stage_id, capacity)| {
if *capacity <= 0.0 {
return None;
}
let tokens = store.rate_limiter_bucket_tokens.get(stage_id)?;
let utilization = 1.0 - (*tokens / *capacity);
Some((*stage_id, utilization.clamp(0.0, 1.0)))
})
.collect();
snapshot.rate_limiter_events_total = store.rate_limiter_events_total.clone();
snapshot.rate_limiter_delayed_total = store.rate_limiter_delayed_total.clone();
snapshot.rate_limiter_tokens_consumed_total =
store.rate_limiter_tokens_consumed_total.clone();
snapshot.rate_limiter_delay_seconds_total = store.rate_limiter_delay_seconds_total.clone();
snapshot.rate_limiter_bucket_tokens = store.rate_limiter_bucket_tokens.clone();
snapshot.rate_limiter_bucket_capacity = store.rate_limiter_bucket_capacity.clone();
snapshot.backpressure_bypass_enabled =
crate::backpressure::BackpressureWriter::is_bypass_enabled();
if let Some(registry) = &self.backpressure_registry {
let bp = registry.metrics_snapshot();
snapshot.backpressure_window = bp.edge_window;
snapshot.backpressure_in_flight = bp.edge_in_flight;
snapshot.backpressure_credits = bp.edge_credits;
snapshot.backpressure_blocked = bp
.stage_blocked
.into_iter()
.map(|(stage_id, blocked)| (stage_id, if blocked { 1.0 } else { 0.0 }))
.collect();
snapshot.backpressure_min_reader_seq = bp.stage_min_reader_seq;
snapshot.backpressure_writer_seq = bp.stage_writer_seq;
snapshot.backpressure_wait_seconds_total = bp
.stage_wait_nanos_total
.into_iter()
.map(|(stage_id, nanos)| (stage_id, nanos as f64 / 1_000_000_000.0))
.collect();
}
snapshot.edge_liveness_state = store.edge_liveness_state.clone();
snapshot.contract_metrics = store.contract_metrics.clone();
let mut http_surface_metrics: Vec<_> =
store.http_surface_metrics.values().cloned().collect();
http_surface_metrics.sort_by(|a, b| {
(
a.surface_name.as_str(),
a.path.as_str(),
a.method.as_str(),
a.status_class.as_str(),
)
.cmp(&(
b.surface_name.as_str(),
b.path.as_str(),
b.method.as_str(),
b.status_class.as_str(),
))
});
snapshot.http_surface_metrics = http_surface_metrics;
snapshot.ingestion_refusal_totals = store.ingestion_refusals_total.clone();
snapshot.http_pull_metrics = store.http_pull_metrics.clone();
snapshot.ai_chunking_metrics = store.ai_chunking_metrics.clone();
snapshot.stage_lifecycle_states = store.stage_lifecycle_states.clone();
snapshot.pipeline_state = store.pipeline_state.clone();
let now = std::time::Instant::now();
let now_utc = chrono::Utc::now();
for (stage_id, metrics) in &store.stage_metrics {
if let Some(first_time) = metrics.first_event_time {
let elapsed_since_first = now.duration_since(first_time);
let first_datetime =
now_utc - chrono::Duration::from_std(elapsed_since_first).unwrap_or_default();
snapshot
.stage_first_event_time
.insert(*stage_id, first_datetime);
}
if let Some(last_time) = metrics.last_event_time {
let elapsed_since_last = now.duration_since(last_time);
let last_datetime =
now_utc - chrono::Duration::from_std(elapsed_since_last).unwrap_or_default();
snapshot
.stage_last_event_time
.insert(*stage_id, last_datetime);
}
if let Some(seq) = store.stage_vector_clocks.get(stage_id) {
snapshot.stage_vector_clocks.insert(*stage_id, *seq);
}
}
use obzenflow_core::metrics::{CompositeMemberHealth, CompositePortTraffic};
snapshot.composite_port_traffic = self
.composite_boundaries
.iter()
.flat_map(|boundary| CompositePortTraffic::project(boundary, store))
.collect();
snapshot.composite_port_traffic.sort_by(|left, right| {
(
left.composite.as_str(),
left.direction.as_str(),
left.port.as_str(),
)
.cmp(&(
right.composite.as_str(),
right.direction.as_str(),
right.port.as_str(),
))
});
snapshot.composite_member_health = self
.composite_boundaries
.iter()
.map(|boundary| CompositeMemberHealth::project(boundary, store))
.collect();
snapshot
.composite_member_health
.sort_by(|left, right| left.composite.cmp(&right.composite));
snapshot.composite_boundary_durations = self.composite_durations.histograms();
snapshot.composite_boundary_duration_invalid = self.composite_durations.invalid_evidence();
use obzenflow_core::metrics::CompositeContract;
snapshot.composite_contracts = self
.composite_boundaries
.iter()
.flat_map(|b| CompositeContract::project(b, &snapshot.contract_metrics))
.collect();
snapshot.composite_contracts.sort_by(|left, right| {
(
left.composite.as_str(),
left.direction.as_str(),
left.port.as_str(),
left.peer,
left.selected_event_type.as_ref(),
left.feed_role.map(|role| role.as_str()),
)
.cmp(&(
right.composite.as_str(),
right.direction.as_str(),
right.port.as_str(),
right.peer,
right.selected_event_type.as_ref(),
right.feed_role.map(|role| role.as_str()),
))
});
snapshot
}
}
impl MetricsStore {
fn fold_sink_operation_failure(&mut self, failure: &SinkOperationFailed) {
*self
.sink_operation_failures
.entry((failure.stage_id, failure.phase, failure.kind.clone()))
.or_insert(0) += 1;
}
fn fold_http_pull_snapshot(&mut self, stage_id: StageId, snapshot: &HttpPullTelemetry) {
let entry = self.http_pull_metrics.entry(stage_id).or_default();
entry.state = snapshot.state.clone();
entry.wait_reason = snapshot.wait_reason.clone();
entry.next_wake_unix_secs = snapshot.next_wake_unix_secs;
entry.last_success_unix_secs = match (
entry.last_success_unix_secs,
snapshot.last_success_unix_secs,
) {
(Some(a), Some(b)) => Some(a.max(b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
};
entry.requests_total = entry.requests_total.max(snapshot.requests_total);
entry.responses_2xx = entry.responses_2xx.max(snapshot.responses_2xx);
entry.responses_4xx = entry.responses_4xx.max(snapshot.responses_4xx);
entry.responses_5xx = entry.responses_5xx.max(snapshot.responses_5xx);
entry.rate_limited_total = entry.rate_limited_total.max(snapshot.rate_limited_total);
entry.retries_total = entry.retries_total.max(snapshot.retries_total);
entry.events_decoded_total = entry
.events_decoded_total
.max(snapshot.events_decoded_total);
entry.wait_seconds_rate_limit = entry
.wait_seconds_rate_limit
.max(snapshot.wait_seconds_rate_limit);
entry.wait_seconds_poll_interval = entry
.wait_seconds_poll_interval
.max(snapshot.wait_seconds_poll_interval);
entry.wait_seconds_backoff = entry
.wait_seconds_backoff
.max(snapshot.wait_seconds_backoff);
}
pub fn mark_known_stages_completed<I>(&mut self, stage_ids: I)
where
I: IntoIterator<Item = StageId>,
{
for stage_id in stage_ids {
let has_terminal_state = self
.stage_lifecycle_states
.get(&(stage_id, "completed".to_string()))
.copied()
.unwrap_or(false)
|| self
.stage_lifecycle_states
.get(&(stage_id, "failed".to_string()))
.copied()
.unwrap_or(false)
|| self
.stage_lifecycle_states
.get(&(stage_id, "cancelled".to_string()))
.copied()
.unwrap_or(false);
if !has_terminal_state {
self.stage_lifecycle_states
.insert((stage_id, "completed".to_string()), true);
}
}
}
pub fn all_stages_terminal(&self, stage_metadata: &HashMap<StageId, StageMetadata>) -> bool {
stage_metadata.keys().all(|stage_id| {
self.stage_lifecycle_states
.get(&(*stage_id, "completed".to_string()))
.copied()
.unwrap_or(false)
|| self
.stage_lifecycle_states
.get(&(*stage_id, "failed".to_string()))
.copied()
.unwrap_or(false)
|| self
.stage_lifecycle_states
.get(&(*stage_id, "cancelled".to_string()))
.copied()
.unwrap_or(false)
})
}
pub fn pipeline_terminal(&self) -> bool {
matches!(
self.pipeline_state.as_str(),
"completed"
| "failed"
| "cancelled"
| "drained"
| "all_stages_completed"
| "stop_requested"
)
}
fn record_circuit_breaker_transition(&mut self, stage_id: StageId, state: &str) {
let Some(next_state) = normalize_circuit_breaker_state_label(state) else {
return;
};
let from_state = self
.circuit_breaker_last_state
.get(&stage_id)
.map(String::as_str)
.unwrap_or("closed");
if from_state != next_state {
let transition_key = (stage_id, from_state.to_string(), next_state.to_string());
let count = self
.circuit_breaker_state_transitions_total
.entry(transition_key)
.or_insert(0);
*count = (*count).saturating_add(1);
}
self.circuit_breaker_last_state
.insert(stage_id, next_state.to_string());
}
fn update_control_metrics_from_runtime_context(
&mut self,
stage_id: StageId,
runtime_ctx: &obzenflow_core::event::context::RuntimeContext,
) {
let cb_present = runtime_ctx.cb_requests_total > 0
|| runtime_ctx.cb_rejections_total > 0
|| runtime_ctx.cb_opened_total > 0
|| runtime_ctx.cb_successes_total > 0
|| runtime_ctx.cb_failures_total > 0
|| runtime_ctx.cb_slow_total > 0
|| runtime_ctx.cb_time_closed_seconds > 0.0
|| runtime_ctx.cb_time_open_seconds > 0.0
|| runtime_ctx.cb_time_half_open_seconds > 0.0;
if cb_present {
let total = self
.circuit_breaker_requests_total
.entry(stage_id)
.or_insert(0);
*total = (*total).max(runtime_ctx.cb_requests_total);
let rejected = self
.circuit_breaker_rejections_total
.entry(stage_id)
.or_insert(0);
*rejected = (*rejected).max(runtime_ctx.cb_rejections_total);
let opened = self
.circuit_breaker_opened_total
.entry(stage_id)
.or_insert(0);
*opened = (*opened).max(runtime_ctx.cb_opened_total);
let successes = self
.circuit_breaker_successes_total
.entry(stage_id)
.or_insert(0);
*successes = (*successes).max(runtime_ctx.cb_successes_total);
let failures = self
.circuit_breaker_failures_total
.entry(stage_id)
.or_insert(0);
*failures = (*failures).max(runtime_ctx.cb_failures_total);
let slow = self.circuit_breaker_slow_total.entry(stage_id).or_insert(0);
*slow = (*slow).max(runtime_ctx.cb_slow_total);
self.circuit_breaker_state
.insert(stage_id, runtime_ctx.cb_state);
self.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "closed".to_string()))
.and_modify(|v| *v = (*v).max(runtime_ctx.cb_time_closed_seconds))
.or_insert(runtime_ctx.cb_time_closed_seconds);
self.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "open".to_string()))
.and_modify(|v| *v = (*v).max(runtime_ctx.cb_time_open_seconds))
.or_insert(runtime_ctx.cb_time_open_seconds);
self.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "half_open".to_string()))
.and_modify(|v| *v = (*v).max(runtime_ctx.cb_time_half_open_seconds))
.or_insert(runtime_ctx.cb_time_half_open_seconds);
}
let rl_present = runtime_ctx.rl_events_total > 0
|| runtime_ctx.rl_delayed_total > 0
|| runtime_ctx.rl_tokens_consumed_total > 0.0
|| runtime_ctx.rl_delay_seconds_total > 0.0;
if rl_present {
let events_total = self.rate_limiter_events_total.entry(stage_id).or_insert(0);
*events_total = (*events_total).max(runtime_ctx.rl_events_total);
let delayed_total = self.rate_limiter_delayed_total.entry(stage_id).or_insert(0);
*delayed_total = (*delayed_total).max(runtime_ctx.rl_delayed_total);
let tokens_consumed = self
.rate_limiter_tokens_consumed_total
.entry(stage_id)
.or_insert(0.0);
*tokens_consumed = (*tokens_consumed).max(runtime_ctx.rl_tokens_consumed_total);
let delay_seconds = self
.rate_limiter_delay_seconds_total
.entry(stage_id)
.or_insert(0.0);
*delay_seconds = (*delay_seconds).max(runtime_ctx.rl_delay_seconds_total);
self.rate_limiter_bucket_tokens
.insert(stage_id, runtime_ctx.rl_bucket_tokens);
self.rate_limiter_bucket_capacity
.insert(stage_id, runtime_ctx.rl_bucket_capacity);
}
}
}
fn normalize_circuit_breaker_state_label(state: &str) -> Option<&'static str> {
let state = state.trim();
if state.eq_ignore_ascii_case("closed") {
Some("closed")
} else if state.eq_ignore_ascii_case("open") {
Some("open")
} else {
let normalized = state.to_ascii_lowercase();
match normalized.as_str() {
"halfopen" | "half_open" | "half-open" => Some("half_open"),
_ => None,
}
}
}
fn edge_liveness_state_gauge_value(state: &obzenflow_core::event::EdgeLivenessState) -> f64 {
match state {
obzenflow_core::event::EdgeLivenessState::Healthy => 1.0,
obzenflow_core::event::EdgeLivenessState::Idle => 0.5,
obzenflow_core::event::EdgeLivenessState::Suspect => 0.25,
obzenflow_core::event::EdgeLivenessState::Stalled => 0.0,
obzenflow_core::event::EdgeLivenessState::Recovered => 1.0,
}
}
#[async_trait::async_trait]
impl FsmAction for MetricsAggregatorAction {
type Context = MetricsAggregatorContext;
async fn execute(&self, ctx: &mut Self::Context) -> Result<(), obzenflow_fsm::FsmError> {
match self {
MetricsAggregatorAction::Initialize => {
tracing::info!("Metrics aggregator initialized");
Ok(())
}
MetricsAggregatorAction::ProcessSystemEvent { envelope } => {
tracing::trace!(
event_id = %envelope.event.id(),
event_type = envelope.event.event_type_name(),
"Metrics aggregator ProcessSystemEvent action"
);
let known_stage_ids = ctx.stage_metadata.keys().copied().collect::<Vec<_>>();
let store = &mut ctx.metrics_store;
if let Some(system_id) = envelope.event.writer_id.as_system() {
let writer_key = envelope.event.writer_id.to_string();
let seq = envelope.vector_clock.get(&writer_key);
let entry = store.system_vector_clocks.entry(*system_id).or_insert(0);
*entry = (*entry).max(seq);
}
match &envelope.event.event {
obzenflow_core::event::SystemEventType::StageLifecycle { stage_id, event } => {
match event {
obzenflow_core::event::StageLifecycleEvent::Running => {
store
.stage_lifecycle_states
.insert((*stage_id, "running".to_string()), true);
tracing::debug!("Stage {:?} transitioned to running", stage_id);
}
obzenflow_core::event::StageLifecycleEvent::Completed { .. } => {
store
.stage_lifecycle_states
.insert((*stage_id, "completed".to_string()), true);
tracing::debug!("Stage {:?} transitioned to completed", stage_id);
}
obzenflow_core::event::StageLifecycleEvent::Cancelled { .. } => {
store
.stage_lifecycle_states
.insert((*stage_id, "cancelled".to_string()), true);
tracing::debug!("Stage {:?} transitioned to cancelled", stage_id);
}
obzenflow_core::event::StageLifecycleEvent::Failed { .. } => {
store
.stage_lifecycle_states
.insert((*stage_id, "failed".to_string()), true);
tracing::debug!("Stage {:?} transitioned to failed", stage_id);
}
_ => {} }
}
obzenflow_core::event::SystemEventType::PipelineLifecycle(event) => {
match event {
obzenflow_core::event::PipelineLifecycleEvent::StopRequested {
..
} => {
if store.pipeline_state.is_empty() {
store.pipeline_state = "stop_requested".to_string();
}
tracing::info!("Pipeline: stop requested (metrics view)");
}
obzenflow_core::event::PipelineLifecycleEvent::AllStagesCompleted {
..
} => {
store.mark_known_stages_completed(known_stage_ids);
if store.pipeline_state.is_empty() {
store.pipeline_state = "all_stages_completed".to_string();
}
tracing::info!("Pipeline: all stages completed (metrics view)");
}
obzenflow_core::event::PipelineLifecycleEvent::Completed { .. } => {
if store.pipeline_state != "failed" {
store.pipeline_state = "completed".to_string();
tracing::info!("Pipeline: completed (metrics view)");
} else {
tracing::info!(
"Pipeline: completed event observed after failed; \
keeping failed as terminal state (metrics view)"
);
}
}
obzenflow_core::event::PipelineLifecycleEvent::Cancelled { .. } => {
if store.pipeline_state != "failed" {
store.pipeline_state = "cancelled".to_string();
tracing::info!("Pipeline: cancelled (metrics view)");
} else {
tracing::info!(
"Pipeline: cancelled event observed after failed; \
keeping failed as terminal state (metrics view)"
);
}
}
obzenflow_core::event::PipelineLifecycleEvent::Failed { .. } => {
if store.pipeline_state != "failed" {
store.pipeline_state = "failed".to_string();
tracing::info!("Pipeline: failed (metrics view)");
}
}
obzenflow_core::event::PipelineLifecycleEvent::Drained => {
match store.pipeline_state.as_str() {
"failed" | "completed" => {
tracing::info!(
"Pipeline: drained event observed after terminal outcome; \
keeping {} as terminal state (metrics view)",
store.pipeline_state
);
}
_ => {
store.pipeline_state = "drained".to_string();
tracing::info!("Pipeline: drained (metrics view)");
}
}
}
_ => {} }
}
obzenflow_core::event::SystemEventType::ContractResult {
upstream,
reader,
selected_event_type,
feed_role,
contract_name,
status,
cause,
reader_seq,
advertised_writer_seq,
} => {
let edge_key = ContractMetricEdgeKey {
upstream: *upstream,
downstream: *reader,
contract: contract_name.clone(),
selected_event_type: selected_event_type.clone(),
feed_role: *feed_role,
};
let result_key = ContractMetricResultKey {
edge: edge_key.clone(),
status: *status,
};
let counter = store
.contract_metrics
.results_total
.entry(result_key)
.or_insert(0);
*counter = (*counter).saturating_add(1);
if let Some(cause) = cause {
let violation_key = ContractMetricViolationKey {
edge: edge_key.clone(),
cause: ContractViolationCauseLabel::from(cause.clone()),
};
let counter = store
.contract_metrics
.violations_total
.entry(violation_key)
.or_insert(0);
*counter = (*counter).saturating_add(1);
}
if let Some(seq) = reader_seq {
let gauge = store
.contract_metrics
.reader_seq
.entry(edge_key.clone())
.or_insert(0);
*gauge = (*gauge).max(seq.0);
}
if let Some(seq) = advertised_writer_seq {
let gauge = store
.contract_metrics
.advertised_writer_seq
.entry(edge_key)
.or_insert(0);
*gauge = (*gauge).max(seq.0);
}
}
obzenflow_core::event::SystemEventType::EdgeLiveness {
upstream,
reader,
state,
..
} => {
store
.edge_liveness_state
.insert((*upstream, *reader), edge_liveness_state_gauge_value(state));
}
obzenflow_core::event::SystemEventType::HttpSurfaceSnapshot { snapshot } => {
for route in &snapshot.routes {
let key = (
route.surface_name.clone(),
route.method,
route.path.clone(),
route.status_class.clone(),
);
let entry = store
.http_surface_metrics
.entry(key)
.or_insert_with(|| route.clone());
entry.requests_total = entry.requests_total.max(route.requests_total);
entry.request_duration_ms_total = entry
.request_duration_ms_total
.max(route.request_duration_ms_total);
entry.request_bytes_total =
entry.request_bytes_total.max(route.request_bytes_total);
entry.response_bytes_total =
entry.response_bytes_total.max(route.response_bytes_total);
}
}
obzenflow_core::event::SystemEventType::IngressRefusal {
ingress_key,
reason,
event_count,
..
} => {
let key = (ingress_key.clone(), reason.as_str().to_string());
let entry = store.ingestion_refusals_total.entry(key).or_insert(0);
*entry = entry.saturating_add(*event_count);
}
_ => {} }
Ok(())
}
MetricsAggregatorAction::UpdateMetrics {
envelope,
journal_kind,
journal_stage,
} => {
tracing::trace!(
event_id = %envelope.event.id(),
event_type = envelope.event.event_type(),
"Metrics aggregator UpdateMetrics action"
);
let event = &envelope.event;
let stage_id = event.flow_context.stage_id;
observe_live_composite_duration(
*journal_kind,
&mut ctx.composite_durations,
&ctx.composite_boundaries,
*journal_stage,
event,
);
let store = &mut ctx.metrics_store;
store.last_event_id = Some(event.id);
let writer_id = *event.writer_id();
let writer_key = writer_id.to_string();
let seq = envelope.vector_clock.get(&writer_key);
let entry = store.stage_vector_clocks.entry(stage_id).or_insert(0);
*entry = (*entry).max(seq);
if let Some(meta) = ctx.stage_metadata.get_mut(&stage_id) {
if meta.flow_id.is_none() {
if let Ok(flow_id) = FlowId::from_str(event.flow_context.flow_id.as_str()) {
meta.flow_id = Some(flow_id);
}
}
if meta.reference_mode.is_none() && meta.stage_type == StageType::Join {
if let Some(runtime_ctx) = &event.runtime_context {
if let Some(mode) =
infer_join_reference_mode_from_fsm_state(&runtime_ctx.fsm_state)
{
meta.reference_mode = Some(mode.to_string());
}
}
}
}
if event.is_system() {
return Ok(());
}
if let ChainEventContent::Observability(ObservabilityPayload::Metrics(
MetricsLifecycle::HttpPullSnapshot { snapshot },
)) = &event.content
{
store.fold_http_pull_snapshot(stage_id, snapshot);
} else if let ChainEventContent::Observability(ObservabilityPayload::Metrics(
MetricsLifecycle::Custom { name, value, .. },
)) = &event.content
{
if name == "ai_chunking.snapshot" {
match serde_json::from_value::<
obzenflow_core::event::observability::AiChunkingSnapshot,
>(value.clone())
{
Ok(snapshot) => {
let entry = store.ai_chunking_metrics.entry(stage_id).or_default();
entry.jobs_total = entry.jobs_total.saturating_add(1);
entry.input_items_total = entry
.input_items_total
.saturating_add(snapshot.input_items_total as u64);
entry.planned_items_total = entry
.planned_items_total
.saturating_add(snapshot.planned_items_total as u64);
entry.excluded_items_total = entry
.excluded_items_total
.saturating_add(snapshot.excluded_items_total as u64);
entry.chunks_emitted_total = entry
.chunks_emitted_total
.saturating_add(snapshot.chunk_count as u64);
entry.rerender_attempts_total = entry
.rerender_attempts_total
.saturating_add(snapshot.rerender_attempts_total);
entry.max_depth_reached = entry
.max_depth_reached
.max(snapshot.max_decomposition_depth_reached);
entry.budget_overhead_tokens = snapshot.budget_overhead_tokens;
}
Err(e) => tracing::warn!(
error = %e,
"Failed to decode ai_chunking.snapshot payload; ignoring"
),
}
}
}
if let ChainEventContent::Data {
event_type,
payload,
} = &event.content
{
if SinkOperationFailed::event_type_matches(event_type) {
match serde_json::from_value::<SinkOperationFailed>(payload.clone()) {
Ok(failure) => store.fold_sink_operation_failure(&failure),
Err(error) => tracing::warn!(
%error,
"Failed to decode SinkOperationFailed metric fact; ignoring"
),
}
}
}
if let ChainEventContent::Observability(ObservabilityPayload::Middleware(
MiddlewareLifecycle::CircuitBreaker(cb),
)) = &event.content
{
match cb {
CircuitBreakerEvent::Opened {
error_rate: _,
failure_count,
trigger,
..
} => {
store.record_circuit_breaker_transition(stage_id, "open");
store.circuit_breaker_state.insert(stage_id, 1.0);
if matches!(trigger, CircuitBreakerOpenTrigger::ConsecutiveFailures) {
store
.circuit_breaker_consecutive_failures
.insert(stage_id, *failure_count as f64);
} else {
store.circuit_breaker_consecutive_failures.remove(&stage_id);
}
}
CircuitBreakerEvent::Closed { .. } => {
store.record_circuit_breaker_transition(stage_id, "closed");
store.circuit_breaker_state.insert(stage_id, 0.0);
store
.circuit_breaker_consecutive_failures
.insert(stage_id, 0.0);
}
CircuitBreakerEvent::HalfOpen { .. } => {
store.record_circuit_breaker_transition(stage_id, "half_open");
store.circuit_breaker_state.insert(stage_id, 0.5);
}
CircuitBreakerEvent::Summary {
requests_processed: _,
requests_rejected: _,
state,
consecutive_failures,
rejection_rate,
successes_total,
failures_total,
opened_total,
time_in_closed_seconds,
time_in_open_seconds,
time_in_half_open_seconds,
..
} => {
store
.circuit_breaker_rejection_rate
.insert(stage_id, *rejection_rate);
store
.circuit_breaker_consecutive_failures
.insert(stage_id, *consecutive_failures as f64);
let opened = store
.circuit_breaker_opened_total
.entry(stage_id)
.or_insert(0);
*opened = (*opened).max(*opened_total);
let successes = store
.circuit_breaker_successes_total
.entry(stage_id)
.or_insert(0);
*successes = (*successes).max(*successes_total);
let failures = store
.circuit_breaker_failures_total
.entry(stage_id)
.or_insert(0);
*failures = (*failures).max(*failures_total);
store
.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "closed".to_string()))
.and_modify(|v| *v = (*v).max(*time_in_closed_seconds))
.or_insert(*time_in_closed_seconds);
store
.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "open".to_string()))
.and_modify(|v| *v = (*v).max(*time_in_open_seconds))
.or_insert(*time_in_open_seconds);
store
.circuit_breaker_time_in_state_seconds_total
.entry((stage_id, "half_open".to_string()))
.and_modify(|v| *v = (*v).max(*time_in_half_open_seconds))
.or_insert(*time_in_half_open_seconds);
let state_norm = state.to_ascii_lowercase();
let state_value = match state_norm.as_str() {
"closed" => Some(0.0),
"open" => Some(1.0),
"halfopen" | "half_open" | "half-open" => Some(0.5),
_ => None,
};
if let Some(val) = state_value {
store.circuit_breaker_state.insert(stage_id, val);
}
}
_ => {}
}
}
let now = std::time::Instant::now();
if store.first_event_time.is_none() {
store.first_event_time = Some(now);
store.flow_start_time = Some(now);
}
store.last_event_time = Some(now);
if event.is_data() || event.is_delivery() {
store.total_events_processed += 1;
}
let stage_id = event.flow_context.stage_id;
{
let metrics = store.stage_metrics.entry(stage_id).or_default();
let now = std::time::Instant::now();
if metrics.first_event_time.is_none() {
metrics.first_event_time = Some(now);
}
metrics.last_event_time = Some(now);
if let Some(runtime_ctx) = &event.runtime_context {
tracing::trace!(
"Runtime context for {:?}: in_flight={}, fsm_state={}",
stage_id,
runtime_ctx.in_flight,
runtime_ctx.fsm_state
);
metrics.merge_runtime_context(runtime_ctx);
}
if let ProcessingStatus::Error { kind, .. } = &event.processing_info.status {
let key = kind.clone().unwrap_or(ErrorKind::Unknown);
*metrics.errors_by_kind.entry(key).or_insert(0) += 1;
}
}
if let Some(runtime_ctx) = &event.runtime_context {
store.update_control_metrics_from_runtime_context(stage_id, runtime_ctx);
}
Ok(())
}
MetricsAggregatorAction::ExportMetrics => {
tracing::debug!("ExportMetrics action triggered");
for (stage_id, data_journal) in &ctx.stage_data_journals {
let error_journal = ctx.stage_error_journals.get(stage_id);
if let Some(snapshot) = tail_read::read_stage_metrics_from_tail(
data_journal,
error_journal,
*stage_id,
)
.await
{
let metrics = ctx
.metrics_store
.stage_metrics
.entry(*stage_id)
.or_default();
metrics.latest_events_processed_total = Some(
metrics
.latest_events_processed_total
.unwrap_or(0)
.max(snapshot.events_processed_total),
);
metrics.latest_events_accumulated_total = Some(
metrics
.latest_events_accumulated_total
.unwrap_or(0)
.max(snapshot.events_accumulated_total),
);
metrics.latest_events_emitted_total = Some(
metrics
.latest_events_emitted_total
.unwrap_or(0)
.max(snapshot.events_emitted_total),
);
metrics.latest_errors_total = Some(
metrics
.latest_errors_total
.unwrap_or(0)
.max(snapshot.errors_total),
);
metrics.errors_by_kind = snapshot.errors_by_kind.clone();
metrics.last_in_flight = Some(snapshot.in_flight);
metrics.snapshot_p50_ms = Some(snapshot.recent_p50_ms);
metrics.snapshot_p90_ms = Some(snapshot.recent_p90_ms);
metrics.snapshot_p95_ms = Some(snapshot.recent_p95_ms);
metrics.snapshot_p99_ms = Some(snapshot.recent_p99_ms);
metrics.snapshot_p999_ms = Some(snapshot.recent_p999_ms);
metrics.processing_time_sum_nanos =
Some(snapshot.processing_time_sum_nanos);
}
if let Some(runtime_ctx) =
tail_read::read_latest_runtime_context_for_stage(data_journal, *stage_id)
.await
{
if let Some(meta) = ctx.stage_metadata.get_mut(stage_id) {
if meta.reference_mode.is_none() && meta.stage_type == StageType::Join {
if let Some(mode) =
infer_join_reference_mode_from_fsm_state(&runtime_ctx.fsm_state)
{
meta.reference_mode = Some(mode.to_string());
}
}
}
if let Some(metrics) = ctx.metrics_store.stage_metrics.get_mut(stage_id) {
metrics.merge_runtime_context(&runtime_ctx);
}
ctx.metrics_store
.update_control_metrics_from_runtime_context(*stage_id, &runtime_ctx);
}
if let Some(error_journal) = error_journal {
if let Some(runtime_ctx) = tail_read::read_latest_runtime_context_for_stage(
error_journal,
*stage_id,
)
.await
{
if let Some(metrics) = ctx.metrics_store.stage_metrics.get_mut(stage_id)
{
metrics.merge_runtime_context(&runtime_ctx);
}
ctx.metrics_store
.update_control_metrics_from_runtime_context(
*stage_id,
&runtime_ctx,
);
}
}
}
if let Some(exporter) = &ctx.exporter {
let snapshot = ctx.build_app_metrics_snapshot();
tracing::debug!("Pushing metrics snapshot to exporter");
if let Err(e) = exporter.update_app_metrics(snapshot) {
tracing::warn!("Failed to export metrics: {}", e);
} else {
tracing::debug!("Successfully exported metrics");
}
}
let mut clocks: std::collections::BTreeMap<String, u64> =
std::collections::BTreeMap::new();
for (stage_id, seq) in &ctx.metrics_store.stage_vector_clocks {
clocks.insert(WriterId::from(*stage_id).to_string(), *seq);
}
for (system_id, seq) in &ctx.metrics_store.system_vector_clocks {
clocks.insert(WriterId::from(*system_id).to_string(), *seq);
}
let export_event = obzenflow_core::event::SystemEvent::new(
WriterId::from(ctx.system_id),
obzenflow_core::event::SystemEventType::MetricsCoordination(
obzenflow_core::event::MetricsCoordinationEvent::Exported {
watermark: obzenflow_core::event::vector_clock::VectorClock { clocks },
},
),
);
if let Err(e) = ctx.system_journal.append(export_event, None).await {
tracing::warn!(
journal_error = %e,
"Failed to publish metrics watermark event; continuing without system journal entry"
);
}
Ok(())
}
MetricsAggregatorAction::PublishDrainComplete { last_event_id } => {
let system_writer_id = WriterId::from(ctx.system_id);
let mut payload = serde_json::json!({});
if let Some(id) = last_event_id {
payload["last_event_id"] = serde_json::json!(id.to_string());
}
let drain_event = obzenflow_core::event::SystemEvent::new(
system_writer_id,
obzenflow_core::event::SystemEventType::MetricsCoordination(
obzenflow_core::event::MetricsCoordinationEvent::Drained,
),
);
ctx.system_journal
.append(drain_event, None)
.await
.map(|_| ())
.map_err(|e| {
obzenflow_fsm::FsmError::HandlerError(format!(
"Failed to publish drain complete event: {e}"
))
})?;
tracing::info!(
"Published metrics drain complete event (last_event_id={:?})",
last_event_id
);
Ok(())
}
}
}
}
pub type MetricsAggregatorFsm = StateMachine<
MetricsAggregatorState,
MetricsAggregatorEvent,
MetricsAggregatorContext,
MetricsAggregatorAction,
>;
pub fn build_metrics_aggregator_fsm() -> MetricsAggregatorFsm {
fsm! {
state: MetricsAggregatorState;
event: MetricsAggregatorEvent;
context: MetricsAggregatorContext;
action: MetricsAggregatorAction;
initial: MetricsAggregatorState::Initializing;
state MetricsAggregatorState::Initializing {
on MetricsAggregatorEvent::StartRunning => |_state: &MetricsAggregatorState, _event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
Box::pin(async move {
Ok(Transition {
next_state: MetricsAggregatorState::Running,
actions: vec![MetricsAggregatorAction::Initialize],
})
})
};
on MetricsAggregatorEvent::Error => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
Box::pin(async move {
let error = match event {
MetricsAggregatorEvent::Error(err) => err.clone(),
_ => {
return Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for Error handler".to_string(),
));
}
};
tracing::error!(error = %error, "Metrics aggregator encountered error");
Ok(Transition {
next_state: MetricsAggregatorState::Failed { error },
actions: vec![],
})
})
};
}
state MetricsAggregatorState::Running {
on MetricsAggregatorEvent::StartDraining => |_state: &MetricsAggregatorState, _event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
Box::pin(async move {
Ok(Transition {
next_state: MetricsAggregatorState::Draining,
actions: vec![],
})
})
};
on MetricsAggregatorEvent::ProcessSystemEvent => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
let last_event_id = ctx.metrics_store.last_event_id;
Box::pin(async move {
match event {
MetricsAggregatorEvent::ProcessSystemEvent { envelope } => {
let pipeline_event = match &envelope.event.event {
obzenflow_core::event::SystemEventType::PipelineLifecycle(event) => {
Some(event)
}
_ => None,
};
let should_drain = matches!(
pipeline_event,
Some(
obzenflow_core::event::PipelineLifecycleEvent::Draining { .. }
| obzenflow_core::event::PipelineLifecycleEvent::AllStagesCompleted { .. }
)
);
let should_finalize = matches!(
pipeline_event,
Some(
obzenflow_core::event::PipelineLifecycleEvent::Completed { .. }
| obzenflow_core::event::PipelineLifecycleEvent::Failed { .. }
| obzenflow_core::event::PipelineLifecycleEvent::Drained
)
);
if should_finalize {
let publish_last_event_id = last_event_id;
return Ok(Transition {
next_state: MetricsAggregatorState::Drained { last_event_id },
actions: vec![
MetricsAggregatorAction::ProcessSystemEvent {
envelope: envelope.clone(),
},
MetricsAggregatorAction::ExportMetrics,
MetricsAggregatorAction::PublishDrainComplete {
last_event_id: publish_last_event_id,
},
],
});
}
if should_drain {
return Ok(Transition {
next_state: MetricsAggregatorState::Draining,
actions: vec![
MetricsAggregatorAction::ProcessSystemEvent {
envelope: envelope.clone(),
},
MetricsAggregatorAction::ExportMetrics,
],
});
}
Ok(Transition {
next_state: MetricsAggregatorState::Running,
actions: vec![MetricsAggregatorAction::ProcessSystemEvent {
envelope: envelope.clone(),
}],
})
}
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for ProcessSystemEvent handler".to_string(),
)),
}
})
};
on MetricsAggregatorEvent::ProcessBatch => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
Box::pin(async move {
match event {
MetricsAggregatorEvent::ProcessBatch {
events,
journal_kind,
journal_stage,
} => {
let actions = events
.iter()
.cloned()
.map(|envelope| MetricsAggregatorAction::UpdateMetrics {
envelope: Box::new(envelope),
journal_kind,
journal_stage,
})
.collect::<Vec<_>>();
Ok(Transition {
next_state: MetricsAggregatorState::Running,
actions,
})
}
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for ProcessBatch handler".to_string(),
)),
}
})
};
on MetricsAggregatorEvent::ExportMetrics => |_state: &MetricsAggregatorState, _event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
Box::pin(async move {
Ok(Transition {
next_state: MetricsAggregatorState::Running,
actions: vec![MetricsAggregatorAction::ExportMetrics],
})
})
};
on MetricsAggregatorEvent::Error => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
Box::pin(async move {
match event {
MetricsAggregatorEvent::Error(error) => Ok(Transition {
next_state: MetricsAggregatorState::Failed {
error: error.clone(),
},
actions: vec![],
}),
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for Error handler".to_string(),
)),
}
})
};
}
state MetricsAggregatorState::Draining {
on MetricsAggregatorEvent::FlowTerminal => |_state: &MetricsAggregatorState, _event: &MetricsAggregatorEvent, ctx: &mut MetricsAggregatorContext| {
Box::pin(async move {
let last_event_id = ctx.metrics_store.last_event_id;
let publish_last_event_id = last_event_id;
Ok(Transition {
next_state: MetricsAggregatorState::Drained { last_event_id },
actions: vec![
MetricsAggregatorAction::ExportMetrics,
MetricsAggregatorAction::PublishDrainComplete {
last_event_id: publish_last_event_id,
},
],
})
})
};
on MetricsAggregatorEvent::ProcessSystemEvent => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
let last_event_id = ctx.metrics_store.last_event_id;
Box::pin(async move {
match event {
MetricsAggregatorEvent::ProcessSystemEvent { envelope } => {
let pipeline_event = match &envelope.event.event {
obzenflow_core::event::SystemEventType::PipelineLifecycle(event) => {
Some(event)
}
_ => None,
};
let should_export = matches!(
pipeline_event,
Some(
obzenflow_core::event::PipelineLifecycleEvent::Draining { .. }
| obzenflow_core::event::PipelineLifecycleEvent::AllStagesCompleted { .. }
)
);
let should_finalize = matches!(
pipeline_event,
Some(
obzenflow_core::event::PipelineLifecycleEvent::Completed { .. }
| obzenflow_core::event::PipelineLifecycleEvent::Failed { .. }
| obzenflow_core::event::PipelineLifecycleEvent::Drained
)
);
if should_finalize {
let publish_last_event_id = last_event_id;
return Ok(Transition {
next_state: MetricsAggregatorState::Drained { last_event_id },
actions: vec![
MetricsAggregatorAction::ProcessSystemEvent {
envelope: envelope.clone(),
},
MetricsAggregatorAction::ExportMetrics,
MetricsAggregatorAction::PublishDrainComplete {
last_event_id: publish_last_event_id,
},
],
});
}
let mut actions = vec![MetricsAggregatorAction::ProcessSystemEvent {
envelope: envelope.clone(),
}];
if should_export {
actions.push(MetricsAggregatorAction::ExportMetrics);
}
Ok(Transition {
next_state: MetricsAggregatorState::Draining,
actions,
})
}
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for ProcessSystemEvent handler in Draining".to_string(),
)),
}
})
};
on MetricsAggregatorEvent::ProcessBatch => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
Box::pin(async move {
match event {
MetricsAggregatorEvent::ProcessBatch {
events,
journal_kind,
journal_stage,
} => {
let actions = events
.iter()
.cloned()
.map(|envelope| MetricsAggregatorAction::UpdateMetrics {
envelope: Box::new(envelope),
journal_kind,
journal_stage,
})
.collect::<Vec<_>>();
Ok(Transition {
next_state: MetricsAggregatorState::Draining,
actions,
})
}
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for ProcessBatch handler in Draining".to_string(),
)),
}
})
};
on MetricsAggregatorEvent::Error => |_state: &MetricsAggregatorState, event: &MetricsAggregatorEvent, _ctx: &mut MetricsAggregatorContext| {
let event = event.clone();
Box::pin(async move {
match event {
MetricsAggregatorEvent::Error(error) => Ok(Transition {
next_state: MetricsAggregatorState::Failed {
error: error.clone(),
},
actions: vec![],
}),
_ => Err(obzenflow_fsm::FsmError::HandlerError(
"Invalid event for Error handler".to_string(),
)),
}
})
};
}
state MetricsAggregatorState::Drained { }
state MetricsAggregatorState::Failed { }
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metrics::instrumentation::StageInstrumentation;
use async_trait::async_trait;
use obzenflow_core::event::context::{CompositeActivationContext, StageType};
use obzenflow_core::event::identity::JournalWriterId;
use obzenflow_core::event::payloads::correlation_payload::CorrelationPayload;
use obzenflow_core::event::payloads::delivery_payload::{DeliveryMethod, DeliveryPayload};
use obzenflow_core::event::status::processing_status::ErrorKind;
use obzenflow_core::event::ChainEventFactory;
use obzenflow_core::event::CorrelationId;
use obzenflow_core::event::JournalEvent;
use obzenflow_core::journal::journal_error::JournalError;
use obzenflow_core::journal::journal_owner::JournalOwner;
use obzenflow_core::journal::journal_reader::JournalReader;
use obzenflow_core::journal::Journal;
use obzenflow_core::metrics::StageMetadata;
use obzenflow_core::{EventEnvelope, JournalId};
use std::collections::VecDeque;
use std::marker::PhantomData;
#[test]
fn sink_operation_metric_projection_folds_only_typed_phase_and_error_kind() {
let stage_id = StageId::new();
let mut store = MetricsStore::default();
for phase in [
SinkOperationPhase::Open,
SinkOperationPhase::Write(obzenflow_core::event::SinkWritePhase::Encode),
SinkOperationPhase::Write(obzenflow_core::event::SinkWritePhase::Acquire),
SinkOperationPhase::Write(obzenflow_core::event::SinkWritePhase::Execute),
SinkOperationPhase::Write(obzenflow_core::event::SinkWritePhase::Commit),
SinkOperationPhase::Flush,
SinkOperationPhase::Drain,
] {
let failure = SinkOperationFailed {
stage_id,
stage_key: "sink".to_string(),
logical_destination: "high.cardinality.destination".to_string(),
causal_event_id: None,
input_position: None,
failed_delivery_event_id: None,
operation_subject_event_id: None,
phase,
kind: ErrorKind::Remote,
destination_error_code: Some(
obzenflow_core::event::SinkDestinationErrorCode::try_new(
"postgresql.sqlstate",
"23505",
)
.unwrap(),
),
detail: "free-form detail must never become a key".to_string(),
};
store.fold_sink_operation_failure(&failure);
store.fold_sink_operation_failure(&failure);
}
assert_eq!(store.sink_operation_failures.len(), 7);
assert!(store
.sink_operation_failures
.iter()
.all(|((metric_stage, _, kind), count)| {
*metric_stage == stage_id && *kind == ErrorKind::Remote && *count == 2
}));
}
#[test]
fn typed_http_pull_snapshots_fold_latest_state_and_monotonic_totals() {
use obzenflow_core::event::observability::{HttpPullState, WaitReason};
let stage_id = StageId::new();
let mut store = MetricsStore::default();
let first = HttpPullTelemetry {
state: HttpPullState::Waiting,
wait_reason: Some(WaitReason::PollInterval),
next_wake_unix_secs: Some(500),
last_success_unix_secs: Some(400),
requests_total: 8,
responses_2xx: 5,
responses_4xx: 2,
responses_5xx: 1,
rate_limited_total: 2,
retries_total: 3,
events_decoded_total: 21,
wait_seconds_rate_limit: 4.0,
wait_seconds_poll_interval: 7.0,
wait_seconds_backoff: 2.0,
};
store.fold_http_pull_snapshot(stage_id, &first);
let latest = HttpPullTelemetry {
state: HttpPullState::Fetching,
wait_reason: None,
next_wake_unix_secs: None,
last_success_unix_secs: Some(399),
requests_total: 7,
responses_2xx: 4,
responses_4xx: 1,
responses_5xx: 0,
rate_limited_total: 1,
retries_total: 2,
events_decoded_total: 20,
wait_seconds_rate_limit: 3.0,
wait_seconds_poll_interval: 6.0,
wait_seconds_backoff: 1.0,
};
store.fold_http_pull_snapshot(stage_id, &latest);
let folded = store
.http_pull_metrics
.get(&stage_id)
.expect("typed snapshot is indexed by stage");
assert!(matches!(folded.state, HttpPullState::Fetching));
assert!(folded.wait_reason.is_none());
assert_eq!(folded.next_wake_unix_secs, None);
assert_eq!(folded.last_success_unix_secs, Some(400));
assert_eq!(folded.requests_total, 8);
assert_eq!(folded.responses_2xx, 5);
assert_eq!(folded.responses_4xx, 2);
assert_eq!(folded.responses_5xx, 1);
assert_eq!(folded.rate_limited_total, 2);
assert_eq!(folded.retries_total, 3);
assert_eq!(folded.events_decoded_total, 21);
assert_eq!(folded.wait_seconds_rate_limit, 4.0);
assert_eq!(folded.wait_seconds_poll_interval, 7.0);
assert_eq!(folded.wait_seconds_backoff, 2.0);
}
#[test]
fn all_stages_completed_reconciles_missing_stage_lifecycle_states() {
let observed = StageId::new();
let missing = StageId::new();
let failed = StageId::new();
let mut stage_metadata = HashMap::new();
for stage_id in [observed, missing, failed] {
stage_metadata.insert(
stage_id,
StageMetadata {
name: format!("stage-{stage_id}"),
stage_type: StageType::Transform,
reference_mode: None,
flow_name: "test_flow".to_string(),
flow_id: None,
},
);
}
let mut store = MetricsStore::default();
store
.stage_lifecycle_states
.insert((observed, "completed".to_string()), true);
store
.stage_lifecycle_states
.insert((failed, "failed".to_string()), true);
assert!(!store.all_stages_terminal(&stage_metadata));
store.mark_known_stages_completed(stage_metadata.keys().copied());
assert!(store.all_stages_terminal(&stage_metadata));
assert_eq!(
store
.stage_lifecycle_states
.get(&(missing, "completed".to_string())),
Some(&true)
);
assert_eq!(
store
.stage_lifecycle_states
.get(&(failed, "completed".to_string())),
None
);
}
#[test]
fn circuit_breaker_snapshot_does_not_suppress_lifecycle_transition() {
let stage_id = StageId::new();
let mut store = MetricsStore::default();
let mut runtime_ctx = StageInstrumentation::new().snapshot();
runtime_ctx.cb_requests_total = 1;
runtime_ctx.cb_slow_total = 3;
runtime_ctx.cb_opened_total = 1;
runtime_ctx.cb_state = 1.0;
store.update_control_metrics_from_runtime_context(stage_id, &runtime_ctx);
assert_eq!(store.circuit_breaker_state.get(&stage_id), Some(&1.0));
assert_eq!(store.circuit_breaker_opened_total.get(&stage_id), Some(&1));
assert_eq!(store.circuit_breaker_slow_total.get(&stage_id), Some(&3));
assert!(
!store.circuit_breaker_last_state.contains_key(&stage_id),
"point-in-time snapshots must not advance lifecycle transition state"
);
store.record_circuit_breaker_transition(stage_id, "open");
assert_eq!(
store.circuit_breaker_state_transitions_total.get(&(
stage_id,
"closed".to_string(),
"open".to_string()
)),
Some(&1)
);
}
#[tokio::test]
async fn historical_prefix_fold_reconstructs_exact_duration_before_tail() {
struct VecReader {
events: VecDeque<EventEnvelope<ChainEvent>>,
position: u64,
}
#[async_trait]
impl JournalReader<ChainEvent> for VecReader {
async fn next(&mut self) -> Result<Option<EventEnvelope<ChainEvent>>, JournalError> {
let next = self.events.pop_front();
if next.is_some() {
self.position += 1;
}
Ok(next)
}
fn position(&self) -> u64 {
self.position
}
fn is_at_end(&self) -> bool {
self.events.is_empty()
}
}
let (entry, exit, peer) = (StageId::new(), StageId::new(), StageId::new());
let composite = obzenflow_core::id::CompositeId::new("saga:checkout");
let boundary = obzenflow_core::metrics::CompositeBoundary {
composite_id: composite.clone(),
members: vec![entry, exit],
ports: vec![
obzenflow_core::metrics::CompositeBoundaryPort {
name: "commands".to_string(),
direction: obzenflow_core::metrics::BoundaryDirection::Inbound,
member: entry,
payload_event_types: vec![EventType::from("checkout.command.v1")],
},
obzenflow_core::metrics::CompositeBoundaryPort {
name: "completed".to_string(),
direction: obzenflow_core::metrics::BoundaryDirection::Outbound,
member: exit,
payload_event_types: vec![EventType::from("checkout.completed.v1")],
},
],
edges: vec![obzenflow_core::metrics::CompositeBoundaryEdge {
port: "completed".to_string(),
direction: obzenflow_core::metrics::BoundaryDirection::Outbound,
member: exit,
peer,
upstream: exit,
downstream: peer,
}],
};
let mut output = ChainEventFactory::data_event(
WriterId::from(exit),
"checkout.completed.v1",
serde_json::json!({}),
);
output.processing_info.event_time = 1_250;
output = output
.try_with_composite_activations(vec![CompositeActivationContext::new(
composite,
EventId::new(),
"commands",
1_000,
)])
.unwrap();
let mut error_rail = CompositeDurationAccumulator::new(vec![0.1, 0.25, 1.0]);
observe_live_composite_duration(
MetricsJournalKind::Error,
&mut error_rail,
std::slice::from_ref(&boundary),
exit,
&output,
);
assert!(error_rail.histograms().is_empty());
let envelope = EventEnvelope::new(JournalWriterId::from(JournalId::new()), output);
let mut reader = VecReader {
events: VecDeque::from([envelope.clone(), envelope]),
position: 0,
};
let mut accumulator = CompositeDurationAccumulator::new(vec![0.1, 0.25, 1.0]);
let position = fold_composite_duration_prefix(
&mut reader,
exit,
std::slice::from_ref(&boundary),
&mut accumulator,
)
.await
.unwrap();
assert_eq!(position, 2);
let histograms = accumulator.histograms();
assert_eq!(histograms.len(), 1);
assert_eq!(histograms[0].count, 1);
assert_eq!(histograms[0].sum_seconds, 0.25);
}
#[tokio::test]
async fn test_delivery_event_preserves_correlation() {
let writer_id = WriterId::from(StageId::new());
let correlation_id = CorrelationId::new();
let mut event = ChainEventFactory::data_event(
writer_id,
"test.event",
serde_json::json!({"data": "test"}),
);
event.set_single_correlation(
correlation_id,
Some(CorrelationPayload::new("test_source", event.id)),
);
let payload = DeliveryPayload::success(DeliveryMethod::Noop, Some(1));
let delivery_event = ChainEventFactory::delivery_event(writer_id, payload)
.with_correlation_from(&event)
.with_cycle_state_from(&event);
let delivery_event = delivery_event
.try_with_composite_activations(event.composite_activations().to_vec())
.unwrap();
assert_eq!(delivery_event.correlation_id(), Some(correlation_id));
assert!(delivery_event.correlation_payload().is_some());
assert_eq!(
delivery_event.correlation_payload().unwrap().entry_stage,
"test_source"
);
}
#[test]
fn build_app_metrics_snapshot_uses_errors_by_kind_from_store() {
let stage_id = StageId::new();
let mut store = MetricsStore::default();
let errors_by_kind = HashMap::from([(ErrorKind::Domain, 2), (ErrorKind::Remote, 1)]);
let stage_metrics = StageMetrics {
errors_by_kind,
latest_events_processed_total: Some(42),
latest_errors_total: Some(3),
latest_data_outputs_by_event_type: HashMap::from([(
EventType::from("checkout.completed.v1"),
5,
)]),
event_loops_total: 10,
event_loops_with_work_total: 7,
..Default::default()
};
store.stage_metrics.insert(stage_id, stage_metrics);
let mut stage_metadata = std::collections::HashMap::new();
stage_metadata.insert(
stage_id,
StageMetadata {
name: "test_stage".to_string(),
stage_type: StageType::Sink,
reference_mode: None,
flow_name: "test_flow".to_string(),
flow_id: None,
},
);
struct NoopJournal<T: JournalEvent> {
id: obzenflow_core::id::JournalId,
owner: Option<JournalOwner>,
_marker: PhantomData<T>,
}
impl<T: JournalEvent> NoopJournal<T> {
fn new(owner: JournalOwner) -> Self {
Self {
id: obzenflow_core::id::JournalId::new(),
owner: Some(owner),
_marker: PhantomData,
}
}
}
struct NoopReader;
#[async_trait]
impl<T: JournalEvent + 'static> Journal<T> for NoopJournal<T> {
fn id(&self) -> &obzenflow_core::id::JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
_event: T,
_parent: Option<&obzenflow_core::EventEnvelope<T>>,
) -> Result<obzenflow_core::EventEnvelope<T>, JournalError> {
Err(JournalError::Implementation {
message: "noop journal".to_string(),
source: "noop".into(),
})
}
async fn read_all_unordered(
&self,
) -> Result<Vec<obzenflow_core::EventEnvelope<T>>, JournalError> {
Ok(Vec::new())
}
async fn read_event(
&self,
_event_id: &obzenflow_core::EventId,
) -> Result<Option<obzenflow_core::EventEnvelope<T>>, JournalError> {
Ok(None)
}
async fn reader_from(
&self,
_position: u64,
) -> Result<Box<dyn JournalReader<T>>, JournalError> {
Ok(Box::new(NoopReader))
}
async fn read_last_n(
&self,
_count: usize,
) -> Result<Vec<obzenflow_core::EventEnvelope<T>>, JournalError> {
Ok(Vec::new())
}
}
#[async_trait]
impl<T: JournalEvent + 'static> JournalReader<T> for NoopReader {
async fn next(
&mut self,
) -> Result<Option<obzenflow_core::EventEnvelope<T>>, JournalError> {
Ok(None)
}
fn position(&self) -> u64 {
0
}
fn is_at_end(&self) -> bool {
true
}
}
let downstream = StageId::new();
let ctx = MetricsAggregatorContext {
system_journal: Arc::new(NoopJournal::<obzenflow_core::event::SystemEvent>::new(
JournalOwner::system(obzenflow_core::SystemId::new()),
)),
stage_data_journals: HashMap::new(),
stage_error_journals: HashMap::new(),
backpressure_registry: None,
include_error_journals: true,
exporter: None,
metrics_store: store,
export_interval_secs: 10,
system_id: obzenflow_core::SystemId::new(),
stage_metadata,
composite_boundaries: vec![obzenflow_core::metrics::CompositeBoundary {
composite_id: obzenflow_core::id::CompositeId::new("saga:checkout"),
members: vec![stage_id],
ports: vec![obzenflow_core::metrics::CompositeBoundaryPort {
name: "completed".to_string(),
direction: obzenflow_core::metrics::BoundaryDirection::Outbound,
member: stage_id,
payload_event_types: vec![EventType::from("checkout.completed.v1")],
}],
edges: vec![obzenflow_core::metrics::CompositeBoundaryEdge {
port: "completed".to_string(),
direction: obzenflow_core::metrics::BoundaryDirection::Outbound,
member: stage_id,
peer: downstream,
upstream: stage_id,
downstream,
}],
}],
composite_durations: CompositeDurationAccumulator::default(),
};
let snapshot = ctx.build_app_metrics_snapshot();
assert_eq!(snapshot.error_counts.get(&stage_id), Some(&3));
let by_kind = snapshot
.error_counts_by_kind
.get(&stage_id)
.expect("per-kind breakdown should be present");
assert_eq!(by_kind.get(&ErrorKind::Domain), Some(&2));
assert_eq!(by_kind.get(&ErrorKind::Remote), Some(&1));
assert!(snapshot.stage_metadata.contains_key(&stage_id));
assert_eq!(snapshot.composite_port_traffic.len(), 1);
assert_eq!(snapshot.composite_port_traffic[0].events_total, 5);
assert_eq!(snapshot.composite_member_health.len(), 1);
assert_eq!(snapshot.composite_member_health[0].member_errors_total, 3);
}
}