use super::*;
pub(super) fn effect_record_from_event(
event: &ChainEvent,
) -> Result<Option<EffectRecord>, EffectError> {
use obzenflow_core::event::payloads::execution_payload::ExecutionPayload;
match &event.payload {
ChainPayload::Execution(
ExecutionPayload::EffectAttemptStarted(_)
| ExecutionPayload::EffectRecoveryAbandoned(_),
) => Ok(None),
ChainPayload::Execution(ExecutionPayload::EffectRecord(record)) => {
let provenance = event.effect_provenance.as_ref().ok_or_else(|| {
EffectError::EffectProvenanceMismatch(
"framework effect record is missing effect provenance".to_string(),
)
})?;
if !provenance.fact_owner.is_framework() {
return Err(EffectError::EffectProvenanceMismatch(
"framework effect record is not framework-owned".to_string(),
));
}
validate_effect_record_provenance(&event.event_type(), record, provenance)?;
Ok(Some(record.clone()))
}
ChainPayload::Fact(_) | ChainPayload::CompositeData(_) => {
let Some(provenance) = event.effect_provenance.as_ref() else {
return Ok(None);
};
if provenance.fact_owner.is_framework() {
return Err(EffectError::EffectProvenanceMismatch(
"framework evidence must use the typed execution contract".to_string(),
));
}
let outcome_fact_ordinal = provenance.outcome_fact_ordinal.ok_or_else(|| {
EffectError::EffectProvenanceMismatch(
"domain effect outcome must set outcome_fact_ordinal".to_string(),
)
})?;
let outcome_fact_count = provenance.outcome_fact_count.ok_or_else(|| {
EffectError::EffectProvenanceMismatch(
"domain effect outcome must set outcome_fact_count".to_string(),
)
})?;
let record = EffectRecord {
cursor: provenance.cursor.clone(),
descriptor_hash: provenance.descriptor_hash.clone(),
descriptor: provenance.descriptor.clone(),
outcome: EffectOutcomePayload::SucceededFact {
event_kind: event.payload.kind(),
event_type: event.event_type().into(),
output: event.payload(),
outcome_fact_ordinal,
outcome_fact_count,
},
origin: provenance.origin.clone(),
};
validate_domain_effect_record_provenance(&record, provenance)?;
Ok(Some(record))
}
ChainPayload::Execution(_) | ChainPayload::Delivery(_) | ChainPayload::FlowControl(_) => {
Ok(None)
}
}
}
fn validate_domain_effect_record_provenance(
record: &EffectRecord,
provenance: &EffectProvenance,
) -> Result<(), EffectError> {
if provenance.cursor != record.cursor {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance cursor does not match domain effect record cursor".to_string(),
));
}
if provenance.descriptor_hash != record.descriptor_hash {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance descriptor_hash does not match domain effect record descriptor_hash"
.to_string(),
));
}
if provenance.descriptor != record.descriptor {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance descriptor does not match domain effect record descriptor"
.to_string(),
));
}
let expected_group_id = effect_outcome_group_id(&record.cursor);
if provenance.group_id.as_ref() != Some(&expected_group_id) {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect_provenance group_id does not match deterministic group id `{expected_group_id}`"
)));
}
let EffectOutcomePayload::SucceededFact {
outcome_fact_ordinal,
..
} = &record.outcome
else {
return Err(EffectError::EffectProvenanceMismatch(
"domain effect outcome facts must use SucceededFact records".to_string(),
));
};
if provenance.outcome_fact_ordinal != Some(*outcome_fact_ordinal) {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance outcome_fact_ordinal does not match domain effect record ordinal"
.to_string(),
));
}
Ok(())
}
fn validate_effect_record_provenance(
event_type: &str,
record: &EffectRecord,
provenance: &EffectProvenance,
) -> Result<(), EffectError> {
let expected_event_type = framework_effect_event_type(&record.descriptor.effect_type);
if event_type != expected_event_type {
return Err(EffectError::EffectProvenanceMismatch(format!(
"reserved framework effect event type `{event_type}` does not match record descriptor `{}` (expected `{expected_event_type}`)",
record.descriptor.effect_type
)));
}
if provenance.cursor != record.cursor {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance cursor does not match effect record cursor".to_string(),
));
}
if provenance.descriptor_hash != record.descriptor_hash {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance descriptor_hash does not match effect record descriptor_hash"
.to_string(),
));
}
if provenance.descriptor != record.descriptor {
return Err(EffectError::EffectProvenanceMismatch(
"effect_provenance descriptor does not match effect record descriptor".to_string(),
));
}
let expected_group_id = effect_outcome_group_id(&record.cursor);
if provenance.group_id.as_ref() != Some(&expected_group_id) {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect_provenance group_id does not match deterministic group id `{expected_group_id}`"
)));
}
if provenance.outcome_fact_ordinal.is_some() {
return Err(EffectError::EffectProvenanceMismatch(
"framework effect record compatibility facts must not set outcome_fact_ordinal"
.to_string(),
));
}
Ok(())
}
pub(super) fn validate_effect_outcome_group(records: &[&EffectRecord]) -> Result<(), EffectError> {
let Some(first) = records.first() else {
return Ok(());
};
let expected = match &first.outcome {
EffectOutcomePayload::Succeeded { .. } | EffectOutcomePayload::Failed { .. } => {
if records.len() != 1 {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {:?} mixes a non-fact outcome with other records",
first.cursor
)));
}
return Ok(());
}
EffectOutcomePayload::SucceededFact {
outcome_fact_count, ..
} => usize::try_from(outcome_fact_count.get()).map_err(|_| {
EffectError::EffectProvenanceMismatch(format!(
"effect outcome fact count for cursor {:?} exceeds usize range",
first.cursor
))
})?,
};
if expected == 0 {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {:?} declares zero facts",
first.cursor
)));
}
let cursor = first.cursor.clone();
let descriptor_hash = first.descriptor_hash.clone();
let descriptor = first.descriptor.clone();
let mut seen = vec![false; expected];
for record in records {
if record.cursor != cursor {
return Err(EffectError::EffectProvenanceMismatch(
"effect outcome group contains multiple cursors".to_string(),
));
}
if record.descriptor_hash != descriptor_hash {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} contains multiple descriptor hashes"
)));
}
if record.descriptor != descriptor {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} contains multiple descriptors"
)));
}
let EffectOutcomePayload::SucceededFact {
outcome_fact_ordinal,
outcome_fact_count,
..
} = &record.outcome
else {
return Err(EffectError::EffectProvenanceMismatch(format!(
"multi-fact effect outcome group for cursor {cursor:?} contains a non-domain-success record"
)));
};
if usize::try_from(outcome_fact_count.get()).unwrap_or(usize::MAX) != expected {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} disagrees on outcome_fact_count"
)));
}
let ordinal = usize::try_from(outcome_fact_ordinal.get()).map_err(|_| {
EffectError::EffectProvenanceMismatch(format!(
"effect outcome fact ordinal for cursor {cursor:?} exceeds usize range"
))
})?;
let Some(slot) = seen.get_mut(ordinal) else {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} has ordinal {outcome_fact_ordinal} beyond declared count {expected}"
)));
};
if *slot {
return Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} has duplicate ordinal {outcome_fact_ordinal}"
)));
}
*slot = true;
}
let present_count = seen.iter().filter(|flag| **flag).count();
if present_count == expected {
return Ok(());
}
if seen[..present_count].iter().all(|flag| *flag) {
return Err(EffectError::IncompleteOutcomeGroup {
cursor,
expected,
present: present_count,
});
}
Err(EffectError::EffectProvenanceMismatch(format!(
"effect outcome group for cursor {cursor:?} is missing interior outcome fact ordinals"
)))
}
pub(super) fn decode_effect_outcome_group<T>(records: &[&EffectRecord]) -> Result<T, EffectError>
where
T: TypedFactSet,
{
validate_effect_outcome_group(records)?;
let [single] = records else {
let mut ordered = records.to_vec();
ordered.sort_by_key(|record| match &record.outcome {
EffectOutcomePayload::SucceededFact {
outcome_fact_ordinal,
..
} => outcome_fact_ordinal.get(),
_ => u32::MAX,
});
let facts = ordered
.iter()
.map(|record| match &record.outcome {
EffectOutcomePayload::SucceededFact {
event_kind,
event_type,
output,
..
} => Ok(TypedFact {
event_type: event_type.clone(),
payload: ChainPayload::decode(*event_kind, event_type.as_str(), output.clone())
.map_err(|error| EffectError::Serialization(error.to_string()))?,
}),
_ => Err(EffectError::EffectProvenanceMismatch(
"multi-fact effect outcome group contains a non-domain-success record"
.to_string(),
)),
})
.collect::<Result<Vec<_>, _>>()?;
return T::try_from_facts(&facts).map_err(effect_fact_set_error);
};
match &single.outcome {
EffectOutcomePayload::SucceededFact {
event_kind,
event_type,
output,
..
} => T::try_from_facts(&[TypedFact {
event_type: event_type.clone(),
payload: ChainPayload::decode(*event_kind, event_type.as_str(), output.clone())
.map_err(|error| EffectError::Serialization(error.to_string()))?,
}])
.map_err(effect_fact_set_error),
EffectOutcomePayload::Succeeded { output } => {
let fact_types = T::fact_types();
let [fact_type] = fact_types.as_slice() else {
return Err(EffectError::EffectProvenanceMismatch(
"legacy single-payload effect success cannot reconstruct a multi-fact output"
.to_string(),
));
};
T::try_from_facts(&[TypedFact {
event_type: fact_type.event_type.clone(),
payload: ChainPayload::Fact(output.clone()),
}])
.map_err(effect_fact_set_error)
}
EffectOutcomePayload::Failed { .. } => recorded_failure_from_outcome(&single.outcome),
}
}
#[derive(Debug)]
pub(super) enum EffectRecordMaterialization {
DomainFacts {
facts: Vec<TypedFact>,
origin: Option<EffectFactOrigin>,
},
FrameworkRecords(Vec<EffectRecord>),
}
fn group_origin(records: &[&EffectRecord]) -> Result<Option<EffectFactOrigin>, EffectError> {
let mut origins = records.iter().map(|record| &record.origin);
let first = origins.next().cloned().flatten();
if origins.any(|origin| origin.as_ref() != first.as_ref()) {
return Err(EffectError::EffectProvenanceMismatch(
"effect outcome group records disagree on fact origin".to_string(),
));
}
Ok(first)
}
pub(super) fn effect_record_group_materialization(
records: &[&EffectRecord],
) -> Result<EffectRecordMaterialization, EffectError> {
validate_effect_outcome_group(records)?;
let [single] = records else {
let origin = group_origin(records)?;
let mut ordered = records.to_vec();
ordered.sort_by_key(|record| match &record.outcome {
EffectOutcomePayload::SucceededFact {
outcome_fact_ordinal,
..
} => outcome_fact_ordinal.get(),
_ => u32::MAX,
});
let mut facts = Vec::new();
for record in ordered {
let EffectOutcomePayload::SucceededFact {
event_kind,
event_type,
output,
..
} = &record.outcome
else {
return Err(EffectError::EffectProvenanceMismatch(
"multi-fact effect outcome group contains a non-domain-success record"
.to_string(),
));
};
facts.push(TypedFact {
event_type: event_type.clone(),
payload: ChainPayload::decode(*event_kind, event_type.as_str(), output.clone())
.map_err(|error| EffectError::Serialization(error.to_string()))?,
});
}
return Ok(EffectRecordMaterialization::DomainFacts { facts, origin });
};
match &single.outcome {
EffectOutcomePayload::SucceededFact {
event_kind,
event_type,
output,
..
} => Ok(EffectRecordMaterialization::DomainFacts {
facts: vec![TypedFact {
event_type: event_type.clone(),
payload: ChainPayload::decode(*event_kind, event_type.as_str(), output.clone())
.map_err(|error| EffectError::Serialization(error.to_string()))?,
}],
origin: single.origin.clone(),
}),
EffectOutcomePayload::Succeeded { .. } | EffectOutcomePayload::Failed { .. } => {
Ok(EffectRecordMaterialization::FrameworkRecords(vec![
(*single).clone(),
]))
}
}
}
pub(super) fn is_routable_output_fact(
output_contract: Option<&StageOutputContract>,
event_type: &str,
) -> bool {
match output_contract {
Some(contract) if !contract.is_empty() => contract.is_routable_event_type(event_type),
_ => false,
}
}
pub(super) fn decode_effect_outcome<T>(outcome: &EffectOutcomePayload) -> Result<T, EffectError>
where
T: DeserializeOwned,
{
match outcome {
EffectOutcomePayload::Succeeded { output } => serde_json::from_value(output.clone())
.map_err(|e| EffectError::Serialization(e.to_string())),
EffectOutcomePayload::SucceededFact { output, .. } => {
serde_json::from_value(output.clone())
.map_err(|e| EffectError::Serialization(e.to_string()))
}
EffectOutcomePayload::Failed {
error_type,
error_message,
retry,
cause,
detail,
} => {
validate_failure_detail(error_type, detail.as_ref())?;
Err(EffectError::RecordedFailure {
error_type: error_type.clone(),
error_message: error_message.clone(),
retry: *retry,
cause: cause.clone(),
detail: detail.clone().map(Box::new),
})
}
}
}
pub(super) fn recorded_failure_from_outcome<T>(
outcome: &EffectOutcomePayload,
) -> Result<T, EffectError> {
match outcome {
EffectOutcomePayload::Failed {
error_type,
error_message,
retry,
cause,
detail,
} => {
validate_failure_detail(error_type, detail.as_ref())?;
Err(EffectError::RecordedFailure {
error_type: error_type.clone(),
error_message: error_message.clone(),
retry: *retry,
cause: cause.clone(),
detail: detail.clone().map(Box::new),
})
}
_ => Err(EffectError::EffectProvenanceMismatch(
"expected recorded effect failure".to_string(),
)),
}
}
fn validate_failure_detail(
error_type: &EffectFailureKind,
detail: Option<&EffectFailureDetail>,
) -> Result<(), EffectError> {
let invariant_type = error_type.as_str() == "target_invariant_violation";
let invariant_detail = matches!(
detail,
Some(EffectFailureDetail::PortBindingInvariantViolation { .. })
);
if invariant_type != invariant_detail {
return Err(EffectError::EffectProvenanceMismatch(
"effect port-binding invariant failure type/detail pairing is invalid".to_string(),
));
}
Ok(())
}
pub(super) fn effect_fact_set_error(error: TypedFactSetError) -> EffectError {
match error {
TypedFactSetError::SerializationFailed(message) => EffectError::Serialization(message),
TypedFactSetError::DeserializationFailed { event_type, error } => {
EffectError::Serialization(format!("{event_type}: {error}"))
}
TypedFactSetError::MissingFact { event_type } => EffectError::EffectProvenanceMismatch(
format!("effect outcome group is missing fact `{event_type}`"),
),
TypedFactSetError::DuplicateFact { event_type } => EffectError::EffectProvenanceMismatch(
format!("effect outcome group has duplicate fact `{event_type}`"),
),
TypedFactSetError::UnexpectedFact { event_type } => EffectError::EffectProvenanceMismatch(
format!("effect outcome group contains unexpected fact `{event_type}`"),
),
}
}