Skip to content

Commit 65349c2

Browse files
usmanabbas7claude
andcommitted
wf(wf-52ref): review R1 — 2 issues fixed (self-join guard, worker_crashed coverage)
- stop() skips self-join when called from worker thread (deadlock guard) - remove pragma; test refresh.worker_crashed guard + re-entrant stop Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent ab93dcc commit 65349c2

2 files changed

Lines changed: 69 additions & 3 deletions

File tree

‎src/convert_sdk/config_loader/refresh.py‎

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -125,11 +125,22 @@ def start(self) -> None:
125125
thread.start()
126126

127127
def stop(self, *, timeout: float = 5.0) -> None:
128-
"""Signal the worker to stop and join the thread (idempotent)."""
128+
"""Signal the worker to stop and join the thread (idempotent).
129+
130+
Safe to call re-entrantly from the worker thread itself (e.g. if a host's
131+
terminal-failure callback calls ``Core.close()``): a thread cannot join
132+
itself, so the self-join is skipped — the daemon loop observes the
133+
``_stopping`` flag and exits on its own. The thread is a daemon, so even
134+
an un-joined worker never blocks interpreter exit.
135+
"""
129136
self._stopping.set()
130137
self._wake.set()
131138
thread = self._thread
132-
if thread is not None and thread.is_alive():
139+
if (
140+
thread is not None
141+
and thread.is_alive()
142+
and thread is not threading.current_thread()
143+
):
133144
thread.join(timeout=timeout)
134145
self._thread = None
135146

@@ -177,7 +188,7 @@ def _run(self) -> None:
177188
break
178189
try:
179190
self._do_refresh()
180-
except Exception: # pragma: no cover - defended; see Task 3 guard
191+
except Exception:
181192
# Outer resilience guard. _do_refresh handles its own transport
182193
# failures (Task 3); reaching here means a logging/clock/callback
183194
# subsystem failure escaped. Emit and keep the worker alive.

‎tests/test_config_refresh.py‎

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -353,6 +353,61 @@ def test_failure_does_not_crash_core_host(self) -> None:
353353
core.close()
354354

355355

356+
class TestWorkerResilience:
357+
def test_swap_callback_exception_does_not_kill_worker(self) -> None:
358+
"""An error escaping _do_refresh (e.g. swap callback) is caught and the
359+
worker survives (refresh.worker_crashed guard)."""
360+
from convert_sdk.config_loader.refresh import ConfigRefresher
361+
362+
def _boom(_snapshot: Any) -> None:
363+
raise RuntimeError("swap blew up")
364+
365+
transport = _FakeTransport([_config_payload(), _config_payload()])
366+
refresher = ConfigRefresher(
367+
config=_remote_config(RefreshConfig(interval_seconds=300)),
368+
transport=transport,
369+
on_snapshot=_boom,
370+
)
371+
refresher.start()
372+
try:
373+
refresher.trigger_now()
374+
assert refresher.wait_for_next_refresh(timeout=5.0)
375+
# Worker absorbed the callback failure and is still running.
376+
assert refresher.is_alive()
377+
# And it can still complete a subsequent cycle.
378+
refresher.trigger_now()
379+
assert refresher.wait_for_next_refresh(timeout=5.0)
380+
finally:
381+
refresher.stop()
382+
383+
def test_stop_is_safe_from_worker_thread(self) -> None:
384+
"""Calling stop() re-entrantly from the worker thread must not deadlock."""
385+
from convert_sdk.config_loader.refresh import ConfigRefresher
386+
from convert_sdk.errors import ConfigLoadError
387+
388+
policy = RefreshConfig(interval_seconds=300, backoff_max_seconds=300)
389+
holder: Dict[str, Any] = {}
390+
391+
def _terminal(_exc: BaseException) -> None:
392+
# Re-entrant stop from inside the worker thread.
393+
holder["refresher"].stop()
394+
395+
transport = _FakeTransport([ConfigLoadError("boom")])
396+
refresher = ConfigRefresher(
397+
config=_remote_config(policy),
398+
transport=transport,
399+
on_snapshot=lambda _s: None,
400+
on_terminal_failure=_terminal,
401+
)
402+
holder["refresher"] = refresher
403+
refresher.start()
404+
refresher.trigger_now()
405+
# If self-join were attempted this would hang; the seam must resolve.
406+
assert refresher.wait_for_next_refresh(timeout=5.0)
407+
refresher.stop()
408+
assert not refresher.is_alive()
409+
410+
356411
class TestConfigUpdatedEvent:
357412
def test_config_updated_emitted_on_successful_swap(self) -> None:
358413
from convert_sdk import LifecycleEvent

0 commit comments

Comments
 (0)