use crate::event::observability::HttpSurfaceMetricsSnapshot;
use crate::event::payloads::flow_control_payload::EofKind;
use crate::event::payloads::observability_payload::MiddlewareLifecycle;
use crate::event::types::{Count, DurationMs, EventId, EventType, SeqNo, WriterId};
use crate::event::vector_clock::VectorClock;
use crate::id::{StageId, StageKey, SystemId};
use crate::ingress::{IngressAttemptSeq, IngressKey, IngressRefusalReason};
use crate::journal::{ArchiveStatus, StatusDerivation};
use crate::metrics::{FlowLifecycleMetricsSnapshot, StageMetricsSnapshot};
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, Serialize, Deserialize)]
pub struct SystemEvent {
pub id: EventId,
pub writer_id: WriterId,
#[serde(flatten)]
pub event: SystemEventType,
pub timestamp: u64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "system_event_type", rename_all = "snake_case")]
pub enum SystemEventType {
#[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 = "stage_heartbeat")]
StageHeartbeat {
stage_id: StageId,
stage_name: String,
heartbeat_seq: SeqNo,
activity: StageActivity,
last_consumed_event_id: Option<EventId>,
last_output_event_id: Option<EventId>,
handler_blocked_ms: Option<DurationMs>,
},
#[serde(rename = "edge_liveness")]
EdgeLiveness {
upstream: StageId,
reader: StageId,
state: EdgeLivenessState,
idle_ms: DurationMs,
#[serde(skip_serializing_if = "Option::is_none")]
last_reader_seq: Option<SeqNo>,
#[serde(skip_serializing_if = "Option::is_none")]
last_event_id: Option<EventId>,
},
#[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: MiddlewareLifecycle,
},
#[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 = "http_surface_snapshot")]
HttpSurfaceSnapshot {
#[serde(flatten)]
snapshot: HttpSurfaceMetricsSnapshot,
},
#[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")]
metrics: Option<StageMetricsSnapshot>,
},
Drained,
Completed {
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<StageMetricsSnapshot>,
},
Cancelled {
reason: String,
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<StageMetricsSnapshot>,
},
Failed {
error: String,
#[serde(skip_serializing_if = "Option::is_none")]
recoverable: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
metrics: Option<StageMetricsSnapshot>,
#[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>,
},
StopRequested {
mode: String,
#[serde(skip_serializing_if = "Option::is_none")]
timeout_ms: Option<DurationMs>,
},
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, 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 },
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum StageActivity {
Polling,
Processing {
event_id: EventId,
elapsed_ms: DurationMs,
},
WaitingOnQuietInput { upstream: Option<StageId> },
Draining,
Completed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum EdgeLivenessState {
Healthy,
Idle,
Suspect,
Stalled,
Recovered,
}
impl SystemEvent {
pub fn new(writer_id: WriterId, event: SystemEventType) -> Self {
Self {
id: EventId::new(),
writer_id,
event,
timestamp: current_timestamp(),
}
}
pub fn stage_running(stage_id: StageId) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Running,
},
)
}
pub fn stage_completed(stage_id: StageId) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Completed { metrics: None },
},
)
}
pub fn stage_cancelled(stage_id: StageId, reason: String) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Cancelled {
reason,
metrics: None,
},
},
)
}
pub fn stage_failed(stage_id: StageId, error: String, recoverable: bool) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Failed {
error,
recoverable: Some(recoverable),
metrics: None,
causal_event_id: None,
},
},
)
}
pub fn stage_draining_with_metrics(stage_id: StageId, metrics: StageMetricsSnapshot) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Draining {
metrics: Some(metrics),
},
},
)
}
pub fn stage_completed_with_metrics(stage_id: StageId, metrics: StageMetricsSnapshot) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Completed {
metrics: Some(metrics),
},
},
)
}
pub fn stage_failed_with_metrics(
stage_id: StageId,
error: String,
recoverable: bool,
metrics: StageMetricsSnapshot,
) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Failed {
error,
recoverable: Some(recoverable),
metrics: Some(metrics),
causal_event_id: None,
},
},
)
}
pub fn stage_failed_with_metrics_causal(
stage_id: StageId,
error: String,
recoverable: bool,
metrics: StageMetricsSnapshot,
causal_event_id: EventId,
) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Failed {
error,
recoverable: Some(recoverable),
metrics: Some(metrics),
causal_event_id: Some(causal_event_id),
},
},
)
}
pub fn stage_cancelled_with_metrics(
stage_id: StageId,
reason: String,
metrics: StageMetricsSnapshot,
) -> Self {
Self::new(
WriterId::from(stage_id),
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Cancelled {
reason,
metrics: Some(metrics),
},
},
)
}
}
fn current_timestamp() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as u64
}
pub struct SystemEventFactory {
writer_id: WriterId,
}
impl SystemEventFactory {
pub fn new(system_id: SystemId) -> Self {
Self {
writer_id: WriterId::from(system_id),
}
}
pub fn stage_running(&self, stage_id: StageId) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Running,
},
)
}
pub fn stage_draining(&self, stage_id: StageId) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Draining { metrics: None },
},
)
}
pub fn stage_drained(&self, stage_id: StageId) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Drained,
},
)
}
pub fn stage_completed(&self, stage_id: StageId) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Completed { metrics: None },
},
)
}
pub fn stage_failed(&self, stage_id: StageId, error: String, recoverable: bool) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Failed {
error,
recoverable: Some(recoverable),
metrics: None,
causal_event_id: None,
},
},
)
}
pub fn stage_cancelled(&self, stage_id: StageId, reason: String) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::StageLifecycle {
stage_id,
event: StageLifecycleEvent::Cancelled {
reason,
metrics: None,
},
},
)
}
pub fn contract_status(
&self,
upstream: StageId,
reader: StageId,
pass: bool,
reader_seq: Option<crate::event::types::SeqNo>,
advertised_writer_seq: Option<crate::event::types::SeqNo>,
reason: Option<crate::event::types::ViolationCause>,
) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::ContractStatus {
upstream,
reader,
selected_event_type: None,
feed_role: None,
pass,
reader_seq,
advertised_writer_seq,
reason,
},
)
}
pub fn pipeline_starting(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Starting),
)
}
pub fn pipeline_running(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Running {
stage_count: None,
}),
)
}
pub fn pipeline_ready_for_run(&self, stage_count: Option<usize>) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::ReadyForRun { stage_count }),
)
}
pub fn pipeline_stop_requested(
&self,
mode: String,
timeout_ms: Option<DurationMs>,
) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::StopRequested {
mode,
timeout_ms,
}),
)
}
pub fn pipeline_all_stages_completed(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::AllStagesCompleted {
metrics: None,
}),
)
}
pub fn pipeline_draining(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Draining { metrics: None }),
)
}
pub fn pipeline_drained(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Drained),
)
}
pub fn pipeline_completed(
&self,
duration_ms: DurationMs,
metrics: FlowLifecycleMetricsSnapshot,
) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Completed {
duration_ms,
metrics,
}),
)
}
pub fn pipeline_failed(
&self,
reason: String,
duration_ms: DurationMs,
metrics: Option<FlowLifecycleMetricsSnapshot>,
failure_cause: Option<crate::event::types::ViolationCause>,
) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Failed {
reason,
duration_ms,
metrics,
failure_cause,
}),
)
}
pub fn pipeline_cancelled(
&self,
reason: String,
duration_ms: DurationMs,
metrics: Option<FlowLifecycleMetricsSnapshot>,
failure_cause: Option<crate::event::types::ViolationCause>,
) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::PipelineLifecycle(PipelineLifecycleEvent::Cancelled {
reason,
duration_ms,
metrics,
failure_cause,
}),
)
}
pub fn metrics_ready(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::MetricsCoordination(MetricsCoordinationEvent::Ready),
)
}
pub fn metrics_drain_requested(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::MetricsCoordination(MetricsCoordinationEvent::DrainRequested),
)
}
pub fn metrics_drained(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::MetricsCoordination(MetricsCoordinationEvent::Drained),
)
}
pub fn metrics_shutdown(&self) -> SystemEvent {
SystemEvent::new(
self.writer_id,
SystemEventType::MetricsCoordination(MetricsCoordinationEvent::Shutdown),
)
}
}
use crate::event::journal_event::{JournalEvent, Sealed};
impl Sealed for SystemEvent {}
impl JournalEvent for SystemEvent {
fn id(&self) -> &EventId {
&self.id
}
fn writer_id(&self) -> &WriterId {
&self.writer_id
}
fn event_type_name(&self) -> &str {
match &self.event {
SystemEventType::SourceCleanupFailed { .. } => "system.source.cleanup_failed",
SystemEventType::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",
},
SystemEventType::PipelineLifecycle(event) => match event {
PipelineLifecycleEvent::Starting => "system.pipeline.starting",
PipelineLifecycleEvent::ReadyForRun { .. } => "system.pipeline.ready_for_run",
PipelineLifecycleEvent::Running { .. } => "system.pipeline.running",
PipelineLifecycleEvent::StopRequested { .. } => "system.pipeline.stop_requested",
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",
},
SystemEventType::ReplayLifecycle(event) => match event {
ReplayLifecycleEvent::Started { .. } => "system.replay.started",
ReplayLifecycleEvent::Completed { .. } => "system.replay.completed",
ReplayLifecycleEvent::ResumedLive { .. } => "system.replay.resumed_live",
},
SystemEventType::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",
},
SystemEventType::StageHeartbeat { .. } => "system.stage.heartbeat",
SystemEventType::EdgeLiveness { .. } => "system.edge.liveness",
SystemEventType::MiddlewareLifecycle { .. } => "system.middleware.lifecycle",
SystemEventType::ContractStatus { pass, .. } => {
if *pass {
"system.contract.pass"
} else {
"system.contract.fail"
}
}
SystemEventType::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",
},
SystemEventType::HttpSurfaceSnapshot { .. } => "system.http_surface.snapshot",
SystemEventType::IngressRefusal { .. } => "system.ingress.refusal",
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::StageId;
use serde_json::json;
#[test]
fn contract_result_feed_fields_are_typed_but_serialize_as_labels() {
let payload = SystemEventType::ContractResult {
upstream: StageId::new(),
reader: StageId::new(),
selected_event_type: Some(EventType::from("test.selected.v1")),
feed_role: Some(SystemFeedRole::Reference),
contract_name: ContractName::from("TransportContract"),
status: ContractResultStatusLabel::Healthy,
cause: None,
reader_seq: Some(SeqNo(3)),
advertised_writer_seq: Some(SeqNo(5)),
};
let serialized = serde_json::to_value(&payload).expect("system event should serialize");
assert_eq!(serialized["selected_event_type"], "test.selected.v1");
assert_eq!(serialized["feed_role"], "reference");
assert_eq!(serialized["contract_name"], "TransportContract");
assert_eq!(serialized["status"], "healthy");
let decoded: SystemEventType = serde_json::from_value(json!({
"system_event_type": "contract_result",
"upstream": serialized["upstream"].clone(),
"reader": serialized["reader"].clone(),
"selected_event_type": "test.selected.v1",
"feed_role": "reference",
"contract_name": "TransportContract",
"status": "healthy",
"reader_seq": 3,
"advertised_writer_seq": 5
}))
.expect("string-label system event should deserialize");
match decoded {
SystemEventType::ContractResult {
selected_event_type,
feed_role,
contract_name,
status,
..
} => {
assert_eq!(
selected_event_type,
Some(EventType::from("test.selected.v1"))
);
assert_eq!(feed_role, Some(SystemFeedRole::Reference));
assert_eq!(contract_name.as_str(), "TransportContract");
assert_eq!(status, ContractResultStatusLabel::Healthy);
}
other => panic!("expected ContractResult, got {other:?}"),
}
}
}