From 5dffd00b65d33d7a6cee04526fd76cafaaaa0229 Mon Sep 17 00:00:00 2001 From: "duanjialing.777" Date: Sat, 26 Sep 2026 14:08:47 +0800 Subject: [PATCH] fix(chat): prevent queue worker lost wakeups Signed-off-by: duanjialing.777 --- loopx/chat_runtime.py | 22 ++++++-- tests/test_chat_session_active_turn.py | 70 ++++++++++++++++++++++++++ 2 files changed, 89 insertions(+), 3 deletions(-) diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 12ac431a58..dfb1896891 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -305,6 +305,7 @@ def __init__( self.session_adapter_locks: dict[str, threading.Lock] = {} self.session_queue_workers: set[str] = set() self.session_queue_threads: dict[str, threading.Thread] = {} + self.session_queue_wakeups: set[str] = set() self.closed = threading.Event() from .chat_loopx_mode import ChatLoopXMode self.loopx_mode = ChatLoopXMode(self) @@ -1073,7 +1074,10 @@ def resume_session_queue( # Admission must not read session files: callers can hold their own # lifecycle fence here. The worker validates the session before effects. with self.lock: - if self.closed.is_set() or session_id in self.session_queue_workers: + if self.closed.is_set(): + return + if session_id in self.session_queue_workers: + self.session_queue_wakeups.add(session_id) return worker = threading.Thread( target=self._drain_session_queue, @@ -1133,8 +1137,17 @@ def _drain_session_queue( work_dir=work_dir, objective=objective, ) + with self.lock: + self.session_queue_wakeups.discard(session_id) turn = self.store.claim_next_queued_turn(session_id) if turn is None: + with self.lock: + if session_id in self.session_queue_wakeups: + continue + worker = self.session_queue_threads.get(session_id) + if worker is threading.current_thread(): + self.session_queue_workers.discard(session_id) + self.session_queue_threads.pop(session_id, None) return turn_id = str(turn["turn_id"]) with self.lock: @@ -1148,8 +1161,11 @@ def _drain_session_queue( ) finally: with self.lock: - self.session_queue_workers.discard(session_id) - self.session_queue_threads.pop(session_id, None) + worker = self.session_queue_threads.get(session_id) + if worker is threading.current_thread(): + self.session_queue_workers.discard(session_id) + self.session_queue_threads.pop(session_id, None) + self.session_queue_wakeups.discard(session_id) def _run_turn( self, diff --git a/tests/test_chat_session_active_turn.py b/tests/test_chat_session_active_turn.py index e9db072509..cc08f2bef0 100644 --- a/tests/test_chat_session_active_turn.py +++ b/tests/test_chat_session_active_turn.py @@ -451,6 +451,76 @@ def create() -> None: assert current["active_turn_id"] is None +def test_enqueue_wakes_worker_after_empty_queue_observation( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + store = ChatSessionStore(tmp_path) + session = store.create_session( + goal_id="goal-one", + agent_id="codex", + executor_endpoint_id="codex", + adapter_kind="codex_app_server", + upstream_thread_id="thread-one", + upstream_mode="chat", + ) + session_id = str(session["session_id"]) + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex") + adapter = _BlockingChatAdapter() + adapter.release.set() + runtime.adapters[session_id] = adapter # type: ignore[assignment] + empty_queue_observed = threading.Event() + release_empty_observation = threading.Event() + original_claim = store.claim_next_queued_turn + claim_count = 0 + + def pause_after_first_empty_claim( + requested_session_id: str, + *, + host_claim_id: str | None = None, + ) -> dict[str, object] | None: + nonlocal claim_count + turn = original_claim( + requested_session_id, + host_claim_id=host_claim_id, + ) + claim_count += 1 + if claim_count == 1: + assert turn is None + empty_queue_observed.set() + assert release_empty_observation.wait(timeout=2) + return turn + + monkeypatch.setattr(store, "claim_next_queued_turn", pause_after_first_empty_claim) + runtime.resume_session_queue( + session_id=session_id, + work_dir=tmp_path, + objective="keep queued work moving", + ) + worker = runtime.session_queue_threads[session_id] + assert empty_queue_observed.wait(timeout=2) + + queued, created = runtime.enqueue_turn( + session_id=session_id, + client_turn_id="enqueue-during-worker-retirement", + message="process after the empty observation", + work_dir=tmp_path, + objective="keep queued work moving", + ) + assert created is True + release_empty_observation.set() + worker.join(timeout=2) + + assert not worker.is_alive() + completed = store.load_turn(session_id, str(queued["turn_id"])) + assert completed is not None + assert completed["status"] == "completed" + assert claim_count >= 2 + assert session_id not in runtime.session_queue_workers + assert session_id not in runtime.session_queue_threads + assert session_id not in runtime.session_queue_wakeups + + def test_queued_turn_rejects_closed_session(tmp_path: Path) -> None: store = ChatSessionStore(tmp_path) session = store.create_session(