From 30beec06b25c5f0279786fd019619b100c2900fe Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 10:48:52 +0000 Subject: [PATCH 1/8] interface: subs_cache: add some comments, type hints --- electrum/interface.py | 18 ++++++++++-------- 1 file changed, 10 insertions(+), 8 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index 918ead8deff5..ef3abc322fa6 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -163,8 +163,8 @@ class NotificationSession(RPCSession): def __init__(self, *args, interface: 'Interface', **kwargs): super(NotificationSession, self).__init__(*args, **kwargs) - self.subscriptions = defaultdict(list) - self.cache = {} + self.subscriptions = defaultdict(list) # type: defaultdict[str, list[asyncio.Queue]] + self.subs_cache = {} # type: dict[str, Any] self._msg_counter = itertools.count(start=1) self.interface = interface self.taskgroup = interface.taskgroup @@ -181,7 +181,7 @@ async def handle_request(self, request): params, result = request.args[:-1], request.args[-1] key = self.get_hashable_key_for_rpc_call(request.method, params) if key in self.subscriptions: - self.cache[key] = result + self.subs_cache[key] = result for queue in self.subscriptions[key]: await queue.put(request.args) else: @@ -223,15 +223,17 @@ def set_default_timeout(self, timeout): self.max_send_delay = timeout async def subscribe(self, method: str, params: List, queue: asyncio.Queue): - # note: until the cache is written for the first time, - # each 'subscribe' call might make a request on the network. key = self.get_hashable_key_for_rpc_call(method, params) + # note: multiple Synchronizers (from different Wallet objects) might sub to the same key, + # hence subscriptions map key->list[queue] self.subscriptions[key].append(queue) - if key in self.cache: - result = self.cache[key] + if key in self.subs_cache: + result = self.subs_cache[key] else: + # note: until subs_cache is written for the first time, + # each 'subscribe' call might make a request on the network. result = await self.send_request(method, params) - self.cache[key] = result + self.subs_cache[key] = result await queue.put(params + [result]) def unsubscribe(self, queue): From df6e92d7c3daee71d0381d38e25206a64d816607 Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 12:16:13 +0000 Subject: [PATCH 2/8] interface: limit subs notification queue sizes to limit memory usage --- electrum/interface.py | 9 ++++++++- electrum/scripts/block_headers.py | 2 +- electrum/synchronizer.py | 2 +- 3 files changed, 10 insertions(+), 3 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index ef3abc322fa6..2dc0147575d9 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -176,6 +176,9 @@ def __init__(self, *args, interface: 'Interface', **kwargs): async def handle_request(self, request): self.maybe_log(f"--> {request}") + # note: size of a single request is bounded by the framer: config.NETWORK_MAX_INCOMING_MSG_SIZE. + # That bound is too relaxed, to allow requesting a large tx. Most messages should + # be MUCH smaller, especially notifications... try: if isinstance(request, Notification): params, result = request.args[:-1], request.args[-1] @@ -183,6 +186,8 @@ async def handle_request(self, request): if key in self.subscriptions: self.subs_cache[key] = result for queue in self.subscriptions[key]: + # note: if queue is full, queue.put will block until a free slot is available. + # This limits memory usage. await queue.put(request.args) else: raise Exception(f'unexpected notification') @@ -223,6 +228,7 @@ def set_default_timeout(self, timeout): self.max_send_delay = timeout async def subscribe(self, method: str, params: List, queue: asyncio.Queue): + assert queue.maxsize > 0, "infinite queue maxsize not allowed" key = self.get_hashable_key_for_rpc_call(method, params) # note: multiple Synchronizers (from different Wallet objects) might sub to the same key, # hence subscriptions map key->list[queue] @@ -234,6 +240,7 @@ async def subscribe(self, method: str, params: List, queue: asyncio.Queue): # each 'subscribe' call might make a request on the network. result = await self.send_request(method, params) self.subs_cache[key] = result + # note: if queue is full, queue.put will block until a free slot is available. This limits memory usage. await queue.put(params + [result]) def unsubscribe(self, queue): @@ -1091,7 +1098,7 @@ async def close(self, *, force_after: int = None): # monitor_connection will cancel tasks async def run_fetch_blocks(self): - header_queue = asyncio.Queue() + header_queue = asyncio.Queue(maxsize=1) # maxsize limits memory usage await self.session.subscribe('blockchain.headers.subscribe', [], header_queue) while True: item = await header_queue.get() diff --git a/electrum/scripts/block_headers.py b/electrum/scripts/block_headers.py index de36173904c5..59ab9bca863f 100755 --- a/electrum/scripts/block_headers.py +++ b/electrum/scripts/block_headers.py @@ -21,7 +21,7 @@ time.sleep(1) print_msg("waiting for network to get connected...") -header_queue = asyncio.Queue() +header_queue = asyncio.Queue(maxsize=1) @log_exceptions async def f(): diff --git a/electrum/synchronizer.py b/electrum/synchronizer.py index ed7369d19ddf..9e4684e343b6 100644 --- a/electrum/synchronizer.py +++ b/electrum/synchronizer.py @@ -71,7 +71,7 @@ def _reset(self): self.scripthash_to_address = {} self._processed_some_notifications = False # so that we don't miss them # Queues - self.status_queue = asyncio.Queue() + self.status_queue = asyncio.Queue(maxsize=15) # maxsize limits memory usage async def _run_tasks(self, *, taskgroup): await super()._run_tasks(taskgroup=taskgroup) From 9252c75a86a8cf7cf49bc288da8f281581f9ed46 Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 14:00:04 +0000 Subject: [PATCH 3/8] interface: subscribe: forbid notifications until initial response --- electrum/interface.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index 2dc0147575d9..11f50c6e033c 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -230,9 +230,6 @@ def set_default_timeout(self, timeout): async def subscribe(self, method: str, params: List, queue: asyncio.Queue): assert queue.maxsize > 0, "infinite queue maxsize not allowed" key = self.get_hashable_key_for_rpc_call(method, params) - # note: multiple Synchronizers (from different Wallet objects) might sub to the same key, - # hence subscriptions map key->list[queue] - self.subscriptions[key].append(queue) if key in self.subs_cache: result = self.subs_cache[key] else: @@ -240,6 +237,12 @@ async def subscribe(self, method: str, params: List, queue: asyncio.Queue): # each 'subscribe' call might make a request on the network. result = await self.send_request(method, params) self.subs_cache[key] = result + # note: multiple Synchronizers (from different Wallet objects) might sub to the same key, + # hence subscriptions map key->list[queue] + # note: only after we got the initial response from the network, we save the subscription. + # This way we disallow force all notifications to arrive after the initial response. + # (as the queue size is bounded, that could even cause a deadlock) + self.subscriptions[key].append(queue) # note: if queue is full, queue.put will block until a free slot is available. This limits memory usage. await queue.put(params + [result]) From c067c264ec19341e8e9e1a8c22d3a1ee418b80c0 Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 12:18:53 +0000 Subject: [PATCH 4/8] interface: limit concurrency of incoming request handling to limit memory usage --- electrum/interface.py | 14 ++++++++++++++ electrum/network.py | 2 ++ 2 files changed, 16 insertions(+) diff --git a/electrum/interface.py b/electrum/interface.py index 11f50c6e033c..99ccea9d02cf 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -161,6 +161,9 @@ class ChainResolutionMode(enum.Enum): class NotificationSession(RPCSession): + INC_REQ_CONCURRENCY_FOR_MAIN_IFACE = 5 + INC_REQ_CONCURRENCY_FOR_OTHER_IFACE = 2 + def __init__(self, *args, interface: 'Interface', **kwargs): super(NotificationSession, self).__init__(*args, **kwargs) self.subscriptions = defaultdict(list) # type: defaultdict[str, list[asyncio.Queue]] @@ -168,7 +171,11 @@ def __init__(self, *args, interface: 'Interface', **kwargs): self._msg_counter = itertools.count(start=1) self.interface = interface self.taskgroup = interface.taskgroup + assert hasattr(self, "cost_hard_limit") # in base class self.cost_hard_limit = 0 # disable aiorpcx resource limits + assert hasattr(self, "initial_concurrent") # in base class + self.initial_concurrent = self.INC_REQ_CONCURRENCY_FOR_OTHER_IFACE + self._incoming_concurrency.set_target(self.INC_REQ_CONCURRENCY_FOR_OTHER_IFACE) # this limits memory usage # To log pre-processed json traffic, uncomment: #self.logger.setLevel(logging.DEBUG) # from aiorpcx @@ -188,6 +195,7 @@ async def handle_request(self, request): for queue in self.subscriptions[key]: # note: if queue is full, queue.put will block until a free slot is available. # This limits memory usage. + # note: there are up to "initial_concurrent" handle_request() calls blocking here. await queue.put(request.args) else: raise Exception(f'unexpected notification') @@ -988,6 +996,12 @@ def is_main_server(self) -> bool: return (self.network.interface == self or self.network.interface is None and self.network.default_server == self.server) + def mark_as_main_server(self) -> None: + """Called when the network switches to this interface.""" + assert self.session + self.session.initial_concurrent = self.session.INC_REQ_CONCURRENCY_FOR_MAIN_IFACE + self.session._incoming_concurrency.set_target(self.session.INC_REQ_CONCURRENCY_FOR_MAIN_IFACE) + async def open_session( self, *, diff --git a/electrum/network.py b/electrum/network.py index c74a0af36752..fc14912a9480 100644 --- a/electrum/network.py +++ b/electrum/network.py @@ -891,6 +891,7 @@ async def switch_to_interface(self, server: ServerAddr): # Stop any current interface in order to terminate subscriptions, # and to cancel tasks in interface.taskgroup. + # This also indirectly undoes i.mark_as_main_server(). if old_server and old_server != server: # don't wait for old_interface to close as that might be slow: await self.taskgroup.spawn(self._close_interface(old_interface)) @@ -907,6 +908,7 @@ async def switch_to_interface(self, server: ServerAddr): self.logger.info(f"switching to {server}") blockchain_updated = i.blockchain != self.blockchain() self.interface = i + i.mark_as_main_server() try: await i.taskgroup.spawn(self._request_server_info(i)) except RuntimeError as e: # see #7677 From 1a9861f7281ba51ab3ba423e0daef9c3c2c9bc3c Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 15:48:52 +0000 Subject: [PATCH 5/8] interface: make sure we don't silently drop notifications --- electrum/interface.py | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index 99ccea9d02cf..d993eff299bc 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -183,6 +183,7 @@ def __init__(self, *args, interface: 'Interface', **kwargs): async def handle_request(self, request): self.maybe_log(f"--> {request}") + # note: the caller enforces a timeout (processing_timeout) on us, after which we will get cancelled # note: size of a single request is bounded by the framer: config.NETWORK_MAX_INCOMING_MSG_SIZE. # That bound is too relaxed, to allow requesting a large tx. Most messages should # be MUCH smaller, especially notifications... @@ -201,9 +202,13 @@ async def handle_request(self, request): raise Exception(f'unexpected notification') else: raise Exception(f'unexpected request. not a notification') - except Exception as e: + except BaseException as e: + # note: we must not silently drop any notification, so we even catch CancelledError(BaseException). + # Instead, make sure we disconnect. self.interface.logger.info(f"error handling request {request}. exc: {repr(e)}") await self.close() + if isinstance(e, asyncio.CancelledError): + raise async def send_request(self, *args, timeout=None, **kwargs): # note: semaphores/timeouts/backpressure etc are handled by @@ -251,7 +256,8 @@ async def subscribe(self, method: str, params: List, queue: asyncio.Queue): # This way we disallow force all notifications to arrive after the initial response. # (as the queue size is bounded, that could even cause a deadlock) self.subscriptions[key].append(queue) - # note: if queue is full, queue.put will block until a free slot is available. This limits memory usage. + # note: if queue is full, queue.put will block until a free slot is available. + # This limits memory usage. FIXME add timeout?? await queue.put(params + [result]) def unsubscribe(self, queue): From 5f849b5cd5e566a470d8b4acf16170f41f14b32d Mon Sep 17 00:00:00 2001 From: SomberNight Date: Thu, 6 Aug 2026 16:32:07 +0000 Subject: [PATCH 6/8] interface: subscriptions: strict checks for response shape part of the motivation is to forbid stuffing extra data inside notifications to inflate their size --- electrum/interface.py | 19 +++++++++++++++---- electrum/synchronizer.py | 11 ++++++++--- 2 files changed, 23 insertions(+), 7 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index d993eff299bc..30412c789a9c 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -126,9 +126,13 @@ def assert_hex_str(val: Any) -> None: raise RequestCorrupted(f'{val!r} should be a hex str') -def assert_dict_contains_field(d: Any, *, field_name: str) -> Any: +def assert_dict(d: Any) -> None: if not isinstance(d, dict): raise RequestCorrupted(f'{d!r} should be a dict') + + +def assert_dict_contains_field(d: Any, *, field_name: str) -> Any: + assert_dict(d) if field_name not in d: raise RequestCorrupted(f'required field {field_name!r} missing from dict') return d[field_name] @@ -1125,10 +1129,17 @@ async def run_fetch_blocks(self): await self.session.subscribe('blockchain.headers.subscribe', [], header_queue) while True: item = await header_queue.get() - raw_header = item[0] - height = raw_header['height'] - header_bytes = bfh(raw_header['hex']) + # parse response + assert len(item) == 1 + resp_header = item[0] + assert_dict(resp_header) + height = assert_dict_contains_field(resp_header, field_name='height') + assert_non_negative_integer(height) + header_hex = assert_dict_contains_field(resp_header, field_name='hex') + header_bytes = bfh(header_hex) + assert len(resp_header) == 2, f"resp_header contains redundant fields. got {resp_header.keys()}" header_dict = blockchain.deserialize_header(header_bytes, height) + # process header self.tip_header = header_dict self.tip = height if self.tip < constants.net.max_checkpoint(): diff --git a/electrum/synchronizer.py b/electrum/synchronizer.py index 9e4684e343b6..20e8653cec27 100644 --- a/electrum/synchronizer.py +++ b/electrum/synchronizer.py @@ -35,7 +35,7 @@ from .util import make_aiohttp_session, NetworkJobOnDefaultServer, random_shuffled_copy, OldTaskGroup from .bitcoin import address_to_scripthash, is_address, neuter_bitcoin_address from .logging import Logger -from .interface import GracefulDisconnect, NetworkTimeout +from .interface import GracefulDisconnect, NetworkTimeout, assert_hash256_str if TYPE_CHECKING: from .network import Network @@ -117,8 +117,13 @@ async def _subscribe_to_address(self, addr): async def handle_status(self): while True: - h, status = await self.status_queue.get() - addr = self.scripthash_to_address[h] + sh, status = await self.status_queue.get() + # basic checks for response + assert_hash256_str(sh) + if status is not None: + assert_hash256_str(status) + # process status + addr = self.scripthash_to_address[sh] self._handling_addr_statuses.add(addr) self.requested_addrs.discard(addr) # ok for addr not to be present await self.taskgroup.spawn(self._on_address_status, addr, status) From 6c7fb405c762c5a884cb4a61241b85067e3d2291 Mon Sep 17 00:00:00 2001 From: SomberNight Date: Fri, 7 Aug 2026 10:53:56 +0000 Subject: [PATCH 7/8] interface: ugly workaround attempt for aiorpcx-internal OOM surface --- electrum/interface.py | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/electrum/interface.py b/electrum/interface.py index 30412c789a9c..2a220773c986 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -214,6 +214,27 @@ async def handle_request(self, request): if isinstance(e, asyncio.CancelledError): raise + async def _throttled_request(self, request): + # If we are already processing too many received messages, pause reading from the transport. + # This limits memory usage. + # FIXME this is a horrible ugly hack, accessing deep internals inside aiorpcx. + # A proper fix instead should be implemented inside aiorpcx. + # A proper fix should also be looking at the cumulative byte size of pending requests, + # instead of request count. + # note: we get called for incoming Requests and Notifications. (not for Responses) + # As we are not server, we don't expect getting Requests: handle_request will raise and close(). + rstransport = self.transport # type: PaddedRSTransport + asyncio_transport = rstransport._asyncio_transport # type: asyncio.Transport + if self._incoming_concurrency._semaphore.locked(): + asyncio_transport.pause_reading() + try: + return await super()._throttled_request(request) + finally: + # Just finished processing one received request. If there are not too many unprocessed ones + # AND the outgoing send buffer has room, we can resume reading: + if not self._incoming_concurrency._semaphore.locked() and rstransport._can_send.is_set(): + asyncio_transport.resume_reading() + async def send_request(self, *args, timeout=None, **kwargs): # note: semaphores/timeouts/backpressure etc are handled by # aiorpcx. the timeout arg here in most cases should not be set From 429738717234457f49e41751547cad9924e2256e Mon Sep 17 00:00:00 2001 From: SomberNight Date: Fri, 7 Aug 2026 12:57:14 +0000 Subject: [PATCH 8/8] tmp1 --- electrum/interface.py | 26 +++++++++++++++++--------- 1 file changed, 17 insertions(+), 9 deletions(-) diff --git a/electrum/interface.py b/electrum/interface.py index 2a220773c986..81cc77a0ab08 100644 --- a/electrum/interface.py +++ b/electrum/interface.py @@ -176,7 +176,7 @@ def __init__(self, *args, interface: 'Interface', **kwargs): self.interface = interface self.taskgroup = interface.taskgroup assert hasattr(self, "cost_hard_limit") # in base class - self.cost_hard_limit = 0 # disable aiorpcx resource limits + self.cost_hard_limit = 0 # disable aiorpcx resource limits # TODO secondaries assert hasattr(self, "initial_concurrent") # in base class self.initial_concurrent = self.INC_REQ_CONCURRENCY_FOR_OTHER_IFACE self._incoming_concurrency.set_target(self.INC_REQ_CONCURRENCY_FOR_OTHER_IFACE) # this limits memory usage @@ -635,6 +635,7 @@ def __init__(self, *, network: 'Network', server: ServerAddr): # Failing verification will get the interface closed. self.tip_header = None # type: Optional[dict] self.tip = 0 + self._tip_unprocessed_evt = asyncio.Event() self._headers_cache = {} # type: Dict[int, bytes] self._rawtx_cache = LRUCache(maxsize=20) # type: LRUCache[str, bytes] # txid->rawtx @@ -1080,7 +1081,8 @@ async def open_session( async with self.taskgroup as group: await group.spawn(self.ping) await group.spawn(self.request_fee_estimates) - await group.spawn(self.run_fetch_blocks) + await group.spawn(self._subscribe_to_headers) + await group.spawn(self._loop_process_header_at_tip) await group.spawn(self.monitor_connection) except aiorpcx.jsonrpc.RPCError as e: if e.code in ( @@ -1145,11 +1147,11 @@ async def close(self, *, force_after: int = None): await self.session.close(force_after=force_after) # monitor_connection will cancel tasks - async def run_fetch_blocks(self): - header_queue = asyncio.Queue(maxsize=1) # maxsize limits memory usage - await self.session.subscribe('blockchain.headers.subscribe', [], header_queue) + async def _subscribe_to_headers(self): + unsanitized_header_queue = asyncio.Queue(maxsize=1) # maxsize limits memory usage + await self.session.subscribe('blockchain.headers.subscribe', [], unsanitized_header_queue) while True: - item = await header_queue.get() + item = await unsanitized_header_queue.get() # parse response assert len(item) == 1 resp_header = item[0] @@ -1166,16 +1168,22 @@ async def run_fetch_blocks(self): if self.tip < constants.net.max_checkpoint(): raise GracefulDisconnect( f"server tip below max checkpoint. ({self.tip} < {constants.net.max_checkpoint()})") - self._mark_ready() self._headers_cache.clear() # tip changed, so assume anything could have happened with chain self._headers_cache[height] = header_bytes + self._tip_unprocessed_evt.set() + + async def _loop_process_header_at_tip(self): + while True: + # wait until tip changes + await self._tip_unprocessed_evt.wait() + self._mark_ready() try: blockchain_updated = await self._process_header_at_tip() finally: self._headers_cache.clear() # to reduce memory usage # header processing done if self.is_main_server() or blockchain_updated: - self.logger.info(f"new chain tip. {height=}") + self.logger.info(f"new chain tip. height={self.tip}") if blockchain_updated: util.trigger_callback('blockchain_updated') self._blockchain_updated.set() @@ -1191,7 +1199,7 @@ async def _process_header_at_tip(self) -> bool: True - new header we didn't have, or reorg """ height, header = self.tip, self.tip_header - async with self.network.bhi_lock: + async with self.network.bhi_lock: # FIXME secondary server can starve main if self.blockchain.height() >= height and self.blockchain.check_header(header): # another interface amended the blockchain return False