From 762f188e4211012e00b7bf2e9b22ead064152e4f Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 17 Sep 2026 14:08:28 +0000 Subject: [PATCH 01/10] solana-autopilot: price auction tokens through CoinGecko --- crates/autopilot-svm/example.toml | 7 + crates/autopilot-svm/src/domain/arbitrator.rs | 18 +- crates/autopilot-svm/src/domain/auction.rs | 17 +- crates/autopilot-svm/src/infra/config.rs | 44 +++ crates/autopilot-svm/src/infra/db.rs | 6 +- crates/autopilot-svm/src/infra/mod.rs | 1 + crates/autopilot-svm/src/infra/prices.rs | 330 ++++++++++++++++++ crates/autopilot-svm/src/infra/provider.rs | 32 +- crates/autopilot-svm/src/run.rs | 10 + crates/autopilot-svm/src/tests.rs | 45 ++- 10 files changed, 487 insertions(+), 23 deletions(-) create mode 100644 crates/autopilot-svm/src/infra/prices.rs diff --git a/crates/autopilot-svm/example.toml b/crates/autopilot-svm/example.toml index 5cb1aa82c6..24738be3b1 100644 --- a/crates/autopilot-svm/example.toml +++ b/crates/autopilot-svm/example.toml @@ -30,6 +30,13 @@ submission-deadline-slots = 25 name = "baseline" url = "http://localhost:11088" +# Native price lookups for auction tokens. The key is optional: without one +# the public CoinGecko API and its rate limits apply. +[native-prices] +endpoint = "https://api.coingecko.com/api/v3/" +# api-key = "..." +ttl = "30s" + [logging] filter = "info,autopilot_svm=debug" diff --git a/crates/autopilot-svm/src/domain/arbitrator.rs b/crates/autopilot-svm/src/domain/arbitrator.rs index 0deff62447..af7a068c1b 100644 --- a/crates/autopilot-svm/src/domain/arbitrator.rs +++ b/crates/autopilot-svm/src/domain/arbitrator.rs @@ -8,10 +8,7 @@ use { }, run_loop::WinnerSelection, }, - chain_types::{ - ChainTypes, - solana::{Pubkey, Solana}, - }, + chain_types::solana::{Pubkey, Solana}, winner_selection::{Arbitrator, AuctionContext}, }; @@ -54,16 +51,9 @@ impl WinnerSelection for SolanaArbitrator { .iter() .map(|order| (order.uid, Vec::new())) .collect(); - // Every auction token is priced at the native denominator, i.e. 1:1 - // to lamports, so scores compare raw surplus. Ranking within one - // token pair is exact, comparisons across pairs are not. - // TODO: replace with native price estimation. - let native_prices = auction - .orders - .iter() - .flat_map(|order| [order.sell_token, order.buy_token]) - .map(|token| (token, Solana::NATIVE_PRICE_DENOMINATOR)) - .collect(); + // A solution trading a token the auction carries no price for gets + // no score and loses to scored ones. + let native_prices = auction.native_prices.clone(); let context = AuctionContext:: { fee_policies, surplus_capturing_jit_order_owners: Default::default(), diff --git a/crates/autopilot-svm/src/domain/auction.rs b/crates/autopilot-svm/src/domain/auction.rs index 0e1eabeaf6..18c41b3c9c 100644 --- a/crates/autopilot-svm/src/domain/auction.rs +++ b/crates/autopilot-svm/src/domain/auction.rs @@ -4,6 +4,7 @@ use { crate::run_loop::AuctionInfo, chain_types::solana::{AppData, IntentHash, Pubkey}, + std::collections::HashMap, }; /// Whether the order sells an exact amount or buys an exact amount. @@ -45,6 +46,10 @@ pub struct Auction { /// cuts by content, and the id is allocated only for a fresh cut. pub id: i64, pub orders: Vec, + /// The lamport value of one atom of each auction token, scaled by 10^9. + /// Tokens without a listing are absent. Excluded from equality like the + /// id: prices refresh between cuts without making the auction new. + pub native_prices: HashMap, } impl PartialEq for Auction { @@ -62,7 +67,7 @@ impl AuctionInfo for Auction { #[cfg(test)] mod tests { use { - super::{Auction, Order, OrderKind}, + super::{Auction, HashMap, Order, OrderKind, Pubkey}, chain_types::solana::AppData, }; @@ -86,19 +91,25 @@ mod tests { } #[test] - fn auction_equality_ignores_id() { + fn auction_equality_ignores_id_and_prices() { let orders = vec![order(10)]; let a = Auction { id: 1, orders: orders.clone(), + native_prices: HashMap::new(), + }; + let b = Auction { + id: 2, + orders, + native_prices: HashMap::from([(Pubkey([0x11; 32]), 7)]), }; - let b = Auction { id: 2, orders }; assert_eq!(a, b); assert_ne!( a, Auction { id: 1, orders: vec![], + native_prices: HashMap::new(), } ); } diff --git a/crates/autopilot-svm/src/infra/config.rs b/crates/autopilot-svm/src/infra/config.rs index dfa7d30e2b..83c23f4f4f 100644 --- a/crates/autopilot-svm/src/infra/config.rs +++ b/crates/autopilot-svm/src/infra/config.rs @@ -62,11 +62,49 @@ pub struct Config { /// sponsored orders: without it their winning solutions dispatch without /// creations and fail at the driver. pub sponsoring: Option, + /// Native price lookups for auction tokens. + #[serde(default)] + pub native_prices: NativePrices, /// Logging configuration. #[serde(default)] pub logging: LoggingConfig, } +/// CoinGecko native price lookups. +#[derive(Debug, Deserialize)] +#[serde(rename_all = "kebab-case", deny_unknown_fields)] +pub struct NativePrices { + /// Base URL of the CoinGecko API. + #[serde(default = "default_prices_endpoint")] + pub endpoint: url::Url, + /// API key sent with every price request, for keyed CoinGecko plans. + #[serde(default)] + pub api_key: Option, + /// How long a fetched price serves auctions before it is refetched. + #[serde(with = "humantime_serde", default = "default_prices_ttl")] + pub ttl: Duration, +} + +impl Default for NativePrices { + fn default() -> Self { + Self { + endpoint: default_prices_endpoint(), + api_key: None, + ttl: default_prices_ttl(), + } + } +} + +fn default_prices_endpoint() -> url::Url { + "https://api.coingecko.com/api/v3/" + .parse() + .expect("valid literal url") +} + +const fn default_prices_ttl() -> Duration { + Duration::from_secs(30) +} + impl Config { /// Build the `observe::Config` for the tracing framework from the logging /// configuration. @@ -176,6 +214,12 @@ mod tests { assert_eq!(config.competition.submission_deadline_slots.get(), 25); assert_eq!(config.max_auction_age, Duration::from_secs(5 * 60)); assert_eq!(config.min_auction_interval, Duration::from_secs(2)); + assert_eq!( + config.native_prices.endpoint.as_str(), + "https://api.coingecko.com/api/v3/" + ); + assert_eq!(config.native_prices.api_key, None); + assert_eq!(config.native_prices.ttl, Duration::from_secs(30)); assert_eq!(config.drivers.len(), 1); assert_eq!(config.drivers[0].name, "baseline"); assert_eq!(config.logging.filter, "info,autopilot_svm=debug"); diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 7f6f34e343..41e54a355b 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -209,7 +209,11 @@ pub async fn cut( block_height: Option, ) -> Result { let orders = orders_from_rows(open_orders(ex, now_unix, block_height).await?); - Ok(Auction { id, orders }) + Ok(Auction { + id, + orders, + native_prices: Default::default(), + }) } /// A row the indexer wrote always converts (on-chain values fit the domain diff --git a/crates/autopilot-svm/src/infra/mod.rs b/crates/autopilot-svm/src/infra/mod.rs index 4338c4f20f..08dceb60fd 100644 --- a/crates/autopilot-svm/src/infra/mod.rs +++ b/crates/autopilot-svm/src/infra/mod.rs @@ -9,6 +9,7 @@ pub mod listen; pub mod observation; pub mod observer; pub mod order_events; +pub mod prices; pub mod provider; pub mod sponsor; pub mod trigger; diff --git a/crates/autopilot-svm/src/infra/prices.rs b/crates/autopilot-svm/src/infra/prices.rs new file mode 100644 index 0000000000..ab07d2c5a9 --- /dev/null +++ b/crates/autopilot-svm/src/infra/prices.rs @@ -0,0 +1,330 @@ +//! Native token prices from CoinGecko, denominated in wSOL. + +use { + crate::infra::config, + anyhow::{Context, Result, anyhow}, + chain_types::{ChainTypes, solana::Solana}, + cow_solana_rpc::SolanaRPC, + serde::Deserialize, + solana_sdk::{program_pack::Pack, pubkey::Pubkey}, + spl_token_interface::state::Mint, + std::{ + collections::{HashMap, HashSet}, + sync::Mutex, + time::{Duration, Instant}, + }, + url::Url, +}; + +/// CoinGecko's authorization header. +const API_KEY_HEADER: &str = "x-cg-pro-api-key"; + +/// Native price lookups for auction tokens, cached per token. +pub struct NativePrices { + client: reqwest::Client, + endpoint: Url, + api_key: Option, + rpc: SolanaRPC, + wrapped_native: Pubkey, + ttl: Duration, + /// Fetched prices by mint. `None` records a mint CoinGecko does not + /// list, so unlisted mints are not refetched every cut. + prices: Mutex)>>, + /// Mint decimals never change, so they are cached forever. + decimals: Mutex>, +} + +/// One priced entry of the CoinGecko `simple/token_price` response. +#[derive(Debug, Deserialize)] +struct Entry { + sol: Option, +} + +impl NativePrices { + pub fn new(config: &config::NativePrices, rpc: SolanaRPC, wrapped_native: Pubkey) -> Self { + Self { + client: reqwest::Client::new(), + endpoint: config.endpoint.clone(), + api_key: config.api_key.clone(), + rpc, + wrapped_native, + ttl: config.ttl, + prices: Mutex::new(HashMap::new()), + decimals: Mutex::new(HashMap::new()), + } + } + + /// A lookup pre-seeded for tests: the given prices never expire and + /// nothing is fetched. + #[cfg(test)] + pub(crate) fn seeded(entries: impl IntoIterator) -> Self { + Self { + client: reqwest::Client::new(), + endpoint: "http://127.0.0.1:1/".parse().expect("literal url"), + api_key: None, + rpc: SolanaRPC::new_mock_with_mocks(Default::default()), + wrapped_native: Pubkey::default(), + ttl: Duration::from_secs(u64::MAX), + prices: Mutex::new( + entries + .into_iter() + .map(|(token, price)| (token, (Instant::now(), Some(price)))) + .collect(), + ), + decimals: Mutex::new(HashMap::new()), + } + } + + /// The lamport value of one atom of each token, scaled by 10^9 like + /// [`ChainTypes::NATIVE_PRICE_DENOMINATOR`]. Tokens CoinGecko does not + /// list are absent from the result. Any lookup failure fails the whole + /// call: a partially priced auction would rank solutions on incomparable + /// scores. + pub async fn prices(&self, tokens: HashSet) -> Result> { + let mut result = HashMap::new(); + let mut fetch = Vec::new(); + let now = Instant::now(); + { + let cache = self.prices.lock().expect("price cache poisoned"); + for token in tokens { + if token == self.wrapped_native { + result.insert(token, Solana::NATIVE_PRICE_DENOMINATOR); + continue; + } + match cache.get(&token) { + Some((fetched, price)) if now.duration_since(*fetched) < self.ttl => { + if let Some(price) = price { + result.insert(token, *price); + } + } + _ => fetch.push(token), + } + } + } + if fetch.is_empty() { + return Ok(result); + } + + let decimals = self.decimals(&fetch).await?; + let quoted = self.fetch(&fetch).await?; + let mut cache = self.prices.lock().expect("price cache poisoned"); + for token in fetch { + let price = quoted + .get(&token.to_string()) + .and_then(|entry| entry.sol) + .and_then(|sol| scale(sol, decimals[&token])); + cache.insert(token, (now, price)); + if let Some(price) = price { + result.insert(token, price); + } + } + Ok(result) + } + + /// Decimals per mint, from the cache or the mint accounts on chain. + async fn decimals(&self, tokens: &[Pubkey]) -> Result> { + let mut result = HashMap::new(); + let mut fetch = Vec::new(); + { + let cache = self.decimals.lock().expect("decimals cache poisoned"); + for token in tokens { + match cache.get(token) { + Some(decimals) => { + result.insert(*token, *decimals); + } + None => fetch.push(*token), + } + } + } + if fetch.is_empty() { + return Ok(result); + } + let accounts = self + .rpc + .multiple_accounts(fetch.iter().copied()) + .await + .context("fetch mint accounts")?; + let mut cache = self.decimals.lock().expect("decimals cache poisoned"); + for token in fetch { + let account = accounts + .get(&token) + .ok_or_else(|| anyhow!("mint {token} does not exist"))?; + let mint = Mint::unpack(&account.data) + .with_context(|| format!("mint {token} does not unpack"))?; + cache.insert(token, mint.decimals); + result.insert(token, mint.decimals); + } + Ok(result) + } + + /// One `simple/token_price` request for the given mints. + async fn fetch(&self, tokens: &[Pubkey]) -> Result> { + let mut url = self + .endpoint + .join("simple/token_price/solana") + .context("price endpoint")?; + let addresses = tokens + .iter() + .map(ToString::to_string) + .collect::>() + .join(","); + url.query_pairs_mut() + .append_pair("contract_addresses", &addresses) + .append_pair("vs_currencies", "sol"); + let mut request = self.client.get(url); + if let Some(key) = &self.api_key { + request = request.header(API_KEY_HEADER, key); + } + let response = request.send().await.context("price request")?; + let status = response.status(); + if !status.is_success() { + let body = response.text().await.unwrap_or_default(); + return Err(anyhow!("price request answered {status}: {body}")); + } + response.json().await.context("price response") + } +} + +/// A whole-token price in SOL converted to the scaled atom price: +/// `price * 10^(18 - decimals)`, saturating at `u64::MAX`. `None` when the +/// price rounds below one, those tokens count as unpriced. +fn scale(price: f64, decimals: u8) -> Option { + let scaled = price * 10f64.powi(18 - i32::from(decimals)); + if scaled.is_nan() || scaled < 1.0 { + return None; + } + Some(if scaled >= u64::MAX as f64 { + u64::MAX + } else { + scaled as u64 + }) +} + +#[cfg(test)] +mod tests { + use { + super::*, + cow_solana_rpc::{Mocks, RpcRequest}, + std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + }; + + fn config(endpoint: Url) -> config::NativePrices { + config::NativePrices { + endpoint, + api_key: None, + ttl: Duration::from_secs(60), + } + } + + /// Serve a fixed CoinGecko response, counting the requests. + async fn coingecko(response: serde_json::Value) -> (Url, Arc) { + let requests = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&requests); + let app = axum::Router::new().route( + "/simple/token_price/solana", + axum::routing::get(move || { + counter.fetch_add(1, Ordering::Relaxed); + let response = response.clone(); + async move { axum::Json(response) } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + (format!("http://{addr}/").parse().unwrap(), requests) + } + + #[test] + fn scales_whole_token_prices_to_atoms() { + // A 6-decimals token at 0.005 SOL: 0.005 * 10^12. + assert_eq!(scale(0.005, 6), Some(5_000_000_000)); + // The native token itself: 1.0 * 10^9. + assert_eq!(scale(1.0, 9), Some(Solana::NATIVE_PRICE_DENOMINATOR)); + assert_eq!(scale(0.0, 6), None); + assert_eq!(scale(f64::NAN, 6), None); + assert_eq!(scale(f64::MAX, 0), Some(u64::MAX)); + } + + /// The wrapped native mint is priced at the denominator without any + /// lookup: the endpoint and the RPC here are dead. + #[tokio::test] + async fn prices_the_native_mint_locally() { + let wrapped = Pubkey::new_unique(); + let prices = NativePrices::new( + &config("http://127.0.0.1:1/".parse().unwrap()), + SolanaRPC::new_mock_with_mocks(Mocks::default()), + wrapped, + ); + let result = prices.prices(HashSet::from([wrapped])).await.unwrap(); + assert_eq!(result[&wrapped], Solana::NATIVE_PRICE_DENOMINATOR); + } + + /// A listed token is priced through its decimals, an unlisted one is + /// absent, and the second call is served from the cache. + #[tokio::test] + async fn prices_and_caches_auction_tokens() { + let listed = Pubkey::new_unique(); + let unlisted = Pubkey::new_unique(); + let (endpoint, requests) = coingecko(serde_json::json!({ + listed.to_string(): { "sol": 0.005 }, + })) + .await; + let mocks = Mocks::from([( + RpcRequest::GetMultipleAccounts, + serde_json::json!({ + "context": {"slot": 1u64, "apiVersion": "2.0.0"}, + "value": [ + crate::tests::mint_account_json(6), + crate::tests::mint_account_json(6), + ], + }), + )]); + let prices = NativePrices::new( + &config(endpoint), + SolanaRPC::new_mock_with_mocks(mocks), + Pubkey::new_unique(), + ); + + let result = prices + .prices(HashSet::from([listed, unlisted])) + .await + .unwrap(); + assert_eq!(result.get(&listed), Some(&5_000_000_000)); + assert_eq!(result.get(&unlisted), None); + + // Both mints are cached, the unlisted one negatively: the RPC mock + // is consumed and the endpoint sees no second request. + let result = prices + .prices(HashSet::from([listed, unlisted])) + .await + .unwrap(); + assert_eq!(result.get(&listed), Some(&5_000_000_000)); + assert_eq!(requests.load(Ordering::Relaxed), 1); + } + + /// A failing price endpoint fails the whole lookup. + #[tokio::test] + async fn fails_closed_on_endpoint_errors() { + let mocks = Mocks::from([( + RpcRequest::GetMultipleAccounts, + serde_json::json!({ + "context": {"slot": 1u64, "apiVersion": "2.0.0"}, + "value": [crate::tests::mint_account_json(6)], + }), + )]); + let prices = NativePrices::new( + &config("http://127.0.0.1:1/".parse().unwrap()), + SolanaRPC::new_mock_with_mocks(mocks), + Pubkey::new_unique(), + ); + assert!( + prices + .prices(HashSet::from([Pubkey::new_unique()])) + .await + .is_err() + ); + } +} diff --git a/crates/autopilot-svm/src/infra/provider.rs b/crates/autopilot-svm/src/infra/provider.rs index ab8d3aeac4..89cde0b04e 100644 --- a/crates/autopilot-svm/src/infra/provider.rs +++ b/crates/autopilot-svm/src/infra/provider.rs @@ -3,10 +3,11 @@ use { crate::{ domain::{auction::Order, cycle::SolanaCycle}, - infra::db, + infra::{db, prices::NativePrices}, run_loop::AuctionProvider, }, async_trait::async_trait, + chain_types::solana::Pubkey as ChainPubkey, cow_solana_rpc::SolanaRPC, solana_sdk::{account::Account, program_pack::Pack, pubkey::Pubkey}, spl_token_interface::state::{Account as TokenAccount, AccountState}, @@ -21,6 +22,7 @@ use { pub struct DbAuctionProvider { pool: PgPool, rpc: SolanaRPC, + prices: NativePrices, /// Last allocated auction id. Ids are unix seconds, bumped past the /// previous allocation when cycles land within the same second. Unique /// only per process: no table allocates auction ids. @@ -30,10 +32,11 @@ pub struct DbAuctionProvider { } impl DbAuctionProvider { - pub fn new(pool: PgPool, rpc: SolanaRPC) -> Self { + pub fn new(pool: PgPool, rpc: SolanaRPC, prices: NativePrices) -> Self { Self { pool, rpc, + prices, last_id: AtomicI64::new(0), } } @@ -129,7 +132,29 @@ impl AuctionProvider for DbAuctionProvider { .map_err(|err| tracing::warn!(?err, "failed to cut the auction")) .ok()?; auction.orders = self.receivable_orders(auction.orders).await; - (!auction.orders.is_empty()).then_some(auction) + if auction.orders.is_empty() { + return None; + } + // A cut without prices would rank solutions on incomparable scores, + // so a failed lookup skips the cycle instead. + let tokens = auction + .orders + .iter() + .flat_map(|order| [order.sell_token, order.buy_token]) + .map(|token| Pubkey::new_from_array(token.0)) + .collect(); + let prices = match self.prices.prices(tokens).await { + Ok(prices) => prices, + Err(err) => { + tracing::warn!(?err, "native price lookup failed, skipping the cut"); + return None; + } + }; + auction.native_prices = prices + .into_iter() + .map(|(token, price)| (ChainPubkey(token.to_bytes()), price)) + .collect(); + Some(auction) } } @@ -188,6 +213,7 @@ mod tests { DbAuctionProvider::new( sqlx::PgPool::connect_lazy("postgresql://").unwrap(), SolanaRPC::new_mock_with_mocks(mocks), + NativePrices::seeded([]), ) } diff --git a/crates/autopilot-svm/src/run.rs b/crates/autopilot-svm/src/run.rs index 426cf00fd5..aaeeba2290 100644 --- a/crates/autopilot-svm/src/run.rs +++ b/crates/autopilot-svm/src/run.rs @@ -12,6 +12,7 @@ use { listen::ListenSession, observation::SettlementWindows, observer::CompetitionObserver, + prices::NativePrices, provider::DbAuctionProvider, sponsor::Sponsor, trigger::SlotTrigger, @@ -129,6 +130,15 @@ async fn run(config: Config) { config.rpc.request_timeout, CommitmentConfig::confirmed(), ), + NativePrices::new( + &config.native_prices, + SolanaRPC::new_with_timeout_and_commitment( + &config.rpc.endpoint, + config.rpc.request_timeout, + CommitmentConfig::confirmed(), + ), + config.contracts.wrapped_native_mint, + ), )), Box::new(DriverCompetition::new( drivers.clone(), diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index a4f6922f84..b32b5ec784 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -10,6 +10,7 @@ use { executor::DriverExecutor, observation::SettlementWindows, observer::CompetitionObserver, + prices::NativePrices, provider::DbAuctionProvider, }, run_loop::{ @@ -87,6 +88,23 @@ async fn spawn_mock_driver(state: MockDriverState) -> SocketAddr { addr } +/// A canned `getMultipleAccounts` entry: an initialized mint of the classic +/// SPL token program with the given decimals. +pub(crate) fn mint_account_json(decimals: u8) -> serde_json::Value { + let mut data = [0u8; 82]; + data[44] = decimals; + // The initialized flag. + data[45] = 1; + serde_json::json!({ + "lamports": 1_461_600u64, + "data": [BASE64_STANDARD.encode(data), "base64"], + "owner": "TokenkegQfeZyiNwAJbNbGKPFXCWuBvf9Ss623VQ5DA", + "executable": false, + "rentEpoch": 0u64, + "space": 82u64, + }) +} + /// A canned `getMultipleAccounts` entry: an initialized account of the /// classic SPL token program holding `mint`. pub(crate) fn token_account_json(mint: [u8; 32]) -> serde_json::Value { @@ -115,6 +133,21 @@ fn mock_rpc() -> SolanaRPC { SolanaRPC::new_mock_with_mocks(Mocks::from([(RpcRequest::GetMultipleAccounts, response)])) } +/// Native prices for the seeded order's pair, at the denominator so scores +/// stay raw surplus. +fn test_prices() -> [(solana_sdk::pubkey::Pubkey, u64); 2] { + [ + ( + solana_sdk::pubkey::Pubkey::new_from_array([0xAA; 32]), + 1_000_000_000, + ), + ( + solana_sdk::pubkey::Pubkey::new_from_array([0xAB; 32]), + 1_000_000_000, + ), + ] +} + async fn seed_open_order(pool: &PgPool, uid: [u8; 32], tip: i64) { crate::test_db::wipe(pool).await; sqlx::query("INSERT INTO solana.indexer_state (slot) VALUES ($1)") @@ -178,7 +211,11 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { // Stage probes: pinpoint the failing phase before driving the loop. { - let provider = DbAuctionProvider::new(pool.clone(), mock_rpc()); + let provider = DbAuctionProvider::new( + pool.clone(), + mock_rpc(), + NativePrices::seeded(test_prices()), + ); let auction = provider.cut_auction(&tip).await.expect("auction cut"); assert_eq!(auction.orders.len(), 1, "open order in the auction"); let competition = DriverCompetition::new(vec![Arc::clone(&driver)], Duration::from_secs(6)); @@ -191,7 +228,11 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { let windows = SettlementWindows::new(pool.clone()); let mut auction_loop = AuctionLoop::new( Box::new(FixedTrigger(tip)), - Box::new(DbAuctionProvider::new(pool.clone(), mock_rpc())), + Box::new(DbAuctionProvider::new( + pool.clone(), + mock_rpc(), + NativePrices::seeded(test_prices()), + )), Box::new(DriverCompetition::new( vec![Arc::clone(&driver)], Duration::from_secs(6), From 1b4d36e1144de41a5788219282187d995ecb1c8b Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 17 Sep 2026 16:58:13 +0000 Subject: [PATCH 02/10] solana-autopilot: persist competitions and hold out in-flight orders --- crates/autopilot-svm/src/domain/arbitrator.rs | 20 ++- crates/autopilot-svm/src/domain/cycle.rs | 4 + crates/autopilot-svm/src/infra/db.rs | 167 ++++++++++++++++-- crates/autopilot-svm/src/infra/executor.rs | 10 +- crates/autopilot-svm/src/infra/observation.rs | 7 +- crates/autopilot-svm/src/infra/observer.rs | 64 +++++-- crates/autopilot-svm/src/infra/provider.rs | 88 +++++---- crates/autopilot-svm/src/test_db.rs | 4 +- crates/autopilot-svm/src/tests.rs | 31 ++++ database/sql-solana/V8__competition.sql | 42 +++++ 10 files changed, 366 insertions(+), 71 deletions(-) create mode 100644 database/sql-solana/V8__competition.sql diff --git a/crates/autopilot-svm/src/domain/arbitrator.rs b/crates/autopilot-svm/src/domain/arbitrator.rs index af7a068c1b..a2d3d5f76e 100644 --- a/crates/autopilot-svm/src/domain/arbitrator.rs +++ b/crates/autopilot-svm/src/domain/arbitrator.rs @@ -70,6 +70,24 @@ impl WinnerSelection for SolanaArbitrator { for winner in inner.winners() { tracing::info!(solver = %winner.solver(), solution = winner.id(), "winner"); } - Ranking { inner, drivers } + // Ranked before filtered, so winner uids come first. + let uids = inner + .ranked + .iter() + .map(|solution| (solution.solver(), solution.id())) + .chain( + inner + .filtered_out + .iter() + .map(|solution| (solution.solver(), solution.id())), + ) + .enumerate() + .map(|(uid, key)| (key, i64::try_from(uid).unwrap_or(i64::MAX))) + .collect(); + Ranking { + inner, + drivers, + uids, + } } } diff --git a/crates/autopilot-svm/src/domain/cycle.rs b/crates/autopilot-svm/src/domain/cycle.rs index 99212ca026..eba71ea948 100644 --- a/crates/autopilot-svm/src/domain/cycle.rs +++ b/crates/autopilot-svm/src/domain/cycle.rs @@ -45,6 +45,10 @@ pub struct Ranking { pub inner: winner_selection::Ranking, /// Driver index per solution, keyed by `(solver, solution id)`. pub drivers: HashMap, + /// Autopilot-generated solution uid per solution, unique within the + /// auction. It disambiguates solver-assigned ids across drivers, and + /// everything persisted references it. + pub uids: HashMap, } impl RankingInfo for Ranking { diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 41e54a355b..594c50ca76 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -1,7 +1,7 @@ //! Database access for the Solana autopilot. use { - crate::domain::auction::{Auction, Order, OrderKind}, + crate::domain::auction::{Order, OrderKind}, anyhow::{Context, Result}, bigdecimal::{BigDecimal, ToPrimitive}, chain_types::solana::{AppData, IntentHash, Pubkey}, @@ -34,10 +34,13 @@ pub struct OrderRow { /// Orders open for solving: unexpired, settleable by a driver, not cancelled /// and not fully filled. A pending sponsored order whose stored creation /// transaction died at `block_height` is excluded, and a `None` height skips -/// that check rather than excluding everything. Settleable means the driver can -/// produce the order PDA: it already exists on chain (an order placed via -/// `CreateOrder` directly), or the driver can create it at settlement time from -/// a signed intent or a presigned transaction. +/// that check rather than excluding everything. An order inside a dispatched +/// settlement whose execution window is still open sits out until the window +/// closes, so an in-flight settlement is not raced by the next auction. +/// Settleable means the driver can produce the order PDA: it already exists +/// on chain (an order placed via `CreateOrder` directly), or the driver can +/// create it at settlement time from a signed intent or a presigned +/// transaction. pub async fn open_orders( ex: impl PgExecutor<'_>, now_unix: i64, @@ -61,6 +64,13 @@ WHERE o.valid_to >= $1 OR o.presigned_transaction IS NULL OR p.order_uid IS NOT NULL OR o.last_valid_block_height >= $2) + AND NOT EXISTS ( + SELECT 1 + FROM solana.settlement_executions se + JOIN solana.proposed_trade_executions pte + ON pte.auction_id = se.auction_id AND pte.solution_uid = se.solution_uid + WHERE pte.order_uid = o.uid AND se.outcome IS NULL + ) AND COALESCE( CASE o.kind WHEN 'sell' THEN p.amount_withdrawn < o.sell_amount @@ -201,19 +211,118 @@ pub async fn open_window_auction_ids(ex: impl PgExecutor<'_>) -> Result .context("read open settlement execution windows") } -/// Cut an auction from the open orders. +/// The solvable orders for a fresh auction cut. pub async fn cut( ex: impl PgExecutor<'_>, - id: i64, now_unix: i64, block_height: Option, -) -> Result { - let orders = orders_from_rows(open_orders(ex, now_unix, block_height).await?); - Ok(Auction { - id, - orders, - native_prices: Default::default(), - }) +) -> Result> { + Ok(orders_from_rows( + open_orders(ex, now_unix, block_height).await?, + )) +} + +/// Replace the current auction and answer the id its identity column +/// allocated, the source of sequential auction ids. +pub async fn replace_current_auction( + pool: &sqlx::PgPool, + tip_slot: i64, + json: &serde_json::Value, +) -> Result { + let mut tx = pool.begin().await.context("begin auction replacement")?; + sqlx::query("DELETE FROM solana.auctions") + .execute(&mut *tx) + .await + .context("delete the previous auction")?; + let id = sqlx::query_scalar( + "INSERT INTO solana.auctions (tip_slot, json) VALUES ($1, $2) RETURNING id", + ) + .bind(tip_slot) + .bind(sqlx::types::Json(json)) + .fetch_one(&mut *tx) + .await + .context("insert the current auction")?; + tx.commit().await.context("commit auction replacement")?; + Ok(id) +} + +/// One execution inside a proposed solution. +pub struct ProposedTrade { + pub order_uid: ByteArray<32>, + pub executed_sell: BigDecimal, + pub executed_buy: BigDecimal, +} + +/// One proposed solution of a competition, with its executions. +pub struct ProposedSolution { + /// Autopilot-generated, unique within the auction. + pub uid: i64, + /// Solver-assigned, unique only within one driver response. + pub id: i64, + pub solver: ByteArray<32>, + pub is_winner: bool, + pub score: BigDecimal, + pub trades: Vec, +} + +/// A competition outcome as persisted after ranking. +pub struct Competition { + pub auction_id: i64, + pub tip_slot: i64, + pub deadline_slot: i64, + pub order_uids: Vec>, + pub price_tokens: Vec>, + pub price_values: Vec, + pub solutions: Vec, +} + +/// Persist a competition: the auction snapshot and every proposed solution +/// with its executions, in one transaction. +pub async fn persist_competition(pool: &sqlx::PgPool, competition: &Competition) -> Result<()> { + let mut tx = pool.begin().await.context("begin competition persist")?; + sqlx::query( + "INSERT INTO solana.competition_auctions (id, tip_slot, deadline_slot, order_uids, \ + price_tokens, price_values) VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(competition.auction_id) + .bind(competition.tip_slot) + .bind(competition.deadline_slot) + .bind(&competition.order_uids) + .bind(&competition.price_tokens) + .bind(&competition.price_values) + .execute(&mut *tx) + .await + .context("insert competition auction")?; + for solution in &competition.solutions { + sqlx::query( + "INSERT INTO solana.proposed_solutions (auction_id, uid, id, solver, is_winner, \ + score) VALUES ($1, $2, $3, $4, $5, $6)", + ) + .bind(competition.auction_id) + .bind(solution.uid) + .bind(solution.id) + .bind(solution.solver) + .bind(solution.is_winner) + .bind(&solution.score) + .execute(&mut *tx) + .await + .context("insert proposed solution")?; + for trade in &solution.trades { + sqlx::query( + "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, \ + order_uid, executed_sell, executed_buy) VALUES ($1, $2, $3, $4, $5)", + ) + .bind(competition.auction_id) + .bind(solution.uid) + .bind(trade.order_uid) + .bind(&trade.executed_sell) + .bind(&trade.executed_buy) + .execute(&mut *tx) + .await + .context("insert proposed trade execution")?; + } + } + tx.commit().await.context("commit competition persist") } /// A row the indexer wrote always converts (on-chain values fit the domain @@ -432,6 +541,36 @@ WHERE uid = $1 assert_eq!(uids(orders), vec![1, 5, 6, 10]); let orders = open_orders(&mut *tx, 1_000, Some(151)).await.unwrap(); assert_eq!(uids(orders), vec![1, 5, 6]); + + // An order inside a dispatched settlement sits out while its window + // is open and returns once the window closes. + sqlx::query( + "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, order_uid, \ + executed_sell, executed_buy) VALUES (77, 0, $1, 10, 20)", + ) + .bind(ByteArray([1u8; 32])) + .execute(&mut *tx) + .await + .unwrap(); + sqlx::query( + "INSERT INTO solana.settlement_executions (auction_id, solver, solution_uid, \ + start_timestamp, start_slot, deadline_slot) VALUES (77, $1, 0, now(), 1, 100)", + ) + .bind(ByteArray([0xEE; 32])) + .execute(&mut *tx) + .await + .unwrap(); + let orders = open_orders(&mut *tx, 1_000, None).await.unwrap(); + assert_eq!(uids(orders), vec![5, 6, 10]); + sqlx::query( + "UPDATE solana.settlement_executions SET outcome = 'timeout', end_slot = 100, \ + end_timestamp = now() WHERE auction_id = 77", + ) + .execute(&mut *tx) + .await + .unwrap(); + let orders = open_orders(&mut *tx, 1_000, None).await.unwrap(); + assert_eq!(uids(orders), vec![1, 5, 6, 10]); } #[tokio::test] diff --git a/crates/autopilot-svm/src/infra/executor.rs b/crates/autopilot-svm/src/infra/executor.rs index d4962e5f71..c4cecc2e8b 100644 --- a/crates/autopilot-svm/src/infra/executor.rs +++ b/crates/autopilot-svm/src/infra/executor.rs @@ -85,10 +85,16 @@ impl SettlementExecutor for DriverExecutor { creations, }; // A window that cannot be opened must not block the settlement, - // the dispatch is the priority. + // the dispatch is the priority. The window carries the generated + // solution uid, the driver request keeps the driver-local id its + // solution cache is keyed by. + let uid = ranking.uids.get(&key).copied().unwrap_or_else(|| { + tracing::error!(solution_id = winner.id(), "winner without a uid"); + i64::MAX + }); if let Err(err) = self .windows - .open_dispatched(auction_id, winner.solver(), winner.id(), *tip, deadline) + .open_dispatched(auction_id, winner.solver(), uid, *tip, deadline) .await { tracing::error!(auction_id, ?err, "failed to open the settlement window"); diff --git a/crates/autopilot-svm/src/infra/observation.rs b/crates/autopilot-svm/src/infra/observation.rs index 66783646b7..4e050521e6 100644 --- a/crates/autopilot-svm/src/infra/observation.rs +++ b/crates/autopilot-svm/src/infra/observation.rs @@ -39,13 +39,12 @@ impl SettlementWindows { } /// Open a window for a dispatched settlement. `solution_uid` is the - /// winner's driver-local solution id until competition persistence - /// allocates uids. + /// autopilot-generated uid the competition persisted. pub async fn open_dispatched( &self, auction_id: i64, solver: Pubkey, - solution_uid: u64, + solution_uid: i64, start_slot: u64, deadline_slot: u64, ) -> Result<()> { @@ -53,7 +52,7 @@ impl SettlementWindows { &self.pool, auction_id, solver, - to_db_integer(solution_uid), + solution_uid, to_db_integer(start_slot), to_db_integer(deadline_slot), ) diff --git a/crates/autopilot-svm/src/infra/observer.rs b/crates/autopilot-svm/src/infra/observer.rs index dafbe02c90..9cbd76aa36 100644 --- a/crates/autopilot-svm/src/infra/observer.rs +++ b/crates/autopilot-svm/src/infra/observer.rs @@ -1,22 +1,19 @@ -//! Competition bookkeeping. Auction progress is written to -//! `solana.order_events`, everything else is logged only: there are no -//! competition tables (auction snapshots, proposed executions) to write. -//! -//! TODO: persist the competition outcome once the tables exist. The -//! settlement attribution (`solana.settlements.solution_uid`) depends on -//! the persisted ranking. +//! Competition bookkeeping: auction progress in `solana.order_events`, the +//! ranked outcome in the competition tables. use { crate::{ domain::{auction::Auction, cycle::Ranking}, - infra::{observation::SettlementWindows, order_events}, + infra::{db, observation::SettlementWindows, order_events}, run_loop::SettlementObserver, }, async_trait::async_trait, + bigdecimal::BigDecimal, chain_types::solana::IntentHash, - database::solana::OrderEventLabel, + database::{byte_array::ByteArray, solana::OrderEventLabel}, sqlx::PgPool, std::{collections::HashSet, sync::Mutex}, + winner_selection::state::RankedItem, }; /// Writes order events, logs the competition phases, and drives the @@ -82,7 +79,7 @@ impl SettlementObserver for CompetitionObserv async fn persist_competition_ranking( &self, - _auction: &Auction, + auction: &Auction, tip: &u64, ranking: &Ranking, deadline: u64, @@ -92,13 +89,58 @@ impl SettlementObserver for CompetitionObserv if let Err(err) = self.windows.expire_past_deadline(*tip).await { tracing::error!(?err, "failed to flag expired settlement windows"); } + let solutions = ranking + .inner + .ranked + .iter() + .chain(ranking.inner.filtered_out.iter()) + .map(|solution| { + let key = (solution.solver(), solution.id()); + db::ProposedSolution { + uid: ranking.uids.get(&key).copied().unwrap_or(i64::MAX), + id: i64::try_from(solution.id()).unwrap_or(i64::MAX), + solver: ByteArray(solution.solver().0), + is_winner: solution.is_winner(), + score: BigDecimal::from(solution.score()), + trades: solution + .orders() + .iter() + .map(|order| db::ProposedTrade { + order_uid: ByteArray(order.uid.0), + executed_sell: BigDecimal::from(order.executed_sell), + executed_buy: BigDecimal::from(order.executed_buy), + }) + .collect(), + } + }) + .collect(); + let (price_tokens, price_values) = auction + .native_prices + .iter() + .map(|(token, price)| (token.0.to_vec(), BigDecimal::from(*price))) + .unzip(); + let competition = db::Competition { + auction_id: auction.id, + tip_slot: i64::try_from(*tip).unwrap_or(i64::MAX), + deadline_slot: i64::try_from(deadline).unwrap_or(i64::MAX), + order_uids: auction + .orders + .iter() + .map(|order| order.uid.0.to_vec()) + .collect(), + price_tokens, + price_values, + solutions, + }; + db::persist_competition(&self.pool, &competition).await?; tracing::info!( + auction_id = auction.id, tip, deadline, winners = ranking.inner.winners().count(), ranked = ranking.inner.ranked.len(), filtered_out = ranking.inner.filtered_out.len(), - "competition ranked" + "competition persisted" ); Ok(()) } diff --git a/crates/autopilot-svm/src/infra/provider.rs b/crates/autopilot-svm/src/infra/provider.rs index 89cde0b04e..c0ae0abd2a 100644 --- a/crates/autopilot-svm/src/infra/provider.rs +++ b/crates/autopilot-svm/src/infra/provider.rs @@ -12,10 +12,7 @@ use { solana_sdk::{account::Account, program_pack::Pack, pubkey::Pubkey}, spl_token_interface::state::{Account as TokenAccount, AccountState}, sqlx::PgPool, - std::{ - sync::atomic::{AtomicI64, Ordering}, - time::{SystemTime, UNIX_EPOCH}, - }, + std::time::{SystemTime, UNIX_EPOCH}, }; /// Cuts auctions from the open orders the indexer persisted. @@ -23,22 +20,11 @@ pub struct DbAuctionProvider { pool: PgPool, rpc: SolanaRPC, prices: NativePrices, - /// Last allocated auction id. Ids are unix seconds, bumped past the - /// previous allocation when cycles land within the same second. Unique - /// only per process: no table allocates auction ids. - /// TODO: allocate from the auctions table sequence once competition - /// persistence lands, like the EVM `auctions.id` bigserial. - last_id: AtomicI64, } impl DbAuctionProvider { pub fn new(pool: PgPool, rpc: SolanaRPC, prices: NativePrices) -> Self { - Self { - pool, - rpc, - prices, - last_id: AtomicI64::new(0), - } + Self { pool, rpc, prices } } /// Drop orders whose buy token account cannot receive the settlement @@ -84,18 +70,6 @@ impl DbAuctionProvider { }) .collect() } - - /// Allocates the next auction id: the current unix second, or one past - /// the previous id when several cycles land within the same second, so - /// ids strictly increase within the process. - fn next_id(&self, now: i64) -> i64 { - let prev = self - .last_id - .update(Ordering::Relaxed, Ordering::Relaxed, |prev| { - now.max(prev + 1) - }); - now.max(prev + 1) - } } fn now_unix() -> i64 { @@ -115,7 +89,7 @@ impl AuctionProvider for DbAuctionProvider { Ok(()) } - async fn cut_auction(&self, _tip: &u64) -> Option { + async fn cut_auction(&self, tip: &u64) -> Option { let now = now_unix(); // A pending sponsored order dies with its creation blockhash, so the // cut drops the dead ones. A failed height fetch keeps them all: they @@ -127,18 +101,17 @@ impl AuctionProvider for DbAuctionProvider { None } }; - let mut auction = db::cut(&self.pool, self.next_id(now), now, block_height) + let orders = db::cut(&self.pool, now, block_height) .await .map_err(|err| tracing::warn!(?err, "failed to cut the auction")) .ok()?; - auction.orders = self.receivable_orders(auction.orders).await; - if auction.orders.is_empty() { + let orders = self.receivable_orders(orders).await; + if orders.is_empty() { return None; } // A cut without prices would rank solutions on incomparable scores, // so a failed lookup skips the cycle instead. - let tokens = auction - .orders + let tokens = orders .iter() .flat_map(|order| [order.sell_token, order.buy_token]) .map(|token| Pubkey::new_from_array(token.0)) @@ -150,14 +123,53 @@ impl AuctionProvider for DbAuctionProvider { return None; } }; - auction.native_prices = prices - .into_iter() - .map(|(token, price)| (ChainPubkey(token.to_bytes()), price)) - .collect(); + let mut auction = crate::domain::auction::Auction { + id: 0, + orders, + native_prices: prices + .into_iter() + .map(|(token, price)| (ChainPubkey(token.to_bytes()), price)) + .collect(), + }; + // The id must be durable before anything references it: windows and + // the competition snapshot key on it, so a failed write skips the + // cycle. + let tip_slot = i64::try_from(*tip).unwrap_or(i64::MAX); + let snapshot = auction_snapshot(*tip, &auction); + auction.id = match db::replace_current_auction(&self.pool, tip_slot, &snapshot).await { + Ok(id) => id, + Err(err) => { + tracing::warn!(?err, "failed to store the auction, skipping the cut"); + return None; + } + }; Some(auction) } } +/// The stored auction body: the solver-facing content without the deadline, +/// which is only known at dispatch. +fn auction_snapshot(tip: u64, auction: &crate::domain::auction::Auction) -> serde_json::Value { + serde_json::json!({ + "tipSlot": tip, + "orders": auction + .orders + .iter() + .map(|order| order.uid.to_string()) + .collect::>(), + "nativePrices": auction + .native_prices + .iter() + .map(|(token, price)| { + ( + Pubkey::new_from_array(token.0).to_string(), + price.to_string(), + ) + }) + .collect::>(), + }) +} + #[derive(prometheus_metric_storage::MetricStorage)] #[metric(subsystem = "auction_provider")] struct Metrics { diff --git a/crates/autopilot-svm/src/test_db.rs b/crates/autopilot-svm/src/test_db.rs index 84879c3498..475e88b90e 100644 --- a/crates/autopilot-svm/src/test_db.rs +++ b/crates/autopilot-svm/src/test_db.rs @@ -11,7 +11,9 @@ pub(crate) async fn pool() -> PgPool { pub(crate) async fn wipe(pool: &PgPool) { sqlx::query( "TRUNCATE solana.trades, solana.settlements, solana.settlement_executions, \ - solana.order_pda, solana.orders, solana.indexer_state, solana.order_events", + solana.order_pda, solana.orders, solana.indexer_state, solana.order_events, \ + solana.auctions, solana.competition_auctions, solana.proposed_solutions, \ + solana.proposed_trade_executions", ) .execute(pool) .await diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index b32b5ec784..2053f1cb07 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -258,6 +258,37 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { .await .unwrap(); assert_eq!(open_windows, 1); + // The competition was persisted: the snapshot row, the proposed solution + // under its generated uid, its execution, and the window keyed by the + // same uid. + let snapshots: i64 = sqlx::query_scalar("SELECT count(*) FROM solana.competition_auctions") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(snapshots, 1); + let (solution_uid, solver_id, is_winner): (i64, i64, bool) = + sqlx::query_as("SELECT uid, id, is_winner FROM solana.proposed_solutions") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!((solution_uid, solver_id, is_winner), (0, 7, true)); + let (executed_sell, executed_buy): (String, String) = sqlx::query_as( + "SELECT executed_sell::text, executed_buy::text FROM solana.proposed_trade_executions \ + WHERE solution_uid = 0", + ) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!( + (executed_sell.as_str(), executed_buy.as_str()), + ("1000", "600") + ); + let window_uid: i64 = + sqlx::query_scalar("SELECT solution_uid FROM solana.settlement_executions") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(window_uid, 0); // The cycle reported the order's auction progress. The writes are detached // from the cycle, so they can land after `run_cycle` returns. let events = tokio::time::timeout(Duration::from_secs(5), async { diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql new file mode 100644 index 0000000000..3dcb307780 --- /dev/null +++ b/database/sql-solana/V8__competition.sql @@ -0,0 +1,42 @@ +-- The current auction, replaced on every cut. Its identity column is the +-- auction id sequence, so ids stay sequential across restarts like the EVM +-- auctions table. +CREATE TABLE solana.auctions ( + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + tip_slot bigint NOT NULL, + json jsonb NOT NULL +); + +-- Auctions that ran a competition, snapshot at ranking time. +CREATE TABLE solana.competition_auctions ( + id bigint PRIMARY KEY, + tip_slot bigint NOT NULL, + deadline_slot bigint NOT NULL, + order_uids bytea[] NOT NULL, + price_tokens bytea[] NOT NULL, + price_values numeric(20,0)[] NOT NULL +); + +-- Every solution proposed during a competition. The autopilot generates +-- `uid` per auction, disambiguating the solver-assigned `id` across drivers. +-- The EVM twin also stores uniform clearing prices, the Solana solve wire +-- carries none. +CREATE TABLE solana.proposed_solutions ( + auction_id bigint NOT NULL, + uid bigint NOT NULL, + id bigint NOT NULL, + solver bytea NOT NULL CHECK (length(solver) = 32), + is_winner boolean NOT NULL, + score numeric(20,0) NOT NULL, + PRIMARY KEY (auction_id, uid) +); + +-- The order executions of every proposed solution. +CREATE TABLE solana.proposed_trade_executions ( + auction_id bigint NOT NULL, + solution_uid bigint NOT NULL, + order_uid bytea NOT NULL CHECK (length(order_uid) = 32), + executed_sell numeric(20,0) NOT NULL, + executed_buy numeric(20,0) NOT NULL, + PRIMARY KEY (auction_id, solution_uid, order_uid) +); From f39d43c69e2991da6bc46b61f0e11fcbc7c6b4ec Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Wed, 23 Sep 2026 19:58:12 +0000 Subject: [PATCH 03/10] solana-autopilot: record which proposed solutions the fairness filter dropped The observer persists ranked and filtered-out solutions alike, so without a flag a debugger cannot tell a filtered solution from one that lost on score. Mirrors the EVM proposed_solutions.filtered_out column. --- crates/autopilot-svm/src/infra/db.rs | 4 +++- crates/autopilot-svm/src/infra/observer.rs | 1 + crates/autopilot-svm/src/tests.rs | 9 ++++++--- database/sql-solana/V8__competition.sql | 14 ++++++++------ 4 files changed, 18 insertions(+), 10 deletions(-) diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index c308a8622c..81198354bd 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -281,6 +281,7 @@ pub struct ProposedSolution { pub id: i64, pub solver: ByteArray<32>, pub is_winner: bool, + pub filtered_out: bool, pub score: BigDecimal, pub trades: Vec, } @@ -316,13 +317,14 @@ pub async fn persist_competition(pool: &sqlx::PgPool, competition: &Competition) for solution in &competition.solutions { sqlx::query( "INSERT INTO solana.proposed_solutions (auction_id, uid, id, solver, is_winner, \ - score) VALUES ($1, $2, $3, $4, $5, $6)", + filtered_out, score) VALUES ($1, $2, $3, $4, $5, $6, $7)", ) .bind(competition.auction_id) .bind(solution.uid) .bind(solution.id) .bind(solution.solver) .bind(solution.is_winner) + .bind(solution.filtered_out) .bind(&solution.score) .execute(&mut *tx) .await diff --git a/crates/autopilot-svm/src/infra/observer.rs b/crates/autopilot-svm/src/infra/observer.rs index 9cbd76aa36..9f7eb795db 100644 --- a/crates/autopilot-svm/src/infra/observer.rs +++ b/crates/autopilot-svm/src/infra/observer.rs @@ -101,6 +101,7 @@ impl SettlementObserver for CompetitionObserv id: i64::try_from(solution.id()).unwrap_or(i64::MAX), solver: ByteArray(solution.solver().0), is_winner: solution.is_winner(), + filtered_out: solution.is_filtered_out(), score: BigDecimal::from(solution.score()), trades: solution .orders() diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index 377cddf86c..95c2cbbad0 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -276,12 +276,15 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { .await .unwrap(); assert_eq!(snapshots, 1); - let (solution_uid, solver_id, is_winner): (i64, i64, bool) = - sqlx::query_as("SELECT uid, id, is_winner FROM solana.proposed_solutions") + let (solution_uid, solver_id, is_winner, filtered_out): (i64, i64, bool, bool) = + sqlx::query_as("SELECT uid, id, is_winner, filtered_out FROM solana.proposed_solutions") .fetch_one(&pool) .await .unwrap(); - assert_eq!((solution_uid, solver_id, is_winner), (0, 7, true)); + assert_eq!( + (solution_uid, solver_id, is_winner, filtered_out), + (0, 7, true, false) + ); let (executed_sell, executed_buy): (String, String) = sqlx::query_as( "SELECT executed_sell::text, executed_buy::text FROM solana.proposed_trade_executions \ WHERE solution_uid = 0", diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql index 3dcb307780..25813250a0 100644 --- a/database/sql-solana/V8__competition.sql +++ b/database/sql-solana/V8__competition.sql @@ -22,12 +22,14 @@ CREATE TABLE solana.competition_auctions ( -- The EVM twin also stores uniform clearing prices, the Solana solve wire -- carries none. CREATE TABLE solana.proposed_solutions ( - auction_id bigint NOT NULL, - uid bigint NOT NULL, - id bigint NOT NULL, - solver bytea NOT NULL CHECK (length(solver) = 32), - is_winner boolean NOT NULL, - score numeric(20,0) NOT NULL, + auction_id bigint NOT NULL, + uid bigint NOT NULL, + id bigint NOT NULL, + solver bytea NOT NULL CHECK (length(solver) = 32), + is_winner boolean NOT NULL, + -- Excluded from the ranking by the fairness filter, never a winner. + filtered_out boolean NOT NULL, + score numeric(20,0) NOT NULL, PRIMARY KEY (auction_id, uid) ); From 008d001d260204ceae9c04f79d94f733628c274b Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 09:39:14 +0000 Subject: [PATCH 04/10] solana-autopilot: derive the in-flight hold-out from the competition rows The cut now holds out orders of winning solutions until their deadline slot passes, the indexer records the solver's settlement, or the execution window closes early, the shape of the EVM fetch_in_flight_orders. The hold no longer depends on a window row being opened or on the deadline sweep running, and it survives a restart. Held orders get a filtered order event, a failed lookup skips the cut. The competition inserts go through QueryBuilder, one statement per table, like the EVM save_solutions. --- crates/autopilot-svm/src/infra/db.rs | 216 +++++++++++++----- crates/autopilot-svm/src/infra/observer.rs | 7 +- .../autopilot-svm/src/infra/order_events.rs | 10 + crates/autopilot-svm/src/infra/provider.rs | 50 ++-- crates/autopilot-svm/src/tests.rs | 29 ++- database/sql-solana/V8__competition.sql | 4 + 6 files changed, 237 insertions(+), 79 deletions(-) diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 81198354bd..c33731fcf2 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -6,7 +6,7 @@ use { bigdecimal::{BigDecimal, ToPrimitive}, chain_types::solana::{AppData, IntentHash, Pubkey}, database::byte_array::ByteArray, - sqlx::PgExecutor, + sqlx::{PgExecutor, Postgres, QueryBuilder}, }; /// The channel the schema's `solana.settlements` trigger notifies. @@ -34,13 +34,10 @@ pub struct OrderRow { /// Orders open for solving: unexpired, settleable by a driver, not cancelled /// and not fully filled. A pending sponsored order whose stored creation /// transaction died at `block_height` is excluded, and a `None` height skips -/// that check rather than excluding everything. An order inside a dispatched -/// settlement whose execution window is still open sits out until the window -/// closes, so an in-flight settlement is not raced by the next auction. -/// Settleable means the driver can produce the order PDA: it already exists -/// on chain (an order placed via `CreateOrder` directly), or the driver can -/// create it at settlement time from a signed intent or a presigned -/// transaction. +/// that check rather than excluding everything. Settleable means the driver can +/// produce the order PDA: it already exists on chain (an order placed via +/// `CreateOrder` directly), or the driver can create it at settlement time from +/// a signed intent or a presigned transaction. pub async fn open_orders( ex: impl PgExecutor<'_>, now_unix: i64, @@ -64,13 +61,6 @@ WHERE o.valid_to >= $1 OR o.presigned_transaction IS NULL OR p.order_uid IS NOT NULL OR o.last_valid_block_height >= $2) - AND NOT EXISTS ( - SELECT 1 - FROM solana.settlement_executions se - JOIN solana.proposed_trade_executions pte - ON pte.auction_id = se.auction_id AND pte.solution_uid = se.solution_uid - WHERE pte.order_uid = o.uid AND se.outcome IS NULL - ) AND COALESCE( CASE o.kind WHEN 'sell' THEN p.amount_withdrawn < o.sell_amount @@ -231,6 +221,38 @@ pub async fn open_window_auction_ids(ex: impl PgExecutor<'_>) -> Result .context("read open settlement execution windows") } +/// Orders inside a winning solution whose settlement may still land: the +/// deadline slot has not passed, the indexer recorded no settlement by that +/// solver for the auction, and the execution window did not close early. A +/// settlement is matched by solver because `settlements.solution_uid` is +/// unattributed, and one solver wins at most one solution per auction. +pub async fn in_flight_orders( + ex: impl PgExecutor<'_>, + tip_slot: i64, +) -> Result>> { + const QUERY: &str = r#" +SELECT DISTINCT pte.order_uid +FROM solana.competition_auctions ca +JOIN solana.proposed_solutions ps ON ps.auction_id = ca.id AND ps.is_winner +JOIN solana.proposed_trade_executions pte + ON pte.auction_id = ca.id AND pte.solution_uid = ps.uid +WHERE ca.deadline_slot >= $1 + AND NOT EXISTS ( + SELECT 1 FROM solana.settlements s + WHERE s.auction_id = ca.id AND s.solver = ps.solver + ) + AND NOT EXISTS ( + SELECT 1 FROM solana.settlement_executions se + WHERE se.auction_id = ca.id AND se.solution_uid = ps.uid AND se.outcome IS NOT NULL + ) + "#; + sqlx::query_scalar(QUERY) + .bind(tip_slot) + .fetch_all(ex) + .await + .context("read in-flight orders") +} + /// The solvable orders for a fresh auction cut. pub async fn cut( ex: impl PgExecutor<'_>, @@ -314,35 +336,53 @@ pub async fn persist_competition(pool: &sqlx::PgPool, competition: &Competition) .execute(&mut *tx) .await .context("insert competition auction")?; - for solution in &competition.solutions { - sqlx::query( + if !competition.solutions.is_empty() { + let mut insert = QueryBuilder::::new( "INSERT INTO solana.proposed_solutions (auction_id, uid, id, solver, is_winner, \ - filtered_out, score) VALUES ($1, $2, $3, $4, $5, $6, $7)", - ) - .bind(competition.auction_id) - .bind(solution.uid) - .bind(solution.id) - .bind(solution.solver) - .bind(solution.is_winner) - .bind(solution.filtered_out) - .bind(&solution.score) - .execute(&mut *tx) - .await - .context("insert proposed solution")?; - for trade in &solution.trades { - sqlx::query( - "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, \ - order_uid, executed_sell, executed_buy) VALUES ($1, $2, $3, $4, $5)", - ) - .bind(competition.auction_id) - .bind(solution.uid) - .bind(trade.order_uid) - .bind(&trade.executed_sell) - .bind(&trade.executed_buy) + filtered_out, score) ", + ); + insert.push_values(&competition.solutions, |mut row, solution| { + row.push_bind(competition.auction_id) + .push_bind(solution.uid) + .push_bind(solution.id) + .push_bind(solution.solver) + .push_bind(solution.is_winner) + .push_bind(solution.filtered_out) + .push_bind(&solution.score); + }); + insert + .build() .execute(&mut *tx) .await - .context("insert proposed trade execution")?; - } + .context("insert proposed solutions")?; + } + let trades: Vec<_> = competition + .solutions + .iter() + .flat_map(|solution| { + solution + .trades + .iter() + .map(move |trade| (solution.uid, trade)) + }) + .collect(); + if !trades.is_empty() { + let mut insert = QueryBuilder::::new( + "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, order_uid, \ + executed_sell, executed_buy) ", + ); + insert.push_values(trades, |mut row, (solution_uid, trade)| { + row.push_bind(competition.auction_id) + .push_bind(solution_uid) + .push_bind(trade.order_uid) + .push_bind(&trade.executed_sell) + .push_bind(&trade.executed_buy); + }); + insert + .build() + .execute(&mut *tx) + .await + .context("insert proposed trade executions")?; } tx.commit().await.context("commit competition persist") } @@ -400,7 +440,7 @@ fn to_amount(value: &BigDecimal) -> Result { #[cfg(test)] mod tests { use { - super::{last_indexed_slot, open_orders}, + super::{in_flight_orders, last_indexed_slot, open_orders}, bigdecimal::BigDecimal, database::byte_array::ByteArray, sqlx::PgTransaction, @@ -563,36 +603,106 @@ WHERE uid = $1 assert_eq!(uids(orders), vec![1, 5, 6, 10]); let orders = open_orders(&mut *tx, 1_000, Some(151)).await.unwrap(); assert_eq!(uids(orders), vec![1, 5, 6]); + } - // An order inside a dispatched settlement sits out while its window - // is open and returns once the window closes. + /// Held: an order of a winning solution through its deadline slot. + /// Released: past the deadline, on a settlement by the winner's solver, + /// or on an execution window closed early. Orders of non-winning + /// solutions are never held. + #[tokio::test] + #[ignore = "needs the solana.* schema applied to the local database"] + async fn solana_db_in_flight_orders_follow_the_winning_settlement() { + let pool = crate::test_db::pool().await; + let mut tx = pool.begin().await.unwrap(); + for table in [ + "settlements", + "settlement_executions", + "proposed_trade_executions", + "proposed_solutions", + "competition_auctions", + ] { + sqlx::query(&format!("DELETE FROM solana.{table}")) + .execute(&mut *tx) + .await + .unwrap(); + } + let winner = ByteArray([0xEE; 32]); sqlx::query( - "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, order_uid, \ - executed_sell, executed_buy) VALUES (77, 0, $1, 10, 20)", + "INSERT INTO solana.competition_auctions (id, tip_slot, deadline_slot, order_uids, \ + price_tokens, price_values) VALUES (77, 1, 100, '{}', '{}', '{}')", ) - .bind(ByteArray([1u8; 32])) .execute(&mut *tx) .await .unwrap(); + for (uid, solver, is_winner, order) in [ + (0i64, winner, true, 1u8), + (1, ByteArray([0xEF; 32]), false, 2), + ] { + sqlx::query( + "INSERT INTO solana.proposed_solutions (auction_id, uid, id, solver, is_winner, \ + filtered_out, score) VALUES (77, $1, 7, $2, $3, false, 1)", + ) + .bind(uid) + .bind(solver) + .bind(is_winner) + .execute(&mut *tx) + .await + .unwrap(); + sqlx::query( + "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, \ + order_uid, executed_sell, executed_buy) VALUES (77, $1, $2, 10, 20)", + ) + .bind(uid) + .bind(ByteArray([order; 32])) + .execute(&mut *tx) + .await + .unwrap(); + } + async fn held(tx: &mut PgTransaction<'_>, tip: i64) -> Vec { + in_flight_orders(&mut **tx, tip) + .await + .unwrap() + .iter() + .map(|uid| uid.0[0]) + .collect() + } + + assert_eq!(held(&mut tx, 100).await, vec![1]); + assert_eq!(held(&mut tx, 101).await, Vec::::new()); + sqlx::query( "INSERT INTO solana.settlement_executions (auction_id, solver, solution_uid, \ start_timestamp, start_slot, deadline_slot) VALUES (77, $1, 0, now(), 1, 100)", ) - .bind(ByteArray([0xEE; 32])) + .bind(winner) .execute(&mut *tx) .await .unwrap(); - let orders = open_orders(&mut *tx, 1_000, None).await.unwrap(); - assert_eq!(uids(orders), vec![5, 6, 10]); + assert_eq!(held(&mut tx, 50).await, vec![1]); sqlx::query( - "UPDATE solana.settlement_executions SET outcome = 'timeout', end_slot = 100, \ + "UPDATE solana.settlement_executions SET outcome = 'rejected', end_slot = 2, \ end_timestamp = now() WHERE auction_id = 77", ) .execute(&mut *tx) .await .unwrap(); - let orders = open_orders(&mut *tx, 1_000, None).await.unwrap(); - assert_eq!(uids(orders), vec![1, 5, 6, 10]); + assert_eq!(held(&mut tx, 50).await, Vec::::new()); + sqlx::query("UPDATE solana.settlement_executions SET outcome = NULL WHERE auction_id = 77") + .execute(&mut *tx) + .await + .unwrap(); + assert_eq!(held(&mut tx, 50).await, vec![1]); + + sqlx::query( + "INSERT INTO solana.settlements (slot, tx_signature, instruction_index, solver, \ + auction_id) VALUES (10, $1, 0, $2, 77)", + ) + .bind([9u8; 64]) + .bind(winner) + .execute(&mut *tx) + .await + .unwrap(); + assert_eq!(held(&mut tx, 50).await, Vec::::new()); } #[tokio::test] diff --git a/crates/autopilot-svm/src/infra/observer.rs b/crates/autopilot-svm/src/infra/observer.rs index 9f7eb795db..9f0663dbb3 100644 --- a/crates/autopilot-svm/src/infra/observer.rs +++ b/crates/autopilot-svm/src/infra/observer.rs @@ -38,12 +38,7 @@ impl CompetitionObserver { /// Store the events without blocking the cycle: a lost event degrades the /// status endpoint, never the competition. fn store_events(&self, uids: Vec, label: OrderEventLabel) { - let pool = self.pool.clone(); - tokio::spawn(async move { - if let Err(err) = order_events::store(&pool, uids, label).await { - tracing::error!(?err, ?label, "failed to store order events"); - } - }); + order_events::store_detached(self.pool.clone(), uids, label); } } diff --git a/crates/autopilot-svm/src/infra/order_events.rs b/crates/autopilot-svm/src/infra/order_events.rs index 73789f0830..de2c420131 100644 --- a/crates/autopilot-svm/src/infra/order_events.rs +++ b/crates/autopilot-svm/src/infra/order_events.rs @@ -14,6 +14,16 @@ const DEDUP_LOCK: &str = "solana_order_events_dedup"; const DEDUP_LOCK_QUERY: &str = "SELECT pg_advisory_xact_lock(hashtextextended($1::text || $2::text, 0))"; +/// Append the events off the caller's path: a failed write is logged, not +/// returned. +pub fn store_detached(pool: PgPool, uids: Vec, label: OrderEventLabel) { + tokio::spawn(async move { + if let Err(err) = store(&pool, uids, label).await { + tracing::error!(?err, ?label, "failed to store order events"); + } + }); +} + /// Append one event per order, skipping a label the order's latest event /// already carries: a looping order marks each state once, not once per cycle. pub async fn store( diff --git a/crates/autopilot-svm/src/infra/provider.rs b/crates/autopilot-svm/src/infra/provider.rs index e6634b01f4..b914971b96 100644 --- a/crates/autopilot-svm/src/infra/provider.rs +++ b/crates/autopilot-svm/src/infra/provider.rs @@ -3,16 +3,20 @@ use { crate::{ domain::{auction::Order, cycle::SolanaCycle}, - infra::{db, inflight::InFlightOrders, prices::NativePrices}, + infra::{db, inflight::InFlightOrders, order_events, prices::NativePrices}, run_loop::AuctionProvider, }, async_trait::async_trait, - chain_types::solana::Pubkey as ChainPubkey, + chain_types::solana::{IntentHash, Pubkey as ChainPubkey}, cow_solana_rpc::SolanaRPC, + database::solana::OrderEventLabel, solana_sdk::{account::Account, program_pack::Pack, pubkey::Pubkey}, spl_token_interface::state::{Account as TokenAccount, AccountState}, sqlx::PgPool, - std::time::{SystemTime, UNIX_EPOCH}, + std::{ + collections::HashSet, + time::{SystemTime, UNIX_EPOCH}, + }, }; /// Cuts auctions from the open orders the indexer persisted. @@ -22,7 +26,7 @@ pub struct DbAuctionProvider { /// Slots the indexer may lag behind the tip before cuts are skipped. max_indexer_lag: u64, /// Orders with a settlement in flight, excluded from cuts until their - /// submission deadline passes. + /// settlement transaction cannot land any more. inflight: InFlightOrders, prices: NativePrices, } @@ -134,22 +138,41 @@ impl AuctionProvider for DbAuctionProvider { None } }; - let mut orders = db::cut(&self.pool, now, block_height) + let orders = db::cut(&self.pool, now, block_height) .await .map_err(|err| tracing::warn!(?err, "failed to cut the auction")) .ok()?; // An order with a settlement in flight stays out until the // settlement cannot land any more: a second winner could - // double-settle it. - let held = self.inflight.held_at(*tip); - let before = orders.len(); - orders.retain(|order| !held.contains(&order.uid)); - let held_out = before - orders.len(); - if held_out > 0 { + // double-settle it. The competition tables hold it through the + // deadline and survive a restart, the in-memory hold covers the + // blockhash lifetime past it. A failed read skips the cut rather + // than cutting without the hold. + let tip_slot = i64::try_from(*tip).unwrap_or(i64::MAX); + let mut held: HashSet = match db::in_flight_orders(&self.pool, tip_slot).await { + Ok(uids) => uids.into_iter().map(|uid| IntentHash(uid.0)).collect(), + Err(err) => { + tracing::warn!(?err, "in-flight order lookup failed, skipping the cut"); + return None; + } + }; + held.extend(self.inflight.held_at(*tip)); + let (orders, held_out): (Vec<_>, Vec<_>) = orders + .into_iter() + .partition(|order| !held.contains(&order.uid)); + if !held_out.is_empty() { metrics() .held_out_orders - .inc_by(u64::try_from(held_out).unwrap_or(u64::MAX)); - tracing::debug!(held_out, "orders held out with settlements in flight"); + .inc_by(u64::try_from(held_out.len()).unwrap_or(u64::MAX)); + tracing::debug!( + held_out = held_out.len(), + "orders held out with settlements in flight" + ); + order_events::store_detached( + self.pool.clone(), + held_out.into_iter().map(|order| order.uid).collect(), + OrderEventLabel::Filtered, + ); } let orders = self.receivable_orders(orders).await; if orders.is_empty() { @@ -180,7 +203,6 @@ impl AuctionProvider for DbAuctionProvider { // The id must be durable before anything references it: windows and // the competition snapshot key on it, so a failed write skips the // cycle. - let tip_slot = i64::try_from(*tip).unwrap_or(i64::MAX); let snapshot = auction_snapshot(*tip, &auction); auction.id = match db::replace_current_auction(&self.pool, tip_slot, &snapshot).await { Ok(id) => id, diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index 95c2cbbad0..6598aede0e 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -306,8 +306,24 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { // open, and past the window until its settlement transaction cannot land // any more: the deadline plus the blockhash lifetime. let expired = tip + 25 + solana_sdk::clock::MAX_PROCESSING_AGE as u64 + 1; - // The lag gate stays out of the way: this provider tests the hold, and - // the watermark is still at the dispatch tip. + // The competition rows alone hold the order through the deadline slot: a + // provider without the in-memory hold stands in for a restart. The lag + // gate stays out of the way, the watermark is still at the dispatch tip. + let restarted_provider = DbAuctionProvider::new( + pool.clone(), + mock_rpc(), + u64::MAX, + InFlightOrders::default(), + NativePrices::seeded(test_prices()), + ); + assert!( + restarted_provider.cut_auction(&(tip + 25)).await.is_none(), + "in-flight order excluded by the persisted competition" + ); + assert!( + restarted_provider.cut_auction(&(tip + 26)).await.is_some(), + "persisted hold ends past the deadline" + ); let held_provider = DbAuctionProvider::new( pool.clone(), mock_rpc(), @@ -330,8 +346,9 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { held_provider.cut_auction(&expired).await.is_some(), "order returns once the transaction cannot land" ); - // The cycle reported the order's auction progress. The writes are detached - // from the cycle, so they can land after `run_cycle` returns. + // The cycle reported the order's auction progress and the later cuts its + // hold-out. The writes are detached, so they can land after the cuts + // return. let events = tokio::time::timeout(Duration::from_secs(5), async { loop { let mut events: Vec = @@ -339,7 +356,7 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { .fetch_all(&pool) .await .unwrap(); - if events.len() == 2 { + if events.len() == 3 { events.sort(); return events; } @@ -348,5 +365,5 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { }) .await .expect("order events written before the timeout"); - assert_eq!(events, ["executing", "ready"]); + assert_eq!(events, ["executing", "filtered", "ready"]); } diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql index 25813250a0..b4c35f78ec 100644 --- a/database/sql-solana/V8__competition.sql +++ b/database/sql-solana/V8__competition.sql @@ -17,6 +17,10 @@ CREATE TABLE solana.competition_auctions ( price_values numeric(20,0)[] NOT NULL ); +-- The in-flight hold-out scans the auctions whose deadline has not passed. +CREATE INDEX solana_competition_auctions_deadline_slot + ON solana.competition_auctions (deadline_slot); + -- Every solution proposed during a competition. The autopilot generates -- `uid` per auction, disambiguating the solver-assigned `id` across drivers. -- The EVM twin also stores uniform clearing prices, the Solana solve wire From f7f1598f86e690050d08516a9920c61fccd99c09 Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 10:04:50 +0000 Subject: [PATCH 05/10] solana-autopilot: drop the in-memory in-flight hold-out The competition rows carry the hold on their own, the EVM shape. The blockhash grace the map provided moves into the query as MAX_PROCESSING_AGE past the deadline slot, and a timed-out window keeps the hold for that span. The map's two release triggers were table state already: an observed settlement is a settlements row, a provable non-send a rejected window. The hold now survives a restart without any seeding. --- crates/autopilot-svm/src/infra/db.rs | 41 ++-- crates/autopilot-svm/src/infra/executor.rs | 19 +- crates/autopilot-svm/src/infra/inflight.rs | 175 ------------------ crates/autopilot-svm/src/infra/mod.rs | 1 - crates/autopilot-svm/src/infra/observation.rs | 26 +-- crates/autopilot-svm/src/infra/provider.rs | 24 +-- crates/autopilot-svm/src/run.rs | 14 +- crates/autopilot-svm/src/tests.rs | 42 +---- 8 files changed, 56 insertions(+), 286 deletions(-) delete mode 100644 crates/autopilot-svm/src/infra/inflight.rs diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index c33731fcf2..25cc3a8dee 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -6,6 +6,7 @@ use { bigdecimal::{BigDecimal, ToPrimitive}, chain_types::solana::{AppData, IntentHash, Pubkey}, database::byte_array::ByteArray, + solana_sdk::clock::MAX_PROCESSING_AGE, sqlx::{PgExecutor, Postgres, QueryBuilder}, }; @@ -221,11 +222,14 @@ pub async fn open_window_auction_ids(ex: impl PgExecutor<'_>) -> Result .context("read open settlement execution windows") } -/// Orders inside a winning solution whose settlement may still land: the -/// deadline slot has not passed, the indexer recorded no settlement by that -/// solver for the auction, and the execution window did not close early. A -/// settlement is matched by solver because `settlements.solution_uid` is -/// unattributed, and one solver wins at most one solution per auction. +/// Orders inside a winning solution whose settlement transaction may still +/// land: a blockhash lifetime past the deadline slot has not run out, the +/// indexer recorded no settlement by that solver for the auction, and the +/// driver did not reject the settlement before sending it. A timed-out +/// window keeps the hold, its transaction may land until the blockhash +/// expires. A settlement is matched by solver because +/// `settlements.solution_uid` is unattributed, and one solver wins at most +/// one solution per auction. pub async fn in_flight_orders( ex: impl PgExecutor<'_>, tip_slot: i64, @@ -243,11 +247,14 @@ WHERE ca.deadline_slot >= $1 ) AND NOT EXISTS ( SELECT 1 FROM solana.settlement_executions se - WHERE se.auction_id = ca.id AND se.solution_uid = ps.uid AND se.outcome IS NOT NULL + WHERE se.auction_id = ca.id AND se.solution_uid = ps.uid AND se.outcome = 'rejected' ) "#; + // The oldest deadline whose transaction can still land at the tip. + let lifetime = i64::try_from(MAX_PROCESSING_AGE).expect("blockhash lifetime fits i64"); + let landable_deadline = tip_slot.saturating_sub(lifetime); sqlx::query_scalar(QUERY) - .bind(tip_slot) + .bind(landable_deadline) .fetch_all(ex) .await .context("read in-flight orders") @@ -605,10 +612,11 @@ WHERE uid = $1 assert_eq!(uids(orders), vec![1, 5, 6]); } - /// Held: an order of a winning solution through its deadline slot. - /// Released: past the deadline, on a settlement by the winner's solver, - /// or on an execution window closed early. Orders of non-winning - /// solutions are never held. + /// Held: an order of a winning solution through its deadline slot plus + /// the blockhash lifetime, a timed-out window included. Released: after + /// that, on a settlement by the winner's solver, or on a window the + /// driver rejected before sending. Orders of non-winning solutions are + /// never held. #[tokio::test] #[ignore = "needs the solana.* schema applied to the local database"] async fn solana_db_in_flight_orders_follow_the_winning_settlement() { @@ -668,7 +676,8 @@ WHERE uid = $1 } assert_eq!(held(&mut tx, 100).await, vec![1]); - assert_eq!(held(&mut tx, 101).await, Vec::::new()); + assert_eq!(held(&mut tx, 250).await, vec![1]); + assert_eq!(held(&mut tx, 251).await, Vec::::new()); sqlx::query( "INSERT INTO solana.settlement_executions (auction_id, solver, solution_uid, \ @@ -692,6 +701,14 @@ WHERE uid = $1 .await .unwrap(); assert_eq!(held(&mut tx, 50).await, vec![1]); + sqlx::query( + "UPDATE solana.settlement_executions SET outcome = 'timeout', end_slot = 100, \ + end_timestamp = now() WHERE auction_id = 77", + ) + .execute(&mut *tx) + .await + .unwrap(); + assert_eq!(held(&mut tx, 150).await, vec![1]); sqlx::query( "INSERT INTO solana.settlements (slot, tx_signature, instruction_index, solver, \ diff --git a/crates/autopilot-svm/src/infra/executor.rs b/crates/autopilot-svm/src/infra/executor.rs index 8897276740..9e08170d38 100644 --- a/crates/autopilot-svm/src/infra/executor.rs +++ b/crates/autopilot-svm/src/infra/executor.rs @@ -5,7 +5,6 @@ use { domain::cycle::{Ranking, SolanaCycle}, infra::{ driver::{Driver, dto}, - inflight::InFlightOrders, observation::SettlementWindows, sponsor::Sponsor, }, @@ -26,9 +25,6 @@ pub struct DriverExecutor { /// without creations, and one containing a pending sponsored order fails /// at the driver. sponsor: Option, - /// Orders dispatched here are held out of auction cuts until their - /// submission deadline passes. - inflight: InFlightOrders, } impl DriverExecutor { @@ -36,13 +32,11 @@ impl DriverExecutor { drivers: Vec>, windows: SettlementWindows, sponsor: Option, - inflight: InFlightOrders, ) -> Self { Self { drivers, windows, sponsor, - inflight, } } } @@ -90,11 +84,6 @@ impl SettlementExecutor for DriverExecutor { submission_deadline_slot: deadline, creations, }; - // Held before the dispatch: the next cut must not re-auction - // these orders while the settlement can still land. - let uids: Vec<_> = winner.orders().iter().map(|order| order.uid).collect(); - self.inflight - .hold(auction_id, winner.solver(), uids.iter().copied(), deadline); // A window that cannot be opened must not block the settlement, // the dispatch is the priority. The window carries the generated // solution uid, the driver request keeps the driver-local id its @@ -110,7 +99,6 @@ impl SettlementExecutor for DriverExecutor { { tracing::error!(auction_id, ?err, "failed to open the settlement window"); } - let inflight = self.inflight.clone(); let windows = self.windows.clone(); let solver = winner.solver(); tokio::spawn(async move { @@ -122,9 +110,9 @@ impl SettlementExecutor for DriverExecutor { tx_signature = %response.tx_signature, "settlement submitted" ), - // No transaction went out, so the orders can re-enter - // the next cut instead of waiting out the hold, and the - // window has nothing left to observe. + // No transaction went out, so the window closes as + // rejected: the orders re-enter the next cut and nothing + // is left to observe. Err(err) if err.settlement_provably_unsent() => { tracing::warn!( driver = %driver.name, @@ -132,7 +120,6 @@ impl SettlementExecutor for DriverExecutor { ?err, "settlement rejected before submission" ); - inflight.release(uids); if let Err(err) = windows.close_rejected(auction_id, solver, uid).await { tracing::error!( auction_id, diff --git a/crates/autopilot-svm/src/infra/inflight.rs b/crates/autopilot-svm/src/infra/inflight.rs deleted file mode 100644 index accdc59134..0000000000 --- a/crates/autopilot-svm/src/infra/inflight.rs +++ /dev/null @@ -1,175 +0,0 @@ -//! In-memory hold-out of orders with a settlement in flight. -//! -//! An order dispatched for settlement must not re-enter the next auction -//! while the first settlement can still land: a second winner would -//! double-settle it. The driver stops waiting at the submission deadline, -//! but the transaction it sent stays landable until its blockhash expires, -//! up to `MAX_PROCESSING_AGE` slots later, so held orders expire only at -//! the deadline plus that lifetime. Two events end a hold early: the -//! settlement is observed on chain (the auction's orders are done), or the -//! driver rejects provably before any send. Everything else, a submit error -//! or an exceeded deadline, may have left the transaction on the wire and -//! keeps the hold. -//! -//! The map lives in memory: a restart forgets it and reopens the window -//! until the entries would have expired. - -use { - chain_types::solana::{IntentHash, Pubkey}, - solana_sdk::clock::MAX_PROCESSING_AGE, - std::{ - collections::{HashMap, HashSet}, - sync::{Arc, Mutex}, - }, -}; - -/// A dispatched settlement, keyed the way the indexer identifies a landed -/// one. -#[derive(Clone, Copy, PartialEq, Eq)] -struct Settlement { - auction_id: i64, - solver: Pubkey, -} - -/// Why an order is held, and until when. -#[derive(Clone, Copy)] -struct Hold { - settlement: Settlement, - /// The last slot the settlement's transaction could still land. - expires_at: u64, -} - -/// Order uids held out of auction cuts until their settlement transaction -/// cannot land any more. -#[derive(Clone, Default)] -pub struct InFlightOrders(Arc>>); - -impl InFlightOrders { - /// Hold the settlement's orders until the deadline slot plus the - /// blockhash lifetime, the last slot its transaction could still land. - /// An order already held keeps the hold that expires last. - pub fn hold( - &self, - auction_id: i64, - solver: Pubkey, - uids: impl IntoIterator, - deadline_slot: u64, - ) { - let hold = Hold { - settlement: Settlement { auction_id, solver }, - expires_at: deadline_slot.saturating_add(MAX_PROCESSING_AGE as u64), - }; - let mut held = self.0.lock().expect("mutex poisoned"); - for uid in uids { - held.entry(uid) - .and_modify(|current| { - if hold.expires_at > current.expires_at { - *current = hold; - } - }) - .or_insert(hold); - } - } - - /// Release the orders: their settlement provably never went out, so no - /// second settlement can collide. - pub fn release(&self, uids: impl IntoIterator) { - let mut held = self.0.lock().expect("mutex poisoned"); - for uid in uids { - held.remove(&uid); - } - } - - /// Release the orders of a settlement observed on chain. The auction is - /// settled for this solver, nothing else can execute these orders. - pub fn release_landed(&self, auction_id: i64, solver: Pubkey) { - let settlement = Settlement { auction_id, solver }; - self.0 - .lock() - .expect("mutex poisoned") - .retain(|_, hold| hold.settlement != settlement); - } - - /// The orders still held at the tip. Expired entries are pruned on the - /// way. - pub fn held_at(&self, tip: u64) -> HashSet { - let mut held = self.0.lock().expect("mutex poisoned"); - held.retain(|_, hold| hold.expires_at >= tip); - held.keys().copied().collect() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - const SOLVER: Pubkey = Pubkey([9; 32]); - - #[test] - fn holds_through_the_blockhash_lifetime_past_the_deadline() { - let inflight = InFlightOrders::default(); - let uid = IntentHash([7; 32]); - inflight.hold(1, SOLVER, [uid], 100); - let expiry = 100 + MAX_PROCESSING_AGE as u64; - - assert!(inflight.held_at(100).contains(&uid)); - assert!(inflight.held_at(expiry).contains(&uid)); - assert!(inflight.held_at(expiry + 1).is_empty()); - } - - #[test] - fn released_orders_re_enter_immediately() { - let inflight = InFlightOrders::default(); - let uid = IntentHash([7; 32]); - let other = IntentHash([8; 32]); - inflight.hold(1, SOLVER, [uid, other], 100); - inflight.release([uid]); - - let held = inflight.held_at(50); - assert!(!held.contains(&uid)); - assert!(held.contains(&other)); - } - - #[test] - fn a_landed_settlement_releases_its_orders() { - let inflight = InFlightOrders::default(); - let uid = IntentHash([7; 32]); - let unrelated = IntentHash([8; 32]); - inflight.hold(1, SOLVER, [uid], 100); - inflight.hold(2, SOLVER, [unrelated], 100); - - inflight.release_landed(1, SOLVER); - let held = inflight.held_at(50); - assert!(!held.contains(&uid)); - assert!(held.contains(&unrelated)); - - // A landing for a solver without a dispatch is a no-op. - inflight.release_landed(3, Pubkey([1; 32])); - assert!(inflight.held_at(50).contains(&unrelated)); - } - - #[test] - fn a_second_dispatch_keeps_the_later_expiry() { - let inflight = InFlightOrders::default(); - let uid = IntentHash([7; 32]); - inflight.hold(1, SOLVER, [uid], 100); - inflight.hold(2, SOLVER, [uid], 90); - - assert!( - inflight - .held_at(100 + MAX_PROCESSING_AGE as u64) - .contains(&uid) - ); - } - - #[test] - fn a_landing_keeps_an_order_the_later_dispatch_still_holds() { - let inflight = InFlightOrders::default(); - let uid = IntentHash([7; 32]); - inflight.hold(1, SOLVER, [uid], 100); - inflight.hold(2, SOLVER, [uid], 200); - - inflight.release_landed(1, SOLVER); - assert!(inflight.held_at(50).contains(&uid)); - } -} diff --git a/crates/autopilot-svm/src/infra/mod.rs b/crates/autopilot-svm/src/infra/mod.rs index b491908e90..08dceb60fd 100644 --- a/crates/autopilot-svm/src/infra/mod.rs +++ b/crates/autopilot-svm/src/infra/mod.rs @@ -5,7 +5,6 @@ pub mod config; pub mod db; pub mod driver; pub mod executor; -pub mod inflight; pub mod listen; pub mod observation; pub mod observer; diff --git a/crates/autopilot-svm/src/infra/observation.rs b/crates/autopilot-svm/src/infra/observation.rs index 137289d31c..81a811ae0d 100644 --- a/crates/autopilot-svm/src/infra/observation.rs +++ b/crates/autopilot-svm/src/infra/observation.rs @@ -11,7 +11,7 @@ //! window. use { - crate::infra::{db, inflight::InFlightOrders, listen::NotifyHandler}, + crate::infra::{db, listen::NotifyHandler}, anyhow::Result, async_trait::async_trait, chain_types::solana::{Pubkey, Signature}, @@ -31,13 +31,11 @@ use { #[derive(Clone)] pub struct SettlementWindows { pool: PgPool, - /// A settlement observed on chain releases its orders from the hold-out. - inflight: InFlightOrders, } impl SettlementWindows { - pub fn new(pool: PgPool, inflight: InFlightOrders) -> Self { - Self { pool, inflight } + pub fn new(pool: PgPool) -> Self { + Self { pool } } /// Open a window for a dispatched settlement. `solution_uid` is the @@ -97,7 +95,6 @@ impl SettlementWindows { tx_signature = %Signature(landed.submitted_signature.0), "settlement observed on chain" ); - self.inflight.release_landed(auction_id, solver); } Ok(()) } @@ -132,8 +129,8 @@ impl NotifyHandler for SettlementWindows { mod tests { use { super::SettlementWindows, - crate::infra::{db, inflight::InFlightOrders, listen::ListenSession}, - chain_types::solana::{IntentHash, Pubkey}, + crate::infra::{db, listen::ListenSession}, + chain_types::solana::Pubkey, sqlx::PgPool, std::time::Duration, }; @@ -171,10 +168,7 @@ VALUES (10, $1, 0, $2, $3, NULL) crate::test_db::wipe(&pool).await; let solver = Pubkey([7; 32]); - let uid = IntentHash([1; 32]); - let inflight = InFlightOrders::default(); - inflight.hold(4242, solver, [uid], 100); - let windows = SettlementWindows::new(pool.clone(), inflight.clone()); + let windows = SettlementWindows::new(pool.clone()); windows .open_dispatched(4242, solver, 1, 90, 100) .await @@ -196,8 +190,6 @@ VALUES (10, $1, 0, $2, $3, NULL) } task.abort(); assert_eq!(outcome(&pool, 4242).await.as_deref(), Some("landed")); - // The observed landing released the held order. - assert!(inflight.held_at(90).is_empty()); let signature: Vec = sqlx::query_scalar( "SELECT submitted_signature FROM solana.settlement_executions WHERE auction_id = 4242", ) @@ -215,7 +207,7 @@ VALUES (10, $1, 0, $2, $3, NULL) let pool = crate::test_db::pool().await; crate::test_db::wipe(&pool).await; - let windows = SettlementWindows::new(pool.clone(), InFlightOrders::default()); + let windows = SettlementWindows::new(pool.clone()); windows .open_dispatched(1, Pubkey([7; 32]), 1, 90, 100) .await @@ -246,7 +238,7 @@ VALUES (10, $1, 0, $2, $3, NULL) let pool = crate::test_db::pool().await; crate::test_db::wipe(&pool).await; - let windows = SettlementWindows::new(pool.clone(), InFlightOrders::default()); + let windows = SettlementWindows::new(pool.clone()); windows .open_dispatched(1, Pubkey([7; 32]), 1, 90, 100) .await @@ -267,7 +259,7 @@ VALUES (10, $1, 0, $2, $3, NULL) let pool = crate::test_db::pool().await; crate::test_db::wipe(&pool).await; - let windows = SettlementWindows::new(pool.clone(), InFlightOrders::default()); + let windows = SettlementWindows::new(pool.clone()); windows .open_dispatched(1, Pubkey([7; 32]), 1, 90, 100) .await diff --git a/crates/autopilot-svm/src/infra/provider.rs b/crates/autopilot-svm/src/infra/provider.rs index b914971b96..5bc8d16a73 100644 --- a/crates/autopilot-svm/src/infra/provider.rs +++ b/crates/autopilot-svm/src/infra/provider.rs @@ -3,7 +3,7 @@ use { crate::{ domain::{auction::Order, cycle::SolanaCycle}, - infra::{db, inflight::InFlightOrders, order_events, prices::NativePrices}, + infra::{db, order_events, prices::NativePrices}, run_loop::AuctionProvider, }, async_trait::async_trait, @@ -25,25 +25,15 @@ pub struct DbAuctionProvider { rpc: SolanaRPC, /// Slots the indexer may lag behind the tip before cuts are skipped. max_indexer_lag: u64, - /// Orders with a settlement in flight, excluded from cuts until their - /// settlement transaction cannot land any more. - inflight: InFlightOrders, prices: NativePrices, } impl DbAuctionProvider { - pub fn new( - pool: PgPool, - rpc: SolanaRPC, - max_indexer_lag: u64, - inflight: InFlightOrders, - prices: NativePrices, - ) -> Self { + pub fn new(pool: PgPool, rpc: SolanaRPC, max_indexer_lag: u64, prices: NativePrices) -> Self { Self { pool, rpc, max_indexer_lag, - inflight, prices, } } @@ -144,19 +134,16 @@ impl AuctionProvider for DbAuctionProvider { .ok()?; // An order with a settlement in flight stays out until the // settlement cannot land any more: a second winner could - // double-settle it. The competition tables hold it through the - // deadline and survive a restart, the in-memory hold covers the - // blockhash lifetime past it. A failed read skips the cut rather - // than cutting without the hold. + // double-settle it. A failed read skips the cut rather than cutting + // without the hold. let tip_slot = i64::try_from(*tip).unwrap_or(i64::MAX); - let mut held: HashSet = match db::in_flight_orders(&self.pool, tip_slot).await { + let held: HashSet = match db::in_flight_orders(&self.pool, tip_slot).await { Ok(uids) => uids.into_iter().map(|uid| IntentHash(uid.0)).collect(), Err(err) => { tracing::warn!(?err, "in-flight order lookup failed, skipping the cut"); return None; } }; - held.extend(self.inflight.held_at(*tip)); let (orders, held_out): (Vec<_>, Vec<_>) = orders .into_iter() .partition(|order| !held.contains(&order.uid)); @@ -311,7 +298,6 @@ mod tests { sqlx::PgPool::connect_lazy("postgresql://").unwrap(), SolanaRPC::new_mock_with_mocks(mocks), 150, - InFlightOrders::default(), NativePrices::seeded([]), ) } diff --git a/crates/autopilot-svm/src/run.rs b/crates/autopilot-svm/src/run.rs index 7a7ce386e5..7492a53c2f 100644 --- a/crates/autopilot-svm/src/run.rs +++ b/crates/autopilot-svm/src/run.rs @@ -9,7 +9,6 @@ use { db, driver::Driver, executor::DriverExecutor, - inflight::InFlightOrders, listen::ListenSession, observation::SettlementWindows, observer::CompetitionObserver, @@ -89,10 +88,7 @@ async fn run(config: Config) { .await .expect("database connection"); - // One shared hold-out: the executor holds into it, the auction cut reads - // it, and the settlement observer releases from it. - let inflight = InFlightOrders::default(); - let windows = SettlementWindows::new(pool.clone(), inflight.clone()); + let windows = SettlementWindows::new(pool.clone()); let listen = ListenSession::spawn( pool.clone(), db::SETTLEMENT_FINALIZED_CHANNEL, @@ -135,7 +131,6 @@ async fn run(config: Config) { CommitmentConfig::confirmed(), ), config.max_indexer_lag_slots, - inflight.clone(), NativePrices::new( &config.native_prices, SolanaRPC::new_with_timeout_and_commitment( @@ -154,12 +149,7 @@ async fn run(config: Config) { config.competition.max_winners.get(), Pubkey(config.contracts.wrapped_native_mint.to_bytes()), )), - Box::new(DriverExecutor::new( - drivers, - windows.clone(), - sponsor, - inflight, - )), + Box::new(DriverExecutor::new(drivers, windows.clone(), sponsor)), Box::new(CompetitionObserver::new(pool, windows)), config.competition.submission_deadline_slots.get(), ); diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index 6598aede0e..980e0d549a 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -8,7 +8,6 @@ use { competition::DriverCompetition, driver::{Driver, dto}, executor::DriverExecutor, - inflight::InFlightOrders, observation::SettlementWindows, observer::CompetitionObserver, prices::NativePrices, @@ -215,7 +214,6 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { pool.clone(), mock_rpc(), 150, - InFlightOrders::default(), NativePrices::seeded(test_prices()), ); let auction = provider.cut_auction(&tip).await.expect("auction cut"); @@ -227,15 +225,13 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { assert_eq!(ranking.winner_count(), 1, "solution won"); } - let inflight = InFlightOrders::default(); - let windows = SettlementWindows::new(pool.clone(), inflight.clone()); + let windows = SettlementWindows::new(pool.clone()); let mut auction_loop = AuctionLoop::new( Box::new(FixedTrigger(tip)), Box::new(DbAuctionProvider::new( pool.clone(), mock_rpc(), 150, - inflight.clone(), NativePrices::seeded(test_prices()), )), Box::new(DriverCompetition::new( @@ -243,12 +239,7 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { Duration::from_secs(6), )), Box::new(SolanaArbitrator::new(1, wrapped_native)), - Box::new(DriverExecutor::new( - vec![driver], - windows.clone(), - None, - inflight.clone(), - )), + Box::new(DriverExecutor::new(vec![driver], windows.clone(), None)), Box::new(CompetitionObserver::new(pool.clone(), windows.clone())), 25, ); @@ -302,37 +293,20 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { .await .unwrap(); assert_eq!(window_uid, 0); - // The dispatched order is held out of the next cut while its window is - // open, and past the window until its settlement transaction cannot land - // any more: the deadline plus the blockhash lifetime. + // The dispatched order is held out of the next cut until its settlement + // transaction cannot land any more: the deadline plus the blockhash + // lifetime. The hold comes from the persisted competition alone, so a + // fresh provider also stands in for a restart. The lag gate stays out of + // the way, the watermark is still at the dispatch tip. let expired = tip + 25 + solana_sdk::clock::MAX_PROCESSING_AGE as u64 + 1; - // The competition rows alone hold the order through the deadline slot: a - // provider without the in-memory hold stands in for a restart. The lag - // gate stays out of the way, the watermark is still at the dispatch tip. - let restarted_provider = DbAuctionProvider::new( - pool.clone(), - mock_rpc(), - u64::MAX, - InFlightOrders::default(), - NativePrices::seeded(test_prices()), - ); - assert!( - restarted_provider.cut_auction(&(tip + 25)).await.is_none(), - "in-flight order excluded by the persisted competition" - ); - assert!( - restarted_provider.cut_auction(&(tip + 26)).await.is_some(), - "persisted hold ends past the deadline" - ); let held_provider = DbAuctionProvider::new( pool.clone(), mock_rpc(), u64::MAX, - inflight.clone(), NativePrices::seeded(test_prices()), ); assert!( - held_provider.cut_auction(&(expired - 1)).await.is_none(), + held_provider.cut_auction(&(tip + 25)).await.is_none(), "in-flight order excluded from the cut" ); // The deadline sweep closes the window as timed out. The hold outlasts From 33df00d34d3a1f56925a6da75df47ba101735ab6 Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:18:08 +0000 Subject: [PATCH 06/10] solana-autopilot: release a held order on its own landed trade, not on any settlement by its solver Winner selection only keeps token pairs disjoint, so one solver can win several solutions of one auction. Matching the landed settlement by solver released all of them at once. Matching the settlement's trades releases exactly the orders that landed. --- crates/autopilot-svm/src/infra/db.rs | 65 +++++++++++++++++++--------- 1 file changed, 44 insertions(+), 21 deletions(-) diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 25cc3a8dee..8b02927afd 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -223,13 +223,12 @@ pub async fn open_window_auction_ids(ex: impl PgExecutor<'_>) -> Result } /// Orders inside a winning solution whose settlement transaction may still -/// land: a blockhash lifetime past the deadline slot has not run out, the -/// indexer recorded no settlement by that solver for the auction, and the -/// driver did not reject the settlement before sending it. A timed-out -/// window keeps the hold, its transaction may land until the blockhash -/// expires. A settlement is matched by solver because -/// `settlements.solution_uid` is unattributed, and one solver wins at most -/// one solution per auction. +/// land: a blockhash lifetime past the deadline slot has not run out, no +/// settlement of the auction traded the order yet, and the driver did not +/// reject the solution before sending it. A timed-out window keeps the hold, +/// its transaction may land until the blockhash expires. Landing is checked +/// per order through the settlement's trades: `settlements.solution_uid` is +/// unattributed, and one solver may win several solutions of one auction. pub async fn in_flight_orders( ex: impl PgExecutor<'_>, tip_slot: i64, @@ -242,8 +241,11 @@ JOIN solana.proposed_trade_executions pte ON pte.auction_id = ca.id AND pte.solution_uid = ps.uid WHERE ca.deadline_slot >= $1 AND NOT EXISTS ( - SELECT 1 FROM solana.settlements s - WHERE s.auction_id = ca.id AND s.solver = ps.solver + SELECT 1 + FROM solana.settlements s + JOIN solana.trades t + ON t.tx_signature = s.tx_signature AND t.instruction_index = s.instruction_index + WHERE s.auction_id = ca.id AND t.order_uid = pte.order_uid ) AND NOT EXISTS ( SELECT 1 FROM solana.settlement_executions se @@ -614,15 +616,17 @@ WHERE uid = $1 /// Held: an order of a winning solution through its deadline slot plus /// the blockhash lifetime, a timed-out window included. Released: after - /// that, on a settlement by the winner's solver, or on a window the - /// driver rejected before sending. Orders of non-winning solutions are - /// never held. + /// that, once a settlement of the auction trades the order, or once the + /// driver rejects its solution before sending. A solver's second winning + /// solution keeps its hold when the first one lands. Orders of + /// non-winning solutions are never held. #[tokio::test] #[ignore = "needs the solana.* schema applied to the local database"] async fn solana_db_in_flight_orders_follow_the_winning_settlement() { let pool = crate::test_db::pool().await; let mut tx = pool.begin().await.unwrap(); for table in [ + "trades", "settlements", "settlement_executions", "proposed_trade_executions", @@ -645,6 +649,7 @@ WHERE uid = $1 for (uid, solver, is_winner, order) in [ (0i64, winner, true, 1u8), (1, ByteArray([0xEF; 32]), false, 2), + (2, winner, true, 3), ] { sqlx::query( "INSERT INTO solana.proposed_solutions (auction_id, uid, id, solver, is_winner, \ @@ -667,16 +672,18 @@ WHERE uid = $1 .unwrap(); } async fn held(tx: &mut PgTransaction<'_>, tip: i64) -> Vec { - in_flight_orders(&mut **tx, tip) + let mut held: Vec = in_flight_orders(&mut **tx, tip) .await .unwrap() .iter() .map(|uid| uid.0[0]) - .collect() + .collect(); + held.sort_unstable(); + held } - assert_eq!(held(&mut tx, 100).await, vec![1]); - assert_eq!(held(&mut tx, 250).await, vec![1]); + assert_eq!(held(&mut tx, 100).await, vec![1, 3]); + assert_eq!(held(&mut tx, 250).await, vec![1, 3]); assert_eq!(held(&mut tx, 251).await, Vec::::new()); sqlx::query( @@ -687,7 +694,7 @@ WHERE uid = $1 .execute(&mut *tx) .await .unwrap(); - assert_eq!(held(&mut tx, 50).await, vec![1]); + assert_eq!(held(&mut tx, 50).await, vec![1, 3]); sqlx::query( "UPDATE solana.settlement_executions SET outcome = 'rejected', end_slot = 2, \ end_timestamp = now() WHERE auction_id = 77", @@ -695,12 +702,12 @@ WHERE uid = $1 .execute(&mut *tx) .await .unwrap(); - assert_eq!(held(&mut tx, 50).await, Vec::::new()); + assert_eq!(held(&mut tx, 50).await, vec![3]); sqlx::query("UPDATE solana.settlement_executions SET outcome = NULL WHERE auction_id = 77") .execute(&mut *tx) .await .unwrap(); - assert_eq!(held(&mut tx, 50).await, vec![1]); + assert_eq!(held(&mut tx, 50).await, vec![1, 3]); sqlx::query( "UPDATE solana.settlement_executions SET outcome = 'timeout', end_slot = 100, \ end_timestamp = now() WHERE auction_id = 77", @@ -708,8 +715,10 @@ WHERE uid = $1 .execute(&mut *tx) .await .unwrap(); - assert_eq!(held(&mut tx, 150).await, vec![1]); + assert_eq!(held(&mut tx, 150).await, vec![1, 3]); + // The solver's settlement trades order 1 only: order 3, its other + // winning solution, stays held until its own trade lands. sqlx::query( "INSERT INTO solana.settlements (slot, tx_signature, instruction_index, solver, \ auction_id) VALUES (10, $1, 0, $2, 77)", @@ -719,7 +728,21 @@ WHERE uid = $1 .execute(&mut *tx) .await .unwrap(); - assert_eq!(held(&mut tx, 50).await, Vec::::new()); + for order in [1u8, 3] { + sqlx::query( + "INSERT INTO solana.trades (tx_signature, instruction_index, order_uid, \ + sell_amount, buy_amount, fee_amount) VALUES ($1, 0, $2, 10, 20, 0)", + ) + .bind([9u8; 64]) + .bind(ByteArray([order; 32])) + .execute(&mut *tx) + .await + .unwrap(); + assert_eq!( + held(&mut tx, 50).await, + if order == 1 { vec![3] } else { vec![] } + ); + } } #[tokio::test] From deacbcb77d5e921fefa2653f9bb60e40d97321ee Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 13:42:49 +0000 Subject: [PATCH 07/10] solana-autopilot: close a settlement window on the settlement that executed its solution A settlement row carries no solution uid, and one solver may hold several windows of one auction. Matching by solver closed all of them on the first landing, with the first signature. The window's solution has persisted trade executions, so the settlement whose trades cover one of them is the one that executed it. --- crates/autopilot-svm/src/infra/db.rs | 23 ++-- crates/autopilot-svm/src/infra/observation.rs | 102 +++++++++++++++++- 2 files changed, 112 insertions(+), 13 deletions(-) diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 8b02927afd..9d3a0bbeef 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -168,14 +168,12 @@ WHERE auction_id = $1 AND solver = $2 AND solution_uid = $3 AND outcome IS NULL Ok(()) } -/// Close the auction's windows against the settlements the indexer recorded, -/// matching each window to its solver's settlement. A window already closed -/// as timed out upgrades to landed: the settlement executed, just late, and -/// lateness stays visible as `end_slot` past `deadline_slot`. -/// -/// A settlement carries no solution uid, so a solver holding several windows -/// of one auction closes all of them on its first settlement. Correct while -/// one solver wins at most one solution per auction. +/// Close the auction's windows against the settlements the indexer recorded. +/// A settlement carries no solution uid, so a window is matched through its +/// solution's trade executions: the solver's settlement that traded one of +/// them is the one that executed it. A window already closed as timed out +/// upgrades to landed: the settlement executed, just late, and lateness stays +/// visible as `end_slot` past `deadline_slot`. pub async fn close_landed_windows( ex: impl PgExecutor<'_>, auction_id: i64, @@ -189,6 +187,15 @@ WHERE e.auction_id = $1 AND s.auction_id = e.auction_id AND s.solver = e.solver AND (e.outcome IS NULL OR e.outcome = 'timeout') + AND EXISTS ( + SELECT 1 + FROM solana.proposed_trade_executions pte + JOIN solana.trades t + ON t.tx_signature = s.tx_signature + AND t.instruction_index = s.instruction_index + AND t.order_uid = pte.order_uid + WHERE pte.auction_id = e.auction_id AND pte.solution_uid = e.solution_uid + ) RETURNING e.solver, e.end_slot, e.submitted_signature "#; sqlx::query_as(QUERY) diff --git a/crates/autopilot-svm/src/infra/observation.rs b/crates/autopilot-svm/src/infra/observation.rs index 81a811ae0d..b9ba306adb 100644 --- a/crates/autopilot-svm/src/infra/observation.rs +++ b/crates/autopilot-svm/src/infra/observation.rs @@ -135,14 +135,44 @@ mod tests { std::time::Duration, }; - async fn insert_settlement(pool: &PgPool, auction_id: i64) { + /// A persisted execution of order `[order; 32]` inside solution `uid` of + /// the auction. + async fn insert_execution(pool: &PgPool, auction_id: i64, uid: i64, order: u8) { + sqlx::query( + "INSERT INTO solana.proposed_trade_executions (auction_id, solution_uid, order_uid, \ + executed_sell, executed_buy) VALUES ($1, $2, $3, 10, 20)", + ) + .bind(auction_id) + .bind(uid) + .bind([order; 32]) + .execute(pool) + .await + .unwrap(); + } + + /// A landed settlement of the auction by solver 7 under signature + /// `[signature; 64]`, trading the given orders. The trades go in first: + /// the indexer commits both together, so the NOTIFY the settlement fires + /// never sees a settlement without its trades. + async fn insert_settlement(pool: &PgPool, auction_id: i64, signature: u8, orders: &[u8]) { + for order in orders { + sqlx::query( + "INSERT INTO solana.trades (tx_signature, instruction_index, order_uid, \ + sell_amount, buy_amount, fee_amount) VALUES ($1, 0, $2, 10, 20, 0)", + ) + .bind([signature; 64]) + .bind([*order; 32]) + .execute(pool) + .await + .unwrap(); + } sqlx::query( r#" INSERT INTO solana.settlements (slot, tx_signature, instruction_index, solver, auction_id, solution_uid) VALUES (10, $1, 0, $2, $3, NULL) "#, ) - .bind([9u8; 64]) + .bind([signature; 64]) .bind([7u8; 32]) .bind(auction_id) .execute(pool) @@ -158,6 +188,21 @@ VALUES (10, $1, 0, $2, $3, NULL) .unwrap() } + /// Outcome and signature of every window of the auction, by solution uid. + async fn windows_of( + pool: &PgPool, + auction_id: i64, + ) -> Vec<(i64, Option, Option>)> { + sqlx::query_as( + "SELECT solution_uid, outcome, submitted_signature FROM solana.settlement_executions \ + WHERE auction_id = $1 ORDER BY solution_uid", + ) + .bind(auction_id) + .fetch_all(pool) + .await + .unwrap() + } + /// The full path: a dispatched settlement opens a window, the trigger's /// NOTIFY (here: a bare INSERT, standing in for the indexer) closes it /// as landed with the settlement's signature. @@ -173,6 +218,7 @@ VALUES (10, $1, 0, $2, $3, NULL) .open_dispatched(4242, solver, 1, 90, 100) .await .unwrap(); + insert_execution(&pool, 4242, 1, 1).await; let task = ListenSession::spawn( pool.clone(), @@ -180,7 +226,7 @@ VALUES (10, $1, 0, $2, $3, NULL) windows.clone(), ); - insert_settlement(&pool, 4242).await; + insert_settlement(&pool, 4242, 9, &[1]).await; for _ in 0..200 { if outcome(&pool, 4242).await.is_some() { @@ -223,7 +269,8 @@ VALUES (10, $1, 0, $2, $3, NULL) // A settlement observed after the timeout upgrades the verdict: it // executed, just late. - insert_settlement(&pool, 1).await; + insert_execution(&pool, 1, 1, 1).await; + insert_settlement(&pool, 1, 9, &[1]).await; crate::infra::db::close_landed_windows(&pool, 1) .await .unwrap(); @@ -265,7 +312,8 @@ VALUES (10, $1, 0, $2, $3, NULL) .await .unwrap(); - insert_settlement(&pool, 1).await; + insert_execution(&pool, 1, 1, 1).await; + insert_settlement(&pool, 1, 9, &[1]).await; crate::infra::db::close_landed_windows(&pool, 1) .await .unwrap(); @@ -274,4 +322,48 @@ VALUES (10, $1, 0, $2, $3, NULL) windows.close_rejected(1, Pubkey([7; 32]), 1).await.unwrap(); assert_eq!(outcome(&pool, 1).await.as_deref(), Some("landed")); } + + /// A solver holding two windows of one auction: the settlement that + /// traded the first solution's order closes only that window, the second + /// window closes on its own settlement with its own signature. + #[tokio::test] + #[ignore = "needs the solana.* schema applied locally, run with --test-threads 1"] + async fn solana_db_a_landing_closes_only_the_window_it_executed() { + let pool = crate::test_db::pool().await; + crate::test_db::wipe(&pool).await; + + let solver = Pubkey([7; 32]); + let windows = SettlementWindows::new(pool.clone()); + for (uid, order) in [(1, 1u8), (2, 2)] { + windows + .open_dispatched(1, solver, uid, 90, 100) + .await + .unwrap(); + insert_execution(&pool, 1, uid, order).await; + } + + insert_settlement(&pool, 1, 9, &[1]).await; + crate::infra::db::close_landed_windows(&pool, 1) + .await + .unwrap(); + assert_eq!( + windows_of(&pool, 1).await, + vec![ + (1, Some("landed".to_string()), Some(vec![9u8; 64])), + (2, None, None), + ] + ); + + insert_settlement(&pool, 1, 8, &[2]).await; + crate::infra::db::close_landed_windows(&pool, 1) + .await + .unwrap(); + assert_eq!( + windows_of(&pool, 1).await, + vec![ + (1, Some("landed".to_string()), Some(vec![9u8; 64])), + (2, Some("landed".to_string()), Some(vec![8u8; 64])), + ] + ); + } } From 2ea6f3292b4f821de968cbd3a2ae9bb0f68933e3 Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 15:00:40 +0000 Subject: [PATCH 08/10] solana-autopilot: persist reference scores and number solutions by rank position Reference scores come from the generic arbitrator, one per winning solver, and go into solana.reference_scores with the solutions, the rewards baseline accounting asked for. The persisted solution uid is now the position in ranked followed by filtered_out, as the EVM autopilot enumerates its ranking. The map keyed by (solver, solution id) collapsed two solutions sharing a key into one uid and failed the whole persist on the primary key. --- crates/autopilot-svm/src/domain/arbitrator.rs | 17 +------ crates/autopilot-svm/src/domain/cycle.rs | 26 ++++++++-- crates/autopilot-svm/src/infra/db.rs | 23 +++++++++ crates/autopilot-svm/src/infra/executor.rs | 10 ++-- crates/autopilot-svm/src/infra/observer.rs | 48 ++++++++++--------- crates/autopilot-svm/src/test_db.rs | 2 +- crates/autopilot-svm/src/tests.rs | 10 ++++ database/sql-solana/V8__competition.sql | 9 ++++ 8 files changed, 96 insertions(+), 49 deletions(-) diff --git a/crates/autopilot-svm/src/domain/arbitrator.rs b/crates/autopilot-svm/src/domain/arbitrator.rs index a2d3d5f76e..2fa5fb14a6 100644 --- a/crates/autopilot-svm/src/domain/arbitrator.rs +++ b/crates/autopilot-svm/src/domain/arbitrator.rs @@ -70,24 +70,11 @@ impl WinnerSelection for SolanaArbitrator { for winner in inner.winners() { tracing::info!(solver = %winner.solver(), solution = winner.id(), "winner"); } - // Ranked before filtered, so winner uids come first. - let uids = inner - .ranked - .iter() - .map(|solution| (solution.solver(), solution.id())) - .chain( - inner - .filtered_out - .iter() - .map(|solution| (solution.solver(), solution.id())), - ) - .enumerate() - .map(|(uid, key)| (key, i64::try_from(uid).unwrap_or(i64::MAX))) - .collect(); + let reference_scores = self.inner.compute_reference_scores(&inner); Ranking { inner, drivers, - uids, + reference_scores, } } } diff --git a/crates/autopilot-svm/src/domain/cycle.rs b/crates/autopilot-svm/src/domain/cycle.rs index eba71ea948..5b27c8a1dd 100644 --- a/crates/autopilot-svm/src/domain/cycle.rs +++ b/crates/autopilot-svm/src/domain/cycle.rs @@ -7,7 +7,7 @@ use { }, chain_types::solana::{IntentHash, Pubkey, Solana}, std::collections::{HashMap, HashSet}, - winner_selection::{Unscored, solution}, + winner_selection::{Unscored, solution, state::Ranked}, }; /// Marker type binding the generic loop to the Solana vocabulary. @@ -45,10 +45,26 @@ pub struct Ranking { pub inner: winner_selection::Ranking, /// Driver index per solution, keyed by `(solver, solution id)`. pub drivers: HashMap, - /// Autopilot-generated solution uid per solution, unique within the - /// auction. It disambiguates solver-assigned ids across drivers, and - /// everything persisted references it. - pub uids: HashMap, + /// Per winning solver, the winners' total score with that solver's + /// solutions removed: the rewards baseline. + pub reference_scores: HashMap, +} + +impl Ranking { + /// Every solution with its persisted uid: the position in `ranked` + /// followed by `filtered_out`, so winners come first. Everything + /// persisted references this uid, it disambiguates solver-assigned ids + /// across drivers. + pub fn enumerated( + &self, + ) -> impl Iterator, Solana>)> { + self.inner + .ranked + .iter() + .chain(self.inner.filtered_out.iter()) + .enumerate() + .map(|(uid, solution)| (i64::try_from(uid).unwrap_or(i64::MAX), solution)) + } } impl RankingInfo for Ranking { diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index 9d3a0bbeef..6bb21ed1eb 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -324,6 +324,13 @@ pub struct ProposedSolution { pub trades: Vec, } +/// A winning solver's reference score: the winners' total score with that +/// solver's solutions removed. +pub struct ReferenceScore { + pub solver: ByteArray<32>, + pub score: BigDecimal, +} + /// A competition outcome as persisted after ranking. pub struct Competition { pub auction_id: i64, @@ -333,6 +340,7 @@ pub struct Competition { pub price_tokens: Vec>, pub price_values: Vec, pub solutions: Vec, + pub reference_scores: Vec, } /// Persist a competition: the auction snapshot and every proposed solution @@ -400,6 +408,21 @@ pub async fn persist_competition(pool: &sqlx::PgPool, competition: &Competition) .await .context("insert proposed trade executions")?; } + if !competition.reference_scores.is_empty() { + let mut insert = QueryBuilder::::new( + "INSERT INTO solana.reference_scores (auction_id, solver, reference_score) ", + ); + insert.push_values(&competition.reference_scores, |mut row, reference| { + row.push_bind(competition.auction_id) + .push_bind(reference.solver) + .push_bind(&reference.score); + }); + insert + .build() + .execute(&mut *tx) + .await + .context("insert reference scores")?; + } tx.commit().await.context("commit competition persist") } diff --git a/crates/autopilot-svm/src/infra/executor.rs b/crates/autopilot-svm/src/infra/executor.rs index 9e08170d38..9eeeed5db9 100644 --- a/crates/autopilot-svm/src/infra/executor.rs +++ b/crates/autopilot-svm/src/infra/executor.rs @@ -12,6 +12,7 @@ use { }, async_trait::async_trait, std::sync::Arc, + winner_selection::state::RankedItem, }; /// Sends `/settle` to each winner's driver. Submission runs detached, the @@ -44,7 +45,10 @@ impl DriverExecutor { #[async_trait] impl SettlementExecutor for DriverExecutor { async fn execute(&self, auction_id: i64, ranking: &Ranking, tip: &u64, deadline: u64) { - for winner in ranking.inner.winners() { + for (uid, winner) in ranking + .enumerated() + .filter(|(_, winner)| winner.is_winner()) + { let key = (winner.solver(), winner.id()); let Some(driver) = ranking .drivers @@ -88,10 +92,6 @@ impl SettlementExecutor for DriverExecutor { // the dispatch is the priority. The window carries the generated // solution uid, the driver request keeps the driver-local id its // solution cache is keyed by. - let uid = ranking.uids.get(&key).copied().unwrap_or_else(|| { - tracing::error!(solution_id = winner.id(), "winner without a uid"); - i64::MAX - }); if let Err(err) = self .windows .open_dispatched(auction_id, winner.solver(), uid, *tip, deadline) diff --git a/crates/autopilot-svm/src/infra/observer.rs b/crates/autopilot-svm/src/infra/observer.rs index 9f0663dbb3..ce73910a86 100644 --- a/crates/autopilot-svm/src/infra/observer.rs +++ b/crates/autopilot-svm/src/infra/observer.rs @@ -85,29 +85,23 @@ impl SettlementObserver for CompetitionObserv tracing::error!(?err, "failed to flag expired settlement windows"); } let solutions = ranking - .inner - .ranked - .iter() - .chain(ranking.inner.filtered_out.iter()) - .map(|solution| { - let key = (solution.solver(), solution.id()); - db::ProposedSolution { - uid: ranking.uids.get(&key).copied().unwrap_or(i64::MAX), - id: i64::try_from(solution.id()).unwrap_or(i64::MAX), - solver: ByteArray(solution.solver().0), - is_winner: solution.is_winner(), - filtered_out: solution.is_filtered_out(), - score: BigDecimal::from(solution.score()), - trades: solution - .orders() - .iter() - .map(|order| db::ProposedTrade { - order_uid: ByteArray(order.uid.0), - executed_sell: BigDecimal::from(order.executed_sell), - executed_buy: BigDecimal::from(order.executed_buy), - }) - .collect(), - } + .enumerated() + .map(|(uid, solution)| db::ProposedSolution { + uid, + id: i64::try_from(solution.id()).unwrap_or(i64::MAX), + solver: ByteArray(solution.solver().0), + is_winner: solution.is_winner(), + filtered_out: solution.is_filtered_out(), + score: BigDecimal::from(solution.score()), + trades: solution + .orders() + .iter() + .map(|order| db::ProposedTrade { + order_uid: ByteArray(order.uid.0), + executed_sell: BigDecimal::from(order.executed_sell), + executed_buy: BigDecimal::from(order.executed_buy), + }) + .collect(), }) .collect(); let (price_tokens, price_values) = auction @@ -127,6 +121,14 @@ impl SettlementObserver for CompetitionObserv price_tokens, price_values, solutions, + reference_scores: ranking + .reference_scores + .iter() + .map(|(solver, score)| db::ReferenceScore { + solver: ByteArray(solver.0), + score: BigDecimal::from(*score), + }) + .collect(), }; db::persist_competition(&self.pool, &competition).await?; tracing::info!( diff --git a/crates/autopilot-svm/src/test_db.rs b/crates/autopilot-svm/src/test_db.rs index 475e88b90e..5f7b662f80 100644 --- a/crates/autopilot-svm/src/test_db.rs +++ b/crates/autopilot-svm/src/test_db.rs @@ -13,7 +13,7 @@ pub(crate) async fn wipe(pool: &PgPool) { "TRUNCATE solana.trades, solana.settlements, solana.settlement_executions, \ solana.order_pda, solana.orders, solana.indexer_state, solana.order_events, \ solana.auctions, solana.competition_auctions, solana.proposed_solutions, \ - solana.proposed_trade_executions", + solana.proposed_trade_executions, solana.reference_scores", ) .execute(pool) .await diff --git a/crates/autopilot-svm/src/tests.rs b/crates/autopilot-svm/src/tests.rs index 980e0d549a..9cea0defdd 100644 --- a/crates/autopilot-svm/src/tests.rs +++ b/crates/autopilot-svm/src/tests.rs @@ -276,6 +276,16 @@ async fn solana_db_mock_cycle_dispatches_the_settlement() { (solution_uid, solver_id, is_winner, filtered_out), (0, 7, true, false) ); + // The sole winner's reference score: the winners' total without it, zero. + let (reference_solver, reference_score): (Vec, String) = + sqlx::query_as("SELECT solver, reference_score::text FROM solana.reference_scores") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!( + (reference_solver.as_slice(), reference_score.as_str()), + (&[0xCC; 32][..], "0") + ); let (executed_sell, executed_buy): (String, String) = sqlx::query_as( "SELECT executed_sell::text, executed_buy::text FROM solana.proposed_trade_executions \ WHERE solution_uid = 0", diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql index b4c35f78ec..2fd2116ffa 100644 --- a/database/sql-solana/V8__competition.sql +++ b/database/sql-solana/V8__competition.sql @@ -46,3 +46,12 @@ CREATE TABLE solana.proposed_trade_executions ( executed_buy numeric(20,0) NOT NULL, PRIMARY KEY (auction_id, solution_uid, order_uid) ); + +-- Per winning solver, the winners' total score with that solver's solutions +-- removed: the rewards baseline. +CREATE TABLE solana.reference_scores ( + auction_id bigint NOT NULL, + solver bytea NOT NULL CHECK (length(solver) = 32), + reference_score numeric(20,0) NOT NULL, + PRIMARY KEY (auction_id, solver) +); From 0b5331c2f6fcf9b2a8999c1fd3185eebfbc43452 Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 15:26:35 +0000 Subject: [PATCH 09/10] solana-autopilot: correct the clearing price note on proposed_solutions The EVM columns exist but are written empty, they are not a feature the Solana table lacks. --- database/sql-solana/V8__competition.sql | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql index 2fd2116ffa..4ff7443aca 100644 --- a/database/sql-solana/V8__competition.sql +++ b/database/sql-solana/V8__competition.sql @@ -23,8 +23,8 @@ CREATE INDEX solana_competition_auctions_deadline_slot -- Every solution proposed during a competition. The autopilot generates -- `uid` per auction, disambiguating the solver-assigned `id` across drivers. --- The EVM twin also stores uniform clearing prices, the Solana solve wire --- carries none. +-- No clearing price columns: the EVM twin writes its own empty, and the +-- Solana solve wire carries none. CREATE TABLE solana.proposed_solutions ( auction_id bigint NOT NULL, uid bigint NOT NULL, From abf73c7c9689e78ad02f147e887d5f8bbc986cf0 Mon Sep 17 00:00:00 2001 From: squadgazzz <22964585+squadgazzz@users.noreply.github.com> Date: Thu, 24 Sep 2026 20:07:47 +0000 Subject: [PATCH 10/10] solana-autopilot: pin the auctions table to one row --- database/sql-solana/V8__competition.sql | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/database/sql-solana/V8__competition.sql b/database/sql-solana/V8__competition.sql index 4ff7443aca..c913214ecb 100644 --- a/database/sql-solana/V8__competition.sql +++ b/database/sql-solana/V8__competition.sql @@ -2,9 +2,12 @@ -- auction id sequence, so ids stay sequential across restarts like the EVM -- auctions table. CREATE TABLE solana.auctions ( - id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, - tip_slot bigint NOT NULL, - json jsonb NOT NULL + id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY, + tip_slot bigint NOT NULL, + json jsonb NOT NULL, + -- Pins the table to one row: a second writer's insert fails instead of + -- leaving two current auctions. + singleton boolean NOT NULL DEFAULT true UNIQUE CHECK (singleton) ); -- Auctions that ran a competition, snapshot at ranking time.