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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
190 changes: 175 additions & 15 deletions fuzz/src/chanmon_consistency.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -49,7 +49,7 @@ use lightning::ln::functional_test_utils::*;
use lightning::offers::invoice::{BlindedPayInfo, UnsignedBolt12Invoice};
use lightning::offers::invoice_request::UnsignedInvoiceRequest;
use lightning::onion_message::messenger::{Destination, MessageRouter, OnionMessagePath};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState};
use lightning::util::test_channel_signer::{TestChannelSigner, EnforcementState, ops};
use lightning::util::errors::APIError;
use lightning::util::logger::Logger;
use lightning::util::config::UserConfig;
Expand All@@ -72,6 +72,8 @@ use std::sync::atomic;
use std::io::Cursor;
use bitcoin::bech32::u5;

#[allow(unused)]
const ASYNC_OPS: u32 = ops::GET_PER_COMMITMENT_POINT | ops::RELEASE_COMMITMENT_SECRET | ops::SIGN_COUNTERPARTY_COMMITMENT;
const MAX_FEE: u32 = 10_000;
struct FuzzEstimator {
ret_val: atomic::AtomicU32,
Expand DownExpand Up@@ -297,7 +299,6 @@ impl SignerProvider for KeyProvider {
inner,
state,
disable_revocation_policy_check: false,
available: Arc::new(Mutex::new(true)),
})
}

Expand DownExpand Up@@ -829,7 +830,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == node_id {
for update_add in update_add_htlcs.iter() {
out.locked_write(format!("Delivering update_add_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_add_htlc to node {} from node {}.\n", idx, $node).as_bytes());
if !$corrupt_forward {
dest.handle_update_add_htlc(&nodes[$node].get_our_node_id(), update_add);
} else {
Expand All@@ -844,19 +845,19 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}
}
for update_fulfill in update_fulfill_htlcs.iter() {
out.locked_write(format!("Delivering update_fulfill_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fulfill_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fulfill_htlc(&nodes[$node].get_our_node_id(), update_fulfill);
}
for update_fail in update_fail_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_htlc(&nodes[$node].get_our_node_id(), update_fail);
}
for update_fail_malformed in update_fail_malformed_htlcs.iter() {
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fail_malformed_htlc to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fail_malformed_htlc(&nodes[$node].get_our_node_id(), update_fail_malformed);
}
if let Some(msg) = update_fee {
out.locked_write(format!("Delivering update_fee to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering update_fee to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_update_fee(&nodes[$node].get_our_node_id(), &msg);
}
let processed_change = !update_add_htlcs.is_empty() || !update_fulfill_htlcs.is_empty() ||
Expand All@@ -873,7 +874,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
} });
break;
}
out.locked_write(format!("Delivering commitment_signed to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering commitment_signed to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_commitment_signed(&nodes[$node].get_our_node_id(), &commitment_signed);
break;
}
Expand All@@ -882,15 +883,15 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
events::MessageSendEvent::SendRevokeAndACK { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering revoke_and_ack to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering revoke_and_ack to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_revoke_and_ack(&nodes[$node].get_our_node_id(), msg);
}
}
},
events::MessageSendEvent::SendChannelReestablish { ref node_id, ref msg } => {
for (idx, dest) in nodes.iter().enumerate() {
if dest.get_our_node_id() == *node_id {
out.locked_write(format!("Delivering channel_reestablish to node {}.\n", idx).as_bytes());
out.locked_write(format!("Delivering channel_reestablish to node {} from node {}.\n", idx, $node).as_bytes());
dest.handle_channel_reestablish(&nodes[$node].get_our_node_id(), msg);
}
}
Expand All@@ -913,7 +914,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
_ => if out.may_fail.load(atomic::Ordering::Acquire) {
return;
} else {
panic!("Unhandled message event {:?}", event)
panic!("Unhandled message event on node {}, {:?}", $node, event)
},
}
if $limit_events != ProcessMessages::AllMessages {
Expand DownExpand Up@@ -1289,6 +1290,118 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
},
0x89 => { fee_est_c.ret_val.store(253, atomic::Ordering::Release); nodes[2].maybe_update_chan_fees(); },

#[cfg(async_signing)]
0xa0 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa1 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa2 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[0].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa3 => {
let states = keys_manager_a.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[0].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa4 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa5 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa6 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xa7 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xa8 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xa9 => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaa => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[1].signer_unblocked(None);
}
#[cfg(async_signing)]
0xab => {
let states = keys_manager_b.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 2);
states.values().last().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[1].signer_unblocked(None);
}

#[cfg(async_signing)]
0xac => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_unavailable(ASYNC_OPS);
}
#[cfg(async_signing)]
0xad => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::GET_PER_COMMITMENT_POINT);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xae => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::RELEASE_COMMITMENT_SECRET);
nodes[2].signer_unblocked(None);
}
#[cfg(async_signing)]
0xaf => {
let states = keys_manager_c.enforcement_states.lock().unwrap();
assert_eq!(states.len(), 1);
states.values().next().unwrap().lock().unwrap().set_signer_available(ops::SIGN_COUNTERPARTY_COMMITMENT);
nodes[2].signer_unblocked(None);
}

