From e25c9df2cfcbf7556f52e0bb80219696f668df66 Mon Sep 17 00:00:00 2001 From: Reynaldo Gil Pons Date: Fri, 24 Jul 2026 13:55:18 +0200 Subject: [PATCH] feat: avoid using Option for supporters --- crates/node/src/coordinator.rs | 2 +- crates/node/src/foreign_chain_policy.rs | 128 +++++++----------- crates/node/src/indexer.rs | 7 +- crates/node/src/indexer/fake.rs | 14 +- crates/node/src/indexer/foreign_chain.rs | 53 +++++--- crates/node/src/indexer/real.rs | 21 ++- crates/node/src/metrics.rs | 3 +- .../node/src/providers/verify_foreign_tx.rs | 5 +- .../src/providers/verify_foreign_tx/sign.rs | 44 ++---- .../src/tests/verify_foreign_tx_gating.rs | 11 +- 10 files changed, 125 insertions(+), 163 deletions(-) diff --git a/crates/node/src/coordinator.rs b/crates/node/src/coordinator.rs index 332bb66ecd..31e0a7b4ec 100644 --- a/crates/node/src/coordinator.rs +++ b/crates/node/src/coordinator.rs @@ -370,7 +370,7 @@ where keyshare_storage: Arc>, running_state: ContractRunningState, chain_txn_sender: TransactionSender, - foreign_chain_supporters_receiver: watch::Receiver>, + foreign_chain_supporters_receiver: watch::Receiver, block_update_receiver: tokio::sync::OwnedMutexGuard< mpsc::UnboundedReceiver, >, diff --git a/crates/node/src/foreign_chain_policy.rs b/crates/node/src/foreign_chain_policy.rs index 8e1b59fe31..7a92b8f8e9 100644 --- a/crates/node/src/foreign_chain_policy.rs +++ b/crates/node/src/foreign_chain_policy.rs @@ -15,20 +15,13 @@ use crate::tracking::{self, AutoAbortTask}; pub(crate) type SupportersByForeignChain = BTreeMap>; /// Resolves the indexer's TLS-key supporters channel against the current -/// participant set, publishing each available chain's supporting participants -/// (filtered by the ForeignTx reconstruction threshold; a `None` threshold -/// means no ForeignTx domain, so no chain is available). The published value -/// is `None` until the upstream delivers its first snapshot. Must be called -/// from a tracked task; dropping the returned [`AutoAbortTask`] stops the -/// resolver, as does dropping the upstream sender or every receiver. +/// participant set. The upstream always holds a real value, so the returned +/// receiver does too. Must be called from a tracked task. pub(crate) fn spawn_supporters_by_foreign_chain( - mut upstream: watch::Receiver>, + mut upstream: watch::Receiver, participants_config: ParticipantsConfig, foreign_tx_reconstruction_threshold: Option, -) -> ( - watch::Receiver>, - AutoAbortTask<()>, -) { +) -> (watch::Receiver, AutoAbortTask<()>) { let init_value = resolve_supporters_by_foreign_chain( &upstream.borrow_and_update(), &participants_config, @@ -69,10 +62,10 @@ pub(crate) fn spawn_supporters_by_foreign_chain( } async fn await_updated_supporters( - upstream: &mut watch::Receiver>, + upstream: &mut watch::Receiver, participants_config: &ParticipantsConfig, foreign_tx_reconstruction_threshold: Option, -) -> anyhow::Result> { +) -> anyhow::Result { upstream .changed() .await @@ -100,27 +93,23 @@ pub(crate) fn foreign_tx_reconstruction_threshold(domains: &[dtos::DomainConfig] /// registrations (prospective or stale ones included) and may come from a /// different block than `participants_config`. Only supporters resolving to /// `participants_config` — the participants this node can sign with — count -/// towards the quorum. -/// Returns `None` while no upstream snapshot exists yet: "not known" must stay -/// distinguishable from "no chain available". +/// towards the quorum. An empty map means no chain is available (either no +/// ForeignTx domain, or no chain reaches the quorum). fn resolve_supporters_by_foreign_chain( - supporters_by_tls_key: &Option, + supporters_by_tls_key: &ForeignChainSupporters, participants_config: &ParticipantsConfig, foreign_tx_reconstruction_threshold: Option, -) -> Option { - let supporters_by_tls_key = supporters_by_tls_key.as_ref()?; +) -> SupportersByForeignChain { let Some(threshold) = foreign_tx_reconstruction_threshold else { - return Some(SupportersByForeignChain::new()); + return SupportersByForeignChain::new(); }; - Some( - supporters_by_tls_key - .iter() - .filter_map(|(chain, tls_keys)| { - let ids = resolve_participant_ids(tls_keys, participants_config); - (ids.len() as u64 >= threshold).then_some((*chain, ids)) - }) - .collect(), - ) + supporters_by_tls_key + .iter() + .filter_map(|(chain, tls_keys)| { + let ids = resolve_participant_ids(tls_keys, participants_config); + (ids.len() as u64 >= threshold).then_some((*chain, ids)) + }) + .collect() } /// Resolves TLS keys to the matching participants' ids; keys not belonging to a @@ -198,10 +187,10 @@ mod tests { make_participant_info(3, &keys[2]), ], ); - let supporters_by_tls_key = Some(BTreeMap::from([( + let supporters_by_tls_key = BTreeMap::from([( dtos::ForeignChain::Bitcoin, keys.iter().map(tls_key_for).collect::>(), - )])); + )]); // When let supporters = resolve_supporters_by_foreign_chain( @@ -213,14 +202,14 @@ mod tests { // Then assert_eq!( supporters, - Some(BTreeMap::from([( + BTreeMap::from([( dtos::ForeignChain::Bitcoin, HashSet::from([ ParticipantId::from_raw(1), ParticipantId::from_raw(2), ParticipantId::from_raw(3), ]), - )])) + )]) ); } @@ -236,10 +225,10 @@ mod tests { make_participant_info(2, &key2), ], ); - let supporters_by_tls_key = Some(BTreeMap::from([( + let supporters_by_tls_key = BTreeMap::from([( dtos::ForeignChain::Bitcoin, BTreeSet::from([tls_key_for(&key1)]), - )])); + )]); // When let supporters = resolve_supporters_by_foreign_chain( @@ -249,7 +238,7 @@ mod tests { ); // Then - assert_eq!(supporters, Some(SupportersByForeignChain::new())); + assert_eq!(supporters, SupportersByForeignChain::new()); } #[test] @@ -257,10 +246,10 @@ mod tests { // Given: a participants threshold far above the single supporter. let key1 = make_signing_key(1); let participants_config = participants(100, vec![make_participant_info(1, &key1)]); - let supporters_by_tls_key = Some(BTreeMap::from([( + let supporters_by_tls_key = BTreeMap::from([( dtos::ForeignChain::Bitcoin, BTreeSet::from([tls_key_for(&key1)]), - )])); + )]); // When let supporters = resolve_supporters_by_foreign_chain( @@ -270,11 +259,7 @@ mod tests { ); // Then: only the ForeignTx domain threshold applies. - assert!( - supporters - .expect("snapshot present") - .contains_key(&dtos::ForeignChain::Bitcoin) - ); + assert!(supporters.contains_key(&dtos::ForeignChain::Bitcoin)); } #[test] @@ -282,30 +267,17 @@ mod tests { // Given: a supported chain but no ForeignTx domain (no threshold). let key1 = make_signing_key(1); let participants_config = participants(1, vec![make_participant_info(1, &key1)]); - let supporters_by_tls_key = Some(BTreeMap::from([( + let supporters_by_tls_key = BTreeMap::from([( dtos::ForeignChain::Bitcoin, BTreeSet::from([tls_key_for(&key1)]), - )])); + )]); // When let supporters = resolve_supporters_by_foreign_chain(&supporters_by_tls_key, &participants_config, None); // Then - assert_eq!(supporters, Some(SupportersByForeignChain::new())); - } - - #[test] - fn resolve_supporters_by_foreign_chain__should_return_none_before_first_upstream_snapshot() { - // Given - let key1 = make_signing_key(1); - let participants_config = participants(1, vec![make_participant_info(1, &key1)]); - - // When - let supporters = resolve_supporters_by_foreign_chain(&None, &participants_config, Some(1)); - - // Then - assert_eq!(supporters, None); + assert_eq!(supporters, SupportersByForeignChain::new()); } #[test] @@ -348,36 +320,32 @@ mod tests { } #[tokio::test] - async fn spawn_supporters_by_foreign_chain__should_publish_resolved_map_on_upstream_change() { + async fn spawn_supporters_by_foreign_chain__should_republish_on_upstream_change() { let (root, _root_handle) = start_root_task("test-root", async move { - // Given: the upstream has not delivered a snapshot yet. + // Given: Bitcoin resolved as available from the first snapshot. let key1 = make_signing_key(1); let participants_config = participants(1, vec![make_participant_info(1, &key1)]); - let (upstream_sender, upstream_receiver) = watch::channel(None); + let (upstream_sender, upstream_receiver) = watch::channel(BTreeMap::from([( + dtos::ForeignChain::Bitcoin, + BTreeSet::from([tls_key_for(&key1)]), + )])); let (mut supporters, _resolver_task) = spawn_supporters_by_foreign_chain(upstream_receiver, participants_config, Some(1)); - assert!(supporters.borrow().is_none()); - - // When: Bitcoin becomes available with the participant registered for it. - upstream_sender - .send(Some(BTreeMap::from([( - dtos::ForeignChain::Bitcoin, - BTreeSet::from([tls_key_for(&key1)]), - )]))) - .unwrap(); + assert!( + supporters + .borrow() + .contains_key(&dtos::ForeignChain::Bitcoin) + ); + + // When: the chain loses its registration upstream. + upstream_sender.send(BTreeMap::new()).unwrap(); tokio::time::timeout(Duration::from_secs(5), supporters.changed()) .await .unwrap() .unwrap(); // Then - assert_eq!( - *supporters.borrow(), - Some(BTreeMap::from([( - dtos::ForeignChain::Bitcoin, - HashSet::from([ParticipantId::from_raw(1)]), - )])) - ); + assert_eq!(*supporters.borrow(), SupportersByForeignChain::new()); }); root.await; } @@ -402,12 +370,12 @@ mod tests { )]); // When - let (_upstream_sender, upstream_receiver) = watch::channel(Some(upstream)); + let (_upstream_sender, upstream_receiver) = watch::channel(upstream); let (supporters, _resolver_task) = spawn_supporters_by_foreign_chain(upstream_receiver, participants_config, Some(2)); // Then: the stranger's key does not count towards the quorum. - assert_eq!(*supporters.borrow(), Some(SupportersByForeignChain::new())); + assert_eq!(*supporters.borrow(), SupportersByForeignChain::new()); }); root.await; } diff --git a/crates/node/src/indexer.rs b/crates/node/src/indexer.rs index 03d1097855..a10384f6f6 100644 --- a/crates/node/src/indexer.rs +++ b/crates/node/src/indexer.rs @@ -584,10 +584,9 @@ pub struct IndexerAPI { pub my_migration_info_receiver: watch::Receiver, /// Watcher that tracks the contract's available foreign chains and their - /// registered supporters (by TLS key). Holds `None` until the first - /// successful read after the indexer syncs. - pub foreign_chain_supporters_receiver: - watch::Receiver>, + /// registered supporters (by TLS key). Seeded with the first successful read + /// before the indexer hands it back, so it always holds a real value. + pub foreign_chain_supporters_receiver: watch::Receiver, pub(crate) attestation_reader: std::sync::Arc, } diff --git a/crates/node/src/indexer/fake.rs b/crates/node/src/indexer/fake.rs index 970774efab..b7af70eaed 100644 --- a/crates/node/src/indexer/fake.rs +++ b/crates/node/src/indexer/fake.rs @@ -545,7 +545,7 @@ struct FakeIndexerCore { /// Broadcasts the contract state to each node. migration_change_sender: broadcast::Sender, /// Mirrors the real indexer's foreign-chain supporters watch channel. - foreign_chain_supporters_sender: watch::Sender>, + foreign_chain_supporters_sender: watch::Sender, /// When the core receives signature response txns, it processes them by sending them through /// this sender. The receiver end of this is in FakeIndexManager to be received by the test @@ -599,10 +599,10 @@ impl FakeIndexerCore { state.foreign_chains_configs(), ); foreign_chain_supporters_sender.send_if_modified(|previous| { - if previous.as_ref() == Some(&supporters) { + if *previous == supporters { false } else { - *previous = Some(supporters); + *previous = supporters; true } }); @@ -881,7 +881,7 @@ pub struct FakeIndexerManager { /// Cloned into each node's `IndexerAPI`; tracks the fake contract's /// foreign-chain supporters. - foreign_chain_supporters_receiver: watch::Receiver>, + foreign_chain_supporters_receiver: watch::Receiver, account_id_by_uid: Arc>>, } @@ -1071,7 +1071,7 @@ impl FakeIndexerManager { let (verify_foreign_tx_response_sender, verify_foreign_tx_response_receiver) = mpsc::unbounded_channel(); let (foreign_chain_supporters_sender, foreign_chain_supporters_receiver) = - watch::channel(None); + watch::channel(ForeignChainSupporters::new()); let contract = Arc::new(tokio::sync::Mutex::new(FakeMpcContractState::new())); let account_id_by_uid = Arc::new(std::sync::Mutex::new(HashMap::new())); let core = FakeIndexerCore { @@ -1134,9 +1134,7 @@ impl FakeIndexerManager { } /// The supporters channel every node's `IndexerAPI` receives. - pub fn subscribe_foreign_chain_supporters( - &self, - ) -> watch::Receiver> { + pub fn subscribe_foreign_chain_supporters(&self) -> watch::Receiver { self.foreign_chain_supporters_receiver.clone() } diff --git a/crates/node/src/indexer/foreign_chain.rs b/crates/node/src/indexer/foreign_chain.rs index fb6666876f..1598ea1f26 100644 --- a/crates/node/src/indexer/foreign_chain.rs +++ b/crates/node/src/indexer/foreign_chain.rs @@ -12,35 +12,46 @@ const FOREIGN_CHAIN_SUPPORTERS_REFRESH_INTERVAL: Duration = Duration::from_secs( /// TLS keys of the nodes whose registered config supports each available chain. pub type ForeignChainSupporters = BTreeMap>; -/// Updates the contract's available chains mapped to their registered -/// supporters in watch channel. -/// The channel holds `None` until the first successful read; afterwards the -/// previously published value stays in effect until viewing new state from -/// contract succeeds. +/// Returns once the first supporters snapshot is read, then keeps it updated in +/// the background. Mirrors `monitor_contract_state`: the receiver always holds a +/// real value, and a failed refresh keeps the previous one. pub async fn monitor_foreign_chain_supporters( - sender: watch::Sender>, indexer_state: Arc, -) { +) -> watch::Receiver { indexer_state.client.wait_for_full_sync().await; - loop { + let initial = loop { match read_supporters(&indexer_state).await { - Ok(supporters) => { - sender.send_if_modified(|previous| { - if previous.as_ref() == Some(&supporters) { - false - } else { - *previous = Some(supporters); - true - } - }); - } + Ok(supporters) => break supporters, Err(e) => { - tracing::error!(target: "mpc", "error reading foreign-chain supporters from chain: {:?}", e) + tracing::error!(target: "mpc", "error reading foreign-chain supporters from chain: {:?}", e); + tokio::time::sleep(FOREIGN_CHAIN_SUPPORTERS_REFRESH_INTERVAL).await; } } - tokio::time::sleep(FOREIGN_CHAIN_SUPPORTERS_REFRESH_INTERVAL).await; - } + }; + + let (sender, receiver) = watch::channel(initial); + tokio::spawn(async move { + loop { + tokio::time::sleep(FOREIGN_CHAIN_SUPPORTERS_REFRESH_INTERVAL).await; + match read_supporters(&indexer_state).await { + Ok(supporters) => { + sender.send_if_modified(|previous| { + if *previous == supporters { + false + } else { + *previous = supporters; + true + } + }); + } + Err(e) => { + tracing::error!(target: "mpc", "error reading foreign-chain supporters from chain: {:?}", e) + } + } + } + }); + receiver } /// The two view calls are not atomic: a change finalized between them yields diff --git a/crates/node/src/indexer/real.rs b/crates/node/src/indexer/real.rs index e0fc91866e..c94b7688bc 100644 --- a/crates/node/src/indexer/real.rs +++ b/crates/node/src/indexer/real.rs @@ -72,6 +72,8 @@ pub fn spawn_real_indexer( ) -> IndexerAPI { let (contract_state_sender_oneshot, contract_state_receiver_oneshot) = oneshot::channel(); let (migration_info_sender_oneshot, migration_info_receiver_oneshot) = oneshot::channel(); + let (foreign_chain_supporters_sender_oneshot, foreign_chain_supporters_receiver_oneshot) = + oneshot::channel(); let (attestation_reader_sender, attestation_reader_receiver) = oneshot::channel(); let (block_update_sender, block_update_receiver) = mpsc::unbounded_channel(); @@ -79,7 +81,6 @@ pub fn spawn_real_indexer( let (allowed_launcher_compose_sender, allowed_launcher_compose_receiver) = watch::channel(vec![]); let (tee_accounts_sender, tee_accounts_receiver) = watch::channel(vec![]); - let (foreign_chain_supporters_sender, foreign_chain_supporters_receiver) = watch::channel(None); let my_near_account_id_clone = my_near_account_id.clone(); let respond_config_clone = respond_config.clone(); @@ -212,10 +213,16 @@ pub fn spawn_real_indexer( indexer_state.clone(), )); - tokio::spawn(monitor_foreign_chain_supporters( - foreign_chain_supporters_sender, - indexer_state.clone(), - )); + let foreign_chain_supporters_receiver = + monitor_foreign_chain_supporters(indexer_state.clone()).await; + if foreign_chain_supporters_sender_oneshot + .send(foreign_chain_supporters_receiver) + .is_err() + { + tracing::error!( + "Indexer thread could not send foreign chain supporters receiver back to main driver." + ) + }; let (foreign_chain_whitelist_sender, foreign_chain_whitelist_receiver) = watch::channel(std::collections::BTreeMap::new()); @@ -321,6 +328,10 @@ pub fn spawn_real_indexer( .blocking_recv() .expect("Migraration info receiver must be returned by indexer."); + let foreign_chain_supporters_receiver = foreign_chain_supporters_receiver_oneshot + .blocking_recv() + .expect("foreign chain supporters receiver must be returned by indexer"); + let attestation_reader = attestation_reader_receiver .blocking_recv() .expect("attestation reader must be returned by indexer"); diff --git a/crates/node/src/metrics.rs b/crates/node/src/metrics.rs index 373689ac13..0578c69810 100644 --- a/crates/node/src/metrics.rs +++ b/crates/node/src/metrics.rs @@ -157,8 +157,7 @@ pub static MPC_NUM_VERIFY_FOREIGN_TX_UNAVAILABLE_CHAIN_REJECTIONS: LazyLock< prometheus::register_int_counter!( "mpc_num_verify_foreign_tx_unavailable_chain_rejections", "Number of gate rejections of verify foreign tx attempts, at most one per node per \ - attempt: the requested chain is not available or the supporters snapshot has not \ - been received yet" + attempt: the requested chain is not available" ) .unwrap() }); diff --git a/crates/node/src/providers/verify_foreign_tx.rs b/crates/node/src/providers/verify_foreign_tx.rs index 9baa70e9a2..8b72f36b3a 100644 --- a/crates/node/src/providers/verify_foreign_tx.rs +++ b/crates/node/src/providers/verify_foreign_tx.rs @@ -148,8 +148,7 @@ impl ForeignChainInspectors { pub struct VerifyForeignTxProvider { config: Arc, inspectors: ForeignChainInspectors, - /// `None` until the indexer delivers its first supporters snapshot. - supporters_by_foreign_chain: watch::Receiver>, + supporters_by_foreign_chain: watch::Receiver, verify_foreign_tx_request_store: Arc, ecdsa_signature_provider: Arc, } @@ -171,7 +170,7 @@ impl From for MpcTaskId { impl VerifyForeignTxProvider { pub fn new( config: Arc, - supporters_by_foreign_chain: watch::Receiver>, + supporters_by_foreign_chain: watch::Receiver, verify_foreign_tx_request_store: Arc, ecdsa_signature_provider: Arc, ) -> anyhow::Result { diff --git a/crates/node/src/providers/verify_foreign_tx/sign.rs b/crates/node/src/providers/verify_foreign_tx/sign.rs index 5004b1539d..5b8d21032e 100644 --- a/crates/node/src/providers/verify_foreign_tx/sign.rs +++ b/crates/node/src/providers/verify_foreign_tx/sign.rs @@ -382,31 +382,25 @@ impl VerifyForeignTxProvider { } #[derive(Debug, thiserror::Error)] -enum ChainAvailabilityError { - #[error("the foreign-chain supporters snapshot has not been received from the contract yet")] - SupportersSnapshotNotReady, - #[error( - "requested chain {requested:?} is not in the list of available foreign chains on the MPC contract" - )] - ChainNotAvailable { requested: dtos::ForeignChain }, +#[error( + "requested chain {requested:?} is not in the list of available foreign chains on the MPC contract" +)] +struct ChainNotAvailableError { + requested: dtos::ForeignChain, } /// A chain counts as available when the supporters map has an entry for it: /// the chain is available on the contract and a signing quorum of current -/// participants supports it. A missing snapshot (`None`) rejects every chain, -/// but distinguishably from a genuinely unavailable one. +/// participants supports it. fn ensure_chain_is_available( - supporters_by_foreign_chain: &Option, + supporters_by_foreign_chain: &SupportersByForeignChain, request: &dtos::ForeignChainRpcRequest, -) -> Result<(), ChainAvailabilityError> { - let Some(supporters_by_foreign_chain) = supporters_by_foreign_chain else { - return Err(ChainAvailabilityError::SupportersSnapshotNotReady); - }; +) -> Result<(), ChainNotAvailableError> { let requested = request.chain(); if supporters_by_foreign_chain.contains_key(&requested) { Ok(()) } else { - Err(ChainAvailabilityError::ChainNotAvailable { requested }) + Err(ChainNotAvailableError { requested }) } } @@ -418,11 +412,11 @@ mod tests { use assert_matches::assert_matches; use std::collections::{BTreeMap, HashSet}; - fn bitcoin_supporters() -> Option { - Some(BTreeMap::from([( + fn bitcoin_supporters() -> SupportersByForeignChain { + BTreeMap::from([( dtos::ForeignChain::Bitcoin, HashSet::from([ParticipantId::from_raw(1)]), - )])) + )]) } fn bitcoin_request() -> dtos::ForeignChainRpcRequest { @@ -458,21 +452,9 @@ mod tests { // When, then assert_matches!( ensure_chain_is_available(&supporters, ðereum_request), - Err(ChainAvailabilityError::ChainNotAvailable { + Err(ChainNotAvailableError { requested: dtos::ForeignChain::Ethereum }) ); } - - #[test] - fn ensure_chain_is_available__should_fail_when_snapshot_not_received_yet() { - // Given: no supporters snapshot from the indexer yet. - let supporters = None; - - // When, then - assert_matches!( - ensure_chain_is_available(&supporters, &bitcoin_request()), - Err(ChainAvailabilityError::SupportersSnapshotNotReady) - ); - } } diff --git a/crates/node/src/tests/verify_foreign_tx_gating.rs b/crates/node/src/tests/verify_foreign_tx_gating.rs index 9f2b607b0d..9e23dd7a00 100644 --- a/crates/node/src/tests/verify_foreign_tx_gating.rs +++ b/crates/node/src/tests/verify_foreign_tx_gating.rs @@ -167,21 +167,16 @@ async fn verify_foreign_tx__should_only_be_served_while_chain_is_available() { assert!(contract.available_foreign_chains().is_empty()); } // Wait for the fake core to publish the post-change snapshot on the shared - // upstream channel; the per-node resolver fan-out from it is in-process, - // covered by the extra second. + // upstream channel; the per-node resolver fan-out from it is in-process and + // subsumed by the response wait below. let mut supporters = setup.indexer.subscribe_foreign_chain_supporters(); tokio::time::timeout(SUPPORTERS_PUBLISH_WAIT, async { - while !supporters - .borrow_and_update() - .as_ref() - .is_some_and(|map| map.is_empty()) - { + while !supporters.borrow_and_update().is_empty() { supporters.changed().await.unwrap(); } }) .await .expect("timed out waiting for the empty supporters snapshot to publish"); - tokio::time::sleep(Duration::from_secs(1)).await; // Then: nodes reject the request against their supporters snapshot, // even though inspection itself would succeed. The distinct tx id keeps