Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/autopilot-svm/src/infra/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
19 changes: 19 additions & 0 deletions crates/cow-solana-rpc/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<bool>, 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(
Expand Down
58 changes: 57 additions & 1 deletion crates/solana-indexer/src/indexer/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ use {
},
recover_discriminator,
},
cow_solana_rpc::SolanaRPC,
solana_sdk::pubkey::Pubkey,
std::collections::BTreeMap,
tokio::sync::mpsc::Receiver,
Expand All @@ -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<StreamUpdate>,

Expand All @@ -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<StreamUpdate>,
settlement_program: Pubkey,
solflow_program: Option<Pubkey>,
) -> Self {
Self {
persistence,
rpc,
rx,
settlement_program,
solflow_program,
Expand All @@ -92,6 +98,7 @@ impl Decoder {
// In-memory mirror of the persisted watermark, spares redundant
// writes and flags late transactions.
let mut watermark: Option<Slot> = None;
let mut finalized_through: Option<Slot> = None;
while let Some(update) = self.rx.recv().await {
let (slot, signature, inner) = match update {
StreamUpdate::Tx {
Expand All @@ -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;
}
};
Expand Down Expand Up @@ -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<Slot>,
) -> 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(());
}
Comment on lines +173 to +176

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Should this branch also call persistence.write_finalized_slot before returning?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The write would be a no-op there. When that branch hits, after came from the persisted watermark (or the in-memory mirror, which only moves after a successful write), so the DB already stores slot or something newer. The one exception is a fresh DB, where after is just the current slot used as a starting point. Persisting that baseline would work too, but the next boot derives the same value again, so it buys nothing.

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?;
Comment on lines +177 to +198

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Risk of mass false-positive reverts when the audit range is large.

after is unbounded below: on the first finalized tick after boot it's the persisted watermark (finalized_slot().unwrap_or(slot)), and after an outage the live geyser stream jumps straight to the current finalized slot, so the first (after, slot] range can span thousands of slots. The same happens whenever audit-RPC failures accumulate and delay finalization.

known_signatures only checks the recent-status cache (~150 slots). Any transaction in that range that finalized normally but has already aged out of the cache comes back as not known → vanished → reverted. So an indexer that catches up over an old range would delete healthy settlements/trades and mark healthy orders is_reorged, exactly the data it just indexed.

The PR description's claim that "an RPC failure … never reverts on its own" holds for a single delayed tick, but not for a range that has grown past the cache window. Consider bounding the audit to the cache window (e.g. don't revert on slot - created_in_slot older than the cache; advance the watermark without auditing those), or pass searchTransactionHistory: true for the aged tail.

Fix this →

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I believe #4877 is the answer for such long missing slot ranges, so the audit introduced in this PR should not be authoritative over them.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

With the history search from 3d7b809, a cache miss no longer reads as vanished, so a large catch-up range cannot mass-revert healthy rows. I see the split as: #4877 recovers data we never saw, the audit removes data the chain dropped, and the backfill cannot do the latter.

}
*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.
Expand Down
27 changes: 26 additions & 1 deletion crates/solana-indexer/src/indexer/decoder/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -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`
Expand Down Expand Up @@ -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),
Expand Down
Loading
Loading