use super::{ContractStatus, ContractTracker, ReaderProgress, UpstreamSubscription};
use crate::messaging::upstream_subscription_policy::{EdgeContext, EdgeContractDecision};
use obzenflow_core::event::payloads::system_payload::{
ContractName, ContractResultStatusLabel, SystemFeedRole, SystemPayload,
};
use obzenflow_core::event::system_event::SystemEvent;
use obzenflow_core::event::types::{
Count, DurationMs, EventType, JournalIndex, JournalPath, SeqNo,
ViolationCause as EventViolationCause,
};
use obzenflow_core::event::{
ChainEventFactory, ConsumptionFinalEventParams, ConsumptionProgressEventParams, JournalEvent,
};
use obzenflow_core::journal::Journal;
use obzenflow_core::{ContractResult, ViolationCause};
use std::sync::Arc;
use tokio::time::Instant;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ContractCheckMode {
Authoritative,
DiagnosticsOnly,
}
fn contract_result_labels_for_emission(
result: &ContractResult,
pending_label: ContractResultStatusLabel,
) -> (ContractResultStatusLabel, Option<String>) {
match result {
ContractResult::Passed(_) => (ContractResultStatusLabel::Passed, None),
ContractResult::Failed(v) => (
ContractResultStatusLabel::Failed,
Some(v.cause.cause_label().to_string()),
),
ContractResult::Pending => (pending_label, None),
}
}
#[derive(Clone, Debug)]
struct DirectFeedContractEvidence {
event_type: EventType,
feed_role: Option<SystemFeedRole>,
reader_seq: SeqNo,
advertised_writer_seq: SeqNo,
pass: bool,
reason: Option<EventViolationCause>,
contract_results: Vec<(ContractName, ContractResult)>,
}
#[derive(Clone, Debug)]
struct DirectFeedProgressEvidence {
feed_index: usize,
event_type: EventType,
feed_role: Option<SystemFeedRole>,
reader_seq: SeqNo,
advertised_writer_seq: Option<SeqNo>,
contract_results: Vec<(ContractName, ContractResult)>,
should_record_heartbeat: bool,
}
impl<T> UpstreamSubscription<T>
where
T: JournalEvent + 'static,
{
fn direct_feed_contract_evidence_for_reader(
&self,
progress: &ReaderProgress,
index: usize,
reader_stage: obzenflow_core::StageId,
) -> Vec<DirectFeedContractEvidence> {
let Some(feed_chains) = self.contract_feed_chains.get(index) else {
return Vec::new();
};
let Some(advertised_by_type) = self.advertised_writer_seq_by_reader_event_type.get(index)
else {
return Vec::new();
};
if advertised_by_type.is_empty() || feed_chains.is_empty() {
return Vec::new();
}
feed_chains
.iter()
.map(|feed_chain| {
let reader_seq = self.selected_reader_seq_for_feed(index, &feed_chain.metadata);
let advertised_writer_seq = self
.advertised_writer_seq_for_feed(index, &feed_chain.metadata)
.unwrap_or(SeqNo(0));
let contract_results = feed_chain.chain.verify_all(progress.stage_id, reader_stage);
let raw_failure = contract_results.iter().find_map(|(contract_name, result)| {
let ContractResult::Failed(violation) = result else {
return None;
};
Some((contract_name.clone(), violation.cause.clone()))
});
let raw_reason = raw_failure.as_ref().map(|(_, cause)| match cause {
ViolationCause::SeqDivergence { advertised, reader } => {
EventViolationCause::SeqDivergence {
advertised: *advertised,
reader: *reader,
}
}
ViolationCause::ContentMismatch { .. } => {
EventViolationCause::Other("content_mismatch".into())
}
ViolationCause::DeliveryMismatch { .. } => {
EventViolationCause::Other("delivery_mismatch".into())
}
ViolationCause::AccountingMismatch { .. } => {
EventViolationCause::Other("accounting_mismatch".into())
}
ViolationCause::Divergence {
predicate,
observed,
threshold,
window_seconds,
} => EventViolationCause::Divergence {
predicate: predicate.clone(),
observed: *observed,
threshold: *threshold,
window_seconds: *window_seconds,
},
ViolationCause::Other(message) => EventViolationCause::Other(message.clone()),
});
let results_only: Vec<ContractResult> = contract_results
.iter()
.map(|(_, result)| result.clone())
.collect();
let edge = EdgeContext {
upstream_stage: progress.stage_id,
downstream_stage: reader_stage,
advertised_writer_seq: Some(advertised_writer_seq),
reader_seq,
};
let decision = self
.contract_policies
.get(index)
.and_then(|policy| policy.as_ref())
.map(|policy| policy.decide(&results_only, &edge));
let (pass, reason) = match decision {
Some(EdgeContractDecision::Pass) => (true, None),
Some(EdgeContractDecision::Fail(cause)) => (false, Some(cause)),
None if raw_failure.is_some() => (false, raw_reason.clone()),
None => (true, None),
};
DirectFeedContractEvidence {
event_type: feed_chain.metadata.event_type().clone(),
feed_role: feed_chain.metadata.system_feed_role(),
reader_seq,
advertised_writer_seq,
pass,
reason,
contract_results,
}
})
.collect()
}
async fn emit_direct_feed_contract_system_events(
&self,
tracker: &ContractTracker,
upstream: obzenflow_core::StageId,
reader: obzenflow_core::StageId,
evidence: &[DirectFeedContractEvidence],
) -> bool {
let Some(system_journal) = &tracker.system_journal else {
return true;
};
let mut append_ok = true;
for feed in evidence {
for (contract_name, result) in &feed.contract_results {
let (status_label, cause_label) =
contract_result_labels_for_emission(result, ContractResultStatusLabel::Pending);
let result_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractResult {
upstream,
reader,
selected_event_type: Some(feed.event_type.clone()),
feed_role: feed.feed_role,
contract_name: contract_name.clone(),
status: status_label,
cause: cause_label,
reader_seq: Some(feed.reader_seq),
advertised_writer_seq: Some(feed.advertised_writer_seq),
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
result_event,
Default::default(),
)
.await
{
append_ok = false;
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?upstream,
reader = ?reader,
selected_event_type = %feed.event_type,
feed_role = ?feed.feed_role,
contract = %contract_name,
error = %e,
"Failed to append direct feed contract result event; skipping emission"
);
}
}
let status_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractStatus {
upstream,
reader,
selected_event_type: Some(feed.event_type.clone()),
feed_role: feed.feed_role,
pass: feed.pass,
reader_seq: Some(feed.reader_seq),
advertised_writer_seq: Some(feed.advertised_writer_seq),
reason: feed.reason.clone(),
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
status_event,
Default::default(),
)
.await
{
append_ok = false;
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?upstream,
reader = ?reader,
selected_event_type = %feed.event_type,
feed_role = ?feed.feed_role,
error = %e,
"Failed to append direct feed contract status event; skipping emission"
);
}
}
append_ok
}
fn direct_feed_progress_evidence_for_reader(
&self,
progress: &ReaderProgress,
index: usize,
reader_stage: obzenflow_core::StageId,
) -> Vec<DirectFeedProgressEvidence> {
let Some(feed_chains) = self.contract_feed_chains.get(index) else {
return Vec::new();
};
if feed_chains.is_empty() {
return Vec::new();
}
feed_chains
.iter()
.enumerate()
.filter_map(|(feed_index, feed_chain)| {
let reader_seq = self.selected_reader_seq_for_feed(index, &feed_chain.metadata);
if reader_seq.0 == 0 {
return None;
}
let contract_results = feed_chain
.chain
.check_progress_all(progress.stage_id, reader_stage);
let has_failure = contract_results
.iter()
.any(|(_, result)| matches!(result, ContractResult::Failed(_)));
let should_emit_healthy = reader_seq != feed_chain.last_contract_result_seq;
if !has_failure && !should_emit_healthy {
return None;
}
Some(DirectFeedProgressEvidence {
feed_index,
event_type: feed_chain.metadata.event_type().clone(),
feed_role: feed_chain.metadata.system_feed_role(),
reader_seq,
advertised_writer_seq: self
.advertised_writer_seq_for_feed(index, &feed_chain.metadata),
contract_results,
should_record_heartbeat: should_emit_healthy,
})
})
.collect()
}
async fn emit_direct_feed_progress_contract_results(
&mut self,
writer_id: obzenflow_core::WriterId,
system_journal: Option<Arc<dyn Journal<SystemEvent>>>,
progress: &ReaderProgress,
index: usize,
reader_stage: obzenflow_core::StageId,
) {
let Some(system_journal) = system_journal else {
return;
};
let evidence = self.direct_feed_progress_evidence_for_reader(progress, index, reader_stage);
if evidence.is_empty() {
return;
}
let mut completed_heartbeats = Vec::new();
for feed in &evidence {
let mut emitted_any_for_feed = false;
for (contract_name, result) in &feed.contract_results {
let (status_label, cause_label) =
contract_result_labels_for_emission(result, ContractResultStatusLabel::Healthy);
let result_event = SystemEvent::new(
writer_id,
SystemPayload::ContractResult {
upstream: progress.stage_id,
reader: reader_stage,
selected_event_type: Some(feed.event_type.clone()),
feed_role: feed.feed_role,
contract_name: contract_name.clone(),
status: status_label,
cause: cause_label,
reader_seq: Some(feed.reader_seq),
advertised_writer_seq: feed.advertised_writer_seq,
},
);
if let Err(e) = crate::supervised_base::publication::append(
&system_journal,
result_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader = ?reader_stage,
reader_index = index,
selected_event_type = %feed.event_type,
feed_role = ?feed.feed_role,
contract = %contract_name,
error = %e,
"Failed to append direct feed progress contract result event; skipping emission"
);
} else {
emitted_any_for_feed = true;
}
}
if emitted_any_for_feed && feed.should_record_heartbeat {
completed_heartbeats.push((feed.feed_index, feed.reader_seq));
}
}
if let Some(feed_chains) = self.contract_feed_chains.get_mut(index) {
for (feed_index, reader_seq) in completed_heartbeats {
if let Some(feed_chain) = feed_chains.get_mut(feed_index) {
feed_chain.last_contract_result_seq = reader_seq;
}
}
}
}
pub async fn check_contracts(
&mut self,
reader_progress: &mut [ReaderProgress],
) -> ContractStatus {
self.check_contracts_with_mode(reader_progress, ContractCheckMode::Authoritative)
.await
}
pub async fn check_contracts_diagnostics_only(
&mut self,
reader_progress: &mut [ReaderProgress],
) -> ContractStatus {
self.check_contracts_with_mode(reader_progress, ContractCheckMode::DiagnosticsOnly)
.await
}
async fn check_contracts_with_mode(
&mut self,
reader_progress: &mut [ReaderProgress],
mode: ContractCheckMode,
) -> ContractStatus {
if self.contract_tracker.is_none() {
return ContractStatus::Healthy;
};
let now = Instant::now();
let mut status = ContractStatus::Healthy;
for (index, progress) in reader_progress.iter_mut().enumerate() {
if progress.final_emitted {
continue;
}
if mode == ContractCheckMode::DiagnosticsOnly && self.state.is_reader_eof(index) {
continue;
}
self.check_progress_contracts_for_reader(progress, index, &mut status)
.await;
let should_emit_progress = self.should_emit_progress(progress, index, now);
if should_emit_progress {
self.emit_progress_for_reader(progress, index, now, &mut status)
.await;
if mode == ContractCheckMode::Authoritative
&& self.state.is_reader_eof(index)
&& !progress.final_emitted
{
self.verify_eof_contracts_for_reader(progress, index, &mut status)
.await;
}
} else {
self.check_stall_for_reader(progress, index, now, &mut status)
.await;
}
}
status
}
async fn check_progress_contracts_for_reader(
&mut self,
progress: &mut ReaderProgress,
index: usize,
status: &mut ContractStatus,
) {
let Some((reader_stage, writer_id, system_journal)) =
self.contract_tracker.as_ref().and_then(|tracker| {
tracker.reader_stage.map(|reader_stage| {
(
reader_stage,
tracker.writer_id,
tracker.system_journal.clone(),
)
})
})
else {
return;
};
if self.state.is_reader_eof(index) {
return;
}
if self
.contract_feed_chains
.get(index)
.is_some_and(|chains| !chains.is_empty())
{
self.emit_direct_feed_progress_contract_results(
writer_id,
system_journal,
progress,
index,
reader_stage,
)
.await;
return;
}
let Some(tracker) = &self.contract_tracker else {
return;
};
let (selected_event_type, feed_role) =
self.unique_selected_feed_for_stage(progress.stage_id);
let Some(chain_slot) = self.contract_chains.get(index).and_then(|c| c.as_ref()) else {
return;
};
let progress_seq = self.progress_seq(progress);
if progress_seq.0 == 0 {
return;
}
let should_emit_healthy = progress_seq != progress.last_contract_result_seq;
let results = chain_slot.check_progress_all(progress.stage_id, reader_stage);
let results_only: Vec<ContractResult> = results.iter().map(|(_, r)| r.clone()).collect();
if let Some(system_journal) = &tracker.system_journal {
let mut emitted_any = false;
for (contract_name, result) in &results {
if matches!(result, ContractResult::Pending) && !should_emit_healthy {
continue;
}
let (status_label, cause_label) =
contract_result_labels_for_emission(result, ContractResultStatusLabel::Healthy);
let result_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractResult {
upstream: progress.stage_id,
reader: reader_stage,
selected_event_type: selected_event_type.clone(),
feed_role,
contract_name: contract_name.clone(),
status: status_label,
cause: cause_label,
reader_seq: Some(progress_seq),
advertised_writer_seq: progress.advertised_writer_seq,
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
result_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader = ?reader_stage,
reader_index = index,
contract = %contract_name,
error = %e,
"Failed to append progress contract result event; skipping emission"
);
} else {
emitted_any = true;
}
}
if emitted_any && should_emit_healthy {
progress.last_contract_result_seq = progress_seq;
}
}
let edge = EdgeContext {
upstream_stage: progress.stage_id,
downstream_stage: reader_stage,
advertised_writer_seq: progress.advertised_writer_seq,
reader_seq: progress.reader_seq,
};
let Some(policy_stack) = self.contract_policies.get(index).and_then(|p| p.as_ref()) else {
return;
};
let decision = policy_stack.decide(&results_only, &edge);
match decision {
EdgeContractDecision::Pass => {}
EdgeContractDecision::Fail(cause) => {
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Violated {
upstream: progress.stage_id,
cause: cause.clone(),
};
}
if let Some(system_journal) = &tracker.system_journal {
let status_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractStatus {
upstream: progress.stage_id,
reader: reader_stage,
selected_event_type: selected_event_type.clone(),
feed_role,
pass: false,
reader_seq: Some(progress.reader_seq),
advertised_writer_seq: progress.advertised_writer_seq,
reason: Some(cause.clone()),
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
status_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader = ?reader_stage,
reader_index = index,
error = %e,
"Failed to append progress contract status; skipping emission"
);
}
}
progress.contract_violated = true;
}
}
}
async fn emit_progress_for_reader(
&mut self,
progress: &mut ReaderProgress,
index: usize,
now: Instant,
status: &mut ContractStatus,
) {
let Some(tracker) = &self.contract_tracker else {
return;
};
let progress_seq = self.progress_seq(progress);
let progress_last_event_id = self.progress_last_event_id(progress);
let progress_vector_clock = self.progress_vector_clock(progress);
let stalled_duration = progress
.stalled_since
.map(|s| DurationMs(now.duration_since(s).as_millis() as u64));
let progress_event =
tracker.with_owner_context(ChainEventFactory::consumption_progress_event(
tracker.writer_id,
ConsumptionProgressEventParams {
reader_seq: progress_seq,
last_event_id: progress_last_event_id,
vector_clock: progress_vector_clock.clone(),
eof_seen: self.state.is_reader_eof(index),
reader_path: JournalPath(progress.stage_id.to_string()),
reader_index: JournalIndex(index as u64),
advertised_writer_seq: progress.advertised_writer_seq,
advertised_vector_clock: progress_vector_clock,
stalled_since: stalled_duration,
},
));
match crate::supervised_base::publication::append(
&tracker.journal,
progress_event,
Default::default(),
)
.await
{
Ok(_) => {
progress.last_progress_seq = progress_seq;
progress.last_progress_instant = Some(now);
progress.stalled_since = None;
progress.consecutive_stall_checks = 0;
if matches!(status, ContractStatus::Healthy) {
*status = ContractStatus::ProgressEmitted;
}
}
Err(e) => {
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append progress event; skipping state update"
);
}
}
}
async fn verify_eof_contracts_for_reader(
&mut self,
progress: &mut ReaderProgress,
index: usize,
status: &mut ContractStatus,
) {
let Some(tracker) = &self.contract_tracker else {
return;
};
let (selected_event_type, feed_role) =
self.unique_selected_feed_for_stage(progress.stage_id);
let progress_seq = self.progress_seq(progress);
let progress_last_event_id = self.progress_last_event_id(progress);
let progress_vector_clock = self.progress_vector_clock(progress);
let direct_feed_evidence = tracker
.reader_stage
.map(|reader_stage| {
self.direct_feed_contract_evidence_for_reader(progress, index, reader_stage)
})
.unwrap_or_default();
let mut pass = true;
let mut failure_reason = None;
let mut aggregate_violation_for_journal = false;
let mut aggregate_failure_reason = None;
if let (Some(reader_stage), Some(chain_slot)) = (
tracker.reader_stage,
self.contract_chains.get(index).and_then(|c| c.as_ref()),
) {
let results = chain_slot.verify_all(progress.stage_id, reader_stage);
let results_only: Vec<ContractResult> =
results.iter().map(|(_, r)| r.clone()).collect();
if let Some(system_journal) = &tracker.system_journal {
for (contract_name, result) in &results {
let (status_label, cause_label) = contract_result_labels_for_emission(
result,
ContractResultStatusLabel::Pending,
);
let result_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractResult {
upstream: progress.stage_id,
reader: reader_stage,
selected_event_type: selected_event_type.clone(),
feed_role,
contract_name: contract_name.clone(),
status: status_label,
cause: cause_label,
reader_seq: Some(progress_seq),
advertised_writer_seq: progress.advertised_writer_seq,
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
result_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader = ?reader_stage,
reader_index = index,
contract = %contract_name,
error = %e,
"Failed to append contract result event; skipping emission"
);
}
}
}
let edge = EdgeContext {
upstream_stage: progress.stage_id,
downstream_stage: reader_stage,
advertised_writer_seq: progress.advertised_writer_seq,
reader_seq: progress.reader_seq,
};
if let Some(policy_stack) = self.contract_policies.get(index).and_then(|p| p.as_ref()) {
let decision = policy_stack.decide(&results_only, &edge);
match decision {
EdgeContractDecision::Pass => {
pass = true;
failure_reason = None;
}
EdgeContractDecision::Fail(cause) => {
pass = false;
failure_reason = Some(cause.clone());
aggregate_violation_for_journal = true;
aggregate_failure_reason = Some(cause.clone());
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Violated {
upstream: progress.stage_id,
cause: cause.clone(),
};
}
if let EventViolationCause::SeqDivergence {
advertised: Some(advertised),
reader,
} = cause
{
if advertised.0 > reader.0 {
let gap_event = tracker.with_owner_context(
ChainEventFactory::consumption_gap_event(
tracker.writer_id,
SeqNo(reader.0 + 1),
advertised,
progress.stage_id,
),
);
if let Err(e) = crate::supervised_base::publication::append(
&tracker.journal,
gap_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append gap event; skipping emission"
);
}
}
}
}
}
}
} else if let Some(advertised) = progress.advertised_writer_seq {
if advertised.0 != progress.reader_seq.0 {
pass = false;
let cause = EventViolationCause::SeqDivergence {
advertised: Some(advertised),
reader: progress.reader_seq,
};
failure_reason = Some(cause.clone());
aggregate_violation_for_journal = true;
aggregate_failure_reason = Some(cause.clone());
if advertised.0 > progress.reader_seq.0 {
let gap_event =
tracker.with_owner_context(ChainEventFactory::consumption_gap_event(
tracker.writer_id,
SeqNo(progress.reader_seq.0 + 1),
advertised,
progress.stage_id,
));
if let Err(e) = crate::supervised_base::publication::append(
&tracker.journal,
gap_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append gap event; skipping emission"
);
}
}
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Violated {
upstream: progress.stage_id,
cause,
};
}
}
}
if let Some(feed_failure) = direct_feed_evidence.iter().find(|feed| !feed.pass) {
let cause = feed_failure
.reason
.clone()
.unwrap_or(EventViolationCause::SeqDivergence {
advertised: Some(feed_failure.advertised_writer_seq),
reader: feed_failure.reader_seq,
});
pass = false;
failure_reason = Some(cause.clone());
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Violated {
upstream: progress.stage_id,
cause,
};
}
}
let status_reason = failure_reason.clone();
if aggregate_violation_for_journal {
if let Some(EventViolationCause::SeqDivergence { advertised, reader }) =
aggregate_failure_reason.clone()
{
let violation_event =
tracker.with_owner_context(ChainEventFactory::at_least_once_violation_event(
tracker.writer_id,
progress.stage_id,
EventViolationCause::SeqDivergence { advertised, reader },
progress.reader_seq,
progress.advertised_writer_seq,
));
if let Err(e) = crate::supervised_base::publication::append(
&tracker.journal,
violation_event,
Default::default(),
)
.await
{
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append at_least_once_violation event; skipping emission"
);
}
}
}
let final_event = tracker.with_owner_context(ChainEventFactory::consumption_final_event(
tracker.writer_id,
ConsumptionFinalEventParams {
pass,
consumed_count: Count(progress_seq.0),
expected_count: None,
eof_seen: true,
last_event_id: progress_last_event_id,
reader_seq: progress_seq,
advertised_writer_seq: progress.advertised_writer_seq,
advertised_vector_clock: progress_vector_clock,
failure_reason,
},
));
let final_append_ok = match crate::supervised_base::publication::append(
&tracker.journal,
final_event,
Default::default(),
)
.await
{
Ok(_) => true,
Err(e) => {
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append final event; skipping state update"
);
false
}
};
let mut status_append_ok = true;
if !direct_feed_evidence.is_empty() {
if let Some(reader_stage) = tracker.reader_stage {
status_append_ok = self
.emit_direct_feed_contract_system_events(
tracker,
progress.stage_id,
reader_stage,
&direct_feed_evidence,
)
.await;
}
} else if let (Some(system_journal), Some(reader_stage)) =
(&tracker.system_journal, tracker.reader_stage)
{
let status_event = SystemEvent::new(
tracker.writer_id,
SystemPayload::ContractStatus {
upstream: progress.stage_id,
reader: reader_stage,
selected_event_type: selected_event_type.clone(),
feed_role,
pass,
reader_seq: Some(progress_seq),
advertised_writer_seq: progress.advertised_writer_seq,
reason: status_reason,
},
);
if let Err(e) = crate::supervised_base::publication::append(
system_journal,
status_event,
Default::default(),
)
.await
{
status_append_ok = false;
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader = ?reader_stage,
reader_index = index,
error = %e,
"Failed to append contract status; skipping state update"
);
}
}
if final_append_ok && status_append_ok {
progress.final_emitted = true;
progress.contract_violated = !pass;
}
}
async fn check_stall_for_reader(
&mut self,
progress: &mut ReaderProgress,
index: usize,
now: Instant,
status: &mut ContractStatus,
) {
let Some(tracker) = &self.contract_tracker else {
return;
};
let Some(last) = progress.last_read_instant else {
return;
};
let elapsed = now.duration_since(last).as_millis() as u64;
if elapsed >= tracker.config.stall_threshold.0 {
progress.consecutive_stall_checks += 1;
if progress.consecutive_stall_checks >= tracker.config.stall_checks_before_emit
&& progress.stalled_since.is_none()
{
if tracker.config.stall_cooloff.0 > 0 {
if let Some(last_emitted) = progress.last_stall_emitted_instant {
let cooloff_elapsed = now.duration_since(last_emitted).as_millis() as u64;
if cooloff_elapsed < tracker.config.stall_cooloff.0 {
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Stalled(progress.stage_id);
}
return;
}
}
}
let stall_since_candidate = Some(last);
let stalled_duration = DurationMs(elapsed);
let stalled_event =
tracker.with_owner_context(ChainEventFactory::reader_stalled_event(
tracker.writer_id,
progress.stage_id,
stalled_duration,
));
let stalled_append_ok = match crate::supervised_base::publication::append(
&tracker.journal,
stalled_event,
Default::default(),
)
.await
{
Ok(_) => true,
Err(e) => {
tracing::error!(
target: "flowip-105",
owner = %self.owner_label,
upstream = ?progress.stage_id,
reader_index = index,
error = %e,
"Failed to append stalled event; skipping state update"
);
false
}
};
if !matches!(status, ContractStatus::Violated { .. }) {
*status = ContractStatus::Stalled(progress.stage_id);
}
if stalled_append_ok {
progress.stalled_since = stall_since_candidate;
progress.last_stall_emitted_instant = Some(now);
}
}
} else {
progress.consecutive_stall_checks = 0;
}
}
fn should_emit_progress(&self, progress: &ReaderProgress, index: usize, now: Instant) -> bool {
let Some(tracker) = &self.contract_tracker else {
return false;
};
let delta_events = self
.progress_seq(progress)
.0
.saturating_sub(progress.last_progress_seq.0);
let time_elapsed = progress
.last_progress_instant
.map(|t| now.duration_since(t).as_millis() as u64)
.unwrap_or(0);
delta_events >= tracker.config.progress_min_events.0
|| time_elapsed >= tracker.config.progress_max_interval.0
|| self.state.is_reader_eof(index)
}
pub fn track_output_event(&mut self) {
if let Some(tracker) = &mut self.contract_tracker {
tracker.output_events_written.0 += 1;
tracing::trace!(
"Tracked output event, total: {}",
tracker.output_events_written.0
);
}
}
pub fn should_check_contracts(&self, reader_progress: &[ReaderProgress]) -> bool {
if let Some(tracker) = &self.contract_tracker {
let now = Instant::now();
for progress in reader_progress {
let Some(last) = progress.last_progress_instant else {
return true;
};
let elapsed = now.duration_since(last).as_millis() as u64;
if elapsed >= tracker.config.progress_max_interval.0 / 2 {
return true;
}
}
}
false
}
pub async fn maybe_check_contracts(
&mut self,
reader_progress: &mut [ReaderProgress],
) -> Option<ContractStatus> {
if self.should_check_contracts(reader_progress) {
Some(self.check_contracts(reader_progress).await)
} else {
None
}
}
pub async fn maybe_check_contracts_tick(
&mut self,
reader_progress: &mut [ReaderProgress],
last_contract_check: &mut Option<Instant>,
) -> Option<ContractStatus> {
let Some(tracker) = &self.contract_tracker else {
return None;
};
let now = Instant::now();
let tick_ms = (tracker.config.progress_max_interval.0 / 2).max(1);
let due = match last_contract_check {
Some(last) => now.duration_since(*last).as_millis() as u64 >= tick_ms,
None => true,
};
if !due {
return None;
}
*last_contract_check = Some(now);
Some(self.check_contracts(reader_progress).await)
}
pub async fn maybe_check_contracts_tick_diagnostics_only(
&mut self,
reader_progress: &mut [ReaderProgress],
last_contract_check: &mut Option<Instant>,
) -> Option<ContractStatus> {
let Some(tracker) = &self.contract_tracker else {
return None;
};
let now = Instant::now();
let tick_ms = (tracker.config.progress_max_interval.0 / 2).max(1);
let due = match last_contract_check {
Some(last) => now.duration_since(*last).as_millis() as u64 >= tick_ms,
None => true,
};
if !due {
return None;
}
*last_contract_check = Some(now);
Some(self.check_contracts_diagnostics_only(reader_progress).await)
}
}