From 205885029395b229ca9d53629db5d32103224c2b Mon Sep 17 00:00:00 2001 From: Shailendra005 Date: Thu, 6 Aug 2026 16:57:12 +0530 Subject: [PATCH 1/2] fix(agent-server): surface _emit_event_from_thread executor failures --- .../openhands/agent_server/event_service.py | 20 +++++++++- tests/agent_server/test_event_service.py | 40 +++++++++++++++++++ 2 files changed, 58 insertions(+), 2 deletions(-) diff --git a/openhands-agent-server/openhands/agent_server/event_service.py b/openhands-agent-server/openhands/agent_server/event_service.py index 141bf318f9..4211351a45 100644 --- a/openhands-agent-server/openhands/agent_server/event_service.py +++ b/openhands-agent-server/openhands/agent_server/event_service.py @@ -5,6 +5,7 @@ from contextlib import nullcontext, suppress from dataclasses import dataclass, field from datetime import datetime +from functools import partial from pathlib import Path from typing import cast from uuid import UUID, uuid4 @@ -92,6 +93,17 @@ class CredentialBindingActivationTooLate(RuntimeError): pass +def _log_emit_failure(event_kind: str, future: "asyncio.Future[None]") -> None: + """Log an event emission that failed inside the executor thread.""" + if future.cancelled(): + return + error = future.exception() + if error is not None: + logger.error( + "Failed to emit %s from thread: %s", event_kind, error, exc_info=error + ) + + def _without_agent_context_secret( agent: AgentBase, secret_name: str, @@ -838,8 +850,12 @@ def locked_on_event(): conversation._on_event(event) # Run the locked callback in an executor to ensure the event is - # both persisted and sent to WebSocket subscribers - main_loop.run_in_executor(None, locked_on_event) + # both persisted and sent to WebSocket subscribers. The future + # needs a done-callback: otherwise it is collected with an unread + # exception, so a dropped stats or LLM-log event leaves no signal + # beyond a GC-time warning. + future = main_loop.run_in_executor(None, locked_on_event) + future.add_done_callback(partial(_log_emit_failure, type(event).__name__)) def _setup_llm_log_streaming(self, agent: AgentBase) -> None: """Configure LLM log callbacks to stream logs via events.""" diff --git a/tests/agent_server/test_event_service.py b/tests/agent_server/test_event_service.py index 3f53ad3331..1862dbf7bf 100644 --- a/tests/agent_server/test_event_service.py +++ b/tests/agent_server/test_event_service.py @@ -1,5 +1,6 @@ import asyncio import contextlib +import logging import shutil import threading import time @@ -3329,6 +3330,9 @@ def record_and_null(*args, **kwargs): # Simulate concurrent close() nulling self._main_loop mid-call object.__setattr__(event_service, "_main_loop", None) captured_calls.append(args) + # run_in_executor always returns a Future; the caller attaches a + # done-callback to surface executor failures. + return MagicMock() mock_loop.run_in_executor.side_effect = record_and_null @@ -3423,3 +3427,39 @@ async def test_event_service_creates_lease_with_custom_ttl(tmp_path: Path) -> No assert service._lease is not None assert service._lease._ttl_seconds == 10.0 assert (tmp_path / stored.id.hex / LEASE_FILE_NAME).exists() + + +class TestEmitEventFromThreadFailures: + """Executor failures must surface instead of dying with the future (#4386).""" + + @pytest.mark.asyncio + async def test_emit_failure_is_logged( + self, sample_stored_conversation, tmp_path, caplog + ): + service = EventService( + stored=sample_stored_conversation, conversations_dir=tmp_path + ) + conversation = MagicMock() + conversation._state = MagicMock() + conversation._state.__enter__ = MagicMock(return_value=None) + conversation._state.__exit__ = MagicMock(return_value=False) + conversation._on_event.side_effect = RuntimeError("disk full") + + service._conversation = conversation + service._main_loop = asyncio.get_running_loop() + + event = MessageEvent( + id="emit-failure", source="agent", llm_message=Message(role="assistant") + ) + + with caplog.at_level(logging.ERROR): + service._emit_event_from_thread(event) + # Let the executor run and the done-callback fire. + for _ in range(20): + await asyncio.sleep(0) + + messages = [record.getMessage() for record in caplog.records] + assert any("Failed to emit MessageEvent from thread" in m for m in messages), ( + f"emission failure was not surfaced; captured={messages}" + ) + assert any("disk full" in m for m in messages) From 4f9e02b17fc39c2393316e2ce1dbf860e78d4e3a Mon Sep 17 00:00:00 2001 From: Shailendra005 Date: Thu, 6 Aug 2026 17:37:53 +0530 Subject: [PATCH 2/2] fix(agent-server): attach the emit done-callback from the loop thread --- .../openhands/agent_server/event_service.py | 10 +++- tests/agent_server/test_event_service.py | 57 ++++++++++++++++++- 2 files changed, 65 insertions(+), 2 deletions(-) diff --git a/openhands-agent-server/openhands/agent_server/event_service.py b/openhands-agent-server/openhands/agent_server/event_service.py index 4211351a45..2838edd94b 100644 --- a/openhands-agent-server/openhands/agent_server/event_service.py +++ b/openhands-agent-server/openhands/agent_server/event_service.py @@ -855,7 +855,15 @@ def locked_on_event(): # exception, so a dropped stats or LLM-log event leaves no signal # beyond a GC-time warning. future = main_loop.run_in_executor(None, locked_on_event) - future.add_done_callback(partial(_log_emit_failure, type(event).__name__)) + # Attach from the loop thread. This helper is called from foreign + # threads, and `add_done_callback` on an already-done future uses + # plain `call_soon`, which neither is thread-safe nor wakes an idle + # loop -- the report could then sit unfired, which is the failure + # this callback exists to prevent. + main_loop.call_soon_threadsafe( + future.add_done_callback, + partial(_log_emit_failure, type(event).__name__), + ) def _setup_llm_log_streaming(self, agent: AgentBase) -> None: """Configure LLM log callbacks to stream logs via events.""" diff --git a/tests/agent_server/test_event_service.py b/tests/agent_server/test_event_service.py index 1862dbf7bf..6fe47a1652 100644 --- a/tests/agent_server/test_event_service.py +++ b/tests/agent_server/test_event_service.py @@ -16,7 +16,10 @@ from openhands.agent_server.conversation_lease import LEASE_FILE_NAME from openhands.agent_server.conversation_service import ConversationService -from openhands.agent_server.event_service import EventService +from openhands.agent_server.event_service import ( + EventService, + _log_emit_failure, +) from openhands.agent_server.models import ( ConfirmationResponseRequest, EventPage, @@ -3463,3 +3466,55 @@ async def test_emit_failure_is_logged( f"emission failure was not surfaced; captured={messages}" ) assert any("disk full" in m for m in messages) + + @pytest.mark.asyncio + async def test_cancelled_future_is_not_reported(self, caplog): + """A cancelled emission is not a failure and must not log or raise.""" + future: asyncio.Future[None] = asyncio.get_running_loop().create_future() + future.cancel() + # Let the cancellation settle so `cancelled()` is True. + with suppress(asyncio.CancelledError): + await future + + with caplog.at_level(logging.ERROR): + _log_emit_failure("MessageEvent", future) + + assert [record.getMessage() for record in caplog.records] == [] + + @pytest.mark.asyncio + async def test_callback_is_attached_from_the_loop_thread( + self, sample_stored_conversation, tmp_path + ): + """`add_done_callback` must be scheduled with call_soon_threadsafe. + + This helper runs on foreign threads, and attaching directly to an + already-done future would use plain `call_soon`, which does not wake an + idle loop. + """ + service = EventService( + stored=sample_stored_conversation, conversations_dir=tmp_path + ) + conversation = MagicMock() + conversation._state = MagicMock() + conversation._state.__enter__ = MagicMock(return_value=None) + conversation._state.__exit__ = MagicMock(return_value=False) + service._conversation = conversation + + real_loop = asyncio.get_running_loop() + mock_loop = MagicMock() + mock_loop.is_running.return_value = True + finished: asyncio.Future[None] = real_loop.create_future() + finished.set_result(None) + mock_loop.run_in_executor.return_value = finished + service._main_loop = mock_loop + + event = MessageEvent( + id="thread-safety", source="agent", llm_message=Message(role="assistant") + ) + service._emit_event_from_thread(event) + + assert mock_loop.call_soon_threadsafe.called, ( + "done-callback must be attached via call_soon_threadsafe" + ) + scheduled = mock_loop.call_soon_threadsafe.call_args.args + assert scheduled[0] == finished.add_done_callback