0xf0 => {
let pending_updates = monitor_a.chain_monitor.list_pending_monitor_updates().remove(&chan_1_funding).unwrap();
if let Some(id) = pending_updates.get(0) {
Expand DownExpand Up@@ -1382,10 +1495,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
// after we resolve all pending events.
// First make sure there are no pending monitor updates, resetting the error state
// and calling force_channel_monitor_updated for each monitor.
*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

out.locked_write(b"Restoring monitors...\n");
if let Some((id, _)) = monitor_a.latest_monitors.lock().unwrap().get(&chan_1_funding) {
monitor_a.chain_monitor.force_channel_monitor_updated(chan_1_funding, *id);
nodes[0].process_monitor_events();
Expand All@@ -1404,7 +1514,10 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
}

// Next, make sure peers are all connected to each other
out.locked_write(b"Reconnecting peers...\n");

if chan_a_disconnected {
out.locked_write(b"Reconnecting node 0 and node 1...\n");
nodes[0].peer_connected(&nodes[1].get_our_node_id(), &Init {
features: nodes[1].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1414,6 +1527,7 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_a_disconnected = false;
}
if chan_b_disconnected {
out.locked_write(b"Reconnecting node 1 and node 2...\n");
nodes[1].peer_connected(&nodes[2].get_our_node_id(), &Init {
features: nodes[2].init_features(), networks: None, remote_network_address: None
}, true).unwrap();
Expand All@@ -1423,8 +1537,33 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
chan_b_disconnected = false;
}

out.locked_write(b"Restoring signers...\n");

*monitor_a.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_b.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;
*monitor_c.persister.update_ret.lock().unwrap() = ChannelMonitorUpdateStatus::Completed;

#[cfg(async_signing)]
{
for state in keys_manager_a.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_b.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
for state in keys_manager_c.enforcement_states.lock().unwrap().values() {
state.lock().unwrap().set_signer_available(!0);
}
nodes[0].signer_unblocked(None);
nodes[1].signer_unblocked(None);
nodes[2].signer_unblocked(None);
}

out.locked_write(b"Running event queues to quiescence...\n");

for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
Expand All@@ -1437,13 +1576,34 @@ pub fn do_test<Out: Output>(data: &[u8], underlying_out: Out, anchors: bool) {
break;
}

out.locked_write(b"All channels restored to normal operation.\n");

// Finally, make sure that at least one end of each channel can make a substantial payment
assert!(
send_payment(&nodes[0], &nodes[1], chan_a, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[1], &nodes[0], chan_a, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 0 and node 1.\n");

assert!(
send_payment(&nodes[1], &nodes[2], chan_b, 10_000_000, &mut payment_id, &mut payment_idx) ||
send_payment(&nodes[2], &nodes[1], chan_b, 10_000_000, &mut payment_id, &mut payment_idx));
out.locked_write(b"Successfully sent a payment between node 1 and node 2.\n");

out.locked_write(b"Flushing pending messages.\n");
for i in 0..std::usize::MAX {
if i == 100 { panic!("It may take may iterations to settle the state, but it should not take forever"); }

// Then, make sure any current forwards make their way to their destination
if process_msg_events!(0, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(1, false, ProcessMessages::AllMessages) { continue; }
if process_msg_events!(2, false, ProcessMessages::AllMessages) { continue; }
// ...making sure any pending PendingHTLCsForwardable events are handled and
// payments claimed.
if process_events!(0, false) { continue; }
if process_events!(1, false) { continue; }
if process_events!(2, false) { continue; }
break;
}

last_htlc_clear_fee_a = fee_est_a.ret_val.load(atomic::Ordering::Acquire);
last_htlc_clear_fee_b = fee_est_b.ret_val.load(atomic::Ordering::Acquire);
Expand Down
4 changes: 1 addition & 3 deletions lightning/src/chain/channelmonitor.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -2944,9 +2944,7 @@ impl<Signer: WriteableEcdsaChannelSigner> ChannelMonitorImpl<Signer> {
},
commitment_txid: htlc.commitment_txid,
per_commitment_number: htlc.per_commitment_number,
per_commitment_point: self.onchain_tx_handler.signer.get_per_commitment_point(
htlc.per_commitment_number, &self.onchain_tx_handler.secp_ctx,
),
per_commitment_point: htlc.per_commitment_point,
feerate_per_kw: 0,
htlc: htlc.htlc,
preimage: htlc.preimage,
Expand Down
3 changes: 3 additions & 0 deletions lightning/src/chain/onchaintx.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -178,6 +178,7 @@ pub(crate) struct ExternalHTLCClaim {
pub(crate) htlc: HTLCOutputInCommitment,
pub(crate) preimage: Option<PaymentPreimage>,
pub(crate) counterparty_sig: Signature,
pub(crate) per_commitment_point: bitcoin::secp256k1::PublicKey,
}

// Represents the different types of claims for which events are yielded externally to satisfy said
Expand DownExpand Up@@ -1177,9 +1178,11 @@ impl<ChannelSigner: WriteableEcdsaChannelSigner> OnchainTxHandler<ChannelSigner>
})
.map(|(htlc_idx, htlc)| {
let counterparty_htlc_sig = holder_commitment.counterparty_htlc_sigs[htlc_idx];

ExternalHTLCClaim {
commitment_txid: trusted_tx.txid(),
per_commitment_number: trusted_tx.commitment_number(),
per_commitment_point: trusted_tx.per_commitment_point(),
htlc: htlc.clone(),
preimage: *preimage,
counterparty_sig: counterparty_htlc_sig,
Expand Down
Loading