diff --git a/agent/relay_runtime.py b/agent/relay_runtime.py index a37432e3e6..67696ae6ca 100644 --- a/agent/relay_runtime.py +++ b/agent/relay_runtime.py @@ -897,6 +897,7 @@ class RelayRuntime: allow_closing: bool = False, failure_label: str = "scope close failed", drain_limit: int = 32, + operation_already_held: bool = False, ) -> str | None: """Pop ``handle``, draining orphaned children in the same session context. @@ -1013,7 +1014,12 @@ class RelayRuntime: error_holder["retry"] = retry_exc try: - self.run_in_session( + run_in_session = ( + self._run_in_session_untracked + if operation_already_held + else self.run_in_session + ) + run_in_session( session, close_with_drain, allow_closing=allow_closing, @@ -1063,6 +1069,7 @@ class RelayRuntime: output={}, allow_closing=True, failure_label="session scope close failed", + operation_already_held=True, ) if failure: failures.append(failure) diff --git a/tests/agent/test_relay_runtime_bounded_scope_ops.py b/tests/agent/test_relay_runtime_bounded_scope_ops.py index 0954ab36da..e1374bd373 100644 --- a/tests/agent/test_relay_runtime_bounded_scope_ops.py +++ b/tests/agent/test_relay_runtime_bounded_scope_ops.py @@ -274,7 +274,10 @@ class TestHealthyPathUnchanged: # Turn scope and session scope both pushed and popped exactly once. assert fake.scope.pushed.count(relay_runtime.TURN_SCOPE) == 1 assert relay_runtime.TURN_SCOPE in fake.scope.popped - assert fake.subscribers.flushed >= 1 + # Session close must not flush process-wide subscribers: another + # session may still own an active publication. Plugin teardown owns + # the final flush after tracked operations drain. + assert fake.subscribers.flushed == 0 def test_healthy_pop_result_propagates_synchronously(self, coordinator): """A healthy pop completes and is observed before end_turn returns."""