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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
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
13 changes: 9 additions & 4 deletions fuzz/fuzz_targets/full_stack_target.rs

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions src/chain/chaininterface.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -78,7 +78,11 @@ pub trait ChainListener: Sync + Send {
fn block_connected(&self, header: &BlockHeader, height: u32, txn_matched: &[&Transaction], indexes_of_txn_matched: &[u32]);
/// Notifies a listener that a block was disconnected.
/// Unlike block_connected, this *must* never be called twice for the same disconnect event.
fn block_disconnected(&self, header: &BlockHeader);
///
/// Provide listeners with filtered txn previously registered as watched because channels are
/// driven by onchain events (tx broadcast, height), a cancel of one of them may conduct to
/// rollback state (ChannelMonitor or Channel).
fn block_disconnected(&self, header: &BlockHeader, height: u32);
}

/// An enum that represents the speed at which we want a transaction to confirm used for feerate
Expand DownExpand Up@@ -279,11 +283,11 @@ impl ChainWatchInterfaceUtil {
}

/// Notify listeners that a block was disconnected.
pub fn block_disconnected(&self, header: &BlockHeader) {
pub fn block_disconnected(&self, block: &Block, height: u32) {
let listeners = self.listeners.lock().unwrap().clone();
for listener in listeners.iter() {
match listener.upgrade() {
Some(arc) => arc.block_disconnected(header),
Some(arc) => arc.block_disconnected(&block.header, height),
None => ()
}
}
Expand Down
90 changes: 75 additions & 15 deletions src/ln/channelmanager.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -332,6 +332,8 @@ pub struct ChannelManager {
channel_state: Mutex<ChannelHolder>,
our_network_key: SecretKey,

channel_closing_waiting_threshold_conf: Mutex<HashMap<u32, Vec<[u8; 32]>>>,

pending_events: Mutex<Vec<events::Event>>,
/// Used when we have to take a BIG lock to make sure everything is self-consistent.
/// Essentially just when we're serializing ourselves out.
Expand DownExpand Up@@ -556,6 +558,8 @@ impl ChannelManager {
}),
our_network_key: keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(HashMap::new()),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),

Expand DownExpand Up@@ -2400,11 +2404,12 @@ impl ChainListener for ChannelManager {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
let mut channel_lock = self.channel_state.lock().unwrap();
let channel_state = channel_lock.borrow_parts();
let short_to_id = channel_state.short_to_id;
let pending_msg_events = channel_state.pending_msg_events;
channel_state.by_id.retain(|_, channel| {
channel_state.by_id.retain(|channel_id, channel| {
let chan_res = channel.block_connected(header, height, txn_matched, indexes_of_txn_matched);
if let Ok(Some(funding_locked)) = chan_res {
pending_msg_events.push(events::MessageSendEvent::SendFundingLocked {
Expand All@@ -2429,20 +2434,24 @@ impl ChainListener for ChannelManager {
for tx in txn_matched {
for inp in tx.input.iter() {
if inp.previous_output == funding_txo.into_bitcoin_outpoint() {
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, closing channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, log_bytes!(channel.channel_id()));
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
log_trace!(self, "Detected channel-closing tx {} spending {}:{}, waiting until {} to close channel {}", tx.txid(), inp.previous_output.txid, inp.previous_output.vout, height + HTLC_FAIL_ANTI_REORG_DELAY - 1, log_bytes!(channel_id[..]));
match channel_closing_lock.entry(height + HTLC_FAIL_ANTI_REORG_DELAY - 1) {
hash_map::Entry::Occupied(mut entry) => {
let mut duplicate = false;
for id in entry.get().iter() {
if *id == *channel_id {
duplicate = true;
break;
}
}
if !duplicate {
entry.get_mut().push(*channel_id);
}
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![*channel_id]);
}
}
return false;
}
}
}
Expand All@@ -2465,6 +2474,25 @@ impl ChainListener for ChannelManager {
}
true
});
if let Some(channel_closings) = channel_closing_lock.remove(&height) {
for channel_id in channel_closings {
log_trace!(self, "Enough confirmations for a broacast commitment tx, channel {} can be closed", log_bytes!(&channel_id[..]));
if let Some(mut channel) = channel_state.by_id.remove(&channel_id) {
if let Some(short_id) = channel.get_short_channel_id() {
short_to_id.remove(&short_id);
}
// It looks like our counterparty went on-chain. We go ahead and
// broadcast our latest local state as well here, just in case its
// some kind of SPV attack, though we expect these to be dropped.
failed_channels.push(channel.force_shutdown());
if let Ok(update) = self.get_channel_update(&channel) {
pending_msg_events.push(events::MessageSendEvent::BroadcastChannelUpdate {
msg: update
});
}
}
}
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
Expand All@@ -2474,7 +2502,7 @@ impl ChainListener for ChannelManager {
}

/// We force-close the channel without letting our counterparty participate in the shutdown
fn block_disconnected(&self, header: &BlockHeader) {
fn block_disconnected(&self, header: &BlockHeader, height: u32) {
let _ = self.total_consistency_lock.read().unwrap();
let mut failed_channels = Vec::new();
{
Expand All@@ -2499,6 +2527,12 @@ impl ChainListener for ChannelManager {
}
});
}
{
let mut channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
if let Some(_) = channel_closing_lock.remove(&(height + HTLC_FAIL_ANTI_REORG_DELAY - 1)) {
// We discard channel_closing there as brooadcast commitment tx has been disconnected, (and may be replaced by a legit closing_signed)
}
}
for failure in failed_channels.drain(..) {
self.finish_force_close_channel(failure);
}
Expand DownExpand Up@@ -2936,6 +2970,15 @@ impl Writeable for ChannelManager {
}
}

let channel_closing_lock = self.channel_closing_waiting_threshold_conf.lock().unwrap();
(channel_closing_lock.len() as u64).write(writer)?;
for (confirmation_height, channel_id) in channel_closing_lock.iter() {
confirmation_height.write(writer)?;
for id in channel_id {
id.write(writer)?;
}
}

Ok(())
}
}
Expand DownExpand Up@@ -3073,6 +3116,21 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
claimable_htlcs.insert(payment_hash, previous_hops);
}

let channel_closing_count: u64 = Readable::read(reader)?;
let mut channel_closing: HashMap<u32, Vec<[u8; 32]>> = HashMap::with_capacity(cmp::min(channel_closing_count as usize, 32));
for _ in 0..channel_closing_count {
let confirmation_height: u32 = Readable::read(reader)?;
let channel_id: [u8; 32] = Readable::read(reader)?;
match channel_closing.entry(confirmation_height) {
hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().push(channel_id);
}
hash_map::Entry::Vacant(entry) => {
entry.insert(vec![channel_id]);
}
}
}

let channel_manager = ChannelManager {
genesis_hash,
fee_estimator: args.fee_estimator,
Expand All@@ -3094,6 +3152,8 @@ impl<'a, R : ::std::io::Read> ReadableArgs<R, ChannelManagerReadArgs<'a>> for (S
}),
our_network_key: args.keys_manager.get_node_secret(),

channel_closing_waiting_threshold_conf: Mutex::new(channel_closing),

pending_events: Mutex::new(Vec::new()),
total_consistency_lock: RwLock::new(()),
keys_manager: args.keys_manager,
Expand Down
Loading