newton-tx-executor 0.7.3

Durable allowlisted transaction executor for Newton submissions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
//! Tracks one persisted nonce through replacement, cancellation, and finality.
//!
//! Every attempt in this module reuses the job's durable nonce. A replacement
//! is persisted before broadcast, and the assignment remains owned until chain
//! effects are projected or a mined cancellation safely requeues the items.

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)?;
        // Negative receipt evidence is consequential only when it can trigger a
        // reorg transition or same-nonce replacement. Ordinary polls query one
        // rotating attempt/provider and simply keep waiting on a miss.
        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;
        }
        // A previously seen receipt disappeared from every provider without a
        // read error, so return the attempt to pending and wait for canonicality.
        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;
            }
            // Deterministic preparation can reproduce the current attempt when
            // market fees have not risen enough. Rebroadcast instead of adding
            // a duplicate journal row.
            if signed.transaction_hash == current.transaction_hash {
                let results = signer.backend.broadcast(&current.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(());
        }
    }
}

/// Classifies the persisted nonce after a watchdog expiry. A consumed nonce is
/// not merely a reason to suppress replacement: its contract effect must be
/// projected so the durable assignment can finish even when receipts were
/// pruned, lost, or unavailable from every configured provider.
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));
    }
}