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
7 changes: 6 additions & 1 deletion SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
121 changes: 103 additions & 18 deletions src/receipt_inventory/backfill.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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?;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Comment thread
rouzwelt marked this conversation as resolved.

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
Expand Down Expand Up @@ -533,11 +540,11 @@ where
)))
}

async fn process_discovery(
async fn read_discovery_balance(
Comment thread
rouzwelt marked this conversation as resolved.
&self,
discovery: ReceiptDiscovery,
read_block: u64,
) -> Result<ProcessOutcome, BackfillError> {
) -> Result<(ReceiptDiscovery, U256), BackfillError> {
let receipt_contract =
Receipt::new(self.receipt_contract, &self.provider);

Expand All @@ -551,6 +558,14 @@ where
.call()
.await?;

Ok((discovery, current_balance))
}

async fn process_discovery(
&self,
discovery: ReceiptDiscovery,
current_balance: U256,
) -> Result<ProcessOutcome, BackfillError> {
if current_balance.is_zero() {
return Ok(ProcessOutcome::ZeroBalance);
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<Log> = (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();
Expand Down
Loading