From 4abb7766c41f826852101f4fe9a29759e80dfb91 Mon Sep 17 00:00:00 2001 From: Michael Xu Date: Thu, 1 Oct 2026 04:36:33 -0500 Subject: [PATCH 1/2] fix(tracing): stop set_tracing_processor_configs deadlocking on its own lock TracingProcessorManager.set_processor_configs takes self.lock and then calls add_processor_config for each config, which takes the same lock again. self.lock was a threading.Lock, which is not reentrant, so the first call to the exported set_tracing_processor_configs() blocked forever: on an ACP server it never finishes startup, on a worker it never reaches worker.run(). Use an RLock so the batch still registers under one lock acquisition and the nested add can re-enter it. add_processor_config is unchanged. The new test registers two configs from a thread and fails on the old lock (the thread is still blocked after 5 s); it passes with the RLock. --- .../core/tracing/tracing_processor_manager.py | 4 +- .../tracing/test_tracing_processor_manager.py | 51 +++++++++++++++++++ 2 files changed, 53 insertions(+), 2 deletions(-) create mode 100644 tests/lib/core/tracing/test_tracing_processor_manager.py diff --git a/src/agentex/lib/core/tracing/tracing_processor_manager.py b/src/agentex/lib/core/tracing/tracing_processor_manager.py index 5227e891c..a1b6a3a37 100644 --- a/src/agentex/lib/core/tracing/tracing_processor_manager.py +++ b/src/agentex/lib/core/tracing/tracing_processor_manager.py @@ -4,7 +4,7 @@ import logging import threading from typing import TYPE_CHECKING -from threading import Lock +from threading import RLock from agentex.lib.types.tracing import TracingProcessorConfig from agentex.lib.core.tracing.processors.sgp_tracing_processor import ( @@ -36,7 +36,7 @@ def __init__(self): # Cache for processors self.sync_processors: list[SyncTracingProcessor] = [] self.async_processors: list[AsyncTracingProcessor] = [] - self.lock = Lock() + self.lock = RLock() self._agentex_registered = False def _ensure_agentex_registered(self): diff --git a/tests/lib/core/tracing/test_tracing_processor_manager.py b/tests/lib/core/tracing/test_tracing_processor_manager.py new file mode 100644 index 000000000..cd7406476 --- /dev/null +++ b/tests/lib/core/tracing/test_tracing_processor_manager.py @@ -0,0 +1,51 @@ +from __future__ import annotations + +import threading +from typing import Any, cast +from dataclasses import dataclass + +from agentex.lib.core.tracing.tracing_processor_manager import TracingProcessorManager + + +@dataclass +class _FakeConfig: + type: str = "fake" + + +class _FakeProcessor: + def __init__(self, config: Any) -> None: + self.config = config + + +def _manager() -> TracingProcessorManager: + manager = TracingProcessorManager() + manager.sync_config_registry["fake"] = cast(Any, _FakeProcessor) + manager.async_config_registry["fake"] = cast(Any, _FakeProcessor) + return manager + + +def _finishes(target: Any, timeout: float = 5.0) -> bool: + thread = threading.Thread(target=target, daemon=True) + thread.start() + thread.join(timeout) + return not thread.is_alive() + + +def test_set_processor_configs_registers_every_config_without_deadlocking() -> None: + manager = _manager() + configs = [_FakeConfig(), _FakeConfig()] + + finished = _finishes(lambda: manager.set_processor_configs(cast(Any, configs))) + + assert finished, "set_processor_configs blocked on the manager's own lock" + assert [cast(Any, p).config for p in manager.get_sync_processors()] == configs + assert [cast(Any, p).config for p in manager.get_async_processors()] == configs + + +def test_add_processor_config_still_registers_one_pair() -> None: + manager = _manager() + + manager.add_processor_config(cast(Any, _FakeConfig())) + + assert len(manager.get_sync_processors()) == 1 + assert len(manager.get_async_processors()) == 1 From ab7f15e791c0973b42ce86ca4558480cf4f44d96 Mon Sep 17 00:00:00 2001 From: Michael Xu Date: Fri, 2 Oct 2026 15:49:48 -0500 Subject: [PATCH 2/2] test(tracing): name the registration deadline and surface worker errors Store the 5 s join deadline as REGISTRATION_DEADLINE_SECONDS, and re-raise an exception from the registration thread after the join so a failing registration reports its own error instead of a later list assertion. --- .../tracing/test_tracing_processor_manager.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/tests/lib/core/tracing/test_tracing_processor_manager.py b/tests/lib/core/tracing/test_tracing_processor_manager.py index cd7406476..d17ac6581 100644 --- a/tests/lib/core/tracing/test_tracing_processor_manager.py +++ b/tests/lib/core/tracing/test_tracing_processor_manager.py @@ -24,10 +24,23 @@ def _manager() -> TracingProcessorManager: return manager -def _finishes(target: Any, timeout: float = 5.0) -> bool: - thread = threading.Thread(target=target, daemon=True) +REGISTRATION_DEADLINE_SECONDS = 5.0 + + +def _finishes(target: Any, timeout: float = REGISTRATION_DEADLINE_SECONDS) -> bool: + errors: list[BaseException] = [] + + def _run() -> None: + try: + target() + except BaseException as exc: + errors.append(exc) + + thread = threading.Thread(target=_run, daemon=True) thread.start() thread.join(timeout) + if errors: + raise errors[0] return not thread.is_alive()