fix: release failed TransferEngine before the next P2PStore retry (#94) - #109
agourakis82 wants to merge 1 commit into
Conversation
…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
left a comment
There was a problem hiding this comment.
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.
|
I re-ran the targeted retry-cleanup test locally on the PR head. I also did a small read-through of the retry path. The cleanup ordering looks right for the stated bug: after a failed This is only a local CPU/stub validation of the retry lifecycle test; I did not exercise real RDMA/mooncake hardware. |
Fixes #94.
What was wrong
P2PStore.__init__retriesTransferEngineinitialization up to 8 times, constructing a newTransferEngine()on every attempt without releasing the previous failed one first:Per the issue's own suggested investigation, I checked the
mooncake-transfer-enginelifecycle API directly against the compiled extension (mooncake/engine.sofrommooncake-transfer-engine==0.3.0b0, since I don't have RDMA hardware to install/run it here). The Python bindings expose noclose/destroy/shutdownmethod — only: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 ashared_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
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.pystubsmooncake.enginewith a fakeTransferEnginethat 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: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:
A note for reviewers on the test itself
My first draft of the
fake_mooncakefixture usedmock.patch.dict(sys.modules, {"mooncake.engine": fake_module}). That crashed a different, unrelated test later in the samepytestrun with a segfault reimportingtorch._C. Root cause:patch.dict(sys.modules, ...).__exit__clears the entiresys.modulesdict and restores only the pre-__enter__snapshot — silently evicting every module imported during thewithblock, includingtorch(pulled in transitively throughcheckpoint_engine). A later test then re-importstorchfrom scratch in the same process, which segfaults since loadingtorch._Ctwice per-process isn't supported. Worked around by setting/popping only the onesys.moduleskey 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 forpatch.dict(sys.modules, ...)around code that does real imports.🤖 Generated with Claude Code