From ebdb9a8be7f39935818f730d1a0e13a6ba047827 Mon Sep 17 00:00:00 2001 From: Sarfaraz Nawaz Date: Tue, 6 Oct 2026 11:36:49 +0530 Subject: [PATCH] refactor(committor): unify transaction execution Combined commits no longer need separate executor implementations for single-stage execution, the first transaction of a split execution, and the subsequent actions or undelegations. Replace their duplicated retry loops and timeout adapters with one transaction executor. - Track the current strategy, an optional pending strategy, the confirmed commit signature, and the current attempt count in one concrete runner. Remove SingleStageExecutor, TwoStageExecutor, Initialized, Committed, the sealed trait, StageExecutor, and its three adapters. - Route selected strategies through execute_strategy. Borrow the existing authority, clients, preparator, callback scheduler, and committed-account list instead of cloning services or rebuilding the account list. - Preserve nonce, action, and undelegation recovery according to whether the commit has already landed. Keep execution-limit fallback splitting after the last combined commit and propagate recovered uniqueness nonces to the pending transaction. - Keep the ten-attempt recovery limit per transaction. Preserve attempts consumed before a timeout, resetting the count only when advancing to another transaction or switching to split execution. - Apply the existing intent-wide callback deadline to the shared loop, including callbacks in the pending transaction. Report each successful transaction's callbacks before preparing the next transaction and retain the existing callback behavior for preparation and recovery failures. - Reuse the confirmed commit signature when removing failed follow-up actions leaves no work. Preserve the failure when undelegation recovery empties the transaction, and prevent whole-intent retries after a known successful commit. - Surrender current, pending, and replaced strategies for cleanup on terminal paths. Record recovery cleanup before awaiting blockhash-cache invalidation so cancellation cannot lose that cleanup information. - Move integration tests to the shared entry point and check callback timing at the next transaction's preparation boundary. Add six focused unit tests covering timeout, retry counts, preparation and nonce-fetch failures, cleanup, and recovery after a confirmed commit. Keep transaction packing thresholds, result and error types, metric labels, persisted status/signature mappings, public intents, and buffer-close policy unchanged. Larger naming and persistence cleanup remains separate work. Validation: - All six transaction-executor unit tests passed. - Clippy with warnings denied passed for all committor-service targets and the test_intent_executor integration target. - Integration test targets compiled successfully. - Nightly formatting checks passed in both workspaces; git diff --check passed. - The devnet test_commit_id_actions_cpi_limit_errors_recovery test compiled but failed during fixture funding, before exercising the executor. Local RPC connectivity and the initial airdrop were unavailable in this session. - No performance benchmark was run. The refactor adds no production RPC calls or transactions; broader runtime coverage remains for CI. --- .../src/intent_executor/mod.rs | 193 ++----- .../intent_executor/single_stage_executor.rs | 242 --------- .../intent_executor/transaction_executor.rs | 342 ++++++++++++ .../transaction_executor/tests.rs | 482 +++++++++++++++++ .../src/intent_executor/two_stage_executor.rs | 502 ------------------ .../src/intent_executor/utils.rs | 172 +----- .../tests/test_intent_executor.rs | 197 ++++--- 7 files changed, 965 insertions(+), 1165 deletions(-) delete mode 100644 magicblock-committor-service/src/intent_executor/single_stage_executor.rs create mode 100644 magicblock-committor-service/src/intent_executor/transaction_executor.rs create mode 100644 magicblock-committor-service/src/intent_executor/transaction_executor/tests.rs delete mode 100644 magicblock-committor-service/src/intent_executor/two_stage_executor.rs diff --git a/magicblock-committor-service/src/intent_executor/mod.rs b/magicblock-committor-service/src/intent_executor/mod.rs index 98eafac7f..d7e492259 100644 --- a/magicblock-committor-service/src/intent_executor/mod.rs +++ b/magicblock-committor-service/src/intent_executor/mod.rs @@ -1,9 +1,8 @@ pub mod error; pub mod intent_execution_client; pub(crate) mod intent_executor_factory; -pub mod single_stage_executor; pub mod task_info_fetcher; -pub mod two_stage_executor; +mod transaction_executor; pub mod utils; use std::{ @@ -26,7 +25,6 @@ use solana_keypair::Keypair; use solana_pubkey::Pubkey; use solana_signature::Signature; use solana_signer::Signer; -use tracing::trace; use crate::{ intent_executor::{ @@ -35,13 +33,8 @@ use crate::{ TransactionStrategyExecutionError, }, intent_execution_client::IntentExecutionClient, - single_stage_executor::SingleStageExecutor, task_info_fetcher::{CacheTaskInfoFetcher, ResetType, TaskInfoFetcher}, - two_stage_executor::TwoStageExecutor, - utils::{ - execute_with_timeout, handle_cpi_limit_error, CommitStage, - FinalizeStage, SingleStage, - }, + transaction_executor::TransactionExecutor, }, persist::{CommitStatus, CommitStatusSignatures, IntentPersister}, tasks::{ @@ -60,14 +53,12 @@ use crate::{ #[derive(Clone, Copy, Debug)] pub enum ExecutionOutput { - // TODO: with arrival of challenge window remove SingleStage - // Protocol requires 2 stage: Commit, Finalize - // SingleStage - optimization for timebeing + /// All tasks completed in one transaction. SingleStage(Signature), TwoStage { - /// Commit stage signature + /// Signature of the transaction containing the combined commits. commit_signature: Signature, - /// Finalize stage signature + /// Signature of the subsequent actions or undelegations. finalize_signature: Signature, }, } @@ -238,9 +229,7 @@ where } if all_committed_pubkeys.is_empty() { - // Build tasks for commit stage - // TODO (snawaz): it's actually MagicBaseIntent::BaseActions scenario, not Commit - // scenario, so the related code needs little bit of refactoring and proper renaming. + // Action-only intents use one transaction. let commit_tasks = TaskBuilderImpl::commit_tasks( &self.task_info_fetcher, &intent_bundle, @@ -256,9 +245,10 @@ where Some(intent_bundle.id), )?; return self - .single_stage_execution_flow( - intent_bundle, - strategy, + .execute_strategy( + intent_bundle.id, + &all_committed_pubkeys, + StrategyExecutionMode::SingleStage(strategy), execution_report, persister, ) @@ -285,168 +275,45 @@ where let uniqueness_nonce = requires_uniqueness_nonce(&commit_tasks) .then_some(intent_bundle.id); - // Build execution strategy - match TaskStrategist::build_execution_strategy( + let strategy = TaskStrategist::build_execution_strategy( commit_tasks, finalize_tasks, &self.authority.pubkey(), persister, uniqueness_nonce, - )? { - StrategyExecutionMode::SingleStage(strategy) => { - trace!("Single stage execution"); - self.single_stage_execution_flow( - intent_bundle, - strategy, - execution_report, - persister, - ) - .await - } - StrategyExecutionMode::TwoStage { - commit_stage, - finalize_stage, - } => { - trace!("Two stage execution"); - self.two_stage_execution_flow( - &all_committed_pubkeys, - commit_stage, - finalize_stage, - execution_report, - persister, - intent_bundle.id, - ) - .await - } - } - } - - fn time_left(&self) -> Option { - self.actions_timeout.checked_sub(self.started_at.elapsed()) - } - - /// Starting execution from single stage - pub async fn single_stage_execution_flow( - &mut self, - base_intent: ScheduledIntentBundle, - transaction_strategy: TransactionStrategy, - execution_report: &mut IntentExecutionReport, - persister: &Option

, - ) -> IntentExecutorResult { - let committed_pubkeys = base_intent.get_all_committed_pubkeys(); - - let mut single_stage_executor = SingleStageExecutor::new( - self.authority.insecure_clone(), - self.intent_client.clone(), - self.task_info_fetcher.clone(), - transaction_strategy, - self.actions_callback_executor.clone(), + )?; + self.execute_strategy( + intent_bundle.id, + &all_committed_pubkeys, + strategy, execution_report, - base_intent.id, - ); - let res = execute_with_timeout( - self.time_left(), - SingleStage { - inner: &mut single_stage_executor, - transaction_preparator: &self.transaction_preparator, - committed_pubkeys: &committed_pubkeys, - }, persister, ) - .await; - - // Here we continue only IF the error is a limit-type execution error - // We can recover that Error by splitting execution - // in 2 stages - commit & finalize - // Otherwise we return error - let execution_err = match res { - Err(IntentExecutorError::FailedToFinalizeError { - err, - commit_signature: _, - finalize_signature: _, - }) if !committed_pubkeys.is_empty() - && single_stage_executor.has_tasks_after_commit() - && err.is_recoverable_by_two_stage() => - { - err - } - res => { - let signature = res.as_ref().ok().copied(); - single_stage_executor - .execute_callbacks(signature, res.as_ref().map(|_| ())); - let transaction_strategy = - single_stage_executor.consume_strategy(); - execution_report.dispose(transaction_strategy); - return res.map(ExecutionOutput::SingleStage); - } - }; - - // With actions, we can't predict num of CPIs - // If we get here we will try to switch from Single stage to Two Stage commit - // Note that this not necessarily will pass at the end due to the same reason - let strategy = single_stage_executor.consume_strategy(); - let (commit_strategy, finalize_strategy, cleanup) = - handle_cpi_limit_error(&self.authority.pubkey(), strategy); - execution_report.dispose(cleanup); - execution_report.add_patched_error(execution_err); - - self.two_stage_execution_flow( - &committed_pubkeys, - commit_strategy, - finalize_strategy, - execution_report, - persister, - base_intent.id, - ) .await } - pub async fn two_stage_execution_flow( + fn time_left(&self) -> Option { + self.actions_timeout.checked_sub(self.started_at.elapsed()) + } + + /// Executes the selected transaction plan, including recovery and callbacks. + pub async fn execute_strategy( &mut self, + intent_id: u64, committed_pubkeys: &[Pubkey], - commit_strategy: TransactionStrategy, - finalize_strategy: TransactionStrategy, + strategy: StrategyExecutionMode, execution_report: &mut IntentExecutionReport, persister: &Option

, - intent_id: u64, ) -> IntentExecutorResult { - let mut executor = TwoStageExecutor::new( - self.authority.insecure_clone(), - commit_strategy, - finalize_strategy, - self.intent_client.clone(), - self.actions_callback_executor.clone(), + TransactionExecutor::new( + self, execution_report, intent_id, - ); - - let commit_signature = execute_with_timeout( - self.time_left(), - CommitStage { - inner: &mut executor, - transaction_preparator: &self.transaction_preparator, - task_info_fetcher: &self.task_info_fetcher, - committed_pubkeys, - }, - persister, - ) - .await?; - - let mut finalize_executor = executor.done(commit_signature); - let finalize_signature = execute_with_timeout( - self.time_left(), - FinalizeStage { - inner: &mut finalize_executor, - transaction_preparator: &self.transaction_preparator, - }, - persister, + committed_pubkeys, + strategy, ) - .await?; - - Ok(ExecutionOutput::TwoStage { - commit_signature, - finalize_signature, - }) + .execute(persister) + .await } /// Flushes result into presistor diff --git a/magicblock-committor-service/src/intent_executor/single_stage_executor.rs b/magicblock-committor-service/src/intent_executor/single_stage_executor.rs deleted file mode 100644 index 634945045..000000000 --- a/magicblock-committor-service/src/intent_executor/single_stage_executor.rs +++ /dev/null @@ -1,242 +0,0 @@ -use std::{ops::ControlFlow, sync::Arc}; - -use magicblock_core::traits::{ActionError, ActionsCallbackScheduler}; -use solana_keypair::Keypair; -use solana_pubkey::Pubkey; -use solana_signature::Signature; -use solana_signer::Signer; -use tracing::{error, instrument}; - -use crate::{ - intent_executor::{ - error::{ - IntentExecutorError, IntentExecutorResult, - TransactionStrategyExecutionError, - }, - intent_execution_client::IntentExecutionClient, - task_info_fetcher::{CacheTaskInfoFetcher, TaskInfoFetcher}, - utils::{ - handle_actions_result, handle_commit_id_error, - handle_undelegation_error, prepare_and_execute_strategy, - }, - IntentExecutionReport, - }, - persist::IntentPersister, - tasks::{task_strategist::TransactionStrategy, BaseTaskImpl}, - transaction_preparator::TransactionPreparator, -}; - -pub struct SingleStageExecutor<'a, F, A> { - current_attempt: u8, - intent_id: u64, - execution_report: &'a mut IntentExecutionReport, - - authority: Keypair, - intent_client: IntentExecutionClient, - task_info_fetcher: Arc>, - callback_scheduler: A, - transaction_strategy: TransactionStrategy, -} - -impl<'a, F, A> SingleStageExecutor<'a, F, A> -where - F: TaskInfoFetcher, - A: ActionsCallbackScheduler, -{ - pub fn new( - authority: Keypair, - intent_client: IntentExecutionClient, - task_info_fetcher: Arc>, - transaction_strategy: TransactionStrategy, - callback_scheduler: A, - execution_report: &'a mut IntentExecutionReport, - intent_id: u64, - ) -> Self { - Self { - current_attempt: 0, - intent_id, - authority, - intent_client, - task_info_fetcher, - transaction_strategy, - callback_scheduler, - execution_report, - } - } - - #[instrument( - skip(self, committed_pubkeys, transaction_preparator, persister), - fields(stage = "single_stage") - )] - pub async fn execute( - &mut self, - committed_pubkeys: &[Pubkey], - transaction_preparator: &T, - persister: &Option

, - ) -> IntentExecutorResult - where - T: TransactionPreparator, - P: IntentPersister, - { - const RECURSION_CEILING: u8 = 10; - - let result = loop { - self.current_attempt += 1; - - // Prepare & execute message - let execution_result = prepare_and_execute_strategy( - &self.intent_client, - &self.authority, - transaction_preparator, - &mut self.transaction_strategy, - persister, - ) - .await - .map_err(IntentExecutorError::FailedFinalizePreparationError)?; - - // Process error: Ok - return, Err - handle further - let execution_err = match execution_result { - // break with result, strategy that was executed at this point has to be returned for cleanup - Ok(value) => { - break Ok(value); - } - Err(err) => err, - }; - - // Attempt patching - let flow = self - .patch_strategy(&execution_err, committed_pubkeys) - .await?; - let cleanup = match flow { - ControlFlow::Continue(cleanup) => cleanup, - ControlFlow::Break(()) => { - break Err(execution_err); - } - }; - self.intent_client.invalidate_cached_blockhash().await; - self.execution_report.dispose(cleanup); - - if self.current_attempt >= RECURSION_CEILING { - error!( - attempt = self.current_attempt, - ceiling = RECURSION_CEILING, - error = ?execution_err, - "Recursion ceiling exceeded" - ); - break Err(execution_err); - } else { - self.execution_report.add_patched_error(execution_err); - } - }; - - result.map_err(|err| { - IntentExecutorError::from_finalize_execution_error( - err, - // TODO(edwin): shall one stage have same signature for commit & finalize - None, - ) - }) - } - - pub fn has_callbacks(&self) -> bool { - self.transaction_strategy.has_actions_callbacks() - } - - pub fn execute_callbacks( - &mut self, - signature: Option, - result: Result<(), impl Into>, - ) { - let junk_strategy = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.transaction_strategy, - signature, - result.map_err(|err| err.into()), - ); - self.execution_report.dispose(junk_strategy); - } - - pub fn consume_strategy(self) -> TransactionStrategy { - self.transaction_strategy - } - - pub(super) fn has_tasks_after_commit(&self) -> bool { - self.transaction_strategy - .optimized_tasks - .iter() - .rposition(|task| matches!(task, BaseTaskImpl::CommitFinalize(_))) - .is_some_and(|index| { - index + 1 < self.transaction_strategy.optimized_tasks.len() - }) - } - - /// Patch the current `transaction_strategy` in response to a recoverable - /// [`TransactionStrategyExecutionError`], optionally preparing cleanup data - /// to be applied after a retry. - /// - /// [`TransactionStrategyExecutionError`], returning either: - /// - `Continue(to_cleanup)` when a retry should be attempted with cleanup metadata, or - /// - `Break(())` when this stage cannot be recovered here. - pub async fn patch_strategy( - &mut self, - err: &TransactionStrategyExecutionError, - committed_pubkeys: &[Pubkey], - ) -> IntentExecutorResult> { - if committed_pubkeys.is_empty() { - // No patching is applicable if intent doesn't commit accounts - return Ok(ControlFlow::Break(())); - } - - match err { - TransactionStrategyExecutionError::ActionsError(err, signature) => { - // Here we patch strategy for it to be retried in next iteration - // & we also record data that has to be cleaned up after patch - let action_error = Err(ActionError::ActionsError(err.clone(), *signature)); - let to_cleanup = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.transaction_strategy, - *signature, - action_error, - ); - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::CommitIDError(_, _) => { - // Here we patch strategy for it to be retried in next iteration - // & we also record data that has to be cleaned up after patch - let to_cleanup = handle_commit_id_error( - &self.authority.pubkey(), - &self.task_info_fetcher, - committed_pubkeys, - &mut self.transaction_strategy, - self.intent_id, - ) - .await?; - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::UndelegationError(_, _) => { - // Here we patch strategy for it to be retried in next iteration - // & we also record data that has to be cleaned up after patch - let to_cleanup = handle_undelegation_error( - &self.authority.pubkey(), - &mut self.transaction_strategy - ); - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::CpiLimitError(_, _) - | TransactionStrategyExecutionError::LoadedAccountsDataSizeExceeded(_, _) - | TransactionStrategyExecutionError::TransactionTooLargeError(_) => { - // Can't be handled in scope of single stage execution - // We signal flow break - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::InternalError(_) => { - // Error that we can't handle - break with cleanup data - Ok(ControlFlow::Break(())) - } - } - } -} diff --git a/magicblock-committor-service/src/intent_executor/transaction_executor.rs b/magicblock-committor-service/src/intent_executor/transaction_executor.rs new file mode 100644 index 000000000..badcb0b70 --- /dev/null +++ b/magicblock-committor-service/src/intent_executor/transaction_executor.rs @@ -0,0 +1,342 @@ +use std::{mem, ops::ControlFlow}; + +use magicblock_core::traits::{ActionError, ActionsCallbackScheduler}; +use solana_pubkey::Pubkey; +use solana_signature::Signature; +use solana_signer::Signer; +use tokio::time::timeout; +use tracing::{error, info}; + +use super::{ + error::{ + IntentExecutorError, IntentExecutorResult, + TransactionStrategyExecutionError, + }, + task_info_fetcher::TaskInfoFetcher, + utils::{ + handle_actions_result, handle_commit_id_error, handle_cpi_limit_error, + handle_undelegation_error, prepare_and_execute_strategy, + }, + ExecutionOutput, IntentExecutionReport, IntentExecutorImpl, +}; +use crate::{ + persist::IntentPersister, + tasks::{ + task_strategist::{StrategyExecutionMode, TransactionStrategy}, + BaseTaskImpl, + }, + transaction_preparator::TransactionPreparator, +}; + +#[cfg(test)] +mod tests; + +/// Runs the selected transaction and any remaining actions or undelegations. +/// A confirmed commit is never included in subsequent recovery attempts. +pub(super) struct TransactionExecutor<'a, T, F, A> { + executor: &'a IntentExecutorImpl, + report: &'a mut IntentExecutionReport, + intent_id: u64, + committed_pubkeys: &'a [Pubkey], + current: TransactionStrategy, + pending: Option, + commit_signature: Option, + current_attempt: u8, +} + +impl<'a, T, F, A> TransactionExecutor<'a, T, F, A> +where + T: TransactionPreparator, + F: TaskInfoFetcher, + A: ActionsCallbackScheduler, +{ + pub(super) fn new( + executor: &'a IntentExecutorImpl, + report: &'a mut IntentExecutionReport, + intent_id: u64, + committed_pubkeys: &'a [Pubkey], + strategy: StrategyExecutionMode, + ) -> Self { + let (current, pending) = match strategy { + StrategyExecutionMode::SingleStage(strategy) => (strategy, None), + StrategyExecutionMode::TwoStage { + commit_stage, + finalize_stage, + } => (commit_stage, Some(finalize_stage)), + }; + Self { + executor, + report, + intent_id, + committed_pubkeys, + current, + pending, + commit_signature: None, + current_attempt: 0, + } + } + + pub(super) async fn execute( + mut self, + persister: &Option

, + ) -> IntentExecutorResult { + loop { + let result = self.execute_with_timeout(persister).await; + let result = match result { + Err(IntentExecutorError::FailedToFinalizeError { + err, .. + }) if self.is_combined() + && !self.committed_pubkeys.is_empty() + && self.has_tasks_after_commit() + && err.is_recoverable_by_two_stage() => + { + let (current, pending, cleanup) = handle_cpi_limit_error( + &self.executor.authority.pubkey(), + mem::take(&mut self.current), + ); + self.current = current; + self.pending = Some(pending); + self.current_attempt = 0; + self.report.dispose(cleanup); + self.report.add_patched_error(err); + continue; + } + result => result, + }; + + // Combined execution reports all terminal errors to callbacks. + // Split execution historically reports transaction outcomes, but + // leaves preparation and nonce-fetch failures to the caller. + if self.is_combined() + || matches!( + &result, + Ok(_) + | Err(IntentExecutorError::FailedToCommitError { .. }) + | Err( + IntentExecutorError::FailedToFinalizeError { .. } + ) + ) + { + self.execute_callbacks( + result.as_ref().ok().copied(), + result.as_ref().map(|_| ()).map_err(ActionError::from), + ); + } + self.report.dispose(mem::take(&mut self.current)); + + match result { + Ok(signature) => { + if let Some(pending) = self.pending.take() { + self.current = pending; + self.commit_signature = Some(signature); + self.current_attempt = 0; + } else { + return Ok(match self.commit_signature { + Some(commit_signature) => { + ExecutionOutput::TwoStage { + commit_signature, + finalize_signature: signature, + } + } + None => ExecutionOutput::SingleStage(signature), + }); + } + } + Err(err) => { + if let Some(pending) = self.pending.take() { + self.report.dispose(pending); + } + return Err(err); + } + } + } + } + + async fn execute_with_timeout( + &mut self, + persister: &Option

, + ) -> IntentExecutorResult { + let has_callbacks = self.current.has_actions_callbacks() + || self + .pending + .as_ref() + .is_some_and(TransactionStrategy::has_actions_callbacks); + if has_callbacks { + if let Some(time_left) = self.executor.time_left() { + if let Ok(result) = + timeout(time_left, self.execute_current(persister)).await + { + return result; + } + } + // A transaction may have landed before confirmation timed out. + // The callback recipient handles that race via TimeoutError. + info!("Intent execution timed out, cleaning up actions"); + self.execute_callbacks(None, Err(ActionError::TimeoutError)); + } + self.execute_current(persister).await + } + + #[tracing::instrument(skip_all, fields(stage = + if self.pending.is_some() { "commit" } + else if self.commit_signature.is_some() { "finalize" } + else { "single_stage" } + ))] + async fn execute_current( + &mut self, + persister: &Option

, + ) -> IntentExecutorResult { + const ATTEMPT_LIMIT: u8 = 10; + let result = loop { + self.current_attempt += 1; + let result = prepare_and_execute_strategy( + &self.executor.intent_client, + &self.executor.authority, + &self.executor.transaction_preparator, + &mut self.current, + persister, + ) + .await + .map_err(|err| { + if self.pending.is_some() { + IntentExecutorError::FailedCommitPreparationError(err) + } else { + IntentExecutorError::FailedFinalizePreparationError(err) + } + })?; + let err = match result { + Ok(signature) => break Ok(signature), + Err(err) => err, + }; + + match self.patch_strategy(&err).await? { + ControlFlow::Continue(cleanup) => self.report.dispose(cleanup), + ControlFlow::Break(()) => break Err(err), + } + self.executor + .intent_client + .invalidate_cached_blockhash() + .await; + + // Failed follow-up actions can leave no work after a landed commit. + // An undelegation failure must still be returned if removing its + // tasks empties the transaction. + if let Some(signature) = self.commit_signature { + if self.current.optimized_tasks.is_empty() { + if matches!( + err, + TransactionStrategyExecutionError::ActionsError(..) + ) { + self.report.add_patched_error(err); + break Ok(signature); + } + break Err(err); + } + } + if self.current_attempt >= ATTEMPT_LIMIT { + error!(attempt = self.current_attempt, error = ?err, "Transaction recovery attempt limit reached"); + break Err(err); + } + self.report.add_patched_error(err); + }; + result.map_err(|err| { + if self.pending.is_some() { + IntentExecutorError::from_commit_execution_error(err) + } else { + IntentExecutorError::from_finalize_execution_error( + err, + self.commit_signature, + ) + } + }) + } + + async fn patch_strategy( + &mut self, + err: &TransactionStrategyExecutionError, + ) -> IntentExecutorResult> { + if self.is_combined() && self.committed_pubkeys.is_empty() { + return Ok(ControlFlow::Break(())); + } + let authority = self.executor.authority.pubkey(); + let cleanup = match err { + TransactionStrategyExecutionError::CommitIDError(..) + if self.commit_signature.is_none() => + { + let cleanup = handle_commit_id_error( + &authority, + &self.executor.task_info_fetcher, + self.committed_pubkeys, + &mut self.current, + self.intent_id, + ) + .await?; + if let Some(pending) = &mut self.pending { + // Re-delegation can reset the nonce to 1. Both transactions + // then need the intent's uniqueness instruction. + if pending.uniqueness_nonce.is_none() { + pending.uniqueness_nonce = + self.current.uniqueness_nonce; + } + } + cleanup + } + TransactionStrategyExecutionError::ActionsError(err, signature) => { + handle_actions_result( + &authority, + &self.executor.actions_callback_executor, + self.report, + &mut self.current, + *signature, + Err(ActionError::ActionsError(err.clone(), *signature)), + ) + } + TransactionStrategyExecutionError::UndelegationError(..) + if self.pending.is_none() => + { + handle_undelegation_error(&authority, &mut self.current) + } + _ => return Ok(ControlFlow::Break(())), + }; + Ok(ControlFlow::Continue(cleanup)) + } + + fn execute_callbacks( + &mut self, + signature: Option, + result: Result<(), ActionError>, + ) { + let cleanup = handle_actions_result( + &self.executor.authority.pubkey(), + &self.executor.actions_callback_executor, + self.report, + &mut self.current, + signature, + result.clone(), + ); + self.report.dispose(cleanup); + if let (Err(_), Some(pending)) = (&result, &mut self.pending) { + let cleanup = handle_actions_result( + &self.executor.authority.pubkey(), + &self.executor.actions_callback_executor, + self.report, + pending, + signature, + result, + ); + self.report.dispose(cleanup); + } + } + + fn is_combined(&self) -> bool { + self.pending.is_none() && self.commit_signature.is_none() + } + + fn has_tasks_after_commit(&self) -> bool { + self.current + .optimized_tasks + .iter() + .rposition(|task| matches!(task, BaseTaskImpl::CommitFinalize(_))) + .is_some_and(|index| index + 1 < self.current.optimized_tasks.len()) + } +} diff --git a/magicblock-committor-service/src/intent_executor/transaction_executor/tests.rs b/magicblock-committor-service/src/intent_executor/transaction_executor/tests.rs new file mode 100644 index 000000000..c55d01018 --- /dev/null +++ b/magicblock-committor-service/src/intent_executor/transaction_executor/tests.rs @@ -0,0 +1,482 @@ +use std::{ + collections::VecDeque, + sync::{Arc, Mutex}, + time::Duration, +}; + +use async_trait::async_trait; +use magicblock_core::{ + intent::{ + types::CommittedAccount, BaseAction, BaseActionCallback, ProgramArgs, + }, + traits::{ActionResult, CallbackScheduleError}, +}; +use magicblock_rpc_client::MagicblockRpcClient; +use solana_account::Account; +use solana_keypair::Keypair; +use solana_message::{Message, VersionedMessage}; +use solana_rpc_client::{ + mock_sender::MocksMap, nonblocking::rpc_client::RpcClient, +}; +use solana_rpc_client_api::request::RpcRequest; + +use super::*; +use crate::{ + intent_executor::{ + task_info_fetcher::{CacheTaskInfoFetcher, RpcTaskInfoFetcher}, + IntentExecutionResult, + }, + persist::IntentPersisterImpl, + tasks::{ + utils::{ + create_action_tasks, create_commit_finalize_task, TransactionUtils, + }, + UndelegateTask, + }, + transaction_preparator::{ + delivery_preparator::{BufferExecutionError, DeliveryPreparatorResult}, + error::{PreparatorResult, TransactionPreparatorError}, + }, + transactions::PreparedMessage, +}; + +const ACCOUNT: Pubkey = Pubkey::new_from_array([7; 32]); +const RESERVED_ALT_KEY: Pubkey = Pubkey::new_from_array([9; 32]); + +enum Preparation { + Ready, + Fail, + Wait, +} + +#[derive(Clone, Default)] +struct RecordingPreparator { + steps: Arc>>, + prepared: Arc>>>, +} + +#[async_trait] +impl TransactionPreparator for RecordingPreparator { + async fn prepare_for_strategy( + &self, + authority: &Keypair, + strategy: &mut TransactionStrategy, + _: &Option

, + ) -> PreparatorResult { + self.prepared + .lock() + .unwrap() + .push(strategy.optimized_tasks.clone()); + // Represent a reservation made before preparation succeeds or fails. + strategy.lookup_tables_keys = vec![RESERVED_ALT_KEY]; + let step = self + .steps + .lock() + .unwrap() + .pop_front() + .unwrap_or(Preparation::Ready); + match step { + Preparation::Fail => { + Err(TransactionPreparatorError::FailedToFitError) + } + Preparation::Wait => futures_util::future::pending().await, + Preparation::Ready => { + let mut instructions = TransactionUtils::budget_instructions( + TransactionUtils::tasks_compute_units( + &strategy.optimized_tasks, + ), + 0, + TransactionUtils::tasks_accounts_size_budget( + &strategy.optimized_tasks, + ), + ) + .to_vec(); + instructions.extend(TransactionUtils::tasks_instructions( + &authority.pubkey(), + &strategy.optimized_tasks, + )); + // The RPC mock decodes legacy messages. Keep the same budget + // instruction offset used by production versioned messages. + Ok(PreparedMessage::Versioned(VersionedMessage::Legacy( + Message::new(&instructions, Some(&authority.pubkey())), + ))) + } + } + } + + async fn cleanup_for_strategy( + &self, + _: &Keypair, + _: &TransactionStrategy, + _: bool, + ) -> DeliveryPreparatorResult<(), BufferExecutionError> { + Ok(()) + } +} + +type CallbackResults = Vec<(Option, ActionResult)>; + +#[derive(Clone, Default)] +struct Callbacks(Arc>); + +impl ActionsCallbackScheduler for Callbacks { + fn schedule( + &self, + callbacks: Vec, + signature: Option, + result: ActionResult, + ) -> Vec> { + callbacks + .into_iter() + .map(|_| { + self.0.lock().unwrap().push((signature, result.clone())); + Ok(Signature::new_unique()) + }) + .collect() + } +} + +type Executor = + IntentExecutorImpl; + +fn executor( + steps: Vec, + errors: &[Option<&str>], + timeout: Duration, +) -> Executor { + let mut mocks = MocksMap::default(); + for error in errors { + let err = error.unwrap_or("null"); + let status = error.map_or_else( + || r#"{"Ok":null}"#.to_string(), + |err| format!(r#"{{"Err":{err}}}"#), + ); + mocks.insert( + RpcRequest::GetSignatureStatuses, + format!( + r#"{{ + "context": {{"slot": 1}}, + "value": [{{"slot": 1, "confirmations": null, "err": {err}, + "status": {status}, "confirmationStatus": "finalized"}}] + }}"# + ) + .parse() + .unwrap(), + ); + } + let rpc = MagicblockRpcClient::from(RpcClient::new_mock_with_mocks_map( + "succeeds", mocks, + )); + magicblock_program::validator::generate_validator_authority_if_needed(); + IntentExecutorImpl::new( + rpc.clone(), + RecordingPreparator { + steps: Arc::new(Mutex::new(steps.into())), + ..Default::default() + }, + Arc::new(CacheTaskInfoFetcher::new(RpcTaskInfoFetcher::new(rpc))), + Callbacks::default(), + timeout, + ) +} + +fn commit() -> BaseTaskImpl { + create_commit_finalize_task( + 2, + true, + CommittedAccount { + pubkey: ACCOUNT, + account: Account::default(), + remote_slot: 0, + }, + None, + ) + .into() +} + +fn undelegate() -> BaseTaskImpl { + BaseTaskImpl::Undelegate(UndelegateTask { + delegated_account: ACCOUNT, + owner_program: Pubkey::new_unique(), + rent_reimbursement: Pubkey::new_unique(), + include_undelegation_request: false, + }) +} + +fn action() -> BaseTaskImpl { + create_action_tasks(&[BaseAction { + id: 0, + destination_program: Pubkey::new_unique(), + source_program: None, + escrow_authority: Pubkey::new_unique(), + account_metas_per_program: vec![], + data_per_program: ProgramArgs { + data: vec![], + escrow_index: 0, + }, + compute_units: 10_000, + callback: Some(BaseActionCallback { + destination_program: Pubkey::new_unique(), + discriminator: vec![], + payload: vec![], + compute_units: 10_000, + account_metas_per_program: vec![], + }), + }]) + .next() + .unwrap() +} + +fn strategy(tasks: Vec) -> TransactionStrategy { + TransactionStrategy { + optimized_tasks: tasks, + ..Default::default() + } +} + +fn split(follow_up: Vec) -> StrategyExecutionMode { + StrategyExecutionMode::TwoStage { + commit_stage: strategy(vec![commit()]), + finalize_stage: strategy(follow_up), + } +} + +fn instruction_error(error: &str) -> String { + format!( + r#"{{"InstructionError":[{},{error}]}}"#, + TransactionUtils::COMPUTE_BUDGET_INSTRUCTION_COUNT + ) +} + +#[tokio::test(start_paused = true)] +async fn pending_callbacks_time_out_while_first_transaction_is_preparing() { + let executor = + executor(vec![Preparation::Wait], &[], Duration::from_millis(10)); + let mut report = IntentExecutionReport::default(); + let output = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + split(vec![undelegate(), action()]), + ) + .execute(&None::) + .await + .unwrap(); + assert!(matches!(output, ExecutionOutput::TwoStage { .. })); + let calls = executor.actions_callback_executor.0.lock().unwrap(); + assert_eq!(calls.len(), 1); + assert!(matches!(calls[0], (None, Err(ActionError::TimeoutError)))); + let prepared = executor.transaction_preparator.prepared.lock().unwrap(); + assert_eq!( + prepared.len(), + 3, + "cancelled preparation, commit retry, then undelegation" + ); + assert!(matches!( + prepared[2].as_slice(), + [BaseTaskImpl::Undelegate(_)] + )); + assert!(report.junk().iter().any(|strategy| strategy + .lookup_tables_keys + .contains(&RESERVED_ALT_KEY))); +} + +#[tokio::test(start_paused = true)] +async fn timeout_keeps_attempts_already_spent_on_current_transaction() { + let executor = + executor(vec![Preparation::Wait], &[], Duration::from_millis(10)); + let mut report = IntentExecutionReport::default(); + let mut runner = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + StrategyExecutionMode::SingleStage(strategy(vec![commit(), action()])), + ); + runner.current_attempt = 8; + runner + .execute_with_timeout(&None::) + .await + .unwrap(); + assert_eq!(runner.current_attempt, 10); + assert_eq!( + executor.actions_callback_executor.0.lock().unwrap().len(), + 1 + ); +} + +#[tokio::test] +async fn preparation_failure_surrenders_current_and_pending_resources() { + let executor = + executor(vec![Preparation::Fail], &[], Duration::from_secs(60)); + let mut report = IntentExecutionReport::default(); + let result = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + split(vec![undelegate(), action()]), + ) + .execute(&None::) + .await; + assert!(matches!( + result, + Err(IntentExecutorError::FailedCommitPreparationError(_)) + )); + assert_eq!(report.junk().len(), 2); + assert_eq!(report.junk()[0].lookup_tables_keys, vec![RESERVED_ALT_KEY]); + assert_eq!(report.junk()[1].optimized_tasks.len(), 2); + assert!(executor + .actions_callback_executor + .0 + .lock() + .unwrap() + .is_empty()); +} + +#[tokio::test] +async fn nonce_fetch_failure_surrenders_both_strategies() { + let nonce_error = instruction_error(&format!( + r#"{{"Custom":{}}}"#, + dlp_api::error::DlpError::NonceOutOfOrder as u32 + )); + let executor = + executor(vec![], &[Some(&nonce_error)], Duration::from_secs(60)); + let mut report = IntentExecutionReport::default(); + // The RPC mock has no delegation record, so nonce recovery fails after + // preparing the first transaction and receiving its on-chain nonce error. + let result = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + split(vec![undelegate()]), + ) + .execute(&None::) + .await; + assert!( + matches!(result, Err(IntentExecutorError::TaskBuilderError(_))), + "{result:?}" + ); + assert_eq!( + executor + .transaction_preparator + .prepared + .lock() + .unwrap() + .len(), + 1 + ); + assert_eq!(report.junk().len(), 2); + assert_eq!(report.junk()[0].lookup_tables_keys, vec![RESERVED_ALT_KEY]); + assert!(matches!( + report.junk()[1].optimized_tasks.as_slice(), + [BaseTaskImpl::Undelegate(_)] + )); +} + +#[tokio::test] +async fn follow_up_preparation_failure_does_not_retry_confirmed_commit() { + let executor = executor( + vec![Preparation::Ready, Preparation::Fail], + &[], + Duration::from_secs(60), + ); + let mut report = IntentExecutionReport::default(); + let result = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + split(vec![undelegate()]), + ) + .execute(&None::) + .await; + assert!(matches!( + result, + Err(IntentExecutorError::FailedFinalizePreparationError(_)) + )); + let result = IntentExecutionResult { + inner: result, + patched_errors: vec![], + callbacks_report: vec![], + }; + assert!(!result.is_retriable(true)); + let prepared = executor.transaction_preparator.prepared.lock().unwrap(); + assert_eq!(prepared.len(), 2); + assert!(matches!( + prepared[0].as_slice(), + [BaseTaskImpl::CommitFinalize(_)] + )); + assert!(matches!( + prepared[1].as_slice(), + [BaseTaskImpl::Undelegate(_)] + )); + assert_eq!( + report + .junk() + .iter() + .filter(|strategy| strategy + .lookup_tables_keys + .contains(&RESERVED_ALT_KEY)) + .count(), + 2 + ); +} + +#[tokio::test] +async fn empty_follow_up_preserves_action_and_undelegation_outcomes() { + for failed_action in [true, false] { + let error = instruction_error(r#"{"Custom":1}"#); + let executor = + executor(vec![], &[None, Some(&error)], Duration::from_secs(60)); + let mut report = IntentExecutionReport::default(); + let follow_up = if failed_action { + action() + } else { + undelegate() + }; + let result = TransactionExecutor::new( + &executor, + &mut report, + 42, + &[ACCOUNT], + split(vec![follow_up]), + ) + .execute(&None::) + .await; + assert_eq!( + executor + .transaction_preparator + .prepared + .lock() + .unwrap() + .len(), + 2, + "must not prepare an empty retry" + ); + if failed_action { + assert!( + matches!(result, Ok(ExecutionOutput::TwoStage { commit_signature, finalize_signature }) if commit_signature == finalize_signature), + "{result:?}" + ); + let calls = executor.actions_callback_executor.0.lock().unwrap(); + assert_eq!(calls.len(), 1); + assert!(matches!(calls[0].1, Err(ActionError::ActionsError(..)))); + assert_eq!(report.patched_errors().len(), 1); + } else { + assert!( + matches!( + result, + Err(IntentExecutorError::FailedToFinalizeError { + commit_signature: Some(_), + .. + }) + ), + "{result:?}" + ); + assert!(!result.unwrap_err().is_transient()); + } + } +} diff --git a/magicblock-committor-service/src/intent_executor/two_stage_executor.rs b/magicblock-committor-service/src/intent_executor/two_stage_executor.rs deleted file mode 100644 index b6db607dd..000000000 --- a/magicblock-committor-service/src/intent_executor/two_stage_executor.rs +++ /dev/null @@ -1,502 +0,0 @@ -use std::{mem, ops::ControlFlow}; - -use magicblock_core::traits::{ActionError, ActionsCallbackScheduler}; -use solana_keypair::Keypair; -use solana_pubkey::Pubkey; -use solana_signature::Signature; -use solana_signer::Signer; -use tracing::{error, instrument, warn}; - -use crate::{ - intent_executor::{ - error::{ - IntentExecutorError, IntentExecutorResult, - TransactionStrategyExecutionError, - }, - intent_execution_client::IntentExecutionClient, - task_info_fetcher::{CacheTaskInfoFetcher, TaskInfoFetcher}, - two_stage_executor::sealed::Sealed, - utils::{ - handle_actions_result, handle_commit_id_error, - handle_undelegation_error, prepare_and_execute_strategy, - }, - IntentExecutionReport, - }, - persist::IntentPersister, - tasks::task_strategist::TransactionStrategy, - transaction_preparator::TransactionPreparator, -}; - -pub struct Initialized { - /// Commit stage strategy - commit_strategy: TransactionStrategy, - /// Finalize stage strategy - finalize_strategy: TransactionStrategy, - - current_attempt: u8, -} - -pub struct Committed { - /// Signature of commit stage - commit_signature: Signature, - /// Finalize stage strategy - finalize_strategy: TransactionStrategy, - - current_attempt: u8, -} - -pub struct TwoStageExecutor<'a, A, S: Sealed> { - state: S, - intent_id: u64, - authority: Keypair, - intent_client: IntentExecutionClient, - callback_scheduler: A, - execution_report: &'a mut IntentExecutionReport, -} - -impl<'a, A> TwoStageExecutor<'a, A, Initialized> -where - A: ActionsCallbackScheduler, -{ - const RECURSION_CEILING: u8 = 10; - - pub fn new( - authority: Keypair, - commit_strategy: TransactionStrategy, - finalize_strategy: TransactionStrategy, - intent_client: IntentExecutionClient, - callback_scheduler: A, - execution_report: &'a mut IntentExecutionReport, - intent_id: u64, - ) -> Self { - Self { - intent_id, - authority, - intent_client, - execution_report, - callback_scheduler, - state: Initialized { - commit_strategy, - finalize_strategy, - current_attempt: 0, - }, - } - } - - #[instrument( - skip( - self, - committed_pubkeys, - transaction_preparator, - task_info_fetcher, - persister - ), - fields(stage = "commit") - )] - pub async fn commit( - &mut self, - committed_pubkeys: &[Pubkey], - transaction_preparator: &T, - task_info_fetcher: &CacheTaskInfoFetcher, - persister: &Option

, - ) -> IntentExecutorResult - where - T: TransactionPreparator, - F: TaskInfoFetcher, - P: IntentPersister, - { - let commit_result = loop { - self.state.current_attempt += 1; - - // Prepare & execute message - let execution_result = prepare_and_execute_strategy( - &self.intent_client, - &self.authority, - transaction_preparator, - &mut self.state.commit_strategy, - persister, - ) - .await - .map_err(IntentExecutorError::FailedCommitPreparationError) - .inspect_err(|_| { - // Preparation may have reserved ALTs or initialized - // buffers already - surrender strategies for cleanup - self.dispose_strategies(); - })?; - - let execution_err = match execution_result { - Ok(value) => break Ok(value), - Err(err) => err, - }; - - let flow = self - .patch_commit_strategy( - &execution_err, - task_info_fetcher, - committed_pubkeys, - ) - .await - .inspect_err(|_| { - self.dispose_strategies(); - })?; - let cleanup = match flow { - ControlFlow::Continue(value) => value, - ControlFlow::Break(()) => { - break Err(execution_err); - } - }; - self.intent_client.invalidate_cached_blockhash().await; - self.execution_report.dispose(cleanup); - - if self.state.current_attempt >= Self::RECURSION_CEILING { - error!("CRITICAL! Recursion ceiling reached"); - break Err(execution_err); - } else { - self.execution_report.add_patched_error(execution_err); - } - }; - - self.execute_callbacks( - commit_result.as_ref().ok().copied(), - commit_result.as_ref().map(|_| ()), - ); - self.execution_report - .dispose(mem::take(&mut self.state.commit_strategy)); - if commit_result.is_err() { - self.execution_report - .dispose(mem::take(&mut self.state.finalize_strategy)); - } - commit_result.map_err(|err| { - IntentExecutorError::from_commit_execution_error(err) - }) - } - - /// Patches Commit stage `transaction_strategy` in response to a recoverable - /// [`TransactionStrategyExecutionError`], optionally preparing cleanup data - /// to be applied after a retry. - /// - /// [`TransactionStrategyExecutionError`], returning either: - /// - `Continue(to_cleanup)` when a retry should be attempted with cleanup metadata, or - /// - `Break(())` when this stage cannot be recovered. - pub async fn patch_commit_strategy( - &mut self, - err: &TransactionStrategyExecutionError, - task_info_fetcher: &CacheTaskInfoFetcher, - committed_pubkeys: &[Pubkey], - ) -> IntentExecutorResult> - where - F: TaskInfoFetcher, - { - match err { - TransactionStrategyExecutionError::CommitIDError(_, _) => { - let to_cleanup = handle_commit_id_error( - &self.authority.pubkey(), - task_info_fetcher, - committed_pubkeys, - &mut self.state.commit_strategy, - self.intent_id, - ) - .await?; - // If recovery re-tagged the commit stage as a first commit, - // the finalize stage aliases the same way and must carry the - // uniqueness noop too. - if self.state.finalize_strategy.uniqueness_nonce.is_none() { - self.state.finalize_strategy.uniqueness_nonce = - self.state.commit_strategy.uniqueness_nonce; - } - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::ActionsError(err, signature) => { - // Intent bundles allow for actions to be in commit stage - let action_error = Err(ActionError::ActionsError(err.clone(), *signature)); - let to_cleanup = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.state.commit_strategy, - *signature, - action_error - ); - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::UndelegationError(_, _) => { - // Unexpected in Two Stage commit - // That would mean that Two Stage executes undelegation in commit phase - error!(error = ?err, "Unexpected error in two stage commit flow"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::CpiLimitError(_, _) - | TransactionStrategyExecutionError::LoadedAccountsDataSizeExceeded( - _, - _, - ) => { - // Can't be handled - error!(error = ?err, "Commit tasks exceeded execution limit"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::TransactionTooLargeError(_) => { - // Can't be handled but also shouldn't occur - error!(strategy = ?self.state.commit_strategy, error = ?err, "Commit tasks do not fit in tx"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::InternalError(_) => { - // Can't be handled - Ok(ControlFlow::Break(())) - } - } - } - - /// Surrenders both strategies to the execution report so the engine - /// cleanup releases partially prepared resources (ALT reservations, - /// buffer accounts) of failed attempts - fn dispose_strategies(&mut self) { - let commit_strategy = mem::take(&mut self.state.commit_strategy); - self.execution_report.dispose(commit_strategy); - let finalize_strategy = mem::take(&mut self.state.finalize_strategy); - self.execution_report.dispose(finalize_strategy); - } - - pub fn has_callbacks(&self) -> bool { - self.state.commit_strategy.has_actions_callbacks() - || self.state.finalize_strategy.has_actions_callbacks() - } - - /// On `Err`: removes actions from both commit and finalize strategies and - /// executes all their callbacks with the error. - /// On `Ok`: removes actions only from commit strategy and executes their - /// callbacks, preserving finalize-stage actions for the finalize phase. - pub fn execute_callbacks( - &mut self, - signature: Option, - result: Result<(), impl Into>, - ) { - let result = result.map(|_| ()).map_err(|err| err.into()); - let junk_strategy = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.state.commit_strategy, - signature, - result.clone(), - ); - self.execution_report.dispose(junk_strategy); - - if result.is_err() { - let junk_strategy = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.state.finalize_strategy, - signature, - result, - ); - self.execution_report.dispose(junk_strategy); - } - } - - /// Transitions to next executor state - pub fn done( - self, - commit_signature: Signature, - ) -> TwoStageExecutor<'a, A, Committed> { - TwoStageExecutor { - intent_id: self.intent_id, - authority: self.authority, - intent_client: self.intent_client, - callback_scheduler: self.callback_scheduler, - execution_report: self.execution_report, - state: Committed { - commit_signature, - finalize_strategy: self.state.finalize_strategy, - current_attempt: 0, - }, - } - } -} - -impl<'a, A> TwoStageExecutor<'a, A, Committed> -where - A: ActionsCallbackScheduler, -{ - const RECURSION_CEILING: u8 = 10; - - #[instrument( - skip(self, transaction_preparator, persister), - fields(stage = "finalize") - )] - pub async fn finalize( - &mut self, - transaction_preparator: &T, - persister: &Option

, - ) -> IntentExecutorResult - where - T: TransactionPreparator, - P: IntentPersister, - { - let finalize_result = loop { - self.state.current_attempt += 1; - - // Prepare & execute message - let execution_result = prepare_and_execute_strategy( - &self.intent_client, - &self.authority, - transaction_preparator, - &mut self.state.finalize_strategy, - persister, - ) - .await - .map_err(IntentExecutorError::FailedFinalizePreparationError) - .inspect_err(|_| { - // Preparation may have reserved ALTs or initialized - // buffers already - surrender strategy for cleanup - self.dispose_strategy(); - })?; - let execution_err = match execution_result { - Ok(value) => break Ok(value), - Err(err) => err, - }; - - let flow = self - .patch_finalize_strategy(&execution_err) - .await - .inspect_err(|_| { - self.dispose_strategy(); - })?; - let cleanup = match flow { - ControlFlow::Continue(cleanup) => cleanup, - ControlFlow::Break(()) => { - break Err(execution_err); - } - }; - self.intent_client.invalidate_cached_blockhash().await; - self.execution_report.dispose(cleanup); - - // Failed actions may be removed after a confirmed combined commit. - // If nothing remains, reuse that commit's signature. An empty - // strategy after an undelegation failure must retain the failure. - if self.state.finalize_strategy.optimized_tasks.is_empty() { - if matches!( - execution_err, - TransactionStrategyExecutionError::ActionsError(_, _) - ) { - self.execution_report.add_patched_error(execution_err); - break Ok(self.state.commit_signature); - } - break Err(execution_err); - } - - if self.state.current_attempt >= Self::RECURSION_CEILING { - error!("CRITICAL! Recursion ceiling reached"); - break Err(execution_err); - } else { - self.execution_report.add_patched_error(execution_err); - } - }; - - // Even if failed - dump finalize into junk - self.execute_callbacks( - finalize_result.as_ref().ok().copied(), - finalize_result.as_ref().map(|_| ()), - ); - self.execution_report - .dispose(mem::take(&mut self.state.finalize_strategy)); - finalize_result.map_err(|err| { - IntentExecutorError::from_finalize_execution_error( - err, - Some(self.state.commit_signature), - ) - }) - } - - /// Surrenders the finalize strategy to the execution report so the engine - /// cleanup releases partially prepared resources of failed attempts - fn dispose_strategy(&mut self) { - let finalize_strategy = mem::take(&mut self.state.finalize_strategy); - self.execution_report.dispose(finalize_strategy); - } - - pub fn has_callbacks(&self) -> bool { - self.state.finalize_strategy.has_actions_callbacks() - } - - /// Removes actions from finalize strategy - /// Executes callbacks - pub fn execute_callbacks( - &mut self, - signature: Option, - result: Result<(), impl Into>, - ) { - let junk_strategy = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.state.finalize_strategy, - signature, - result.map_err(|err| err.into()), - ); - self.execution_report.dispose(junk_strategy); - } - - /// Patches Finalize stage `transaction_strategy` in response to a recoverable - /// [`TransactionStrategyExecutionError`], optionally preparing cleanup data - /// to be applied after a retry. - /// - /// [`TransactionStrategyExecutionError`], returning either: - /// - `Continue(to_cleanup)` when a retry should be attempted with cleanup metadata, or - /// - `Break(())` when this stage cannot be recovered. - pub async fn patch_finalize_strategy( - &mut self, - err: &TransactionStrategyExecutionError, - ) -> IntentExecutorResult> { - match err { - TransactionStrategyExecutionError::CommitIDError(_, _) => { - // Unexpected error in Two Stage commit - error!(error = ?err, "Unexpected error in two stage finalize flow"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::ActionsError(err, signature) => { - // Here we patch strategy for it to be retried in next iteration - // & we also record data that has to be cleaned up after patch - let action_error = Err(ActionError::ActionsError(err.clone(), *signature)); - let to_cleanup = handle_actions_result( - &self.authority.pubkey(), - &self.callback_scheduler, - self.execution_report, - &mut self.state.finalize_strategy, - *signature, - action_error - ); - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::UndelegationError(_, _) => { - // Here we patch strategy for it to be retried in next iteration - // & we also record data that has to be cleaned up after patch - let to_cleanup = - handle_undelegation_error(&self.authority.pubkey(), &mut self.state.finalize_strategy); - Ok(ControlFlow::Continue(to_cleanup)) - } - TransactionStrategyExecutionError::CpiLimitError(_, _) - | TransactionStrategyExecutionError::LoadedAccountsDataSizeExceeded(_, _) => { - // Can't be handled - warn!(error = ?err, "Finalization tasks exceeded execution limit"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::TransactionTooLargeError(_) => { - // Can't be handled but also shouldn't occur - error!(strategy = ?self.state.finalize_strategy, error = ?err, "Finalization tasks do not fit in tx"); - Ok(ControlFlow::Break(())) - } - TransactionStrategyExecutionError::InternalError(_) => { - // Can't be handled - Ok(ControlFlow::Break(())) - } - } - } -} - -mod sealed { - pub trait Sealed {} - - impl Sealed for super::Initialized {} - impl Sealed for super::Committed {} -} diff --git a/magicblock-committor-service/src/intent_executor/utils.rs b/magicblock-committor-service/src/intent_executor/utils.rs index 38482399a..645782e0d 100644 --- a/magicblock-committor-service/src/intent_executor/utils.rs +++ b/magicblock-committor-service/src/intent_executor/utils.rs @@ -1,22 +1,15 @@ -use std::{collections::HashMap, time::Duration}; +use std::collections::HashMap; -use async_trait::async_trait; -use magicblock_core::traits::{ - ActionError, ActionResult, ActionsCallbackScheduler, -}; +use magicblock_core::traits::{ActionResult, ActionsCallbackScheduler}; use solana_keypair::Keypair; use solana_pubkey::Pubkey; use solana_signature::Signature; -use tokio::time::timeout; -use tracing::info; use crate::{ intent_executor::{ error::{IntentExecutorResult, TransactionStrategyExecutionError}, intent_execution_client::IntentExecutionClient, - single_stage_executor::SingleStageExecutor, task_info_fetcher::{CacheTaskInfoFetcher, ResetType, TaskInfoFetcher}, - two_stage_executor::{Committed, Initialized, TwoStageExecutor}, IntentExecutionReport, }, persist::IntentPersister, @@ -262,167 +255,6 @@ where junk } -pub(in crate::intent_executor) async fn execute_with_timeout< - P: IntentPersister, ->( - time_left: Option, - mut executor: impl StageExecutor, - persister: &Option

, -) -> IntentExecutorResult { - if executor.has_callbacks() { - if let Some(time_left) = time_left { - match timeout(time_left, executor.execute(persister)).await { - Ok(res) => return res, - Err(_) => { - // The race between callback and intent txn is handled - // on the user smart contract side via TimeoutError. - // We must respect the timeout contract. - info!("Intent execution timed out, cleaning up actions"); - executor.execute_callbacks( - None, - Err(ActionError::TimeoutError), - ); - } - } - } else { - // Already timed out; see comment above. - executor.execute_callbacks(None, Err(ActionError::TimeoutError)); - } - } - - executor.execute(persister).await -} - -#[async_trait] -pub(in crate::intent_executor) trait StageExecutor { - fn has_callbacks(&self) -> bool; - async fn execute( - &mut self, - persister: &Option

, - ) -> IntentExecutorResult; - fn execute_callbacks( - &mut self, - signature: Option, - result: ActionResult, - ); -} - -pub(in crate::intent_executor) struct SingleStage<'a, 'e, A, T, F> { - pub(in crate::intent_executor) inner: &'a mut SingleStageExecutor<'e, F, A>, - pub(in crate::intent_executor) transaction_preparator: &'a T, - pub(in crate::intent_executor) committed_pubkeys: &'a [Pubkey], -} - -#[async_trait] -impl<'a, 'e, A, T, F> StageExecutor for SingleStage<'a, 'e, A, T, F> -where - A: ActionsCallbackScheduler, - T: TransactionPreparator, - F: TaskInfoFetcher, -{ - fn has_callbacks(&self) -> bool { - self.inner.has_callbacks() - } - - async fn execute( - &mut self, - persister: &Option

, - ) -> IntentExecutorResult { - self.inner - .execute( - self.committed_pubkeys, - self.transaction_preparator, - persister, - ) - .await - } - - fn execute_callbacks( - &mut self, - signature: Option, - result: ActionResult, - ) { - self.inner.execute_callbacks(signature, result) - } -} - -pub(in crate::intent_executor) struct CommitStage<'a, 'e, A, T, F> { - pub(in crate::intent_executor) inner: - &'a mut TwoStageExecutor<'e, A, Initialized>, - pub(in crate::intent_executor) transaction_preparator: &'a T, - pub(in crate::intent_executor) task_info_fetcher: - &'a CacheTaskInfoFetcher, - pub(in crate::intent_executor) committed_pubkeys: &'a [Pubkey], -} - -#[async_trait] -impl<'a, 'e, A, T, F> StageExecutor for CommitStage<'a, 'e, A, T, F> -where - A: ActionsCallbackScheduler, - T: TransactionPreparator, - F: TaskInfoFetcher, -{ - fn has_callbacks(&self) -> bool { - self.inner.has_callbacks() - } - - async fn execute( - &mut self, - persister: &Option

, - ) -> IntentExecutorResult { - self.inner - .commit( - self.committed_pubkeys, - self.transaction_preparator, - self.task_info_fetcher, - persister, - ) - .await - } - - fn execute_callbacks( - &mut self, - signature: Option, - result: ActionResult, - ) { - self.inner.execute_callbacks(signature, result) - } -} - -pub(in crate::intent_executor) struct FinalizeStage<'a, 'e, A, T> { - pub(in crate::intent_executor) inner: - &'a mut TwoStageExecutor<'e, A, Committed>, - pub(in crate::intent_executor) transaction_preparator: &'a T, -} - -#[async_trait] -impl<'a, 'e, A, T> StageExecutor for FinalizeStage<'a, 'e, A, T> -where - A: ActionsCallbackScheduler, - T: TransactionPreparator, -{ - fn has_callbacks(&self) -> bool { - self.inner.has_callbacks() - } - - async fn execute( - &mut self, - persister: &Option

, - ) -> IntentExecutorResult { - self.inner - .finalize(self.transaction_preparator, persister) - .await - } - - fn execute_callbacks( - &mut self, - signature: Option, - result: ActionResult, - ) { - self.inner.execute_callbacks(signature, result) - } -} - #[cfg(test)] mod tests { use std::collections::HashMap; diff --git a/test-integration/test-committor-service/tests/test_intent_executor.rs b/test-integration/test-committor-service/tests/test_intent_executor.rs index eacac4c42..facc75504 100644 --- a/test-integration/test-committor-service/tests/test_intent_executor.rs +++ b/test-integration/test-committor-service/tests/test_intent_executor.rs @@ -19,22 +19,25 @@ use magicblock_committor_service::{ CacheTaskInfoFetcher, RpcTaskInfoFetcher, TaskInfoFetcher, TaskInfoFetcherError, }, - two_stage_executor::{Initialized, TwoStageExecutor}, utils::prepare_and_execute_strategy, ExecutionOutput, IntentExecutionReport, IntentExecutionResult, IntentExecutor, IntentExecutorImpl, }, - persist::IntentPersisterImpl, + persist::{IntentPersister, IntentPersisterImpl}, tasks::{ task_builder::{TaskBuilderError, TaskBuilderImpl}, task_strategist::{ - TaskStrategist, TaskStrategistError, TransactionStrategy, + StrategyExecutionMode, TaskStrategist, TaskStrategistError, + TransactionStrategy, }, BaseTask, }, transaction_preparator::{ + delivery_preparator::{BufferExecutionError, DeliveryPreparatorResult}, + error::PreparatorResult, TransactionPreparator, TransactionPreparatorImpl, }, + transactions::PreparedMessage, DEFAULT_ACTIONS_TIMEOUT, }; use magicblock_core::{ @@ -838,9 +841,10 @@ async fn test_cpi_limits_error_recovery() { .await; let mut execution_report = IntentExecutionReport::default(); let execution_result = intent_executor - .single_stage_execution_flow( - scheduled_intent, - strategy, + .execute_strategy( + scheduled_intent.id, + &scheduled_intent.get_all_committed_pubkeys(), + StrategyExecutionMode::SingleStage(strategy), &mut execution_report, &None::, ) @@ -967,9 +971,10 @@ async fn test_commit_id_actions_cpi_limit_errors_recovery() { intent_executor.started_at = std::time::Instant::now(); let mut execution_report = IntentExecutionReport::default(); let res = intent_executor - .single_stage_execution_flow( - scheduled_intent, - strategy, + .execute_strategy( + scheduled_intent.id, + &scheduled_intent.get_all_committed_pubkeys(), + StrategyExecutionMode::SingleStage(strategy), &mut execution_report, &None::, ) @@ -1187,6 +1192,7 @@ async fn test_action_callback_fired_on_timeout() { async fn test_two_stage_action_failure_keeps_combined_commit() { let TestEnv { fixture, + mut intent_executor, task_info_fetcher, callback_executor, .. @@ -1211,34 +1217,30 @@ async fn test_two_stage_action_failure_keeps_combined_commit() { base_actions: actions, }, )); - let preparator = fixture.create_transaction_preparator(); let mut report = IntentExecutionReport::default(); - let mut executor = create_two_stage_executor( - &fixture, - &callback_executor, - &intent, - &task_info_fetcher, - &mut report, - ) - .await; - let signature = executor - .commit( - &[pubkey], - &preparator, - &task_info_fetcher, + let strategy = + create_two_transaction_strategy(&fixture, &intent, &task_info_fetcher) + .await; + let result = intent_executor + .execute_strategy( + intent.id, + &intent.get_all_committed_pubkeys(), + strategy, + &mut report, &None::, ) .await .unwrap(); - let mut finalize = executor.done(signature); - let final_signature = finalize - .finalize(&preparator, &None::) - .await - .unwrap(); - drop(finalize); + let ExecutionOutput::TwoStage { + commit_signature, + finalize_signature, + } = result + else { + panic!("expected two-transaction execution"); + }; assert_eq!( - final_signature, signature, - "no empty second-stage transaction should be sent" + finalize_signature, commit_signature, + "no empty follow-up transaction should be sent" ); assert!(matches!( report.patched_errors().as_slice(), @@ -1299,68 +1301,92 @@ async fn test_callbacks_fired_in_two_stage() { standalone_actions: vec![commit_base_action], ..Default::default() }); - let committed_pubkeys = intent.get_all_committed_pubkeys(); - - let transaction_preparator = fixture.create_transaction_preparator(); - let mut execution_report = IntentExecutionReport::default(); - let mut executor = create_two_stage_executor( - &fixture, - &callback_executor, - &intent, - &task_info_fetcher, - &mut execution_report, - ) - .await; - - // Execute commit stage - let commit_sig = executor - .commit( - &committed_pubkeys, - &transaction_preparator, - &task_info_fetcher, + let transaction_preparator = CallbackCheckingPreparator { + inner: fixture.create_transaction_preparator(), + callback_executor: callback_executor.clone(), + first_callback: expected_commit_callback.clone(), + preparations: AtomicU64::new(0), + }; + let mut executor = IntentExecutorImpl::new( + fixture.rpc_client.clone(), + transaction_preparator, + task_info_fetcher.clone(), + callback_executor.clone(), + DEFAULT_ACTIONS_TIMEOUT, + ); + let strategy = + create_two_transaction_strategy(&fixture, &intent, &task_info_fetcher) + .await; + let mut report = IntentExecutionReport::default(); + executor + .execute_strategy( + intent.id, + &intent.get_all_committed_pubkeys(), + strategy, + &mut report, &None::, ) .await - .expect("commit must succeed"); + .expect("both transactions must succeed"); - // standalone_actions land in the commit strategy as a BaseAction let calls = callback_executor.calls(); - assert_eq!( - calls.len(), - 1, - "commit-stage callback must be fired by execute_callbacks" - ); + assert_eq!(calls.len(), 2, "each transaction reports its own callbacks"); assert_eq!(calls[0].0[0], expected_commit_callback); assert!(calls[0].1.is_ok()); - - // Execute finalize stage - let mut finalize_executor = executor.done(commit_sig); - finalize_executor - .finalize(&transaction_preparator, &None::) - .await - .expect("finalize must succeed"); - - // Expect 2 actions to be executed in finalize stage - let calls = callback_executor.calls(); - assert_eq!( - calls.len(), - 2, - "finalize-stage callback must be fired by execute_callbacks" - ); assert_eq!(calls[1].0[0], expected_finalize_callback); assert!(calls[1].1.is_ok()); } -/// Builds a [`TwoStageExecutor`] directly from an intent by constructing the -/// commit and finalize strategies independently, without going through -/// `execute_inner` or any CPI-limit recovery path. -async fn create_two_stage_executor<'a>( +/// Checks callback timing before preparing the next transaction. +struct CallbackCheckingPreparator { + inner: TransactionPreparatorImpl, + callback_executor: MockActionsCallbackExecutor, + first_callback: BaseActionCallback, + preparations: AtomicU64, +} + +#[async_trait::async_trait] +impl TransactionPreparator for CallbackCheckingPreparator { + async fn prepare_for_strategy( + &self, + authority: &Keypair, + strategy: &mut TransactionStrategy, + persister: &Option

, + ) -> PreparatorResult { + if self.preparations.fetch_add(1, Ordering::Relaxed) == 1 { + let calls = self.callback_executor.calls(); + assert_eq!( + calls.len(), + 1, + "only the first transaction's callback is due" + ); + assert_eq!(calls[0].0[0], self.first_callback); + assert!(calls[0].1.is_ok()); + } + self.inner + .prepare_for_strategy(authority, strategy, persister) + .await + } + + async fn cleanup_for_strategy( + &self, + authority: &Keypair, + strategy: &TransactionStrategy, + close_buffers: bool, + ) -> DeliveryPreparatorResult<(), BufferExecutionError> { + self.inner + .cleanup_for_strategy(authority, strategy, close_buffers) + .await + } +} + +/// Builds two transactions directly, bypassing the packing decision so tests +/// can exercise follow-up execution without depending on size thresholds. +async fn create_two_transaction_strategy( fixture: &TestFixture, - callback_executor: &MockActionsCallbackExecutor, intent: &ScheduledIntentBundle, task_info_fetcher: &Arc>, - execution_report: &'a mut IntentExecutionReport, -) -> TwoStageExecutor<'a, MockActionsCallbackExecutor, Initialized> { +) -> StrategyExecutionMode { let authority = &fixture.authority.pubkey(); let commit_tasks = TaskBuilderImpl::commit_tasks( task_info_fetcher, @@ -1387,15 +1413,10 @@ async fn create_two_stage_executor<'a>( None, ) .unwrap(); - TwoStageExecutor::new( - fixture.authority.insecure_clone(), - commit_strategy, - finalize_strategy, - IntentExecutionClient::new(fixture.rpc_client.clone()), - callback_executor.clone(), - execution_report, - intent.id, - ) + StrategyExecutionMode::TwoStage { + commit_stage: commit_strategy, + finalize_stage: finalize_strategy, + } } fn succeeding_commit_action(