use obzenflow_core::event::{ChainPayload, JournalRecord};
use obzenflow_core::EventId;
use std::collections::{HashMap, HashSet};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeliveredOrderRow {
pub upstream_stage_key: String,
pub per_input_ordinal: u64,
pub event_type: String,
pub payload: String,
pub parent_event_id: EventId,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DeliveredOrderProjection {
pub rows: Vec<DeliveredOrderRow>,
}
impl DeliveredOrderProjection {
pub fn from_envelopes(
stage_outputs: &[JournalRecord<ChainPayload>],
upstreams: &[(String, Vec<JournalRecord<ChainPayload>>)],
) -> Self {
let mut parent_index: HashMap<EventId, &str> = HashMap::new();
for (stage_key, envelopes) in upstreams {
for envelope in envelopes {
if envelope.consumes_data_credit() {
parent_index.insert(envelope.envelope.provenance.event.id, stage_key.as_str());
}
}
}
let mut ordinals: HashMap<&str, u64> = HashMap::new();
let mut assigned: HashMap<EventId, (String, u64)> = HashMap::new();
let mut rows = Vec::new();
for envelope in stage_outputs {
let event = &envelope.authored();
if !event.consumes_data_credit() && !event.is_delivery() {
continue;
}
let Some(parent_id) = event.causality.parent_ids.first().copied() else {
continue;
};
let Some(stage_key) = parent_index.get(&parent_id).copied() else {
continue;
};
let (upstream_stage_key, per_input_ordinal) = assigned
.entry(parent_id)
.or_insert_with(|| {
let next = ordinals.entry(stage_key).or_insert(0);
*next += 1;
(stage_key.to_string(), *next)
})
.clone();
let payload = if event.consumes_data_credit() {
serde_json::to_string(&event.payload.contract_body().expect("typed payload"))
.expect("JSON body")
} else {
String::new()
};
rows.push(DeliveredOrderRow {
upstream_stage_key,
per_input_ordinal,
event_type: event.event_type(),
payload,
parent_event_id: parent_id,
});
}
Self { rows }
}
pub fn consumption_sequence(&self) -> Vec<(String, u64)> {
let mut sequence = Vec::new();
let mut seen = HashSet::new();
for row in &self.rows {
if seen.insert(row.parent_event_id) {
sequence.push((row.upstream_stage_key.clone(), row.per_input_ordinal));
}
}
sequence
}
pub fn assert_equal(&self, other: &Self) {
let comparable = |projection: &Self| {
projection
.rows
.iter()
.map(|row| {
(
row.upstream_stage_key.clone(),
row.per_input_ordinal,
row.event_type.clone(),
row.payload.clone(),
)
})
.collect::<Vec<_>>()
};
let left = comparable(self);
let right = comparable(other);
for (index, (a, b)) in left.iter().zip(right.iter()).enumerate() {
assert_eq!(
a, b,
"delivered-order projections diverge at row {index}: {a:?} vs {b:?}"
);
}
assert_eq!(
left.len(),
right.len(),
"delivered-order projections differ in length: {} vs {}",
left.len(),
right.len()
);
}
}