Skip to content
Closed
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
31 changes: 30 additions & 1 deletion keeper/src/accessor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,27 @@ impl<'a> TransactionsAccessor<'a> {
Ok(true)
}

/// Serves waiters on `signature` the outcome of the submission that already
/// owns it, after a duplicate of it was dropped.
///
/// A duplicate is not a failure — it *is* the transaction already in flight
/// or already settled, and its submitter is entitled to that transaction's
/// result. Which is what it gets while the original is in flight: the
/// signature's channel still exists, so both submitters share it and
/// `commit_execution` reaches both. Once the original settles, that channel
/// is gone; a duplicate arriving afterwards subscribes to a freshly created
/// one that nothing will ever write to, and blocks until its caller's own
/// deadline. Replaying the cached status closes that case.
pub fn notify_duplicate(&self, signature: Signature) {
// A cached `None` means still in flight, so the original's own commit
// will serve both. Absent entirely means it has aged out of the dedup
// window, in which case `append` would not have rejected this one.
let Some(Some(status)) = self.keeper.caches.signatures.get(&signature) else {
return;
};
self.keeper.subscriptions.signatures.send(&signature, &status, true);
}

/// Commits execution metadata and publishes resulting account changes.
pub fn commit_execution(&self, mut txn: FullTransaction) -> Result<()> {
let subs = &self.keeper.subscriptions;
Expand Down Expand Up @@ -242,8 +263,16 @@ impl<'a> TransactionsAccessor<'a> {
}
// Clear TLS unconditionally so unsent messages cannot leak into the next transaction.
TlsManager::clear();
// Cache before broadcasting, not after. The broadcast is terminal: it
// drops the signature's channel, so from that instant onwards the cache
// is the only way to learn this outcome. Updating it second leaves a
// window in which a late subscriber finds neither — a fresh channel
// nobody will ever write to, and a cache entry still reading `None`.
self.keeper
.caches
.signatures
.update(&commit.signature, Some(commit.status.clone()));
subs.signatures.send(&commit.signature, &commit.status, true);
self.keeper.caches.signatures.update(&commit.signature, Some(commit.status));
Ok(())
}

Expand Down
67 changes: 67 additions & 0 deletions keeper/src/tests/subscriptions.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
//! Subscription fanout primitives, transaction-append dedup

use ledger::request::TransactionStatus;

use super::{TestKeeper, signed_tx};
use crate::subscriptions::Subscribers;

Expand Down Expand Up @@ -61,3 +63,68 @@ async fn append_dedup_and_status_sentinel() {

keeper.close().await;
}

// A duplicate submitted after the original settled must still be told the
// original's outcome.
//
// Waiters are keyed by signature and the settling broadcast is terminal, so it
// takes the channel with it. A duplicate arriving afterwards subscribes to a
// freshly created channel that nothing will ever write to again, and without
// `notify_duplicate` it simply blocks until its caller's own deadline — which is
// how a transaction that executed perfectly well gets reported upstream as a
// timeout.
#[tokio::test]
async fn duplicate_after_settlement_is_served_the_original_status() {
let keeper = TestKeeper::new().await;
let (signature, txn) = signed_tx();
assert!(keeper.transactions().append(&txn).await.unwrap());

// Settle it the way `commit_execution` does: cache first, then the terminal
// broadcast that drops the channel.
let settled = TransactionStatus {
result: Ok(()),
slot: 7,
};
keeper.caches.signatures.update(&signature, Some(settled.clone()));
keeper.subscriptions.signatures.send(&signature, &settled, true);

// The late duplicate: subscribes to a channel created after the fact, so
// only the cache can serve it.
let mut rx = keeper.transactions().subscribe_signature(signature).await;
assert!(
!keeper.transactions().append(&txn).await.unwrap(),
"duplicate is still deduplicated"
);
keeper.transactions().notify_duplicate(signature);

let received = rx.try_recv().expect("late duplicate is served the settled status");
assert_eq!(received.slot, settled.slot);
assert!(received.result.is_ok());

keeper.close().await;
}

// While the original is still in flight there is nothing to replay, and
// replaying a sentinel would settle the waiter on a status that does not exist
// yet. The shared channel is what serves both submitters in that case.
#[tokio::test]
async fn duplicate_while_in_flight_is_left_to_the_shared_channel() {
let keeper = TestKeeper::new().await;
let (signature, txn) = signed_tx();
assert!(keeper.transactions().append(&txn).await.unwrap());

let mut rx = keeper.transactions().subscribe_signature(signature).await;
keeper.transactions().notify_duplicate(signature);
assert!(rx.try_recv().is_err(), "nothing is published for an unsettled signature");

// ...and when it does settle, that same channel delivers.
let settled = TransactionStatus {
result: Ok(()),
slot: 9,
};
keeper.caches.signatures.update(&signature, Some(settled.clone()));
keeper.subscriptions.signatures.send(&signature, &settled, true);
assert_eq!(rx.try_recv().expect("in-flight waiter is served").slot, 9);

keeper.close().await;
}
6 changes: 6 additions & 0 deletions processor/src/sequencer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,12 @@ impl Sequencer {
};
if !self.replay {
if !self.state.transactions().append(&txn).await? {
// Dropping a duplicate is right — it must not execute twice —
// but dropping it *silently* strands whoever submitted it,
// because waiters are keyed by signature and no further status
// will ever be published for one that has already settled. Hand
// them the original's result instead.
self.state.transactions().notify_duplicate(txn.signatures()[0]);
metrics::failed_transaction(FailureKind::SequencerDrop);
return Ok(());
}
Expand Down