use crate::event::observability::HttpPullTelemetry;
use crate::event::payloads::effect_payload::EffectCursor;
use crate::id::StageId;
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "observability_type", rename_all = "snake_case")]
pub enum ObservabilityPayload {
Stage(StageLifecycle),
Metrics(MetricsLifecycle),
Middleware(MiddlewareLifecycle),
Backpressure(BackpressureEvent),
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "stage_state", rename_all = "snake_case")]
pub enum StageLifecycle {
Running {
stage_id: StageId,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<Value>,
},
Draining {
stage_id: StageId,
#[serde(skip_serializing_if = "Option::is_none")]
reason: Option<String>,
},
Drained {
stage_id: StageId,
#[serde(skip_serializing_if = "Option::is_none")]
events_processed: Option<u64>,
},
Completed {
stage_id: StageId,
#[serde(skip_serializing_if = "Option::is_none")]
final_metrics: Option<Value>,
},
Failed {
stage_id: StageId,
error: String,
#[serde(skip_serializing_if = "Option::is_none")]
recoverable: Option<bool>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "metrics_event", rename_all = "snake_case")]
pub enum MetricsLifecycle {
Ready {
#[serde(skip_serializing_if = "Option::is_none")]
exporter_count: Option<usize>,
},
StateSnapshot {
metrics: Value,
#[serde(skip_serializing_if = "Option::is_none")]
window_duration_ms: Option<u64>,
},
ResourceUsage {
cpu_percent: f64,
memory_bytes: u64,
#[serde(skip_serializing_if = "Option::is_none")]
thread_count: Option<u32>,
},
HttpPullSnapshot {
snapshot: HttpPullTelemetry,
},
Custom {
name: String,
value: Value,
#[serde(skip_serializing_if = "Option::is_none")]
tags: Option<Value>,
},
DrainRequested,
Drained {
#[serde(skip_serializing_if = "Option::is_none")]
final_flush_count: Option<u64>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(
tag = "middleware_event",
content = "details",
rename_all = "snake_case"
)]
pub enum MiddlewareLifecycle {
CircuitBreaker(CircuitBreakerEvent),
RateLimiter(RateLimiterEvent),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CircuitBreakerOpenTrigger {
ConsecutiveFailures,
FailureRate,
SlowCallRate,
FailureAndSlowCallRate,
HalfOpenProbeFailure,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum CircuitBreakerEvent {
Opened {
error_rate: f64,
failure_count: u64,
trigger: CircuitBreakerOpenTrigger,
observed_calls: u64,
#[serde(skip_serializing_if = "Option::is_none")]
slow_call_rate: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
slow_call_count: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
last_error: Option<String>,
},
Closed {
success_count: u64,
recovery_duration_ms: u64,
},
Rejected {
#[serde(default)]
reason: CircuitBreakerRejectionReason,
#[serde(default, skip_serializing_if = "Option::is_none")]
cooldown_remaining_ms: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
circuit_open_duration_ms: Option<u64>,
},
HalfOpen {
test_request_count: u32,
},
AttemptSettled {
cursor: EffectCursor,
attempt: u32,
health_classification: CircuitBreakerHealthClassification,
slow: bool,
dependency_elapsed_ms: u64,
admission_wait_ms: u64,
},
RetryScheduled {
cursor: EffectCursor,
next_attempt: u32,
delay_ms: u64,
},
RetrySucceeded {
cursor: EffectCursor,
total_attempts: u32,
terminal_classification: CircuitBreakerHealthClassification,
},
RetryExhausted {
cursor: EffectCursor,
total_attempts: u32,
reason: CircuitBreakerRetryStopReason,
},
RetryStoppedNonRetryable {
cursor: EffectCursor,
total_attempts: u32,
},
RecoveryCompleted {
cursor: EffectCursor,
total_attempts: u32,
backoff_elapsed_ms: u64,
recovery_elapsed_ms: u64,
},
Summary {
window_duration_s: u64,
requests_processed: u64,
requests_rejected: u64,
state: String,
consecutive_failures: usize,
rejection_rate: f64,
#[serde(default)]
successes_total: u64,
#[serde(default)]
failures_total: u64,
#[serde(default)]
opened_total: u64,
#[serde(default)]
time_in_closed_seconds: f64,
#[serde(default)]
time_in_open_seconds: f64,
#[serde(default)]
time_in_half_open_seconds: f64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CircuitBreakerHealthClassification {
Success,
TransientFailure,
PermanentFailure,
RateLimited,
Ignored,
NoObservation,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CircuitBreakerRetryStopReason {
AttemptLimit,
AttemptStartWindow,
CircuitNoLongerClosed,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CircuitBreakerRejectionReason {
CircuitOpen,
ProbeInProgress,
#[default]
Unknown,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum RateLimiterEvent {
Delayed {
delay_ms: u64,
current_rate: f64,
limit_rate: f64,
},
ActivityPulse {
window_ms: u64,
delayed_events: u64,
delay_ms_total: u64,
delay_ms_max: u64,
limit_rate: f64,
},
ModeChange {
mode_from: String,
mode_to: String,
limit_rate: f64,
},
WindowUtilization {
utilization_percent: f64,
events_in_window: u64,
window_size_ms: u64,
},
ConfigChanged {
old_rate: f64,
new_rate: f64,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum BackpressureEvent {
ActivityPulse {
window_ms: u64,
delayed_events: u64,
delay_ms_total: u64,
delay_ms_max: u64,
#[serde(skip_serializing_if = "Option::is_none")]
min_credit: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
limiting_downstream_stage_id: Option<StageId>,
},
Stalled {
upstream: StageId,
downstream: StageId,
window: u64,
stall_timeout_ms: u64,
elapsed_ms: u64,
in_flight: u64,
},
}