use super::{ExecutorError, ManagedSigner};
use crate::{BackendError, EffectObservation, EffectStatus, TransportEffectStatus};
use newton_submission_protocol::{EffectOutcome, ExecutionOutcome};
use newton_submission_service::{ExecutableIntent, JobRecord, SignerRole, Store};
use newton_task_submission::{BatchIntentItem, SubmissionId, TaskOperation};
use tracing::info;
pub(super) async fn classify_preflight_failure(
store: &Store,
signer: &ManagedSigner,
job: &JobRecord,
nonce: u64,
batch_error: BackendError,
) -> Result<(), ExecutorError> {
if job.intent.signer_role() == SignerRole::Transporter {
let status = signer.backend.classify_transport_effect(&job.intent).await?;
let outcome = transport_preflight_outcome(status, &batch_error);
store.complete_execution(job.job_id, &outcome).await?;
return Ok(());
}
let observations = match signer.backend.classify_onchain_effects(&job.intent).await {
Ok(observations) => observations,
Err(error) if error.is_transient() => {
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, job.job_id)
.await?;
return Err(error.into());
}
Err(error) => return Err(error.into()),
};
let Some(observations) = ordered_effect_observations(&job.intent, &observations) else {
return Err(ExecutorError::EffectCoverage);
};
let (contract_role, items, original_operation) = match &job.intent {
ExecutableIntent::BatchCreateAndRespond { contract_role, items } => {
(contract_role, items, TaskOperation::CombinedCreateAndRespond)
}
ExecutableIntent::BatchRespond { contract_role, items } => (contract_role, items, TaskOperation::RespondOnly),
_ => return Err(BackendError::Policy("unsupported task preflight intent".to_string()).into()),
};
let mut effects = Vec::with_capacity(items.len());
for observation in observations {
let effect = match observation.status {
EffectStatus::Verified => EffectOutcome::Succeeded,
EffectStatus::TaskOnly => EffectOutcome::PartiallySucceeded,
EffectStatus::Conflict => EffectOutcome::Conflict,
EffectStatus::Missing => {
let item = items
.iter()
.find(|item| item.submission_id == observation.submission_id)
.ok_or(ExecutorError::EffectCoverage)?;
let single = match original_operation {
TaskOperation::CombinedCreateAndRespond => ExecutableIntent::BatchCreateAndRespond {
contract_role: contract_role.clone(),
items: vec![item.clone()],
},
TaskOperation::RespondOnly => ExecutableIntent::BatchRespond {
contract_role: contract_role.clone(),
items: vec![item.clone()],
},
};
match signer.backend.prepare(&single, nonce, None).await {
Ok(_) | Err(BackendError::SimulationRetryable(_)) => EffectOutcome::Missing,
Err(error) if error.is_transient() => {
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, job.job_id)
.await?;
return Err(error.into());
}
Err(error) => {
EffectOutcome::Failed(format!("item_preflight:{error};batch_preflight:{batch_error}"))
}
}
}
};
effects.push(effect);
}
let outcome = ExecutionOutcome::Effects(effects);
store.complete_execution(job.job_id, &outcome).await?;
log_effects(items, &outcome, "preflight");
Ok(())
}
pub(super) async fn classify_and_project(
store: &Store,
signer: &ManagedSigner,
job: &JobRecord,
receipt_succeeded: Option<bool>,
) -> Result<(), ExecutorError> {
if job.intent.signer_role() == SignerRole::Transporter {
if receipt_succeeded == Some(false) {
return complete_transport_effect(store, signer, job, "transaction_reverted").await;
}
if receipt_succeeded == Some(true) {
store
.complete_execution(job.job_id, &ExecutionOutcome::Effects(vec![EffectOutcome::Succeeded]))
.await?;
return Ok(());
}
return complete_transport_effect(store, signer, job, "consumed_nonce_without_verifiable_effect").await;
}
let observations = signer.backend.classify_onchain_effects(&job.intent).await?;
let observations = ordered_effect_observations(&job.intent, &observations).ok_or(ExecutorError::EffectCoverage)?;
let effects = observations
.into_iter()
.map(|observation| match observation.status {
EffectStatus::Verified => EffectOutcome::Succeeded,
EffectStatus::Missing => EffectOutcome::Missing,
EffectStatus::TaskOnly => EffectOutcome::PartiallySucceeded,
EffectStatus::Conflict => EffectOutcome::Conflict,
})
.collect();
let outcome = ExecutionOutcome::Effects(effects);
store.complete_execution(job.job_id, &outcome).await?;
log_effects(
job.intent.task_items(),
&outcome,
match receipt_succeeded {
Some(true) => "successful_receipt",
Some(false) => "reverted_receipt",
None => "consumed_nonce",
},
);
Ok(())
}
fn log_effects(items: &[BatchIntentItem], outcome: &ExecutionOutcome, evidence: &'static str) {
let ExecutionOutcome::Effects(effects) = outcome else {
return;
};
for (item, effect) in items.iter().zip(effects) {
info!(
submission_id = %item.submission_id,
task_id = %item.task.taskId,
effect = effect_label(effect),
effect_error = effect_error(effect),
evidence,
"on-chain task effect classified"
);
}
}
const fn effect_label(effect: &EffectOutcome) -> &'static str {
match effect {
EffectOutcome::Succeeded => "succeeded",
EffectOutcome::Missing => "missing",
EffectOutcome::PartiallySucceeded => "partially_succeeded",
EffectOutcome::Conflict => "conflict",
EffectOutcome::Failed(_) => "failed",
}
}
fn effect_error(effect: &EffectOutcome) -> &str {
match effect {
EffectOutcome::Failed(error) => error.as_str(),
_ => "",
}
}
async fn complete_transport_effect(
store: &Store,
signer: &ManagedSigner,
job: &JobRecord,
missing_error: &'static str,
) -> Result<(), ExecutorError> {
let status = signer.backend.classify_transport_effect(&job.intent).await?;
let outcome = transport_observation_outcome(status, missing_error);
store.complete_execution(job.job_id, &outcome).await?;
Ok(())
}
fn transport_preflight_outcome(status: TransportEffectStatus, error: &BackendError) -> ExecutionOutcome {
match status {
TransportEffectStatus::Verified => ExecutionOutcome::Effects(vec![EffectOutcome::Succeeded]),
TransportEffectStatus::Missing if error.is_transport_destination_behind() => {
ExecutionOutcome::RetryableFailure("destination_clock_behind".to_string())
}
TransportEffectStatus::Missing => ExecutionOutcome::PermanentFailure(format!("preflight:{error}")),
TransportEffectStatus::Conflict => ExecutionOutcome::PermanentFailure("onchain_effect_conflict".to_string()),
}
}
fn transport_observation_outcome(status: TransportEffectStatus, missing_error: &'static str) -> ExecutionOutcome {
match status {
TransportEffectStatus::Verified => ExecutionOutcome::Effects(vec![EffectOutcome::Succeeded]),
TransportEffectStatus::Missing => ExecutionOutcome::RetryableFailure(missing_error.to_string()),
TransportEffectStatus::Conflict => ExecutionOutcome::PermanentFailure("onchain_effect_conflict".to_string()),
}
}
pub(super) async fn requeue_after_cancellation(store: &Store, job: &JobRecord) -> Result<(), ExecutorError> {
store
.complete_execution(
job.job_id,
&ExecutionOutcome::RetryableFailure("transaction_cancelled".to_string()),
)
.await?;
Ok(())
}
pub(super) async fn finalize_preflight_failure(
store: &Store,
job: &JobRecord,
error: &BackendError,
) -> Result<(), ExecutorError> {
store
.complete_execution(
job.job_id,
&ExecutionOutcome::PermanentFailure(format!("preflight:{error}")),
)
.await?;
Ok(())
}
fn ordered_effect_observations(
intent: &ExecutableIntent,
observations: &[EffectObservation],
) -> Option<Vec<EffectObservation>> {
let items = match intent {
ExecutableIntent::BatchCreateAndRespond { items, .. } | ExecutableIntent::BatchRespond { items, .. } => items,
_ => return None,
};
let submission_ids = items.iter().map(|item| item.submission_id).collect::<Vec<_>>();
order_observations(&submission_ids, observations)
}
fn order_observations(
submission_ids: &[SubmissionId],
observations: &[EffectObservation],
) -> Option<Vec<EffectObservation>> {
if observations.len() != submission_ids.len() {
return None;
}
submission_ids
.iter()
.map(|submission_id| {
let mut matches = observations
.iter()
.filter(|observation| observation.submission_id == *submission_id);
let observation = *matches.next()?;
matches.next().is_none().then_some(observation)
})
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn transport_preflight_matrix_distinguishes_clock_lag_from_deterministic_failure() {
let destination_behind = BackendError::Simulation("GlobalTableRootInFuture".to_string());
assert_eq!(
transport_preflight_outcome(TransportEffectStatus::Missing, &destination_behind),
ExecutionOutcome::RetryableFailure("destination_clock_behind".to_string())
);
let invalid_root = BackendError::Simulation("InvalidGlobalTableRoot".to_string());
assert_eq!(
transport_preflight_outcome(TransportEffectStatus::Missing, &invalid_root),
ExecutionOutcome::PermanentFailure(
"preflight:transaction simulation failed: InvalidGlobalTableRoot".to_string()
)
);
assert_eq!(
transport_preflight_outcome(TransportEffectStatus::Verified, &invalid_root),
ExecutionOutcome::Effects(vec![EffectOutcome::Succeeded])
);
assert_eq!(
transport_preflight_outcome(TransportEffectStatus::Conflict, &destination_behind),
ExecutionOutcome::PermanentFailure("onchain_effect_conflict".to_string())
);
}
#[test]
fn transport_post_nonce_matrix_retries_only_missing_effects() {
assert_eq!(
transport_observation_outcome(
TransportEffectStatus::Missing,
"consumed_nonce_without_verifiable_effect"
),
ExecutionOutcome::RetryableFailure("consumed_nonce_without_verifiable_effect".to_string())
);
assert_eq!(
transport_observation_outcome(TransportEffectStatus::Verified, "transaction_reverted"),
ExecutionOutcome::Effects(vec![EffectOutcome::Succeeded])
);
assert_eq!(
transport_observation_outcome(TransportEffectStatus::Conflict, "transaction_reverted"),
ExecutionOutcome::PermanentFailure("onchain_effect_conflict".to_string())
);
}
#[test]
fn effect_observations_are_restored_to_request_order() {
let first = SubmissionId::new();
let second = SubmissionId::new();
let ordered = order_observations(
&[first, second],
&[
EffectObservation {
submission_id: second,
status: EffectStatus::Conflict,
},
EffectObservation {
submission_id: first,
status: EffectStatus::Verified,
},
],
)
.expect("exact effect coverage");
assert_eq!(ordered[0].submission_id, first);
assert_eq!(ordered[0].status, EffectStatus::Verified);
assert_eq!(ordered[1].submission_id, second);
assert_eq!(ordered[1].status, EffectStatus::Conflict);
}
#[test]
fn duplicate_effect_observations_are_rejected() {
let first = SubmissionId::new();
let second = SubmissionId::new();
assert!(order_observations(
&[first, second],
&[
EffectObservation {
submission_id: first,
status: EffectStatus::Verified,
},
EffectObservation {
submission_id: first,
status: EffectStatus::Missing,
},
],
)
.is_none());
}
}