From 9f64b314e913e1d2232ab2cb23d6647b1c4857ed Mon Sep 17 00:00:00 2001 From: rouzwelt Date: Fri, 2 Oct 2026 03:47:37 +0000 Subject: [PATCH] fix: RAI-2835 - startup backfill fails with an aggregate conflict when one pass finds many receipts --- SPEC.md | 7 +- src/receipt_inventory/backfill.rs | 121 +++++++++++++++++++++++++----- 2 files changed, 109 insertions(+), 19 deletions(-) diff --git a/SPEC.md b/SPEC.md index 126a1e1e..d95ff0a3 100644 --- a/SPEC.md +++ b/SPEC.md @@ -3230,7 +3230,12 @@ would move past its block for good. A node that does not have the read block yet fails the pass, so the checkpoint stays and the next pass retries. At startup, a failed pass stops startup, as any backfill RPC error does, and the service restarts. The fresh head keeps the reads at recent state after a long scan, -which a node that is not an archive node still serves. Live monitoring processes +which a node that is not an archive node still serves. A pass reads its balances +concurrently, but records its discoveries one at a time, in the order it +collected them: every discovery writes to the vault's one inventory aggregate, +and concurrent writes race on its optimistic concurrency check. A pass with many +discoveries (the restart after a rollback finds every returned receipt) could +otherwise lose that race on every retry and fail. Live monitoring processes observed logs opportunistically but does not advance the durable checkpoint, because WebSocket logs can arrive out of order within or across blocks. This prevents long-running services from restarting with a stale receipt checkpoint diff --git a/src/receipt_inventory/backfill.rs b/src/receipt_inventory/backfill.rs index d40f9602..be1f4b7d 100644 --- a/src/receipt_inventory/backfill.rs +++ b/src/receipt_inventory/backfill.rs @@ -1,12 +1,12 @@ use alloy::eips::BlockId; -use alloy::primitives::{Address, Bytes, TxHash}; +use alloy::primitives::{Address, Bytes, TxHash, U256}; use alloy::rpc::types::{Filter, Log}; use alloy::sol_types::SolEvent; use alloy::transports::{RpcError, TransportErrorKind}; use async_trait::async_trait; use cqrs_es::AggregateError; use event_sorcery::Store; -use futures::{StreamExt, stream}; +use futures::{StreamExt, TryStreamExt, stream}; use itertools::Itertools; use sqlx::{Pool, Sqlite}; use std::sync::Arc; @@ -266,20 +266,27 @@ where .flatten() .unique_by(|discovery| discovery.receipt_id); - let results: Vec<_> = stream::iter(discoveries) - .map(|discovery| self.process_discovery(discovery, read_block)) - .buffer_unordered(MAX_CONCURRENT_BALANCE_CHECKS) - .collect() - .await; + // Balance reads are network calls, so they run concurrently. Every + // discovery then writes to this vault's one inventory aggregate, so + // the writes run one at a time, in the order the pass collected the + // logs: concurrent writers race on the aggregate's optimistic + // concurrency check, and a pass with many discoveries can lose that + // race on every retry and fail. A rollback's restart finds every + // returned receipt in one pass. + let readings: Vec<_> = stream::iter(discoveries) + .map(|discovery| self.read_discovery_balance(discovery, read_block)) + .buffered(MAX_CONCURRENT_BALANCE_CHECKS) + .try_collect() + .await?; - let (processed_count, skipped_zero_balance) = results - .into_iter() - .try_fold((0u64, 0u64), |(processed, zero), result| { - result.map(|outcome| match outcome { - ProcessOutcome::Processed => (processed + 1, zero), - ProcessOutcome::ZeroBalance => (processed, zero + 1), - }) - })?; + let mut processed_count = 0u64; + let mut skipped_zero_balance = 0u64; + for (discovery, current_balance) in readings { + match self.process_discovery(discovery, current_balance).await? { + ProcessOutcome::Processed => processed_count += 1, + ProcessOutcome::ZeroBalance => skipped_zero_balance += 1, + } + } // Process reconciliation events (Withdraw + outbound transfers). // Once a recorded migration moved this vault's custody away from the @@ -533,11 +540,11 @@ where ))) } - async fn process_discovery( + async fn read_discovery_balance( &self, discovery: ReceiptDiscovery, read_block: u64, - ) -> Result { + ) -> Result<(ReceiptDiscovery, U256), BackfillError> { let receipt_contract = Receipt::new(self.receipt_contract, &self.provider); @@ -551,6 +558,14 @@ where .call() .await?; + Ok((discovery, current_balance)) + } + + async fn process_discovery( + &self, + discovery: ReceiptDiscovery, + current_balance: U256, + ) -> Result { if current_balance.is_zero() { return Ok(ProcessOutcome::ZeroBalance); } @@ -714,7 +729,7 @@ where /// Queries the on-chain balance of a receipt at `read_block` and /// reconciles the aggregate. /// - /// The read is pinned for the same reason as in `process_discovery`. A + /// The read is pinned for the reason `backfill_receipts` gives. A /// burn of ours that lands after `read_block` can settle in inventory /// before this read, so the read can briefly restore the shares that /// burn consumed. The next pass scans the burn's block and reconciles @@ -1057,6 +1072,76 @@ mod tests { ); } + /// Every discovery in a pass writes to the vault's one inventory + /// aggregate. Concurrent writers race on its optimistic concurrency check, + /// and a pass this size could lose the race on every retry and fail, which + /// at startup stops the service. The restart after a rollback finds every + /// returned receipt in one pass, so the pass must never race itself. + #[tokio::test] + #[traced_test] + async fn backfill_records_many_discoveries_without_racing_itself() { + let (receipt_contract, bot_wallet, vault) = test_addresses(); + let (store, pool) = setup_store().await; + let balance = U256::from(1000); + let receipt_count = 65u64; + + let deposit_logs: Vec = (1..=receipt_count) + .map(|receipt_id| { + create_deposit_log(DepositLogParams { + vault, + sender: bot_wallet, + owner: bot_wallet, + assets: balance, + shares: balance, + id: U256::from(receipt_id), + receipt_information: Bytes::new(), + tx_hash: B256::from(U256::from(receipt_id)), + block_number: 100 + receipt_id, + }) + }) + .collect(); + + let asserter = Asserter::new(); + asserter.push_success(&deposit_logs); + push_empty_non_deposit_logs(&asserter); + // eth_blockNumber (the read block) + asserter.push_success(&U256::from(200u64)); + for _ in 0..receipt_count { + asserter.push_success(&balance.to_be_bytes::<32>()); + } + + let provider = ProviderBuilder::new() + .wallet(EthereumWallet::from(PrivateKeySigner::random())) + .connect_mocked_client(asserter); + + let backfiller = ReceiptBackfiller::new(ReceiptBackfillDeps { + provider, + receipt_contract, + bot_wallet, + chain_id: ANVIL_CHAIN_ID, + network: Network::Base, + vault, + store: store.clone(), + pool, + handler: NoOpItnHandler, + }); + + let result = backfiller.backfill_receipts(0, 200).await.unwrap(); + + assert_eq!(result.processed_count, receipt_count); + let inventory = + load_inventory(&store, ANVIL_CHAIN_ID, &vault).await.unwrap(); + assert_eq!( + inventory.receipts_with_balance().len(), + usize::try_from(receipt_count).unwrap(), + "inventory must track every receipt the pass found" + ); + assert!( + !logs_contain("optimistic-concurrency conflict"), + "one pass must never race itself on the inventory aggregate" + ); + } + #[tokio::test] async fn backfill_is_idempotent() { let (receipt_contract, bot_wallet, vault) = test_addresses();