diff --git a/crates/autopilot-svm/src/infra/db.rs b/crates/autopilot-svm/src/infra/db.rs index fa76871166..3b6af458cd 100644 --- a/crates/autopilot-svm/src/infra/db.rs +++ b/crates/autopilot-svm/src/infra/db.rs @@ -43,6 +43,7 @@ SELECT o.uid, o.owner, o.sell_token, o.buy_token, o.sell_token_account, FROM solana.orders o LEFT JOIN solana.order_pda p ON p.order_uid = o.uid WHERE o.valid_to >= $1 + AND NOT COALESCE(p.is_reorged, false) AND (o.valid_from IS NULL OR o.valid_from <= $1) AND (o.intent_signature IS NOT NULL OR o.presigned_transaction IS NOT NULL diff --git a/crates/cow-solana-rpc/src/lib.rs b/crates/cow-solana-rpc/src/lib.rs index c44cd488fb..33223df5e1 100644 --- a/crates/cow-solana-rpc/src/lib.rs +++ b/crates/cow-solana-rpc/src/lib.rs @@ -113,6 +113,25 @@ impl SolanaRPC { }) } + /// Whether each signature exists on the node's chain. The lookup falls + /// back to the ledger history when a signature is past the recent-status + /// cache, so `false` means the transaction never landed or was rolled + /// back, not merely that it aged out. A status of any commitment level + /// counts as existing: the question here is presence, not finality. + pub async fn known_signatures(&self, signatures: &[Signature]) -> Result, Error> { + // getSignatureStatuses rejects calls with more than 256 signatures. + const CHUNK: usize = 256; + let mut known = Vec::with_capacity(signatures.len()); + for chunk in signatures.chunks(CHUNK) { + let response = self + .inner + .get_signature_statuses_with_history(chunk) + .await?; + known.extend(response.value.into_iter().map(|status| status.is_some())); + } + Ok(known) + } + /// Simulate a versioned transaction without sending it. Returns the /// simulation result including logs and any error. pub async fn simulate_transaction( diff --git a/crates/solana-indexer/src/indexer/decoder.rs b/crates/solana-indexer/src/indexer/decoder.rs index de25cff242..855e1f6543 100644 --- a/crates/solana-indexer/src/indexer/decoder.rs +++ b/crates/solana-indexer/src/indexer/decoder.rs @@ -38,6 +38,7 @@ use { }, recover_discriminator, }, + cow_solana_rpc::SolanaRPC, solana_sdk::pubkey::Pubkey, std::collections::BTreeMap, tokio::sync::mpsc::Receiver, @@ -48,6 +49,9 @@ pub(crate) struct Decoder { /// Persistence layer. pub persistence: Postgres, + /// RPC client the finalization audit checks signatures against. + pub rpc: SolanaRPC, + /// Incoming `StreamUpdate` from the ingester. pub rx: Receiver, @@ -63,12 +67,14 @@ impl Decoder { /// Construct a new decoder. The caller owns the channel capacity decision. pub fn new( persistence: Postgres, + rpc: SolanaRPC, rx: Receiver, settlement_program: Pubkey, solflow_program: Option, ) -> Self { Self { persistence, + rpc, rx, settlement_program, solflow_program, @@ -92,6 +98,7 @@ impl Decoder { // In-memory mirror of the persisted watermark, spares redundant // writes and flags late transactions. let mut watermark: Option = None; + let mut finalized_through: Option = None; while let Some(update) = self.rx.recv().await { let (slot, signature, inner) = match update { StreamUpdate::Tx { @@ -105,7 +112,7 @@ impl Decoder { continue; } StreamUpdate::Finalized { slot } => { - self.persistence.write_finalized_slot(slot).await?; + self.finalize(slot, &mut finalized_through).await?; continue; } }; @@ -145,6 +152,55 @@ impl Decoder { Ok(()) } + /// Advance the finalized watermark to `slot`, first auditing the rows + /// that finalization would freeze: every transaction indexed in the + /// newly final range must still exist on chain, and one that vanished + /// was rolled back, so its rows are reverted in the same SQL transaction + /// that advances the watermark. An audit RPC failure leaves the + /// watermark untouched, the next finalized status retries the range. + async fn finalize( + &self, + slot: Slot, + finalized_through: &mut Option, + ) -> Result<(), PersistenceError> { + let after = match *finalized_through { + Some(after) => after, + // First finalized status since boot: resume the audit from the + // persisted watermark, or from this very slot on a fresh + // database, where nothing older was indexed. + None => self.persistence.finalized_slot().await?.unwrap_or(slot), + }; + if slot <= after { + *finalized_through = Some(after); + return Ok(()); + } + let signatures = self.persistence.unfinalized_signatures(after, slot).await?; + let vanished = if signatures.is_empty() { + Vec::new() + } else { + match self.rpc.known_signatures(&signatures).await { + Ok(known) => signatures + .into_iter() + .zip(known) + .filter_map(|(signature, known)| (!known).then_some(signature)) + .collect(), + Err(err) => { + tracing::warn!(?err, "signature audit failed, finalization delayed"); + *finalized_through = Some(after); + return Ok(()); + } + } + }; + if vanished.is_empty() { + self.persistence.write_finalized_slot(slot).await?; + } else { + tracing::warn!(?vanished, "rolled-back transactions reverted"); + self.persistence.finalize_through(slot, &vanished).await?; + } + *finalized_through = Some(slot); + Ok(()) + } + /// Flush every buffer at or below the confirmed slot, then advance the /// watermark to it: the slot's transactions all arrived before its /// status, and quiet slots below it have nothing to wait for. diff --git a/crates/solana-indexer/src/indexer/decoder/tests.rs b/crates/solana-indexer/src/indexer/decoder/tests.rs index bf9700f92f..0be56c464b 100644 --- a/crates/solana-indexer/src/indexer/decoder/tests.rs +++ b/crates/solana-indexer/src/indexer/decoder/tests.rs @@ -40,6 +40,7 @@ use { data::intent::{Flags, OrderIntent, OrderKind as IntentOrderKind}, pda::order::find_order_pda, }, + cow_solana_rpc::{Mocks, RpcRequest, SolanaRPC}, futures::StreamExt, solana_sdk::pubkey::Pubkey, std::sync::{Arc, atomic::AtomicU64}, @@ -396,7 +397,30 @@ fn create_order_tx() -> (SubscribeUpdateTransactionInfo, CreatedOrder) { fn pure_decoder(settlement: Pubkey, solflow: Pubkey) -> Decoder { let pool = sqlx::PgPool::connect_lazy("postgresql://").unwrap(); let (_sender, rx) = tokio::sync::mpsc::channel(1); - Decoder::new(Postgres::new(pool), rx, settlement, Some(solflow)) + Decoder::new( + Postgres::new(pool), + SolanaRPC::new_mock_with_mocks(Default::default()), + rx, + settlement, + Some(solflow), + ) +} + +/// The finalization audit asks for two signatures (the dead letter and the +/// healthy create), both still known to the chain. +fn mock_rpc_with_signature_statuses() -> SolanaRPC { + let status = serde_json::json!({ + "slot": 43u64, + "confirmations": null, + "err": null, + "status": { "Ok": null }, + "confirmationStatus": "finalized", + }); + let statuses = serde_json::json!({ + "context": { "slot": 43u64, "apiVersion": "2.0.0" }, + "value": [status.clone(), status], + }); + SolanaRPC::new_mock_with_mocks(Mocks::from([(RpcRequest::GetSignatureStatuses, statuses)])) } /// `decode` wraps settlement events as `DecodedEvent::Settlement` for `run` @@ -655,6 +679,7 @@ async fn solana_db_ingester_to_decoder_persists_decoded_events() { let mut ingester = Ingester::new(geyser_stream, sender, Arc::new(AtomicU64::new(0))); let mut decoder = Decoder::new( Postgres::new(pool.clone()), + mock_rpc_with_signature_statuses(), receiver, settlement, Some(solflow), diff --git a/crates/solana-indexer/src/persistence.rs b/crates/solana-indexer/src/persistence.rs index 35f0e5c5d8..6aca8c0ed9 100644 --- a/crates/solana-indexer/src/persistence.rs +++ b/crates/solana-indexer/src/persistence.rs @@ -82,10 +82,11 @@ impl Postgres { async fn apply( tx: &mut PgTransaction<'_>, event: DecodedEvent, + slot: Slot, ) -> Result<(), PersistenceError> { match event { DecodedEvent::Settlement(SettlementEvent::OrderCreated(order)) => { - Self::apply_order_created(tx, &order).await + Self::apply_order_created(tx, &order, slot).await } DecodedEvent::Settlement(SettlementEvent::SettlementFinalized(settlement)) => { Self::apply_settlement_finalized(tx, settlement).await @@ -104,16 +105,20 @@ impl Postgres { async fn apply_order_created( tx: &mut PgTransaction<'_>, order: &CreatedOrder, + slot: Slot, ) -> Result<(), PersistenceError> { sqlx::query( r#" -INSERT INTO solana.order_pda (order_uid, created_by) -VALUES ($1, $2) -ON CONFLICT (order_uid) DO NOTHING +INSERT INTO solana.order_pda (order_uid, created_by, created_by_tx, created_in_slot) +VALUES ($1, $2, $3, $4) +ON CONFLICT (order_uid) DO UPDATE SET is_reorged = false + WHERE order_pda.is_reorged "#, ) .bind(order.order_uid.0) .bind(order.created_by.to_bytes()) + .bind(order.signature.as_ref()) + .bind(to_db_slot(slot)) .execute(&mut **tx) .await?; // creation_timestamp is the indexing time, the stream carries no @@ -305,12 +310,103 @@ WHERE indexer_state.slot < EXCLUDED.slot ) -> Result<(), PersistenceError> { let mut tx = self.pool.begin().await?; for event in events { - Self::apply(&mut tx, event).await?; + Self::apply(&mut tx, event, last_indexed_slot).await?; } Self::upsert_last_indexed_slot(&mut *tx, last_indexed_slot).await?; Ok(tx.commit().await?) } + /// The finalized watermark, `None` before the first flush writes the + /// state row. + pub(crate) async fn finalized_slot(&self) -> Result, PersistenceError> { + let slot: Option = + sqlx::query_scalar("SELECT finalized_slot FROM solana.indexer_state") + .fetch_optional(&self.pool) + .await?; + Ok(slot.map(from_db_slot)) + } + + /// The distinct transaction signatures of rows with slots in + /// `(after, through]`, the rows about to become final. + pub(crate) async fn unfinalized_signatures( + &self, + after: Slot, + through: Slot, + ) -> Result, PersistenceError> { + let rows: Vec> = sqlx::query_scalar( + r#" +SELECT tx_signature FROM solana.settlements WHERE slot > $1 AND slot <= $2 +UNION +SELECT tx_signature FROM solana.dead_letter WHERE slot > $1 AND slot <= $2 +UNION +SELECT created_by_tx FROM solana.order_pda + WHERE created_in_slot > $1 AND created_in_slot <= $2 AND created_by_tx IS NOT NULL + "#, + ) + .bind(to_db_slot(after)) + .bind(to_db_slot(through)) + .fetch_all(&self.pool) + .await?; + Ok(rows + .into_iter() + .filter_map(|bytes| Signature::try_from(bytes.as_slice()).ok()) + .collect()) + } + + /// Advance the finalized watermark past slots whose transactions all + /// still exist, reverting the rows of the `vanished` ones in the same + /// SQL transaction: their trades (with the fill sums the trades added), + /// settlements, and dead letters are deleted, and the orders they + /// created are marked `is_reorged`. + pub(crate) async fn finalize_through( + &self, + finalized: Slot, + vanished: &[Signature], + ) -> Result<(), PersistenceError> { + let mut tx = self.pool.begin().await?; + for signature in vanished { + let signature = signature.as_ref(); + // Aggregated per order: one transaction can carry several + // settlements, and Postgres applies only one FROM row per target. + sqlx::query( + r#" +UPDATE solana.order_pda AS pda +SET amount_withdrawn = pda.amount_withdrawn - deltas.sell, + amount_received = pda.amount_received - deltas.buy +FROM ( + SELECT order_uid, SUM(sell_amount) AS sell, SUM(buy_amount) AS buy + FROM solana.trades WHERE tx_signature = $1 GROUP BY order_uid +) AS deltas +WHERE pda.order_uid = deltas.order_uid + "#, + ) + .bind(signature) + .execute(&mut *tx) + .await?; + for statement in [ + "DELETE FROM solana.trades WHERE tx_signature = $1", + "DELETE FROM solana.settlements WHERE tx_signature = $1", + "DELETE FROM solana.dead_letter WHERE tx_signature = $1", + // Orders are marked, not deleted: the rows are the audit + // trail, and the stream re-delivering a re-landed creation + // clears the flag. + "UPDATE solana.order_pda SET is_reorged = true WHERE created_by_tx = $1", + ] { + sqlx::query(statement) + .bind(signature) + .execute(&mut *tx) + .await?; + } + } + sqlx::query( + "UPDATE solana.indexer_state SET finalized_slot = GREATEST(finalized_slot, $1)", + ) + .bind(to_db_slot(finalized)) + .execute(&mut *tx) + .await?; + Ok(tx.commit().await?) + } + /// Advance the finalized watermark. Update-only: before the first flush /// there is no state row and nothing indexed to finalize. A backward /// write is a no-op. @@ -381,6 +477,143 @@ mod tests { sqlx::{PgPool, Row}, }; + /// Reverting a vanished transaction deletes its settlement, trades, and + /// dead letter, subtracts the fills its trades added, and marks the + /// order it created. Rows of surviving transactions stay. + #[tokio::test] + #[ignore = "needs the solana.* schema applied locally, run with --test-threads 1"] + async fn solana_db_finalize_through_reverts_vanished_transactions() { + let pool = pool().await; + wipe(&pool).await; + let postgres = Postgres::new(pool.clone()); + let survivor = Signature::from([1; 64]); + let vanished = Signature::from([2; 64]); + + // Order O: created by the surviving transaction, traded by both. + SeedOrder::new([0x01; 32]).insert(&pool).await; + // Order P: created by the vanished transaction. + SeedOrder::new([0x02; 32]).insert(&pool).await; + for (uid, creation, slot, withdrawn, received) in [ + ([0x01u8; 32], survivor, 40i64, 500i64, 700i64), + ([0x02u8; 32], vanished, 41, 0, 0), + ] { + sqlx::query( + "INSERT INTO solana.order_pda + (order_uid, created_by, created_by_tx, created_in_slot, + amount_withdrawn, amount_received) + VALUES ($1, $1, $2, $3, $4, $5)", + ) + .bind(uid) + .bind(creation.as_ref()) + .bind(slot) + .bind(withdrawn) + .bind(received) + .execute(&pool) + .await + .unwrap(); + } + // The vanished transaction carries two settlements trading order O, + // so the revert must aggregate their deltas. The survivor trades O + // too, its rows and sums must stay. + for (signature, instruction_index, slot, sell, buy) in [ + (survivor, 0i32, 40i64, 100i64, 200i64), + (vanished, 0, 41, 300, 400), + (vanished, 5, 41, 100, 100), + ] { + sqlx::query( + "INSERT INTO solana.settlements + (slot, tx_signature, instruction_index, solver, auction_id) + VALUES ($1, $2, $3, $4, 7)", + ) + .bind(slot) + .bind(signature.as_ref()) + .bind(instruction_index) + .bind([0xCCu8; 32]) + .execute(&pool) + .await + .unwrap(); + sqlx::query( + "INSERT INTO solana.trades + (tx_signature, instruction_index, order_uid, sell_amount, + buy_amount, fee_amount) + VALUES ($1, $2, $3, $4, $5, 0)", + ) + .bind(signature.as_ref()) + .bind(instruction_index) + .bind([0x01u8; 32]) + .bind(sell) + .bind(buy) + .execute(&pool) + .await + .unwrap(); + } + sqlx::query( + "INSERT INTO solana.dead_letter (tx_signature, slot, reason) + VALUES ($1, 41, 'decoder_error')", + ) + .bind(vanished.as_ref()) + .execute(&pool) + .await + .unwrap(); + postgres.write_last_indexed_slot(Slot(45)).await.unwrap(); + + // Both signatures sit in the audit window. The query promises no + // order, so sort before comparing. + let mut audited = postgres + .unfinalized_signatures(Slot(39), Slot(45)) + .await + .unwrap(); + audited.sort(); + assert_eq!(audited, vec![survivor, vanished]); + + postgres + .finalize_through(Slot(45), &[vanished]) + .await + .unwrap(); + + let reorged: Vec<(Vec, bool)> = + sqlx::query_as("SELECT order_uid, is_reorged FROM solana.order_pda ORDER BY order_uid") + .fetch_all(&pool) + .await + .unwrap(); + assert_eq!( + reorged, + vec![(vec![0x01; 32], false), (vec![0x02; 32], true)] + ); + let sums: Vec<(Vec, i64, i64)> = sqlx::query_as( + "SELECT order_uid, amount_withdrawn::bigint, amount_received::bigint + FROM solana.order_pda ORDER BY order_uid", + ) + .fetch_all(&pool) + .await + .unwrap(); + assert_eq!( + sums, + vec![(vec![0x01; 32], 100, 200), (vec![0x02; 32], 0, 0)] + ); + let trades: Vec<(Vec, i32)> = + sqlx::query_as("SELECT tx_signature, instruction_index FROM solana.trades") + .fetch_all(&pool) + .await + .unwrap(); + assert_eq!(trades, vec![(survivor.as_ref().to_vec(), 0)]); + let settlements: i64 = sqlx::query_scalar("SELECT count(*) FROM solana.settlements") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(settlements, 1); + let dead: i64 = sqlx::query_scalar("SELECT count(*) FROM solana.dead_letter") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(dead, 0); + let finalized: i64 = sqlx::query_scalar("SELECT finalized_slot FROM solana.indexer_state") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(finalized, 45); + } + /// The finalized watermark only moves forward and needs an existing /// state row: before the first flush the update is a no-op. #[tokio::test] diff --git a/crates/solana-indexer/src/run.rs b/crates/solana-indexer/src/run.rs index a4f5993584..b1843e3824 100644 --- a/crates/solana-indexer/src/run.rs +++ b/crates/solana-indexer/src/run.rs @@ -12,6 +12,7 @@ use { yellowstone, }, clap::Parser, + cow_solana_rpc::{CommitmentConfig, SolanaRPC}, observe::metrics::{DEFAULT_METRICS_PORT, LivenessChecking, serve_metrics}, sqlx::{PgPool, postgres::PgPoolOptions}, std::{ @@ -63,11 +64,23 @@ async fn run(config: Config, start_slot: Option) { .expect("database connection"); let mut metrics = serve_probes(pool.clone()); let persistence = Postgres::new(pool); + // Confirmed commitment, matching the stream subscription. + let rpc = SolanaRPC::new_with_timeout_and_commitment( + &config.rpc.endpoint, + config.rpc.request_timeout, + CommitmentConfig::confirmed(), + ); let (tx, rx) = mpsc::channel(INGEST_TO_DECODER_CAPACITY); let settlement_program = config.chain.settlement_program_id; let solflow_program = config.chain.solflow_program_id; - let mut decoder = Decoder::new(persistence.clone(), rx, settlement_program, solflow_program); + let mut decoder = Decoder::new( + persistence.clone(), + rpc, + rx, + settlement_program, + solflow_program, + ); let mut decoder_task = tokio::spawn(async move { decoder.run().await }); let latest_chain_slot = Arc::new(AtomicU64::default()); diff --git a/crates/solana-orderbook/src/infra/db.rs b/crates/solana-orderbook/src/infra/db.rs index a79bb3f5c7..72dc8a1e91 100644 --- a/crates/solana-orderbook/src/infra/db.rs +++ b/crates/solana-orderbook/src/infra/db.rs @@ -45,7 +45,7 @@ SELECT o.uid, o.owner, o.sell_token, o.buy_token, o.sell_token_account, p.cancellation_timestamp FROM solana.orders o LEFT JOIN solana.order_pda p ON p.order_uid = o.uid -WHERE o.uid = $1 +WHERE o.uid = $1 AND NOT COALESCE(p.is_reorged, false) "#; sqlx::query_as(QUERY) .bind(ByteArray(uid)) diff --git a/database/sql-solana/V3__order_provenance.sql b/database/sql-solana/V3__order_provenance.sql new file mode 100644 index 0000000000..319a657dd2 --- /dev/null +++ b/database/sql-solana/V3__order_provenance.sql @@ -0,0 +1,14 @@ +-- Where the order PDA was created, so a creation that rolls back before its +-- slot finalizes can be marked `is_reorged`. Readers skip marked orders, a +-- re-landed creation clears the flag, and NULL provenance (rows predating +-- this migration) counts as already final. +ALTER TABLE solana.order_pda + ADD COLUMN created_by_tx bytea CHECK (created_by_tx IS NULL OR length(created_by_tx) = 64), + ADD COLUMN created_in_slot bigint, + ADD COLUMN is_reorged boolean NOT NULL DEFAULT false; + +-- The finalization audit scans each table by slot range. +CREATE INDEX solana_order_pda_created_in_slot ON solana.order_pda (created_in_slot) + WHERE created_in_slot IS NOT NULL; +CREATE INDEX solana_settlements_slot ON solana.settlements (slot); +CREATE INDEX solana_dead_letter_slot ON solana.dead_letter (slot);