Skip to content

fix: release failed TransferEngine before the next P2PStore retry (#94) - #109

Open
agourakis82 wants to merge 1 commit into
MoonshotAI:mainfrom
agourakis82:fix/p2p-store-cleanup-failed-transfer-engine-retries
Open

agourakis82 wants to merge 1 commit into
MoonshotAI:mainfrom
agourakis82:fix/p2p-store-cleanup-failed-transfer-engine-retries

Conversation

@agourakis82

Copy link
Copy Markdown

Fixes #94.

What was wrong

P2PStore.__init__ retries TransferEngine initialization up to 8 times, constructing a new TransferEngine() on every attempt without releasing the previous failed one first:

for i in range(retry_count):
    self.engine = TransferEngine()
    ret = self.engine.initialize(...)
    if ret == 0:
        break
    # sleep, retry

Per the issue's own suggested investigation, I checked the mooncake-transfer-engine lifecycle API directly against the compiled extension (mooncake/engine.so from mooncake-transfer-engine==0.3.0b0, since I don't have RDMA hardware to install/run it here). The Python bindings expose no close/destroy/shutdown method — only:

initialize, initialize_ext, allocate_managed_buffer, free_managed_buffer,
transfer_sync_write, transfer_sync_read, transfer_sync,
write_bytes_to_buffer, read_bytes_from_buffer, unregister_memory,
get_first_buffer_address

So the only way to release a failed engine's native resources (sockets, ports, fds, shared memory) is to drop the last Python reference to it and let refcounting collect it — the C++ side does own a freeEngine(), but it's only reachable via the object's destructor (the Python wrapper holds a shared_ptr<TransferEngine> with the default deleter), not exposed directly.

Without that, the retry loop's next self.engine = TransferEngine() constructs and initializes the new engine while the old failed one is still alive. On a port-conflict failure specifically, this means the new attempt can collide with the very port the still-referenced old instance is holding onto — turning what should be a transient, recoverable conflict into a self-inflicted repeat failure.

Fix

if ret == 0:
    break
del self.engine
sleep_ms = ...

One line: drop the reference right after a failed initialize(), before the backoff sleep and the next attempt.

Testing

Also no RDMA hardware here, so tests/test_p2p_store_retry.py stubs mooncake.engine with a fake TransferEngine that tracks how many instances are alive concurrently (via a class-level counter in __init__/__del__), and asserts at most one is ever alive across several failing attempts before a success:

$ pytest tests/test_p2p_store_retry.py -v
test_failed_transfer_engine_is_freed_before_next_retry PASSED
test_raises_after_exhausting_retries PASSED

Checked against the pre-fix code too — both fail there (the mocked engine count reaches 2), confirming they catch the actual regression rather than passing vacuously.

Also ran, on top of this fix:

$ pytest tests/ -q -k "not gpu"
109 passed, 1 skipped, 12 deselected
$ ruff check .
All checks passed!
$ ruff format . --check
34 files already formatted

A note for reviewers on the test itself

My first draft of the fake_mooncake fixture used mock.patch.dict(sys.modules, {"mooncake.engine": fake_module}). That crashed a different, unrelated test later in the same pytest run with a segfault reimporting torch._C. Root cause: patch.dict(sys.modules, ...).__exit__ clears the entire sys.modules dict and restores only the pre-__enter__ snapshot — silently evicting every module imported during the with block, including torch (pulled in transitively through checkpoint_engine). A later test then re-imports torch from scratch in the same process, which segfaults since loading torch._C twice per-process isn't supported. Worked around by setting/popping only the one sys.modules key this test actually needs, leaving everything else untouched. Flagging this in case it's useful elsewhere in the test suite if any other test reaches for patch.dict(sys.modules, ...) around code that does real imports.

🤖 Generated with Claude Code

…onshotAI#94)

P2PStore.__init__ retries TransferEngine initialization up to 8 times,
constructing a new TransferEngine() on every attempt without releasing the
previous failed one first. Checked mooncake-transfer-engine's compiled
extension (0.3.0b0): the Python bindings expose no close/destroy/shutdown
method, only initialize/initialize_ext/allocate_managed_buffer/
free_managed_buffer/transfer_sync*/(un)register_memory/*_bytes_to_buffer/
get_first_buffer_address -- so the only way to release a failed engine's
native resources (sockets, ports, fds, shared memory) is to drop the last
Python reference to it and let refcounting collect it immediately.

Without that, the retry loop's next `self.engine = TransferEngine()`
constructs and initializes the new engine *before* the old failed one is
dereferenced, so on a port-conflict failure the new attempt can collide
with the very port the still-alive old instance is holding -- turning a
transient conflict into a self-inflicted repeat failure.

Fix: `del self.engine` right after a failed initialize(), before the
backoff sleep and the next attempt.

Adds tests/test_p2p_store_retry.py: since mooncake needs RDMA hardware
this can't be exercised for real here, so mooncake.engine is stubbed with
a fake TransferEngine that tracks how many instances are alive
concurrently, and the test asserts at most one is ever alive across
several failing attempts. Both new tests fail against the pre-fix code
(mocked engine count reaches 2) and pass after.

Note for reviewers: I initially wrote this fixture with
mock.patch.dict(sys.modules, ...), which crashed a different, unrelated
test later in the same pytest run (segfault reimporting torch._C). Root
cause: patch.dict(sys.modules).__exit__ clears the whole sys.modules dict
and restores only the pre-__enter__ snapshot, silently evicting every
module imported during the `with` block, including torch (pulled in
transitively through checkpoint_engine). Worked around by setting/popping
only the one key this test needs.

Fixes MoonshotAI#94

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>

@koriyoshi2041 koriyoshi2041 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.

Reviewed at eb19844. The failed instance is dereferenced before the backoff and before the next constructor, while the successful instance remains owned by the store; exhausted retries also leave no failed engine alive. Independently ran the focused CPU regression on Python 3.12: 2 passed. The single-owner lifecycle change looks correct.

@agourakis82

Copy link
Copy Markdown
Author

I re-ran the targeted retry-cleanup test locally on the PR head.

HEAD=eb198442802d8d19263e1068af3677cb84c5c557
Python 3.14.3
python -m pytest tests/test_p2p_store_retry.py -v

tests/test_p2p_store_retry.py::test_failed_transfer_engine_is_freed_before_next_retry PASSED
tests/test_p2p_store_retry.py::test_raises_after_exhausting_retries PASSED

2 passed in 6.78s

I also did a small read-through of the retry path. The cleanup ordering looks right for the stated bug: after a failed initialize(), del self.engine drops the last Python reference before the backoff sleep and before the next TransferEngine() is constructed. That is the important point for avoiding a new retry colliding with native state still held by the previous failed instance.

This is only a local CPU/stub validation of the retry lifecycle test; I did not exercise real RDMA/mooncake hardware.

This branch has not been deployed

No deployments
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.

Clean up failed TransferEngine instances during P2PStore retries

2 participants