Skip to content

sdk/python: return yamux backlog slots on close, bound every stream open - #505

Merged
aojea merged 2 commits into
google:mainfrom
aojea:fix/sdk-python-yamux-backlog
Sep 25, 2026
Merged

aojea merged 2 commits into
google:mainfrom
aojea:fix/sdk-python-yamux-backlog

Conversation

@aojea

@aojea aojea commented Sep 25, 2026

Copy link
Copy Markdown
Collaborator

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_stream takes 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) every open_stream blocks forever with no error. A trio task dump injected into the live canary:

=== TASK agent_mesh.session.MeshSession._reservation_loop
  agent_mesh/relay.py:63 reserve_relay
  libp2p/host/basic_host.py:777 new_stream
  libp2p/stream_muxer/yamux/yamux.py:831 open_stream
  trio/_sync.py:501 acquire        <- parked on stream_backlog_semaphore
=== TASK agent_mesh.session.MeshSession._provide_loop
  ... same

py-libp2p main releases 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.py replaces Yamux.open_stream with one that releases the slot the first time the local side closes or resets the stream. A subclass through muxer_opt does not work: MuxerMultistream.new_conn constructs Yamux by name. Gated on libp2p < 0.8; test_the_stream_slot_workaround_retires_with_libp2p_0_8 fails on the pin bump as the reminder to delete it.
  • host.open_stream(host, peer, protocol, timeout) bounds new_stream with 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.
  • Test: 300 open/close cycles on one connection between two SDK hosts. With stock 0.7 the 257th open hits the deadline (verified by running the test with the patch disabled).

hub runs v0.1.0-rc.4, which has neither #502 nor this; it will stay red until the next release.

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.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread sdk/python/src/agent_mesh/host.py Outdated
Comment on lines +62 to +88
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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

critical

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 stream

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread sdk/python/src/agent_mesh/host.py Outdated
Comment on lines +45 to +47
def _libp2p_version() -> tuple[int, int]:
major, minor = importlib.metadata.version("libp2p").split(".")[:2]
return int(major), int(minor)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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.

Suggested change
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

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.
@aojea
aojea merged commit cd21719 into google:main Sep 25, 2026
18 of 19 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant