mod outcomes;
mod tracking;
use self::{
outcomes::{classify_and_project, classify_preflight_failure, finalize_preflight_failure},
tracking::{maximum_transaction_cost, track_attempt},
};
use crate::{BackendError, ChainBackend, ObservationScope, TransactionObservation};
use alloy::primitives::{Address, U256};
use newton_submission_protocol::SignerId;
use newton_submission_service::{
ActiveAssignment, AttemptKind, AttemptRecord, AttemptState, ExecutableIntent, JobRecord, PreparedAttempt,
RuntimeHealth, RuntimeTask, SignerLifecycle, SignerRecord, SignerRole, Store, StoreError,
};
use newton_task_submission::{submission_ids_field, task_ids_field};
use std::{
sync::Arc,
time::{Duration, SystemTime, UNIX_EPOCH},
};
use tokio::{task::JoinSet, time::MissedTickBehavior};
use tokio_util::sync::CancellationToken;
use tracing::{error, info, info_span, warn, Instrument, Span};
#[derive(Clone)]
pub struct ManagedSigner {
pub signer_id: SignerId,
pub chain_id: u64,
pub role: SignerRole,
pub block_time: Duration,
pub backend: Arc<dyn ChainBackend>,
}
impl std::fmt::Debug for ManagedSigner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ManagedSigner")
.field("signer_id", &self.signer_id)
.field("chain_id", &self.chain_id)
.field("role", &self.role)
.field("block_time", &self.block_time)
.field("address", &self.backend.address())
.finish()
}
}
impl ManagedSigner {
fn address(&self) -> Address {
self.backend.address()
}
fn record(&self) -> SignerRecord {
SignerRecord {
signer_id: self.signer_id.clone(),
chain_id: self.chain_id,
role: self.role,
address: self.address(),
}
}
}
#[derive(Debug, Clone)]
pub struct ExecutorConfig {
pub poll_interval: Duration,
pub receipt_poll_interval: Duration,
pub receipt_poll_max_interval: Duration,
pub watchdog_timeout: Duration,
pub cancel_after_bumps: u32,
pub underfunded_poll_interval: Duration,
}
#[derive(Debug)]
pub struct Executor {
store: Arc<Store>,
signers: Vec<ManagedSigner>,
config: ExecutorConfig,
health: Arc<RuntimeHealth>,
}
impl Executor {
pub fn new(
store: Arc<Store>,
signers: Vec<ManagedSigner>,
config: ExecutorConfig,
health: Arc<RuntimeHealth>,
) -> Self {
Self {
store,
signers,
config,
health,
}
}
pub async fn run(&self, cancellation: CancellationToken) -> Result<(), ExecutorError> {
self.health.heartbeat(RuntimeTask::Executor);
let configured_signers = self.signers.iter().map(ManagedSigner::record).collect::<Vec<_>>();
self.store.register_signers(&configured_signers).await?;
let lane_cancellation = cancellation.child_token();
let mut lanes = JoinSet::new();
for signer in self.signers.clone() {
let store = self.store.clone();
let config = self.config.clone();
let cancellation = lane_cancellation.clone();
let health = self.health.clone();
lanes.spawn(async move {
if recover_lane(&store, &signer, &config, &cancellation, &health).await? == LaneRecovery::Stopped {
return Ok(());
}
worker_loop(store, signer, config, cancellation, health).await
});
}
let mut heartbeat = tokio::time::interval(self.config.poll_interval.max(Duration::from_secs(1)));
heartbeat.set_missed_tick_behavior(MissedTickBehavior::Skip);
let failure = loop {
tokio::select! {
_ = cancellation.cancelled() => break None,
_ = heartbeat.tick() => self.health.heartbeat(RuntimeTask::Executor),
joined = lanes.join_next() => match joined {
Some(Ok(Ok(()))) if cancellation.is_cancelled() => break None,
Some(Ok(Ok(()))) => break Some(ExecutorError::WorkerExited),
Some(Ok(Err(error))) => break Some(error),
Some(Err(error)) => break Some(ExecutorError::WorkerPanicked(error.to_string())),
None => break Some(ExecutorError::WorkerExited),
},
}
};
lane_cancellation.cancel();
while let Some(joined) = lanes.join_next().await {
if let Err(error) = joined {
return Err(ExecutorError::WorkerPanicked(error.to_string()));
}
}
if let Some(error) = failure {
return Err(error);
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LaneRecovery {
Ready,
Stopped,
}
async fn recover_lane(
store: &Store,
signer: &ManagedSigner,
config: &ExecutorConfig,
cancellation: &CancellationToken,
health: &RuntimeHealth,
) -> Result<LaneRecovery, ExecutorError> {
loop {
match reconcile(store, signer, config, cancellation.clone()).await {
Ok(()) => {
health.heartbeat(RuntimeTask::Executor);
health.record_rpc(signer.chain_id, Ok(()));
newton_metric::inc_signer_recoveries_total(signer.chain_id, "ready");
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), true);
return Ok(LaneRecovery::Ready);
}
Err(error) if error.waits_for_chain() => {
health.heartbeat(RuntimeTask::Executor);
if matches!(error, ExecutorError::UnownedPendingTransaction { .. }) {
store
.mark_signer_waiting(&signer.signer_id, signer.chain_id, &error.to_string())
.await?;
}
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), false);
warn!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
%error,
"signer lane is waiting for authoritative chain state"
);
}
Err(error) if error.is_underfunded() => {
health.heartbeat(RuntimeTask::Executor);
record_backend_error(health, signer.chain_id, &error);
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), false);
warn!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
%error,
retry_after = ?config.underfunded_poll_interval,
"signer lane is underfunded; waiting before balance recheck"
);
if sleep_or_cancel(config.underfunded_poll_interval, cancellation).await {
return Ok(LaneRecovery::Stopped);
}
continue;
}
Err(error) if error.is_transient() => {
health.heartbeat(RuntimeTask::Executor);
let message = error.to_string();
if error.is_provider_failure() {
health.record_rpc(signer.chain_id, Err(&message));
}
warn!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
%error,
provider_failure = error.is_provider_failure(),
"transient signer recovery failure; retrying"
);
}
Err(error) if error.requires_quarantine() => {
if store.signer_lifecycle(&signer.signer_id, signer.chain_id).await? != SignerLifecycle::Quarantined {
store
.quarantine_signer(&signer.signer_id, signer.chain_id, &error.to_string())
.await?;
newton_metric::inc_signer_recoveries_total(signer.chain_id, "quarantined");
error!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
%error,
"signer safety invariant failed; quarantining nonce lane"
);
}
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), false);
}
Err(error) => return Err(error),
}
if sleep_or_cancel(config.poll_interval, cancellation).await {
return Ok(LaneRecovery::Stopped);
}
}
}
async fn reconcile(
store: &Store,
signer: &ManagedSigner,
config: &ExecutorConfig,
cancellation: CancellationToken,
) -> Result<(), ExecutorError> {
if store.signer_lifecycle(&signer.signer_id, signer.chain_id).await? == SignerLifecycle::Underfunded {
let required = store
.signer_required_balance(&signer.signer_id, signer.chain_id)
.await?
.ok_or(ExecutorError::MissingRequiredBalance)?;
let balance = signer
.backend
.balance()
.await
.map_err(ExecutorError::UnderfundedBalanceRead)?;
if balance < required {
return Err(ExecutorError::InsufficientSignerBalance { balance, required });
}
info!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
%balance,
%required,
"signer balance recovered"
);
store.mark_signer_ready(&signer.signer_id, signer.chain_id).await?;
}
let latest = signer.backend.latest_transaction_count().await?;
let pending = signer.backend.pending_transaction_count().await?;
newton_metric::set_signer_nonce_gap(signer.chain_id, signer.signer_id.as_str(), latest, pending);
let assignment = store.active_assignment(signer.address(), signer.chain_id).await?;
if assignment.is_some() {
store.restore_signer_busy(&signer.signer_id, signer.chain_id).await?;
}
match assignment {
None => {
if latest != pending {
return Err(ExecutorError::UnownedPendingTransaction {
signer: signer.backend.address(),
latest,
pending,
});
}
if latest > 0 && !store.signer_has_journal(signer.address(), signer.chain_id).await? {
info!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
nonce = latest,
"derived signer enrollment baseline from RPC"
);
}
store.mark_signer_ready(&signer.signer_id, signer.chain_id).await?;
}
Some(assignment) if assignment.attempts.is_empty() => {
if latest != pending {
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, assignment.job.job_id)
.await?;
return Err(ExecutorError::UnownedPendingTransaction {
signer: signer.backend.address(),
latest,
pending,
});
}
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, assignment.job.job_id)
.await?;
}
Some(assignment) => {
let span = submission_job_span(signer, &assignment.job);
async {
info!(
attempt_count = assignment.attempts.len(),
"resuming durable submission job"
);
recover_prepared_assignment(store, signer, config, assignment, latest, pending, cancellation).await
}
.instrument(span)
.await?;
}
}
info!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
latest,
"signer reconciled"
);
Ok(())
}
async fn worker_loop(
store: Arc<Store>,
signer: ManagedSigner,
config: ExecutorConfig,
cancellation: CancellationToken,
health: Arc<RuntimeHealth>,
) -> Result<(), ExecutorError> {
const MIN_IDLE_FALLBACK: Duration = Duration::from_secs(5);
let fallback = config.poll_interval.max(MIN_IDLE_FALLBACK);
loop {
health.heartbeat(RuntimeTask::Executor);
match store.reserve_job(&signer.record()).await {
Ok(Some(job)) => {
let job_id = job.job_id;
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), false);
let span = submission_job_span(&signer, &job);
let execution = async {
log_job_reserved(&job);
execute_job(&store, &signer, &config, job, cancellation.clone()).await
}
.instrument(span)
.await;
match execution {
Ok(()) => {
health.record_rpc(signer.chain_id, Ok(()));
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), true);
}
Err(error) if error.is_retryable() => {
record_backend_error(&health, signer.chain_id, &error);
warn!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
signer_role = signer.role.as_str(),
chain_id = signer.chain_id,
%job_id,
%error,
"retryable job execution failure"
);
if error.is_transient() && sleep_or_cancel(config.poll_interval, &cancellation).await {
return Ok(());
}
if recover_lane(&store, &signer, &config, &cancellation, &health).await?
== LaneRecovery::Stopped
{
return Ok(());
}
}
Err(error) if error.requires_quarantine() => {
record_backend_error(&health, signer.chain_id, &error);
error!(
signer_id = %signer.signer_id,
signer_address = %signer.address(),
signer_role = signer.role.as_str(),
chain_id = signer.chain_id,
%job_id,
%error,
"permanent job failure; quarantining only this signer lane"
);
if !error.assignment_already_finalized() {
match store
.release_safe_reservation(&signer.signer_id, signer.chain_id, job_id)
.await
{
Ok(()) | Err(StoreError::UnsafeRelease) => {}
Err(store_error) => return Err(store_error.into()),
}
}
store
.quarantine_signer(&signer.signer_id, signer.chain_id, &error.to_string())
.await?;
newton_metric::set_signer_ready(signer.chain_id, signer.signer_id.as_str(), false);
if recover_lane(&store, &signer, &config, &cancellation, &health).await?
== LaneRecovery::Stopped
{
return Ok(());
}
}
Err(error) => {
record_backend_error(&health, signer.chain_id, &error);
return Err(error);
}
}
}
Ok(None) => {
tokio::select! {
_ = cancellation.cancelled() => return Ok(()),
_ = store.wait_for_job(fallback) => {}
}
}
Err(error) => {
return Err(error.into());
}
}
}
}
fn submission_job_span(signer: &ManagedSigner, job: &JobRecord) -> Span {
let items = job.intent.task_items();
info_span!(
"submission_job",
job_id = %job.job_id,
chain_id = job.chain_id,
signer_id = %signer.signer_id,
signer_address = %signer.address(),
signer_role = job.signer_role.as_str(),
intent_kind = job.intent.kind(),
item_count = job.intent.item_count(),
contract_role = %job.contract_role,
submission_ids = %submission_ids_field(items),
task_ids = %task_ids_field(items),
)
}
fn log_job_reserved(job: &JobRecord) {
match &job.intent {
ExecutableIntent::BatchCreateAndRespond { items, .. } | ExecutableIntent::BatchRespond { items, .. } => {
info!(submission_count = items.len(), "task batch reserved for submission");
}
ExecutableIntent::ConfirmGlobalTableRoot {
root,
reference_timestamp,
reference_block_number,
..
} => {
info!(
%root,
reference_timestamp,
reference_block_number,
"transporter root confirmation reserved for submission"
);
}
ExecutableIntent::UpdateOperatorTable {
root,
reference_timestamp,
operator_set_index,
expected_leaf,
..
} => {
info!(
%root,
reference_timestamp,
operator_set_index,
%expected_leaf,
"transporter table update reserved for submission"
);
}
}
}
fn record_backend_error(health: &RuntimeHealth, chain_id: u64, error: &ExecutorError) {
if matches!(
error,
ExecutorError::Backend(_) | ExecutorError::UnderfundedBalanceRead(_)
) {
let message = error.to_string();
health.record_rpc(chain_id, Err(&message));
}
}
async fn execute_job(
store: &Store,
signer: &ManagedSigner,
config: &ExecutorConfig,
job: newton_submission_service::JobRecord,
cancellation: CancellationToken,
) -> Result<(), ExecutorError> {
let latest = match signer.backend.latest_transaction_count().await {
Ok(latest) => latest,
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 pending = match signer.backend.pending_transaction_count().await {
Ok(pending) => pending,
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()),
};
newton_metric::set_signer_nonce_gap(signer.chain_id, signer.signer_id.as_str(), latest, pending);
if latest != pending {
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, job.job_id)
.await?;
return Err(ExecutorError::UnownedPendingTransaction {
signer: signer.backend.address(),
latest,
pending,
});
}
let signed = match signer.backend.prepare(&job.intent, latest, None).await {
Ok(signed) => signed,
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 @ BackendError::Simulation(_)) => {
classify_preflight_failure(store, signer, &job, latest, error).await?;
return Ok(());
}
Err(error) => {
finalize_preflight_failure(store, &job, &error).await?;
return Err(ExecutorError::FinalizedPreflight(error));
}
};
if signed.nonce != latest {
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() => {
store
.release_safe_reservation(&signer.signer_id, signer.chain_id, job.job_id)
.await?;
return Err(error.into());
}
Err(error) => return Err(error.into()),
};
if balance < required {
store
.release_underfunded_reservation(&signer.signer_id, signer.chain_id, job.job_id, balance, required)
.await?;
return Err(ExecutorError::InsufficientSignerBalance { balance, required });
}
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: AttemptKind::Canonical,
replaces_attempt_number: None,
})
.await?;
info!(
attempt_number,
attempt_kind = "canonical",
nonce = signed.nonce,
tx_hash = %signed.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,
"transaction prepared and durably recorded"
);
let results = signer.backend.broadcast(&signed.raw_transaction).await;
let any_success = results.iter().any(|result| result.result.is_ok());
let accepted_provider = results
.iter()
.find(|result| result.result.is_ok())
.map(|result| result.provider.as_str());
let encoded_results = rmp_serde::to_vec(&results)?;
store
.record_broadcast(job.job_id, attempt_number, &encoded_results)
.await?;
let broadcast_at_ms = unix_time_ms();
if any_success {
info!(
attempt_number,
attempt_kind = "canonical",
nonce = signed.nonce,
tx_hash = %signed.transaction_hash,
accepted_provider = accepted_provider.unwrap_or("unknown"),
provider_attempts = results.len(),
"transaction broadcast accepted"
);
} else {
warn!(
job_id = %job.job_id,
tx_hash = %signed.transaction_hash,
"all providers rejected broadcast; exact bytes remain recoverable"
);
}
track_attempt(
store,
signer,
config,
&job,
vec![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: AttemptKind::Canonical,
state: AttemptState::Broadcast,
receipt_provider: None,
updated_at_ms: broadcast_at_ms,
}],
cancellation,
)
.await
}
async fn recover_prepared_assignment(
store: &Store,
signer: &ManagedSigner,
config: &ExecutorConfig,
mut assignment: ActiveAssignment,
latest: u64,
pending: u64,
cancellation: CancellationToken,
) -> Result<(), ExecutorError> {
let attempt = assignment
.attempts
.last()
.cloned()
.ok_or(ExecutorError::MissingAttempt)?;
if assignment
.attempts
.iter()
.any(|candidate| candidate.nonce != attempt.nonce)
{
return Err(ExecutorError::AttemptNonceMismatch);
}
if latest > attempt.nonce {
let mut known_receipt = false;
for candidate in assignment.attempts.iter().rev() {
match signer
.backend
.observe(
candidate.transaction_hash,
candidate.receipt_provider.as_deref(),
ObservationScope::Exhaustive,
)
.await?
{
TransactionObservation::Pending | TransactionObservation::PendingInconclusive { .. } => {}
TransactionObservation::Mined { .. } | TransactionObservation::Reverted { .. } => {
known_receipt = true;
break;
}
}
}
if known_receipt {
return track_attempt(
store,
signer,
config,
&assignment.job,
assignment.attempts,
cancellation,
)
.await;
}
warn!(
signer_id = %signer.signer_id,
job_id = %assignment.job.job_id,
nonce = attempt.nonce,
latest,
"chain consumed the persisted nonce without a known receipt; reconciling task effects"
);
return classify_and_project(store, signer, &assignment.job, None).await;
}
if pending < attempt.nonce {
return Err(ExecutorError::NonceJournalAhead {
journal: attempt.nonce,
pending,
});
}
if attempt.state == AttemptState::Prepared {
let results = signer.backend.broadcast(&attempt.raw_transaction).await;
store
.record_broadcast(
assignment.job.job_id,
attempt.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.attempt_number,
attempt_kind = attempt.kind.as_str(),
nonce = attempt.nonce,
tx_hash = %attempt.transaction_hash,
accepted_provider,
provider_attempts = results.len(),
"recovered prepared transaction broadcast accepted"
);
} else {
warn!(
attempt_number = attempt.attempt_number,
attempt_kind = attempt.kind.as_str(),
nonce = attempt.nonce,
tx_hash = %attempt.transaction_hash,
provider_attempts = results.len(),
"all providers rejected recovered transaction broadcast; exact bytes remain durable"
);
}
if let Some(latest_attempt) = assignment.attempts.last_mut() {
latest_attempt.state = AttemptState::Broadcast;
latest_attempt.updated_at_ms = unix_time_ms();
}
}
track_attempt(
store,
signer,
config,
&assignment.job,
assignment.attempts,
cancellation,
)
.await
}
async fn sleep_or_cancel(duration: Duration, cancellation: &CancellationToken) -> bool {
tokio::select! {
_ = cancellation.cancelled() => true,
_ = tokio::time::sleep(duration) => false,
}
}
fn unix_time_ms() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.try_into()
.unwrap_or(i64::MAX)
}
#[derive(Debug, thiserror::Error)]
pub enum ExecutorError {
#[error(transparent)]
Store(#[from] newton_submission_service::StoreError),
#[error(transparent)]
Backend(#[from] BackendError),
#[error(transparent)]
Encode(#[from] rmp_serde::encode::Error),
#[error("signer {signer} has unowned pending activity: latest={latest}, pending={pending}")]
UnownedPendingTransaction {
signer: Address,
latest: u64,
pending: u64,
},
#[error("prepared signer assignment has no durable transaction attempt")]
MissingAttempt,
#[error("durable attempts for one job use different nonces")]
AttemptNonceMismatch,
#[error("journaled nonce {journal} is ahead of pending transaction count {pending}")]
NonceJournalAhead {
journal: u64,
pending: u64,
},
#[error("on-chain effect classification did not cover every batch item")]
EffectCoverage,
#[error("signer native balance {balance} is below required transaction cost {required}")]
InsufficientSignerBalance {
balance: U256,
required: U256,
},
#[error("failed to recheck underfunded signer balance: {0}")]
UnderfundedBalanceRead(BackendError),
#[error("underfunded signer is missing its required balance threshold")]
MissingRequiredBalance,
#[error("signer worker exited unexpectedly")]
WorkerExited,
#[error("signer worker panicked: {0}")]
WorkerPanicked(String),
#[error("job finalized after permanent preflight failure: {0}")]
FinalizedPreflight(BackendError),
}
impl ExecutorError {
fn is_process_fatal(&self) -> bool {
matches!(self, Self::Store(_) | Self::WorkerExited | Self::WorkerPanicked(_))
}
fn is_transient(&self) -> bool {
match self {
Self::Backend(error) => error.is_transient(),
Self::Store(error) => error.is_transient(),
_ => false,
}
}
fn is_provider_failure(&self) -> bool {
matches!(self, Self::Backend(_))
}
fn waits_for_chain(&self) -> bool {
matches!(
self,
Self::UnownedPendingTransaction { .. } | Self::NonceJournalAhead { .. } | Self::EffectCoverage
)
}
fn is_underfunded(&self) -> bool {
matches!(
self,
Self::InsufficientSignerBalance { .. } | Self::UnderfundedBalanceRead(_)
)
}
fn is_retryable(&self) -> bool {
self.is_transient() || self.waits_for_chain() || self.is_underfunded()
}
fn requires_quarantine(&self) -> bool {
!self.is_process_fatal() && !self.is_retryable()
}
fn assignment_already_finalized(&self) -> bool {
matches!(self, Self::FinalizedPreflight(_))
}
}
#[cfg(test)]
mod tests;