Merge pull request #10984 from f321x/fix_onchain_backup_discovery
What changed, and why it matters
This commit fixes how Electrum's Lightning wallet discovers and watches 'on-chain channel backups'—recovery records embedded in funding transactions. Previously, backups could be missed if the wallet learned about a transaction in stages, or if two channels reused the same on-chain funding address. The fix makes the wallet re-check transactions when more of their inputs are identified as belonging to the wallet, and switches watcher callbacks from being keyed by address to being keyed by the unique funding outpoint. A malicious or buggy remote peer could reuse a funding public key, causing two channels to share the same address; the old code might then overwrite the watcher callback for one channel with the other, leaving a backup unwatched and funds unrecoverable if the channel closed.
Treat this as a security fix and include it in the next maintenance release. Users who rely on Lightning channel backups, especially restored wallets, should upgrade. No immediate emergency response is indicated because exploitation requires a malicious or uncooperative remote peer and a specific address-reuse scenario, but the fix prevents a real loss-of-funds vector.
Security signals we found
Loss of funds due to missed channel backup discovery
Watcher callback collision when funding address is reused
Insufficient replay protection for on-chain backup OP_RETURN data
Race/ordering issue between transaction addition and wallet input discovery
Regression test added for same-funding-address channel backup scenario
Evidence from the diff
The patch addresses two related bugs in Electrum’s on-chain backup discovery and LNWatcher callback management. First, maybe_add_backup_from_tx was only called from wallet.on_event_adb_added_tx, which fires once when a transaction is first added. If the wallet later learned that additional inputs of that transaction were ‘is_mine’ (e.g., during gap-limit rolling), the backup check would not re-run. The fix introduces an adb_updated_tx callback, emitted by AddressSynchronizer.add_transaction when new tx inputs are discovered, and wires LNWatcher to call maybe_add_backup_from_tx on both adb_added_tx and adb_updated_tx (and adb_added_verified_tx). Second, LNWatcher callbacks were keyed by funding address. Because a remote peer can reuse the same funding_pubkey, two distinct channels/backups can share the same P2WSH address, causing callback collisions and loss of monitoring. The patch changes LNWatcher and the watchtower plugin to key callbacks by the unique funding outpoint (txid:n). It also adds a replay check in maybe_add_backup_from_tx requiring at least one funding-tx input to be wallet-owned, and prevents overriding existing channels/backups. A new regression test verifies two channels with the same funding address are both watched.
Changed components
electrum/address_synchronizer.pyelectrum/lnutil.pyelectrum/lnwatcher.pyelectrum/lnworker.pyelectrum/plugins/watchtower/watchtower.pyelectrum/submarine_swaps.pyelectrum/wallet.pyelectrum/wallet_db.pytests/test_lnwallet.pyInspect captured patch +117 / −35
### electrum/address_synchronizer.py
@@ -333,7 +333,9 @@ def add_transaction(self, tx: Transaction, *, allow_unrelated=False, is_new=True
for tx_hash2 in conflicting_txns:
self.remove_transaction(tx_hash2)
# add inputs
+ txi_changed = False
def add_value_from_prev_output():
+ nonlocal txi_changed
# note: this takes linear time in num is_mine outputs of prev_tx
addr = self.get_txin_address(txi)
if addr and self.is_mine(addr):
@@ -343,7 +345,7 @@ def add_value_from_prev_output():
except KeyError:
pass
else:
- self.db.add_txi_addr(tx_hash, addr, ser, v)
+ txi_changed |= self.db.add_txi_addr(tx_hash, addr, ser, v)
self.invalidate_cache()
for txi in tx.inputs():
if txi.is_coinbase_input():
@@ -366,15 +368,22 @@ def add_value_from_prev_output():
# give v to txi that spends me
next_tx = self.db.get_spent_outpoint(tx_hash, n)
if next_tx is not None:
- self.db.add_txi_addr(next_tx, addr, ser, v)
+ is_new_txi = self.db.add_txi_addr(next_tx, addr, ser, v)
self._add_tx_to_local_history(next_tx)
+ if is_new_txi:
+ spender_tx = self.db.get_transaction(next_tx)
+ assert spender_tx
+ util.trigger_callback('adb_updated_tx', self, next_tx, spender_tx)
# add to local history
self._add_tx_to_local_history(tx_hash)
# save
self.db.add_transaction(tx_hash, tx)
self.db.add_num_inputs_to_tx(tx_hash, len(tx.inputs()))
if is_new:
util.trigger_callback('adb_added_tx', self, tx_hash, tx)
+ elif txi_changed:
+ # adb_updated_tx: we learned of more is_mine inputs of an already known tx
+ util.trigger_callback('adb_updated_tx', self, tx_hash, tx)
return True
@with_lock
### electrum/lnutil.py
@@ -336,7 +336,7 @@ class ChannelBackupStorage:
def funding_outpoint(self):
return Outpoint(self.funding_txid, self.funding_index)
- def channel_id(self):
+ def channel_id(self) -> bytes:
chan_id, _ = channel_id_from_funding_tx(self.funding_txid, self.funding_index)
return chan_id
### electrum/lnwatcher.py
@@ -12,7 +12,7 @@
)
from .transaction import Transaction, TxOutpoint
from .logging import Logger
-from .address_synchronizer import TX_HEIGHT_LOCAL
+from .address_synchronizer import TX_HEIGHT_LOCAL, AddressSynchronizer
from .lnutil import REDEEM_AFTER_DOUBLE_SPENT_DELAY
from .lnsweep import KeepWatchingTXO, SweepInfo, MaybeSweepInfo
@@ -31,7 +31,7 @@ def __init__(self, lnworker: 'LNWallet'):
Logger.__init__(self)
self.adb = lnworker.wallet.adb
self.config = lnworker.config
- self.callbacks = {} # type: Dict[str, Callable[[], Awaitable[None]]] # address -> lambda function
+ self.callbacks = {} # type: Dict[str, Callable[[], Awaitable[None]]] # address/txo -> lambda function
self.network = None
self.register_callbacks()
self._pending_force_closes = {} # type: Dict['AbstractChannel', int] # chan -> lowest remote htlc timeout height
@@ -72,10 +72,10 @@ async def _callback_loop(self):
await self.trigger_callbacks()
await asyncio.sleep(self.CALLBACK_LOOP_POLL_INTERVAL_SEC)
- def remove_callback(self, address: str) -> None:
- self.callbacks.pop(address, None)
+ def remove_callback(self, key: str) -> None: # address or outpoint
+ self.callbacks.pop(key, None)
- def add_callback(
+ def add_address_callback(
self,
address: str,
callback: Callable[[], Awaitable[None]],
@@ -90,17 +90,31 @@ def add_callback(
# subscribed to, and will have adb.synchronizer sub to them again.
# (even for old redeemed channels and old swaps)
self.adb.add_address(address)
+ assert address not in self.callbacks, f"not overriding existing callback: {address=}"
self.callbacks[address] = callback
+ def add_outpoint_callback(
+ self,
+ outpoint: str,
+ callback: Callable[[], Awaitable[None]],
+ *,
+ address: str,
+ subscribe: bool = True,
+ ) -> None:
+ if subscribe: # FIXME: see add_address_callback
+ self.adb.add_address(address)
+ assert outpoint not in self.callbacks, f"not overriding existing callback: {outpoint=}"
+ self.callbacks[outpoint] = callback
+
async def trigger_callbacks(self, *, requires_synchronizer: bool = True):
if requires_synchronizer and not self.adb.synchronizer:
self.logger.debug("synchronizer not set yet")
return
- for address, callback in list(self.callbacks.items()):
+ for key, callback in list(self.callbacks.items()):
try:
await callback()
except Exception:
- self.logger.exception(f"LNWatcher callback failed {address=}")
+ self.logger.exception(f"LNWatcher callback failed {key=}")
# send callback to GUI
util.trigger_callback('wallet_updated', self.lnworker.wallet)
self._last_callback_trigger_ts = now()
@@ -110,16 +124,28 @@ async def on_event_blockchain_updated(self, *args):
await self.trigger_callbacks()
@event_listener
- async def on_event_adb_added_tx(self, adb, tx_hash, tx):
- # called if we add local tx
+ async def on_event_adb_added_tx(self, adb: AddressSynchronizer, tx_hash: str, tx: Transaction):
+ # called for every tx added to the adb: by the synchronizer, or if we add a local tx
if adb != self.adb:
return
+ if adb.db.is_in_verified_tx(tx_hash):
+ self.lnworker.maybe_add_backup_from_tx(tx)
await self.trigger_callbacks()
@event_listener
- async def on_event_adb_added_verified_tx(self, adb, tx_hash):
+ async def on_event_adb_updated_tx(self, adb: AddressSynchronizer, tx_hash: str, tx: Transaction):
+ if adb != self.adb:
+ return
+ if adb.db.is_in_verified_tx(tx_hash):
+ self.lnworker.maybe_add_backup_from_tx(tx)
+
+ @event_listener
+ async def on_event_adb_added_verified_tx(self, adb: AddressSynchronizer, tx_hash: str):
if adb != self.adb:
return
+ if tx := adb.db.get_transaction(tx_hash):
+ # if we don't have the tx yet, the backup will be added once the tx gets added
+ self.lnworker.maybe_add_backup_from_tx(tx)
await self.trigger_callbacks()
@event_listener
@@ -132,7 +158,9 @@ def add_channel(self, chan: 'AbstractChannel') -> None:
outpoint = chan.funding_outpoint.to_str()
address = chan.get_funding_address()
callback = lambda: self.check_onchain_situation(address, outpoint)
- self.add_callback(address, callback, subscribe=chan.need_to_subscribe())
+ # keyed by outpoint, as we cannot fully enforce uniqueness of the funding address.
+ # a remote peer might reuse the funding_pubkey, causing the address to conflict with a channel backup.
+ self.add_outpoint_callback(outpoint, callback, address=address, subscribe=chan.need_to_subscribe())
@ignore_exceptions
@log_exceptions
@@ -152,7 +180,7 @@ async def check_onchain_situation(self, address: str, funding_outpoint: str) ->
if closing_tx:
keep_watching = await self.sweep_commitment_transaction(funding_outpoint, closing_tx)
if not keep_watching:
- self.remove_callback(address)
+ self.remove_callback(funding_outpoint)
else:
self.logger.info(f"channel {funding_outpoint} closed by {closing_txid}. still waiting for tx itself...")
keep_watching = True
### electrum/lnworker.py
@@ -1839,6 +1839,9 @@ def decrypt_cb_data(self, encrypted_data: bytes, funding_address: str) -> bytes:
def encrypt_cb_data(self, data: bytes, funding_address: str) -> bytes:
funding_scripthash = bytes.fromhex(address_to_scripthash(funding_address))
nonce = funding_scripthash[0:12]
+ # note: would have been nice, if besides the funding_script, we also committed to
+ # the funding_amount (sats), and maybe all the inputs of the funding_tx. Without that,
+ # the OP_RETURN can be replayed.
# note: we are only using chacha20 instead of chacha20+poly1305 to save onchain space
# (not have the 16 byte MAC). Otherwise, the latter would be preferable.
return chacha20_encrypt(key=self.backup_key, data=data, nonce=nonce)
@@ -3760,6 +3763,7 @@ def remove_channel(self, chan_id):
with self.lock:
self._channels.pop(chan_id)
self.db.get('channels').pop(chan_id.hex())
+ self.lnwatcher.remove_callback(chan.funding_outpoint.to_str())
self.wallet.set_reserved_addresses_for_chan(chan, reserved=False)
util.trigger_callback('channels_updated', self.wallet)
@@ -3899,6 +3903,7 @@ def import_channel_backup(self, data):
self.wallet.set_reserved_addresses_for_chan(cb, reserved=True)
self.wallet.save_db()
util.trigger_callback('channels_updated', self.wallet)
+ self.lnwatcher.remove_callback(cb.funding_outpoint.to_str())
self.lnwatcher.add_channel(cb)
if not cb.can_sweep_their_ctx_to_remote():
# the user has lost their channel state and cannot locally force close. If they'd request a remote fclose
@@ -3933,6 +3938,7 @@ def remove_channel_backup(self, channel_id):
raise Exception('Channel not found')
with self.lock:
self._channel_backups.pop(channel_id)
+ self.lnwatcher.remove_callback(chan.funding_outpoint.to_str())
self.wallet.set_reserved_addresses_for_chan(chan, reserved=False)
self.wallet.save_db()
util.trigger_callback('channels_updated', self.wallet)
@@ -3993,34 +3999,45 @@ async def _request_fclose(addresses):
if not success:
raise Exception('failed to connect')
- def maybe_add_backup_from_tx(self, tx):
+ def maybe_add_backup_from_tx(self, tx: Transaction):
+ """note: currently no support for batched channel opens"""
+ assert self.wallet.adb.db.is_in_verified_tx(tx.txid())
+ if not any(self.wallet.is_mine(self.wallet.adb.get_txin_address(txin)) for txin in tx.inputs()):
+ # only allow funding tx with inputs of our wallet to prevent replay of the channel backup.
+ # note: is_mine can be false during initial wallet synchronization: we might learn
+ # of more-and-more inputs of being is_mine, as we roll the gap_limit forward.
+ # Hence maybe_add_backup_from_tx also needs to be called on adb_updated_tx.
+ # note: if the channel was funded with wallet-external UTXOs we won't detect the backup (we don't do this).
+ return
+ funding_txid = tx.txid()
+ if any(funding_txid == c.funding_outpoint.txid for c in self.get_channel_objects().values()):
+ # Check we don't override imported backups or full channels.
+ return
funding_address = None
node_id_prefix = None
+ # note: loop is quadratic but that's ok as we require at least 1 tx input to be is_mine
for i, o in enumerate(tx.outputs()):
script_type = get_script_type_from_output_script(o.scriptpubkey)
if script_type == 'p2wsh':
- funding_index = i
- funding_address = o.address
for o2 in tx.outputs():
if o2.scriptpubkey.startswith(bytes([opcodes.OP_RETURN])):
encrypted_data = o2.scriptpubkey[2:]
- data = self.decrypt_cb_data(encrypted_data, funding_address)
+ data = self.decrypt_cb_data(encrypted_data, o.address)
if data.startswith(CB_MAGIC_BYTES):
+ funding_index = i
+ funding_address = o.address
node_id_prefix = data[len(CB_MAGIC_BYTES):]
if node_id_prefix is None:
return
- funding_txid = tx.txid()
cb_storage = OnchainChannelBackupStorage(
node_id_prefix=node_id_prefix,
funding_txid=funding_txid,
funding_index=funding_index,
funding_address=funding_address,
is_initiator=True)
- channel_id = cb_storage.channel_id().hex()
- if channel_id in self.db.get_dict("channels"):
- return
self.logger.info(f"adding backup from tx")
d = self.db.get_dict("onchain_channel_backups")
+ channel_id: str = cb_storage.channel_id().hex()
d[channel_id] = cb_storage
cb = ChannelBackup(cb_storage, lnworker=self)
self.wallet.set_reserved_addresses_for_chan(cb, reserved=True)
### electrum/plugins/watchtower/watchtower.py
@@ -70,19 +70,19 @@ def __init__(self, network: 'Network'):
wallet_db = WalletDB('', storage=None, upgrade=True)
self.adb = AddressSynchronizer(wallet_db, self.config, name=self.diagnostic_name())
self.adb.start_network(network)
- self.callbacks = {} # address -> lambda function
+ self.callbacks = {} # outpoint -> lambda function
self.register_callbacks()
# status gets populated when we run
self.channel_status = {}
self.network = network
self.sweepstore = SweepStore(os.path.join(self.config.path, "watchtower_db"), network)
- def remove_callback(self, address):
- self.callbacks.pop(address, None)
+ def remove_callback(self, outpoint):
+ self.callbacks.pop(outpoint, None)
- def add_callback(self, address, callback):
+ def add_callback(self, outpoint, callback, *, address):
self.adb.add_address(address)
- self.callbacks[address] = callback
+ self.callbacks[outpoint] = callback
@event_listener
async def on_event_blockchain_updated(self, *args):
@@ -112,7 +112,7 @@ async def trigger_callbacks(self):
if not self.adb.synchronizer:
self.logger.info("synchronizer not set yet")
return
- for address, callback in list(self.callbacks.items()):
+ for outpoint, callback in list(self.callbacks.items()):
await callback()
async def stop(self):
@@ -121,7 +121,7 @@ async def stop(self):
def add_channel(self, outpoint: str, address: str) -> None:
callback = lambda: self.check_onchain_situation(address, outpoint)
- self.add_callback(address, callback)
+ self.add_callback(outpoint, callback, address=address)
def diagnostic_name(self):
return "watchtower"
@@ -231,7 +231,7 @@ async def broadcast_or_log(self, funding_outpoint: str, tx: Transaction):
return txid
async def get_ctn(self, outpoint, addr):
- if addr not in self.callbacks.keys():
+ if outpoint not in self.callbacks.keys():
self.logger.info(f'watching new channel: {outpoint} {addr}')
self.add_channel(outpoint, addr)
return await self.sweepstore.get_ctn(outpoint, addr)
@@ -252,6 +252,7 @@ async def f():
return self.network.run_from_another_thread(f())
async def unwatch_channel(self, address, funding_outpoint):
+ self.remove_callback(funding_outpoint)
await self.sweepstore.remove_sweep_tx(funding_outpoint)
await self.sweepstore.remove_channel(funding_outpoint)
### electrum/submarine_swaps.py
@@ -787,7 +787,7 @@ def get_swap(self, payment_hash: bytes) -> Optional[SwapData]:
def add_lnwatcher_callback(self, swap: SwapData) -> None:
callback = lambda: self._claim_swap(swap)
- self.lnwatcher.add_callback(swap.lockup_address, callback)
+ self.lnwatcher.add_address_callback(swap.lockup_address, callback)
async def hold_invoice_callback(self, payment_hash: bytes) -> None:
# note: this assumes the wallet has been unlocked
### electrum/wallet.py
@@ -654,8 +654,6 @@ def on_event_adb_added_tx(self, adb, tx_hash: str, tx: Transaction):
if not self.tx_is_related(tx):
return
self.clear_tx_parents_cache()
- if self.lnworker:
- self.lnworker.maybe_add_backup_from_tx(tx)
self._update_invoices_and_reqs_touched_by_tx(tx)
util.trigger_callback('new_transaction', self, tx)
### electrum/wallet_db.py
@@ -1702,7 +1702,8 @@ def get_txo_addr(self, tx_hash: str, address: str) -> Dict[int, Tuple[int, bool]
return {int(n): (v, cb) for (n, (v, cb)) in d.items()}
@modifier
- def add_txi_addr(self, tx_hash: str, addr: str, ser: str, v: int) -> None:
+ def add_txi_addr(self, tx_hash: str, addr: str, ser: str, v: int) -> bool:
+ """Returns True if the item was newly added to the DB, or False if it was already there."""
assert isinstance(tx_hash, str)
assert isinstance(addr, str)
assert isinstance(ser, str)
@@ -1712,7 +1713,10 @@ def add_txi_addr(self, tx_hash: str, addr: str, ser: str, v: int) -> None:
d = self.txi[tx_hash]
if addr not in d:
d[addr] = {}
+ if d[addr].get(ser) == v:
+ return False
d[addr][ser] = v
+ return True
@modifier
def add_txo_addr(self, tx_hash: str, addr: str, n: Union[int, str], v: int, is_coinbase: bool) -> None:
### tests/test_lnwallet.py
@@ -1,3 +1,4 @@
+import dataclasses
import logging
import os
import asyncio
@@ -970,3 +971,27 @@ async def test_imported_channel_backup_upgrade(self):
# now try importing an older version again, this should not work
with self.assertRaises(UserFacingException):
alice.lnworker.import_channel_backup(backup_v2)
+
+ async def test_channel_backup_with_same_funding_address_as_channel(self):
+ """
+ We restore from seed.
+ The remote peer reuses their funding_pubkey to open a channel to us.
+ We import a channel backup of a different channel using the same funding_pubkey.
+ LNWatcher should watch both channels."""
+ alice = self.create_deterministic_wallet(self.alice_instance)
+ chan = await self.fund_and_open_channel(alice, anchors=True)
+ lnwatcher = alice.lnworker.lnwatcher
+
+ # alice imports a backup of another channel, with the same funding address
+ backup = alice.lnworker.export_channel_backup(chan.channel_id)
+ cb_storage = ImportedChannelBackupStorage.from_encrypted_str(backup, password=alice.get_fingerprint())
+ cb_storage = dataclasses.replace(cb_storage, funding_txid=os.urandom(32).hex())
+ alice.lnworker.import_channel_backup('channel_backup:' + pw_encode_with_version_and_mac(cb_storage.to_bytes(), alice.get_fingerprint()))
+ cb = alice.lnworker.channel_backups[cb_storage.channel_id()]
+ self.assertEqual(chan.get_funding_address(), cb.get_funding_address())
+ self.assertEqual({chan.funding_outpoint.to_str(), cb.funding_outpoint.to_str()}, set(lnwatcher.callbacks))
+
+ # wallet load
+ lnwatcher.callbacks.clear()
+ alice.lnworker.subscribe_to_channels()
+ self.assertEqual({chan.funding_outpoint.to_str(), cb.funding_outpoint.to_str()}, set(lnwatcher.callbacks))Why this scored 59/100
Community notes
Notes can correct, qualify, or add evidence to the AI analysis. Every note shown here has been validated by a human moderator.
The AI analysis stands alone for now. Submit a note if you can add evidence or important context.