1use 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
82const LOCAL_EXECUTION_TIMEOUT: Duration = Duration::from_secs(10);
85
86const WAIT_FOR_FINALITY_TIMEOUT: Duration = Duration::from_secs(30);
87
88enum Driver<A: Clone> {
91 Quorum(Arc<QuorumDriverHandler<A>>),
93 Transaction(Arc<TransactionDriver<A>>),
95}
96
97pub struct TransactionOrchestrator<A: Clone> {
108 quorum_driver: Option<Arc<QuorumDriverHandler<A>>>,
109 transaction_driver: Arc<TransactionDriver<A>>,
110 validator_state: Arc<AuthorityState>,
111 _local_executor_handle: Option<JoinHandle<()>>,
114 pending_tx_log: Arc<WritePathPendingTransactionLog>,
115 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 #[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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 #[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 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 #[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 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 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 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 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 #[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 #[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 pub fn clone_authority_aggregator(&self) -> Arc<AuthorityAggregator<A>> {
1286 self.transaction_driver.authority_aggregator().load_full()
1287 }
1288
1289 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 #[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 let tx = tx.into_inner();
1357 let tx_digest = *tx.digest();
1358 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 });
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 #[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
1405fn 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
1425fn 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
1444fn 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 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 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 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
1519fn 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
1575async 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#[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 skip_effect_cert_events_cache_miss: GenericCounter<AtomicU64>,
1615
1616 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
1631impl 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(®istry)
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 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
1859fn 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
1910type InFlightSubmissionResult = Result<Arc<QuorumTransactionResponse>, QuorumDriverError>;
1917
1918type InFlightTransactions =
1922 Arc<Mutex<HashMap<TransactionDigest, watch::Sender<Option<InFlightSubmissionResult>>>>>;
1923
1924enum TransactionSubmission {
1929 Driving(TransactionSubmissionGuard),
1930 AlreadyInFlight(watch::Receiver<Option<InFlightSubmissionResult>>),
1931}
1932
1933struct 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 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 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 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 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 #[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 #[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 assert!(
2161 orchestrator
2162 .select_driver_for_testing(&state.epoch_store_for_testing())
2163 .is_ok()
2164 );
2165
2166 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 #[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 #[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 #[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 #[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 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}