use super::{
CompositeEntrySpec, ContractConfig, ContractStatus, ContractsWiring, DeliveredCount,
FeedIdentity, PollResult, ReaderProgress, ReaderSelectionPolicy, SelectedFeedMetadata,
SelectedFeedRole, StageInputPosition, UpstreamSubscription,
};
use crate::control_plane::{ControlPlaneProvider, NoControlPlane};
use async_trait::async_trait;
use obzenflow_core::chrono::Utc;
use obzenflow_core::event::identity::JournalWriterId;
use obzenflow_core::event::journal_event::JournalEvent;
use obzenflow_core::event::journal_record::JournalRecord;
use obzenflow_core::event::payloads::delivery_payload::{DeliveryMethod, DeliveryPayload};
use obzenflow_core::event::payloads::effect_payload::{EffectFactOwner, EffectProvenance};
use obzenflow_core::event::payloads::execution_payload::ExecutionPayload;
use obzenflow_core::event::payloads::flow_control_payload::FlowControlPayload;
use obzenflow_core::event::payloads::system_payload::{ContractResultStatusLabel, SystemPayload};
use obzenflow_core::event::provenance::causality_context::CausalityContext;
use obzenflow_core::event::provenance::JournalProvenance;
use obzenflow_core::event::system_event::SystemEvent;
use obzenflow_core::event::types::{
Count, DurationMs, SeqNo, ViolationCause as EventViolationCause,
};
use obzenflow_core::event::vector_clock::VectorClock;
use obzenflow_core::event::{ChainEvent, ChainEventFactory, ChainPayload};
use obzenflow_core::id::{CompositeId, JournalId};
use obzenflow_core::journal::journal_error::JournalError;
use obzenflow_core::journal::journal_owner::JournalOwner;
use obzenflow_core::journal::reader::JournalReader;
use obzenflow_core::journal::AppendOptions;
use obzenflow_core::journal::Journal;
use obzenflow_core::{
AdmissionSeq, ContractResult, DeliveryContract, EventId, EventType, ReaderGeneration, StageId,
TransportContract, WriterId,
};
use serde_json::json;
use std::collections::HashMap;
use std::io;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use tokio::time::Instant;
fn committed_input(event: ChainEvent, vector_clock: VectorClock) -> JournalRecord<ChainPayload> {
JournalRecord::commit_event(
event,
JournalProvenance {
journal_writer_id: JournalWriterId::new(),
vector_clock,
timestamp: Utc::now(),
journal_group_id: None,
journal_group_member: None,
},
)
.unwrap()
}
fn contract_flow_context(stage_id: StageId) -> obzenflow_core::event::provenance::FlowContext {
crate::stages::common::supervision::flow_context_factory::make_flow_context(
"contract_flow",
"contract_run",
"consumer",
stage_id,
obzenflow_core::event::context::StageType::Join,
)
}
fn assert_contract_owner(
event: &ChainEvent,
owner: &obzenflow_core::event::provenance::FlowContext,
) {
assert_eq!(event.writer_id, WriterId::from(owner.stage_id));
assert_eq!(
serde_json::to_value(&event.flow_context).unwrap(),
serde_json::to_value(owner).unwrap()
);
}
#[tokio::test]
async fn contract_owner_context_survives_fan_in_append_failure_and_retry() {
let upstream_a = StageId::new();
let upstream_b = StageId::new();
let upstreams = [upstream_a, upstream_b].map(|id| {
(
id,
id.to_string(),
Arc::new(TestJournal::new(JournalOwner::stage(id))) as Arc<dyn Journal<ChainEvent>>,
)
});
let owner = contract_flow_context(StageId::new());
let journal: Arc<dyn Journal<ChainEvent>> = Arc::new(ControlledJournal::new(
JournalOwner::stage(owner.stage_id),
Arc::new(|_: &ChainEvent, call| call == 0),
));
let mut sub = UpstreamSubscription::new_with_names("consumer", &upstreams)
.await
.unwrap()
.with_contracts(ContractsWiring {
writer_id: WriterId::from(owner.stage_id),
contract_journal: journal.clone(),
config: ContractConfig::default(),
system_journal: None,
reader_stage: Some(owner.stage_id),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
})
.with_contract_flow_context(owner.clone());
let mut progress = [
ReaderProgress::new(upstream_a),
ReaderProgress::new(upstream_b),
];
progress[0].reader_seq = SeqNo(1);
progress[1].reader_seq = SeqNo(2);
assert!(matches!(
sub.check_contracts(&mut progress).await,
ContractStatus::ProgressEmitted
));
assert_eq!(progress[0].last_progress_seq, SeqNo(0));
assert_eq!(progress[1].last_progress_seq, SeqNo(2));
assert!(matches!(
sub.check_contracts(&mut progress).await,
ContractStatus::ProgressEmitted
));
assert_eq!(progress[0].last_progress_seq, SeqNo(1));
assert_eq!(progress[1].last_progress_seq, SeqNo(2));
let events = journal.read_causally_ordered().await.unwrap();
assert_eq!(events.len(), 2);
let mut observed = HashMap::new();
for env in events {
assert_contract_owner(&env.authored(), &owner);
let ChainPayload::FlowControl(FlowControlPayload::ConsumptionProgress {
reader_index,
reader_path,
reader_seq,
..
}) = env.payload
else {
panic!("expected progress")
};
observed.insert(reader_index.0, (reader_path.0, reader_seq));
}
assert_eq!(observed.get(&0), Some(&(upstream_a.to_string(), SeqNo(1))));
assert_eq!(observed.get(&1), Some(&(upstream_b.to_string(), SeqNo(2))));
}
#[tokio::test]
async fn contract_owner_context_covers_both_gap_paths_violation_and_final() {
for legacy in [false, true] {
let (mut sub, journal, _, upstream, consumer) =
build_upstream_with_seq_divergence(Arc::new(NoControlPlane)).await;
let owner = contract_flow_context(consumer);
sub = sub.with_contract_flow_context(owner.clone());
if legacy {
sub.contract_chains.clear();
}
let mut progress = [ReaderProgress::new(upstream)];
drive_subscription_to_eof(&mut sub, &mut progress).await;
assert!(matches!(
sub.check_contracts(&mut progress).await,
ContractStatus::Violated { .. }
));
let events = journal.read_causally_ordered().await.unwrap();
let (mut gap, mut violation, mut final_record) = (false, false, false);
for env in events {
assert_contract_owner(&env.authored(), &owner);
match env.payload {
ChainPayload::FlowControl(FlowControlPayload::ConsumptionGap { .. }) => gap = true,
ChainPayload::FlowControl(FlowControlPayload::AtLeastOnceViolation {
upstream: id,
..
}) => {
assert_eq!(id, upstream);
violation = true;
}
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal {
pass, ..
}) => {
assert!(!pass);
final_record = true;
}
_ => {}
}
}
assert!(gap && violation && final_record, "legacy={legacy}");
}
}
struct TestJournal<T: JournalEvent> {
id: JournalId,
owner: Option<JournalOwner>,
events: Arc<Mutex<Vec<JournalRecord<T::Payload>>>>,
}
impl<T: JournalEvent> TestJournal<T> {
fn new(owner: JournalOwner) -> Self {
Self {
id: JournalId::new(),
owner: Some(owner),
events: Arc::new(Mutex::new(Vec::new())),
}
}
}
type AppendFailurePredicate<T> = dyn Fn(&T, usize) -> bool + Send + Sync;
struct ControlledJournal<T: JournalEvent> {
id: JournalId,
owner: Option<JournalOwner>,
events: Arc<Mutex<Vec<JournalRecord<T::Payload>>>>,
append_calls: AtomicUsize,
should_fail: Arc<AppendFailurePredicate<T>>,
}
impl<T: JournalEvent> ControlledJournal<T> {
fn new(owner: JournalOwner, should_fail: Arc<AppendFailurePredicate<T>>) -> Self {
Self {
id: JournalId::new(),
owner: Some(owner),
events: Arc::new(Mutex::new(Vec::new())),
append_calls: AtomicUsize::new(0),
should_fail,
}
}
}
struct TestJournalReader<T: JournalEvent> {
events: Vec<JournalRecord<T::Payload>>,
pos: usize,
}
#[async_trait]
impl<T: JournalEvent + 'static> Journal<T> for TestJournal<T> {
fn id(&self) -> &JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
event: T,
mut options: AppendOptions<'_, T>,
) -> std::result::Result<JournalRecord<T::Payload>, JournalError> {
let event = options.capture.prepare(0, event);
let envelope = JournalRecord::new(JournalWriterId::from(self.id), event);
let mut guard = self.events.lock().unwrap();
guard.push(envelope.clone());
Ok(envelope)
}
async fn read_all_unordered(
&self,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().unwrap();
Ok(guard.clone())
}
async fn read_event(
&self,
_event_id: &obzenflow_core::EventId,
) -> std::result::Result<Option<JournalRecord<T::Payload>>, JournalError> {
Ok(None)
}
async fn reader_from(
&self,
position: u64,
) -> std::result::Result<Box<dyn JournalReader<T>>, JournalError> {
let guard = self.events.lock().unwrap();
Ok(Box::new(TestJournalReader {
events: guard.clone(),
pos: position as usize,
}))
}
async fn read_last_n(
&self,
count: usize,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().unwrap();
let len = guard.len();
let start = len.saturating_sub(count);
Ok(guard[start..].iter().rev().cloned().collect())
}
}
#[async_trait]
impl<T: JournalEvent + 'static> Journal<T> for ControlledJournal<T> {
fn id(&self) -> &JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
event: T,
mut options: AppendOptions<'_, T>,
) -> std::result::Result<JournalRecord<T::Payload>, JournalError> {
let event = options.capture.prepare(0, event);
let call_index = self.append_calls.fetch_add(1, Ordering::Relaxed);
if (self.should_fail)(&event, call_index) {
return Err(JournalError::Implementation {
message: "append failed".to_string(),
source: "append failed".into(),
});
}
let envelope = JournalRecord::new(JournalWriterId::from(self.id), event);
let mut guard = self.events.lock().unwrap();
guard.push(envelope.clone());
Ok(envelope)
}
async fn read_all_unordered(
&self,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().unwrap();
Ok(guard.clone())
}
async fn read_event(
&self,
_event_id: &obzenflow_core::EventId,
) -> std::result::Result<Option<JournalRecord<T::Payload>>, JournalError> {
Ok(None)
}
async fn reader_from(
&self,
position: u64,
) -> std::result::Result<Box<dyn JournalReader<T>>, JournalError> {
let guard = self.events.lock().unwrap();
Ok(Box::new(TestJournalReader {
events: guard.clone(),
pos: position as usize,
}))
}
async fn read_last_n(
&self,
count: usize,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().unwrap();
let len = guard.len();
let start = len.saturating_sub(count);
Ok(guard[start..].iter().rev().cloned().collect())
}
}
#[async_trait]
impl<T: JournalEvent + 'static> JournalReader<T> for TestJournalReader<T> {
async fn next(
&mut self,
) -> std::result::Result<Option<JournalRecord<T::Payload>>, JournalError> {
if self.pos >= self.events.len() {
Ok(None)
} else {
let envelope = self.events.get(self.pos).cloned();
self.pos += 1;
Ok(envelope)
}
}
fn position(&self) -> u64 {
self.pos as u64
}
fn is_at_end(&self) -> bool {
self.pos >= self.events.len()
}
}
#[cfg(unix)]
struct EmfileJournal<T: JournalEvent> {
id: JournalId,
owner: Option<JournalOwner>,
_phantom: std::marker::PhantomData<T>,
}
#[cfg(unix)]
impl<T: JournalEvent> EmfileJournal<T> {
fn new(owner: JournalOwner) -> Self {
Self {
id: JournalId::new(),
owner: Some(owner),
_phantom: std::marker::PhantomData,
}
}
}
#[cfg(unix)]
#[async_trait]
impl<T: JournalEvent + 'static> Journal<T> for EmfileJournal<T> {
fn id(&self) -> &JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
_event: T,
_options: AppendOptions<'_, T>,
) -> std::result::Result<JournalRecord<T::Payload>, JournalError> {
Err(JournalError::Implementation {
message: "append not supported".to_string(),
source: "append not supported".into(),
})
}
async fn read_all_unordered(
&self,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
Ok(Vec::new())
}
async fn read_event(
&self,
_event_id: &obzenflow_core::EventId,
) -> std::result::Result<Option<JournalRecord<T::Payload>>, JournalError> {
Ok(None)
}
async fn reader(&self) -> std::result::Result<Box<dyn JournalReader<T>>, JournalError> {
Err(JournalError::Implementation {
message: "open failed".to_string(),
source: Box::new(io::Error::from_raw_os_error(libc::EMFILE)),
})
}
async fn reader_from(
&self,
_position: u64,
) -> std::result::Result<Box<dyn JournalReader<T>>, JournalError> {
self.reader().await
}
async fn read_last_n(
&self,
_count: usize,
) -> std::result::Result<Vec<JournalRecord<T::Payload>>, JournalError> {
Ok(Vec::new())
}
}
#[tokio::test]
#[cfg(unix)]
async fn fails_fast_on_too_many_open_files() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(EmfileJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let err = UpstreamSubscription::<ChainEvent>::new_with_names_from_positions(
"downstream",
&upstreams,
&[0u64],
)
.await
.err()
.expect("Expected Too many open files error")
.to_string();
assert!(err.contains("Too many open files"));
}
#[tokio::test]
async fn progress_append_failure_does_not_advance_progress_state() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(ControlledJournal::new(
contract_owner,
Arc::new(|_event: &ChainEvent, _call| true),
));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].reader_seq = SeqNo(1);
reader_progress[0].last_progress_seq = SeqNo(0);
let _status = subscription.check_contracts(&mut reader_progress).await;
assert_eq!(reader_progress[0].last_progress_seq, SeqNo(0));
assert!(reader_progress[0].last_progress_instant.is_none());
}
#[tokio::test]
async fn final_append_failure_keeps_final_emitted_false() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(ControlledJournal::new(
contract_owner,
Arc::new(|event: &ChainEvent, _call| {
matches!(
&event.payload,
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal { .. })
)
}),
));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
subscription.state.mark_reader_eof(0);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let _status = subscription.check_contracts(&mut reader_progress).await;
assert!(!reader_progress[0].final_emitted);
assert!(!reader_progress[0].contract_violated);
let events = contract_journal.read_causally_ordered().await.unwrap();
assert!(
!events.iter().any(|env| matches!(
&env.payload,
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal { .. })
)),
"expected final event append to have failed"
);
}
#[tokio::test]
async fn diagnostics_only_eof_check_does_not_emit_final_or_latch_state() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: true,
cycle_guard_config: None,
});
subscription.state.mark_reader_eof(0);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].reader_seq = SeqNo(1);
reader_progress[0].receipted_seq = SeqNo(1);
let _status = subscription
.check_contracts_diagnostics_only(&mut reader_progress)
.await;
assert!(!reader_progress[0].final_emitted);
assert!(!reader_progress[0].contract_violated);
let events = contract_journal.read_causally_ordered().await.unwrap();
assert!(
!events.iter().any(|env| matches!(
&env.payload,
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal { .. })
)),
"diagnostics-only checks must not emit final contract evidence"
);
}
#[tokio::test]
async fn diagnostics_only_eof_check_does_not_emit_stall() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
let reader_stage = StageId::new();
let system_owner = JournalOwner::stage(reader_stage);
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(system_owner));
let config = ContractConfig {
progress_min_events: Count(100),
progress_max_interval: DurationMs(10_000),
stall_threshold: DurationMs(100),
stall_cooloff: DurationMs(0),
stall_checks_before_emit: 1,
};
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config,
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: true,
cycle_guard_config: None,
});
subscription.state.mark_reader_eof(0);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].reader_seq = SeqNo(1);
reader_progress[0].receipted_seq = SeqNo(1);
reader_progress[0].advertised_writer_seq = Some(SeqNo(1));
reader_progress[0].last_read_instant =
Some(Instant::now() - std::time::Duration::from_millis(250));
let status = subscription
.check_contracts_diagnostics_only(&mut reader_progress)
.await;
assert!(
matches!(status, ContractStatus::Healthy),
"EOF readers must not be stall-checked in diagnostics-only mode"
);
assert_eq!(reader_progress[0].consecutive_stall_checks, 0);
assert!(reader_progress[0].stalled_since.is_none());
assert!(!reader_progress[0].contract_violated);
assert!(contract_journal
.read_causally_ordered()
.await
.expect("read contract journal")
.is_empty());
assert!(system_journal
.read_causally_ordered()
.await
.expect("read system journal")
.is_empty());
}
#[tokio::test]
async fn contract_status_append_failure_keeps_final_emitted_false() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
let reader_stage = StageId::new();
let system_owner = JournalOwner::stage(reader_stage);
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(ControlledJournal::new(
system_owner,
Arc::new(|event: &SystemEvent, _call| {
matches!(&event.payload, SystemPayload::ContractStatus { .. })
}),
));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
subscription.contract_chains = (0..subscription.readers.len()).map(|_| None).collect();
subscription.contract_policies = (0..subscription.readers.len()).map(|_| None).collect();
subscription.state.mark_reader_eof(0);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].advertised_writer_seq = Some(SeqNo(0));
reader_progress[0].reader_seq = SeqNo(0);
let _status = subscription.check_contracts(&mut reader_progress).await;
assert!(!reader_progress[0].final_emitted);
assert!(!reader_progress[0].contract_violated);
let contract_events = contract_journal.read_causally_ordered().await.unwrap();
assert!(
contract_events.iter().any(|env| matches!(
&env.payload,
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal { .. })
)),
"expected final event to be persisted even when ContractStatus append fails"
);
}
#[tokio::test]
async fn progress_contract_heartbeats_are_suppressed_until_data_observed() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
let reader_stage = StageId::new();
let system_owner = JournalOwner::stage(reader_stage);
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(system_owner));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].last_event_id = Some(EventId::new());
reader_progress[0].reader_seq = SeqNo(0);
let _status = subscription.check_contracts(&mut reader_progress).await;
let events = system_journal.read_causally_ordered().await.unwrap();
assert!(
!events.iter().any(|env| matches!(
&env.payload,
SystemPayload::ContractResult { .. } | SystemPayload::ContractStatus { .. }
)),
"expected progress contract heartbeats to be suppressed before any data is observed"
);
reader_progress[0].reader_seq = SeqNo(1);
let _status = subscription.check_contracts(&mut reader_progress).await;
let events = system_journal.read_causally_ordered().await.unwrap();
assert!(
events.iter().any(|env| matches!(
&env.payload,
SystemPayload::ContractResult { contract_name, status, cause, .. }
if contract_name.as_str() == TransportContract::NAME
&& *status == ContractResultStatusLabel::Healthy
&& cause.is_none()
)),
"expected a healthy TransportContract ContractResult heartbeat once data is observed"
);
}
#[tokio::test]
async fn progress_emission_uses_receipt_watermark_when_delivery_contract_enabled() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: true,
cycle_guard_config: None,
});
let read_event_id = EventId::new();
let receipted_event_id = EventId::new();
let mut read_clock = VectorClock::new();
read_clock.clocks.insert("upstream".to_string(), 3);
let mut receipted_clock = VectorClock::new();
receipted_clock.clocks.insert("upstream".to_string(), 1);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].reader_seq = SeqNo(3);
reader_progress[0].receipted_seq = SeqNo(1);
reader_progress[0].last_event_id = Some(read_event_id);
reader_progress[0].last_vector_clock = Some(read_clock);
reader_progress[0].last_receipted_event_id = Some(receipted_event_id);
reader_progress[0].last_receipted_vector_clock = Some(receipted_clock.clone());
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(matches!(status, ContractStatus::ProgressEmitted));
assert_eq!(reader_progress[0].last_progress_seq, SeqNo(1));
let events = contract_journal.read_causally_ordered().await.unwrap();
let progress = events
.iter()
.find_map(|env| match &env.payload {
ChainPayload::FlowControl(FlowControlPayload::ConsumptionProgress {
reader_seq,
last_event_id,
vector_clock,
..
}) => Some((*reader_seq, *last_event_id, vector_clock.clone())),
_ => None,
})
.expect("expected ConsumptionProgress event");
assert_eq!(progress.0, SeqNo(1));
assert_eq!(progress.1, Some(receipted_event_id));
assert_eq!(progress.2, Some(receipted_clock));
}
#[tokio::test]
async fn record_delivery_receipt_advances_only_when_receipts_become_contiguous() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: true,
cycle_guard_config: None,
});
let writer_id = WriterId::from(upstream_stage);
let first = ChainEventFactory::data_event(writer_id, "test.event", json!({"seq": 1}));
let second = ChainEventFactory::data_event(writer_id, "test.event", json!({"seq": 2}));
let mut clock_1 = VectorClock::new();
clock_1.clocks.insert("upstream".to_string(), 1);
let mut clock_2 = VectorClock::new();
clock_2.clocks.insert("upstream".to_string(), 2);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].reader_seq = SeqNo(1);
reader_progress[0]
.track_pending_delivery_input(committed_input(first.clone(), clock_1.clone()));
reader_progress[0].track_pending_receipt(first.id, clock_1);
reader_progress[0].reader_seq = SeqNo(2);
reader_progress[0]
.track_pending_delivery_input(committed_input(second.clone(), clock_2.clone()));
reader_progress[0].track_pending_receipt(second.id, clock_2.clone());
let second_receipt = ChainEventFactory::delivery_event(
WriterId::from(contract_stage),
DeliveryPayload::success(DeliveryMethod::Noop, None),
)
.with_causality(CausalityContext::with_parent(second.id));
assert!(subscription
.record_delivery_receipt(&second_receipt, &mut reader_progress)
.is_none());
assert_eq!(reader_progress[0].receipted_seq, SeqNo(0));
let first_receipt = ChainEventFactory::delivery_event(
WriterId::from(contract_stage),
DeliveryPayload::success(DeliveryMethod::Noop, None),
)
.with_causality(CausalityContext::with_parent(first.id));
let watermark = subscription
.record_delivery_receipt(&first_receipt, &mut reader_progress)
.expect("expected contiguous receipts to advance watermark");
assert_eq!(watermark.0, SeqNo(2));
assert_eq!(watermark.1, second.id);
assert_eq!(watermark.2, clock_2);
assert_eq!(reader_progress[0].receipted_seq, SeqNo(2));
assert!(reader_progress[0].pending_delivery_inputs.is_empty());
assert!(reader_progress[0].pending_receipts.is_empty());
assert!(reader_progress[0].committed_out_of_order.is_empty());
}
#[tokio::test]
async fn forwarded_sink_input_settles_without_entering_authored_delivery_contract() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let sink_stage = StageId::new();
let contract_owner = JournalOwner::stage(sink_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(sink_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: None,
reader_stage: Some(sink_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: true,
cycle_guard_config: None,
});
let source_stage = StageId::new();
let forwarded = ChainEventFactory::data_event(
WriterId::from(source_stage),
"forwarded.event",
json!({"value": 1}),
);
let mut clock = VectorClock::new();
clock.clocks.insert("source".to_string(), 1);
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].track_pending_delivery_input(committed_input(forwarded.clone(), clock));
let receipt = ChainEventFactory::delivery_event(
WriterId::from(sink_stage),
DeliveryPayload::success(DeliveryMethod::Noop, None),
)
.with_causality(CausalityContext::with_parent(forwarded.id));
assert!(subscription
.record_delivery_receipt(&receipt, &mut reader_progress)
.is_none());
assert!(reader_progress[0].pending_delivery_inputs.is_empty());
assert_eq!(reader_progress[0].receipted_seq, SeqNo(0));
let delivery_result = subscription.contract_chains[0]
.as_ref()
.expect("delivery contract chain")
.verify_all(upstream_stage, sink_stage)
.into_iter()
.find_map(|(name, result)| (name.as_str() == DeliveryContract::NAME).then_some(result))
.expect("delivery contract result");
assert!(matches!(delivery_result, ContractResult::Passed(_)));
}
#[tokio::test]
async fn stall_append_failure_does_not_set_stalled_since() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(ControlledJournal::new(
contract_owner,
Arc::new(|event: &ChainEvent, _call| {
matches!(
&event.payload,
ChainPayload::FlowControl(FlowControlPayload::ReaderStalled { .. })
)
}),
));
let config = ContractConfig {
progress_min_events: Count(100),
progress_max_interval: DurationMs(10_000),
stall_threshold: DurationMs(100),
stall_cooloff: DurationMs(0),
stall_checks_before_emit: 1,
};
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal,
config,
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].last_read_instant =
Some(Instant::now() - std::time::Duration::from_millis(250));
reader_progress[0].last_progress_seq = reader_progress[0].reader_seq;
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(status, ContractStatus::Stalled(s) if s == upstream_stage),
"expected stall status even when append fails"
);
assert!(reader_progress[0].stalled_since.is_none());
assert!(!reader_progress[0].contract_violated);
}
#[tokio::test]
async fn stall_cooloff_suppresses_repeat_stalled_emission() {
tokio::time::pause();
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
let config = ContractConfig {
progress_min_events: Count(100),
progress_max_interval: DurationMs(500),
stall_threshold: DurationMs(100),
stall_cooloff: DurationMs(1_000),
stall_checks_before_emit: 1,
};
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config,
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let owner = contract_flow_context(contract_stage);
subscription = subscription.with_contract_flow_context(owner.clone());
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
reader_progress[0].last_read_instant = Some(Instant::now());
reader_progress[0].last_progress_instant = Some(Instant::now());
reader_progress[0].last_progress_seq = reader_progress[0].reader_seq;
tokio::time::advance(std::time::Duration::from_millis(200)).await;
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(status, ContractStatus::Stalled(s) if s == upstream_stage),
"expected initial stalled status"
);
tokio::time::advance(std::time::Duration::from_millis(400)).await;
let _ = subscription.check_contracts(&mut reader_progress).await;
tokio::time::advance(std::time::Duration::from_millis(50)).await;
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(status, ContractStatus::Stalled(s) if s == upstream_stage),
"expected stalled status during cooloff"
);
let envelopes = contract_journal
.read_causally_ordered()
.await
.expect("read contract journal");
for env in &envelopes {
assert_contract_owner(&env.authored(), &owner);
}
let stalled_count = envelopes
.iter()
.filter(|e| {
matches!(
&e.payload,
ChainPayload::FlowControl(FlowControlPayload::ReaderStalled { .. })
)
})
.count();
assert_eq!(
stalled_count, 1,
"expected ReaderStalled to be emitted once within cooloff"
);
}
#[tokio::test]
async fn idle_reader_without_any_reads_does_not_emit_stall() {
tokio::time::pause();
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(contract_owner));
let config = ContractConfig {
progress_min_events: Count(100),
progress_max_interval: DurationMs(10_000),
stall_threshold: DurationMs(100),
stall_cooloff: DurationMs(0),
stall_checks_before_emit: 1,
};
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal: contract_journal.clone(),
config,
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
for _ in 0..4 {
tokio::time::advance(std::time::Duration::from_millis(150)).await;
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(status, ContractStatus::Healthy),
"idle reader should stay healthy before any reads"
);
}
assert!(reader_progress[0].last_read_instant.is_none());
assert!(reader_progress[0].stalled_since.is_none());
assert!(!reader_progress[0].contract_violated);
assert!(contract_journal
.read_causally_ordered()
.await
.expect("read contract journal")
.is_empty());
}
#[tokio::test]
async fn multi_reader_progress_isolated_under_partial_append_failure() {
let upstream_a = StageId::new();
let upstream_b = StageId::new();
let journal_a: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(upstream_a)));
let journal_b: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(upstream_b)));
let upstreams = [
(upstream_a, "upstream_a".to_string(), journal_a),
(upstream_b, "upstream_b".to_string(), journal_b),
];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_stage = StageId::new();
let contract_owner = JournalOwner::stage(contract_stage);
let contract_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(ControlledJournal::new(
contract_owner,
Arc::new(|event: &ChainEvent, _call| match &event.payload {
ChainPayload::FlowControl(FlowControlPayload::ConsumptionProgress {
reader_index,
..
}) => reader_index.0 == 0,
_ => false,
}),
));
subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(contract_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: None,
reader_stage: None,
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
let mut reader_progress = [
ReaderProgress::new(upstream_a),
ReaderProgress::new(upstream_b),
];
reader_progress[0].reader_seq = SeqNo(1);
reader_progress[0].last_progress_seq = SeqNo(0);
reader_progress[1].reader_seq = SeqNo(1);
reader_progress[1].last_progress_seq = SeqNo(0);
let status = subscription.check_contracts(&mut reader_progress[..]).await;
assert!(matches!(status, ContractStatus::ProgressEmitted));
assert_eq!(reader_progress[0].last_progress_seq, SeqNo(0));
assert!(reader_progress[0].last_progress_instant.is_none());
assert_eq!(reader_progress[1].last_progress_seq, SeqNo(1));
assert!(reader_progress[1].last_progress_instant.is_some());
}
async fn build_upstream_with_seq_divergence(
control_plane: Arc<dyn ControlPlaneProvider>,
) -> (
UpstreamSubscription<ChainEvent>,
Arc<dyn Journal<ChainEvent>>,
Arc<dyn Journal<SystemEvent>>,
StageId,
StageId,
) {
let upstream_stage = StageId::new();
let reader_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let reader_owner = JournalOwner::stage(reader_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(reader_owner.clone()));
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(reader_owner));
let writer_id = WriterId::Stage(upstream_stage);
let data_event = ChainEventFactory::data_event(writer_id, "test.event", json!({}));
upstream_journal
.append(data_event, Default::default())
.await
.unwrap();
let mut eof_event = ChainEventFactory::eof_event(writer_id, true);
if let ChainPayload::FlowControl(FlowControlPayload::Eof {
writer_id: writer_id_field,
writer_seq,
..
}) = &mut eof_event.payload
{
*writer_id_field = Some(writer_id);
*writer_seq = Some(SeqNo(3));
}
upstream_journal
.append(eof_event, Default::default())
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap();
let contract_config = ContractConfig::default();
let writer_id_for_contracts = WriterId::from(reader_stage);
subscription = subscription.with_contracts(ContractsWiring {
writer_id: writer_id_for_contracts,
contract_journal: contract_journal.clone(),
config: contract_config,
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane,
include_delivery_contract: false,
cycle_guard_config: None,
});
(
subscription,
contract_journal,
system_journal,
upstream_stage,
reader_stage,
)
}
async fn drive_subscription_to_eof(
subscription: &mut UpstreamSubscription<ChainEvent>,
reader_progress: &mut [ReaderProgress],
) {
loop {
match subscription
.poll_next_with_state("test_fsm", Some(reader_progress))
.await
{
PollResult::Event(_env) => continue,
PollResult::CursorAdvanced { .. } => continue,
PollResult::NoEvents => break,
PollResult::Error(e) => {
panic!("poll_next_with_state returned error: {e:?}");
}
}
}
}
#[tokio::test]
async fn strict_mode_produces_seq_divergence_and_gap_event() {
let (mut subscription, contract_journal, system_journal, upstream_stage, reader_stage) =
build_upstream_with_seq_divergence(Arc::new(NoControlPlane)).await;
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
drive_subscription_to_eof(&mut subscription, &mut reader_progress).await;
let status = subscription.check_contracts(&mut reader_progress).await;
match status {
ContractStatus::Violated { upstream, cause } => {
assert_eq!(upstream, upstream_stage);
match cause {
EventViolationCause::SeqDivergence { advertised, reader } => {
assert_eq!(advertised, Some(SeqNo(3)));
assert_eq!(reader, SeqNo(1));
}
other => panic!("expected SeqDivergence cause, got {other:?}"),
}
}
other => panic!("expected violated status, got {other:?}"),
}
let events = contract_journal.read_causally_ordered().await.unwrap();
let mut final_found = false;
let mut gap_found = false;
let mut violation_found = false;
for env in &events {
match &env.payload {
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal {
pass,
reader_seq,
advertised_writer_seq,
failure_reason,
..
}) => {
final_found = true;
assert!(!pass);
assert_eq!(*reader_seq, SeqNo(1));
assert_eq!(*advertised_writer_seq, Some(SeqNo(3)));
match failure_reason {
Some(EventViolationCause::SeqDivergence { advertised, reader }) => {
assert_eq!(*advertised, Some(SeqNo(3)));
assert_eq!(*reader, SeqNo(1));
}
other => panic!("expected SeqDivergence failure_reason, got {other:?}"),
}
}
ChainPayload::FlowControl(FlowControlPayload::ConsumptionGap {
from_seq,
to_seq,
upstream,
}) => {
gap_found = true;
assert_eq!(*from_seq, SeqNo(2));
assert_eq!(*to_seq, SeqNo(3));
assert_eq!(*upstream, upstream_stage);
}
ChainPayload::FlowControl(FlowControlPayload::AtLeastOnceViolation {
upstream,
reason,
reader_seq,
advertised_writer_seq,
}) => {
violation_found = true;
assert_eq!(*upstream, upstream_stage);
assert_eq!(*reader_seq, SeqNo(1));
assert_eq!(*advertised_writer_seq, Some(SeqNo(3)));
match reason {
EventViolationCause::SeqDivergence { advertised, reader } => {
assert_eq!(*advertised, Some(SeqNo(3)));
assert_eq!(*reader, SeqNo(1));
}
other => panic!(
"expected SeqDivergence reason in AtLeastOnceViolation, got {other:?}"
),
}
}
_ => {}
}
}
assert!(final_found, "expected a ConsumptionFinal event");
assert!(gap_found, "expected a ConsumptionGap event");
assert!(
violation_found,
"expected an AtLeastOnceViolation event for SeqDivergence"
);
let system_events = system_journal.read_causally_ordered().await.unwrap();
let mut status_found = false;
for env in &system_events {
if let SystemPayload::ContractStatus {
upstream,
reader,
pass,
reason,
..
} = &env.payload
{
if *pass {
continue;
}
status_found = true;
assert_eq!(*upstream, upstream_stage);
assert_eq!(*reader, reader_stage);
match reason {
Some(EventViolationCause::SeqDivergence { .. }) => {}
other => {
panic!("expected SeqDivergence reason in ContractStatus, got {other:?}")
}
}
}
}
assert!(status_found, "expected ContractStatus system event");
}
#[tokio::test]
async fn transport_only_skips_observability_events() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let writer_id = WriterId::Stage(upstream_stage);
upstream_journal
.append(
ChainEventFactory::stage_running(writer_id, upstream_stage),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::stage_running(writer_id, upstream_stage),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.event", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::stage_running(writer_id, upstream_stage),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::eof_event(writer_id, true),
Default::default(),
)
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
for _ in 0..2 {
match subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await
{
PollResult::CursorAdvanced {
upstream,
completed_data_rows,
} => {
assert_eq!(upstream, upstream_stage);
assert_eq!(completed_data_rows, 0, "observability is not Data");
}
other => panic!("expected bounded observability cursor progress, got {other:?}"),
}
}
let data = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match data {
PollResult::Event(env) => match env.payload {
ChainPayload::Fact(_) => {}
other => panic!("expected first delivered event to be Data, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
match subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await
{
PollResult::CursorAdvanced {
upstream,
completed_data_rows,
} => {
assert_eq!(upstream, upstream_stage);
assert_eq!(completed_data_rows, 0);
}
other => panic!("expected trailing observability cursor progress, got {other:?}"),
}
let eof = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match eof {
PollResult::Event(env) => match env.payload {
ChainPayload::FlowControl(FlowControlPayload::Eof { .. }) => {}
other => panic!("expected second delivered event to be EOF, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
let outcome = subscription
.take_last_eof_outcome()
.expect("expected subscription to mark authoritative EOF");
assert!(outcome.is_final);
assert_eq!(outcome.stage_id, upstream_stage);
assert_eq!(outcome.reader_index, 0);
assert_eq!(outcome.eof_count, 1);
assert_eq!(outcome.total_readers, 1);
}
#[tokio::test]
async fn transport_only_filters_unselected_data_and_reconciles_selected_writer_seq() {
let upstream_stage = StageId::new();
let reader_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let reader_owner = JournalOwner::stage(reader_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(reader_owner.clone()));
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(reader_owner));
let writer_id = WriterId::Stage(upstream_stage);
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.ignored.v1", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.selected.v1", json!({"n": 2})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.selected", json!({"n": 3})),
Default::default(),
)
.await
.unwrap();
let mut eof_event = ChainEventFactory::eof_event(writer_id, true);
if let ChainPayload::FlowControl(FlowControlPayload::Eof {
writer_seq,
writer_seq_by_event_type,
..
}) = &mut eof_event.payload
{
*writer_seq = Some(SeqNo(3));
writer_seq_by_event_type.insert("test.ignored.v1".into(), SeqNo(1));
writer_seq_by_event_type.insert("test.selected.v1".into(), SeqNo(1));
writer_seq_by_event_type.insert("test.selected".into(), SeqNo(1));
}
upstream_journal
.append(eof_event, Default::default())
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut selected_feeds = HashMap::new();
selected_feeds.insert(
upstream_stage,
vec![
SelectedFeedMetadata::new(EventType::from("test.selected.v1"), SelectedFeedRole::Input),
SelectedFeedMetadata::new(
EventType::from("test.second_selected.v1"),
SelectedFeedRole::Stream,
),
],
);
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.with_selected_feeds(selected_feeds)
.with_contracts(ContractsWiring {
writer_id: WriterId::from(reader_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
})
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let first = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match first {
PollResult::CursorAdvanced {
upstream,
completed_data_rows,
} => {
assert_eq!(upstream, upstream_stage);
assert_eq!(completed_data_rows, 1);
}
other => panic!("expected physical cursor completion, got {other:?}"),
}
assert_eq!(reader_progress[0].reader_seq, SeqNo(0));
assert_eq!(reader_progress[0].receipted_seq, SeqNo(0));
assert!(reader_progress[0].last_event_id.is_none());
assert_eq!(subscription.delivered_data_count(), 0);
assert_eq!(subscription.delivered_counts()[0], DeliveredCount(0));
assert_eq!(subscription.last_delivered_stage_input_position(), None);
let selected = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match selected {
PollResult::Event(env) => match env.payload {
ChainPayload::Fact(_) => {
assert_eq!(env.event_type(), "test.selected.v1");
}
other => panic!("expected selected Data event, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
assert_eq!(reader_progress[0].reader_seq, SeqNo(1));
assert_eq!(
subscription
.last_delivered_stage_input_position()
.expect("selected data should receive a stage input position")
.0,
1
);
drive_subscription_to_eof(&mut subscription, &mut reader_progress).await;
assert_eq!(reader_progress[0].reader_seq, SeqNo(2));
assert_eq!(reader_progress[0].advertised_writer_seq, Some(SeqNo(2)));
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(
status,
ContractStatus::ProgressEmitted | ContractStatus::Healthy
),
"selected feed should reconcile successfully, got {status:?}"
);
assert!(!reader_progress[0].contract_violated);
let contract_events = contract_journal.read_causally_ordered().await.unwrap();
let mut final_found = false;
for env in &contract_events {
match &env.payload {
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal {
pass,
reader_seq,
advertised_writer_seq,
failure_reason,
..
}) => {
final_found = true;
assert!(*pass);
assert_eq!(*reader_seq, SeqNo(2));
assert_eq!(*advertised_writer_seq, Some(SeqNo(2)));
assert!(failure_reason.is_none());
}
ChainPayload::FlowControl(
FlowControlPayload::ConsumptionGap { .. }
| FlowControlPayload::AtLeastOnceViolation { .. },
) => panic!("selected feed should not emit gap or violation events"),
_ => {}
}
}
assert!(final_found, "expected a selected-feed ConsumptionFinal");
let system_events = system_journal.read_causally_ordered().await.unwrap();
assert!(system_events.iter().any(|env| matches!(
&env.payload,
SystemPayload::ContractStatus {
upstream,
reader,
selected_event_type,
feed_role,
pass,
reader_seq,
advertised_writer_seq,
reason,
} if *upstream == upstream_stage
&& *reader == reader_stage
&& selected_event_type.as_ref().map(|event_type| event_type.as_str())
== Some("test.selected.v1")
&& feed_role.as_ref().map(|role| role.as_str()) == Some("input")
&& *pass
&& *reader_seq == Some(SeqNo(2))
&& *advertised_writer_seq == Some(SeqNo(2))
&& reason.is_none()
)));
}
#[tokio::test]
async fn contract_prefix_resolves_replay_alias_and_excludes_forwarded_rows_symmetrically() {
use obzenflow_core::event::payloads::flow_control_payload::EofKind;
use obzenflow_core::event::provenance::replay_context::ReplayContext;
use obzenflow_core::event::status::processing_status::ProcessingStatus;
let upstream_stage = StageId::new();
let archived_upstream_stage = StageId::new();
let lineage_root_stage = StageId::new();
let foreign_author = StageId::new();
let reader_stage = StageId::new();
let upstream_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(upstream_stage)));
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(reader_stage)));
for n in 0..2 {
let mut replayed = ChainEventFactory::data_event(
WriterId::Stage(lineage_root_stage),
"test.joined.v1",
json!({"n": n}),
);
replayed.replay_context = Some(ReplayContext {
original_event_id: replayed.id,
original_flow_id: "flow_parent".to_string(),
original_stage_id: archived_upstream_stage,
});
upstream_journal
.append(replayed, Default::default())
.await
.unwrap();
}
let mut forwarded_error = ChainEventFactory::data_event(
WriterId::Stage(foreign_author),
"test.joined.v1",
json!({"foreign": true}),
);
forwarded_error.processing.status = ProcessingStatus::error("forwarded pre-error row");
upstream_journal
.append(forwarded_error, Default::default())
.await
.unwrap();
let mut forwarded_eof =
ChainEventFactory::eof_event_with_kind(WriterId::Stage(foreign_author), EofKind::Natural);
if let ChainPayload::FlowControl(FlowControlPayload::Eof { writer_seq, .. }) =
&mut forwarded_eof.payload
{
*writer_seq = Some(SeqNo(3));
}
upstream_journal
.append(forwarded_eof, Default::default())
.await
.unwrap();
let mut local_terminal =
ChainEventFactory::eof_event_with_kind(WriterId::Stage(upstream_stage), EofKind::Poison);
if let ChainPayload::FlowControl(FlowControlPayload::Eof {
writer_seq,
writer_seq_by_event_type,
..
}) = &mut local_terminal.payload
{
*writer_seq = Some(SeqNo(2));
writer_seq_by_event_type.insert("test.joined.v1".into(), SeqNo(2));
}
upstream_journal
.append(local_terminal, Default::default())
.await
.unwrap();
let upstreams = [(upstream_stage, "join".to_string(), upstream_journal)];
let mut selected_feeds = HashMap::new();
selected_feeds.insert(
upstream_stage,
vec![SelectedFeedMetadata::new(
EventType::from("test.joined.v1"),
SelectedFeedRole::Input,
)],
);
let mut subscription = UpstreamSubscription::new_with_names("consumer", &upstreams)
.await
.unwrap()
.with_archived_stage_ids(HashMap::from([(upstream_stage, archived_upstream_stage)]))
.with_selected_feeds(selected_feeds)
.with_contracts(ContractsWiring {
writer_id: WriterId::from(reader_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: None,
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
})
.transport_only();
let mut progress = [ReaderProgress::new(upstream_stage)];
let mut saw_foreign_error = false;
let mut saw_forwarded_eof = false;
loop {
match subscription
.poll_next_with_state("test_fsm", Some(&mut progress))
.await
{
PollResult::Event(envelope) => {
saw_foreign_error |= envelope.envelope.provenance.event.writer_id
== WriterId::Stage(foreign_author)
&& envelope.consumes_data_credit();
saw_forwarded_eof |= envelope.envelope.provenance.event.writer_id
== WriterId::Stage(foreign_author)
&& matches!(
envelope.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
);
}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(error) => panic!("unexpected subscription error: {error}"),
}
}
assert!(
saw_foreign_error,
"forwarded pre-error Data remains deliverable"
);
assert!(saw_forwarded_eof, "forwarded EOF remains deliverable");
assert_eq!(progress[0].reader_seq, SeqNo(2));
assert_eq!(progress[0].advertised_writer_seq, Some(SeqNo(2)));
let status = subscription.check_contracts(&mut progress).await;
assert!(
matches!(
status,
ContractStatus::ProgressEmitted | ContractStatus::Healthy
),
"owner-authored prefix should reconcile, got {status:?}"
);
assert!(!progress[0].contract_violated);
let contract_events = contract_journal.read_causally_ordered().await.unwrap();
let final_event = contract_events.iter().find(|envelope| {
matches!(
envelope.payload,
ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal { .. })
)
});
let Some(final_event) = final_event else {
panic!("expected ConsumptionFinal")
};
let ChainPayload::FlowControl(FlowControlPayload::ConsumptionFinal {
pass,
reader_seq,
advertised_writer_seq,
failure_reason,
..
}) = &final_event.payload
else {
unreachable!()
};
assert!(*pass);
assert_eq!(*reader_seq, SeqNo(2));
assert_eq!(*advertised_writer_seq, Some(SeqNo(2)));
assert!(failure_reason.is_none());
}
#[tokio::test]
async fn matching_input_boundary_delivery_stamps_exact_replayable_activation() {
let upstream_stage = StageId::new();
let upstream_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(upstream_stage)));
let writer_id = WriterId::Stage(upstream_stage);
let mut ignored = ChainEventFactory::data_event(writer_id, "checkout.other.v1", json!({}));
ignored.processing.event_time = 100;
upstream_journal
.append(ignored, Default::default())
.await
.unwrap();
let mut admitted = ChainEventFactory::data_event(writer_id, "checkout.command.v1", json!({}));
admitted.processing.event_time = 123;
let admitted_id = admitted.id;
upstream_journal
.append(admitted, Default::default())
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut entries = HashMap::new();
entries.insert(
upstream_stage,
vec![CompositeEntrySpec {
composite_id: CompositeId::new("saga:checkout"),
port_name: "commands".to_string(),
event_types: vec![EventType::from("checkout.command.v1")],
}],
);
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.with_composite_entries(entries);
let mut progress = [ReaderProgress::new(upstream_stage)];
let PollResult::Event(ignored) = subscription
.poll_next_with_state("test_fsm", Some(&mut progress))
.await
else {
panic!("expected ignored Data delivery");
};
assert!(ignored.composite_activations().is_empty());
let PollResult::Event(admitted) = subscription
.poll_next_with_state("test_fsm", Some(&mut progress))
.await
else {
panic!("expected boundary Data delivery");
};
assert_eq!(admitted.composite_activations().len(), 1);
let activation = &admitted.composite_activations()[0];
assert_eq!(activation.composite_id, CompositeId::new("saga:checkout"));
assert_eq!(activation.activation, admitted_id);
assert_eq!(activation.entry_port, "commands");
assert_eq!(activation.entered_at_ms, 123);
}
#[tokio::test]
async fn multi_selected_feeds_emit_direct_contract_status_per_feed() {
let upstream_stage = StageId::new();
let reader_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let reader_owner = JournalOwner::stage(reader_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(reader_owner.clone()));
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(reader_owner));
let writer_id = WriterId::Stage(upstream_stage);
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.first.v1", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.second.v1", json!({"n": 2})),
Default::default(),
)
.await
.unwrap();
let mut eof_event = ChainEventFactory::eof_event(writer_id, true);
if let ChainPayload::FlowControl(FlowControlPayload::Eof {
writer_seq,
writer_seq_by_event_type,
..
}) = &mut eof_event.payload
{
*writer_seq = Some(SeqNo(2));
writer_seq_by_event_type.insert("test.first.v1".into(), SeqNo(2));
writer_seq_by_event_type.insert("test.second.v1".into(), SeqNo(0));
}
upstream_journal
.append(eof_event, Default::default())
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut selected_feeds = HashMap::new();
selected_feeds.insert(
upstream_stage,
vec![
SelectedFeedMetadata::new(EventType::from("test.first.v1"), SelectedFeedRole::Input),
SelectedFeedMetadata::new(
EventType::from("test.second.v1"),
SelectedFeedRole::Reference,
),
],
);
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.with_selected_feeds(selected_feeds)
.with_contracts(ContractsWiring {
writer_id: WriterId::from(reader_stage),
contract_journal: contract_journal.clone(),
config: ContractConfig::default(),
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
})
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
drive_subscription_to_eof(&mut subscription, &mut reader_progress).await;
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(
matches!(status, ContractStatus::Violated { upstream, .. } if upstream == upstream_stage),
"per-feed divergence should fail even when aggregate selected counts reconcile, got {status:?}"
);
let system_events = system_journal.read_causally_ordered().await.unwrap();
let mut first_status = None;
let mut second_status = None;
let mut aggregate_status_found = false;
for env in &system_events {
if let SystemPayload::ContractStatus {
upstream,
reader,
selected_event_type,
feed_role,
pass,
reader_seq,
advertised_writer_seq,
reason,
} = &env.payload
{
if *upstream != upstream_stage || *reader != reader_stage {
continue;
}
match (
selected_event_type
.as_ref()
.map(|event_type| event_type.as_str()),
feed_role.as_ref().map(|role| role.as_str()),
) {
(Some("test.first.v1"), Some("input")) => {
first_status =
Some((*pass, *reader_seq, *advertised_writer_seq, reason.clone()));
}
(Some("test.second.v1"), Some("reference")) => {
second_status =
Some((*pass, *reader_seq, *advertised_writer_seq, reason.clone()));
}
(None, None) => {
aggregate_status_found = true;
}
_ => {}
}
}
}
let first_status = first_status.expect("expected first feed ContractStatus");
assert!(!first_status.0);
assert_eq!(first_status.1, Some(SeqNo(1)));
assert_eq!(first_status.2, Some(SeqNo(2)));
assert!(matches!(
first_status.3,
Some(EventViolationCause::SeqDivergence {
advertised: Some(SeqNo(2)),
reader: SeqNo(1),
})
));
let second_status = second_status.expect("expected second feed ContractStatus");
assert!(!second_status.0);
assert_eq!(second_status.1, Some(SeqNo(1)));
assert_eq!(second_status.2, Some(SeqNo(0)));
assert!(matches!(
second_status.3,
Some(EventViolationCause::SeqDivergence {
advertised: Some(SeqNo(0)),
reader: SeqNo(1),
})
));
assert!(
!aggregate_status_found,
"direct feed status should suppress ambiguous aggregate ContractStatus"
);
}
#[tokio::test]
async fn multi_selected_feeds_emit_midflight_contract_results_per_feed() {
let upstream_stage = StageId::new();
let reader_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let reader_owner = JournalOwner::stage(reader_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(reader_owner.clone()));
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(TestJournal::new(reader_owner));
let writer_id = WriterId::Stage(upstream_stage);
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.first.v1", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.second.v1", json!({"n": 2})),
Default::default(),
)
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut selected_feeds = HashMap::new();
selected_feeds.insert(
upstream_stage,
vec![
SelectedFeedMetadata::new(EventType::from("test.first.v1"), SelectedFeedRole::Input),
SelectedFeedMetadata::new(
EventType::from("test.second.v1"),
SelectedFeedRole::Reference,
),
],
);
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.with_selected_feeds(selected_feeds)
.with_contracts(ContractsWiring {
writer_id: WriterId::from(reader_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: Some(system_journal.clone()),
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
})
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let first = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
assert!(matches!(
first,
PollResult::Event(ref record) if record.is_fact() && record.event_type() == "test.first.v1"
));
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(matches!(
status,
ContractStatus::ProgressEmitted | ContractStatus::Healthy
));
let events_after_first = system_journal.read_causally_ordered().await.unwrap();
let first_feed_results = events_after_first
.iter()
.filter(|env| {
matches!(
&env.payload,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type,
feed_role,
contract_name,
status,
reader_seq,
advertised_writer_seq,
..
} if *upstream == upstream_stage
&& *reader == reader_stage
&& selected_event_type.as_ref().map(|event_type| event_type.as_str())
== Some("test.first.v1")
&& feed_role.as_ref().map(|role| role.as_str()) == Some("input")
&& contract_name.as_str() == TransportContract::NAME
&& *status == ContractResultStatusLabel::Healthy
&& *reader_seq == Some(SeqNo(1))
&& advertised_writer_seq.is_none()
)
})
.count();
let second_feed_results_after_first = events_after_first
.iter()
.filter(|env| {
matches!(
&env.payload,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type,
feed_role,
..
} if *upstream == upstream_stage
&& *reader == reader_stage
&& selected_event_type.as_ref().map(|event_type| event_type.as_str())
== Some("test.second.v1")
&& feed_role.as_ref().map(|role| role.as_str()) == Some("reference")
)
})
.count();
let aggregate_results_after_first = events_after_first
.iter()
.filter(|env| {
matches!(
&env.payload,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type: None,
feed_role: None,
..
} if *upstream == upstream_stage && *reader == reader_stage
)
})
.count();
assert_eq!(first_feed_results, 1);
assert_eq!(second_feed_results_after_first, 0);
assert_eq!(aggregate_results_after_first, 0);
let second = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
assert!(matches!(
second,
PollResult::Event(ref record) if record.is_fact() && record.event_type() == "test.second.v1"
));
let status = subscription.check_contracts(&mut reader_progress).await;
assert!(matches!(
status,
ContractStatus::ProgressEmitted | ContractStatus::Healthy
));
let events_after_second = system_journal.read_causally_ordered().await.unwrap();
let first_feed_results_after_second = events_after_second
.iter()
.filter(|env| {
matches!(
&env.payload,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type,
feed_role,
..
} if *upstream == upstream_stage
&& *reader == reader_stage
&& selected_event_type.as_ref().map(|event_type| event_type.as_str())
== Some("test.first.v1")
&& feed_role.as_ref().map(|role| role.as_str()) == Some("input")
)
})
.count();
let second_feed_results = events_after_second
.iter()
.filter(|env| {
matches!(
&env.payload,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type,
feed_role,
contract_name,
status,
reader_seq,
advertised_writer_seq,
..
} if *upstream == upstream_stage
&& *reader == reader_stage
&& selected_event_type.as_ref().map(|event_type| event_type.as_str())
== Some("test.second.v1")
&& feed_role.as_ref().map(|role| role.as_str()) == Some("reference")
&& contract_name.as_str() == TransportContract::NAME
&& *status == ContractResultStatusLabel::Healthy
&& *reader_seq == Some(SeqNo(1))
&& advertised_writer_seq.is_none()
)
})
.count();
assert_eq!(first_feed_results_after_second, 1);
assert_eq!(second_feed_results, 1);
}
#[tokio::test]
async fn transport_only_skips_framework_effect_data_without_stage_input_position() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let writer_id = WriterId::Stage(upstream_stage);
let effect_record = obzenflow_core::event::payloads::effect_payload::EffectRecord {
cursor: obzenflow_core::event::payloads::effect_payload::EffectCursor::new(
"recorded-flow",
"gateway",
1,
0,
),
descriptor_hash: "hash".into(),
descriptor: obzenflow_core::event::payloads::effect_payload::EffectDescriptor::new(
"test.effect",
"test",
1,
"1",
"input-hash",
),
outcome: obzenflow_core::event::payloads::effect_payload::EffectOutcomePayload::Succeeded {
output: json!({"ok": true}),
},
origin: None,
};
upstream_journal
.append(
ChainEventFactory::derived_event(
writer_id,
&ChainEventFactory::data_event(writer_id, "test.parent", json!({})),
ChainPayload::Execution(ExecutionPayload::EffectRecord(effect_record.clone())),
obzenflow_core::config::LineagePolicy::default(),
)
.with_effect_provenance(EffectProvenance::from_record(
&effect_record,
EffectFactOwner::Framework,
)),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(writer_id, "test.event", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let filtered = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match filtered {
PollResult::CursorAdvanced {
upstream,
completed_data_rows,
} => {
assert_eq!(upstream, upstream_stage);
assert_eq!(completed_data_rows, 1);
}
other => panic!("expected framework effect physical completion, got {other:?}"),
}
assert_eq!(subscription.last_delivered_stage_input_position(), None);
assert_eq!(reader_progress[0].reader_seq, SeqNo(0));
let first = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match first {
PollResult::Event(env) => match env.payload {
ChainPayload::Fact(_) => {}
other => {
panic!("expected framework effect Data to be skipped and domain Data delivered, got {other:?}")
}
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
assert_eq!(
subscription
.last_delivered_stage_input_position()
.expect("data event should receive a stage input position")
.0,
1
);
let second = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
assert!(matches!(second, PollResult::NoEvents));
}
#[tokio::test]
async fn forwarded_eof_with_missing_writer_is_not_terminal() {
let upstream_stage = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
let upstream_writer_id = WriterId::Stage(upstream_stage);
let forwarded_writer_id = WriterId::Stage(StageId::new());
upstream_journal
.append(
ChainEventFactory::data_event(upstream_writer_id, "test.event", json!({"n": 1})),
Default::default(),
)
.await
.unwrap();
let mut forwarded_eof = ChainEventFactory::eof_event(forwarded_writer_id, true);
if let ChainPayload::FlowControl(FlowControlPayload::Eof { writer_id, .. }) =
&mut forwarded_eof.payload
{
*writer_id = None;
}
upstream_journal
.append(forwarded_eof, Default::default())
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::data_event(upstream_writer_id, "test.event", json!({"n": 2})),
Default::default(),
)
.await
.unwrap();
upstream_journal
.append(
ChainEventFactory::eof_event(upstream_writer_id, true),
Default::default(),
)
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let first = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match first {
PollResult::Event(env) => match env.payload {
ChainPayload::Fact(_) => {}
other => panic!("expected first delivered event to be Data, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
let second = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match second {
PollResult::Event(env) => match env.payload {
ChainPayload::FlowControl(FlowControlPayload::Eof { .. }) => {}
other => panic!("expected second delivered event to be EOF, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
assert!(
subscription.take_last_eof_outcome().is_none(),
"expected forwarded EOF with missing writer_id to be non-terminal"
);
let third = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match third {
PollResult::Event(env) => match env.payload {
ChainPayload::Fact(_) => {}
other => panic!("expected third delivered event to be Data, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
let fourth = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
match fourth {
PollResult::Event(env) => match env.payload {
ChainPayload::FlowControl(FlowControlPayload::Eof { .. }) => {}
other => panic!("expected fourth delivered event to be EOF, got {other:?}"),
},
other => panic!("expected PollResult::Event, got {other:?}"),
}
let outcome = subscription
.take_last_eof_outcome()
.expect("expected subscription to mark authoritative EOF");
assert!(outcome.is_final);
assert_eq!(outcome.stage_id, upstream_stage);
assert_eq!(outcome.eof_count, 1);
assert_eq!(outcome.total_readers, 1);
}
struct SharedTestJournal {
id: JournalId,
owner: Option<JournalOwner>,
events: Arc<Mutex<Vec<JournalRecord<ChainPayload>>>>,
}
impl SharedTestJournal {
fn new(owner: JournalOwner) -> Self {
Self {
id: JournalId::new(),
owner: Some(owner),
events: Arc::new(Mutex::new(Vec::new())),
}
}
fn append_with_clock(&self, event: ChainEvent, vector_clock: VectorClock) {
let envelope = JournalRecord::commit_event(
event,
JournalProvenance {
journal_writer_id: JournalWriterId::from(self.id),
vector_clock,
timestamp: chrono::Utc::now(),
journal_group_id: None,
journal_group_member: None,
},
)
.expect("valid committed fixture");
self.events.lock().unwrap().push(envelope);
}
}
struct SharedTestJournalReader {
events: Arc<Mutex<Vec<JournalRecord<ChainPayload>>>>,
pos: usize,
}
#[async_trait]
impl JournalReader<ChainEvent> for SharedTestJournalReader {
async fn next(
&mut self,
) -> std::result::Result<Option<JournalRecord<ChainPayload>>, JournalError> {
let guard = self.events.lock().unwrap();
if self.pos < guard.len() {
let envelope = guard[self.pos].clone();
self.pos += 1;
Ok(Some(envelope))
} else {
Ok(None)
}
}
fn position(&self) -> u64 {
self.pos as u64
}
fn is_at_end(&self) -> bool {
self.pos >= self.events.lock().unwrap().len()
}
}
#[async_trait]
impl Journal<ChainEvent> for SharedTestJournal {
fn id(&self) -> &JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
event: ChainEvent,
mut options: AppendOptions<'_, ChainEvent>,
) -> std::result::Result<JournalRecord<ChainPayload>, JournalError> {
let event = options.capture.prepare(0, event);
let envelope = JournalRecord::new(JournalWriterId::from(self.id), event);
self.events.lock().unwrap().push(envelope.clone());
Ok(envelope)
}
async fn read_all_unordered(
&self,
) -> std::result::Result<Vec<JournalRecord<ChainPayload>>, JournalError> {
Ok(self.events.lock().unwrap().clone())
}
async fn read_event(
&self,
_event_id: &obzenflow_core::EventId,
) -> std::result::Result<Option<JournalRecord<ChainPayload>>, JournalError> {
Ok(None)
}
async fn reader_from(
&self,
position: u64,
) -> std::result::Result<Box<dyn JournalReader<ChainEvent>>, JournalError> {
Ok(Box::new(SharedTestJournalReader {
events: self.events.clone(),
pos: position as usize,
}))
}
async fn read_last_n(
&self,
count: usize,
) -> std::result::Result<Vec<JournalRecord<ChainPayload>>, JournalError> {
let guard = self.events.lock().unwrap();
let len = guard.len();
let start = len.saturating_sub(count);
Ok(guard[start..].iter().rev().cloned().collect())
}
}
fn merge_data(writer: StageId, event_type: &str) -> ChainEvent {
ChainEventFactory::data_event(WriterId::Stage(writer), event_type, json!({}))
}
fn merge_authored_eof(writer: StageId) -> ChainEvent {
ChainEventFactory::eof_event(WriterId::Stage(writer), true)
}
fn merge_consumption_progress(writer: StageId, seq: u64) -> ChainEvent {
ChainEventFactory::consumption_progress_event(
WriterId::Stage(writer),
obzenflow_core::event::ConsumptionProgressEventParams {
reader_seq: SeqNo(seq),
last_event_id: None,
vector_clock: None,
eof_seen: false,
reader_path: obzenflow_core::event::types::JournalPath("upstream".to_string()),
reader_index: obzenflow_core::event::types::JournalIndex(0),
advertised_writer_seq: None,
advertised_vector_clock: None,
stalled_since: None,
},
)
}
fn merge_clock(entries: &[(&str, u64)]) -> VectorClock {
let mut clock = VectorClock::new();
for (writer, seq) in entries {
for _ in 0..*seq {
obzenflow_core::event::vector_clock::CausalOrderingService::increment(
&mut clock, writer,
);
}
}
clock
}
fn merge_catch_up(writer: StageId, generation: u64, stage_key: &str) -> ChainEvent {
ChainEventFactory::source_event(
WriterId::Stage(writer),
ChainPayload::FlowControl(FlowControlPayload::CatchUpComplete {
generation: ReaderGeneration(generation),
stage_key: obzenflow_core::StageKey::from(stage_key),
}),
)
}
async fn canonical_pair(
name_a: &str,
name_b: &str,
) -> (
UpstreamSubscription<ChainEvent>,
(StageId, Arc<SharedTestJournal>),
(StageId, Arc<SharedTestJournal>),
) {
let stage_a = StageId::new();
let stage_b = StageId::new();
let journal_a = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_a)));
let journal_b = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_b)));
let upstreams = [
(
stage_a,
name_a.to_string(),
journal_a.clone() as Arc<dyn Journal<ChainEvent>>,
),
(
stage_b,
name_b.to_string(),
journal_b.clone() as Arc<dyn Journal<ChainEvent>>,
),
];
let subscription = UpstreamSubscription::new_with_names("merge_owner", &upstreams)
.await
.unwrap()
.with_reader_selection(ReaderSelectionPolicy::CanonicalMerge);
(subscription, (stage_a, journal_a), (stage_b, journal_b))
}
async fn canonical_pair_transport_only(
name_a: &str,
name_b: &str,
) -> (
UpstreamSubscription<ChainEvent>,
(StageId, Arc<SharedTestJournal>),
(StageId, Arc<SharedTestJournal>),
) {
let (subscription, a, b) = canonical_pair(name_a, name_b).await;
(subscription.transport_only(), a, b)
}
async fn canonical_single(
name: &str,
) -> (
UpstreamSubscription<ChainEvent>,
StageId,
Arc<SharedTestJournal>,
) {
let stage = StageId::new();
let journal = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage)));
let upstreams = [(
stage,
name.to_string(),
journal.clone() as Arc<dyn Journal<ChainEvent>>,
)];
let subscription = UpstreamSubscription::new_with_names("merge_owner", &upstreams)
.await
.unwrap()
.with_reader_selection(ReaderSelectionPolicy::CanonicalMerge);
(subscription, stage, journal)
}
async fn expect_delivery(
subscription: &mut UpstreamSubscription<ChainEvent>,
) -> JournalRecord<ChainPayload> {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => envelope,
other => panic!(
"expected PollResult::Event, got {other:?} (delivered={:?}, merge_wait={:?})",
subscription.delivered_counts(),
subscription.merge_wait(),
),
}
}
async fn expect_poll_error(subscription: &mut UpstreamSubscription<ChainEvent>) -> String {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Error(error) => error.to_string(),
other => panic!("expected PollResult::Error, got {other:?}"),
}
}
fn delivered(subscription: &UpstreamSubscription<ChainEvent>) -> Vec<u64> {
subscription
.delivered_counts()
.iter()
.map(|count| count.0)
.collect()
}
#[tokio::test]
async fn canonical_merge_waits_while_any_input_is_quiet() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a.event"), VectorClock::new());
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while an input is quiet, got {other:?}"),
}
let wait = subscription.merge_wait().expect("merge wait recorded");
assert_eq!(wait.quiet_inputs, vec![(stage_b, "upstream_b".to_string())]);
assert_eq!(delivered(&subscription), [0, 0]);
journal_b.append_with_clock(merge_data(stage_b, "b.event"), VectorClock::new());
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "a.event");
assert_eq!(subscription.last_delivered_upstream_stage(), Some(stage_a));
assert!(subscription.merge_wait().is_none());
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while A is quiet, got {other:?}"),
}
let wait = subscription.merge_wait().expect("merge wait recorded");
assert_eq!(wait.quiet_inputs, vec![(stage_a, "upstream_a".to_string())]);
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
let second = expect_delivery(&mut subscription).await;
assert_eq!(second.event_type(), "b.event");
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while B is quiet, got {other:?}"),
}
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let third = expect_delivery(&mut subscription).await;
assert!(matches!(
third.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
));
let fourth = expect_delivery(&mut subscription).await;
assert!(matches!(
fourth.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
));
assert_eq!(delivered(&subscription), [2, 2]);
assert!(subscription.all_readers_eof());
}
#[tokio::test]
async fn canonical_merge_alternates_by_ordinal_then_stage_key_and_never_waits_on_sealed_inputs() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
for label in ["a1", "a2", "a3"] {
journal_a.append_with_clock(merge_data(stage_a, label), VectorClock::new());
}
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
for label in ["b1", "b2", "b3"] {
journal_b.append_with_clock(merge_data(stage_b, label), VectorClock::new());
}
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut delivered_from = Vec::new();
for _ in 0..8 {
let _ = expect_delivery(&mut subscription).await;
delivered_from.push(subscription.last_delivered_upstream_stage().unwrap());
}
assert_eq!(
delivered_from,
vec![stage_a, stage_b, stage_a, stage_b, stage_a, stage_b, stage_a, stage_b],
"ordinal-balanced tiebreak alternates fairly, stage key breaks ties"
);
assert_eq!(delivered(&subscription), [4, 4]);
assert!(subscription.all_readers_eof());
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents after exhaustion, got {other:?}"),
}
}
#[tokio::test]
async fn canonical_merge_excludes_happened_before_heads() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("z_upstream", "a_upstream").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), merge_clock(&[("wa", 1)]));
journal_a.append_with_clock(merge_data(stage_a, "a2"), merge_clock(&[("wa", 2)]));
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
journal_b.append_with_clock(
merge_data(stage_b, "b1"),
merge_clock(&[("wa", 2), ("wb", 1)]),
);
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut order = Vec::new();
for _ in 0..3 {
let envelope = expect_delivery(&mut subscription).await;
order.push(envelope.event_type());
}
assert_eq!(
order,
vec!["a1".to_string(), "a2".to_string(), "b1".to_string()],
"causal ancestors deliver before the head that derives from them"
);
}
#[tokio::test]
async fn canonical_merge_authored_eof_orders_by_tiebreak_and_exhausts_reader() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
for label in ["b1", "b2", "b3"] {
journal_b.append_with_clock(merge_data(stage_b, label), VectorClock::new());
}
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut delivered_from = Vec::new();
for _ in 0..6 {
let _ = expect_delivery(&mut subscription).await;
delivered_from.push(subscription.last_delivered_upstream_stage().unwrap());
}
assert_eq!(
delivered_from,
vec![stage_a, stage_b, stage_a, stage_b, stage_b, stage_b]
);
assert!(subscription.all_readers_eof());
let outcome = subscription
.take_last_eof_outcome()
.expect("final EOF outcome recorded");
assert!(outcome.is_final);
assert_eq!(outcome.stage_id, stage_b);
}
fn merge_authored_eof_with_kind(
writer: StageId,
kind: obzenflow_core::event::payloads::flow_control_payload::EofKind,
) -> ChainEvent {
ChainEventFactory::eof_event_with_kind(WriterId::Stage(writer), kind)
}
#[tokio::test]
async fn eof_kind_fold_is_worst_wins_under_both_arrival_orders() {
use obzenflow_core::event::payloads::flow_control_payload::EofKind;
for (kind_a, kind_b) in [
(EofKind::Natural, EofKind::Truncated),
(EofKind::Truncated, EofKind::Natural),
(EofKind::Natural, EofKind::Poison),
(EofKind::Truncated, EofKind::Poison),
(EofKind::Poison, EofKind::Truncated),
(EofKind::Natural, EofKind::Natural),
] {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(
merge_authored_eof_with_kind(stage_a, kind_a),
VectorClock::new(),
);
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(
merge_authored_eof_with_kind(stage_b, kind_b),
VectorClock::new(),
);
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(_) => {}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert!(subscription.all_readers_eof());
let outcome = subscription
.take_last_eof_outcome()
.expect("final EOF outcome recorded");
assert!(outcome.is_final);
assert_eq!(
outcome.worst_kind,
Some(kind_a.worst(kind_b)),
"fold of {kind_a:?} and {kind_b:?} must be the worst-wins join"
);
}
}
#[tokio::test]
async fn eof_kind_fold_single_input_is_the_identity() {
use obzenflow_core::event::payloads::flow_control_payload::EofKind;
let (mut subscription, stage, journal) = canonical_single("upstream").await;
journal.append_with_clock(merge_data(stage, "a1"), VectorClock::new());
journal.append_with_clock(
merge_authored_eof_with_kind(stage, EofKind::Truncated),
VectorClock::new(),
);
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(_) => {}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
let outcome = subscription
.take_last_eof_outcome()
.expect("final EOF outcome recorded");
assert_eq!(outcome.worst_kind, Some(EofKind::Truncated));
}
#[tokio::test]
async fn eof_kind_fold_keeps_the_worst_on_duplicate_eofs_from_one_input() {
use obzenflow_core::event::payloads::flow_control_payload::EofKind;
let (mut subscription, stage, journal) = canonical_single("upstream").await;
journal.append_with_clock(
merge_authored_eof_with_kind(stage, EofKind::Poison),
VectorClock::new(),
);
journal.append_with_clock(
merge_authored_eof_with_kind(stage, EofKind::Natural),
VectorClock::new(),
);
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(_) => {}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
let outcome = subscription
.take_last_eof_outcome()
.expect("EOF outcome recorded");
assert_eq!(outcome.worst_kind, Some(EofKind::Poison));
}
#[tokio::test]
async fn canonical_merge_forwarded_control_takes_ordinal_without_exhausting() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
let foreign_stage = StageId::new();
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(foreign_stage), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut deliveries = 0;
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(_) => deliveries += 1,
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert_eq!(deliveries, 6);
assert_eq!(delivered(&subscription), [4, 2]);
assert!(subscription.all_readers_eof());
}
#[tokio::test]
async fn canonical_merge_contract_read_accounting_fires_at_delivery_not_at_hold() {
let (subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
let reader_stage = StageId::new();
let contract_journal: Arc<dyn Journal<ChainEvent>> =
Arc::new(TestJournal::new(JournalOwner::stage(reader_stage)));
let mut subscription = subscription.with_contracts(ContractsWiring {
writer_id: WriterId::from(reader_stage),
contract_journal,
config: ContractConfig::default(),
system_journal: None,
reader_stage: Some(reader_stage),
control_plane: Arc::new(NoControlPlane),
include_delivery_contract: false,
cycle_guard_config: None,
});
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
let mut progress = [ReaderProgress::new(stage_a), ReaderProgress::new(stage_b)];
match subscription
.poll_next_with_state("test_fsm", Some(&mut progress[..]))
.await
{
PollResult::NoEvents => {}
other => panic!("expected NoEvents while A is quiet, got {other:?}"),
}
assert_eq!(progress[1].reader_seq, SeqNo(0));
assert!(
progress[1].last_read_instant.is_some(),
"read instant stamps at head acquisition so the held edge stays healthy"
);
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
let first = match subscription
.poll_next_with_state("test_fsm", Some(&mut progress[..]))
.await
{
PollResult::Event(envelope) => envelope,
other => panic!("expected delivery, got {other:?}"),
};
assert!(matches!(
first.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
));
match subscription
.poll_next_with_state("test_fsm", Some(&mut progress[..]))
.await
{
PollResult::Event(_) => {}
other => panic!("expected delivery, got {other:?}"),
}
assert_eq!(progress[1].reader_seq, SeqNo(1));
}
#[tokio::test]
async fn canonical_merge_filtered_events_take_no_ordinals() {
let (subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
let mut selected = HashMap::new();
selected.insert(
stage_a,
[EventType::from("sel.event")].into_iter().collect(),
);
let mut subscription = subscription
.with_selected_event_types(selected)
.transport_only();
journal_a.append_with_clock(merge_data(stage_a, "sel.event"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "other.event"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "sel.event"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b.event"), VectorClock::new());
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut deliveries = 0;
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => {
assert_ne!(
envelope.event_type(),
"other.event",
"unselected events must never deliver"
);
deliveries += 1;
}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert_eq!(deliveries, 5);
assert_eq!(delivered(&subscription), [3, 2]);
}
#[tokio::test]
async fn long_filtered_run_yields_bounded_cursor_progress_before_delivery() {
const FILTERED_ROWS: usize = 64;
let upstream = StageId::new();
let journal = Arc::new(SharedTestJournal::new(JournalOwner::stage(upstream)));
for _ in 0..FILTERED_ROWS {
journal.append_with_clock(merge_data(upstream, "test.unselected"), VectorClock::new());
}
journal.append_with_clock(merge_data(upstream, "test.selected"), VectorClock::new());
let upstreams = [(
upstream,
"upstream".to_string(),
journal as Arc<dyn Journal<ChainEvent>>,
)];
let mut selected = HashMap::new();
selected.insert(
upstream,
vec![SelectedFeedMetadata::new(
EventType::from("test.selected"),
SelectedFeedRole::Input,
)],
);
let mut subscription = UpstreamSubscription::new_with_names("bounded_filter", &upstreams)
.await
.expect("subscription")
.with_selected_feeds(selected)
.transport_only();
for completed in 0..FILTERED_ROWS {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::CursorAdvanced {
upstream: completed_upstream,
completed_data_rows,
} => {
assert_eq!(completed_upstream, upstream);
assert_eq!(completed_data_rows, 1);
}
other => panic!(
"filtered poll {completed} should yield cursor progress before scanning on: {other:?}"
),
}
assert_eq!(subscription.delivered_data_count(), 0);
assert_eq!(subscription.last_delivered_stage_input_position(), None);
}
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => {
assert_eq!(envelope.event_type(), "test.selected");
}
other => panic!("selected row should eventually deliver: {other:?}"),
}
assert_eq!(subscription.delivered_data_count(), 1);
}
#[tokio::test]
async fn round_robin_stage_input_positions_skip_control_events() {
let stage_a = StageId::new();
let journal_a = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_a)));
let foreign_stage = StageId::new();
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(foreign_stage), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
let upstreams = [(
stage_a,
"upstream_a".to_string(),
journal_a as Arc<dyn Journal<ChainEvent>>,
)];
let mut subscription = UpstreamSubscription::new_with_names("rr_owner", &upstreams)
.await
.unwrap();
let expected_positions = [Some(1u64), None, Some(2u64), None];
for expected in expected_positions {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(_) => {}
other => panic!("expected delivery, got {other:?}"),
}
assert_eq!(
subscription
.last_delivered_stage_input_position()
.map(|position| position.0),
expected
);
}
}
#[tokio::test]
async fn canonical_merge_skips_reader_telemetry_without_taking_ordinals() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair_transport_only("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_consumption_progress(stage_a, 1), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b2"), VectorClock::new());
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let mut order = Vec::new();
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => order.push(envelope.event_type()),
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert_eq!(order.len(), 6);
assert_eq!(&order[..4], ["a1", "b1", "a2", "b2"]);
assert_eq!(delivered(&subscription), [3, 3]);
assert!(subscription.all_readers_eof());
}
#[tokio::test]
async fn round_robin_transport_only_skips_reader_telemetry() {
let stage_a = StageId::new();
let journal_a = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_a)));
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_consumption_progress(stage_a, 1), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
let upstreams = [(
stage_a,
"upstream_a".to_string(),
journal_a as Arc<dyn Journal<ChainEvent>>,
)];
let mut subscription = UpstreamSubscription::new_with_names("rr_owner", &upstreams)
.await
.unwrap()
.transport_only();
let mut order = Vec::new();
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => order.push(envelope.event_type()),
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert_eq!(order.len(), 3);
assert_eq!(&order[..2], ["a1", "a2"]);
assert_eq!(delivered(&subscription), [3]);
assert!(subscription.all_readers_eof());
}
#[tokio::test]
async fn interleaved_telemetry_does_not_perturb_merged_order() {
async fn drain_order(
subscription: &mut UpstreamSubscription<ChainEvent>,
) -> (Vec<String>, Vec<u64>) {
let mut order = Vec::new();
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => order.push(envelope.event_type()),
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
let ordinals = delivered(subscription);
(order, ordinals)
}
let (mut clean, (clean_a, clean_journal_a), (clean_b, clean_journal_b)) =
canonical_pair_transport_only("upstream_a", "upstream_b").await;
let (mut noisy, (noisy_a, noisy_journal_a), (noisy_b, noisy_journal_b)) =
canonical_pair_transport_only("upstream_a", "upstream_b").await;
for (stage_a, journal_a, telemetry) in [
(clean_a, &clean_journal_a, false),
(noisy_a, &noisy_journal_a, true),
] {
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
if telemetry {
journal_a.append_with_clock(merge_consumption_progress(stage_a, 1), VectorClock::new());
}
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a3"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
}
for (stage_b, journal_b) in [(clean_b, &clean_journal_b), (noisy_b, &noisy_journal_b)] {
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b2"), VectorClock::new());
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
}
let (clean_order, clean_ordinals) = drain_order(&mut clean).await;
let (noisy_order, noisy_ordinals) = drain_order(&mut noisy).await;
assert_eq!(
clean_order, noisy_order,
"an interleaved telemetry row must not perturb the merged order"
);
assert_eq!(
clean_ordinals, noisy_ordinals,
"per-reader delivered ordinals must be a pure function of stream content"
);
}
#[tokio::test]
#[should_panic(expected = "CanonicalMerge requires reader")]
async fn canonical_merge_rejects_tail_start_readers() {
let stage_a = StageId::new();
let journal_a = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_a)));
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
let upstreams = [(
stage_a,
"upstream_a".to_string(),
journal_a as Arc<dyn Journal<ChainEvent>>,
)];
let _ = UpstreamSubscription::new_at_tail("tail_owner", &upstreams, &[1])
.await
.unwrap()
.with_reader_selection(ReaderSelectionPolicy::CanonicalMerge);
}
#[tokio::test]
async fn feed_role_distinguishes_tiebreak_identity() {
let stage_a = StageId::new();
let stage_b = StageId::new();
let journal_a = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_a)));
let journal_b = Arc::new(SharedTestJournal::new(JournalOwner::stage(stage_b)));
let upstreams = [
(
stage_a,
"shared_upstream".to_string(),
journal_a as Arc<dyn Journal<ChainEvent>>,
),
(
stage_b,
"shared_upstream".to_string(),
journal_b as Arc<dyn Journal<ChainEvent>>,
),
];
let mut selected_feeds = HashMap::new();
selected_feeds.insert(
stage_a,
vec![SelectedFeedMetadata::new(
EventType::from("shared.event"),
SelectedFeedRole::Reference,
)],
);
selected_feeds.insert(
stage_b,
vec![SelectedFeedMetadata::new(
EventType::from("shared.event"),
SelectedFeedRole::Stream,
)],
);
let subscription = UpstreamSubscription::new_with_names("role_owner", &upstreams)
.await
.unwrap()
.with_selected_feeds(selected_feeds);
let keys = &subscription.reader_tiebreak_keys;
assert_eq!(
keys[0].stage_key, keys[1].stage_key,
"stage keys deliberately collide"
);
assert_ne!(
keys[0].feed_identity, keys[1].feed_identity,
"the role qualifier must keep the tiebreak key total"
);
assert_eq!(
keys[0].feed_identity,
FeedIdentity::from_feeds(&[SelectedFeedMetadata::new(
EventType::from("shared.event"),
SelectedFeedRole::Reference,
)])
);
assert_eq!(
keys[1].feed_identity,
FeedIdentity::from_feeds(&[SelectedFeedMetadata::new(
EventType::from("shared.event"),
SelectedFeedRole::Stream,
)])
);
}
#[tokio::test]
async fn delivered_upstream_identity_ignores_event_writer() {
let upstream_stage = StageId::new();
let foreign_author = StageId::new();
let upstream_owner = JournalOwner::stage(upstream_stage);
let upstream_journal: Arc<dyn Journal<ChainEvent>> = Arc::new(TestJournal::new(upstream_owner));
upstream_journal
.append(
ChainEventFactory::data_event(
WriterId::Stage(foreign_author),
"test.forwarded",
json!({"n": 1}),
),
Default::default(),
)
.await
.unwrap();
let upstreams = [(upstream_stage, "upstream".to_string(), upstream_journal)];
let mut subscription = UpstreamSubscription::new_with_names("test_owner", &upstreams)
.await
.unwrap()
.transport_only();
let mut reader_progress = [ReaderProgress::new(upstream_stage)];
let polled = subscription
.poll_next_with_state("test_fsm", Some(&mut reader_progress[..]))
.await;
let envelope = match polled {
PollResult::Event(envelope) => envelope,
other => panic!("expected delivered data event, got {other:?}"),
};
assert_eq!(
envelope.envelope.provenance.event.writer_id,
WriterId::Stage(foreign_author),
"premise: the forwarded row preserves its original author"
);
assert_eq!(
subscription.last_delivered_upstream_stage(),
Some(upstream_stage),
"edge identity must be the reader slot's stage, not the event author"
);
}
#[tokio::test]
async fn catch_up_watermark_advances_reader_generation_at_delivery_only() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_catch_up(stage_a, 1, "upstream_a"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b2"), VectorClock::new());
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "a1");
assert_eq!(
subscription.last_delivered_generation(),
Some(ReaderGeneration(0))
);
let second = expect_delivery(&mut subscription).await;
assert_eq!(second.event_type(), "b1");
assert!(
!subscription.all_readers_caught_up(ReaderGeneration(1)),
"a merely held watermark must not advance its reader"
);
assert_eq!(
subscription.last_delivered_stage_input_position(),
Some(StageInputPosition(2))
);
let third = expect_delivery(&mut subscription).await;
assert!(matches!(
third.payload,
ChainPayload::FlowControl(FlowControlPayload::CatchUpComplete { .. })
));
assert_eq!(
subscription.last_delivered_generation(),
Some(ReaderGeneration(0)),
"the watermark delivers at the generation it closes"
);
assert_eq!(
subscription.last_delivered_stage_input_position(),
None,
"the watermark is control: it reports no data-input position"
);
let fourth = expect_delivery(&mut subscription).await;
assert_eq!(fourth.event_type(), "b2");
assert_eq!(
subscription.last_delivered_stage_input_position(),
Some(StageInputPosition(3))
);
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while B is quiet, got {other:?}"),
}
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let fifth = expect_delivery(&mut subscription).await;
assert!(matches!(
fifth.payload,
ChainPayload::FlowControl(FlowControlPayload::Eof { .. })
));
let sixth = expect_delivery(&mut subscription).await;
assert_eq!(sixth.event_type(), "a2");
assert_eq!(
subscription.last_delivered_generation(),
Some(ReaderGeneration(1))
);
assert_eq!(
subscription.last_delivered_stage_input_position(),
Some(StageInputPosition(4))
);
assert_eq!(delivered(&subscription), [3, 3]);
}
#[tokio::test]
async fn recorded_heads_deliver_before_any_live_head() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_catch_up(stage_a, 1, "upstream_a"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "a2"), VectorClock::new());
for label in ["b1", "b2", "b3"] {
journal_b.append_with_clock(merge_data(stage_b, label), VectorClock::new());
}
let mut order = Vec::new();
for _ in 0..5 {
let envelope = expect_delivery(&mut subscription).await;
order.push(envelope.event_type());
}
assert_eq!(order, ["a1", "b1", "control.catch_up_complete", "b2", "b3"]);
assert_eq!(delivered(&subscription), [2, 3]);
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while B is quiet, got {other:?}"),
}
let wait = subscription.merge_wait().expect("merge wait recorded");
assert_eq!(wait.quiet_inputs, vec![(stage_b, "upstream_b".to_string())]);
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
let sixth = expect_delivery(&mut subscription).await;
assert_eq!(sixth.event_type(), "control.eof");
let seventh = expect_delivery(&mut subscription).await;
assert_eq!(seventh.event_type(), "a2");
assert_eq!(
subscription.last_delivered_generation(),
Some(ReaderGeneration(1))
);
}
#[tokio::test]
async fn watermark_generation_skip_fails_closed() {
let (mut subscription, stage_a, journal_a) = canonical_single("upstream_a").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
journal_a.append_with_clock(merge_catch_up(stage_a, 2, "upstream_a"), VectorClock::new());
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "a1");
let message = expect_poll_error(&mut subscription).await;
assert!(
message.contains("FLOWIP-120n F11"),
"error must name F11: {message}"
);
assert_eq!(
subscription.last_delivered_generation(),
Some(ReaderGeneration(0)),
"the rejected delivery stamps nothing"
);
assert!(
!subscription.all_readers_caught_up(ReaderGeneration(1)),
"the reader's generation must not advance on a rejected boundary"
);
}
#[tokio::test]
async fn forwarded_watermark_stage_key_mismatch_fails_closed() {
let (mut subscription, stage_a, journal_a) = canonical_single("upstream_a").await;
journal_a.append_with_clock(merge_catch_up(stage_a, 1, "upstream_b"), VectorClock::new());
let message = expect_poll_error(&mut subscription).await;
assert!(
message.contains("FLOWIP-120n F8"),
"error must name F8: {message}"
);
}
#[tokio::test]
async fn all_readers_caught_up_counts_eof_as_crossed() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_catch_up(stage_a, 1, "upstream_a"), VectorClock::new());
journal_a.append_with_clock(merge_authored_eof(stage_a), VectorClock::new());
journal_b.append_with_clock(merge_data(stage_b, "b1"), VectorClock::new());
journal_b.append_with_clock(merge_authored_eof(stage_b), VectorClock::new());
assert!(
!subscription.all_readers_caught_up(ReaderGeneration(1)),
"no reader has crossed before any delivery"
);
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "control.catch_up_complete");
assert!(
!subscription.all_readers_caught_up(ReaderGeneration(1)),
"B has neither crossed nor delivered its EOF"
);
let second = expect_delivery(&mut subscription).await;
assert_eq!(second.event_type(), "b1");
assert!(!subscription.all_readers_caught_up(ReaderGeneration(1)));
let third = expect_delivery(&mut subscription).await;
assert_eq!(third.event_type(), "control.eof");
assert!(subscription.all_readers_caught_up(ReaderGeneration(1)));
let (mut lagging, (lag_a, lag_journal_a), (lag_b, lag_journal_b)) =
canonical_pair("upstream_a", "upstream_b").await;
lag_journal_a.append_with_clock(merge_data(lag_a, "a1"), VectorClock::new());
lag_journal_a.append_with_clock(merge_data(lag_a, "a2"), VectorClock::new());
lag_journal_b.append_with_clock(merge_authored_eof(lag_b), VectorClock::new());
let first = expect_delivery(&mut lagging).await;
assert_eq!(first.event_type(), "a1");
let second = expect_delivery(&mut lagging).await;
assert_eq!(second.event_type(), "control.eof");
assert!(
!lagging.all_readers_caught_up(ReaderGeneration(1)),
"only B is EOF-exhausted; A has not crossed"
);
}
#[tokio::test]
async fn resume_of_resume_watermarks_stack() {
let (mut subscription, stage_a, journal_a) = canonical_single("upstream_a").await;
journal_a.append_with_clock(merge_data(stage_a, "d1"), VectorClock::new());
journal_a.append_with_clock(merge_catch_up(stage_a, 1, "upstream_a"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "d2"), VectorClock::new());
journal_a.append_with_clock(merge_catch_up(stage_a, 2, "upstream_a"), VectorClock::new());
journal_a.append_with_clock(merge_data(stage_a, "d3"), VectorClock::new());
let expected = [
("d1", ReaderGeneration(0)),
("control.catch_up_complete", ReaderGeneration(0)),
("d2", ReaderGeneration(1)),
("control.catch_up_complete", ReaderGeneration(1)),
("d3", ReaderGeneration(2)),
];
for (event_type, generation) in expected {
let envelope = expect_delivery(&mut subscription).await;
assert_eq!(envelope.event_type(), event_type);
assert_eq!(subscription.last_delivered_generation(), Some(generation));
}
assert!(subscription.all_readers_caught_up(ReaderGeneration(2)));
assert!(!subscription.all_readers_caught_up(ReaderGeneration(3)));
}
fn with_seq(mut event: ChainEvent, seq: u64) -> ChainEvent {
event.admission_seq = Some(AdmissionSeq(seq));
event
}
async fn seq_pair_with_entered(
name_a: &str,
name_b: &str,
entered_generation: u64,
) -> (
UpstreamSubscription<ChainEvent>,
(StageId, Arc<SharedTestJournal>),
(StageId, Arc<SharedTestJournal>),
) {
let (subscription, a, b) = canonical_pair(name_a, name_b).await;
(
subscription
.with_seq_ordered(true)
.with_entered_generation(ReaderGeneration(entered_generation)),
a,
b,
)
}
async fn seq_pair(
name_a: &str,
name_b: &str,
) -> (
UpstreamSubscription<ChainEvent>,
(StageId, Arc<SharedTestJournal>),
(StageId, Arc<SharedTestJournal>),
) {
seq_pair_with_entered(name_a, name_b, 0).await
}
#[tokio::test]
async fn seq_merge_delivers_available_head_while_sibling_is_quiet() {
let (mut subscription, (stage_a, journal_a), (_stage_b, journal_b)) =
seq_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a1"), 1), VectorClock::new());
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a2"), 2), VectorClock::new());
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "a1");
let second = expect_delivery(&mut subscription).await;
assert_eq!(second.event_type(), "a2");
assert!(subscription.merge_wait().is_none());
assert_eq!(delivered(&subscription), [2, 0]);
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while both are quiet, got {other:?}"),
}
assert!(subscription.merge_wait().is_none());
journal_b.append_with_clock(with_seq(merge_data(_stage_b, "b1"), 3), VectorClock::new());
let third = expect_delivery(&mut subscription).await;
assert_eq!(third.event_type(), "b1");
}
#[tokio::test]
async fn seq_merge_orders_by_admission_seq_not_arrival() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
seq_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a1"), 2), VectorClock::new());
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a2"), 4), VectorClock::new());
journal_a.append_with_clock(with_seq(merge_authored_eof(stage_a), 5), VectorClock::new());
journal_b.append_with_clock(with_seq(merge_data(stage_b, "b1"), 1), VectorClock::new());
journal_b.append_with_clock(with_seq(merge_data(stage_b, "b2"), 3), VectorClock::new());
journal_b.append_with_clock(with_seq(merge_authored_eof(stage_b), 6), VectorClock::new());
let mut order = Vec::new();
let mut delivering_stage = Vec::new();
loop {
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::Event(envelope) => {
order.push(envelope.event_type());
delivering_stage.push(subscription.last_delivered_upstream_stage().unwrap());
}
PollResult::CursorAdvanced { .. } => {}
PollResult::NoEvents => break,
PollResult::Error(e) => panic!("unexpected poll error: {e:?}"),
}
}
assert_eq!(
order,
["b1", "a1", "b2", "control.eof", "a2", "control.eof"],
"delivery follows admission sequence with position-inherited control"
);
assert_eq!(
delivering_stage,
[stage_b, stage_a, stage_b, stage_b, stage_a, stage_a],
);
assert!(subscription.all_readers_eof());
}
#[tokio::test]
async fn seq_merge_orders_re_authored_control_by_journal_position_not_stamp() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
seq_pair("upstream_a", "upstream_b").await;
let contract = ChainEventFactory::source_contract_event(
WriterId::Stage(stage_a),
obzenflow_core::event::SourceContractEventParams {
expected_count: None,
source_id: stage_a,
route: None,
journal_path: obzenflow_core::event::types::JournalPath("upstream_a".to_string()),
journal_index: obzenflow_core::event::types::JournalIndex(0),
writer_seq: None,
vector_clock: None,
},
);
journal_a.append_with_clock(with_seq(contract, 999), VectorClock::new());
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a1"), 1), VectorClock::new());
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a2"), 3), VectorClock::new());
journal_b.append_with_clock(with_seq(merge_data(stage_b, "b1"), 2), VectorClock::new());
journal_b.append_with_clock(with_seq(merge_data(stage_b, "b2"), 4), VectorClock::new());
let expected = ["control.source_contract", "a1", "b1", "a2", "b2"];
for event_type in expected {
let envelope = expect_delivery(&mut subscription).await;
assert_eq!(envelope.event_type(), event_type);
}
}
#[tokio::test]
async fn seq_merge_fails_closed_on_sequence_less_head() {
let (mut subscription, (stage_a, journal_a), (_stage_b, _journal_b)) =
seq_pair("upstream_a", "upstream_b").await;
journal_a.append_with_clock(merge_data(stage_a, "a1"), VectorClock::new());
let message = expect_poll_error(&mut subscription).await;
assert!(
message.contains("FLOWIP-120n F18"),
"error must name F18: {message}"
);
assert!(
message.contains("upstream_a"),
"error must name the reader's stage key: {message}"
);
}
#[tokio::test]
async fn seq_reader_below_entered_generation_keeps_kahn_wait_until_crossing() {
let (mut subscription, (stage_a, journal_a), (stage_b, journal_b)) =
seq_pair_with_entered("upstream_a", "upstream_b", 1).await;
journal_a.append_with_clock(with_seq(merge_data(stage_a, "a1"), 1), VectorClock::new());
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while pre-crossing B is quiet, got {other:?}"),
}
let wait = subscription.merge_wait().expect("blocking wait recorded");
assert_eq!(wait.quiet_inputs, vec![(stage_b, "upstream_b".to_string())]);
journal_b.append_with_clock(
with_seq(merge_catch_up(stage_b, 1, "upstream_b"), 2),
VectorClock::new(),
);
let first = expect_delivery(&mut subscription).await;
assert_eq!(first.event_type(), "a1");
match subscription.poll_next_with_state("test_fsm", None).await {
PollResult::NoEvents => {}
other => panic!("expected NoEvents while pre-crossing A is quiet, got {other:?}"),
}
let wait = subscription.merge_wait().expect("blocking wait recorded");
assert_eq!(wait.quiet_inputs, vec![(stage_a, "upstream_a".to_string())]);
journal_a.append_with_clock(
with_seq(merge_catch_up(stage_a, 1, "upstream_a"), 3),
VectorClock::new(),
);
let second = expect_delivery(&mut subscription).await;
assert_eq!(second.event_type(), "control.catch_up_complete");
let third = expect_delivery(&mut subscription).await;
assert_eq!(third.event_type(), "control.catch_up_complete");
assert!(subscription.all_readers_caught_up(ReaderGeneration(1)));
journal_b.append_with_clock(with_seq(merge_data(stage_b, "b1"), 4), VectorClock::new());
let fourth = expect_delivery(&mut subscription).await;
assert_eq!(fourth.event_type(), "b1");
assert!(subscription.merge_wait().is_none());
}