use obzenflow_core::event::ChainPayload;
mod construction;
mod contract_checking;
mod polling;
mod types;
#[cfg(test)]
mod tests;
pub use super::subscription_poller::{PollResult, SubscriptionPoller};
use types::{AdvertisedWriterSeqByEventType, SelectedDataSeqByEventType};
pub use types::{
CompositeEntrySpec, ContractConfig, ContractStatus, ContractTracker, ContractsWiring,
DeliveredCount, DeliveredOrdinal, EofOutcome, FeedIdentity, MergeCandidateStatus,
MergeWaitState, ReaderProgress, ReaderSelectionPolicy, ReaderTiebreakKey, SelectedFeedMetadata,
SelectedFeedRole, StageInputPosition, StageKey, SubscriptionState,
};
use crate::contracts::ContractChain;
use crate::control_plane::ControlPlaneProvider;
use crate::feed_plan::declared_event_type_matches;
use crate::messaging::upstream_subscription_policy::ContractPolicyStack;
use obzenflow_core::event::payloads::delivery_payload::DeliveryResult;
use obzenflow_core::event::types::SeqNo;
use obzenflow_core::event::vector_clock::VectorClock;
use obzenflow_core::event::{ChainEvent, JournalEvent, JournalRecord};
use obzenflow_core::journal::reader::JournalReader;
use obzenflow_core::{AdmissionSeq, EventId, EventType, ReaderGeneration, StageId};
use std::collections::{BTreeMap, HashMap, HashSet};
use std::sync::Arc;
use tokio::time::Instant;
pub(crate) struct ReceiptSettlement {
owner_label: String,
receipt_aware: bool,
reader_stage: Option<StageId>,
upstreams: Vec<StageId>,
chains: Vec<Option<ContractChain>>,
progress: Vec<ReaderProgress>,
}
impl ReceiptSettlement {
pub(crate) fn record(&mut self, receipt: &ChainEvent) -> Option<(SeqNo, EventId, VectorClock)> {
record_receipt(
receipt,
&self.owner_label,
self.receipt_aware,
self.reader_stage,
&self.upstreams,
&mut self.chains,
&mut self.progress,
)
}
}
fn record_receipt(
receipt: &ChainEvent,
owner_label: &str,
receipt_aware: bool,
reader_stage: Option<StageId>,
upstreams: &[StageId],
chains: &mut [Option<ContractChain>],
reader_progress: &mut [ReaderProgress],
) -> Option<(SeqNo, EventId, VectorClock)> {
if !receipt_aware {
return None;
}
let Some(parent_id) = receipt.causality.parent_ids.first().copied() else {
tracing::warn!(
owner = %owner_label,
receipt_id = %receipt.id,
"record_delivery_receipt: receipt missing parent causality"
);
return None;
};
let Some(index) = reader_progress
.iter()
.enumerate()
.find_map(|(index, progress)| {
progress
.pending_delivery_inputs
.contains_key(&parent_id)
.then_some(index)
})
else {
tracing::warn!(
owner = %owner_label,
receipt_id = %receipt.id,
?parent_id,
"record_delivery_receipt: no pending delivered input for parent"
);
return None;
};
let upstream_stage = reader_progress[index].stage_id;
let ChainPayload::Delivery(payload) = &receipt.payload else {
tracing::warn!(
owner = %owner_label,
receipt_id = %receipt.id,
?upstream_stage,
?parent_id,
"record_delivery_receipt: non-delivery event passed to receipt recorder"
);
return None;
};
let is_accounted_receipt = reader_progress[index]
.pending_receipts
.contains_key(&parent_id);
if is_accounted_receipt {
if let Some(reader_stage) = reader_stage {
if let Some(slot) = upstreams.iter().position(|id| *id == upstream_stage) {
if let Some(Some(chain)) = chains.get_mut(slot) {
chain.on_write(receipt, reader_stage, SeqNo(0));
}
}
}
}
if matches!(&payload.result, DeliveryResult::Buffered { .. }) {
return None;
}
reader_progress[index]
.pending_delivery_inputs
.remove(&parent_id);
if !is_accounted_receipt {
tracing::debug!(
owner = %owner_label,
?upstream_stage,
?parent_id,
"record_delivery_receipt: terminal forwarded input has no receipt-watermark position"
);
return None;
}
let previous_seq = reader_progress[index].receipted_seq;
if reader_progress[index].mark_receipted(parent_id) {
reader_progress[index].last_read_instant = Some(Instant::now());
if reader_progress[index].receipted_seq != previous_seq {
if let (Some(event_id), Some(vector_clock)) = (
reader_progress[index].last_receipted_event_id,
reader_progress[index].last_receipted_vector_clock.clone(),
) {
return Some((reader_progress[index].receipted_seq, event_id, vector_clock));
}
}
} else {
tracing::debug!(
owner = %owner_label,
?upstream_stage,
?parent_id,
"record_delivery_receipt: parent was not pending when receipt arrived"
);
}
None
}
struct FeedContractChain {
metadata: SelectedFeedMetadata,
chain: ContractChain,
last_contract_result_seq: SeqNo,
}
pub(super) struct ReaderSlot<T: JournalEvent> {
pub(super) stage_id: StageId,
pub(super) stage_key: StageKey,
pub(super) reader: Box<dyn JournalReader<T>>,
}
pub(super) struct HeldHead<T: JournalEvent> {
pub(super) envelope: JournalRecord<T::Payload>,
pub(super) is_authored_eof: bool,
pub(super) is_drain: bool,
pub(super) catch_up: Option<ReaderGeneration>,
pub(super) orders_by_own_seq: bool,
}
pub struct MergeCandidateMeta<'a> {
pub generation: ReaderGeneration,
pub ordinal: DeliveredOrdinal,
pub key: &'a ReaderTiebreakKey,
pub vector_clock: &'a VectorClock,
pub is_authored_eof: bool,
pub admission_seq: Option<AdmissionSeq>,
}
impl FeedContractChain {
fn new(metadata: SelectedFeedMetadata, chain: ContractChain) -> Self {
Self {
metadata,
chain,
last_contract_result_seq: SeqNo(0),
}
}
}
pub struct UpstreamSubscription<T>
where
T: JournalEvent,
{
delivery_filter: DeliveryFilter,
owner_label: String,
readers: Vec<ReaderSlot<T>>,
archived_stage_ids_by_current: HashMap<StageId, StageId>,
selected_event_types_by_stage: HashMap<StageId, HashSet<EventType>>,
selected_feeds_by_stage: HashMap<StageId, Vec<SelectedFeedMetadata>>,
composite_entries_by_stage: HashMap<StageId, Vec<CompositeEntrySpec>>,
selected_data_seq_by_reader: Vec<SeqNo>,
selected_data_seq_by_reader_event_type: Vec<SelectedDataSeqByEventType>,
advertised_writer_seq_by_reader_event_type: Vec<AdvertisedWriterSeqByEventType>,
state: SubscriptionState,
contract_tracker: Option<ContractTracker>,
contract_chains: Vec<Option<ContractChain>>,
contract_feed_chains: Vec<Vec<FeedContractChain>>,
contract_policies: Vec<Option<ContractPolicyStack>>,
control_plane: Arc<dyn ControlPlaneProvider>,
last_eof_outcome: Option<EofOutcome>,
last_delivered_upstream_stage: Option<StageId>,
next_stage_input_position: u64,
last_delivered_stage_input_position: Option<StageInputPosition>,
reader_selection: types::ReaderSelectionPolicy,
seq_ordered: bool,
entered_generation: ReaderGeneration,
held_heads: Vec<Option<HeldHead<T>>>,
delivered_count_by_reader: Vec<DeliveredCount>,
reader_tiebreak_keys: Vec<ReaderTiebreakKey>,
generation_by_reader: Vec<ReaderGeneration>,
last_positional_seq: Vec<AdmissionSeq>,
last_delivered_generation: Option<ReaderGeneration>,
last_merge_wait: Option<types::MergeWaitState>,
merge_candidate_index: Option<usize>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum DeliveryFilter {
All,
TransportOnly,
}
impl<T> UpstreamSubscription<T>
where
T: JournalEvent + 'static,
{
pub fn last_delivered_upstream_stage(&self) -> Option<StageId> {
self.last_delivered_upstream_stage
}
pub fn last_delivered_stage_input_position(&self) -> Option<StageInputPosition> {
self.last_delivered_stage_input_position
}
pub fn merge_wait(&self) -> Option<&types::MergeWaitState> {
self.last_merge_wait.as_ref()
}
pub fn delivered_counts(&self) -> &[DeliveredCount] {
&self.delivered_count_by_reader
}
pub fn reader_selection(&self) -> types::ReaderSelectionPolicy {
self.reader_selection
}
pub fn seq_ordered(&self) -> bool {
self.seq_ordered
}
pub(crate) fn held_head_count(&self) -> usize {
self.held_heads.iter().filter(|head| head.is_some()).count()
}
fn uses_receipt_watermark(&self) -> bool {
self.contract_tracker
.as_ref()
.map(|tracker| tracker.receipt_aware_progress)
.unwrap_or(false)
}
fn progress_seq(&self, progress: &ReaderProgress) -> SeqNo {
if self.uses_receipt_watermark() {
progress.receipted_seq
} else {
progress.reader_seq
}
}
fn progress_last_event_id(&self, progress: &ReaderProgress) -> Option<EventId> {
if self.uses_receipt_watermark() {
progress.last_receipted_event_id
} else {
progress.last_event_id
}
}
fn progress_vector_clock(&self, progress: &ReaderProgress) -> Option<VectorClock> {
if self.uses_receipt_watermark() {
progress.last_receipted_vector_clock.clone()
} else {
progress.last_vector_clock.clone()
}
}
fn has_selected_event_type_filter(&self, stage_id: StageId) -> bool {
self.selected_event_types_by_stage
.get(&stage_id)
.is_some_and(|selected| !selected.is_empty())
}
fn data_event_selected_for_stage(&self, stage_id: StageId, event_type: &str) -> bool {
self.selected_event_types_by_stage
.get(&stage_id)
.filter(|selected| !selected.is_empty())
.map(|selected| {
selected.iter().any(|selected_event_type| {
declared_event_type_matches(selected_event_type.as_str(), event_type, None)
})
})
.unwrap_or(true)
}
fn selected_writer_seq_for_reader(&self, reader_index: usize, stage_id: StageId) -> SeqNo {
if self.has_selected_event_type_filter(stage_id) {
self.selected_data_seq_by_reader
.get(reader_index)
.copied()
.unwrap_or(SeqNo(0))
} else {
SeqNo(0)
}
}
fn selected_writer_seq_from_eof_map(
&self,
stage_id: StageId,
writer_seq_by_event_type: &BTreeMap<EventType, SeqNo>,
) -> Option<SeqNo> {
let selected = self.selected_event_types_by_stage.get(&stage_id)?;
if selected.is_empty() || writer_seq_by_event_type.is_empty() {
return None;
}
let selected_total = writer_seq_by_event_type
.iter()
.filter(|(actual_event_type, _)| {
selected.iter().any(|selected_event_type| {
declared_event_type_matches(
selected_event_type.as_str(),
actual_event_type.as_str(),
None,
)
})
})
.fold(0u64, |total, (_, seq)| total.saturating_add(seq.0));
Some(SeqNo(selected_total))
}
fn selected_feed_matches_event_type(feed: &SelectedFeedMetadata, event_type: &str) -> bool {
feed.matches_event_type(event_type)
}
fn selected_reader_seq_for_feed(
&self,
reader_index: usize,
feed: &SelectedFeedMetadata,
) -> SeqNo {
self.selected_data_seq_by_reader_event_type
.get(reader_index)
.map(|reader_by_type| reader_by_type.seq_for_feed(feed))
.unwrap_or(SeqNo(0))
}
fn advertised_writer_seq_for_feed(
&self,
reader_index: usize,
feed: &SelectedFeedMetadata,
) -> Option<SeqNo> {
self.advertised_writer_seq_by_reader_event_type
.get(reader_index)
.and_then(|advertised_by_type| advertised_by_type.seq_for_feed(feed))
}
fn unique_selected_feed_for_stage(
&self,
stage_id: StageId,
) -> (
Option<obzenflow_core::EventType>,
Option<obzenflow_core::event::payloads::system_payload::SystemFeedRole>,
) {
let Some(feeds) = self.selected_feeds_by_stage.get(&stage_id) else {
return (None, None);
};
let mut unique_feeds = feeds.iter();
let Some(first) = unique_feeds.next() else {
return (None, None);
};
if unique_feeds.next().is_some() {
return (None, None);
}
(Some(first.event_type().clone()), first.system_feed_role())
}
pub fn notify_delivery_receipt(&mut self, receipt: &ChainEvent, upstream_stage: StageId) {
let Some(reader_stage) = self.contract_tracker.as_ref().and_then(|t| t.reader_stage) else {
return;
};
let Some(index) = self
.readers
.iter()
.position(|slot| slot.stage_id == upstream_stage)
else {
tracing::warn!(
owner = %self.owner_label,
?upstream_stage,
"notify_delivery_receipt: no reader slot for upstream stage"
);
return;
};
let Some(chain_slot) = self.contract_chains.get_mut(index) else {
return;
};
let Some(chain) = chain_slot.as_mut() else {
return;
};
chain.on_write(receipt, reader_stage, SeqNo(0));
}
pub fn record_delivery_receipt(
&mut self,
receipt: &ChainEvent,
reader_progress: &mut [ReaderProgress],
) -> Option<(SeqNo, EventId, VectorClock)> {
record_receipt(
receipt,
&self.owner_label,
self.uses_receipt_watermark(),
self.contract_tracker
.as_ref()
.and_then(|tracker| tracker.reader_stage),
&self
.readers
.iter()
.map(|slot| slot.stage_id)
.collect::<Vec<_>>(),
&mut self.contract_chains,
reader_progress,
)
}
pub(crate) fn take_receipt_settlement(
&mut self,
progress: &mut Vec<ReaderProgress>,
) -> ReceiptSettlement {
ReceiptSettlement {
owner_label: self.owner_label.clone(),
receipt_aware: self.uses_receipt_watermark(),
reader_stage: self
.contract_tracker
.as_ref()
.and_then(|tracker| tracker.reader_stage),
upstreams: self.readers.iter().map(|slot| slot.stage_id).collect(),
chains: std::mem::take(&mut self.contract_chains),
progress: std::mem::take(progress),
}
}
pub(crate) fn restore_receipt_settlement(
&mut self,
progress: &mut Vec<ReaderProgress>,
settlement: ReceiptSettlement,
) {
self.contract_chains = settlement.chains;
*progress = settlement.progress;
}
pub fn pending_receipt_envelope(
&self,
parent_event_id: EventId,
reader_progress: &[ReaderProgress],
) -> Option<(StageId, JournalRecord<ChainPayload>)> {
reader_progress.iter().find_map(|progress| {
progress
.pending_delivery_inputs
.get(&parent_event_id)
.map(|pending| {
let envelope = pending.clone();
(progress.stage_id, envelope)
})
})
}
pub fn take_last_eof_outcome(&mut self) -> Option<EofOutcome> {
self.last_eof_outcome.take()
}
pub fn last_eof_outcome(&self) -> Option<&EofOutcome> {
self.last_eof_outcome.as_ref()
}
pub fn has_pending(&self) -> bool {
self.state.has_pending()
}
pub fn upstream_count(&self) -> usize {
self.readers.len()
}
pub fn all_readers_eof(&self) -> bool {
self.state.eof_count() == self.readers.len()
}
pub fn all_readers_logically_eof(&self) -> bool {
self.state.logical_eof_count() == self.readers.len()
}
pub fn last_delivered_generation(&self) -> Option<ReaderGeneration> {
self.last_delivered_generation
}
pub fn delivered_data_count(&self) -> u64 {
self.next_stage_input_position - 1
}
pub fn all_readers_caught_up(&self, target: ReaderGeneration) -> bool {
(0..self.readers.len()).all(|index| {
self.generation_by_reader[index] >= target || self.state.is_reader_eof(index)
})
}
pub fn max_reader_generation(&self) -> ReaderGeneration {
self.generation_by_reader
.iter()
.copied()
.max()
.unwrap_or_default()
}
pub fn has_upstream(&self) -> bool {
!self.readers.is_empty()
}
}
#[async_trait::async_trait]
impl<T> SubscriptionPoller for UpstreamSubscription<T>
where
T: JournalEvent + 'static,
{
type Event = T;
async fn poll_next(&mut self) -> PollResult<Self::Event> {
UpstreamSubscription::poll_next(self).await
}
fn name(&self) -> &str {
"upstream_subscription"
}
}