-
Notifications
You must be signed in to change notification settings - Fork 192
solana-indexer: revert rolled-back transactions when their slots finalize #4875
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
56bbde4
61b1fa3
1bbfdab
be99aa8
7bad48c
5be0190
ac363cf
8e7c905
e2ec78a
64cf65c
70fd3cd
3d7b809
9352df7
a101938
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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<StreamUpdate>, | ||
|
|
||
|
|
@@ -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, | ||
|
|
@@ -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 { | ||
|
|
@@ -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<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(()); | ||
| } | ||
| 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
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Risk of mass false-positive reverts when the audit range is large.
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
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. |
||
| } | ||
| *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. | ||
|
|
||
There was a problem hiding this comment.
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_slotbefore returning?There was a problem hiding this comment.
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,
aftercame 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, whereafteris 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.