sdk/python: return yamux backlog slots on close, bound every stream open - #505
Conversation
py-libp2p 0.7 takes one of the 256 slots of a connection's stream_backlog_semaphore for every outbound stream and gives it back only when sending the SYN fails. After 256 streams on one connection every open_stream blocks forever, with no error. A member keeps one connection to its router for as long as it runs and opens a short stream on it every few minutes, so the python canaries reached that after about six hours: a trio task dump from the live process showed the reservation renewal and the DHT provide loop both parked in Yamux.open_stream on the semaphore, and the router dropped the reservation an hour later. Until py-libp2p releases the fix on main (libp2p/py-libp2p#1426), host.py replaces Yamux.open_stream with one that releases the slot the first time the local side closes or resets the stream. A subclass would not do: MuxerMultistream.new_conn constructs Yamux by name, whatever muxer_opt says. The patch and its test are gated on libp2p < 0.8 and go away with the pin bump. Every stream the SDK opens now goes through host.open_stream, which bounds new_stream with the caller's deadline: py-libp2p bounds the protocol negotiation but not the muxer, and a muxer that cannot open a stream must fail the call, so the renewal loop logs a warning and reconnects instead of freezing. Test: 300 open/close cycles on one connection between two SDK hosts; with stock 0.7 the 257th open hits the deadline.
There was a problem hiding this comment.
Code Review
This pull request introduces a workaround for a stream slot leak in py-libp2p versions prior to 0.8 by wrapping stream creation in a custom open_stream function that enforces timeouts and ensures proper slot release. It also adds comprehensive tests to verify this behavior. The review feedback highlights two important improvements: first, handling BaseException (such as trio.Cancelled) during stream creation to prevent slot leaks on timeouts, and second, making the libp2p version parsing more robust to handle pre-release or development version strings without raising a ValueError.
| async def open_stream_releasing_slot(self: Yamux) -> YamuxStream: | ||
| stream = await open_stream(self) | ||
| release = self.stream_backlog_semaphore.release | ||
| close, reset = stream.close, stream.reset | ||
| released = False | ||
|
|
||
| def release_once() -> None: | ||
| nonlocal released | ||
| if not released: | ||
| released = True | ||
| release() | ||
|
|
||
| async def close_and_release() -> None: | ||
| try: | ||
| await close() | ||
| finally: | ||
| release_once() | ||
|
|
||
| async def reset_and_release() -> None: | ||
| try: | ||
| await reset() | ||
| finally: | ||
| release_once() | ||
|
|
||
| stream.close = close_and_release # type: ignore[method-assign] | ||
| stream.reset = reset_and_release # type: ignore[method-assign] | ||
| return stream |
There was a problem hiding this comment.
When open_stream is called within a trio.fail_after block (as done in the open_stream wrapper), a timeout will trigger a cancellation, raising trio.Cancelled (which is a BaseException, not an Exception).
In py-libp2p's Yamux.open_stream, the cleanup block is except Exception as e: self.stream_backlog_semaphore.release(). Because trio.Cancelled inherits from BaseException and not Exception, this cleanup is bypassed when a timeout occurs during stream creation (e.g., during send_packet). This results in a leaked backlog slot on every timeout, eventually blocking the connection forever.
To fix this, we should catch BaseException in our wrapper, check if a stream was added to self.streams, clean it up, and release the semaphore if the exception is not a standard Exception (which the original except Exception block already handles).
async def open_stream_releasing_slot(self: Yamux) -> YamuxStream:
streams_before = set(self.streams.keys())
try:
stream = await open_stream(self)
except BaseException as err:
added_streams = set(self.streams.keys()) - streams_before
if added_streams:
for stream_id in added_streams:
self.streams.pop(stream_id, None)
if not isinstance(err, Exception):
self.stream_backlog_semaphore.release()
raise
release = self.stream_backlog_semaphore.release
close, reset = stream.close, stream.reset
released = False
def release_once() -> None:
nonlocal released
if not released:
released = True
release()
async def close_and_release() -> None:
try:
await close()
finally:
release_once()
async def reset_and_release() -> None:
try:
await reset()
finally:
release_once()
stream.close = close_and_release # type: ignore[method-assign]
stream.reset = reset_and_release # type: ignore[method-assign]
return streamThere was a problem hiding this comment.
Correct, and worse than described: the slot is also kept when the cancellation lands on the streams_lock acquire, before any stream is added, so the suggested added_streams check would still leak that case. Fixed in acb1e00 by replacing open_stream with the same steps and one release_once() on any BaseException after the acquire, instead of wrapping it. test_a_stream_open_cut_by_its_deadline_returns_the_slot stalls the SYN, cuts three opens at their deadline and checks the slot count is back where it was; it fails against stock 0.7.
| def _libp2p_version() -> tuple[int, int]: | ||
| major, minor = importlib.metadata.version("libp2p").split(".")[:2] | ||
| return int(major), int(minor) |
There was a problem hiding this comment.
Parsing the version string by splitting on . and converting directly to int can fail with a ValueError if the installed libp2p package is a pre-release or development version (e.g., 0.8a1 or 0.8.dev0), where the minor version string contains non-digit characters.
Using a try-except block and filtering out non-digit characters from the minor version makes this check much more robust.
| def _libp2p_version() -> tuple[int, int]: | |
| major, minor = importlib.metadata.version("libp2p").split(".")[:2] | |
| return int(major), int(minor) | |
| def _libp2p_version() -> tuple[int, int]: | |
| try: | |
| major, minor = importlib.metadata.version("libp2p").split(".")[:2] | |
| return int(major), int("".join(filter(str.isdigit, minor))) | |
| except Exception: | |
| return 0, 0 |
There was a problem hiding this comment.
Right for 0.8a1, where PEP 440 puts the pre-release directly after the minor. Fixed in acb1e00 with re.match(r"(\d+)\.(\d+)", ...), which reads 0.8a1, 0.8.0rc1 and 0.8.dev0 as (0, 8); an unparseable string maps to (0, 0), so the workaround stays on rather than off. No broad except: a missing distribution should still fail loudly at import.
…ase versions 0.7's open_stream releases the backlog slot only for an Exception. The deadline every SDK stream open now runs under raises trio.Cancelled, a BaseException, so a SYN that waited on a stalled connection until the deadline kept its slot; every such timeout was one slot fewer. Replace open_stream with the same steps and a release on any failure after the acquire, instead of wrapping it. The version check reads the first two numbers with a regex, so a pre-release such as 0.8a1 does not raise at import. Test: three opens cut by their deadline while the SYN stalls, then the slot count is back where it was and a stream opens.
Follow-up to #502. The python canary on bananas ran the reservation renewal for five hours and then went silent again; the router dropped its reservation at 03:45 UTC and both probe CronJobs have failed since.
Cause: a py-libp2p 0.7 bug.
Yamux.open_streamtakes one of 256 per-connection backlog slots for every outbound stream and releases it only when sending the SYN fails. After 256 streams on the router connection (a DHT provide every 10 min, a renewal every hour, gossipsub, auth) everyopen_streamblocks forever with no error. A trio task dump injected into the live canary:py-libp2p
mainreleases the slot on close (libp2p/py-libp2p#1426) but there is no release with it yet; a release request goes upstream separately.Workaround, self-retiring
host.pyreplacesYamux.open_streamwith one that releases the slot the first time the local side closes or resets the stream. A subclass throughmuxer_optdoes not work:MuxerMultistream.new_connconstructsYamuxby name. Gated onlibp2p < 0.8;test_the_stream_slot_workaround_retires_with_libp2p_0_8fails on the pin bump as the reminder to delete it.host.open_stream(host, peer, protocol, timeout)boundsnew_streamwith the caller's deadline and every stream the SDK opens goes through it. py-libp2p bounds protocol negotiation but not the muxer; with this, a muxer that cannot open a stream fails the call, so the renewal loop warns and reconnects instead of freezing.hubrunsv0.1.0-rc.4, which has neither #502 nor this; it will stay red until the next release.