newton-tx-executor 0.7.3

Durable allowlisted transaction executor for Newton submissions
//! Classifies authoritative chain state into producer-neutral submission outcomes.
//!
//! Transaction receipts prove nonce consumption. Domain planners project the
//! resulting effect report into their own aggregates after the executor stores
//! it through the common execution lifecycle.

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 {
        // The worker's quarantine branch owns release for non-terminal hard
        // failures. Releasing here as well would turn the second release into
        // a store state race and incorrectly stop the whole executor.
        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) {
            // A reverted receipt proves only that the nonce was consumed. A
            // concurrent writer may already have produced the requested
            // idempotent effect, so classify authoritative destination state.
            return complete_transport_effect(store, signer, job, "transaction_reverted").await;
        }
        // Root visibility is a producer-owned workflow barrier. A successful
        // finalized receipt is the common transaction boundary.
        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(())
}

/// Reports what authoritative chain state proved about each batch member.
///
/// The batch verdict a receipt carries is not the per-task verdict: an atomic
/// batch can mine successfully while an individual effect is absent or in
/// conflict. These lines run inside the job span, so they inherit `job_id`,
/// the signer lane, and the batch identities. They follow the durable write, so
/// a failed completion is retried rather than reported twice.
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())
        }
        // A deterministic simulation failure with a still-missing effect cannot
        // be repaired by replaying the same immutable intent.
        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]),
        // After nonce consumption, absence is retryable because a fresh
        // reconciliation can safely emit the idempotent effect again.
        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());
    }
}