use bitcoin::amount::Amount;
use bitcoin::constants::ChainHash;
use bitcoin::TxOut;
use bitcoin::hex::DisplayHex;
use crate::ln::chan_utils::make_funding_redeemscript_from_slices;
use crate::ln::msgs::{self, ErrorAction, LightningError, MessageSendEvent};
use crate::routing::gossip::{NetworkGraph, NodeId};
use crate::util::logger::{Level, Logger};
use crate::util::wakers::Notifier;
use crate::prelude::*;
use crate::sync::{LockTestExt, Mutex};
use alloc::sync::{Arc, Weak};
use core::ops::Deref;
#[derive(Clone, Debug)]
pub enum UtxoLookupError {
UnknownChain,
UnknownTx,
}
#[derive(Clone)]
pub enum UtxoResult {
Sync(Result<TxOut, UtxoLookupError>),
Async(UtxoFuture),
}
pub trait UtxoLookup {
fn get_utxo(
&self, chain_hash: &ChainHash, short_channel_id: u64,
async_completion_notifier: Arc<Notifier>,
) -> UtxoResult;
}
impl<T: UtxoLookup + ?Sized, U: Deref<Target = T>> UtxoLookup for U {
fn get_utxo(
&self, chain_hash: &ChainHash, short_channel_id: u64,
async_completion_notifier: Arc<Notifier>,
) -> UtxoResult {
self.deref().get_utxo(chain_hash, short_channel_id, async_completion_notifier)
}
}
enum ChannelAnnouncement {
Full(msgs::ChannelAnnouncement),
Unsigned(msgs::UnsignedChannelAnnouncement),
}
impl ChannelAnnouncement {
fn node_id_1(&self) -> &NodeId {
match self {
ChannelAnnouncement::Full(msg) => &msg.contents.node_id_1,
ChannelAnnouncement::Unsigned(msg) => &msg.node_id_1,
}
}
}
enum NodeAnnouncement {
Full(msgs::NodeAnnouncement),
Unsigned(msgs::UnsignedNodeAnnouncement),
}
impl NodeAnnouncement {
fn timestamp(&self) -> u32 {
match self {
NodeAnnouncement::Full(msg) => msg.contents.timestamp,
NodeAnnouncement::Unsigned(msg) => msg.timestamp,
}
}
}
enum ChannelUpdate {
Full(msgs::ChannelUpdate),
Unsigned(msgs::UnsignedChannelUpdate),
}
impl ChannelUpdate {
fn timestamp(&self) -> u32 {
match self {
ChannelUpdate::Full(msg) => msg.contents.timestamp,
ChannelUpdate::Unsigned(msg) => msg.timestamp,
}
}
}
struct UtxoMessages {
notifier: Arc<Notifier>,
complete: Option<Result<TxOut, UtxoLookupError>>,
channel_announce: Option<ChannelAnnouncement>,
latest_node_announce_a: Option<NodeAnnouncement>,
latest_node_announce_b: Option<NodeAnnouncement>,
latest_channel_update_a: Option<ChannelUpdate>,
latest_channel_update_b: Option<ChannelUpdate>,
}
#[derive(Clone)]
pub struct UtxoFuture {
state: Arc<Mutex<UtxoMessages>>,
}
pub(crate) struct UtxoResolver(Result<TxOut, UtxoLookupError>);
impl UtxoLookup for UtxoResolver {
fn get_utxo(&self, _hash: &ChainHash, _scid: u64, _notifier: Arc<Notifier>) -> UtxoResult {
UtxoResult::Sync(self.0.clone())
}
}
impl UtxoFuture {
pub fn new(notifier: Arc<Notifier>) -> Self {
Self {
state: Arc::new(Mutex::new(UtxoMessages {
notifier,
complete: None,
channel_announce: None,
latest_node_announce_a: None,
latest_node_announce_b: None,
latest_channel_update_a: None,
latest_channel_update_b: None,
})),
}
}
pub fn resolve(&self, result: Result<TxOut, UtxoLookupError>) {
let mut state = self.state.lock().unwrap();
state.complete = Some(result);
state.notifier.notify();
}
}
struct PendingChecksContext {
pending_states: Vec<Arc<Mutex<UtxoMessages>>>,
channels: HashMap<u64, Weak<Mutex<UtxoMessages>>>,
nodes: HashMap<NodeId, Vec<Weak<Mutex<UtxoMessages>>>>,
}
pub(super) struct PendingChecks {
internal: Mutex<PendingChecksContext>,
pub(super) completion_notifier: Arc<Notifier>,
}
impl PendingChecks {
pub(super) fn new() -> Self {
PendingChecks {
internal: Mutex::new(PendingChecksContext {
pending_states: Vec::new(),
channels: new_hash_map(),
nodes: new_hash_map(),
}),
completion_notifier: Arc::new(Notifier::new()),
}
}
pub(super) fn check_hold_pending_channel_update(
&self, msg: &msgs::UnsignedChannelUpdate, full_msg: Option<&msgs::ChannelUpdate>,
) -> Result<(), LightningError> {
let mut pending_checks = self.internal.lock().unwrap();
if let hash_map::Entry::Occupied(e) = pending_checks.channels.entry(msg.short_channel_id) {
let is_from_a = (msg.channel_flags & 1) == 1;
match Weak::upgrade(e.get()) {
Some(msgs_ref) => {
let mut messages = msgs_ref.lock().unwrap();
let latest_update = if is_from_a {
&mut messages.latest_channel_update_a
} else {
&mut messages.latest_channel_update_b
};
if latest_update.is_none()
|| latest_update.as_ref().unwrap().timestamp() < msg.timestamp
{
*latest_update = Some(if let Some(msg) = full_msg {
ChannelUpdate::Full(msg.clone())
} else {
ChannelUpdate::Unsigned(msg.clone())
});
}
return Err(LightningError {
err: "Awaiting channel_announcement validation to accept channel_update"
.to_owned(),
action: ErrorAction::IgnoreAndLog(Level::Gossip),
});
},
None => {
e.remove();
},
}
}
Ok(())
}
pub(super) fn check_hold_pending_node_announcement(
&self, msg: &msgs::UnsignedNodeAnnouncement, full_msg: Option<&msgs::NodeAnnouncement>,
) -> Result<(), LightningError> {
let mut pending_checks = self.internal.lock().unwrap();
if let hash_map::Entry::Occupied(mut e) = pending_checks.nodes.entry(msg.node_id) {
let mut found_at_least_one_chan = false;
e.get_mut().retain(|node_msgs| match Weak::upgrade(&node_msgs) {
Some(chan_mtx) => {
let mut chan_msgs = chan_mtx.lock().unwrap();
if let Some(chan_announce) = &chan_msgs.channel_announce {
let latest_announce = if *chan_announce.node_id_1() == msg.node_id {
&mut chan_msgs.latest_node_announce_a
} else {
&mut chan_msgs.latest_node_announce_b
};
if latest_announce.is_none()
|| latest_announce.as_ref().unwrap().timestamp() < msg.timestamp
{
*latest_announce = Some(if let Some(msg) = full_msg {
NodeAnnouncement::Full(msg.clone())
} else {
NodeAnnouncement::Unsigned(msg.clone())
});
}
found_at_least_one_chan = true;
true
} else {
debug_assert!(
false,
"channel_announce is set before struct is added to node map"
);
false
}
},
None => false,
});
if e.get().is_empty() {
e.remove();
}
if found_at_least_one_chan {
return Err(LightningError {
err: "Awaiting channel_announcement validation to accept node_announcement"
.to_owned(),
action: ErrorAction::IgnoreAndLog(Level::Gossip),
});
}
}
Ok(())
}
fn check_replace_previous_entry(
msg: &msgs::UnsignedChannelAnnouncement, full_msg: Option<&msgs::ChannelAnnouncement>,
replacement: Option<Weak<Mutex<UtxoMessages>>>,
pending_channels: &mut HashMap<u64, Weak<Mutex<UtxoMessages>>>,
) -> Result<(), msgs::LightningError> {
match pending_channels.entry(msg.short_channel_id) {
hash_map::Entry::Occupied(mut e) => {
match Weak::upgrade(&e.get()) {
Some(pending_msgs) => {
let pending_state = pending_msgs.unsafe_well_ordered_double_lock_self();
let pending_matches = match &pending_state.channel_announce {
Some(ChannelAnnouncement::Full(pending_msg)) => {
Some(pending_msg) == full_msg
},
Some(ChannelAnnouncement::Unsigned(pending_msg)) => pending_msg == msg,
None => {
debug_assert!(
pending_state.complete.is_none(),
"channel_announce is None but complete is still pending"
);
false
},
};
drop(pending_state);
if pending_matches {
return Err(LightningError {
err: "Channel announcement is already being checked".to_owned(),
action: ErrorAction::IgnoreDuplicateGossip,
});
} else {
if let Some(item) = replacement {
*e.get_mut() = item;
}
}
},
None => {
if let Some(item) = replacement {
*e.get_mut() = item;
} else {
e.remove();
}
},
}
},
hash_map::Entry::Vacant(v) => {
if let Some(item) = replacement {
v.insert(item);
}
},
}
Ok(())
}
pub(super) fn check_channel_announcement<U: UtxoLookup>(
&self, utxo_lookup: &Option<U>, msg: &msgs::UnsignedChannelAnnouncement,
full_msg: Option<&msgs::ChannelAnnouncement>,
) -> Result<Option<Amount>, msgs::LightningError> {
let handle_result = |res| match res {
Ok(TxOut { value, script_pubkey }) => {
let expected_script = make_funding_redeemscript_from_slices(
msg.bitcoin_key_1.as_array(),
msg.bitcoin_key_2.as_array(),
)
.to_p2wsh();
if script_pubkey != expected_script {
return Err(LightningError {
err: format!(
"Channel announcement key ({}) didn't match on-chain script ({})",
expected_script.to_hex_string(),
script_pubkey.to_hex_string()
),
action: ErrorAction::IgnoreError,
});
}
Ok(Some(value))
},
Err(UtxoLookupError::UnknownChain) => Err(LightningError {
err: format!(
"Channel announced on an unknown chain ({})",
msg.chain_hash.to_bytes().as_hex()
),
action: ErrorAction::IgnoreError,
}),
Err(UtxoLookupError::UnknownTx) => Err(LightningError {
err: "Channel announced without corresponding UTXO entry".to_owned(),
action: ErrorAction::IgnoreError,
}),
};
Self::check_replace_previous_entry(
msg,
full_msg,
None,
&mut self.internal.lock().unwrap().channels,
)?;
match utxo_lookup {
&None => {
Ok(None)
},
&Some(ref utxo_lookup) => {
let notifier = Arc::clone(&self.completion_notifier);
match utxo_lookup.get_utxo(&msg.chain_hash, msg.short_channel_id, notifier) {
UtxoResult::Sync(res) => handle_result(res),
UtxoResult::Async(future) => {
let mut pending_checks = self.internal.lock().unwrap();
let mut async_messages = future.state.lock().unwrap();
if let Some(res) = async_messages.complete.take() {
handle_result(res)
} else {
let pending_states = &mut pending_checks.pending_states;
if pending_states
.iter()
.find(|s| Arc::ptr_eq(s, &future.state))
.is_none()
{
pending_states.push(Arc::clone(&future.state));
}
Self::check_replace_previous_entry(
msg,
full_msg,
Some(Arc::downgrade(&future.state)),
&mut pending_checks.channels,
)?;
async_messages.channel_announce = Some(if let Some(msg) = full_msg {
ChannelAnnouncement::Full(msg.clone())
} else {
ChannelAnnouncement::Unsigned(msg.clone())
});
pending_checks
.nodes
.entry(msg.node_id_1)
.or_default()
.push(Arc::downgrade(&future.state));
pending_checks
.nodes
.entry(msg.node_id_2)
.or_default()
.push(Arc::downgrade(&future.state));
Err(LightningError {
err: "Channel being checked async".to_owned(),
action: ErrorAction::IgnoreAndLog(Level::Gossip),
})
}
},
}
},
}
}
const MAX_PENDING_LOOKUPS: usize = 32;
pub(super) fn too_many_checks_pending(&self) -> bool {
let mut pending_checks = self.internal.lock().unwrap();
if pending_checks.channels.len() > Self::MAX_PENDING_LOOKUPS {
pending_checks.channels.retain(|_, chan| Weak::upgrade(&chan).is_some());
pending_checks.nodes.retain(|_, channels| {
channels.retain(|chan| Weak::upgrade(&chan).is_some());
!channels.is_empty()
});
pending_checks.channels.len() > Self::MAX_PENDING_LOOKUPS
} else {
false
}
}
fn resolve_single_future<L: Logger>(
&self, graph: &NetworkGraph<L>, entry: Arc<Mutex<UtxoMessages>>,
new_messages: &mut Vec<MessageSendEvent>,
) {
let (announcement, result, announce_a, announce_b, update_a, update_b);
{
let mut state = entry.lock().unwrap();
announcement = if let Some(announcement) = state.channel_announce.take() {
announcement
} else {
return;
};
result = if let Some(result) = state.complete.take() {
result
} else {
debug_assert!(false, "Future should have been resolved");
return;
};
announce_a = state.latest_node_announce_a.take();
announce_b = state.latest_node_announce_b.take();
update_a = state.latest_channel_update_a.take();
update_b = state.latest_channel_update_b.take();
}
let resolver = UtxoResolver(result);
let (node_id_1, node_id_2) = match &announcement {
ChannelAnnouncement::Full(signed_msg) => {
(signed_msg.contents.node_id_1, signed_msg.contents.node_id_2)
},
ChannelAnnouncement::Unsigned(msg) => (msg.node_id_1, msg.node_id_2),
};
match announcement {
ChannelAnnouncement::Full(signed_msg) => {
if graph.update_channel_from_announcement(&signed_msg, &Some(&resolver)).is_ok() {
new_messages.push(MessageSendEvent::BroadcastChannelAnnouncement {
msg: signed_msg,
update_msg: None,
});
}
},
ChannelAnnouncement::Unsigned(msg) => {
let _ = graph.update_channel_from_unsigned_announcement(&msg, &Some(&resolver));
},
}
for announce in [announce_a, announce_b] {
match announce {
Some(NodeAnnouncement::Full(signed_msg)) => {
if graph.update_node_from_announcement(&signed_msg).is_ok() {
new_messages
.push(MessageSendEvent::BroadcastNodeAnnouncement { msg: signed_msg });
}
},
Some(NodeAnnouncement::Unsigned(msg)) => {
let _ = graph.update_node_from_unsigned_announcement(&msg);
},
None => {},
}
}
for update in [update_a, update_b] {
match update {
Some(ChannelUpdate::Full(signed_msg)) => {
if graph.update_channel(&signed_msg).is_ok() {
new_messages.push(MessageSendEvent::BroadcastChannelUpdate {
msg: signed_msg,
node_id_1,
node_id_2,
});
}
},
Some(ChannelUpdate::Unsigned(msg)) => {
let _ = graph.update_channel_unsigned(&msg);
},
None => {},
}
}
}
pub(super) fn check_resolved_futures<L: Logger>(
&self, graph: &NetworkGraph<L>,
) -> Vec<MessageSendEvent> {
let mut completed_states = Vec::new();
{
let mut lck = self.internal.lock().unwrap();
lck.pending_states.retain(|state| {
if state.lock().unwrap().complete.is_some() {
completed_states.push(Arc::clone(&state));
false
} else {
if Arc::strong_count(state) == 1 {
false
} else {
true
}
}
});
lck.channels.retain(|_, state| {
if let Some(state) = state.upgrade() {
if state.lock().unwrap().complete.is_some() {
completed_states.push(state);
false
} else {
true
}
} else {
false
}
});
lck.nodes.retain(|_, lookups| {
lookups.retain(|state| {
if let Some(state) = state.upgrade() {
if state.lock().unwrap().complete.is_some() {
completed_states.push(state);
false
} else {
true
}
} else {
false
}
});
!lookups.is_empty()
});
}
let mut res = Vec::with_capacity(completed_states.len() * 5);
for state in completed_states {
self.resolve_single_future(graph, state, &mut res);
}
res
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::routing::gossip::tests::*;
use crate::util::test_utils::{TestChainSource, TestLogger};
use bitcoin::amount::Amount;
use bitcoin::secp256k1::{Secp256k1, SecretKey};
use core::sync::atomic::Ordering;
fn get_network() -> (TestChainSource, NetworkGraph<Box<TestLogger>>) {
let logger = Box::new(TestLogger::new());
let chain_source = TestChainSource::new(bitcoin::Network::Testnet);
let network_graph = NetworkGraph::new(bitcoin::Network::Testnet, logger);
(chain_source, network_graph)
}
fn get_test_objects() -> (
msgs::ChannelAnnouncement,
TestChainSource,
NetworkGraph<Box<TestLogger>>,
bitcoin::ScriptBuf,
msgs::NodeAnnouncement,
msgs::NodeAnnouncement,
msgs::ChannelUpdate,
msgs::ChannelUpdate,
msgs::ChannelUpdate,
) {
let secp_ctx = Secp256k1::new();
let (chain_source, network_graph) = get_network();
let good_script = get_channel_script(&secp_ctx);
let node_1_privkey = &SecretKey::from_slice(&[42; 32]).unwrap();
let node_2_privkey = &SecretKey::from_slice(&[41; 32]).unwrap();
let valid_announcement =
get_signed_channel_announcement(|_| {}, node_1_privkey, node_2_privkey, &secp_ctx);
let node_a_announce = get_signed_node_announcement(|_| {}, node_1_privkey, &secp_ctx);
let node_b_announce = get_signed_node_announcement(|_| {}, node_2_privkey, &secp_ctx);
(
valid_announcement,
chain_source,
network_graph,
good_script,
node_a_announce,
node_b_announce,
get_signed_channel_update(|msg| msg.channel_flags = 0, node_1_privkey, &secp_ctx),
get_signed_channel_update(|msg| msg.channel_flags = 1, node_2_privkey, &secp_ctx),
get_signed_channel_update(
|msg| {
msg.channel_flags = 1;
msg.timestamp += 1;
},
node_2_privkey,
&secp_ctx,
),
)
}
#[test]
fn test_fast_async_lookup() {
let (valid_announcement, chain_source, network_graph, good_script, ..) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
future
.resolve(Ok(TxOut { value: Amount::from_sat(1_000_000), script_pubkey: good_script }));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap();
assert!(network_graph.read_only().channels().get(&scid).is_some());
}
#[test]
fn test_async_lookup() {
let (
valid_announcement,
chain_source,
network_graph,
good_script,
node_a_announce,
node_b_announce,
..,
) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let node_id_1 = valid_announcement.contents.node_id_1;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert!(network_graph.read_only().channels().get(&scid).is_none());
future.resolve(Ok(TxOut { value: Amount::ZERO, script_pubkey: good_script }));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
network_graph.read_only().channels().get(&scid).unwrap();
network_graph.read_only().channels().get(&scid).unwrap();
#[rustfmt::skip]
let is_node_a_announced = network_graph.read_only().nodes().get(&node_id_1).unwrap()
.announcement_info.is_some();
assert!(!is_node_a_announced);
network_graph.update_node_from_announcement(&node_a_announce).unwrap();
network_graph.update_node_from_announcement(&node_b_announce).unwrap();
#[rustfmt::skip]
let is_node_a_announced = network_graph.read_only().nodes().get(&node_id_1).unwrap()
.announcement_info.is_some();
assert!(is_node_a_announced);
}
#[test]
fn test_invalid_async_lookup() {
let (valid_announcement, chain_source, network_graph, ..) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert!(network_graph.read_only().channels().get(&scid).is_none());
let value = Amount::from_sat(1_000_000);
future.resolve(Ok(TxOut { value, script_pubkey: bitcoin::ScriptBuf::new() }));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert!(network_graph.read_only().channels().get(&scid).is_none());
}
#[test]
fn test_failing_async_lookup() {
let (valid_announcement, chain_source, network_graph, ..) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert!(network_graph.read_only().channels().get(&scid).is_none());
future.resolve(Err(UtxoLookupError::UnknownTx));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert!(network_graph.read_only().channels().get(&scid).is_none());
}
#[test]
fn test_updates_async_lookup() {
let (
valid_announcement,
chain_source,
network_graph,
good_script,
node_a_announce,
node_b_announce,
chan_update_a,
chan_update_b,
..,
) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert!(network_graph.read_only().channels().get(&scid).is_none());
assert_eq!(
network_graph.update_node_from_announcement(&node_a_announce).unwrap_err().err,
"Awaiting channel_announcement validation to accept node_announcement"
);
assert_eq!(
network_graph.update_node_from_announcement(&node_b_announce).unwrap_err().err,
"Awaiting channel_announcement validation to accept node_announcement"
);
assert_eq!(
network_graph.update_channel(&chan_update_a).unwrap_err().err,
"Awaiting channel_announcement validation to accept channel_update"
);
assert_eq!(
network_graph.update_channel(&chan_update_b).unwrap_err().err,
"Awaiting channel_announcement validation to accept channel_update"
);
assert!(!notifier.notify_pending());
future
.resolve(Ok(TxOut { value: Amount::from_sat(1_000_000), script_pubkey: good_script }));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert!(network_graph.read_only().channels().get(&scid).unwrap().one_to_two.is_some());
assert!(network_graph.read_only().channels().get(&scid).unwrap().two_to_one.is_some());
assert!(network_graph
.read_only()
.nodes()
.get(&valid_announcement.contents.node_id_1)
.unwrap()
.announcement_info
.is_some());
assert!(network_graph
.read_only()
.nodes()
.get(&valid_announcement.contents.node_id_2)
.unwrap()
.announcement_info
.is_some());
}
#[test]
fn test_latest_update_async_lookup() {
let (
valid_announcement,
chain_source,
network_graph,
good_script,
_,
_,
chan_update_a,
chan_update_b,
chan_update_c,
..,
) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert!(network_graph.read_only().channels().get(&scid).is_none());
assert_eq!(
network_graph.update_channel(&chan_update_a).unwrap_err().err,
"Awaiting channel_announcement validation to accept channel_update"
);
assert_eq!(
network_graph.update_channel(&chan_update_b).unwrap_err().err,
"Awaiting channel_announcement validation to accept channel_update"
);
assert_eq!(
network_graph.update_channel(&chan_update_c).unwrap_err().err,
"Awaiting channel_announcement validation to accept channel_update"
);
assert!(!notifier.notify_pending());
future
.resolve(Ok(TxOut { value: Amount::from_sat(1_000_000), script_pubkey: good_script }));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert_eq!(chan_update_a.contents.timestamp, chan_update_b.contents.timestamp);
let graph_lock = network_graph.read_only();
#[rustfmt::skip]
let one_to_two_update =
graph_lock.channels().get(&scid).as_ref().unwrap().one_to_two.as_ref().unwrap().last_update;
#[rustfmt::skip]
let two_to_one_update =
graph_lock.channels().get(&scid).as_ref().unwrap().two_to_one.as_ref().unwrap().last_update;
assert!(one_to_two_update != two_to_one_update);
}
#[test]
fn test_no_double_lookups() {
let (valid_announcement, chain_source, network_graph, good_script, ..) = get_test_objects();
let scid = valid_announcement.contents.short_channel_id;
let notifier_a = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier_a));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert_eq!(chain_source.get_utxo_call_count.load(Ordering::Relaxed), 1);
let notifier_b = Arc::new(Notifier::new());
let future_b = UtxoFuture::new(Arc::clone(¬ifier_b));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future_b.clone());
assert_eq!(
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel announcement is already being checked"
);
assert_eq!(chain_source.get_utxo_call_count.load(Ordering::Relaxed), 1);
let secp_ctx = Secp256k1::new();
let replacement_pk_1 = &SecretKey::from_slice(&[99; 32]).unwrap();
let replacement_pk_2 = &SecretKey::from_slice(&[98; 32]).unwrap();
let invalid_announcement =
get_signed_channel_announcement(|_| {}, replacement_pk_1, replacement_pk_2, &secp_ctx);
assert_eq!(
network_graph
.update_channel_from_announcement(&invalid_announcement, &Some(&chain_source))
.unwrap_err()
.err,
"Channel being checked async"
);
assert_eq!(chain_source.get_utxo_call_count.load(Ordering::Relaxed), 2);
future
.resolve(Ok(TxOut { value: Amount::from_sat(1_000_000), script_pubkey: good_script }));
assert!(notifier_a.notify_pending());
assert!(!notifier_b.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
#[rustfmt::skip]
let is_test_feature_set =
network_graph.read_only().channels().get(&scid).unwrap().announcement_message
.as_ref().unwrap().contents.features.supports_unknown_test_feature();
assert!(!is_test_feature_set);
}
#[test]
fn test_checks_backpressure() {
let secp_ctx = Secp256k1::new();
let (chain_source, network_graph) = get_network();
let notifier = Arc::new(Notifier::new());
let future = UtxoFuture::new(Arc::clone(¬ifier));
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(future.clone());
let node_1_privkey = &SecretKey::from_slice(&[42; 32]).unwrap();
let node_2_privkey = &SecretKey::from_slice(&[41; 32]).unwrap();
for i in 0..PendingChecks::MAX_PENDING_LOOKUPS {
let valid_announcement = get_signed_channel_announcement(
|msg| msg.short_channel_id += 1 + i as u64,
node_1_privkey,
node_2_privkey,
&secp_ctx,
);
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err();
assert!(!network_graph.pending_checks.too_many_checks_pending());
}
let valid_announcement =
get_signed_channel_announcement(|_| {}, node_1_privkey, node_2_privkey, &secp_ctx);
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err();
assert!(network_graph.pending_checks.too_many_checks_pending());
future.resolve(Err(UtxoLookupError::UnknownTx));
assert!(notifier.notify_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert!(!network_graph.pending_checks.too_many_checks_pending());
}
#[test]
fn test_checks_backpressure_drop() {
let secp_ctx = Secp256k1::new();
let (chain_source, network_graph) = get_network();
let notifier = Arc::new(Notifier::new());
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Async(UtxoFuture::new(notifier));
let node_1_privkey = &SecretKey::from_slice(&[42; 32]).unwrap();
let node_2_privkey = &SecretKey::from_slice(&[41; 32]).unwrap();
for i in 0..PendingChecks::MAX_PENDING_LOOKUPS {
let valid_announcement = get_signed_channel_announcement(
|msg| msg.short_channel_id += 1 + i as u64,
node_1_privkey,
node_2_privkey,
&secp_ctx,
);
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err();
assert!(!network_graph.pending_checks.too_many_checks_pending());
}
let valid_announcement =
get_signed_channel_announcement(|_| {}, node_1_privkey, node_2_privkey, &secp_ctx);
network_graph
.update_channel_from_announcement(&valid_announcement, &Some(&chain_source))
.unwrap_err();
assert!(network_graph.pending_checks.too_many_checks_pending());
*chain_source.utxo_ret.lock().unwrap() = UtxoResult::Sync(Err(UtxoLookupError::UnknownTx));
assert!(network_graph.pending_checks.too_many_checks_pending());
network_graph.pending_checks.check_resolved_futures(&network_graph);
assert!(!network_graph.pending_checks.too_many_checks_pending());
}
}