use crate::event::payloads::execution_payload::MiddlewareFact;
use crate::event::payloads::flow_control_payload::EofKind;
use crate::event::provenance::ExecutionAccounting;
use crate::event::types::{Count, DurationMs, EventId, EventType, SeqNo};
use crate::event::vector_clock::VectorClock;
use crate::id::{StageId, StageKey};
use crate::ingress::{IngressAttemptSeq, IngressKey, IngressRefusalReason};
use crate::journal::{ArchiveStatus, StatusDerivation};
use crate::metrics::FlowLifecycleMetricsSnapshot;
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
use std::str::FromStr;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MiddlewareEventOrigin {
pub event_id: EventId,
pub writer_key: String,
pub seq: SeqNo,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct ContractName(String);
impl ContractName {
pub fn new(value: impl Into<String>) -> Self {
Self(value.into())
}
pub fn as_str(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for ContractName {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
impl From<&str> for ContractName {
fn from(value: &str) -> Self {
Self::new(value)
}
}
impl From<String> for ContractName {
fn from(value: String) -> Self {
Self::new(value)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SystemFeedRole {
Input,
Reference,
Stream,
}
impl SystemFeedRole {
pub const fn as_str(self) -> &'static str {
match self {
Self::Input => "input",
Self::Reference => "reference",
Self::Stream => "stream",
}
}
}
impl std::fmt::Display for SystemFeedRole {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
impl FromStr for SystemFeedRole {
type Err = ();
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"input" => Ok(Self::Input),
"reference" => Ok(Self::Reference),
"stream" => Ok(Self::Stream),
_ => Err(()),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CommandDiscardDisposition {
ObsoleteControl,
UnexpectedError,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "system_event_type", rename_all = "snake_case")]
pub enum SystemPayload {
SupervisorCommandDiscarded {
supervisor: String,
terminal_state: String,
command: String,
disposition: CommandDiscardDisposition,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<String>,
},
#[serde(rename = "source_cleanup_failed")]
SourceCleanupFailed {
stage_id: StageId,
stage_name: String,
error: String,
},
#[serde(rename = "stage_lifecycle")]
StageLifecycle {
stage_id: StageId,
#[serde(flatten)]
event: StageLifecycleEvent,
},
#[serde(rename = "pipeline_lifecycle")]
PipelineLifecycle(PipelineLifecycleEvent),
#[serde(rename = "replay_lifecycle")]
ReplayLifecycle(ReplayLifecycleEvent),
#[serde(rename = "metrics_coordination")]
MetricsCoordination(MetricsCoordinationEvent),
#[serde(rename = "middleware_lifecycle")]
MiddlewareLifecycle {
stage_id: StageId,
#[serde(skip_serializing_if = "Option::is_none")]
stage_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
flow_name: Option<String>,
origin: MiddlewareEventOrigin,
middleware: MiddlewareFact,
},
#[serde(rename = "contract_status")]
ContractStatus {
upstream: StageId,
reader: StageId,
#[serde(default, skip_serializing_if = "Option::is_none")]
selected_event_type: Option<EventType>,
#[serde(default, skip_serializing_if = "Option::is_none")]
feed_role: Option<SystemFeedRole>,
pass: bool,
#[serde(skip_serializing_if = "Option::is_none")]
reader_seq: Option<crate::event::types::SeqNo>,
#[serde(skip_serializing_if = "Option::is_none")]
advertised_writer_seq: Option<crate::event::types::SeqNo>,
#[serde(skip_serializing_if = "Option::is_none")]
reason: Option<crate::event::types::ViolationCause>,
},
#[serde(rename = "contract_result")]
ContractResult {
upstream: StageId,
reader: StageId,
#[serde(default, skip_serializing_if = "Option::is_none")]
selected_event_type: Option<EventType>,
#[serde(default, skip_serializing_if = "Option::is_none")]
feed_role: Option<SystemFeedRole>,
contract_name: ContractName,
status: ContractResultStatusLabel,
#[serde(skip_serializing_if = "Option::is_none")]
cause: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
reader_seq: Option<crate::event::types::SeqNo>,
#[serde(skip_serializing_if = "Option::is_none")]
advertised_writer_seq: Option<crate::event::types::SeqNo>,
},
#[serde(rename = "ingress_refusal")]
IngressRefusal {
ingress_key: IngressKey,
stage_id: StageId,
stage_key: StageKey,
reason: IngressRefusalReason,
attempt_seq: IngressAttemptSeq,
request_count: u64,
event_count: u64,
batch_count: u64,
http_status: u16,
#[serde(skip_serializing_if = "Option::is_none")]
retry_after_ms_bucket: Option<u64>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ContractResultStatusLabel {
Passed,
Failed,
Pending,
Healthy,
}
impl ContractResultStatusLabel {
pub const fn as_str(self) -> &'static str {
match self {
Self::Passed => "passed",
Self::Failed => "failed",
Self::Pending => "pending",
Self::Healthy => "healthy",
}
}
}
impl std::fmt::Display for ContractResultStatusLabel {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
impl std::str::FromStr for ContractResultStatusLabel {
type Err = ();
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s {
"passed" => Ok(Self::Passed),
"failed" => Ok(Self::Failed),
"pending" => Ok(Self::Pending),
"healthy" => Ok(Self::Healthy),
_ => Err(()),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "lifecycle_event", rename_all = "snake_case")]
pub enum StageLifecycleEvent {
Running,
Draining {
#[serde(skip_serializing_if = "Option::is_none")]
accounting: Option<ExecutionAccounting>,
},
Drained,
Completed {
#[serde(skip_serializing_if = "Option::is_none")]
accounting: Option<ExecutionAccounting>,
},
Cancelled {
reason: String,
#[serde(skip_serializing_if = "Option::is_none")]
accounting: Option<ExecutionAccounting>,
},
Failed {
error: String,
#[serde(skip_serializing_if = "Option::is_none")]
recoverable: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
accounting: Option<ExecutionAccounting>,
#[serde(default, skip_serializing_if = "Option::is_none")]
causal_event_id: Option<EventId>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "pipeline_event", rename_all = "snake_case")]
pub enum PipelineLifecycleEvent {
Starting,
ReadyForRun {
#[serde(skip_serializing_if = "Option::is_none")]
stage_count: Option<usize>,
},
Running {
#[serde(skip_serializing_if = "Option::is_none")]
stage_count: Option<usize>,
},
StopAdmitted {
admission: PipelineStopAdmission,
},
NotStarted,
Draining {
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<FlowLifecycleMetricsSnapshot>,
},
AllStagesCompleted {
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<FlowLifecycleMetricsSnapshot>,
},
Drained,
Completed {
duration_ms: DurationMs,
metrics: FlowLifecycleMetricsSnapshot,
},
Failed {
reason: String,
duration_ms: DurationMs,
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<FlowLifecycleMetricsSnapshot>,
#[serde(skip_serializing_if = "Option::is_none")]
failure_cause: Option<crate::event::types::ViolationCause>,
},
Cancelled {
reason: String,
duration_ms: DurationMs,
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<FlowLifecycleMetricsSnapshot>,
#[serde(skip_serializing_if = "Option::is_none")]
failure_cause: Option<crate::event::types::ViolationCause>,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "mode", rename_all = "snake_case")]
pub enum PipelineStopAdmission {
Graceful { timeout_ms: DurationMs },
Cancel { cause: PipelineCancellationCause },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PipelineCancellationCause {
Requested,
GracefulTimeout,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "replay_event", rename_all = "snake_case")]
pub enum ReplayLifecycleEvent {
Started {
archive_path: PathBuf,
archive_flow_id: String,
archive_status: ArchiveStatus,
archive_status_derivation: StatusDerivation,
allow_incomplete: bool,
source_stages: Vec<String>,
},
Completed {
replayed_count: Count,
skipped_count: Count,
duration_ms: DurationMs,
#[serde(default, skip_serializing_if = "Option::is_none")]
synthesized_eof_kind: Option<EofKind>,
},
ResumedLive {
archive_flow_id: String,
replayed_count: Count,
generation: u64,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "metrics_event", rename_all = "snake_case")]
pub enum MetricsCoordinationEvent {
Ready,
DrainRequested,
Drained,
Shutdown,
Exported {
watermark: VectorClock,
},
}
impl SystemPayload {
pub fn event_type(&self) -> &'static str {
match self {
SystemPayload::SupervisorCommandDiscarded { .. } => {
"system.supervisor.command_discarded"
}
SystemPayload::SourceCleanupFailed { .. } => "system.source.cleanup_failed",
SystemPayload::StageLifecycle { event, .. } => match event {
StageLifecycleEvent::Running => "system.stage.running",
StageLifecycleEvent::Draining { .. } => "system.stage.draining",
StageLifecycleEvent::Drained => "system.stage.drained",
StageLifecycleEvent::Completed { .. } => "system.stage.completed",
StageLifecycleEvent::Failed { .. } => "system.stage.failed",
StageLifecycleEvent::Cancelled { .. } => "system.stage.cancelled",
},
SystemPayload::PipelineLifecycle(event) => match event {
PipelineLifecycleEvent::Starting => "system.pipeline.starting",
PipelineLifecycleEvent::ReadyForRun { .. } => "system.pipeline.ready_for_run",
PipelineLifecycleEvent::Running { .. } => "system.pipeline.running",
PipelineLifecycleEvent::StopAdmitted { .. } => "system.pipeline.stop_admitted",
PipelineLifecycleEvent::NotStarted => "system.pipeline.not_started",
PipelineLifecycleEvent::AllStagesCompleted { .. } => {
"system.pipeline.all_stages_completed"
}
PipelineLifecycleEvent::Draining { .. } => "system.pipeline.draining",
PipelineLifecycleEvent::Drained => "system.pipeline.drained",
PipelineLifecycleEvent::Completed { .. } => "system.pipeline.completed",
PipelineLifecycleEvent::Failed { .. } => "system.pipeline.failed",
PipelineLifecycleEvent::Cancelled { .. } => "system.pipeline.cancelled",
},
SystemPayload::ReplayLifecycle(event) => match event {
ReplayLifecycleEvent::Started { .. } => "system.replay.started",
ReplayLifecycleEvent::Completed { .. } => "system.replay.completed",
ReplayLifecycleEvent::ResumedLive { .. } => "system.replay.resumed_live",
},
SystemPayload::MetricsCoordination(event) => match event {
MetricsCoordinationEvent::Ready => "system.metrics.ready",
MetricsCoordinationEvent::DrainRequested => "system.metrics.drain_requested",
MetricsCoordinationEvent::Drained => "system.metrics.drained",
MetricsCoordinationEvent::Shutdown => "system.metrics.shutdown",
MetricsCoordinationEvent::Exported { .. } => "system.metrics.exported",
},
SystemPayload::MiddlewareLifecycle { .. } => "system.middleware.lifecycle",
SystemPayload::ContractStatus { pass, .. } => {
if *pass {
"system.contract.pass"
} else {
"system.contract.fail"
}
}
SystemPayload::ContractResult { status, .. } => match status {
ContractResultStatusLabel::Passed => "system.contract.result.passed",
ContractResultStatusLabel::Failed => "system.contract.result.failed",
ContractResultStatusLabel::Pending => "system.contract.result.pending",
ContractResultStatusLabel::Healthy => "system.contract.result",
},
SystemPayload::IngressRefusal { .. } => "system.ingress.refusal",
}
}
}