use super::factory::ChainEventFactory;
use crate::event::envelope::AuthoredEnvelope;
use crate::event::observability::{ObservabilityContext, RuntimeSnapshot};
use crate::event::payloads::correlation_payload::CorrelationPayload;
use crate::event::payloads::effect_payload::EffectProvenance;
use crate::event::payloads::flow_control_payload::FlowControlPayload;
use crate::event::payloads::JournalPayload;
use crate::event::provenance::causality_context::CausalityContext;
use crate::event::provenance::{ChainEventProvenance, FlowContext, RuntimeProvenance};
use crate::event::status::processing_status::{ErrorKind, ProcessingStatus};
use crate::event::types::CorrelationId;
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)]
pub struct ChainEvent {
pub envelope: AuthoredEnvelope<ChainEventProvenance>,
pub payload: ChainPayload,
}
pub use crate::event::payloads::chain_payload::ChainPayload;
impl std::ops::Deref for ChainEvent {
type Target = ChainEventProvenance;
fn deref(&self) -> &Self::Target {
&self.envelope.provenance.event
}
}
impl std::ops::DerefMut for ChainEvent {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.envelope.provenance.event
}
}
impl Serialize for ChainEvent {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
use serde::ser::{Error, SerializeStruct};
JournalPayload::validate(&self.payload, &self.envelope.provenance.event)
.map_err(S::Error::custom)?;
let mut event = serializer.serialize_struct("ChainEvent", 2)?;
event.serialize_field("envelope", &self.envelope)?;
event.serialize_field("payload", &self.payload)?;
event.end()
}
}
impl<'de> Deserialize<'de> for ChainEvent {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
use serde::de::Error;
let raw = crate::event::record_serde::deserialize::<
_,
AuthoredEnvelope<ChainEventProvenance>,
Value,
>(deserializer)?;
let payload = ChainPayload::decode(
raw.envelope.provenance.event.event_kind,
&raw.envelope.provenance.event.event_type,
raw.payload,
)
.map_err(D::Error::custom)?;
JournalPayload::validate(&payload, &raw.envelope.provenance.event)
.map_err(D::Error::custom)?;
Ok(Self {
envelope: raw.envelope,
payload,
})
}
}
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.envelope.observability = Some(observability);
self
}
pub fn with_effect_provenance(mut self, provenance: EffectProvenance) -> Self {
self.effect_provenance = Some(provenance);
self
}
pub fn with_runtime_provenance(mut self, ctx: RuntimeProvenance) -> Self {
self.runtime = Some(ctx);
self
}
pub fn with_runtime_snapshot(mut self, snapshot: RuntimeSnapshot) -> Self {
let capture = snapshot.capture;
self.envelope
.observability
.get_or_insert_with(|| ObservabilityContext::new(capture))
.runtime_snapshot = Some(snapshot);
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.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
)
}
pub fn is_control(&self) -> bool {
matches!(self.payload, ChainPayload::FlowControl(_))
}
pub fn is_system(&self) -> bool {
false
}
pub fn is_fact(&self) -> bool {
matches!(self.payload, ChainPayload::Fact(_))
}
pub fn consumes_data_credit(&self) -> bool {
self.payload.consumes_data_credit()
}
pub fn is_typed_input(&self) -> bool {
matches!(
self.payload,
ChainPayload::Fact(_) | ChainPayload::CompositeData(_)
)
}
pub fn is_transport_excluded_execution(&self) -> bool {
matches!(&self.payload, ChainPayload::Execution(p) if !p.consumes_data_credit())
}
pub fn typed_payload(&self) -> Option<Value> {
match &self.payload {
ChainPayload::Fact(value) => Some(value.clone()),
ChainPayload::CompositeData(value) => serde_json::to_value(value).ok(),
ChainPayload::Execution(_)
| ChainPayload::FlowControl(_)
| ChainPayload::Delivery(_) => None,
}
}
pub fn is_delivery(&self) -> bool {
matches!(self.payload, ChainPayload::Delivery(_))
}
pub fn is_lifecycle(&self) -> bool {
matches!(&self.payload, ChainPayload::Execution(execution) if !execution.consumes_data_credit())
}
pub fn replay_disposition(&self) -> ReplayDisposition {
self.payload.replay_disposition()
}
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.status = ProcessingStatus::error_with_kind(reason.into(), Some(kind));
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 {
self.envelope.provenance.event.event_type.clone()
}
pub fn payload(&self) -> Value {
serde_json::to_value(&self.payload).expect("closed payloads serialize to JSON")
}
}