use super::{sleep_or_cancel, ExecutorConfig, ExecutorError, ManagedSigner};
use crate::{ObservationScope, SignedAttempt, TransactionObservation};
use alloy::primitives::U256;
use newton_submission_service::{AttemptKind, AttemptRecord, AttemptState, JobRecord, PreparedAttempt, Store};
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use super::outcomes::{classify_and_project, requeue_after_cancellation};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum NonceDisposition {
Replaceable,
Consumed,
Inconclusive,
}
pub(super) async fn track_attempt(
store: &Store,
signer: &ManagedSigner,
config: &ExecutorConfig,
job: &JobRecord,
mut attempts: Vec<AttemptRecord>,
cancellation: CancellationToken,
) -> Result<(), ExecutorError> {
if attempts.is_empty() {
return Err(ExecutorError::MissingAttempt);
}
let mut mined_observed = attempts.iter().any(|candidate| candidate.state == AttemptState::Mined);
let mut observation_round = 0_u64;
let mut candidate_cursor = 0_usize;
let mut replacement_count = attempts
.iter()
.filter(|candidate| candidate.kind == AttemptKind::Replacement)
.count()
.try_into()
.unwrap_or(u32::MAX);
let mut cancellation_prepared = attempts
.iter()
.any(|candidate| candidate.kind == AttemptKind::Cancellation);
let required_confirmations = signer.backend.finality_confirmations().max(1);
let current = attempts.last().ok_or(ExecutorError::MissingAttempt)?;
info!(
attempt_count = attempts.len(),
attempt_number = current.attempt_number,
attempt_kind = current.kind.as_str(),
nonce = current.nonce,
tx_hash = %current.transaction_hash,
required_confirmations,
"tracking transaction receipt and finality"
);
let mut finality_logged = false;
'tracking: loop {
if cancellation.is_cancelled() {
return Ok(());
}
let mut pending_age = attempts
.last()
.map(|attempt| elapsed_since_ms(attempt.updated_at_ms))
.ok_or(ExecutorError::MissingAttempt)?;
let exhaustive = mined_observed || pending_age >= config.watchdog_timeout;
let mut observed_receipt = false;
let mut pending_receipt_evidence = false;
let mut receipt_evidence_incomplete = false;
let candidate_indices = if exhaustive {
(0..attempts.len()).rev().collect::<Vec<_>>()
} else {
let candidate_index = attempts.len() - 1 - (candidate_cursor % attempts.len());
candidate_cursor = candidate_cursor.wrapping_add(1);
vec![candidate_index]
};
for candidate_index in candidate_indices {
let candidate = attempts[candidate_index].clone();
let observation = match signer
.backend
.observe(
candidate.transaction_hash,
candidate.receipt_provider.as_deref(),
if exhaustive {
ObservationScope::Exhaustive
} else {
ObservationScope::Economical
},
)
.await
{
Ok(observation) => observation,
Err(error) if error.is_transient() => {
warn!(
signer_id = %signer.signer_id,
transaction_hash = %candidate.transaction_hash,
%error,
"transient receipt observation failure"
);
receipt_evidence_incomplete = true;
continue;
}
Err(error) => return Err(error.into()),
};
let receipt_succeeded = matches!(&observation, TransactionObservation::Mined { .. });
let (provider, block_number, block_hash, confirmations, receipt) = match observation {
TransactionObservation::Pending => {
pending_receipt_evidence = true;
continue;
}
TransactionObservation::PendingInconclusive {
pending_providers,
failed_providers,
} => {
pending_receipt_evidence = true;
receipt_evidence_incomplete = true;
warn!(
signer_id = %signer.signer_id,
transaction_hash = %candidate.transaction_hash,
pending_providers,
failed_providers,
"partial pending receipt evidence; replacement remains eligible but reorg decisions are deferred"
);
continue;
}
TransactionObservation::Reverted {
provider,
block_number,
block_hash,
confirmations,
receipt,
}
| TransactionObservation::Mined {
provider,
block_number,
block_hash,
confirmations,
receipt,
} => (provider, block_number, block_hash, confirmations, receipt),
};
observed_receipt = true;
mined_observed = true;
let first_mined_observation = candidate.state != AttemptState::Mined;
store
.record_mined(
job.job_id,
candidate.attempt_number,
&receipt,
&provider,
block_number,
block_hash,
)
.await?;
if first_mined_observation {
info!(
attempt_number = candidate.attempt_number,
attempt_kind = candidate.kind.as_str(),
nonce = candidate.nonce,
tx_hash = %candidate.transaction_hash,
receipt_provider = %provider,
block_number,
%block_hash,
confirmations,
required_confirmations,
receipt_succeeded,
"transaction receipt observed"
);
}
attempts[candidate_index].state = AttemptState::Mined;
attempts[candidate_index].receipt_provider = Some(provider);
attempts[candidate_index].updated_at_ms = super::unix_time_ms();
store.record_confirmations(job.job_id, confirmations).await?;
if confirmations >= required_confirmations {
if !finality_logged {
info!(
attempt_number = candidate.attempt_number,
attempt_kind = candidate.kind.as_str(),
nonce = candidate.nonce,
tx_hash = %candidate.transaction_hash,
block_number,
confirmations,
required_confirmations,
receipt_succeeded,
"transaction reached required confirmations"
);
finality_logged = true;
}
let projection = if candidate.kind == AttemptKind::Cancellation {
requeue_after_cancellation(store, job).await
} else {
classify_and_project(store, signer, job, Some(receipt_succeeded)).await
};
match projection {
Ok(()) => {
info!(
attempt_number = candidate.attempt_number,
attempt_kind = candidate.kind.as_str(),
nonce = candidate.nonce,
tx_hash = %candidate.transaction_hash,
"confirmed transaction outcome committed"
);
return Ok(());
}
Err(error) if error.is_retryable() => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
%error,
"transient finalized on-chain effect classification failure"
);
if sleep_or_cancel(config.receipt_poll_interval, &cancellation).await {
return Ok(());
}
continue 'tracking;
}
Err(error) => return Err(error),
}
}
break;
}
if exhaustive && !observed_receipt && mined_observed && !receipt_evidence_incomplete {
store.record_receipt_reorg(job.job_id).await?;
mined_observed = false;
let now = super::unix_time_ms();
for attempt in &mut attempts {
if attempt.state == AttemptState::Mined {
attempt.state = AttemptState::Broadcast;
attempt.receipt_provider = None;
attempt.updated_at_ms = now;
}
}
pending_age = Duration::ZERO;
}
let nonce_disposition = if !observed_receipt && pending_age >= config.watchdog_timeout {
observe_nonce_disposition(signer, job, attempts.last().ok_or(ExecutorError::MissingAttempt)?.nonce).await?
} else {
NonceDisposition::Inconclusive
};
if nonce_disposition == NonceDisposition::Consumed {
match classify_and_project(store, signer, job, None).await {
Ok(()) => {
info!(
nonce = attempts.last().ok_or(ExecutorError::MissingAttempt)?.nonce,
"submission outcome recovered from on-chain effects without a receipt"
);
return Ok(());
}
Err(error) if error.is_retryable() => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
%error,
"consumed nonce effect classification remains inconclusive"
);
}
Err(error) => return Err(error),
}
}
if !observed_receipt && pending_receipt_evidence && nonce_disposition == NonceDisposition::Replaceable {
let current = attempts.last().ok_or(ExecutorError::MissingAttempt)?;
let previous_fees = (current.max_fee_per_gas, current.max_priority_fee_per_gas);
let (prepared, kind) = if !cancellation_prepared && replacement_count < config.cancel_after_bumps {
(
signer
.backend
.prepare(&job.intent, current.nonce, Some(previous_fees))
.await,
AttemptKind::Replacement,
)
} else {
(
signer.backend.prepare_cancellation(current.nonce, previous_fees).await,
AttemptKind::Cancellation,
)
};
let signed = match prepared {
Ok(signed) => signed,
Err(error) if error.is_transient() => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
%error,
"transient replacement preparation failure"
);
if sleep_or_cancel(config.receipt_poll_interval, &cancellation).await {
return Ok(());
}
continue 'tracking;
}
Err(error) => return Err(error.into()),
};
if signed.nonce != current.nonce {
return Err(ExecutorError::AttemptNonceMismatch);
}
let required = maximum_transaction_cost(&signed);
let balance = match signer.backend.balance().await {
Ok(balance) => balance,
Err(error) if error.is_transient() => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
%error,
"transient balance read failure before replacement persistence"
);
if sleep_or_cancel(config.receipt_poll_interval, &cancellation).await {
return Ok(());
}
continue 'tracking;
}
Err(error) => return Err(error.into()),
};
if balance < required {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
%balance,
%required,
retry_after = ?config.underfunded_poll_interval,
"busy signer cannot cover replacement; retaining assignment and nonce"
);
if sleep_or_cancel(config.underfunded_poll_interval, &cancellation).await {
return Ok(());
}
continue 'tracking;
}
if kind == AttemptKind::Cancellation {
cancellation_prepared = true;
}
if signed.transaction_hash == current.transaction_hash {
let results = signer.backend.broadcast(¤t.raw_transaction).await;
store
.record_broadcast(job.job_id, current.attempt_number, &rmp_serde::to_vec(&results)?)
.await?;
let accepted_provider = results
.iter()
.find(|result| result.result.is_ok())
.map(|result| result.provider.as_str());
if let Some(accepted_provider) = accepted_provider {
info!(
attempt_number = current.attempt_number,
attempt_kind = current.kind.as_str(),
nonce = current.nonce,
tx_hash = %current.transaction_hash,
accepted_provider,
provider_attempts = results.len(),
"same-nonce transaction rebroadcast accepted"
);
} else {
warn!(
attempt_number = current.attempt_number,
attempt_kind = current.kind.as_str(),
nonce = current.nonce,
tx_hash = %current.transaction_hash,
provider_attempts = results.len(),
"all providers rejected same-nonce rebroadcast; exact bytes remain durable"
);
}
if let Some(current) = attempts.last_mut() {
current.updated_at_ms = super::unix_time_ms();
}
candidate_cursor = 0;
continue;
}
let attempt_number = store
.record_prepared(&PreparedAttempt {
job_id: job.job_id,
signer_id: signer.signer_id.clone(),
signer_address: signer.backend.address(),
chain_id: signer.chain_id,
nonce: signed.nonce,
gas_limit: signed.gas_limit,
max_fee_per_gas: signed.max_fee_per_gas,
max_priority_fee_per_gas: signed.max_priority_fee_per_gas,
raw_transaction: signed.raw_transaction.clone(),
transaction_hash: signed.transaction_hash,
kind,
replaces_attempt_number: Some(current.attempt_number),
})
.await?;
info!(
attempt_number,
attempt_kind = kind.as_str(),
replaces_attempt_number = current.attempt_number,
nonce = signed.nonce,
tx_hash = %signed.transaction_hash,
previous_tx_hash = %current.transaction_hash,
gas_limit = signed.gas_limit,
max_fee_per_gas = signed.max_fee_per_gas,
max_priority_fee_per_gas = signed.max_priority_fee_per_gas,
"same-nonce transaction attempt prepared and durably recorded"
);
let results = signer.backend.broadcast(&signed.raw_transaction).await;
store
.record_broadcast(job.job_id, attempt_number, &rmp_serde::to_vec(&results)?)
.await?;
let accepted_provider = results
.iter()
.find(|result| result.result.is_ok())
.map(|result| result.provider.as_str());
if let Some(accepted_provider) = accepted_provider {
info!(
attempt_number,
attempt_kind = kind.as_str(),
nonce = signed.nonce,
tx_hash = %signed.transaction_hash,
accepted_provider,
provider_attempts = results.len(),
"same-nonce transaction attempt broadcast accepted"
);
} else {
warn!(
attempt_number,
attempt_kind = kind.as_str(),
nonce = signed.nonce,
tx_hash = %signed.transaction_hash,
provider_attempts = results.len(),
"all providers rejected same-nonce transaction attempt; exact bytes remain durable"
);
}
attempts.push(AttemptRecord {
attempt_number,
nonce: signed.nonce,
gas_limit: signed.gas_limit,
max_fee_per_gas: signed.max_fee_per_gas,
max_priority_fee_per_gas: signed.max_priority_fee_per_gas,
raw_transaction: signed.raw_transaction,
transaction_hash: signed.transaction_hash,
kind,
state: AttemptState::Broadcast,
receipt_provider: None,
updated_at_ms: super::unix_time_ms(),
});
if kind == AttemptKind::Replacement {
replacement_count = replacement_count.saturating_add(1);
}
candidate_cursor = 0;
pending_age = Duration::ZERO;
}
let delay = receipt_poll_delay(
config,
signer.block_time,
if mined_observed { Duration::ZERO } else { pending_age },
attempts.last().ok_or(ExecutorError::MissingAttempt)?.transaction_hash,
observation_round,
);
observation_round = observation_round.wrapping_add(1);
if sleep_or_cancel(delay, &cancellation).await {
return Ok(());
}
}
}
async fn observe_nonce_disposition(
signer: &ManagedSigner,
job: &JobRecord,
nonce: u64,
) -> Result<NonceDisposition, ExecutorError> {
match signer.backend.latest_transaction_count().await {
Ok(latest) if latest == nonce => Ok(NonceDisposition::Replaceable),
Ok(latest) if latest > nonce => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
nonce,
latest,
"receipt is unavailable but the nonce is consumed; reconciling contract effects"
);
Ok(NonceDisposition::Consumed)
}
Ok(latest) => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
nonce,
latest,
"latest nonce is behind the durable assignment; deferring replacement"
);
Ok(NonceDisposition::Inconclusive)
}
Err(error) if error.is_transient() => {
warn!(
signer_id = %signer.signer_id,
job_id = %job.job_id,
nonce,
%error,
"cannot verify nonce before replacement; deferring"
);
Ok(NonceDisposition::Inconclusive)
}
Err(error) => Err(error.into()),
}
}
fn elapsed_since_ms(timestamp_ms: i64) -> Duration {
let elapsed_ms = u64::try_from(super::unix_time_ms().saturating_sub(timestamp_ms).max(0)).unwrap_or_default();
Duration::from_millis(elapsed_ms)
}
fn receipt_poll_delay(
config: &ExecutorConfig,
block_time: Duration,
pending_age: Duration,
transaction_hash: alloy::primitives::B256,
round: u64,
) -> Duration {
let chain_interval = (block_time / 2).min(Duration::from_secs(10));
let base = config.receipt_poll_interval.max(chain_interval);
let maximum = config.receipt_poll_max_interval.max(base);
let multiplier = if pending_age < base.saturating_mul(3) {
1
} else if pending_age < base.saturating_mul(6) {
2
} else if pending_age < base.saturating_mul(12) {
4
} else {
8
};
let delay = base.saturating_mul(multiplier).min(maximum);
jitter(delay, maximum, transaction_hash, round)
}
fn jitter(delay: Duration, maximum: Duration, transaction_hash: alloy::primitives::B256, round: u64) -> Duration {
if delay < Duration::from_millis(100) {
return delay;
}
let bytes = transaction_hash.as_slice();
let seed = u64::from(u16::from_be_bytes([bytes[0], bytes[1]])) ^ round;
let percent = 90_u128 + u128::from(seed % 21);
let jittered_ms = delay.as_millis().saturating_mul(percent).saturating_div(100);
Duration::from_millis(u64::try_from(jittered_ms).unwrap_or(u64::MAX)).min(maximum)
}
pub(super) fn maximum_transaction_cost(attempt: &SignedAttempt) -> U256 {
U256::from(attempt.gas_limit)
.checked_mul(U256::from(attempt.max_fee_per_gas))
.unwrap_or(U256::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
use alloy::primitives::B256;
fn config() -> ExecutorConfig {
ExecutorConfig {
poll_interval: Duration::from_secs(5),
receipt_poll_interval: Duration::from_secs(2),
receipt_poll_max_interval: Duration::from_secs(30),
watchdog_timeout: Duration::from_secs(120),
cancel_after_bumps: 5,
underfunded_poll_interval: Duration::from_secs(900),
}
}
#[test]
fn receipt_polling_is_chain_aware_and_bounded() {
let config = config();
let hash = B256::ZERO;
let initial = receipt_poll_delay(&config, Duration::from_secs(12), Duration::ZERO, hash, 0);
assert!((Duration::from_millis(5_400)..=Duration::from_millis(6_600)).contains(&initial));
let old = receipt_poll_delay(&config, Duration::from_secs(12), Duration::from_secs(300), hash, 0);
assert!(old <= Duration::from_secs(30));
assert!(old >= Duration::from_secs(27));
}
#[test]
fn configured_floor_is_not_reduced_for_fast_chains() {
let config = config();
let delay = receipt_poll_delay(&config, Duration::from_secs(2), Duration::ZERO, B256::ZERO, 0);
assert!(delay >= Duration::from_millis(1_800));
}
}