Skip to content
178 changes: 167 additions & 11 deletions packages/rs-platform-wallet/src/changeset/core_bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@

use std::collections::{BTreeMap, HashMap, HashSet};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::sync::{Arc, Mutex};

use dashcore::blockdata::transaction::{txout::TxOut, OutPoint};
use dashcore::ScriptBuf;
Expand Down Expand Up @@ -294,10 +294,16 @@ async fn run_wallet_event_adapter<P>(
P: PlatformWalletPersistence + 'static,
{
tracing::debug!("wallet-event adapter task started");
let mut fault = AdapterFaultState::default();
// Both live behind handles rather than as locals because the commit runs
// on a blocking thread (see the `spawn_blocking` below) and has to be able
// to carry its state across drains. Moving them into the closure by value
// would lose a wallet's frozen watermark if that thread ever panicked —
// un-freezing a wallet that failed verification is the one outcome the
// fail-closed guard exists to prevent.
let fault = Arc::new(Mutex::new(AdapterFaultState::default()));
// One-shot latch so the hard "watermark frozen" line hits logcat exactly
// once per session rather than once per faulted batch.
let mut freeze_logged = false;
let freeze_logged = Arc::new(AtomicBool::new(false));

loop {
// Block for the first event of a batch. Everything already sitting in
Expand Down Expand Up @@ -360,14 +366,84 @@ async fn run_wallet_event_adapter<P>(
// Commit the folded batch. The channel is lossless, so the only way a
// watermark is held back is a rejected `store()` (the fail-closed
// backstop inside `commit_batch`).
let diag = commit_batch(
&*persister,
batch,
folded,
&mut fault,
&sync_fault,
&mut freeze_logged,
);
// Commit on a blocking thread, never on the async worker.
//
// `store()` is synchronous and, for the SQLite backend, commits a real
// transaction per call — its own docs warn that a slow write blocks
// every other wallet accessor for its duration. Called inline here it
// blocked a tokio worker instead: a field restore showed one drain of
// 512 folded events park the runtime long enough for the metrics tick
// covering it to report a 1.4s mean poll, with the whole sync stalled
// for minutes at a time and the durable watermark left hundreds of
// thousands of blocks behind the chain tip.
//
// The handle is awaited rather than raced against `cancel`: a store
// that has started must be allowed to finish, and dropping the handle
// would not stop the thread anyway. Shutdown is observed at the next
// `recv` instead.
// Captured before the batch moves into the closure: if the commit
// thread panics, these are the wallets whose rows have an unknown fate
// and whose watermark must therefore be frozen.
let batch_wallet_ids: Vec<WalletId> = batch.keys().copied().collect();
let persister_for_commit = Arc::clone(&persister);
let sync_fault_for_commit = Arc::clone(&sync_fault);
let fault_for_commit = Arc::clone(&fault);
let freeze_for_commit = Arc::clone(&freeze_logged);
let committed = tokio::task::spawn_blocking(move || {
// The lock is uncontended by construction — this task is the only
// writer, and one drain commits at a time — so it never blocks;
// it exists to carry the state, not to arbitrate.
let mut fault = fault_for_commit
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let mut freeze_logged = freeze_for_commit.load(Ordering::Relaxed);
let diag = commit_batch(
&*persister_for_commit,
batch,
folded,
&mut fault,
&sync_fault_for_commit,
&mut freeze_logged,
);
freeze_for_commit.store(freeze_logged, Ordering::Relaxed);
diag
})
.await;
Comment thread
romchornyi marked this conversation as resolved.

let diag = match committed {
Ok(diag) => diag,
// The commit thread panicked, so `commit_batch` never reached the
// `store()` rejection arm that would have frozen the affected
// wallets. Freeze them here instead.
//
// Before this call moved off the runtime a panic unwound the whole
// adapter task, which stopped every later watermark advance by
// killing the writer. `spawn_blocking` turns that into a recoverable
// `JoinError`, and simply continuing would let the NEXT batch
// persist a higher `synced_height` for a wallet whose rows from this
// batch may never have landed — the exact hole the fail-closed rule
// exists to prevent. Faulting per wallet rather than stopping the
// adapter keeps the existing design: a wallet whose commit is in
// doubt freezes, its siblings keep syncing.
Err(join_error) => {
{
let mut fault = fault
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
for wallet_id in &batch_wallet_ids {
fault.fault_wallet(*wallet_id, &sync_fault);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
Comment thread
shumkov marked this conversation as resolved.
Outdated
}
tracing::error!(
error = %join_error,
folded,
wallets = batch_wallet_ids.len(),
"wallet-event commit thread failed; freezing the batch's \
wallets because their rows have an unknown outcome"
);
continue;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Comment thread
romchornyi marked this conversation as resolved.
Comment thread
romchornyi marked this conversation as resolved.
}
};

// One structured line per drain via the `log` facade so a tester
// logcat is unambiguous about whether the watermark is advancing.
Expand Down Expand Up @@ -1880,18 +1956,26 @@ mod tests {
struct ProbePersister {
obs: UnboundedSender<StoreObserved>,
fail_once: Mutex<HashSet<WalletId>>,
/// Wallets whose NEXT `store()` panics instead of returning. Models a
/// backend that dies mid-write — the case that used to unwind the whole
/// adapter task and now surfaces as a `JoinError`.
panic_once: Mutex<HashSet<WalletId>>,
}

impl ProbePersister {
fn new(obs: UnboundedSender<StoreObserved>) -> Self {
Self {
obs,
fail_once: Mutex::new(HashSet::new()),
panic_once: Mutex::new(HashSet::new()),
}
}
fn fail_next(&self, wallet_id: WalletId) {
self.fail_once.lock().unwrap().insert(wallet_id);
}
fn panic_next(&self, wallet_id: WalletId) {
self.panic_once.lock().unwrap().insert(wallet_id);
}
}

impl PlatformWalletPersistence for ProbePersister {
Expand All @@ -1901,6 +1985,9 @@ mod tests {
changeset: PlatformWalletChangeSet,
) -> Result<(), PersistenceError> {
let core = changeset.core.as_ref();
if self.panic_once.lock().unwrap().remove(&wallet_id) {
panic!("probe persister: store panicked for {wallet_id:?}");
}
let rejected = self.fail_once.lock().unwrap().remove(&wallet_id);
let _ = self.obs.send(StoreObserved {
wallet_id,
Expand Down Expand Up @@ -2253,6 +2340,75 @@ mod tests {
);
}

/// (h) SAFETY INVARIANT under a commit-thread PANIC: a wallet whose
/// `store()` panicked must be frozen just as if the store had been
/// rejected, because its rows have an unknown fate.
///
/// This is a regression guard on the move to `spawn_blocking`. Before it,
/// a panic unwound the adapter task itself, which stopped every later
/// watermark advance by killing the writer outright. `spawn_blocking`
/// turns that into a recoverable `JoinError` — and merely logging it would
/// let the NEXT batch persist a higher `synced_height` for a wallet whose
/// earlier rows may never have landed, which is exactly the hole
/// dashpay/platform#4069 closed.
#[tokio::test]
async fn a_panicking_commit_freezes_the_batch_wallets() {
let wallet_id = [0xEEu8; 32];
let (tx, rx) = unbounded_channel::<WalletEvent>();

let (obs_tx, mut obs_rx) = unbounded_channel();
let persister = Arc::new(ProbePersister::new(obs_tx));
persister.panic_next(wallet_id);
let sync_fault = Arc::new(AtomicBool::new(false));
let cancel = CancellationToken::new();
let handle = tokio::spawn(run_wallet_event_adapter(
test_manager(),
Arc::clone(&persister),
rx,
Arc::clone(&sync_fault),
cancel.clone(),
));

// The store for this batch panics: no observation is emitted, and the
// adapter must fault the wallet rather than carry on unaffected.
tx.send(block_processed_event(wallet_id, 10)).unwrap();
// Bounded, so a regression fails the test instead of hanging it: with
// the fault-on-panic path removed, `sync_fault` is simply never raised
// and an unbounded spin would wedge CI with no diagnosis.
tokio::time::timeout(std::time::Duration::from_secs(5), async {
while !sync_fault.load(Ordering::Relaxed) {
tokio::task::yield_now().await;
}
})
.await
.expect("a panicked commit must raise the hard-fault signal");

// The adapter must still be alive — the point of moving the commit off
// the runtime is that one bad batch does not take the writer with it.
assert!(
!handle.is_finished(),
"a panicked commit must not kill the adapter"
);

// A later watermark for the same wallet must not reach the store.
tx.send(block_processed_event(wallet_id, 60)).unwrap();
tx.send(sync_height_event(wallet_id, 900)).unwrap();

let post = obs_rx
.recv()
.await
.expect("the record-bearing event must still persist while faulted");
assert_eq!(post.wallet_id, wallet_id);
assert_eq!(
post.synced_height, None,
"a wallet whose commit panicked must not advance its durable watermark"
);

cancel.cancel();
drop(tx);
handle.await.unwrap();
}

/// (g) SAFETY INVARIANT under a fault: once the per-wallet fault latch is
/// set (here by a rejected `store()`), no later changeset for that wallet
/// may advance the durable `synced_height` — whether the watermark arrives
Expand Down
Loading