use std::collections::BTreeMap;
use metrics::{counter, histogram};
use num_traits::ToPrimitive;
use tracing::{trace, Level};
use super::{
is_rankable, Order, OrderQuote, OrderResponses, OrderSide, QuoteOptions, QuoteStatus,
SolveError, WorkerPoolQuote,
};
use crate::{bps, simulation::deviation::deviation_bps, SimulationResult};
const QUOTE_COMPARISON_TARGET: &str = "fynd::quote_comparison";
#[derive(Clone, Copy)]
struct OrderCoverage {
ranked_candidates: usize,
responders: usize,
}
pub(super) fn record_quote_comparison(
order: &Order,
responses: &OrderResponses,
options: &QuoteOptions,
) {
let mut ranked: Vec<&WorkerPoolQuote> = responses
.quotes
.iter()
.filter(|wq| is_rankable(&wq.quote, options))
.collect();
ranked.sort_by(|a, b| {
b.quote
.amount_out_net_gas()
.cmp(a.quote.amount_out_net_gas())
});
let baseline_net = ranked
.last()
.and_then(|wq| wq.quote.amount_out_net_gas().to_f64());
let best_pool = ranked
.first()
.map(|wq| wq.worker_pool.as_str());
let coverage = OrderCoverage {
ranked_candidates: ranked.len(),
responders: responses.quotes.len() + responses.failed_solvers.len(),
};
let log_enabled = tracing::enabled!(target: QUOTE_COMPARISON_TARGET, Level::TRACE);
for worker_quote in &responses.quotes {
let improvement = ranked
.iter()
.any(|wq| wq.worker_pool == worker_quote.worker_pool)
.then(|| {
improvement_bps(
baseline_net,
worker_quote
.quote
.amount_out_net_gas()
.to_f64(),
)
})
.flatten();
if let Some(bps) = improvement {
histogram!(
"quote_improvement_bps",
"pool" => worker_quote.worker_pool.clone(),
"algorithm" => worker_quote.quote.algorithm().to_string()
)
.record(bps);
}
if log_enabled {
log_quote(
order,
worker_quote,
coverage,
improvement,
best_pool == Some(worker_quote.worker_pool.as_str()),
);
}
}
if log_enabled {
for (worker_pool, error) in &responses.failed_solvers {
log_failure(order, worker_pool, error, coverage);
}
}
}
fn log_quote(
order: &Order,
worker_quote: &WorkerPoolQuote,
coverage: OrderCoverage,
improvement_bps: Option<f64>,
is_best: bool,
) {
let quote = &worker_quote.quote;
trace!(
target: QUOTE_COMPARISON_TARGET,
parent: None,
"quote_comparison order_id={} block={} token_in={} token_out={} side={} amount_in={} \
pool={} algorithm={} status={} amount_out={} amount_out_net_gas={} gas_estimate={} \
solve_time_ms={} improvement_bps={} is_best={} ranked_candidates={} responders={}",
order.id(),
quote.block().number(),
order.token_in(),
order.token_out(),
order_side_label(order.side()),
quote.amount_in(),
worker_quote.worker_pool,
quote.algorithm(),
quote_status_label(quote.status()),
quote.amount_out(),
quote.amount_out_net_gas(),
quote.gas_estimate(),
worker_quote.solve_time_ms,
improvement_bps
.map(|bps| format!("{bps:.4}"))
.unwrap_or_default(),
is_best,
coverage.ranked_candidates,
coverage.responders,
);
}
fn log_failure(order: &Order, worker_pool: &str, error: &SolveError, coverage: OrderCoverage) {
let solve_time_ms = match error {
SolveError::Timeout { elapsed_ms } => elapsed_ms.to_string(),
_ => String::new(),
};
trace!(
target: QUOTE_COMPARISON_TARGET,
parent: None,
"quote_comparison order_id={} block= token_in={} token_out={} side={} amount_in={} \
pool={} algorithm= status={} amount_out= amount_out_net_gas= gas_estimate= \
solve_time_ms={} improvement_bps= is_best=false ranked_candidates={} responders={}",
order.id(),
order.token_in(),
order.token_out(),
order_side_label(order.side()),
order.amount(),
worker_pool,
solver_error_label(error),
solve_time_ms,
coverage.ranked_candidates,
coverage.responders,
);
}
fn quote_status_label(status: QuoteStatus) -> &'static str {
match status {
QuoteStatus::Success => "success",
QuoteStatus::NoRouteFound => "no_route",
QuoteStatus::InsufficientLiquidity => "insufficient_liquidity",
QuoteStatus::Timeout => "timeout",
QuoteStatus::NotReady => "not_ready",
QuoteStatus::PriceCheckFailed => "price_check_failed",
QuoteStatus::EncodingFailed => "encoding_failed",
}
}
fn order_side_label(side: OrderSide) -> &'static str {
match side {
OrderSide::Sell => "sell",
}
}
pub(super) fn solver_error_label(error: &SolveError) -> &'static str {
match error {
SolveError::Timeout { .. } => "timeout",
SolveError::NoRouteFound { .. } => "no_route",
SolveError::RouteRejected { .. } => "route_rejected",
SolveError::InsufficientLiquidity { .. } => "insufficient_liquidity",
SolveError::QueueFull => "queue_full",
SolveError::Internal(_) => "internal",
SolveError::InvalidWorkerPools(_) => "invalid_worker_pools",
SolveError::PriceCheckFailed { .. } => "price_check_failed",
SolveError::AlgorithmError(_) => "algorithm_error",
SolveError::MarketDataStale { .. } => "market_data_stale",
SolveError::InvalidOrder(_) => "invalid_order",
SolveError::NotReady(_) => "not_ready",
SolveError::ComputationFailed(_) => "computation_failed",
SolveError::FailedEncoding(_) => "encoding_failed",
SolveError::EncodingUnavailable(_) => "encoding_unavailable",
SolveError::MaxGasExceeded => "max_gas_exceeded",
SolveError::MissingData(_) => "missing_data",
SolveError::SimulationFailed(_) => "simulation_failed",
}
}
fn improvement_bps(baseline_net: Option<f64>, net: Option<f64>) -> Option<f64> {
let (baseline, net) = (baseline_net?, net?);
if baseline <= 0.0 || !baseline.is_finite() || !net.is_finite() {
return None;
}
Some((net - baseline) / baseline * f64::from(bps::DENOMINATOR))
}
const WINNING_PROTOCOLS_TARGET: &str = "fynd::winning_protocols";
const SHORTFALL_SCALE: f64 = 100.0;
pub(super) fn record_winning_protocols(quote: &OrderQuote) {
if quote.status() != QuoteStatus::Success {
return;
}
let Some(route) = quote.route() else {
return;
};
let mut swaps_per_protocol: BTreeMap<&str, usize> = BTreeMap::new();
for swap in route.swaps() {
*swaps_per_protocol
.entry(swap.protocol())
.or_default() += 1;
}
let (outcome, deviation) = simulation_outcome(quote);
count_swaps_per_protocol(&swaps_per_protocol, outcome);
if let Some(bps) = deviation {
record_shortfall_per_protocol(&swaps_per_protocol, bps);
}
log_winning_protocols(quote, &swaps_per_protocol, outcome, deviation);
}
fn count_swaps_per_protocol(swaps_per_protocol: &BTreeMap<&str, usize>, outcome: &'static str) {
let simulated = if outcome.is_empty() { "none" } else { outcome };
for (protocol, swaps) in swaps_per_protocol {
counter!(
"winning_quote_swaps_total",
"protocol" => protocol.to_string(),
"simulated" => simulated
)
.increment(*swaps as u64);
}
}
fn record_shortfall_per_protocol(swaps_per_protocol: &BTreeMap<&str, usize>, bps: f64) {
if bps >= 0.0 {
return;
}
let shortfall = (-bps * SHORTFALL_SCALE).round() as u64;
for protocol in swaps_per_protocol.keys() {
counter!("winning_quote_shortfall_centibps_total", "protocol" => protocol.to_string())
.increment(shortfall);
counter!("winning_quote_shortfall_routes_total", "protocol" => protocol.to_string())
.increment(1);
}
}
fn simulation_outcome(quote: &OrderQuote) -> (&'static str, Option<f64>) {
match quote.simulation_result() {
None => ("", None),
Some(SimulationResult::Failure { .. }) => ("failed", None),
Some(SimulationResult::Success { amount_out, .. }) => {
("success", deviation_bps(quote, amount_out))
}
}
}
fn log_winning_protocols(
quote: &OrderQuote,
swaps_per_protocol: &BTreeMap<&str, usize>,
outcome: &'static str,
deviation: Option<f64>,
) {
if !tracing::enabled!(target: WINNING_PROTOCOLS_TARGET, Level::TRACE) {
return;
}
let swaps = quote
.route()
.map_or(0, |route| route.swaps().len());
let deviation = deviation.map_or_else(String::new, |bps| format!("{bps:.4}"));
trace!(
target: WINNING_PROTOCOLS_TARGET,
parent: None,
"winning_protocols order_id={} block={} pool={} algorithm={} swaps={} simulated={} \
deviation_bps={} protocols={}",
quote.order_id(),
quote.block().number(),
quote.worker_pool(),
quote.algorithm(),
swaps,
outcome,
deviation,
serde_json::to_string(swaps_per_protocol).unwrap_or_default(),
);
}
#[cfg(test)]
mod tests {
use num_bigint::BigUint;
use rstest::rstest;
use tycho_simulation::tycho_common::{models::Address, Bytes};
use super::*;
use crate::{BlockInfo, OrderQuote};
fn timed_worker_quote(
worker_pool: &str,
quote: OrderQuote,
solve_time_ms: u64,
) -> WorkerPoolQuote {
WorkerPoolQuote { worker_pool: worker_pool.to_string(), quote, solve_time_ms }
}
fn make_address(byte: u8) -> Address {
Address::from([byte; 20])
}
fn comparison_quote(worker_pool: &str, net: u64, solve_time_ms: u64) -> WorkerPoolQuote {
timed_worker_quote(
worker_pool,
OrderQuote::new(
"o1".to_string(),
QuoteStatus::Success,
BigUint::from(1_000u64),
BigUint::from(net + 10),
BigUint::from(10u64),
BigUint::from(net),
BlockInfo::new(42, "0xabc".to_string(), 0),
format!("{worker_pool}_algo"),
Bytes::default(),
Bytes::default(),
"1".to_string(),
),
solve_time_ms,
)
}
fn comparison_order() -> Order {
Order::new(
make_address(0xAA),
make_address(0xBB),
BigUint::from(1_000u64),
OrderSide::Sell,
make_address(0xCC),
)
.with_id("o1".to_string())
}
fn comparison_payloads(responses: &OrderResponses) -> Vec<String> {
capture_comparison(&comparison_order(), responses, &QuoteOptions::default())
}
fn capture_comparison(
order: &Order,
responses: &OrderResponses,
options: &QuoteOptions,
) -> Vec<String> {
crate::worker_pool_router::log_capture::capture_payloads("quote_comparison ", || {
record_quote_comparison(order, responses, options);
})
}
fn responses_with(quotes: Vec<WorkerPoolQuote>) -> OrderResponses {
OrderResponses { order_id: "o1".to_string(), quotes, failed_solvers: vec![] }
}
#[test]
fn test_payload_carries_no_enclosing_span_fields() {
let responses = responses_with(vec![comparison_quote("winner", 1_000, 3)]);
let payloads =
crate::worker_pool_router::log_capture::capture_payloads("quote_comparison ", || {
let request =
tracing::info_span!("HTTP request", request_id = "abcd", http.method = "POST");
let _entered = request.enter();
record_quote_comparison(&comparison_order(), &responses, &QuoteOptions::default());
});
assert_eq!(payloads.len(), 1);
for leaked in ["request_id", "http.method", "abcd"] {
assert!(!payloads[0].contains(leaked), "{leaked} leaked into: {}", payloads[0]);
}
}
#[test]
fn test_comparison_payload_has_no_ansi_escapes() {
let responses = responses_with(vec![comparison_quote("winner", 1_000, 3)]);
let payloads =
capture_comparison(&comparison_order(), &responses, &QuoteOptions::default());
assert_eq!(payloads.len(), 1);
assert!(!payloads[0].contains('\u{1b}'), "escape sequence in payload: {}", payloads[0]);
}
#[test]
fn test_comparison_payload_is_logfmt_tokenised() {
let line = comparison_line_for("winner", 1_000, 3);
for token in line.split_whitespace() {
assert!(token.contains('='), "token `{token}` is not key=value in: {line}");
}
}
fn comparison_line_for(worker_pool: &str, net: u64, solve_time_ms: u64) -> String {
let responses = responses_with(vec![comparison_quote(worker_pool, net, solve_time_ms)]);
comparison_payloads(&responses)
.pop()
.expect("one line per quote")
}
#[test]
fn test_improvement_measured_against_weakest_quote() {
let responses = responses_with(vec![
comparison_quote("winner", 1_000, 3),
comparison_quote("laggard", 900, 11),
]);
let payloads = comparison_payloads(&responses);
assert_eq!(payloads.len(), 2);
assert!(payloads[0].contains("pool=winner"), "{}", payloads[0]);
assert!(payloads[0].contains("improvement_bps=1111.1111"), "{}", payloads[0]);
assert!(payloads[0].contains("is_best=true"), "{}", payloads[0]);
assert!(payloads[1].contains("improvement_bps=0.0000"), "{}", payloads[1]);
assert!(payloads[1].contains("is_best=false"), "{}", payloads[1]);
}
fn recorded_improvements(
responses: &OrderResponses,
) -> Vec<(String, Vec<String>, metrics_util::debugging::DebugValue)> {
let recorder = metrics_util::debugging::DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
record_quote_comparison(&comparison_order(), responses, &QuoteOptions::default());
});
crate::tests::metrics::recorded_metrics(&snapshotter)
.into_iter()
.filter(|(name, ..)| name == "quote_improvement_bps")
.collect()
}
#[test]
fn test_improvement_histogram_is_recorded_without_the_log() {
let responses = responses_with(vec![
comparison_quote("winner", 1_000, 3),
comparison_quote("laggard", 900, 11),
]);
let recorded = recorded_improvements(&responses);
assert_eq!(recorded.len(), 2, "one series per pool: {recorded:?}");
let (_, labels, value) = recorded
.iter()
.find(|(_, labels, _)| labels.contains(&"pool=winner".to_string()))
.expect("the winning pool records its improvement");
assert!(
labels.contains(&"algorithm=winner_algo".to_string()),
"the algorithm is what the comparison groups by: {labels:?}"
);
assert!(
matches!(value, metrics_util::debugging::DebugValue::Histogram(values)
if values.len() == 1 && (values[0].into_inner() - 1_111.111_111_111_111).abs() < 1e-6),
"{value:?}"
);
}
#[test]
fn test_unrankable_quote_records_no_improvement() {
let mut unranked = comparison_quote("no_route", 900, 4);
unranked.quote = OrderQuote::new(
"o1".to_string(),
QuoteStatus::NoRouteFound,
BigUint::from(1_000u64),
BigUint::ZERO,
BigUint::ZERO,
BigUint::ZERO,
BlockInfo::new(42, "0xabc".to_string(), 0),
"no_route_algo".to_string(),
Bytes::default(),
Bytes::default(),
"1".to_string(),
);
let responses = responses_with(vec![comparison_quote("winner", 1_000, 3), unranked]);
let recorded = recorded_improvements(&responses);
assert!(
!recorded
.iter()
.any(|(_, labels, _)| labels.contains(&"pool=no_route".to_string())),
"{recorded:?}"
);
}
#[test]
fn test_solve_time_is_reported_per_worker_pool() {
let responses = responses_with(vec![
comparison_quote("fast", 1_000, 3),
comparison_quote("slow", 900, 417),
]);
let payloads = comparison_payloads(&responses);
assert!(payloads[0].contains("solve_time_ms=3"), "{}", payloads[0]);
assert!(payloads[1].contains("solve_time_ms=417"), "{}", payloads[1]);
}
#[test]
fn test_failed_pool_line_shares_the_schema() {
let responses = OrderResponses {
order_id: "o1".to_string(),
quotes: vec![comparison_quote("winner", 1_000, 3)],
failed_solvers: vec![("slowpoke".to_string(), SolveError::Timeout { elapsed_ms: 500 })],
};
let payloads = comparison_payloads(&responses);
assert_eq!(payloads.len(), 2);
let keys = |payload: &str| -> Vec<String> {
payload
.split_whitespace()
.filter_map(|t| t.split_once('='))
.map(|(k, _)| k.to_string())
.collect()
};
assert_eq!(keys(&payloads[0]), keys(&payloads[1]), "both line kinds need one schema");
assert!(payloads[1].contains("pool=slowpoke"), "{}", payloads[1]);
assert!(payloads[1].contains("status=timeout"), "{}", payloads[1]);
assert!(payloads[1].contains("solve_time_ms=500"), "{}", payloads[1]);
assert!(payloads[1].contains("improvement_bps= "), "{}", payloads[1]);
assert!(payloads[1].contains("responders=2"), "{}", payloads[1]);
}
#[test]
fn test_quote_rejected_by_max_gas_reports_no_improvement() {
let responses = responses_with(vec![
comparison_quote("cheap", 900, 3),
comparison_quote("expensive", 1_000, 4),
]);
let options = QuoteOptions::default().with_max_gas(BigUint::from(1u64));
let payloads = capture_comparison(&comparison_order(), &responses, &options);
assert_eq!(payloads.len(), 2);
for payload in &payloads {
assert!(payload.contains("improvement_bps= "), "{payload}");
assert!(payload.contains("ranked_candidates=0"), "{payload}");
assert!(payload.contains("responders=2"), "{payload}");
}
}
#[test]
fn test_responders_counts_answers_and_failures() {
let responses = OrderResponses {
order_id: "o1".to_string(),
quotes: vec![comparison_quote("a", 1_000, 3), comparison_quote("b", 900, 4)],
failed_solvers: vec![("c".to_string(), SolveError::QueueFull)],
};
let payloads = comparison_payloads(&responses);
for payload in &payloads {
assert!(payload.contains("responders=3"), "{payload}");
assert!(payload.contains("ranked_candidates=2"), "{payload}");
}
}
#[test]
fn test_non_success_quote_reports_no_improvement() {
let mut quote = comparison_quote("stale", 900, 5);
quote.quote = OrderQuote::new(
"o1".to_string(),
QuoteStatus::NotReady,
BigUint::from(1_000u64),
BigUint::ZERO,
BigUint::ZERO,
BigUint::ZERO,
BlockInfo::new(42, "0xabc".to_string(), 0),
"stale_algo".to_string(),
Bytes::default(),
Bytes::default(),
"1".to_string(),
);
let responses = responses_with(vec![comparison_quote("ok", 1_000, 3), quote]);
let payloads = comparison_payloads(&responses);
assert!(payloads[1].contains("status=not_ready"), "{}", payloads[1]);
assert!(payloads[1].contains("improvement_bps= "), "{}", payloads[1]);
}
#[rstest]
#[case(Some(900.0), Some(900.0), Some(0.0))]
#[case(Some(900.0), Some(1_000.0), Some(1_111.111_111_111_111))]
#[case(None, Some(900.0), None)]
#[case(Some(900.0), None, None)]
#[case(Some(0.0), Some(900.0), None)]
#[case(Some(f64::INFINITY), Some(900.0), None)]
#[case(Some(900.0), Some(f64::NAN), None)]
#[case(Some(1_000.0), Some(900.0), Some(-1_000.0))]
fn test_improvement_bps(
#[case] baseline_net: Option<f64>,
#[case] net: Option<f64>,
#[case] expected: Option<f64>,
) {
match (improvement_bps(baseline_net, net), expected) {
(Some(got), Some(want)) => assert!((got - want).abs() < 1e-9, "got {got}, want {want}"),
(got, want) => assert_eq!(got, want),
}
}
#[rstest]
#[case(SolveError::Timeout { elapsed_ms: 7 }, "timeout")]
#[case(SolveError::NoRouteFound { order_id: "o1".to_string(), reason: None }, "no_route")]
#[case(SolveError::QueueFull, "queue_full")]
#[case(SolveError::MaxGasExceeded, "max_gas_exceeded")]
#[case(SolveError::AlgorithmError("boom".to_string()), "algorithm_error")]
#[case(SolveError::MissingData("gas".to_string()), "missing_data")]
#[case(SolveError::SimulationFailed("revert".to_string()), "simulation_failed")]
#[case(SolveError::NotReady("derived".to_string()), "not_ready")]
fn test_solver_error_label(#[case] error: SolveError, #[case] expected: &str) {
assert_eq!(solver_error_label(&error), expected);
}
#[rstest]
#[case(QuoteStatus::Success, "success")]
#[case(QuoteStatus::NoRouteFound, "no_route")]
#[case(QuoteStatus::Timeout, "timeout")]
#[case(QuoteStatus::NotReady, "not_ready")]
#[case(QuoteStatus::InsufficientLiquidity, "insufficient_liquidity")]
#[case(QuoteStatus::PriceCheckFailed, "price_check_failed")]
fn test_quote_status_label(#[case] status: QuoteStatus, #[case] expected: &str) {
assert_eq!(quote_status_label(status), expected);
}
#[test]
fn test_status_vocabularies_agree() {
assert_eq!(
quote_status_label(QuoteStatus::Timeout),
solver_error_label(&SolveError::Timeout { elapsed_ms: 1 })
);
assert_eq!(
quote_status_label(QuoteStatus::NoRouteFound),
solver_error_label(&SolveError::NoRouteFound {
order_id: "o1".to_string(),
reason: None
})
);
}
}