use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
use degenbot_arbitrage::{
candidate_is_stale, dispatch_profitable_results, DispatchOutcome, FeeOnTransferRegistry,
PoolDivergence, SimulateContext,
};
use degenbot_executor::composers::{EncodeOptions, PathInfo};
use degenbot_rpc::provider::AlloyProvider;
use degenbot_submission::{
dispatch_and_submit, Dispatcher, NonceLane, PathSuppression, PipelineFailure, ReceiptProbe,
SimFuture, SimLeaf, SimSubmitPipeline, SubmissionTarget, SubmitLeaf, SubmitRecord, TxSigner,
};
use degenbot_substrate::state_lock::{LockSite, StateLock};
use degenbot_substrate::BotState;
use parking_lot::RwLock;
use crate::assembly::{
assemble_batch, join_sim_result, AssemblyError, AssemblyInputs, PathResolver,
};
use crate::policy::{ExecutorPolicy, ExecutorRuntime};
use crate::record::{fold_counters, AssemblyVerdict, BatchOutcome, SimulateVerdict, SubmitVerdict};
use crate::row::{PayloadRow, RawResult};
pub struct ExecutorConfig {
pub(crate) sim_concurrency: usize,
pub(crate) max_candidates: usize,
pub(crate) min_profit_margin_bps: u64,
pub(crate) opts: EncodeOptions,
pub(crate) resolver: Arc<dyn PathResolver>,
pub(crate) suppression: Arc<Mutex<PathSuppression>>,
pub(crate) divergence: Arc<Mutex<PoolDivergence>>,
pub(crate) fot: Arc<Mutex<FeeOnTransferRegistry>>,
pub(crate) provider: Arc<AlloyProvider>,
pub(crate) executor_owner: alloy::primitives::Address,
pub(crate) executor_address: alloy::primitives::Address,
pub(crate) weth_address: alloy::primitives::Address,
pub(crate) pool_manager_address: alloy::primitives::Address,
pub(crate) multicall3_address: alloy::primitives::Address,
pub(crate) inject_code: bool,
pub(crate) injected_address: Option<alloy::primitives::Address>,
pub(crate) runtime_bytecode: alloy::primitives::Bytes,
pub(crate) warmup: degenbot_executor::WarmupSlots,
pub(crate) bot_state: Option<Arc<StateLock<BotState>>>,
pub(crate) warm_cache: Option<Arc<RwLock<degenbot_simulation::WarmCodeCacheInner>>>,
pub(crate) dispatcher: Arc<Mutex<Dispatcher>>,
pub(crate) signer: Arc<TxSigner>,
pub(crate) probe: Arc<dyn ReceiptProbe + Send + Sync>,
pub(crate) nonce_lane: Arc<NonceLane>,
pub(crate) dry_run: bool,
pub(crate) inject_code_guard: bool,
pub(crate) extra_broadcast: Vec<Arc<AlloyProvider>>,
pub(crate) target: SubmissionTarget,
}
impl ExecutorConfig {
#[must_use]
pub fn from_verdict(verdict: °enbot_config::BotConfig, runtime: ExecutorRuntime) -> Self {
let policy = ExecutorPolicy::from(verdict);
let max_candidates = if runtime.max_candidates == 0 {
usize::MAX
} else {
runtime.max_candidates
};
Self {
sim_concurrency: policy.sim_concurrency,
max_candidates,
min_profit_margin_bps: policy.min_profit_margin_bps,
opts: EncodeOptions {
erc6909_profit: policy.erc6909_profit,
use_v4_batch: runtime.use_v4_batch,
..EncodeOptions::default()
},
resolver: runtime.resolver,
suppression: runtime.suppression,
divergence: runtime.divergence,
fot: runtime.fot,
provider: runtime.provider,
executor_owner: runtime.executor_owner,
executor_address: runtime.executor_address,
weth_address: runtime.weth_address,
pool_manager_address: runtime.pool_manager_address,
multicall3_address: runtime.multicall3_address,
inject_code: policy.inject_code,
injected_address: runtime.injected_address,
runtime_bytecode: runtime.runtime_bytecode,
warmup: runtime.warmup,
bot_state: runtime.bot_state,
warm_cache: runtime.warm_cache,
dispatcher: runtime.dispatcher,
signer: runtime.signer,
probe: runtime.probe,
nonce_lane: runtime.nonce_lane,
dry_run: runtime.dry_run,
inject_code_guard: policy.inject_code_guard,
extra_broadcast: runtime.extra_broadcast,
target: runtime.target,
}
}
}
#[derive(Debug, Clone)]
pub struct BatchWork {
pub rows: Vec<RawResult>,
pub payloads: Vec<PayloadRow>,
pub current_block: u64,
pub base_fee_next: u128,
pub block_timestamp: u64,
pub block_priority_fees: Option<degenbot_rpc::BlockPriorityFees>,
}
struct SimStageOutput {
outcomes: Vec<BatchOutcome>,
submit_candidates: Vec<degenbot_submission::SubmitCandidate>,
}
#[derive(Debug, Clone)]
pub struct BatchOutcomeSet {
pub records: Vec<BatchOutcome>,
pub submit_records: Vec<SubmitRecord>,
}
struct EngineLeaves {
config: Arc<ExecutorConfig>,
outcome_tx: mpsc::WeakUnboundedSender<BatchOutcomeSet>,
}
impl EngineLeaves {
#[expect(
clippy::too_many_lines,
reason = "the leaf reads best as one ordered stage pass"
)]
fn run_stages_1_2(&self, work: &BatchWork) -> Result<SimStageOutput, String> {
let cfg = &self.config;
let served: HashSet<u64> = work.payloads.iter().map(|p| p.path_id).collect();
let mut batch = {
let mut suppression = cfg
.suppression
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
assemble_batch(&AssemblyInputs {
rows: &work.rows,
payloads: &work.payloads,
resolver: cfg.resolver.as_ref(),
opts: cfg.opts,
payload_served: &served,
suppression: &mut suppression,
divergence: &cfg.divergence,
fot: &cfg.fot,
current_block: work.current_block,
min_profit_margin_bps: cfg.min_profit_margin_bps,
max_candidates: cfg.max_candidates,
executor_address: cfg.executor_address,
})
.map_err(|e: AssemblyError| e.to_string())?
};
if let Some(arc) = &cfg.bot_state {
let stale_ids: Vec<u64> = {
let guard = arc.read_at(LockSite::Sim);
batch
.candidates
.iter()
.filter(|c| candidate_is_stale(&guard, c))
.map(|c| c.path_id)
.collect()
};
if !stale_ids.is_empty() {
batch.candidates.retain(|c| !stale_ids.contains(&c.path_id));
batch
.outcomes
.retain(|o| !(o.simulate.is_none() && stale_ids.contains(&o.path_id)));
}
}
let provider = Arc::clone(&cfg.provider);
let ctx = SimulateContext {
provider: &provider,
executor_owner: cfg.executor_owner,
executor_address: cfg.executor_address,
weth_address: cfg.weth_address,
pool_manager_address: cfg.pool_manager_address,
multicall3_address: cfg.multicall3_address,
inject_code: cfg.inject_code,
injected_address: cfg.injected_address,
runtime_bytecode: cfg.runtime_bytecode.clone(),
warmup: cfg.warmup,
base_fee_next: work.base_fee_next,
current_block: work.current_block,
block_timestamp: work.block_timestamp,
block_priority_fees: work.block_priority_fees.clone(),
};
let outcome: DispatchOutcome = {
let span = degenbot_bot::telemetry::simulate_dispatch_span(
work.current_block,
batch.candidates.len(),
);
let _guard = span.enter();
dispatch_profitable_results(
std::mem::take(&mut batch.candidates),
&ctx,
&cfg.suppression,
work.current_block,
cfg.min_profit_margin_bps,
&cfg.divergence,
&cfg.fot,
cfg.bot_state.clone(),
cfg.warm_cache.clone(),
)
};
let profitable: HashMap<u64, degenbot_arbitrage::SimResult> = outcome
.gas_profitable
.iter()
.map(|r| (r.path_id, r.clone()))
.collect();
let unprofitable: HashMap<u64, degenbot_arbitrage::SimResult> = outcome
.gas_unprofitable
.iter()
.map(|r| (r.path_id, r.clone()))
.collect();
let failures: HashMap<u64, degenbot_arbitrage::SimFailure> = outcome
.failures
.iter()
.map(|f| (f.path_id, f.clone()))
.collect();
let mut submit_candidates: Vec<degenbot_submission::SubmitCandidate> = Vec::new();
let empty_path = PathInfo::new(Vec::new());
for c in std::mem::take(&mut batch.candidates) {
let path_info = batch.path_info_by_id.get(&c.path_id);
let simulate = if let Some(r) = profitable.get(&c.path_id) {
submit_candidates.push(join_sim_result(r, path_info, cfg.executor_address));
Some(SimulateVerdict::Profitable(crate::record::SimReceipt {
gross_profit: r.gross_profit,
net_profit: r.net_profit,
gas_used: r.gas_used,
priority_fee: r.priority_fee,
}))
} else if let Some(r) = unprofitable.get(&c.path_id) {
Some(SimulateVerdict::GasUnprofitable(
crate::record::SimReceipt {
gross_profit: r.gross_profit,
net_profit: r.net_profit,
gas_used: r.gas_used,
priority_fee: r.priority_fee,
},
))
} else if let Some(f) = failures.get(&c.path_id) {
Some(SimulateVerdict::Failed(Box::new(f.clone())))
} else {
Some(SimulateVerdict::Exception)
};
batch.outcomes.push(BatchOutcome {
path_id: c.path_id,
block: work.current_block,
assembly: AssemblyVerdict::Assembled,
simulate,
submit: None,
path_info: crate::record::PathInfoView::from(path_info.unwrap_or(&empty_path)),
});
}
submit_candidates.append(&mut batch.payload_submits);
Ok(SimStageOutput {
outcomes: batch.outcomes,
submit_candidates,
})
}
}
impl SimLeaf<BatchWork, SimStageOutput> for EngineLeaves {
fn simulate<'a>(&'a self, work: &'a BatchWork) -> SimFuture<'a, SimStageOutput> {
Box::pin(async move { self.run_stages_1_2(work).map(Some) })
}
}
fn log_submit_arm(candidates: &[degenbot_submission::SubmitCandidate], solve_block: u64) {
for c in candidates {
degenbot_core::op_info!(
domain = exec,
"path={} solve_block={} net_wei={} gas={} calldata={}",
c.path_id,
solve_block,
c.net_profit,
c.gas_used,
alloy::hex::encode(&c.execute_calldata),
);
}
}
impl SubmitLeaf<BatchWork, SimStageOutput> for EngineLeaves {
fn submit<'a>(
&'a self,
work: &'a BatchWork,
outcome: &'a SimStageOutput,
) -> degenbot_submission::SubmitFuture<'a> {
let cfg = Arc::clone(&self.config);
let tx = self.outcome_tx.clone();
let outcomes = outcome.outcomes.clone();
let candidates = outcome.submit_candidates.clone();
let current_block = work.current_block;
Box::pin(async move {
log_submit_arm(&candidates, current_block);
let submit_outcome = dispatch_and_submit(
candidates,
&cfg.dispatcher,
&cfg.provider,
&cfg.signer,
Arc::clone(&cfg.probe),
&cfg.nonce_lane,
current_block,
cfg.dry_run,
cfg.inject_code_guard,
&cfg.extra_broadcast,
cfg.target.clone(),
)
.await
.map_err(|e| e.to_string())?;
let mut stamped = outcomes;
for record in &submit_outcome.records {
let (path_id, verdict) = match record {
SubmitRecord::Submitted { path_id, .. } => (*path_id, SubmitVerdict::Submitted),
SubmitRecord::Skipped { path_id, reason } => (
*path_id,
SubmitVerdict::Failed(crate::record::FailureKind::from_skip_reason(reason)),
),
};
if let Some(slot) = stamped
.iter_mut()
.find(|o| o.path_id == path_id && o.submit.is_none())
{
slot.submit = Some(verdict);
}
}
if let Some(tx) = tx.upgrade() {
let _ = tx.send(BatchOutcomeSet {
records: stamped,
submit_records: submit_outcome.records.clone(),
});
}
Ok(())
})
}
}
pub struct BatchDrain {
rx: Arc<tokio::sync::Mutex<mpsc::UnboundedReceiver<BatchOutcomeSet>>>,
}
impl BatchDrain {
pub async fn next(&mut self) -> Option<BatchOutcomeSet> {
self.rx.lock().await.recv().await
}
}
pub struct BatchExecutor {
pipeline: SimSubmitPipeline<BatchWork, SimStageOutput>,
outcome_tx: std::sync::Mutex<Option<mpsc::UnboundedSender<BatchOutcomeSet>>>,
outcome_rx: Arc<tokio::sync::Mutex<mpsc::UnboundedReceiver<BatchOutcomeSet>>>,
}
impl BatchExecutor {
#[must_use]
pub fn new(config: ExecutorConfig) -> Self {
let (outcome_tx, outcome_rx) = mpsc::unbounded_channel();
let config = Arc::new(config);
let leaves = Arc::new(EngineLeaves {
config: Arc::clone(&config),
outcome_tx: outcome_tx.downgrade(),
});
Self::from_leaves(
config.sim_concurrency,
Arc::clone(&leaves) as Arc<dyn SimLeaf<BatchWork, SimStageOutput>>,
Arc::clone(&leaves) as Arc<dyn SubmitLeaf<BatchWork, SimStageOutput>>,
outcome_tx,
outcome_rx,
)
}
#[must_use]
pub fn drain_handle(&self) -> BatchDrain {
BatchDrain {
rx: Arc::clone(&self.outcome_rx),
}
}
fn from_leaves(
concurrency: usize,
sim: Arc<dyn SimLeaf<BatchWork, SimStageOutput>>,
submit: Arc<dyn SubmitLeaf<BatchWork, SimStageOutput>>,
outcome_tx: mpsc::UnboundedSender<BatchOutcomeSet>,
outcome_rx: mpsc::UnboundedReceiver<BatchOutcomeSet>,
) -> Self {
let pipeline = SimSubmitPipeline::new(concurrency, sim, submit);
Self {
pipeline,
outcome_tx: std::sync::Mutex::new(Some(outcome_tx)),
outcome_rx: Arc::new(tokio::sync::Mutex::new(outcome_rx)),
}
}
pub fn enqueue(&self, work: BatchWork) {
self.pipeline.enqueue(work);
}
pub async fn next_outcome(&self) -> Option<BatchOutcomeSet> {
self.outcome_rx.lock().await.recv().await
}
#[must_use]
pub async fn try_next_outcome(&self) -> Option<BatchOutcomeSet> {
self.outcome_rx.lock().await.try_recv().ok()
}
#[must_use]
pub fn enqueued(&self) -> u64 {
self.pipeline.enqueued()
}
#[must_use]
pub fn submitted(&self) -> u64 {
self.pipeline.submitted()
}
pub fn raise_if_failed(&self) -> Result<(), PipelineFailure> {
self.pipeline.raise_if_failed()
}
pub async fn shutdown(&self) -> Result<(), PipelineFailure> {
let result = self.pipeline.shutdown().await;
*self
.outcome_tx
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = None;
result
}
}
#[must_use]
pub fn batch_counters(batch: &[BatchOutcome]) -> crate::record::BatchCounters {
fold_counters(batch)
}
#[cfg(test)]
#[expect(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
reason = "tests drive known-valid fixtures and assert on the outcomes"
)]
mod tests {
use super::*;
use std::time::Duration;
use tracing_subscriber::layer::SubscriberExt;
use crate::record::{FailureDetail, PathInfoView, SimReceipt};
const BASE_BLOCK: u64 = 100;
fn work(block: u64) -> BatchWork {
BatchWork {
rows: Vec::new(),
payloads: Vec::new(),
current_block: block,
base_fee_next: 0,
block_timestamp: 0,
block_priority_fees: None,
}
}
fn receipt() -> SimReceipt {
SimReceipt {
gross_profit: alloy::primitives::U256::from(600_u64),
net_profit: alloy::primitives::U256::from(500_u64),
gas_used: 300_000,
priority_fee: 2,
}
}
struct ReversedFinishSim {
fail_block: Option<u64>,
}
impl SimLeaf<BatchWork, SimStageOutput> for ReversedFinishSim {
fn simulate<'a>(&'a self, work: &'a BatchWork) -> SimFuture<'a, SimStageOutput> {
let batch = work.current_block - BASE_BLOCK;
let fail_block = self.fail_block;
Box::pin(async move {
tokio::time::sleep(Duration::from_millis(40 - 10 * batch)).await;
if fail_block == Some(work.current_block) {
return Err(format!("sim leaf failed at block {}", work.current_block));
}
let outcome = BatchOutcome {
path_id: batch,
block: work.current_block,
assembly: AssemblyVerdict::Assembled,
simulate: Some(SimulateVerdict::Profitable(receipt())),
submit: None,
path_info: PathInfoView::empty(),
};
Ok(Some(SimStageOutput {
outcomes: vec![outcome],
submit_candidates: Vec::new(),
}))
})
}
}
struct StampingSubmit {
order: Arc<std::sync::Mutex<Vec<u64>>>,
tx: mpsc::WeakUnboundedSender<BatchOutcomeSet>,
}
impl SubmitLeaf<BatchWork, SimStageOutput> for StampingSubmit {
fn submit<'a>(
&'a self,
work: &'a BatchWork,
outcome: &'a SimStageOutput,
) -> degenbot_submission::SubmitFuture<'a> {
let order = Arc::clone(&self.order);
let tx = self.tx.clone();
let block = work.current_block;
let mut records = outcome.outcomes.clone();
Box::pin(async move {
order.lock().unwrap().push(block);
for record in &mut records {
record.submit = Some(SubmitVerdict::Submitted);
}
if let Some(tx) = tx.upgrade() {
let _ = tx.send(BatchOutcomeSet {
records,
submit_records: Vec::new(),
});
}
Ok(())
})
}
}
fn executor_with(
sim: ReversedFinishSim,
order: Arc<std::sync::Mutex<Vec<u64>>>,
) -> BatchExecutor {
let (tx, rx) = mpsc::unbounded_channel();
BatchExecutor::from_leaves(
4,
Arc::new(sim),
Arc::new(StampingSubmit {
order,
tx: tx.downgrade(),
}),
tx,
rx,
)
}
#[tokio::test]
async fn submit_lane_runs_in_nonce_order_not_sim_completion_order() {
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
let executor = executor_with(ReversedFinishSim { fail_block: None }, Arc::clone(&order));
for batch in 0_u64..4 {
executor.enqueue(work(BASE_BLOCK + batch));
}
let failure = executor.shutdown().await;
assert!(failure.is_ok(), "no leaf failed: {failure:?}");
assert_eq!(executor.submitted(), 4);
assert_eq!(
*order.lock().unwrap(),
vec![BASE_BLOCK, BASE_BLOCK + 1, BASE_BLOCK + 2, BASE_BLOCK + 3]
);
let mut drained = Vec::new();
while let Some(records) = executor.next_outcome().await {
drained.extend(records.records.iter().map(|r| r.block));
}
assert_eq!(
drained,
vec![BASE_BLOCK, BASE_BLOCK + 1, BASE_BLOCK + 2, BASE_BLOCK + 3]
);
}
#[tokio::test]
async fn sim_leaf_failure_is_reraised_in_the_caller_frame() {
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
let executor = executor_with(
ReversedFinishSim {
fail_block: Some(BASE_BLOCK),
},
Arc::clone(&order),
);
executor.enqueue(work(BASE_BLOCK));
executor.enqueue(work(BASE_BLOCK + 1));
let failure = executor.shutdown().await;
let failure = failure.expect_err("the failing batch's detail is re-raised");
assert_eq!(failure.detail, "sim leaf failed at block 100");
assert!(order.lock().unwrap().is_empty());
assert_eq!(executor.submitted(), 0);
}
#[tokio::test]
async fn raise_if_failed_surfaces_the_stored_failure_once() {
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
let executor = executor_with(
ReversedFinishSim {
fail_block: Some(BASE_BLOCK + 1),
},
order,
);
executor.enqueue(work(BASE_BLOCK));
executor.enqueue(work(BASE_BLOCK + 1));
let _ = executor.try_next_outcome().await;
while executor.submitted() < 1 {
tokio::time::sleep(Duration::from_millis(5)).await;
}
let failure = executor.raise_if_failed().expect_err("stored failure");
assert_eq!(failure.detail, "sim leaf failed at block 101");
assert!(executor.raise_if_failed().is_ok());
executor
.shutdown()
.await
.expect("the failure was already re-raised");
}
#[tokio::test]
async fn block_only_drain_reads_just_the_heartbeat() {
let order = Arc::new(std::sync::Mutex::new(Vec::new()));
let executor = executor_with(ReversedFinishSim { fail_block: None }, order);
for batch in 0_u64..3 {
executor.enqueue(work(BASE_BLOCK + batch));
}
executor.shutdown().await.expect("no failure");
let mut heartbeats = Vec::new();
while let Some(records) = executor.next_outcome().await {
for record in records.records {
heartbeats.push(record.block);
}
}
assert_eq!(heartbeats, vec![BASE_BLOCK, BASE_BLOCK + 1, BASE_BLOCK + 2]);
}
#[expect(
clippy::too_many_lines,
reason = "the fixture spells every stage's record verbatim"
)]
#[tokio::test]
async fn renderer_drain_reads_every_stage_and_folds_counters() {
use degenbot_executor::composers::V2HopInfo;
use std::collections::HashSet;
struct MixedSim;
impl SimLeaf<BatchWork, SimStageOutput> for MixedSim {
fn simulate<'a>(&'a self, work: &'a BatchWork) -> SimFuture<'a, SimStageOutput> {
Box::pin(async move {
let pool = degenbot_core::address_utils::parse_address(
"0x1111111111111111111111111111111111111111",
)
.unwrap();
let view =
PathInfoView::from(°enbot_executor::composers::PathInfo::new(vec![
degenbot_executor::composers::HopInfo::V2(V2HopInfo {
pool_address: pool,
token0_address: pool,
token1_address: pool,
fee: 30,
zfo: true,
}),
]));
let failure = FailureDetail {
path_id: 3,
bucket: "no-profit".to_string(),
fail_index: None,
revert_data: alloy::primitives::Bytes::new(),
reverting_frame: None,
captured_swaps: Vec::new(),
log_full_count: 0,
reverted_swaps: Vec::new(),
optimal_input: 0,
hop_outputs: Vec::new(),
call_trace: Vec::new(),
weth_before: 0,
weth_after: 0,
eth_before: 0,
eth_after: 0,
erc6909_before: 0,
erc6909_after: 0,
};
let outcomes = vec![
BatchOutcome {
path_id: 1,
block: work.current_block,
assembly: AssemblyVerdict::Assembled,
simulate: Some(SimulateVerdict::Profitable(receipt())),
submit: None,
path_info: view.clone(),
},
BatchOutcome {
path_id: 2,
block: work.current_block,
assembly: AssemblyVerdict::SkipSuppressed,
simulate: None,
submit: None,
path_info: PathInfoView::empty(),
},
BatchOutcome {
path_id: 3,
block: work.current_block,
assembly: AssemblyVerdict::Assembled,
simulate: Some(SimulateVerdict::Failed(Box::new(failure))),
submit: None,
path_info: view.clone(),
},
BatchOutcome {
path_id: 4,
block: work.current_block,
assembly: AssemblyVerdict::Assembled,
simulate: Some(SimulateVerdict::GasUnprofitable(receipt())),
submit: None,
path_info: view,
},
];
Ok(Some(SimStageOutput {
outcomes,
submit_candidates: Vec::new(),
}))
})
}
}
struct PartialSubmit {
tx: mpsc::WeakUnboundedSender<BatchOutcomeSet>,
}
impl SubmitLeaf<BatchWork, SimStageOutput> for PartialSubmit {
fn submit<'a>(
&'a self,
_work: &'a BatchWork,
outcome: &'a SimStageOutput,
) -> degenbot_submission::SubmitFuture<'a> {
let tx = self.tx.clone();
let mut records = outcome.outcomes.clone();
Box::pin(async move {
for record in &mut records {
record.submit = match record.path_id {
1 => Some(SubmitVerdict::Submitted),
4 => Some(SubmitVerdict::Failed(crate::record::FailureKind::RpcFailed)),
_ => None,
};
}
if let Some(tx) = tx.upgrade() {
let _ = tx.send(BatchOutcomeSet {
records,
submit_records: Vec::new(),
});
}
Ok(())
})
}
}
let (tx, rx) = mpsc::unbounded_channel();
let executor = BatchExecutor::from_leaves(
1,
Arc::new(MixedSim),
Arc::new(PartialSubmit { tx: tx.downgrade() }),
tx,
rx,
);
executor.enqueue(work(BASE_BLOCK));
executor.shutdown().await.expect("no failure");
let mut records = executor.next_outcome().await.expect("one batch").records;
assert_eq!(records.len(), 4);
let by_id: std::collections::HashMap<u64, BatchOutcome> =
records.drain(..).map(|r| (r.path_id, r)).collect();
let profitable = &by_id[&1];
assert_eq!(profitable.assembly.label(), "assembled");
let Some(SimulateVerdict::Profitable(sim)) = &profitable.simulate else {
panic!("path 1 is profitable")
};
assert_eq!(sim.net_profit, alloy::primitives::U256::from(500_u64));
assert_eq!(sim.gas_used, 300_000);
assert_eq!(profitable.submit, Some(SubmitVerdict::Submitted));
assert_eq!(profitable.path_info.path_type, "V2");
assert_eq!(profitable.path_info.hops.len(), 1);
assert_eq!(by_id[&2].assembly.label(), "skip-suppressed");
assert!(by_id[&2].simulate.is_none() && by_id[&2].submit.is_none());
let Some(SimulateVerdict::Failed(detail)) = &by_id[&3].simulate else {
panic!("path 3 failed")
};
assert_eq!(detail.bucket, "no-profit");
assert!(by_id[&3].submit.is_none());
assert_eq!(
by_id[&4].submit,
Some(SubmitVerdict::Failed(crate::record::FailureKind::RpcFailed))
);
let counters = batch_counters(&by_id.into_values().collect::<Vec<_>>());
assert_eq!(counters.candidate_count, 3);
assert_eq!(counters.profitable_count, 1);
assert_eq!(counters.gas_unprofitable_count, 1);
assert_eq!(counters.fail_count, 1);
assert_eq!(counters.suppressed_count, 1);
assert_eq!(counters.fail_buckets.get("no-profit"), Some(&1));
assert_eq!(
counters
.fail_buckets
.keys()
.cloned()
.collect::<HashSet<_>>(),
std::iter::once("no-profit".to_string()).collect::<HashSet<_>>()
);
}
fn forensic_candidate(
path_id: u64,
net_profit: u128,
calldata: &[u8],
) -> degenbot_submission::SubmitCandidate {
degenbot_submission::SubmitCandidate {
path_id,
gross_profit: alloy::primitives::U256::from(net_profit + 1_000_000_000_u128),
net_profit: alloy::primitives::U256::from(net_profit),
gas_used: 200_000,
priority_fee: 1_000_000_000_u128,
base_fee_next: 1_000_000_000_u128,
execute_calldata: alloy::primitives::Bytes::copy_from_slice(calldata),
executor_address: degenbot_core::address_utils::parse_address(
"0x1111111111111111111111111111111111111111",
)
.unwrap(),
access_list: None,
path_pools: std::iter::once(degenbot_submission::PoolKey::new("0xpool")).collect(),
}
}
struct EventCapture {
lines: Arc<Mutex<Vec<(String, String)>>>,
}
struct MessageVisitor(String);
impl tracing::field::Visit for MessageVisitor {
fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
if field.name() == "message" {
self.0 = format!("{value:?}");
}
}
}
impl<S> tracing_subscriber::Layer<S> for EventCapture
where
S: tracing::Subscriber,
{
fn on_event(
&self,
event: &tracing::Event<'_>,
_ctx: tracing_subscriber::layer::Context<'_, S>,
) {
let mut visitor = MessageVisitor(String::new());
event.record(&mut visitor);
let message = visitor.0;
let message = message
.strip_prefix('"')
.and_then(|s| s.strip_suffix('"'))
.unwrap_or(&message)
.to_string();
self.lines
.lock()
.unwrap()
.push((event.metadata().target().to_string(), message));
}
}
#[test]
fn submit_arm_forensic_line_captures_calldata_and_economics_per_candidate() {
let lines = Arc::new(Mutex::new(Vec::new()));
let subscriber = tracing_subscriber::registry().with(EventCapture {
lines: Arc::clone(&lines),
});
let _guard = tracing::subscriber::set_default(subscriber);
log_submit_arm(
&[
forensic_candidate(7, 123, &[0xde, 0xad]),
forensic_candidate(9, 456, &[0xbe, 0xef]),
],
4242,
);
let lines = lines.lock().unwrap().clone();
assert_eq!(lines.len(), 2);
assert_eq!(
lines[0],
(
"degenbot::exec".to_string(),
"path=7 solve_block=4242 net_wei=123 gas=200000 calldata=dead".to_string()
)
);
assert_eq!(
lines[1],
(
"degenbot::exec".to_string(),
"path=9 solve_block=4242 net_wei=456 gas=200000 calldata=beef".to_string()
)
);
}
}