Skip to main content

iota_core/
transaction_orchestrator.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5// Transaction Orchestrator is a Node component that utilizes Quorum Driver or
6// TransactionDriver (selected per request by the P-COOL protocol flag) to
7// submit transactions to validators for finality, and proactively executes
8// finalized transactions locally.
9
10use std::{
11    collections::{BTreeMap, HashMap, hash_map::Entry},
12    net::SocketAddr,
13    ops::Deref,
14    path::Path,
15    sync::Arc,
16    time::Duration,
17};
18
19use futures::{
20    FutureExt,
21    future::{Either, Future, select},
22};
23use iota_common::{debug_fatal, sync::notify_read::NotifyRead};
24use iota_config::NodeConfig;
25use iota_metrics::{
26    TX_TYPE_SHARED_OBJ_TX, TX_TYPE_SINGLE_WRITER_TX, add_server_timing,
27    spawn_logged_monitored_task, spawn_monitored_task,
28};
29use iota_sdk_types::{Transaction, TransactionDigest};
30use iota_storage::write_path_pending_tx_log::WritePathPendingTransactionLog;
31use iota_types::{
32    effects::TransactionEffectsAPI,
33    error::{IotaError, IotaResult},
34    iota_system_state::IotaSystemState,
35    messages_checkpoint::CheckpointSequenceNumber,
36    quorum_driver_types::{
37        EffectsFinalityInfo, ExecuteTransactionRequestType, ExecuteTransactionRequestV1,
38        ExecuteTransactionResponseV1, FinalizedEffects, GroupedErrors,
39        IsTransactionExecutedLocally, QuorumDriverEffectsQueueResult, QuorumDriverError,
40        QuorumDriverResponse, QuorumDriverResult,
41    },
42    transaction::{SenderSignedTransactionAPI, VerifiedTransaction},
43    transaction_driver_types::{
44        EffectsFinalityInfo as TdEffectsFinalityInfo, FinalizedEffects as TdFinalizedEffects,
45    },
46    transaction_executor::{SimulateTransactionResult, VmChecks},
47};
48use parking_lot::Mutex;
49use prometheus_filtered::{
50    Histogram, MetricLevel, Registry,
51    core::{AtomicI64, AtomicU64, GenericCounter, GenericGauge},
52    register_histogram_vec_with_registry, register_int_counter_vec_with_registry,
53    register_int_counter_with_registry, register_int_gauge_vec_with_registry,
54    register_int_gauge_with_registry,
55};
56use tokio::{
57    sync::{
58        broadcast::{Receiver, error::RecvError},
59        watch,
60    },
61    task::JoinHandle,
62    time::timeout,
63};
64use tracing::{Instrument, debug, error, info, instrument, trace_span, warn};
65
66use crate::{
67    authority::{AuthorityState, authority_per_epoch_store::AuthorityPerEpochStore},
68    authority_aggregator::AuthorityAggregator,
69    authority_client::{AuthorityAPI, NetworkAuthorityClient},
70    quorum_driver::{
71        QuorumDriverHandler, QuorumDriverHandlerBuilder, QuorumDriverMetrics,
72        reconfig_observer::{OnsiteReconfigObserver, ReconfigObserver},
73    },
74    transaction_driver::{
75        AggregatedRequestErrors, QuorumTransactionResponse, SubmitTransactionOptions,
76        TransactionDriver, TransactionDriverError, TransactionDriverMetrics,
77        reconfig_observer::OnsiteReconfigObserver as TdOnsiteReconfigObserver,
78    },
79    validator_client_monitor::ValidatorClientMetrics,
80};
81
82// How long to wait for local execution (including parents) before a timeout
83// is returned to client.
84const LOCAL_EXECUTION_TIMEOUT: Duration = Duration::from_secs(10);
85
86const WAIT_FOR_FINALITY_TIMEOUT: Duration = Duration::from_secs(30);
87
88/// The submission flow for a transaction, selected per request by the
89/// epoch's P-COOL flag.
90enum Driver<A: Clone> {
91    /// Certificate-based flow (P-COOL disabled).
92    Quorum(Arc<QuorumDriverHandler<A>>),
93    /// Direct-to-consensus P-COOL flow.
94    Transaction(Arc<TransactionDriver<A>>),
95}
96
97/// Transaction Orchestrator is a Node component that supports both QuorumDriver
98/// and TransactionDriver for submitting transactions to validators for
99/// finality. It adds inflight deduplication, waiting for local execution,
100/// recovery, and epoch change handling.
101///
102/// The epoch's P-COOL flag selects the flow serving a request. The
103/// TransactionDriver always exists, so it tracks epochs before the flag
104/// enables it. The QuorumDriver exists only when the node booted with
105/// P-COOL disabled and ran WAL recovery then; after a rollback a node
106/// booted under P-COOL must be restarted.
107pub struct TransactionOrchestrator<A: Clone> {
108    quorum_driver: Option<Arc<QuorumDriverHandler<A>>>,
109    transaction_driver: Arc<TransactionDriver<A>>,
110    validator_state: Arc<AuthorityState>,
111    /// Handle to the pending-tx-log cleanup loop; present only with the
112    /// quorum driver.
113    _local_executor_handle: Option<JoinHandle<()>>,
114    pending_tx_log: Arc<WritePathPendingTransactionLog>,
115    /// Digests currently being driven to finality by the TransactionDriver;
116    /// used to deduplicate concurrent submissions of the same transaction,
117    /// with a channel per digest through which the driving submission
118    /// publishes its outcome to concurrent duplicates. Kept in memory only:
119    /// the driver path is best-effort, so there is nothing to recover after
120    /// a restart. The QuorumDriver path tracks its submissions in
121    /// `pending_tx_log` instead.
122    in_flight_transactions: InFlightTransactions,
123    notifier: Arc<NotifyRead<TransactionDigest, QuorumDriverResult>>,
124    metrics: Arc<TransactionOrchestratorMetrics>,
125}
126
127impl TransactionOrchestrator<NetworkAuthorityClient> {
128    pub fn new_with_auth_aggregator(
129        validators: Arc<AuthorityAggregator<NetworkAuthorityClient>>,
130        validator_state: Arc<AuthorityState>,
131        reconfig_channel: Receiver<IotaSystemState>,
132        parent_path: &Path,
133        prometheus_registry: &Registry,
134        node_config: Option<&NodeConfig>,
135    ) -> Self {
136        let td_reconfig_observer = TdOnsiteReconfigObserver::new(
137            reconfig_channel.resubscribe(),
138            validator_state.get_object_cache_reader().clone(),
139            validator_state.clone_committee_store(),
140            validators.safe_client_metrics_base.clone(),
141        );
142
143        let qd_reconfig_observer = OnsiteReconfigObserver::new(
144            reconfig_channel.resubscribe(),
145            validator_state.get_object_cache_reader().clone(),
146            validator_state.clone_committee_store(),
147            validators.safe_client_metrics_base.clone(),
148            validators.metrics.deref().clone(),
149        );
150
151        TransactionOrchestrator::new(
152            validators,
153            validator_state,
154            parent_path,
155            prometheus_registry,
156            qd_reconfig_observer,
157            td_reconfig_observer,
158            node_config,
159        )
160    }
161}
162
163impl<A> TransactionOrchestrator<A>
164where
165    A: AuthorityAPI + Send + Sync + 'static + Clone,
166    OnsiteReconfigObserver: ReconfigObserver<A>,
167    TdOnsiteReconfigObserver: crate::transaction_driver::reconfig_observer::ReconfigObserver<A>,
168{
169    pub fn new(
170        validators: Arc<AuthorityAggregator<A>>,
171        validator_state: Arc<AuthorityState>,
172        parent_path: &Path,
173        prometheus_registry: &Registry,
174        reconfig_observer: OnsiteReconfigObserver,
175        td_reconfig_observer: TdOnsiteReconfigObserver,
176        node_config: Option<&NodeConfig>,
177    ) -> Self {
178        let epoch_store = validator_state.load_epoch_store_one_call_per_task();
179        let use_transaction_driver = epoch_store.protocol_config().enable_pcool_flow();
180
181        let notifier = Arc::new(NotifyRead::new());
182        let metrics = Arc::new(TransactionOrchestratorMetrics::new(prometheus_registry));
183        let pending_tx_log = Arc::new(WritePathPendingTransactionLog::new(
184            parent_path.join("fullnode_pending_transactions"),
185        ));
186
187        // Registered for both flows even when the quorum driver is not
188        // built, so metric presence does not depend on the boot mode.
189        let quorum_driver_metrics = Arc::new(QuorumDriverMetrics::new(prometheus_registry));
190        let transaction_driver_metrics =
191            Arc::new(TransactionDriverMetrics::new(prometheus_registry));
192        let client_metrics = Arc::new(ValidatorClientMetrics::new(prometheus_registry));
193
194        let (quorum_driver, _local_executor_handle) = if use_transaction_driver {
195            (None, None)
196        } else {
197            let quorum_driver = Arc::new(
198                QuorumDriverHandlerBuilder::new(validators.clone(), quorum_driver_metrics)
199                    .with_notifier(notifier.clone())
200                    .with_reconfig_observer(Arc::new(reconfig_observer))
201                    .start(),
202            );
203            // The cleanup loop must exist before WAL recovery runs, so a
204            // recovered transaction cannot complete before its receiver
205            // exists.
206            let effects_receiver = quorum_driver.subscribe_to_effects();
207            let pending_tx_log_clone = pending_tx_log.clone();
208            let local_executor_handle = spawn_monitored_task!(async move {
209                Self::loop_pending_transaction_log(effects_receiver, pending_tx_log_clone).await;
210            });
211            Self::schedule_txes_in_log(pending_tx_log.clone(), quorum_driver.clone());
212            (Some(quorum_driver), Some(local_executor_handle))
213        };
214
215        // `Weak` so detached driver tasks cannot pin the authority state.
216        let pcool_flow_enabled: Arc<dyn Fn() -> bool + Send + Sync> = {
217            let validator_state = Arc::downgrade(&validator_state);
218            Arc::new(move || {
219                validator_state.upgrade().is_some_and(|state| {
220                    state
221                        .load_epoch_store_one_call_per_task()
222                        .protocol_config()
223                        .enable_pcool_flow()
224                })
225            })
226        };
227        let transaction_driver = TransactionDriver::new(
228            validators,
229            Arc::new(td_reconfig_observer),
230            transaction_driver_metrics,
231            node_config.and_then(|config| config.validator_client_monitor_config.clone()),
232            client_metrics,
233            pcool_flow_enabled,
234        );
235
236        Self {
237            quorum_driver,
238            transaction_driver,
239            validator_state,
240            _local_executor_handle,
241            pending_tx_log,
242            in_flight_transactions: Default::default(),
243            notifier,
244            metrics,
245        }
246    }
247}
248
249impl<A> TransactionOrchestrator<A>
250where
251    A: AuthorityAPI + Send + Sync + 'static + Clone,
252{
253    /// Returns the flow selected by `epoch_store`'s P-COOL flag. Call with
254    /// the request's own epoch store snapshot so the flag and the submission
255    /// see the same epoch. Errors when the quorum driver is selected on a
256    /// node that booted under P-COOL: such a node never ran WAL recovery and
257    /// must be restarted.
258    fn select_driver(
259        &self,
260        epoch_store: &AuthorityPerEpochStore,
261    ) -> Result<Driver<A>, QuorumDriverError> {
262        if epoch_store.protocol_config().enable_pcool_flow() {
263            return Ok(Driver::Transaction(self.transaction_driver.clone()));
264        }
265        self.quorum_driver
266            .clone()
267            .map(Driver::Quorum)
268            .ok_or_else(|| {
269                error!(
270                    "This fullnode started while P-COOL was enabled and must be restarted to \
271                     serve the certificate-based flow"
272                );
273                QuorumDriverError::QuorumDriverInternal(IotaError::UnsupportedFeature {
274                    error: "this fullnode started while P-COOL was enabled and must be \
275                            restarted to serve the certificate-based flow"
276                        .to_string(),
277                })
278            })
279    }
280
281    #[instrument(name = "tx_orchestrator_execute_transaction_block", level = "trace", skip_all,
282        fields(
283        tx_digest = ?request.transaction.digest(),
284        tx_type = ?request_type,
285        ),
286        err)]
287    pub async fn execute_transaction_block(
288        &self,
289        request: ExecuteTransactionRequestV1,
290        request_type: ExecuteTransactionRequestType,
291        client_addr: Option<SocketAddr>,
292    ) -> Result<(ExecuteTransactionResponseV1, IsTransactionExecutedLocally), QuorumDriverError>
293    {
294        let epoch_store = self.validator_state.load_epoch_store_one_call_per_task();
295
296        let transaction = epoch_store
297            .verify_transaction(request.transaction.clone())
298            .map_err(QuorumDriverError::InvalidUserSignature)?;
299
300        // Captured before `request` moves so the skip-cert reconcile reads
301        // caller intent, not whatever the submitter happened to return — a
302        // Byzantine submitter could otherwise censor a field by returning
303        // `None`.
304        let include_events = request.include_events;
305        let include_input_objects = request.include_input_objects;
306        let include_output_objects = request.include_output_objects;
307
308        let tx_digest = *transaction.digest();
309
310        // A resubmission of an already-executed transaction is answered from
311        // the local cache instead of being driven through the validators
312        // again.
313        if let Some(response) = Self::build_response_from_local_effects(
314            &self.validator_state,
315            &tx_digest,
316            include_events,
317            include_input_objects,
318            include_output_objects,
319        )? {
320            self.metrics.early_cached_response.inc();
321            debug!(
322                ?tx_digest,
323                "Returning cached results for already-executed transaction"
324            );
325            return Ok((response, true));
326        }
327
328        // Reject malformed transactions before either driver inspects shared
329        // inputs or `MoveAuthenticator`. Runs after the cache lookup so that,
330        // as on the upstream flow, a resubmission of an executed transaction
331        // gets its cached results even if it no longer passes the current
332        // epoch's checks (e.g. its expiration epoch has passed).
333        transaction
334            .validity_check(&epoch_store.tx_validity_check_context())
335            .map_err(QuorumDriverError::InvalidTransaction)?;
336
337        let wait_for_local_execution = matches!(
338            request_type,
339            ExecuteTransactionRequestType::WaitForLocalExecution
340        );
341        let (mut response, seq) =
342            match (self.select_driver(&epoch_store)?, wait_for_local_execution) {
343                (Driver::Transaction(td), true) => {
344                    let in_flight_transactions = self.in_flight_transactions.clone();
345                    let validator_state = self.validator_state.clone();
346                    let metrics = self.metrics.clone();
347                    // Detached so a client disconnect (this future dropped) does
348                    // not cancel a submission that may already be in consensus;
349                    // the task drives the transaction to finality on its own.
350                    join_submission_task(spawn_monitored_task!(Self::submit_with_checkpoint_race(
351                        td,
352                        in_flight_transactions,
353                        validator_state,
354                        metrics,
355                        request,
356                        client_addr,
357                        tx_digest,
358                    )))
359                    .await?
360                }
361                (Driver::Transaction(td), false) => {
362                    let in_flight_transactions = self.in_flight_transactions.clone();
363                    let validator_state = self.validator_state.clone();
364                    // Detached for the same reason as above.
365                    let result = join_submission_task(spawn_monitored_task!(
366                        Self::submit_with_transaction_driver(
367                            td,
368                            in_flight_transactions,
369                            validator_state,
370                            request,
371                            client_addr,
372                            false,
373                        )
374                    ))
375                    .await?;
376                    (Some(result), None)
377                }
378                (Driver::Quorum(qd), _) => {
379                    let qd_resp = self
380                        .execute_transaction_impl(
381                            &qd,
382                            &epoch_store,
383                            request,
384                            transaction.clone(),
385                            client_addr,
386                        )
387                        .await?;
388                    (Some(quorum_driver_response_to_v1(qd_resp)), None)
389                }
390            };
391
392        // `needs_cache_rebuild` is derived from finality, not caller intent:
393        // the QD fallback path returns `Certified` and a duplicate
394        // submission inheriting the outcome of an in-flight certifying
395        // submission returns `QuorumExecuted` — neither needs a rebuild —
396        // even when the caller asked for `WaitForLocalExecution`, while only
397        // the TD skip-cert engine produces `UncertifiedSingleValidator`. The
398        // checkpoint sequence comes from `submit_with_checkpoint_race`, which
399        // relies on `executed_transactions_to_checkpoint` being written
400        // strictly after every tx's effects — so a `Some(seq)` here implies
401        // the cache has authoritative effects.
402        let needs_cache_rebuild = matches!(
403            response.as_ref().map(|r| &r.effects.finality_info),
404            None | Some(EffectsFinalityInfo::UncertifiedSingleValidator(_)),
405        );
406
407        let executed_locally = if !wait_for_local_execution {
408            false
409        } else if needs_cache_rebuild {
410            let Some(seq) = seq else {
411                // Timed out waiting for the tx to land in a local checkpoint.
412                // In this branch `response` is either `None` (recovery) or
413                // `UncertifiedSingleValidator` (TD skip-cert) — both must
414                // surface as `TimeoutBeforeFinality` rather than leaking
415                // uncorroborated single-validator effects to the client.
416                return Err(QuorumDriverError::TimeoutBeforeFinality);
417            };
418            match response.as_mut() {
419                Some(existing) => Self::reconcile_effects_from_cache(
420                    &self.validator_state,
421                    tx_digest,
422                    seq,
423                    include_events,
424                    include_input_objects,
425                    include_output_objects,
426                    existing,
427                    &self.metrics,
428                )?,
429                None => {
430                    response = Some(Self::build_response_from_cache(
431                        &self.validator_state,
432                        tx_digest,
433                        seq,
434                        include_events,
435                        include_input_objects,
436                        include_output_objects,
437                    )?);
438                }
439            }
440            true
441        } else {
442            // The response is already 2f+1 certified — from the QuorumDriver,
443            // or inherited by a duplicate submission from an in-flight
444            // certifying submission — so just confirm local execution
445            // finished.
446            let ok = Self::wait_for_finalized_tx_executed_locally_with_timeout(
447                &self.validator_state,
448                &transaction,
449                &self.metrics,
450            )
451            .await
452            .is_ok();
453            add_server_timing("local_execution");
454            ok
455        };
456
457        let response = response.expect("response must be populated before return");
458
459        // Safety guard: `UncertifiedSingleValidator` finality carries effects
460        // from the single submitting validator only — they MUST NOT reach the
461        // client without first being corroborated against the local cache. The
462        // reachable paths today all either upgrade finality via
463        // `reconcile_effects_from_cache` / `build_response_from_cache`, or
464        // branch to `TimeoutBeforeFinality`; this guard is the last-chance
465        // fallback for a future refactor that forgets to reconcile. Do not
466        // remove as dead code.
467        if matches!(
468            response.effects.finality_info,
469            EffectsFinalityInfo::UncertifiedSingleValidator(_)
470        ) {
471            debug_fatal!(
472                "Uncertified effects (UncertifiedSingleValidator) about to be returned \
473                 to the client for tx {:?}",
474                response.effects.effects.transaction_digest()
475            );
476            return Err(QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(
477                "internal error: transaction effects not finalized".to_string(),
478            )));
479        }
480
481        Ok((response, executed_locally))
482    }
483
484    /// Replace the response's effects, events, and input/output objects with
485    /// the authoritative copies derived from the local cache — the local
486    /// checkpoint executor has processed the tx, so the cache has the real
487    /// data and the TD-returned (single-validator) copies can be discarded.
488    ///
489    /// `tx_digest` must be the digest of the caller's original transaction,
490    /// not the digest carried in `response.effects.effects` — a byzantine
491    /// submitter could set the latter to an unrelated (already-executed) tx
492    /// so we'd read unrelated effects from the cache.
493    ///
494    /// The caller must have obtained `checkpoint_seq` from
495    /// `wait_for_checkpoint_inclusion` (not just `get_transaction_checkpoint`),
496    /// because that function guarantees both the effects write and the
497    /// checkpoint-mapping write have landed — it's the only way to avoid the
498    /// race between `notify_read_executed_effects_digests` (fires per-tx) and
499    /// `insert_finalized_transactions` (fires per-checkpoint, after the
500    /// `CheckpointExecutor` has awaited every tx in that checkpoint).
501    ///
502    /// Upgrades the finality info to `Checkpointed(epoch, checkpoint_seq)`. A
503    /// warning is logged if the TD-returned effects digest diverges from the
504    /// cache digest, or the submitter claimed events the cache doesn't have
505    /// (byzantine submitter or bug).
506    fn reconcile_effects_from_cache(
507        validator_state: &Arc<AuthorityState>,
508        tx_digest: TransactionDigest,
509        checkpoint_seq: CheckpointSequenceNumber,
510        include_events: bool,
511        include_input_objects: bool,
512        include_output_objects: bool,
513        response: &mut ExecuteTransactionResponseV1,
514        metrics: &TransactionOrchestratorMetrics,
515    ) -> Result<(), QuorumDriverError> {
516        let rebuilt = Self::build_response_from_cache(
517            validator_state,
518            tx_digest,
519            checkpoint_seq,
520            include_events,
521            include_input_objects,
522            include_output_objects,
523        )?;
524
525        let td_digest = response.effects.effects.digest();
526        let cache_digest = rebuilt.effects.effects.digest();
527        if td_digest != cache_digest {
528            warn!(
529                ?tx_digest,
530                ?td_digest,
531                ?cache_digest,
532                "reconcile_effects_from_cache: TransactionDriver and local cache disagree \
533                 on effects digest — replacing with cache (possible byzantine submitter)"
534            );
535        }
536        if include_events && response.events.is_some() && rebuilt.events.is_none() {
537            warn!(
538                ?tx_digest,
539                "reconcile_effects_from_cache: submitter claimed events but cache has \
540                 none — discarding (possible byzantine submitter)"
541            );
542            metrics.skip_effect_cert_events_cache_miss.inc();
543        }
544        *response = rebuilt;
545        Ok(())
546    }
547
548    /// Build a skip-effect-certification response entirely from the local
549    /// cache. The caller must have already obtained `checkpoint_seq` via
550    /// `wait_for_checkpoint_inclusion`, which is supposed to guarantee both
551    /// the effects write and the checkpoint-mapping write have landed. A
552    /// missing cache entry here would mean a transient races we observed in
553    /// practice; mapped to `TimeoutBeforeFinality` so the client retries
554    /// rather than seeing a misleading `QuorumDriverInternal`.
555    fn build_response_from_cache(
556        validator_state: &Arc<AuthorityState>,
557        tx_digest: TransactionDigest,
558        checkpoint_seq: CheckpointSequenceNumber,
559        include_events: bool,
560        include_input_objects: bool,
561        include_output_objects: bool,
562    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
563        let cached = read_cached_transaction_data(
564            validator_state,
565            &tx_digest,
566            include_events,
567            include_input_objects,
568            include_output_objects,
569        )
570        .map_err(|e| {
571            QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(format!(
572                "failed to read cached tx data for {tx_digest:?}: {e:?}"
573            )))
574        })?
575        .ok_or_else(|| {
576            // Checkpoint inclusion is supposed to guarantee the cache has
577            // effects, but we've seen transient misses; surface as a retriable
578            // timeout rather than an internal error.
579            warn!(
580                ?tx_digest,
581                "effects missing from cache after checkpoint inclusion — surfacing as \
582                 TimeoutBeforeFinality"
583            );
584            QuorumDriverError::TimeoutBeforeFinality
585        })?;
586        let iota_types::transaction_executor::CachedTransactionData {
587            effects,
588            events,
589            input_objects,
590            output_objects,
591        } = cached;
592
593        let epoch = effects.epoch();
594        Ok(ExecuteTransactionResponseV1 {
595            effects: FinalizedEffects {
596                effects,
597                finality_info: EffectsFinalityInfo::Checkpointed(epoch, checkpoint_seq),
598            },
599            events,
600            input_objects,
601            output_objects,
602            auxiliary_data: None,
603        })
604    }
605
606    /// Build a response from the local cache for a transaction that has
607    /// already been executed on this node. Returns `Ok(None)` when the
608    /// transaction has not been executed locally. Unlike
609    /// `build_response_from_cache`, no checkpoint sequence is required: local
610    /// effects only exist for finalized transactions, so the response is
611    /// tagged `QuorumExecuted`.
612    fn build_response_from_local_effects(
613        validator_state: &Arc<AuthorityState>,
614        tx_digest: &TransactionDigest,
615        include_events: bool,
616        include_input_objects: bool,
617        include_output_objects: bool,
618    ) -> Result<Option<ExecuteTransactionResponseV1>, QuorumDriverError> {
619        let Some(cached) = read_cached_transaction_data(
620            validator_state,
621            tx_digest,
622            include_events,
623            include_input_objects,
624            include_output_objects,
625        )
626        .map_err(QuorumDriverError::QuorumDriverInternal)?
627        else {
628            return Ok(None);
629        };
630        let iota_types::transaction_executor::CachedTransactionData {
631            effects,
632            events,
633            input_objects,
634            output_objects,
635        } = cached;
636
637        let epoch = effects.epoch();
638        Ok(Some(ExecuteTransactionResponseV1 {
639            effects: FinalizedEffects {
640                effects,
641                finality_info: EffectsFinalityInfo::QuorumExecuted(epoch),
642            },
643            events,
644            input_objects,
645            output_objects,
646            auxiliary_data: None,
647        }))
648    }
649
650    // Utilize the handle_certificate_v1 validator api to request input/output
651    // objects
652    #[instrument(name = "tx_orchestrator_execute_transaction_v1", level = "trace", skip_all,
653        fields(tx_digest = ?request.transaction.digest()))]
654    pub async fn execute_transaction_v1(
655        &self,
656        request: ExecuteTransactionRequestV1,
657        skip_certification: bool,
658        client_addr: Option<SocketAddr>,
659    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
660        let epoch_store = self.validator_state.load_epoch_store_one_call_per_task();
661
662        let transaction = epoch_store
663            .verify_transaction(request.transaction.clone())
664            .map_err(QuorumDriverError::InvalidUserSignature)?;
665        let tx_digest = *transaction.digest();
666
667        // A resubmission of an already-executed transaction is answered from
668        // the local cache instead of being driven through the validators
669        // again.
670        if let Some(response) = Self::build_response_from_local_effects(
671            &self.validator_state,
672            &tx_digest,
673            request.include_events,
674            request.include_input_objects,
675            request.include_output_objects,
676        )? {
677            self.metrics.early_cached_response.inc();
678            debug!(
679                ?tx_digest,
680                "Returning cached results for already-executed transaction"
681            );
682            return Ok(response);
683        }
684
685        // Reject malformed transactions before either driver inspects shared
686        // inputs or `MoveAuthenticator`. Runs after the cache lookup so that,
687        // as on the upstream flow, a resubmission of an executed transaction
688        // gets its cached results even if it no longer passes the current
689        // epoch's checks (e.g. its expiration epoch has passed).
690        transaction
691            .validity_check(&epoch_store.tx_validity_check_context())
692            .map_err(QuorumDriverError::InvalidTransaction)?;
693
694        match self.select_driver(&epoch_store)? {
695            Driver::Transaction(td) => {
696                let in_flight_transactions = self.in_flight_transactions.clone();
697                let validator_state = self.validator_state.clone();
698                // v1 does not do an internal wait; callers (e.g. the gRPC
699                // execution service) are responsible for their own
700                // `wait_for_checkpoint_inclusion` when they need it, and will
701                // reconcile the response from the cache there.
702                //
703                // Detached so a client disconnect does not cancel a submission
704                // that may already be in consensus.
705                join_submission_task(spawn_monitored_task!(Self::submit_with_transaction_driver(
706                    td,
707                    in_flight_transactions,
708                    validator_state,
709                    request,
710                    client_addr,
711                    skip_certification,
712                )))
713                .await
714            }
715            Driver::Quorum(qd) => {
716                let qd_resp = self
717                    .execute_transaction_impl(&qd, &epoch_store, request, transaction, client_addr)
718                    .await?;
719                Ok(quorum_driver_response_to_v1(qd_resp))
720            }
721        }
722    }
723
724    /// Submit on the skip-effect-certification path while concurrently
725    /// waiting for local checkpoint inclusion. The race is asymmetric:
726    ///
727    /// - If the **checkpoint** future resolves first (slow driver, e.g. stuck corroborating a
728    ///   Byzantine validator's rejection), the driver future is dropped and the caller rebuilds the
729    ///   response from the local cache.
730    /// - If the **driver** returns first, its result is taken and the checkpoint future is awaited
731    ///   to completion (up to the shared `WAIT_FOR_FINALITY_TIMEOUT`) before returning, so the
732    ///   caller has a checkpoint sequence to reconcile against.
733    ///
734    /// Returns `(response, seq)` where `response` is `Some` when the driver
735    /// returned a result (which may carry `UncertifiedSingleValidator`
736    /// finality requiring rebuild) and `seq` is the checkpoint sequence if
737    /// either future yielded it.
738    ///
739    /// Run inside a detached task so a client disconnect cannot cancel the
740    /// race before the checkpoint-sequence bookkeeping completes.
741    #[instrument(name = "tx_orchestrator_submit_with_checkpoint_race", level = "trace", skip_all,
742        fields(tx_digest = ?tx_digest))]
743    async fn submit_with_checkpoint_race(
744        td: Arc<TransactionDriver<A>>,
745        in_flight_transactions: InFlightTransactions,
746        validator_state: Arc<AuthorityState>,
747        metrics: Arc<TransactionOrchestratorMetrics>,
748        request: ExecuteTransactionRequestV1,
749        client_addr: Option<SocketAddr>,
750        tx_digest: TransactionDigest,
751    ) -> Result<
752        (
753            Option<ExecuteTransactionResponseV1>,
754            Option<CheckpointSequenceNumber>,
755        ),
756        QuorumDriverError,
757    > {
758        let digests = [tx_digest];
759        let checkpoint_inclusion =
760            validator_state.wait_for_checkpoint_inclusion(&digests, WAIT_FOR_FINALITY_TIMEOUT);
761        tokio::pin!(checkpoint_inclusion);
762        let driver = Self::submit_with_transaction_driver(
763            td,
764            in_flight_transactions,
765            validator_state.clone(),
766            request,
767            client_addr,
768            true,
769        );
770
771        let seq_for_tx = |inclusion_map: BTreeMap<_, (CheckpointSequenceNumber, _)>| {
772            inclusion_map.get(&tx_digest).map(|&(seq, _)| seq)
773        };
774
775        let result = tokio::select! {
776            biased;
777            // `SubmittedButFetchFailed` is retriable (`ErrorCategory::Unavailable`)
778            // so the driver's outer loop reissues submission internally and
779            // only returns here as `Ok`, `TimeoutWithLastRetriableError`, or
780            // a non-retriable error like `RejectedByValidators`.
781            driver_result = driver => {
782                let response = Some(driver_result?);
783                let seq = (&mut checkpoint_inclusion).await.ok().and_then(seq_for_tx);
784                (response, seq)
785            }
786            checkpoint_result = &mut checkpoint_inclusion => {
787                metrics.skip_effect_cert_checkpoint_overrode_driver.inc();
788                // Dropping the cancelled driver closes the in-flight outcome
789                // channel; duplicate submissions fall back to waiting for
790                // checkpoint inclusion, which this race winning guarantees
791                // resolves immediately.
792                let seq = checkpoint_result.ok().and_then(seq_for_tx);
793                (None, seq)
794            }
795        };
796        add_server_timing("local_execution");
797        Ok(result)
798    }
799
800    /// Submit a transaction via the TransactionDriver (P-COOL flow).
801    ///
802    /// With `skip_certification = true` the driver may return
803    /// `UncertifiedSingleValidator` effects without a 2f+1 broadcast. The
804    /// caller (gRPC `execute_transactions` or `execute_transaction_block`)
805    /// is then responsible for `wait_for_checkpoint_inclusion` and the
806    /// cache-rebuild gate that replaces those single-validator effects with
807    /// authoritative data — uncertified data must never reach the client.
808    /// See `corroborate_single_validator_error` for the per-submission
809    /// fetch-failure recovery flow inside the driver.
810    ///
811    /// Run inside a detached task so a client disconnect cannot cancel a
812    /// `drive_transaction` call that may already be in consensus.
813    #[instrument(name = "tx_orchestrator_submit_with_td", level = "trace", skip_all,
814        fields(tx_digest = ?request.transaction.digest()))]
815    async fn submit_with_transaction_driver(
816        td: Arc<TransactionDriver<A>>,
817        in_flight_transactions: InFlightTransactions,
818        validator_state: Arc<AuthorityState>,
819        request: ExecuteTransactionRequestV1,
820        client_addr: Option<SocketAddr>,
821        skip_certification: bool,
822    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
823        let tx_digest = *request.transaction.digest();
824
825        // Deduplicate concurrent submissions of the same digest: only the
826        // first caller drives the committee-wide submission and publishes its
827        // outcome; the rest await that outcome. The guard removes the digest
828        // from the in-flight map on every exit path (success, error, timeout,
829        // or cancellation) when it is dropped.
830        let guard = match TransactionSubmissionGuard::acquire(in_flight_transactions, tx_digest) {
831            TransactionSubmission::Driving(guard) => guard,
832            TransactionSubmission::AlreadyInFlight(receiver) => {
833                debug!(
834                    ?tx_digest,
835                    "transaction already in flight; awaiting its outcome instead of driving a \
836                     duplicate submission"
837                );
838                return Self::await_in_flight_transaction(
839                    receiver,
840                    &td,
841                    &validator_state,
842                    tx_digest,
843                    &request,
844                    client_addr,
845                    skip_certification,
846                )
847                .await;
848            }
849        };
850
851        // This call runs inside a task detached from the caller, so the
852        // outcome is logged here rather than left to the caller — a
853        // disconnected client's continuation never runs and would
854        // otherwise never observe it.
855        let td_response = match td
856            .drive_transaction(
857                Some(request.transaction.clone()),
858                SubmitTransactionOptions {
859                    forwarded_client_addr: client_addr,
860                    ..Default::default()
861                },
862                Some(WAIT_FOR_FINALITY_TIMEOUT),
863                skip_certification,
864            )
865            .await
866        {
867            Ok(response) => response,
868            Err(e) => {
869                warn!(?tx_digest, "TransactionDriver submission failed: {e}");
870                let error = map_td_error_to_qd(e);
871                guard.publish(Err(error.clone()));
872                return Err(error);
873            }
874        };
875
876        debug!(?tx_digest, "TransactionDriver submission succeeded");
877
878        let td_response = Arc::new(td_response);
879        guard.publish(Ok(td_response.clone()));
880        // Dropping the guard closes the channel, releasing its copy of the
881        // response unless a duplicate submission still holds a receiver — in
882        // the common no-duplicate case the response is then moved into the
883        // reply instead of cloned.
884        drop(guard);
885        let td_response = Arc::try_unwrap(td_response).unwrap_or_else(|shared| (*shared).clone());
886
887        Ok(Self::response_from_driver_response(td_response, &request))
888    }
889
890    /// Build a caller-specific response from a driver response, honoring the
891    /// caller's include flags.
892    fn response_from_driver_response(
893        td_response: QuorumTransactionResponse,
894        request: &ExecuteTransactionRequestV1,
895    ) -> ExecuteTransactionResponseV1 {
896        let QuorumTransactionResponse {
897            effects,
898            events,
899            input_objects,
900            output_objects,
901            auxiliary_data,
902        } = td_response;
903        ExecuteTransactionResponseV1 {
904            effects: convert_td_to_qd_effects(effects),
905            events: request.include_events.then_some(events).flatten(),
906            input_objects: request
907                .include_input_objects
908                .then_some(input_objects)
909                .flatten(),
910            output_objects: request
911                .include_output_objects
912                .then_some(output_objects)
913                .flatten(),
914            auxiliary_data: request
915                .include_auxiliary_data
916                .then_some(auxiliary_data)
917                .flatten(),
918        }
919    }
920
921    /// Await the outcome of an already in-flight submission of `tx_digest`
922    /// instead of starting a second committee-wide submission for the same
923    /// transaction. Resolves to that submission's outcome — running the
924    /// effects-certification step first if the outcome does not satisfy this
925    /// caller — falls back to waiting for checkpoint inclusion if the
926    /// driving submission went away without publishing one (checkpoint-race
927    /// cancellation, panic, or shutdown), and returns
928    /// `TimeoutBeforeFinality` if nothing is published within
929    /// `WAIT_FOR_FINALITY_TIMEOUT` or the follow-up effects certification
930    /// does not complete within another `WAIT_FOR_FINALITY_TIMEOUT`.
931    async fn await_in_flight_transaction(
932        mut receiver: watch::Receiver<Option<InFlightSubmissionResult>>,
933        td: &Arc<TransactionDriver<A>>,
934        validator_state: &Arc<AuthorityState>,
935        tx_digest: TransactionDigest,
936        request: &ExecuteTransactionRequestV1,
937        client_addr: Option<SocketAddr>,
938        skip_certification: bool,
939    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
940        // The `Ref` returned by `wait_for` is a read guard and must not be
941        // held across an await, so the outcome is cloned out before
942        // branching. `wait_for` only returns a value matching its predicate,
943        // so `Some` is guaranteed on success; a closed channel yields `None`.
944        let published = tokio::time::timeout(
945            WAIT_FOR_FINALITY_TIMEOUT,
946            receiver.wait_for(|outcome| outcome.is_some()),
947        )
948        .await
949        .map_err(|_elapsed| QuorumDriverError::TimeoutBeforeFinality)?
950        .ok()
951        .and_then(|outcome_ref| outcome_ref.clone());
952
953        let Some(outcome) = published else {
954            // Channel closed without an outcome: the driving submission went
955            // away without publishing — routinely because its checkpoint
956            // race observed the transaction in a local checkpoint and
957            // cancelled it, exceptionally on panic or shutdown. Checkpoint
958            // inclusion is the remaining signal of the outcome.
959            return Self::response_from_checkpoint_inclusion(validator_state, tx_digest, request)
960                .await;
961        };
962        let td_response = outcome?;
963
964        let uncertified = matches!(
965            td_response.effects.finality_info,
966            TdEffectsFinalityInfo::UncertifiedSingleValidator(_)
967        );
968        if uncertified && !skip_certification {
969            // The in-flight submission already drove the transaction into
970            // consensus; only the 2f+1 effects certification is missing for
971            // this caller. Certify the effects directly instead of starting
972            // a second committee-wide submission or waiting for a checkpoint
973            // inclusion the caller never asked for. `certify_transaction` is
974            // internally bounded only by committee size times its per-request
975            // timeout, so cap it to the same client-facing budget the driving
976            // submission gets for its whole `drive_transaction` call.
977            let certified = tokio::time::timeout(
978                WAIT_FOR_FINALITY_TIMEOUT,
979                td.certify_transaction(
980                    tx_digest,
981                    SubmitTransactionOptions {
982                        forwarded_client_addr: client_addr,
983                        ..Default::default()
984                    },
985                ),
986            )
987            .await
988            .map_err(|_elapsed| QuorumDriverError::TimeoutBeforeFinality)?
989            .map_err(map_td_error_to_qd)?;
990            return Ok(Self::response_from_driver_response(certified, request));
991        }
992
993        Ok(Self::response_from_driver_response(
994            (*td_response).clone(),
995            request,
996        ))
997    }
998
999    /// Wait for `tx_digest` to reach a local checkpoint and build the
1000    /// response from the authoritative cache. The fallback outcome signal
1001    /// for a duplicate whose driving submission went away without
1002    /// publishing; the result carries `Checkpointed` finality, so it
1003    /// satisfies every caller. Times out with `TimeoutBeforeFinality` if
1004    /// the transaction does not get checkpointed in time.
1005    async fn response_from_checkpoint_inclusion(
1006        validator_state: &Arc<AuthorityState>,
1007        tx_digest: TransactionDigest,
1008        request: &ExecuteTransactionRequestV1,
1009    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
1010        let digests = [tx_digest];
1011        // The caller has typically already waited a full timeout on the
1012        // outcome channel, but this wait is still required: it is what
1013        // yields the checkpoint sequence and guarantees the checkpoint
1014        // mapping write has landed (see `reconcile_effects_from_cache`).
1015        // When the transaction is already checkpointed — the routine reason
1016        // the fallback fires — it resolves immediately; only after a
1017        // driving-task death does it actually wait, as the last remaining
1018        // signal of the outcome.
1019        let seq = validator_state
1020            .wait_for_checkpoint_inclusion(&digests, WAIT_FOR_FINALITY_TIMEOUT)
1021            .await
1022            .ok()
1023            .and_then(|inclusion| inclusion.get(&tx_digest).map(|&(seq, _)| seq))
1024            .ok_or(QuorumDriverError::TimeoutBeforeFinality)?;
1025        Self::build_response_from_cache(
1026            validator_state,
1027            tx_digest,
1028            seq,
1029            request.include_events,
1030            request.include_input_objects,
1031            request.include_output_objects,
1032        )
1033    }
1034
1035    /// Submit a transaction via the QuorumDriver. `transaction` must be the
1036    /// signature-verified form of `request.transaction`, and the caller must
1037    /// have run `validity_check` on it beforehand.
1038    #[instrument(level = "trace", skip_all, fields(tx_digest = ?request.transaction.digest()))]
1039    async fn execute_transaction_impl(
1040        &self,
1041        quorum_driver: &Arc<QuorumDriverHandler<A>>,
1042        epoch_store: &Arc<AuthorityPerEpochStore>,
1043        request: ExecuteTransactionRequestV1,
1044        transaction: VerifiedTransaction,
1045        client_addr: Option<SocketAddr>,
1046    ) -> Result<QuorumDriverResponse, QuorumDriverError> {
1047        let (_in_flight_metrics_guards, good_response_metrics) = self.update_metrics(&transaction);
1048        let tx_digest = *transaction.digest();
1049        debug!(?tx_digest, "TO Received transaction execution request.");
1050
1051        let (_e2e_latency_timer, _txn_finality_timer) = if transaction.contains_shared_object() {
1052            (
1053                self.metrics.request_latency_shared_obj.start_timer(),
1054                self.metrics
1055                    .wait_for_finality_latency_shared_obj
1056                    .start_timer(),
1057            )
1058        } else {
1059            (
1060                self.metrics.request_latency_single_writer.start_timer(),
1061                self.metrics
1062                    .wait_for_finality_latency_single_writer
1063                    .start_timer(),
1064            )
1065        };
1066
1067        // TODO: refactor all the gauge and timer metrics with `monitored_scope`
1068        let wait_for_finality_gauge = self.metrics.wait_for_finality_in_flight.clone();
1069        wait_for_finality_gauge.inc();
1070        let _wait_for_finality_gauge = scopeguard::guard(wait_for_finality_gauge, |in_flight| {
1071            in_flight.dec();
1072        });
1073
1074        let ticket = self
1075            .submit(
1076                quorum_driver,
1077                epoch_store.clone(),
1078                transaction.clone(),
1079                request,
1080                client_addr,
1081            )
1082            .await
1083            .map_err(|e| {
1084                warn!(?tx_digest, "QuorumDriverInternalError: {e:?}");
1085                QuorumDriverError::QuorumDriverInternal(e)
1086            })?;
1087
1088        let Ok(result) = timeout(WAIT_FOR_FINALITY_TIMEOUT, ticket).await else {
1089            debug!(?tx_digest, "Timeout waiting for transaction finality.");
1090            self.metrics.wait_for_finality_timeout.inc();
1091            return Err(QuorumDriverError::TimeoutBeforeFinality);
1092        };
1093        add_server_timing("wait_for_finality");
1094
1095        drop(_txn_finality_timer);
1096        drop(_wait_for_finality_gauge);
1097        self.metrics.wait_for_finality_finished.inc();
1098
1099        match result {
1100            Err(err) => {
1101                warn!(?tx_digest, "QuorumDriverInternalError: {err:?}");
1102                Err(QuorumDriverError::QuorumDriverInternal(err))
1103            }
1104            Ok(Err(err)) => Err(err),
1105            Ok(Ok(response)) => {
1106                good_response_metrics.inc();
1107                Ok(response)
1108            }
1109        }
1110    }
1111
1112    /// Submits the transaction to Quorum Driver for execution.
1113    /// Returns an awaitable Future.
1114    #[instrument(name = "tx_orchestrator_submit", level = "trace", skip_all)]
1115    async fn submit(
1116        &self,
1117        quorum_driver: &Arc<QuorumDriverHandler<A>>,
1118        epoch_store: Arc<AuthorityPerEpochStore>,
1119        transaction: VerifiedTransaction,
1120        request: ExecuteTransactionRequestV1,
1121        client_addr: Option<SocketAddr>,
1122    ) -> IotaResult<impl Future<Output = IotaResult<QuorumDriverResult>> + '_> {
1123        let tx_digest = *transaction.digest();
1124        let ticket = self.notifier.register_one(&tx_digest);
1125        // TODO(william) need to also write client adr to pending tx log below
1126        // so that we can re-execute with this client addr if we restart
1127        if self
1128            .pending_tx_log
1129            .write_pending_transaction_maybe(&transaction)
1130            .await?
1131        {
1132            debug!(?tx_digest, "no pending request in flight, submitting.");
1133            quorum_driver
1134                .submit_transaction_no_ticket(request.clone(), client_addr)
1135                .await?;
1136        }
1137        // It's possible that the transaction effects is already stored in DB at this
1138        // point. So we also subscribe to that. If we hear from `effects_await`
1139        // first, it means the ticket misses the previous notification, and we
1140        // want to ask quorum driver to form a certificate for us again, to
1141        // serve this request.
1142        let cache_reader = self.validator_state.get_transaction_cache_reader().clone();
1143        let qd = quorum_driver.clone();
1144        Ok(async move {
1145            let digests = [tx_digest];
1146            let effects_await =
1147                epoch_store.within_alive_epoch(cache_reader.try_notify_read_executed_effects(
1148                    "TransactionOrchestrator::notify_read_submit_with_qd",
1149                    &digests,
1150                ));
1151            // let-and-return necessary to satisfy borrow checker.
1152            let res = match select(ticket, effects_await.boxed()).await {
1153                Either::Left((quorum_driver_response, _)) => Ok(quorum_driver_response),
1154                Either::Right((_, unfinished_quorum_driver_task)) => {
1155                    debug!(
1156                        ?tx_digest,
1157                        "Effects are available in DB, use quorum driver to get a certificate"
1158                    );
1159                    qd.submit_transaction_no_ticket(request, client_addr)
1160                        .await?;
1161                    Ok(unfinished_quorum_driver_task.await)
1162                }
1163            };
1164            res
1165        })
1166    }
1167
1168    #[instrument(
1169        name = "tx_orchestrator_wait_for_finalized_tx_executed_locally_with_timeout",
1170        level = "debug",
1171        skip_all,
1172        fields(tx_digest = ?transaction.digest()),
1173        err
1174    )]
1175    async fn wait_for_finalized_tx_executed_locally_with_timeout(
1176        validator_state: &Arc<AuthorityState>,
1177        transaction: &VerifiedTransaction,
1178        metrics: &TransactionOrchestratorMetrics,
1179    ) -> IotaResult {
1180        let tx_digest = *transaction.digest();
1181        metrics.local_execution_in_flight.inc();
1182        let _metrics_guard =
1183            scopeguard::guard(metrics.local_execution_in_flight.clone(), |in_flight| {
1184                in_flight.dec();
1185            });
1186
1187        let _guard = if transaction.contains_shared_object() {
1188            metrics.local_execution_latency_shared_obj.start_timer()
1189        } else {
1190            metrics.local_execution_latency_single_writer.start_timer()
1191        };
1192        debug!(
1193            ?tx_digest,
1194            "Waiting for finalized tx to be executed locally."
1195        );
1196        match timeout(
1197            LOCAL_EXECUTION_TIMEOUT,
1198            validator_state
1199                .get_transaction_cache_reader()
1200                .try_notify_read_executed_effects_digests(
1201                    "TransactionOrchestrator::notify_read_wait_for_local_execution",
1202                    &[tx_digest],
1203                ),
1204        )
1205        .instrument(trace_span!("local_execution"))
1206        .await
1207        {
1208            Err(_elapsed) => {
1209                debug!(
1210                    ?tx_digest,
1211                    "Waiting for finalized tx to be executed locally timed out within {:?}.",
1212                    LOCAL_EXECUTION_TIMEOUT
1213                );
1214                metrics.local_execution_timeout.inc();
1215                Err(IotaError::Timeout)
1216            }
1217            Ok(Err(err)) => {
1218                debug!(
1219                    ?tx_digest,
1220                    "Waiting for finalized tx to be executed locally failed with error: {:?}", err
1221                );
1222                metrics.local_execution_failure.inc();
1223                Err(IotaError::TransactionOrchestratorLocalExecution {
1224                    error: err.to_string(),
1225                })
1226            }
1227            Ok(Ok(_)) => {
1228                metrics.local_execution_success.inc();
1229                Ok(())
1230            }
1231        }
1232    }
1233
1234    // TODO: Potentially cleanup this function and pending transaction log.
1235    async fn loop_pending_transaction_log(
1236        mut effects_receiver: Receiver<QuorumDriverEffectsQueueResult>,
1237        pending_transaction_log: Arc<WritePathPendingTransactionLog>,
1238    ) {
1239        loop {
1240            match effects_receiver.recv().await {
1241                Ok(Ok((transaction, ..))) => {
1242                    let tx_digest = transaction.digest();
1243                    if let Err(err) = pending_transaction_log.finish_transaction(tx_digest) {
1244                        error!(
1245                            ?tx_digest,
1246                            "Failed to finish transaction in pending transaction log: {err}"
1247                        );
1248                    }
1249                }
1250                Ok(Err((tx_digest, _err))) => {
1251                    if let Err(err) = pending_transaction_log.finish_transaction(&tx_digest) {
1252                        error!(
1253                            ?tx_digest,
1254                            "Failed to finish transaction in pending transaction log: {err}"
1255                        );
1256                    }
1257                }
1258                Err(RecvError::Closed) => {
1259                    error!("Sender of effects subscriber queue has been dropped!");
1260                    return;
1261                }
1262                Err(RecvError::Lagged(skipped_count)) => {
1263                    warn!("Skipped {skipped_count} transasctions in effects subscriber queue.");
1264                }
1265            }
1266        }
1267    }
1268
1269    /// Returns the quorum driver, or `None` when the node booted under
1270    /// P-COOL. Test-only: submissions must go through the per-request
1271    /// driver selection.
1272    #[cfg(any(test, feature = "test-utils"))]
1273    pub fn quorum_driver(&self) -> Option<&Arc<QuorumDriverHandler<A>>> {
1274        self.quorum_driver.as_ref()
1275    }
1276
1277    /// Owned variant of [`Self::quorum_driver`].
1278    #[cfg(any(test, feature = "test-utils"))]
1279    pub fn clone_quorum_driver(&self) -> Option<Arc<QuorumDriverHandler<A>>> {
1280        self.quorum_driver.clone()
1281    }
1282
1283    /// Returns the transaction driver's aggregator; it always exists and its
1284    /// reconfig observer keeps it current.
1285    pub fn clone_authority_aggregator(&self) -> Arc<AuthorityAggregator<A>> {
1286        self.transaction_driver.authority_aggregator().load_full()
1287    }
1288
1289    /// Returns an effects receiver only while the quorum driver is the
1290    /// currently selected flow and this node can serve it; the P-COOL flow
1291    /// has no effects broadcast.
1292    pub fn subscribe_to_effects_queue(&self) -> Option<Receiver<QuorumDriverEffectsQueueResult>> {
1293        let epoch_store = self.validator_state.load_epoch_store_one_call_per_task();
1294        if epoch_store.protocol_config().enable_pcool_flow() {
1295            return None;
1296        }
1297        self.quorum_driver
1298            .as_ref()
1299            .map(|quorum_driver| quorum_driver.subscribe_to_effects())
1300    }
1301
1302    /// Runs driver selection for `epoch_store` and reports its outcome.
1303    #[cfg(any(test, feature = "test-utils"))]
1304    pub fn select_driver_for_testing(
1305        &self,
1306        epoch_store: &AuthorityPerEpochStore,
1307    ) -> Result<(), QuorumDriverError> {
1308        self.select_driver(epoch_store).map(|_| ())
1309    }
1310
1311    fn update_metrics(
1312        &'_ self,
1313        transaction: &VerifiedTransaction,
1314    ) -> (impl Drop, &'_ GenericCounter<AtomicU64>) {
1315        let (in_flight, good_response) = if transaction.contains_shared_object() {
1316            self.metrics.total_req_received_shared_object.inc();
1317            (
1318                self.metrics.req_in_flight_shared_object.clone(),
1319                &self.metrics.good_response_shared_object,
1320            )
1321        } else {
1322            self.metrics.total_req_received_single_writer.inc();
1323            (
1324                self.metrics.req_in_flight_single_writer.clone(),
1325                &self.metrics.good_response_single_writer,
1326            )
1327        };
1328        in_flight.inc();
1329        (
1330            scopeguard::guard(in_flight, |in_flight| {
1331                in_flight.dec();
1332            }),
1333            good_response,
1334        )
1335    }
1336
1337    fn schedule_txes_in_log(
1338        pending_tx_log: Arc<WritePathPendingTransactionLog>,
1339        quorum_driver: Arc<QuorumDriverHandler<A>>,
1340    ) {
1341        if std::env::var("SKIP_LOADING_FROM_PENDING_TX_LOG").is_ok() {
1342            info!("Skipping loading pending transactions from pending_tx_log.");
1343            return;
1344        }
1345        spawn_logged_monitored_task!(async move {
1346            let pending_txes = pending_tx_log
1347                .load_all_pending_transactions()
1348                .expect("failed to load all pending transactions");
1349            info!(
1350                "Recovering {} pending transactions from pending_tx_log.",
1351                pending_txes.len()
1352            );
1353            for (i, tx) in pending_txes.into_iter().enumerate() {
1354                // TODO: ideally pending_tx_log would not contain VerifiedTransaction, but that
1355                // requires a migration.
1356                let tx = tx.into_inner();
1357                let tx_digest = *tx.digest();
1358                // It's not impossible we fail to enqueue a task but that's not the end of
1359                // world. TODO(william) correctly extract client_addr from logs
1360                if let Err(err) = quorum_driver
1361                    .submit_transaction_no_ticket(
1362                        ExecuteTransactionRequestV1 {
1363                            transaction: tx,
1364                            include_events: true,
1365                            include_input_objects: false,
1366                            include_output_objects: false,
1367                            include_auxiliary_data: false,
1368                        },
1369                        None,
1370                    )
1371                    .await
1372                {
1373                    warn!(
1374                        ?tx_digest,
1375                        "Failed to enqueue transaction from pending_tx_log, err: {err:?}"
1376                    );
1377                } else {
1378                    debug!(?tx_digest, "Enqueued transaction from pending_tx_log");
1379                    if (i + 1) % 1000 == 0 {
1380                        info!("Enqueued {} transactions from pending_tx_log.", i + 1);
1381                    }
1382                }
1383            }
1384            // Transactions will be cleaned up in
1385            // loop_execute_finalized_tx_locally() after they
1386            // produce effects.
1387        });
1388    }
1389
1390    pub fn load_all_pending_transactions(&self) -> IotaResult<Vec<VerifiedTransaction>> {
1391        self.pending_tx_log.load_all_pending_transactions()
1392    }
1393
1394    /// Reports whether a driver submission of `tx_digest` is in flight, and
1395    /// if so how many duplicate submissions are awaiting its outcome.
1396    #[cfg(any(test, feature = "test-utils"))]
1397    pub fn in_flight_duplicates_for_testing(&self, tx_digest: &TransactionDigest) -> Option<usize> {
1398        self.in_flight_transactions
1399            .lock()
1400            .get(tx_digest)
1401            .map(|sender| sender.receiver_count())
1402    }
1403}
1404
1405/// Convert a `QuorumDriverResponse` (contains
1406/// `VerifiedCertifiedTransactionEffects`) to the V1 response format that uses
1407/// `FinalizedEffects`.
1408fn quorum_driver_response_to_v1(response: QuorumDriverResponse) -> ExecuteTransactionResponseV1 {
1409    let QuorumDriverResponse {
1410        effects_cert,
1411        events,
1412        input_objects,
1413        output_objects,
1414        auxiliary_data,
1415    } = response;
1416    ExecuteTransactionResponseV1 {
1417        effects: FinalizedEffects::new_from_effects_cert(effects_cert.into()),
1418        events,
1419        input_objects,
1420        output_objects,
1421        auxiliary_data,
1422    }
1423}
1424
1425/// Convert a `transaction_driver_types::FinalizedEffects` into a
1426/// `quorum_driver_types::FinalizedEffects`.
1427fn convert_td_to_qd_effects(td: TdFinalizedEffects) -> FinalizedEffects {
1428    let finality_info = match td.finality_info {
1429        TdEffectsFinalityInfo::Certified(sig) => EffectsFinalityInfo::Certified(sig),
1430        TdEffectsFinalityInfo::Checkpointed(epoch, seq) => {
1431            EffectsFinalityInfo::Checkpointed(epoch, seq)
1432        }
1433        TdEffectsFinalityInfo::QuorumExecuted(epoch) => EffectsFinalityInfo::QuorumExecuted(epoch),
1434        TdEffectsFinalityInfo::UncertifiedSingleValidator(epoch) => {
1435            EffectsFinalityInfo::UncertifiedSingleValidator(epoch)
1436        }
1437    };
1438    FinalizedEffects {
1439        effects: td.effects,
1440        finality_info,
1441    }
1442}
1443
1444/// Map a `TransactionDriverError` to a `QuorumDriverError` for client
1445/// reporting. The variant choice signals retriability: clients retry on
1446/// `QuorumDriverInternal`, `FailedWithTransientErrorAfterMaximumAttempts`,
1447/// and `TimeoutBeforeFinality`, but treat `InvalidTransaction`,
1448/// `InvalidUserSignature`, and `RejectedByValidators` as terminal.
1449/// An overload-dominated timeout maps to the retriable `SystemOverload` or
1450/// `SystemOverloadRetryAfter` instead of `TimeoutBeforeFinality`.
1451/// Submission-time rejections that cannot succeed on resubmission must
1452/// therefore not be reported as internal.
1453fn map_td_error_to_qd(e: TransactionDriverError) -> QuorumDriverError {
1454    use TransactionDriverError::*;
1455    match e {
1456        ValidationFailed { error } => {
1457            QuorumDriverError::InvalidUserSignature(IotaError::InvalidSignature { error })
1458        }
1459        TimeoutWithLastRetriableError { last_error, .. } => last_error
1460            .and_then(|e| map_overload_to_qd(&e))
1461            .unwrap_or(QuorumDriverError::TimeoutBeforeFinality),
1462        RejectedByValidators {
1463            submission_non_retriable_errors,
1464            ..
1465        } => {
1466            // f+1 stake of validators returned non-retriable errors during
1467            // submission (bad signature, malformed tx, lock conflict, ...).
1468            // Some verdicts depend on validator-local state, so this does not
1469            // prove the transaction can never execute; the variant is kept
1470            // distinct from `InvalidTransaction` so clients can tell the two
1471            // apart.
1472            let representative = submission_non_retriable_errors
1473                .errors
1474                .into_iter()
1475                .next()
1476                .map(|(msg, _, _, _)| msg)
1477                .unwrap_or_else(|| "transaction rejected as invalid during submission".to_string());
1478            QuorumDriverError::RejectedByValidators(IotaError::Unknown(format!(
1479                "Transaction was rejected as invalid by more than 1/3 of validator stake \
1480                 during submission (non-retriable): {representative}"
1481            )))
1482        }
1483        Aborted {
1484            submission_retriable_errors,
1485            submission_non_retriable_errors,
1486            ..
1487        } => {
1488            // Driver exhausted the validator list without reaching the f+1
1489            // non-retriable threshold — most failures were transient
1490            // (validator down, network, overload). Surface as retriable so
1491            // the client can resubmit.
1492            let attempts = count_validator_attempts(&submission_retriable_errors)
1493                + count_validator_attempts(&submission_non_retriable_errors);
1494            QuorumDriverError::FailedWithTransientErrorAfterMaximumAttempts {
1495                total_attempts: attempts,
1496            }
1497        }
1498        other @ ForkedExecution { .. } => {
1499            // Validators disagree on effects digests — a protocol-level
1500            // invariant violation, never a client retry case. Log loud so
1501            // on-call sees it; surface as internal.
1502            let msg = other.to_string();
1503            error!("TransactionDriver observed forked execution: {msg}");
1504            QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(msg))
1505        }
1506        other @ ClientInternal { .. } => {
1507            let msg = other.to_string();
1508            warn!("TransactionDriver client-internal error: {msg}");
1509            QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(msg))
1510        }
1511        other @ SubmittedButFetchFailed { .. } => {
1512            let msg = other.to_string();
1513            warn!("TransactionDriver submitted transaction but failed to fetch effects: {msg}");
1514            QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(msg))
1515        }
1516    }
1517}
1518
1519/// Maps an overload-dominated aborted attempt to an overload error, so the
1520/// validators' retry hints reach the client instead of a generic timeout.
1521fn map_overload_to_qd(e: &TransactionDriverError) -> Option<QuorumDriverError> {
1522    use iota_types::{base_types::ConciseableName, error::ErrorCategory};
1523
1524    if !e.is_overload_dominated() {
1525        return None;
1526    }
1527    let TransactionDriverError::Aborted {
1528        submission_retriable_errors,
1529        submission_non_retriable_errors,
1530        ..
1531    } = e
1532    else {
1533        return None;
1534    };
1535    let mut overloaded_stake = 0;
1536    let mut errors: GroupedErrors = Vec::new();
1537    for (msg, names, stake, category) in submission_retriable_errors
1538        .errors
1539        .iter()
1540        .chain(submission_non_retriable_errors.errors.iter())
1541    {
1542        if *category != ErrorCategory::ValidatorOverloaded {
1543            continue;
1544        }
1545        overloaded_stake += *stake;
1546        errors.push((
1547            IotaError::Unknown(msg.clone()),
1548            *stake,
1549            names.iter().map(|n| n.concise_owned()).collect(),
1550        ));
1551    }
1552    Some(
1553        match submission_retriable_errors.median_retry_after_secs() {
1554            Some(retry_after_secs) => QuorumDriverError::SystemOverloadRetryAfter {
1555                overload_stake: overloaded_stake,
1556                errors,
1557                retry_after_secs,
1558            },
1559            None => QuorumDriverError::SystemOverload {
1560                overloaded_stake,
1561                errors,
1562            },
1563        },
1564    )
1565}
1566
1567fn count_validator_attempts(errors: &AggregatedRequestErrors) -> u32 {
1568    errors
1569        .errors
1570        .iter()
1571        .map(|(_, authorities, _, _)| authorities.len() as u32)
1572        .sum()
1573}
1574
1575/// Await a detached submission task, surfacing a task panic as an internal
1576/// error.
1577async fn join_submission_task<T>(
1578    handle: tokio::task::JoinHandle<Result<T, QuorumDriverError>>,
1579) -> Result<T, QuorumDriverError> {
1580    handle.await.unwrap_or_else(|e| {
1581        Err(QuorumDriverError::QuorumDriverInternal(IotaError::Unknown(
1582            format!("transaction submission task panicked: {e}"),
1583        )))
1584    })
1585}
1586
1587/// Prometheus metrics which can be displayed in Grafana, queried and alerted on
1588#[derive(Clone)]
1589pub struct TransactionOrchestratorMetrics {
1590    total_req_received_single_writer: GenericCounter<AtomicU64>,
1591    total_req_received_shared_object: GenericCounter<AtomicU64>,
1592
1593    good_response_single_writer: GenericCounter<AtomicU64>,
1594    good_response_shared_object: GenericCounter<AtomicU64>,
1595
1596    req_in_flight_single_writer: GenericGauge<AtomicI64>,
1597    req_in_flight_shared_object: GenericGauge<AtomicI64>,
1598
1599    wait_for_finality_in_flight: GenericGauge<AtomicI64>,
1600    wait_for_finality_finished: GenericCounter<AtomicU64>,
1601    wait_for_finality_timeout: GenericCounter<AtomicU64>,
1602
1603    local_execution_in_flight: GenericGauge<AtomicI64>,
1604    local_execution_success: GenericCounter<AtomicU64>,
1605    local_execution_timeout: GenericCounter<AtomicU64>,
1606    local_execution_failure: GenericCounter<AtomicU64>,
1607
1608    early_cached_response: GenericCounter<AtomicU64>,
1609
1610    // Bumped when the skip-effect-certification path reconciles against the
1611    // local cache but the cache has no events for a tx the single submitter
1612    // claimed had events. Uncertified events are rejected and the request
1613    // fails via the safety guard.
1614    skip_effect_cert_events_cache_miss: GenericCounter<AtomicU64>,
1615
1616    // Bumped when local checkpoint inclusion completes before the TD
1617    // skip-effect-certification call returns. Indicates the driver was slow
1618    // (e.g., corroborating a single-validator rejection) and the checkpoint
1619    // race cancelled the in-flight driver work in favor of rebuilding from
1620    // the local cache.
1621    skip_effect_cert_checkpoint_overrode_driver: GenericCounter<AtomicU64>,
1622
1623    request_latency_single_writer: Histogram,
1624    request_latency_shared_obj: Histogram,
1625    wait_for_finality_latency_single_writer: Histogram,
1626    wait_for_finality_latency_shared_obj: Histogram,
1627    local_execution_latency_single_writer: Histogram,
1628    local_execution_latency_shared_obj: Histogram,
1629}
1630
1631// Note that labeled-metrics are stored upfront individually
1632// to mitigate the perf hit by MetricsVec.
1633// See https://github.com/tikv/rust-prometheus/tree/master/static-metric
1634impl TransactionOrchestratorMetrics {
1635    pub fn new(registry: &Registry) -> Self {
1636        let total_req_received = register_int_counter_vec_with_registry!(
1637            "tx_orchestrator_total_req_received",
1638            "Total number of executions request Transaction Orchestrator receives, group by tx type",
1639            &["tx_type"],
1640            registry;
1641            MetricLevel::Warn,
1642        )
1643            .unwrap();
1644
1645        let total_req_received_single_writer =
1646            total_req_received.with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]);
1647        let total_req_received_shared_object =
1648            total_req_received.with_label_values(&[TX_TYPE_SHARED_OBJ_TX]);
1649
1650        let good_response = register_int_counter_vec_with_registry!(
1651            "tx_orchestrator_good_response",
1652            "Total number of good responses Transaction Orchestrator generates, group by tx type",
1653            &["tx_type"],
1654            registry;
1655            MetricLevel::Warn,
1656        )
1657        .unwrap();
1658
1659        let good_response_single_writer =
1660            good_response.with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]);
1661        let good_response_shared_object = good_response.with_label_values(&[TX_TYPE_SHARED_OBJ_TX]);
1662
1663        let req_in_flight = register_int_gauge_vec_with_registry!(
1664            "tx_orchestrator_req_in_flight",
1665            "Number of requests in flights Transaction Orchestrator processes, group by tx type",
1666            &["tx_type"],
1667            registry;
1668            MetricLevel::Warn,
1669        )
1670        .unwrap();
1671
1672        let req_in_flight_single_writer =
1673            req_in_flight.with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]);
1674        let req_in_flight_shared_object = req_in_flight.with_label_values(&[TX_TYPE_SHARED_OBJ_TX]);
1675
1676        let request_latency = register_histogram_vec_with_registry!(
1677            "tx_orchestrator_request_latency",
1678            "Time spent in processing one Transaction Orchestrator request",
1679            &["tx_type"],
1680            iota_metrics::COARSE_LATENCY_SEC_BUCKETS.to_vec(),
1681            registry;
1682            MetricLevel::Warn,
1683        )
1684        .unwrap();
1685        let wait_for_finality_latency = register_histogram_vec_with_registry!(
1686            "tx_orchestrator_wait_for_finality_latency",
1687            "Time spent in waiting for one Transaction Orchestrator request gets finalized",
1688            &["tx_type"],
1689            iota_metrics::COARSE_LATENCY_SEC_BUCKETS.to_vec(),
1690            registry;
1691            MetricLevel::Warn,
1692        )
1693        .unwrap();
1694        let local_execution_latency = register_histogram_vec_with_registry!(
1695            "tx_orchestrator_local_execution_latency",
1696            "Time spent in waiting for one Transaction Orchestrator gets locally executed",
1697            &["tx_type"],
1698            iota_metrics::COARSE_LATENCY_SEC_BUCKETS.to_vec(),
1699            registry;
1700            MetricLevel::Warn,
1701        )
1702        .unwrap();
1703
1704        Self {
1705            total_req_received_single_writer,
1706            total_req_received_shared_object,
1707            good_response_single_writer,
1708            good_response_shared_object,
1709            req_in_flight_single_writer,
1710            req_in_flight_shared_object,
1711            wait_for_finality_in_flight: register_int_gauge_with_registry!(
1712                "tx_orchestrator_wait_for_finality_in_flight",
1713                "Number of in flight txns Transaction Orchestrator are waiting for finality for",
1714                registry;
1715                MetricLevel::Warn,
1716            )
1717                .unwrap(),
1718            wait_for_finality_finished: register_int_counter_with_registry!(
1719                "tx_orchestrator_wait_for_finality_finished",
1720                "Total number of txns Transaction Orchestrator gets responses from Quorum Driver before timeout, either success or failure",
1721                registry;
1722                MetricLevel::Warn,
1723            )
1724                .unwrap(),
1725            wait_for_finality_timeout: register_int_counter_with_registry!(
1726                "tx_orchestrator_wait_for_finality_timeout",
1727                "Total number of txns timing out in waiting for finality Transaction Orchestrator handles",
1728                registry;
1729                MetricLevel::Warn,
1730            )
1731                .unwrap(),
1732            local_execution_in_flight: register_int_gauge_with_registry!(
1733                "tx_orchestrator_local_execution_in_flight",
1734                "Number of local execution txns in flights Transaction Orchestrator handles",
1735                registry;
1736                MetricLevel::Warn,
1737            )
1738                .unwrap(),
1739            local_execution_success: register_int_counter_with_registry!(
1740                "tx_orchestrator_local_execution_success",
1741                "Total number of successful local execution txns Transaction Orchestrator handles",
1742                registry;
1743                MetricLevel::Warn,
1744            )
1745                .unwrap(),
1746            local_execution_timeout: register_int_counter_with_registry!(
1747                "tx_orchestrator_local_execution_timeout",
1748                "Total number of timed-out local execution txns Transaction Orchestrator handles",
1749                registry;
1750                MetricLevel::Warn,
1751            )
1752                .unwrap(),
1753            local_execution_failure: register_int_counter_with_registry!(
1754                "tx_orchestrator_local_execution_failure",
1755                "Total number of failed local execution txns Transaction Orchestrator handles",
1756                registry;
1757                MetricLevel::Warn,
1758            )
1759                .unwrap(),
1760            early_cached_response: register_int_counter_with_registry!(
1761                "tx_orchestrator_early_cached_response",
1762                "Total number of requests returning cached results for already-executed transactions",
1763                registry,
1764            )
1765                .unwrap(),
1766            skip_effect_cert_events_cache_miss: register_int_counter_with_registry!(
1767                "tx_orchestrator_skip_effect_cert_events_cache_miss",
1768                "Number of skip-effect-certification responses rejected because the \
1769                 single submitter claimed to have events but the local cache did not \
1770                 corroborate them",
1771                registry,
1772            )
1773                .unwrap(),
1774            skip_effect_cert_checkpoint_overrode_driver: register_int_counter_with_registry!(
1775                "tx_orchestrator_skip_effect_cert_checkpoint_overrode_driver",
1776                "Number of skip-effect-certification requests where local checkpoint \
1777                 inclusion completed before the TransactionDriver call returned; the \
1778                 driver future was cancelled and the response was rebuilt from cache",
1779                registry,
1780            )
1781                .unwrap(),
1782            request_latency_single_writer: request_latency
1783                .with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]),
1784            request_latency_shared_obj: request_latency.with_label_values(&[TX_TYPE_SHARED_OBJ_TX]),
1785            wait_for_finality_latency_single_writer: wait_for_finality_latency
1786                .with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]),
1787            wait_for_finality_latency_shared_obj: wait_for_finality_latency
1788                .with_label_values(&[TX_TYPE_SHARED_OBJ_TX]),
1789            local_execution_latency_single_writer: local_execution_latency
1790                .with_label_values(&[TX_TYPE_SINGLE_WRITER_TX]),
1791            local_execution_latency_shared_obj: local_execution_latency
1792                .with_label_values(&[TX_TYPE_SHARED_OBJ_TX]),
1793        }
1794    }
1795
1796    pub fn new_for_tests() -> Self {
1797        let registry = Registry::new();
1798        Self::new(&registry)
1799    }
1800}
1801
1802#[async_trait::async_trait]
1803impl<A> iota_types::transaction_executor::TransactionExecutor for TransactionOrchestrator<A>
1804where
1805    A: AuthorityAPI + Send + Sync + 'static + Clone,
1806{
1807    async fn execute_transaction(
1808        &self,
1809        request: ExecuteTransactionRequestV1,
1810        skip_certification: bool,
1811        client_addr: Option<std::net::SocketAddr>,
1812    ) -> Result<ExecuteTransactionResponseV1, QuorumDriverError> {
1813        self.execute_transaction_v1(request, skip_certification, client_addr)
1814            .await
1815    }
1816
1817    fn simulate_transaction(
1818        &self,
1819        transaction: Transaction,
1820        checks: VmChecks,
1821    ) -> Result<SimulateTransactionResult, IotaError> {
1822        self.validator_state
1823            .simulate_transaction(transaction, checks)
1824    }
1825
1826    /// Wait for the given transactions to be included in a checkpoint.
1827    ///
1828    /// Returns a mapping from transaction digest to
1829    /// `(checkpoint_sequence_number, checkpoint_timestamp_ms)`.
1830    /// On timeout, returns partial results for any transactions that were
1831    /// already checkpointed.
1832    async fn wait_for_checkpoint_inclusion(
1833        &self,
1834        digests: &[TransactionDigest],
1835        timeout: Duration,
1836    ) -> Result<BTreeMap<TransactionDigest, (CheckpointSequenceNumber, u64)>, IotaError> {
1837        self.validator_state
1838            .wait_for_checkpoint_inclusion(digests, timeout)
1839            .await
1840    }
1841
1842    fn read_transaction_from_cache(
1843        &self,
1844        digest: &TransactionDigest,
1845        include_events: bool,
1846        include_input_objects: bool,
1847        include_output_objects: bool,
1848    ) -> Result<Option<iota_types::transaction_executor::CachedTransactionData>, IotaError> {
1849        read_cached_transaction_data(
1850            &self.validator_state,
1851            digest,
1852            include_events,
1853            include_input_objects,
1854            include_output_objects,
1855        )
1856    }
1857}
1858
1859/// Read a transaction's authoritative data from the local cache. Returns
1860/// `Ok(None)` if the tx hasn't been executed locally yet. Shared by the
1861/// orchestrator's skip-cert response builder and the `TransactionExecutor`
1862/// trait method consumed by the gRPC handler.
1863fn read_cached_transaction_data(
1864    validator_state: &Arc<AuthorityState>,
1865    digest: &TransactionDigest,
1866    include_events: bool,
1867    include_input_objects: bool,
1868    include_output_objects: bool,
1869) -> Result<Option<iota_types::transaction_executor::CachedTransactionData>, IotaError> {
1870    let cache = validator_state.get_transaction_cache_reader();
1871    let Some(effects) = cache.try_get_executed_effects(digest)? else {
1872        return Ok(None);
1873    };
1874
1875    let events = if include_events && effects.events_digest().is_some() {
1876        Some(validator_state.get_transaction_events(digest)?)
1877    } else {
1878        None
1879    };
1880
1881    let input_objects = if include_input_objects {
1882        Some(
1883            validator_state
1884                .get_transaction_input_objects(&effects)
1885                .map_err(|e| IotaError::Unknown(format!("input objects: {e:?}")))?,
1886        )
1887    } else {
1888        None
1889    };
1890    let output_objects = if include_output_objects {
1891        Some(
1892            validator_state
1893                .get_transaction_output_objects(&effects)
1894                .map_err(|e| IotaError::Unknown(format!("output objects: {e:?}")))?,
1895        )
1896    } else {
1897        None
1898    };
1899
1900    Ok(Some(
1901        iota_types::transaction_executor::CachedTransactionData {
1902            effects,
1903            events,
1904            input_objects,
1905            output_objects,
1906        },
1907    ))
1908}
1909
1910/// Successful outcome of an in-flight driver submission: the unfiltered
1911/// driver response, or the local checkpoint the transaction was observed in
1912/// when the checkpoint race cancelled the driver (the cache is then the
1913/// authoritative source of the effects).
1914/// Outcome of an in-flight driver submission, shared with concurrent
1915/// submissions of the same digest.
1916type InFlightSubmissionResult = Result<Arc<QuorumTransactionResponse>, QuorumDriverError>;
1917
1918/// Digests currently being driven to finality by the TransactionDriver,
1919/// each with a channel through which the driving submission publishes its
1920/// outcome to concurrent duplicates.
1921type InFlightTransactions =
1922    Arc<Mutex<HashMap<TransactionDigest, watch::Sender<Option<InFlightSubmissionResult>>>>>;
1923
1924/// Result of trying to register a submission of a digest in the in-flight
1925/// map: either this caller drives the committee-wide submission, or another
1926/// submission of the same digest is already in flight and this caller should
1927/// await its published outcome instead.
1928enum TransactionSubmission {
1929    Driving(TransactionSubmissionGuard),
1930    AlreadyInFlight(watch::Receiver<Option<InFlightSubmissionResult>>),
1931}
1932
1933/// Tracks a transaction that is being submitted to finality so that
1934/// concurrent submissions of the same digest deduplicate.
1935///
1936/// Held only by the driving submission, which must `publish` its outcome so
1937/// concurrent duplicates can return it. Dropping the guard removes the
1938/// digest from the in-flight map on every exit path (success, error,
1939/// timeout, and cancellation); receivers subscribed before removal still
1940/// observe a published outcome, and if the entry is removed without any
1941/// outcome (checkpoint-race cancellation, panic, or shutdown) the closed
1942/// channel tells duplicates to fall back to checkpoint inclusion.
1943struct TransactionSubmissionGuard {
1944    in_flight_transactions: InFlightTransactions,
1945    tx_digest: TransactionDigest,
1946}
1947
1948impl TransactionSubmissionGuard {
1949    fn acquire(
1950        in_flight_transactions: InFlightTransactions,
1951        tx_digest: TransactionDigest,
1952    ) -> TransactionSubmission {
1953        {
1954            let mut in_flight = in_flight_transactions.lock();
1955            match in_flight.entry(tx_digest) {
1956                Entry::Occupied(entry) => {
1957                    return TransactionSubmission::AlreadyInFlight(entry.get().subscribe());
1958                }
1959                Entry::Vacant(entry) => {
1960                    let (sender, _initial_receiver) = watch::channel(None);
1961                    entry.insert(sender);
1962                    debug!(?tx_digest, "added transaction to in-flight map");
1963                }
1964            }
1965        }
1966        TransactionSubmission::Driving(Self {
1967            in_flight_transactions,
1968            tx_digest,
1969        })
1970    }
1971
1972    /// Publish the submission outcome to concurrent duplicate submissions.
1973    /// The outcome is stored in the channel even when nobody is subscribed
1974    /// yet, so a duplicate that subscribes after this call but before the
1975    /// entry is removed still reads it instead of a closed channel.
1976    fn publish(&self, result: InFlightSubmissionResult) {
1977        if let Some(sender) = self.in_flight_transactions.lock().get(&self.tx_digest) {
1978            sender.send_replace(Some(result));
1979        }
1980    }
1981}
1982
1983impl Drop for TransactionSubmissionGuard {
1984    fn drop(&mut self) {
1985        self.in_flight_transactions.lock().remove(&self.tx_digest);
1986    }
1987}
1988
1989#[cfg(test)]
1990mod tests {
1991    use super::*;
1992
1993    fn acquire_driving(
1994        in_flight: &InFlightTransactions,
1995        tx_digest: TransactionDigest,
1996    ) -> TransactionSubmissionGuard {
1997        match TransactionSubmissionGuard::acquire(in_flight.clone(), tx_digest) {
1998            TransactionSubmission::Driving(guard) => guard,
1999            TransactionSubmission::AlreadyInFlight(_) => {
2000                panic!("expected to acquire the driving submission")
2001            }
2002        }
2003    }
2004
2005    fn acquire_duplicate(
2006        in_flight: &InFlightTransactions,
2007        tx_digest: TransactionDigest,
2008    ) -> watch::Receiver<Option<InFlightSubmissionResult>> {
2009        match TransactionSubmissionGuard::acquire(in_flight.clone(), tx_digest) {
2010            TransactionSubmission::Driving(_) => {
2011                panic!("expected the digest to already be in flight")
2012            }
2013            TransactionSubmission::AlreadyInFlight(receiver) => receiver,
2014        }
2015    }
2016
2017    #[tokio::test]
2018    async fn duplicate_submission_receives_published_outcome() {
2019        let in_flight = InFlightTransactions::default();
2020        let tx_digest = TransactionDigest::random();
2021
2022        let guard = acquire_driving(&in_flight, tx_digest);
2023        let mut receiver = acquire_duplicate(&in_flight, tx_digest);
2024
2025        guard.publish(Err(QuorumDriverError::TimeoutBeforeFinality));
2026        drop(guard);
2027
2028        // The published outcome must survive the guard drop for receivers
2029        // subscribed before the entry was removed.
2030        let outcome = receiver
2031            .wait_for(|outcome| outcome.is_some())
2032            .await
2033            .expect("outcome was published before the sender dropped")
2034            .clone()
2035            .expect("wait_for only returns once the outcome is Some");
2036        assert!(matches!(
2037            outcome,
2038            Err(QuorumDriverError::TimeoutBeforeFinality)
2039        ));
2040        assert!(
2041            in_flight.lock().is_empty(),
2042            "guard drop must remove the in-flight entry"
2043        );
2044    }
2045
2046    #[tokio::test]
2047    async fn duplicate_subscribing_after_publish_receives_outcome() {
2048        let in_flight = InFlightTransactions::default();
2049        let tx_digest = TransactionDigest::random();
2050
2051        let guard = acquire_driving(&in_flight, tx_digest);
2052        guard.publish(Err(QuorumDriverError::TimeoutBeforeFinality));
2053
2054        // Subscribing between the publish and the entry removal must still
2055        // resolve to the outcome; falling back to checkpoint inclusion here
2056        // would cost the duplicate a full finality timeout.
2057        let mut receiver = acquire_duplicate(&in_flight, tx_digest);
2058        drop(guard);
2059
2060        let outcome = receiver
2061            .wait_for(|outcome| outcome.is_some())
2062            .await
2063            .expect("the outcome is stored in the channel regardless of subscribers")
2064            .clone()
2065            .expect("wait_for only returns once the outcome is Some");
2066        assert!(matches!(
2067            outcome,
2068            Err(QuorumDriverError::TimeoutBeforeFinality)
2069        ));
2070    }
2071
2072    #[tokio::test]
2073    async fn dropped_guard_without_outcome_closes_channel() {
2074        let in_flight = InFlightTransactions::default();
2075        let tx_digest = TransactionDigest::random();
2076
2077        let guard = acquire_driving(&in_flight, tx_digest);
2078        let mut receiver = acquire_duplicate(&in_flight, tx_digest);
2079        drop(guard);
2080
2081        receiver
2082            .wait_for(|outcome| outcome.is_some())
2083            .await
2084            .expect_err("dropping the guard without publishing must close the channel");
2085        assert!(in_flight.lock().is_empty());
2086
2087        // The digest can be driven again once the entry is gone.
2088        let _guard = acquire_driving(&in_flight, tx_digest);
2089    }
2090
2091    async fn build_orchestrator_with_pcool(
2092        enable_pcool: bool,
2093    ) -> (
2094        Arc<AuthorityState>,
2095        TransactionOrchestrator<NetworkAuthorityClient>,
2096        tempfile::TempDir,
2097        tokio::sync::broadcast::Sender<IotaSystemState>,
2098    ) {
2099        use iota_protocol_config::{Chain, ProtocolConfig, ProtocolVersion};
2100
2101        use crate::{
2102            authority::test_authority_builder::TestAuthorityBuilder,
2103            authority_aggregator::AuthorityAggregatorBuilder,
2104        };
2105
2106        telemetry_subscribers::init_for_testing();
2107        let network_config =
2108            iota_swarm_config::network_config_builder::ConfigBuilder::new_with_temp_dir().build();
2109
2110        let mut protocol_config =
2111            ProtocolConfig::get_for_version(ProtocolVersion::MAX, Chain::Unknown);
2112        protocol_config.set_enable_pcool_flow_for_testing(enable_pcool);
2113        let state = TestAuthorityBuilder::new()
2114            .with_network_config(&network_config, 0)
2115            .with_protocol_config(protocol_config)
2116            .build()
2117            .await;
2118
2119        let (aggregator, _clients) =
2120            AuthorityAggregatorBuilder::from_genesis(&network_config.genesis)
2121                .build_network_clients();
2122        let (reconfig_tx, reconfig_rx) = tokio::sync::broadcast::channel(16);
2123        let tempdir = tempfile::tempdir().unwrap();
2124        let orchestrator = TransactionOrchestrator::new_with_auth_aggregator(
2125            Arc::new(aggregator),
2126            state.clone(),
2127            reconfig_rx,
2128            tempdir.path(),
2129            &Registry::new(),
2130            None,
2131        );
2132        (state, orchestrator, tempdir, reconfig_tx)
2133    }
2134
2135    /// A flag-off boot builds the quorum driver and runs WAL recovery at
2136    /// construction.
2137    #[tokio::test(flavor = "multi_thread")]
2138    async fn qd_recovery_eager_on_flag_off_boot() {
2139        let (state, orchestrator, _tempdir, _reconfig_tx) =
2140            build_orchestrator_with_pcool(false).await;
2141        assert!(orchestrator.quorum_driver().is_some());
2142        assert!(
2143            orchestrator
2144                .select_driver_for_testing(&state.epoch_store_for_testing())
2145                .is_ok()
2146        );
2147    }
2148
2149    /// A node booted under P-COOL has no quorum driver. After a rollback it
2150    /// must reject quorum-driver selection until restarted.
2151    #[tokio::test(flavor = "multi_thread")]
2152    async fn flag_on_boot_rejects_selection_after_rollback() {
2153        use iota_protocol_config::{Chain, ProtocolConfig, ProtocolVersion};
2154
2155        let (state, orchestrator, _tempdir, _reconfig_tx) =
2156            build_orchestrator_with_pcool(true).await;
2157        assert!(orchestrator.quorum_driver().is_none());
2158
2159        // Selection under the flag serves the P-COOL flow.
2160        assert!(
2161            orchestrator
2162                .select_driver_for_testing(&state.epoch_store_for_testing())
2163                .is_ok()
2164        );
2165
2166        // Epoch 1 with P-COOL off: selection and the effects queue must both
2167        // report the missing quorum driver.
2168        let mut protocol_config =
2169            ProtocolConfig::get_for_version(ProtocolVersion::MAX, Chain::Unknown);
2170        protocol_config.set_enable_pcool_flow_for_testing(false);
2171        state
2172            .reconfigure_for_testing_with_protocol_config(protocol_config)
2173            .await;
2174        let epoch_store = state.epoch_store_for_testing();
2175        assert_eq!(epoch_store.epoch(), 1);
2176
2177        assert!(matches!(
2178            orchestrator.select_driver_for_testing(&epoch_store),
2179            Err(QuorumDriverError::QuorumDriverInternal(_))
2180        ));
2181        assert!(orchestrator.subscribe_to_effects_queue().is_none());
2182    }
2183
2184    /// Validator rejections must map to `RejectedByValidators`, not
2185    /// `InvalidTransaction`: the verdicts can be validator-local, so clients
2186    /// must be able to tell them apart from a locally proven invalid
2187    /// transaction.
2188    #[test]
2189    fn map_validator_rejections_to_distinct_variant() {
2190        use iota_types::error::ErrorCategory;
2191
2192        use crate::transaction_driver::AggregatedRequestErrors;
2193
2194        let td_error = TransactionDriverError::RejectedByValidators {
2195            submission_non_retriable_errors: AggregatedRequestErrors {
2196                errors: vec![(
2197                    "Object lock conflict".to_string(),
2198                    vec![],
2199                    3334,
2200                    ErrorCategory::LockConflict,
2201                )],
2202                total_stake: 3334,
2203                stake_requested_retry_after: Default::default(),
2204            },
2205            submission_retriable_errors: AggregatedRequestErrors::default(),
2206        };
2207
2208        let qd_error = map_td_error_to_qd(td_error);
2209        let QuorumDriverError::RejectedByValidators(inner) = &qd_error else {
2210            panic!("expected RejectedByValidators, got {qd_error:?}");
2211        };
2212        assert!(inner.to_string().contains("Object lock conflict"));
2213        assert_eq!(qd_error.reason(), "rejected_by_validators");
2214    }
2215
2216    fn timeout_with_last_error(last_error: Option<TransactionDriverError>) -> QuorumDriverError {
2217        map_td_error_to_qd(TransactionDriverError::TimeoutWithLastRetriableError {
2218            last_error: last_error.map(Box::new),
2219            attempts: 3,
2220            timeout: std::time::Duration::from_secs(60),
2221        })
2222    }
2223
2224    fn aborted_with_retriable(retriable_errors: AggregatedRequestErrors) -> TransactionDriverError {
2225        TransactionDriverError::Aborted {
2226            submission_non_retriable_errors: AggregatedRequestErrors::default(),
2227            submission_retriable_errors: retriable_errors,
2228            observed_effects_digests: crate::transaction_driver::AggregatedEffectsDigests {
2229                digests: vec![],
2230            },
2231        }
2232    }
2233
2234    fn bucket(
2235        msg: &str,
2236        stake: u64,
2237        category: iota_types::error::ErrorCategory,
2238    ) -> (
2239        String,
2240        Vec<iota_types::base_types::AuthorityName>,
2241        u64,
2242        iota_types::error::ErrorCategory,
2243    ) {
2244        (msg.to_string(), vec![], stake, category)
2245    }
2246
2247    /// Overload split across delays still beats a larger single bucket. The
2248    /// 50/50 hint split resolves to the longer delay.
2249    #[test]
2250    fn map_overloaded_timeout_to_retry_after() {
2251        use std::collections::BTreeMap;
2252
2253        use iota_types::error::ErrorCategory;
2254
2255        let qd_error =
2256            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2257                errors: vec![
2258                    bucket("validator unavailable", 4000, ErrorCategory::Unavailable),
2259                    bucket("overloaded 10s", 2500, ErrorCategory::ValidatorOverloaded),
2260                    bucket("overloaded 30s", 2500, ErrorCategory::ValidatorOverloaded),
2261                ],
2262                total_stake: 9000,
2263                stake_requested_retry_after: BTreeMap::from([(10, 2500), (30, 2500)]),
2264            })));
2265        let QuorumDriverError::SystemOverloadRetryAfter {
2266            overload_stake,
2267            retry_after_secs,
2268            ..
2269        } = qd_error
2270        else {
2271            panic!("expected SystemOverloadRetryAfter, got {qd_error:?}");
2272        };
2273        assert_eq!(overload_stake, 5000);
2274        assert_eq!(retry_after_secs, 30);
2275    }
2276
2277    /// A low-stake validator must not dictate the retry hint.
2278    #[test]
2279    fn map_overloaded_timeout_ignores_low_stake_outlier() {
2280        use std::collections::BTreeMap;
2281
2282        use iota_types::error::ErrorCategory;
2283
2284        let qd_error =
2285            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2286                errors: vec![
2287                    bucket("overloaded 10s", 7000, ErrorCategory::ValidatorOverloaded),
2288                    bucket("overloaded 3600s", 1000, ErrorCategory::ValidatorOverloaded),
2289                ],
2290                total_stake: 8000,
2291                stake_requested_retry_after: BTreeMap::from([(10, 7000), (3600, 1000)]),
2292            })));
2293        assert!(
2294            matches!(
2295                qd_error,
2296                QuorumDriverError::SystemOverloadRetryAfter {
2297                    retry_after_secs: 10,
2298                    ..
2299                }
2300            ),
2301            "expected the 10s majority hint, got {qd_error:?}"
2302        );
2303    }
2304
2305    /// Overload without hints maps to `SystemOverload`. Non-dominant or
2306    /// absent overload keeps mapping to `TimeoutBeforeFinality`.
2307    #[test]
2308    fn map_overloaded_timeout_without_hint_and_non_overload() {
2309        use iota_types::error::ErrorCategory;
2310
2311        let qd_error =
2312            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2313                errors: vec![bucket(
2314                    "too many transactions pending",
2315                    5000,
2316                    ErrorCategory::ValidatorOverloaded,
2317                )],
2318                total_stake: 5000,
2319                stake_requested_retry_after: Default::default(),
2320            })));
2321        assert!(
2322            matches!(qd_error, QuorumDriverError::SystemOverload { overloaded_stake, .. } if overloaded_stake == 5000),
2323            "expected SystemOverload, got {qd_error:?}"
2324        );
2325
2326        // A hint backed by a minority of the overloaded stake is not exposed.
2327        let qd_error =
2328            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2329                errors: vec![
2330                    bucket(
2331                        "too many transactions pending",
2332                        6000,
2333                        ErrorCategory::ValidatorOverloaded,
2334                    ),
2335                    bucket("overloaded 3600s", 100, ErrorCategory::ValidatorOverloaded),
2336                ],
2337                total_stake: 6100,
2338                stake_requested_retry_after: std::collections::BTreeMap::from([(3600, 100)]),
2339            })));
2340        assert!(
2341            matches!(qd_error, QuorumDriverError::SystemOverload { overloaded_stake, .. } if overloaded_stake == 6100),
2342            "expected SystemOverload without a hint, got {qd_error:?}"
2343        );
2344
2345        let qd_error =
2346            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2347                errors: vec![
2348                    bucket("overloaded 10s", 3500, ErrorCategory::ValidatorOverloaded),
2349                    bucket("timed out submitting", 3000, ErrorCategory::Unavailable),
2350                    bucket(
2351                        "timed out getting effects",
2352                        2500,
2353                        ErrorCategory::Unavailable,
2354                    ),
2355                ],
2356                total_stake: 9000,
2357                stake_requested_retry_after: Default::default(),
2358            })));
2359        assert!(matches!(qd_error, QuorumDriverError::TimeoutBeforeFinality));
2360
2361        let qd_error =
2362            timeout_with_last_error(Some(aborted_with_retriable(AggregatedRequestErrors {
2363                errors: vec![bucket(
2364                    "validator unavailable",
2365                    5000,
2366                    ErrorCategory::Unavailable,
2367                )],
2368                total_stake: 5000,
2369                stake_requested_retry_after: Default::default(),
2370            })));
2371        assert!(matches!(qd_error, QuorumDriverError::TimeoutBeforeFinality));
2372
2373        let qd_error = timeout_with_last_error(None);
2374        assert!(matches!(qd_error, QuorumDriverError::TimeoutBeforeFinality));
2375    }
2376}