use super::factory::ChainEventFactory;
use crate::event::context::causality_context::CausalityContext;
use crate::event::context::observability_context::ObservabilityContext;
use crate::event::context::{
FlowContext, IntentContext, ProcessingContext, ReplayContext, RuntimeContext,
};
use crate::event::payloads::correlation_payload::CorrelationPayload;
use crate::event::payloads::delivery_payload::DeliveryPayload;
use crate::event::payloads::effect_payload::{is_framework_effect_event_type, EffectProvenance};
use crate::event::payloads::flow_control_payload::FlowControlPayload;
use crate::event::payloads::observability_payload::{
MetricsLifecycle, MiddlewareLifecycle, ObservabilityPayload, StageLifecycle,
};
use crate::event::status::processing_status::{ErrorKind, ProcessingStatus};
use crate::event::types::{AdmissionSeq, CorrelationId, EventId, WriterId};
use crate::id::{CycleDepth, SccId};
use crate::ingress::IngressContext;
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CorrelationContext {
pub ids: Vec<CorrelationId>,
#[serde(default, skip_serializing_if = "is_false")]
pub truncated: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub payload: Option<CorrelationPayload>,
}
impl CorrelationContext {
pub fn single(id: CorrelationId, payload: Option<CorrelationPayload>) -> Self {
Self {
ids: vec![id],
truncated: false,
payload,
}
}
pub fn sample(ids: Vec<CorrelationId>, truncated: bool) -> Self {
Self {
ids,
truncated,
payload: None,
}
}
pub fn single_id(&self) -> Option<CorrelationId> {
if self.ids.len() == 1 && !self.truncated {
self.ids.first().copied()
} else {
None
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChainEvent {
pub id: EventId,
pub writer_id: WriterId,
pub content: ChainEventContent,
pub causality: CausalityContext,
pub flow_context: FlowContext,
pub processing_info: ProcessingContext,
pub intent: Option<IntentContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub correlation: Option<CorrelationContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub replay_context: Option<ReplayContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub ingress_context: Option<IngressContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cycle_depth: Option<CycleDepth>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cycle_scc_id: Option<SccId>,
#[serde(skip_serializing_if = "Option::is_none")]
pub runtime_context: Option<RuntimeContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub observability: Option<ObservabilityContext>,
#[serde(skip_serializing_if = "Option::is_none")]
pub effect_provenance: Option<EffectProvenance>,
#[serde(skip_serializing_if = "Option::is_none")]
pub admission_seq: Option<AdmissionSeq>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "content_type", rename_all = "snake_case")]
pub enum ChainEventContent {
#[serde(rename = "data")]
Data {
event_type: String, payload: Value,
},
#[serde(rename = "flow_signal")]
FlowControl(FlowControlPayload),
#[serde(rename = "delivery")]
Delivery(DeliveryPayload),
#[serde(rename = "lifecycle")]
Observability(ObservabilityPayload),
}
fn is_false(value: &bool) -> bool {
!*value
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReplayDisposition {
ReAdmit,
ReAuthor,
}
impl ChainEvent {
pub fn with_observability_context(mut self, observability: ObservabilityContext) -> Self {
self.observability = Some(observability);
self
}
pub fn with_effect_provenance(mut self, provenance: EffectProvenance) -> Self {
self.effect_provenance = Some(provenance);
self
}
pub fn with_runtime_context(mut self, ctx: RuntimeContext) -> Self {
self.runtime_context = Some(ctx);
self
}
pub fn with_ingress_context(mut self, ctx: IngressContext) -> Self {
self.ingress_context = Some(ctx);
self
}
pub fn with_flow_context(mut self, ctx: FlowContext) -> Self {
self.flow_context = ctx;
self
}
pub fn with_causality(mut self, causality: CausalityContext) -> Self {
self.causality = causality;
self
}
pub fn is_eof(&self) -> bool {
matches!(
self.content,
ChainEventContent::FlowControl(FlowControlPayload::Eof { .. })
)
}
pub fn is_control(&self) -> bool {
matches!(self.content, ChainEventContent::FlowControl(_))
}
pub fn is_system(&self) -> bool {
false
}
pub fn is_data(&self) -> bool {
matches!(self.content, ChainEventContent::Data { .. })
}
pub fn is_delivery(&self) -> bool {
matches!(self.content, ChainEventContent::Delivery(_))
}
pub fn is_lifecycle(&self) -> bool {
matches!(self.content, ChainEventContent::Observability(_))
}
pub fn replay_disposition(&self) -> ReplayDisposition {
match &self.content {
ChainEventContent::Data { event_type, .. } => {
let is_framework_effect_record = self
.effect_provenance
.as_ref()
.is_some_and(|provenance| provenance.fact_owner.is_framework())
&& is_framework_effect_event_type(event_type);
if is_framework_effect_record {
ReplayDisposition::ReAuthor
} else {
ReplayDisposition::ReAdmit
}
}
ChainEventContent::FlowControl(payload) => match payload {
FlowControlPayload::Watermark { .. }
| FlowControlPayload::CatchUpComplete { .. } => ReplayDisposition::ReAdmit,
FlowControlPayload::Eof { .. }
| FlowControlPayload::Checkpoint { .. }
| FlowControlPayload::Drain
| FlowControlPayload::PipelineAbort { .. }
| FlowControlPayload::SourceContract { .. }
| FlowControlPayload::ConsumptionProgress { .. }
| FlowControlPayload::ConsumptionGap { .. }
| FlowControlPayload::ConsumptionFinal { .. }
| FlowControlPayload::ReaderStalled { .. }
| FlowControlPayload::AtLeastOnceViolation { .. } => ReplayDisposition::ReAuthor,
},
ChainEventContent::Delivery(_) => ReplayDisposition::ReAuthor,
ChainEventContent::Observability(_) => ReplayDisposition::ReAuthor,
}
}
pub fn is_source_replayable(&self) -> bool {
self.replay_disposition() == ReplayDisposition::ReAdmit
}
pub fn mark_as_error(mut self, reason: impl Into<String>, kind: ErrorKind) -> Self {
self.processing_info.status = ProcessingStatus::error_with_kind(reason.into(), Some(kind));
self.processing_info.error_hops_remaining = Some(1);
self
}
pub fn mark_as_validation_error(self, reason: impl Into<String>) -> Self {
self.mark_as_error(reason, ErrorKind::Validation)
}
pub fn mark_as_infra_error(self, reason: impl Into<String>) -> Self {
self.mark_as_error(reason, ErrorKind::Remote)
}
pub fn derive_error_event(
&self,
event_type: impl Into<String>,
payload: Value,
reason: impl Into<String>,
kind: ErrorKind,
lineage: crate::config::LineagePolicy,
) -> ChainEvent {
let reason_str = reason.into();
ChainEventFactory::derived_data_event(self.writer_id, self, event_type, payload, lineage)
.mark_as_error(reason_str, kind)
}
pub fn event_type(&self) -> String {
match &self.content {
ChainEventContent::Data { event_type, .. } => event_type.clone(),
ChainEventContent::FlowControl(signal) => match signal {
FlowControlPayload::Eof { .. } => "control.eof".into(),
FlowControlPayload::Watermark { .. } => "control.watermark".into(),
FlowControlPayload::CatchUpComplete { .. } => "control.catch_up_complete".into(),
FlowControlPayload::Checkpoint { .. } => "control.checkpoint".into(),
FlowControlPayload::Drain => "control.drain".into(),
FlowControlPayload::PipelineAbort { .. } => "control.pipeline_abort".into(),
FlowControlPayload::SourceContract { .. } => "control.source_contract".into(),
FlowControlPayload::ConsumptionProgress { .. } => {
"control.consumption_progress".into()
}
FlowControlPayload::ConsumptionGap { .. } => "control.consumption_gap".into(),
FlowControlPayload::ConsumptionFinal { .. } => "control.consumption_final".into(),
FlowControlPayload::ReaderStalled { .. } => "control.reader_stalled".into(),
FlowControlPayload::AtLeastOnceViolation { .. } => {
"control.at_least_once_violation".into()
}
},
ChainEventContent::Delivery(_) => "sink.delivery".into(),
ChainEventContent::Observability(obs) => match obs {
ObservabilityPayload::Stage(stage) => match stage {
StageLifecycle::Running { .. } => "lifecycle.stage.running".into(),
StageLifecycle::Draining { .. } => "lifecycle.stage.draining".into(),
StageLifecycle::Drained { .. } => "lifecycle.stage.drained".into(),
StageLifecycle::Completed { .. } => "lifecycle.stage.completed".into(),
StageLifecycle::Failed { .. } => "lifecycle.stage.failed".into(),
},
ObservabilityPayload::Metrics(metrics) => match metrics {
MetricsLifecycle::Ready { .. } => "lifecycle.metrics.ready".into(),
MetricsLifecycle::StateSnapshot { .. } => "lifecycle.metrics.state".into(),
MetricsLifecycle::ResourceUsage { .. } => "lifecycle.metrics.resource".into(),
MetricsLifecycle::HttpPullSnapshot { .. } => {
"lifecycle.metrics.http_pull_snapshot".into()
}
MetricsLifecycle::Custom { .. } => "lifecycle.metrics.custom".into(),
MetricsLifecycle::DrainRequested => "lifecycle.metrics.drain".into(),
MetricsLifecycle::Drained { .. } => "lifecycle.metrics.drained".into(),
},
ObservabilityPayload::Middleware(mw) => match mw {
MiddlewareLifecycle::CircuitBreaker(_) => {
"lifecycle.middleware.circuit_breaker".into()
}
MiddlewareLifecycle::RateLimiter(_) => {
"lifecycle.middleware.rate_limiter".into()
}
},
ObservabilityPayload::Backpressure(_) => "lifecycle.backpressure".into(),
},
}
}
pub fn payload(&self) -> Value {
match &self.content {
ChainEventContent::Data { payload, .. } => payload.clone(),
ChainEventContent::FlowControl(signal) => {
serde_json::to_value(signal).unwrap_or_default()
}
ChainEventContent::Delivery(delivery) => {
serde_json::to_value(delivery).unwrap_or_default()
}
ChainEventContent::Observability(lifecycle) => {
serde_json::to_value(lifecycle).unwrap_or_default()
}
}
}
}