From d7550e61cc5449e838cacb458826fd57bdd4df1b Mon Sep 17 00:00:00 2001 From: Chao Wang <26245345+ChaoWao@users.noreply.github.com> Date: Mon, 17 Aug 2026 19:17:59 -0700 Subject: [PATCH] Fix: drop a dead guard and state two facts release_buffer left implicit MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `_abandoned_run_handles`' scan guarded `handle._resources` against `None` three lines below a loop that dereferences the same field directly. `RunHandle.__init__` substitutes a fresh `_RunResources()` when its argument is `None` and nothing else ever assigns the field, so the guard could not fire — while the asymmetry told a reader one of the two paths could. Both loops now dereference it the same way. Two facts the surrounding docstrings asserted only one side of: - `release_buffer` states that a failed close never tells a descendant to drop a mapping the owner still considers live. The converse is what actually needs saying: `Buffer.close()` unlinks from its `finally`, so a raising `shm.close()` still removes the name, and the backing can be nameless while descendants keep mappings this call never told them to drop. Recorded together with what makes that window diagnosable — the named `FileNotFoundError` from `ImportRegistry.materialize`. - `_record_touched_identities` runs before `_admit_task_submission`, which can raise, so the touched set can name a task that never went out. The group entry points already document the mirrored case (a rejected member dispatches nothing, so nothing is recorded) while this one, the reason the ordering is chosen at all, was unstated. Recording after admission would invert the direction of the error into the unsafe one. `release_buffer`'s docstring is also split into paragraphs and its one mid-sentence line break repaired; no wording in the pre-existing text changed. --- python/simpler/worker.py | 54 ++++++++++++++++++++++++++-------------- 1 file changed, 35 insertions(+), 19 deletions(-) diff --git a/python/simpler/worker.py b/python/simpler/worker.py index 1624e52841..58437c78f0 100644 --- a/python/simpler/worker.py +++ b/python/simpler/worker.py @@ -10092,6 +10092,14 @@ def _record_touched_identities(self, args: Any) -> None: ``submit_next_level_group``, ``submit_sub``, ``submit_sub_group``), per the ``current_resources = self._building_run_resources; if ... is not None`` idiom used elsewhere for the same "attach to the open run, if any" shape. + + The set over-approximates in one direction on purpose: every caller records before + ``_admit_task_submission``, which can itself raise on a sticky ordered-cleanup failure, so an + identity can be recorded for a task that never went out. That costs a spurious + ``release_buffer`` refusal, recoverable through ``close()``, where recording after admission + would leave a dispatched task's identity unrecorded and let a release unlink under it. The + group entry points additionally validate every member before recording any, since a rejected + member dispatches nothing at all. """ resources = self._building_run_resources if resources is None: @@ -10508,29 +10516,38 @@ def release_buffer(self, buffer: Buffer) -> None: sent it — a Buffer never goes away while a dispatched task still names it. All three dispatch paths retain: a SUB task maps the identity into a sub-worker process just as a NEXT_LEVEL task maps it into a child, so unlinking the backing under either one faults the - consumer on a segment that no longer has a name. The L3+ check takes - ``_submit_mu`` first: a handle is visible in ``_accepted_run_handles`` before its - orchestration callback (where ``touched_identities`` gets populated) has run, and that - callback is what ``_submit_mu`` already serializes graph construction against, so taking it - here means the check only ever runs between callbacks, never mid-callback with a + consumer on a segment that no longer has a name. + + The L3+ check takes ``_submit_mu`` first: a handle is visible in ``_accepted_run_handles`` + before its orchestration callback (where ``touched_identities`` gets populated) has run, and + that callback is what ``_submit_mu`` already serializes graph construction against, so taking + it here means the check only ever runs between callbacks, never mid-callback with a not-yet-complete touched set. ``_abandoned_run_handles`` is scanned in the same block and without the ``_cleanup_published`` test: ``_publish_abandoned_run`` sets that flag and drops the handle from the accepted set while the run itself stays retained until native teardown drains it, so an abandoned run is exactly the case where the flag stops describing whether the device is done with the backing. Such a buffer therefore stops being releasable through this API for the Worker's remaining life, which strands nothing: ``close()`` reclaims it via - ``_release_all_buffers`` calling ``Buffer.close()`` directly. The L2 check is independent (a - separate run-id namespace with - no callback to serialize against — ``_chip_run_touched_identities`` is written atomically - alongside ``_chip_runs`` under ``_registry_lock`` instead, see ``_submit_l2_locked``), so the - two checks run sequentially rather than under one shared lock. Neither is checked once - ``buffer`` is already closed, matching ``Buffer.close()``'s own idempotency. The entry - survives a failed close, so ``_release_all_buffers`` still reports the leak at close() rather - than losing it here — the import-cache broadcast only fires once close() has actually - succeeded, so a failed release never tells a descendant to drop a mapping the owner still - considers live. The slot is dropped only when it still holds *this* buffer: a buffer_id - minted elsewhere can collide with a registry key, and evicting the live entry it names would - strand that backing.""" + ``_release_all_buffers`` calling ``Buffer.close()`` directly. + + The L2 check is independent (a separate run-id namespace with no callback to serialize + against — ``_chip_run_touched_identities`` is written atomically alongside ``_chip_runs`` + under ``_registry_lock`` instead, see ``_submit_l2_locked``), so the two checks run + sequentially rather than under one shared lock. Neither is checked once ``buffer`` is already + closed, matching ``Buffer.close()``'s own idempotency. + + The entry survives a failed close, so ``_release_all_buffers`` still reports the leak at + close() rather than losing it here — the import-cache broadcast only fires once close() has + actually succeeded, so a failed release never tells a descendant to drop a mapping the owner + still considers live. The converse does not hold, and is the one asymmetry here: a + ``Buffer.close()`` whose ``shm.close()`` raises has still unlinked the name (its ``finally`` + runs the unlink), so the backing can be nameless while descendants keep mappings this call + never told them to drop. A descendant materializing that identity afterwards gets the named + ``FileNotFoundError`` from ``ImportRegistry.materialize``, not a silent bad mapping. + + The slot is dropped only when it still holds *this* buffer: a buffer_id minted elsewhere can + collide with a registry key, and evicting the live entry it names would strand that + backing.""" if not buffer.closed: # Exclusive, not shared: `shared()` would already exclude admission and so # satisfy the "never mid-callback" argument above, but this keeps the @@ -10541,8 +10558,7 @@ def release_buffer(self, buffer: Buffer) -> None: if not handle._cleanup_published and buffer.identity in handle._resources.touched_identities: raise RuntimeError(f"release_buffer: {buffer.identity} is still referenced by an in-flight run") for handle in self._abandoned_run_handles: - resources = handle._resources - if resources is not None and buffer.identity in resources.touched_identities: + if buffer.identity in handle._resources.touched_identities: raise RuntimeError( f"release_buffer: {buffer.identity} is still referenced by an abandoned run " f"whose native teardown has not completed"