From 1aaeb88eb3ccc8086dbc8fad22c8257beedc642a Mon Sep 17 00:00:00 2001 From: octobocto Date: Sun, 27 Sep 2026 07:43:51 -0700 Subject: [PATCH] RPC: expose the two way peg events of a block --- app/rpc_server.rs | 10 +++ cli/lib.rs | 8 +++ lib/node/error.rs | 4 +- lib/node/mod.rs | 42 +++++++++--- lib/state/mod.rs | 26 +++++++- lib/state/two_way_peg_data.rs | 117 ++++++++++++++++++++++++++++++++++ rpc-api/lib.rs | 14 +++- types/state.rs | 17 ++++- 8 files changed, 223 insertions(+), 15 deletions(-) diff --git a/app/rpc_server.rs b/app/rpc_server.rs index 82b0392f..69f85c8f 100644 --- a/app/rpc_server.rs +++ b/app/rpc_server.rs @@ -281,6 +281,16 @@ impl rpc_api::node::RpcServer Ok(res) } + async fn get_two_way_peg_events( + &self, + block_hash: thunder::types::BlockHash, + ) -> RpcResult> { + self.app + .node + .get_two_way_peg_events(block_hash) + .map_err(custom_err) + } + async fn get_utxos( &self, addresses: HashSet
, diff --git a/cli/lib.rs b/cli/lib.rs index 085eb600..431c2b38 100644 --- a/cli/lib.rs +++ b/cli/lib.rs @@ -102,6 +102,10 @@ pub enum Command { }, /// Get transaction by txid GetTransaction { txid: Txid }, + /// Get the coin movements that a block applied outside its body + GetTwoWayPegEvents { + block_hash: thunder_types::BlockHash, + }, /// Get utxos for addresses GetUtxos { #[arg(required = true)] @@ -298,6 +302,10 @@ where let tx_info = rpc_client.get_transaction(txid).await?; serde_json::to_string_pretty(&tx_info)? } + Command::GetTwoWayPegEvents { block_hash } => { + let events = rpc_client.get_two_way_peg_events(block_hash).await?; + serde_json::to_string_pretty(&events)? + } Command::GetUtxos { addresses } => { let addresses = addresses.into_iter().collect(); let utxos = rpc_client.get_utxos(addresses).await?; diff --git a/lib/node/error.rs b/lib/node/error.rs index 1d30f085..61cfde80 100644 --- a/lib/node/error.rs +++ b/lib/node/error.rs @@ -4,7 +4,7 @@ use transitive::Transitive; use crate::{ archive, mempool, net, state, - types::{AmountOverflowError, AmountUnderflowError, proto}, + types::{AmountOverflowError, AmountUnderflowError, BlockHash, proto}, }; pub mod mainchain_task { @@ -176,6 +176,8 @@ pub enum Error { Net(#[from] Box), #[error("net task error")] NetTask(#[source] Box), + #[error("block {block_hash} is not in the active chain")] + NotInActiveChain { block_hash: BlockHash }, #[error("peer info stream closed")] PeerInfoRxClosed, #[error("Receive mainchain task response cancelled")] diff --git a/lib/node/mod.rs b/lib/node/mod.rs index 0319c17f..fcbf0bed 100644 --- a/lib/node/mod.rs +++ b/lib/node/mod.rs @@ -9,7 +9,7 @@ use std::{ use bitcoin::amount::CheckedSum; use fallible_iterator::{FallibleIterator, IteratorExt}; use futures::Stream; -use sneed::{DbError, Env, EnvError, RwTxnError}; +use sneed::{DbError, Env, EnvError, RoTxn, RwTxnError}; use tokio::sync::Mutex; use tonic::transport::Channel; @@ -27,7 +27,7 @@ use crate::{ authorization::{BatchVerificationContext, rand_core::CryptoRng}, net::{Peer, PeerAddress, ResolvedPeerAddress}, proto::{self, mainchain}, - state::WithdrawalBundleInfo, + state::{TwoWayPegEvent, WithdrawalBundleInfo}, }, util::Watchable, }; @@ -392,22 +392,20 @@ where Ok(self.archive.get_header(&txn, block_hash)?) } - /// Get the block hash at the specified height in the active chain, - /// if it exists - pub fn try_get_block_hash( + fn try_get_block_hash_at( &self, + rotxn: &RoTxn, height: u32, ) -> Result, Error> { - let rotxn = self.env.read_txn().map_err(EnvError::from)?; - let Some(tip) = self.state.try_get_tip(&rotxn)? else { + let Some(tip) = self.state.try_get_tip(rotxn)? else { return Ok(None); }; - let Some(tip_height) = self.state.try_get_height(&rotxn)? else { + let Some(tip_height) = self.state.try_get_height(rotxn)? else { return Ok(None); }; if tip_height >= height { self.archive - .ancestors(&rotxn, tip) + .ancestors(rotxn, tip) .nth((tip_height - height) as usize) .map_err(Error::from) } else { @@ -415,6 +413,32 @@ where } } + /// Get the block hash at the specified height in the active chain, + /// if it exists + pub fn try_get_block_hash( + &self, + height: u32, + ) -> Result, Error> { + let rotxn = self.env.read_txn().map_err(EnvError::from)?; + self.try_get_block_hash_at(&rotxn, height) + } + + /// Get the coin movements that a block applied outside its body, in the + /// order the node applied them + pub fn get_two_way_peg_events( + &self, + block_hash: BlockHash, + ) -> Result, Error> { + let rotxn = self.env.read_txn().map_err(EnvError::from)?; + let height = self.archive.get_height(&rotxn, block_hash)?; + // The events are keyed by height, so a block off the active chain + // would read another block's events. + if self.try_get_block_hash_at(&rotxn, height)? != Some(block_hash) { + return Err(Error::NotInActiveChain { block_hash }); + } + Ok(self.state.get_two_way_peg_events(&rotxn, height)?) + } + pub fn try_get_body( &self, block_hash: BlockHash, diff --git a/lib/state/mod.rs b/lib/state/mod.rs index 7284c61d..679d4f08 100644 --- a/lib/state/mod.rs +++ b/lib/state/mod.rs @@ -21,7 +21,7 @@ use crate::{ WithdrawalBundle, WithdrawalBundleStatus, authorization::{self, BatchVerificationContext}, proto::mainchain::TwoWayPegData, - state::WithdrawalBundleInfo, + state::{TwoWayPegEvent, WithdrawalBundleInfo}, }, util::Watchable, }; @@ -79,12 +79,16 @@ pub struct State { SerdeBincode, SerdeBincode<(bitcoin::BlockHash, u32)>, >, + /// Coin movements that no block body carries, keyed by the height that + /// applied them, in the order the node applied them + two_way_peg_events: + DatabaseUnique, SerdeBincode>>, pub utreexo_accumulator: DatabaseUnique>, _version: DatabaseUnique>, } impl State { - pub const NUM_DBS: u32 = 11; + pub const NUM_DBS: u32 = 12; pub fn new(env: &sneed::Env) -> Result { let mut rwtxn = env.write_txn().map_err(EnvError::from)?; @@ -120,6 +124,9 @@ impl State { "withdrawal_bundle_event_blocks", ) .map_err(EnvError::from)?; + let two_way_peg_events = + DatabaseUnique::create(env, &mut rwtxn, "two_way_peg_events") + .map_err(EnvError::from)?; let utreexo_accumulator = DatabaseUnique::create(env, &mut rwtxn, "utreexo_accumulator") .map_err(EnvError::from)?; @@ -153,11 +160,26 @@ impl State { withdrawal_bundles, deposit_blocks, withdrawal_bundle_event_blocks, + two_way_peg_events, utreexo_accumulator, _version: version, }) } + /// Coin movements that the block at this height applied outside its body, + /// in the order the node applied them + pub fn get_two_way_peg_events( + &self, + rotxn: &RoTxn, + height: u32, + ) -> Result, Error> { + let events = self + .two_way_peg_events + .try_get(rotxn, &height)? + .unwrap_or_default(); + Ok(events) + } + pub fn try_get_tip( &self, rotxn: &RoTxn, diff --git a/lib/state/two_way_peg_data.rs b/lib/state/two_way_peg_data.rs index fe1aaf3e..484b7166 100644 --- a/lib/state/two_way_peg_data.rs +++ b/lib/state/two_way_peg_data.rs @@ -16,6 +16,7 @@ use crate::{ PointedOutputRef, SpentOutput, WithdrawalBundle, WithdrawalBundleEvent, WithdrawalBundleEventStatus, WithdrawalBundleStatus, hash, proto::mainchain::{BlockEvent, TwoWayPegData}, + state::TwoWayPegEvent, }, }; @@ -118,6 +119,7 @@ fn connect_withdrawal_bundle_submitted( rwtxn: &mut RwTxn, block_height: u32, accumulator_diff: &mut AccumulatorDiff, + two_way_peg_events: &mut Vec, event_block_hash: &bitcoin::BlockHash, m6id: M6id, ) -> Result<(), error::ConnectWithdrawalBundleSubmitted> { @@ -160,6 +162,10 @@ fn connect_withdrawal_bundle_submitted( inpoint: InPoint::Withdrawal { m6id }, }; state.stxos.put(rwtxn, &key, &spent_output)?; + two_way_peg_events.push(TwoWayPegEvent::BundleSpend { + outpoint: *outpoint, + m6id, + }); } assert_eq!( bundle_status.latest().value, @@ -284,6 +290,7 @@ fn connect_withdrawal_bundle_confirmed( rwtxn: &mut RwTxn, block_height: u32, accumulator_diff: &mut AccumulatorDiff, + two_way_peg_events: &mut Vec, event_block_hash: &bitcoin::BlockHash, m6id: M6id, ) -> Result<(), Error> { @@ -341,6 +348,10 @@ fn connect_withdrawal_bundle_confirmed( output: &spent_output.output, }); accumulator_diff.remove(utxo_hash.into()); + two_way_peg_events.push(TwoWayPegEvent::BundleSpend { + outpoint: *outpoint, + m6id, + }); } state.utxos.clear(rwtxn).map_err(DbError::from)?; bundle = WithdrawalBundleInfo::UnknownConfirmed { @@ -387,6 +398,10 @@ fn connect_withdrawal_bundle_confirmed( output: &spent_output.output, }); accumulator_diff.remove(utxo_hash.into()); + two_way_peg_events.push(TwoWayPegEvent::BundleSpend { + outpoint: *outpoint, + m6id, + }); } } } @@ -406,6 +421,7 @@ fn connect_withdrawal_bundle_failed( rwtxn: &mut RwTxn, block_height: u32, accumulator_diff: &mut AccumulatorDiff, + two_way_peg_events: &mut Vec, m6id: M6id, ) -> Result<(), Error> { tracing::debug!( @@ -448,6 +464,11 @@ fn connect_withdrawal_bundle_failed( output: output.clone(), }); accumulator_diff.insert(utxo_hash.into()); + two_way_peg_events.push(TwoWayPegEvent::BundleReturn { + outpoint: *outpoint, + output: output.clone(), + m6id, + }); } let latest_failed_m6id = if let Some(mut latest_failed_m6id) = state .latest_failed_withdrawal_bundle @@ -482,6 +503,7 @@ fn connect_withdrawal_bundle_event( rwtxn: &mut RwTxn, block_height: u32, accumulator_diff: &mut AccumulatorDiff, + two_way_peg_events: &mut Vec, event_block_hash: &bitcoin::BlockHash, event: &WithdrawalBundleEvent, ) -> Result<(), Error> { @@ -492,6 +514,7 @@ fn connect_withdrawal_bundle_event( rwtxn, block_height, accumulator_diff, + two_way_peg_events, event_block_hash, event.m6id, ) @@ -503,6 +526,7 @@ fn connect_withdrawal_bundle_event( rwtxn, block_height, accumulator_diff, + two_way_peg_events, event_block_hash, event.m6id, ) @@ -513,6 +537,7 @@ fn connect_withdrawal_bundle_event( rwtxn, block_height, accumulator_diff, + two_way_peg_events, event.m6id, ) } @@ -525,6 +550,7 @@ fn connect_event( rwtxn: &mut RwTxn, block_height: u32, accumulator_diff: &mut AccumulatorDiff, + two_way_peg_events: &mut Vec, latest_deposit_block_hash: &mut Option, latest_withdrawal_bundle_event_block_hash: &mut Option, event_block_hash: bitcoin::BlockHash, @@ -540,6 +566,10 @@ fn connect_event( .map_err(DbError::from)?; let utxo_hash = hash(&PointedOutputRef { outpoint, output }); accumulator_diff.insert(utxo_hash.into()); + two_way_peg_events.push(TwoWayPegEvent::Deposit { + outpoint, + output: output.clone(), + }); *latest_deposit_block_hash = Some(event_block_hash); } BlockEvent::WithdrawalBundle(withdrawal_bundle_event) => { @@ -548,6 +578,7 @@ fn connect_event( rwtxn, block_height, accumulator_diff, + two_way_peg_events, &event_block_hash, withdrawal_bundle_event, )?; @@ -570,6 +601,7 @@ pub fn connect( .map_err(DbError::from)? .unwrap_or_default(); let mut accumulator_diff = AccumulatorDiff::default(); + let mut two_way_peg_events = Vec::new(); let mut latest_deposit_block_hash = None; let mut latest_withdrawal_bundle_event_block_hash = None; for (event_block_hash, event_block_info) in &two_way_peg_data.block_info { @@ -579,6 +611,7 @@ pub fn connect( rwtxn, block_height, &mut accumulator_diff, + &mut two_way_peg_events, &mut latest_deposit_block_hash, &mut latest_withdrawal_bundle_event_block_hash, *event_block_hash, @@ -586,6 +619,12 @@ pub fn connect( )?; } } + if !two_way_peg_events.is_empty() { + state + .two_way_peg_events + .put(rwtxn, &block_height, &two_way_peg_events) + .map_err(DbError::from)?; + } // Handle deposits. if let Some(latest_deposit_block_hash) = latest_deposit_block_hash { let deposit_block_seq_idx = state @@ -1023,6 +1062,10 @@ pub fn disconnect( let mut accumulator_diff = AccumulatorDiff::default(); let mut latest_deposit_block_hash = None; let mut latest_withdrawal_bundle_event_block_hash = None; + state + .two_way_peg_events + .delete(rwtxn, &block_height) + .map_err(DbError::from)?; // Restore pending withdrawal bundle for (event_block_hash, event_block_info) in two_way_peg_data.block_info.iter().rev() @@ -1143,6 +1186,7 @@ mod test { WithdrawalBundleEvent, WithdrawalBundleEventStatus, WithdrawalBundleStatus, proto::mainchain::{BlockEvent, BlockInfo, Deposit, TwoWayPegData}, + state::TwoWayPegEvent, }, }; @@ -1373,6 +1417,79 @@ mod test { } // connecting a deposit then disconnecting it on a reorg must round-trip + /// An unknown bundle that a mainchain block confirms at height 0 spends + /// every output the state holds, including a deposit that the same + /// mainchain block created. + #[test] + fn a_bundle_spends_a_deposit_of_the_same_block() -> anyhow::Result<()> { + let (_temp_dir, env, state) = + fresh_state("a_bundle_spends_a_deposit_of_the_same_block")?; + let m6id = M6id(bitcoin::Txid::from_byte_array([5; 32])); + let deposit_outpoint = bitcoin::OutPoint { + txid: bitcoin::Txid::from_byte_array([7; 32]), + vout: 0, + }; + let output = value_output(Address::ALL_ZEROS, 5000); + let mut block_info = LinkedHashMap::new(); + block_info.insert( + bitcoin::BlockHash::from_byte_array([9; 32]), + BlockInfo { + bmm_commitment: None, + events: vec![ + BlockEvent::Deposit(Deposit { + tx_index: 0, + outpoint: deposit_outpoint, + output: output.clone(), + }), + BlockEvent::WithdrawalBundle(WithdrawalBundleEvent { + m6id, + status: WithdrawalBundleEventStatus::Confirmed, + }), + ], + }, + ); + let two_way_peg_data = TwoWayPegData { block_info }; + { + let mut rwtxn = env.write_txn()?; + state.height.put(&mut rwtxn, &(), &0)?; + state.withdrawal_bundles.put( + &mut rwtxn, + &m6id, + &( + WithdrawalBundleInfo::Unknown, + RollBack::new(WithdrawalBundleStatus::Submitted, 0), + ), + )?; + let () = connect(&state, &mut rwtxn, &two_way_peg_data)?; + rwtxn.commit()?; + } + let expected = vec![ + TwoWayPegEvent::Deposit { + outpoint: OutPoint::Deposit(deposit_outpoint), + output, + }, + TwoWayPegEvent::BundleSpend { + outpoint: OutPoint::Deposit(deposit_outpoint), + m6id, + }, + ]; + { + let rotxn = env.read_txn()?; + let events = state.get_two_way_peg_events(&rotxn, 0)?; + anyhow::ensure!( + events == expected, + "the events must keep the order the node applied: {events:?}" + ); + } + + // A disconnect drops the events, so the block that takes the height + // finds none of them. + let mut rwtxn = env.write_txn()?; + let () = disconnect(&state, &mut rwtxn, &two_way_peg_data)?; + anyhow::ensure!(state.get_two_way_peg_events(&rwtxn, 0)?.is_empty()); + Ok(()) + } + #[test] fn deposit_reorg_round_trips() -> anyhow::Result<()> { use crate::types::{ diff --git a/rpc-api/lib.rs b/rpc-api/lib.rs index 37a718e7..db6e3e36 100644 --- a/rpc-api/lib.rs +++ b/rpc-api/lib.rs @@ -33,7 +33,7 @@ pub mod node { PointedOutput, SpentOutput, Transaction, Txid, WithdrawalBundle, WithdrawalBundleStatus, net::{Peer, PeerAddress, PeerConnectionStatus}, - state::WithdrawalBundleInfo, + state::{TwoWayPegEvent, WithdrawalBundleInfo}, }; use typewit::const_marker::Bool; use utoipa::ToSchema; @@ -264,7 +264,7 @@ pub mod node { ref_schemas[ Address, Authorization, BlockHash, Body, Header, InPoint, M6id, MerkleRoot, OutPoint, Output, OutputContent, PeerConnectionStatus, - SpentOutput, Transaction, Txid, WithdrawalBundle, + SpentOutput, Transaction, TwoWayPegEvent, Txid, WithdrawalBundle, WithdrawalBundleInfo, WithdrawalBundleStatus, schema::BitcoinAddr, schema::BitcoinBlockHash, schema::BitcoinOutPoint, schema::BitcoinTransaction, schema::SocketAddr, @@ -351,6 +351,16 @@ pub mod node { txid: Txid, ) -> RpcResult>; + /// Get the coin movements that a block applied outside its body: a + /// mainchain deposit, a withdrawal bundle spend, and the outputs a + /// failed bundle returned. The list keeps the order the node applied. + #[open_api_method(output_schema(ToSchema))] + #[method(name = "get_two_way_peg_events")] + async fn get_two_way_peg_events( + &self, + block_hash: thunder_types::BlockHash, + ) -> RpcResult>; + /// Get utxos for addresses #[method(name = "get_utxos")] async fn get_utxos( diff --git a/types/state.rs b/types/state.rs index dfd0effe..aff5488e 100644 --- a/types/state.rs +++ b/types/state.rs @@ -3,7 +3,7 @@ use std::collections::BTreeMap; use serde::{Deserialize, Serialize}; use utoipa::ToSchema; -use crate::{OutPoint, Output, WithdrawalBundle}; +use crate::{M6id, OutPoint, Output, WithdrawalBundle}; /// Information we have regarding a withdrawal bundle #[derive(Clone, Debug, Deserialize, Serialize, ToSchema)] @@ -18,3 +18,18 @@ pub enum WithdrawalBundleInfo { spend_utxos: BTreeMap, }, } + +/// A coin movement that a block applied outside its body +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize, ToSchema)] +pub enum TwoWayPegEvent { + /// A mainchain deposit created this output + Deposit { outpoint: OutPoint, output: Output }, + /// A withdrawal bundle spent this output + BundleSpend { outpoint: OutPoint, m6id: M6id }, + /// A failed withdrawal bundle returned this output to the UTXO set + BundleReturn { + outpoint: OutPoint, + output: Output, + m6id: M6id, + }, +}