From 7b1c6f677221162f7ed21bb2762f23ff9729f24f Mon Sep 17 00:00:00 2001 From: Leo Vasanko Date: Tue, 11 Aug 2026 04:57:20 +0000 Subject: [PATCH] Fix preview pool permanently wedging after worker stderr pipe fills MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Production symptom: previews of all types (pdf/pyvips/onlyoffice) start hitting the 10s timeout and never recover until server restart, while the rest of the server stays healthy. Root cause (reproduced on Python 3.12): workers were spawned with stderr=PIPE that nothing drained after startup. Once the OS pipe buffer filled from accumulated worker tracebacks and library warnings, asyncio flow control stopped the parent reading it and the worker blocked forever mid-request on a stderr write. The 10s timeout then fired, but _replace_worker hung forever in proc.wait() even after kill() — the flow-control-paused pipe transport never sees EOF — permanently wedging one dispatcher per stuck worker. Once all dispatchers were stuck, every preview request timed out. Restart cleared it. Fixes: - Spawn workers with inherited stderr (stderr=None) so worker diagnostics go straight to the server log and no undrained pipe can exist. - Bound proc.wait() in worker kill() with a 5s grace timeout so a wedged transport can never hang a dispatcher; log the worker pid instead. - Guard the dispatch loop with an outer exception handler so a dispatcher can never die silently and shrink pool capacity. - Retry failed worker replacement spawns with 1s-30s backoff instead of silently shrinking the pool. - Fix latent crash: except-tuple referenced msgspec.json.DecodeError, which does not exist in the installed msgspec; any protocol failure would itself raise AttributeError. Use msgspec.DecodeError. - Worker: redirect Python-level sys.stdout to stderr in persistent mode and keep the raw buffer solely for the binary protocol, so a library print() can never corrupt the command channel again (cf. the pymupdf deprecation warning that crashed workers at startup). - Worker: close the pymupdf document explicitly in process_pdf. - Log worker pid on timeout/protocol/checksum failures, and log failed kills and replacement retries, for future production diagnostics. Add tests/test_preview_pool.py with an end-to-end regression recreating the wedged-worker setup (piped, undrained stderr) plus kill-grace, respawn-retry and dispatcher-survival tests. --- cista/preview.py | 275 +++++++++++++++++++++---------------- cista/preview_worker.py | 46 ++++--- pyproject.toml | 2 +- tests/test_preview_pool.py | 175 +++++++++++++++++++++++ 4 files changed, 360 insertions(+), 138 deletions(-) create mode 100644 tests/test_preview_pool.py diff --git a/cista/preview.py b/cista/preview.py index 8c92335..a17da9f 100644 --- a/cista/preview.py +++ b/cista/preview.py @@ -77,6 +77,9 @@ _preview_cache = PreviewCache(capacity=500) PREVIEW_TIMEOUT = 10.0 # seconds until preview subprocess is killed PREVIEW_WORKERS = max(2, min(8, cpu_count())) +WORKER_KILL_GRACE = 5.0 # max seconds to wait for a killed worker to be reaped +WORKER_RESPAWN_DELAY = 1.0 # initial delay before retrying a failed worker spawn +WORKER_RESPAWN_DELAY_MAX = 30.0 _active_procs: set[asyncio.subprocess.Process] = set() _preview_pool = None _preview_pool_lock = asyncio.Lock() @@ -144,14 +147,25 @@ class _PreviewWorker: return payload or None, resp async def kill(self) -> None: - if self.proc.returncode is None: - # Safe to hard-kill: the worker is stateless per request, and its - # subprocesses (ffmpeg) use stdin=DEVNULL so they never hold the - # worker's pipes open — proc.wait() cannot hang on pipe EOF. - with contextlib.suppress(ProcessLookupError): - self.proc.kill() - await self.proc.wait() - _active_procs.discard(self.proc) + try: + if self.proc.returncode is None: + # Safe to hard-kill: the worker is stateless per request. + # proc.wait() must not be awaited unaided: if a pipe + # transport is flow-control paused (e.g. an undrained stderr + # pipe), asyncio may never resolve wait() even after SIGKILL, + # which would permanently wedge the calling dispatcher. + with contextlib.suppress(ProcessLookupError): + self.proc.kill() + try: + await asyncio.wait_for(self.proc.wait(), timeout=WORKER_KILL_GRACE) + except TimeoutError: + logger.error( + "Preview worker pid=%s not reaped within %ds of kill", + self.proc.pid, + int(WORKER_KILL_GRACE), + ) + finally: + _active_procs.discard(self.proc) class _PreviewWorkerPool: @@ -166,22 +180,19 @@ class _PreviewWorkerPool: self._seq = 0 self._closed = False - async def _read_startup_stderr(self, proc: asyncio.subprocess.Process) -> str: - if proc.stderr is None: - return "" - with contextlib.suppress(TimeoutError): - data = await asyncio.wait_for(proc.stderr.read(), timeout=0.5) - return data.decode(errors="replace").strip() - return "" - async def _spawn_worker(self) -> _PreviewWorker: + # stderr is inherited, not piped: a piped stderr that nobody drains + # eventually fills its OS buffer, blocking the worker mid-request, + # and its flow-control-paused transport makes proc.wait() hang even + # after kill() — together this used to permanently wedge the pool. + # Inheriting sends worker diagnostics straight to the server log. proc = await asyncio.create_subprocess_exec( sys.executable, "-m", "cista.preview_worker", stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.PIPE, + stderr=None, ) _active_procs.add(proc) try: @@ -191,21 +202,14 @@ class _PreviewWorkerPool: proc.kill() with contextlib.suppress(Exception): await proc.wait() - stderr = await self._read_startup_stderr(proc) - if stderr: - raise WorkerProtocolError( - "preview worker failed to become ready: " + stderr.splitlines()[-1] - ) from err - raise WorkerProtocolError("preview worker failed to become ready") from err + raise WorkerProtocolError( + "preview worker failed to become ready" + " (worker stderr goes to the server log)" + ) from err except asyncio.IncompleteReadError as err: - stderr = await self._read_startup_stderr(proc) - if stderr: - raise WorkerProtocolError( - "preview worker exited before signalling readiness: " - + stderr.splitlines()[-1] - ) from err raise WorkerProtocolError( "preview worker exited before signalling readiness" + " (worker stderr goes to the server log)" ) from err if ready != b"\x01": raise WorkerProtocolError(f"preview worker ready signal invalid: {ready!r}") @@ -218,106 +222,135 @@ class _PreviewWorkerPool: async def _replace_worker(self, worker: _PreviewWorker) -> None: self._workers.discard(worker) - await worker.kill() - if self._closed: - return try: - await self._add_worker() + await worker.kill() except Exception: - logger.exception("Failed to replace preview worker") - - async def _dispatch_loop(self) -> None: - while True: + logger.exception("Failed to kill preview worker pid=%s", worker.proc.pid) + # Keep retrying until a replacement is up: a pool that silently + # shrinks degrades all preview traffic to timeouts. + delay = WORKER_RESPAWN_DELAY + while not self._closed: try: - _priority, _seq, future, args = await self._pending.get() - except asyncio.CancelledError: + await self._add_worker() + except Exception: + logger.exception( + "Failed to replace preview worker (pool %d/%d); retrying in %ds", + len(self._workers), + self.size, + int(delay), + ) + await asyncio.sleep(delay) + delay = min(delay * 2, WORKER_RESPAWN_DELAY_MAX) + else: return - if future.cancelled(): - continue - + async def _dispatch_loop(self) -> None: + # Nothing may escape the loop body: a dispatcher that dies silently + # permanently shrinks pool capacity and degrades all preview + # traffic to timeouts. + while True: try: - worker = await asyncio.wait_for( - self._idle.get(), timeout=PREVIEW_TIMEOUT - ) - except TimeoutError: - logger.warning( - "Preview worker unavailable (%ds) for %s", - int(PREVIEW_TIMEOUT), - args[0].name, - ) - if not future.done(): - future.set_exception( - PreviewTimeoutError( - args[0].name, - backend=_expected_preview_backend(args[0]), - ) - ) - continue - - filepath = args[0] - replace = False - try: - out, resp = await asyncio.wait_for( - worker.request(*args), - timeout=PREVIEW_TIMEOUT, - ) - if not future.done(): - future.set_result((out, resp)) - except TimeoutError: - replace = True - if not future.done(): - future.set_exception( - PreviewTimeoutError( - filepath.name, - backend=_expected_preview_backend(filepath), - ) - ) - except WorkerChecksumError: - replace = True - logger.error("Preview checksum mismatch for %s", filepath.name) - if not future.done(): - future.set_exception( - PreviewError(f"worker checksum mismatch for {filepath.name}") - ) - except PreviewError as e: - if not future.done(): - future.set_exception(e) - except ( - WorkerProtocolError, - asyncio.IncompleteReadError, - BrokenPipeError, - ConnectionResetError, - OSError, - ValueError, - msgspec.json.DecodeError, - ) as e: - replace = True - logger.warning( - "Preview worker protocol failure for %s: %s", filepath.name, e - ) - if not future.done(): - future.set_exception( - PreviewError( - f"worker protocol failure for {filepath.name}: {e}" - ) - ) + await self._dispatch_one() + except asyncio.CancelledError: + return except Exception: - replace = True - logger.exception( - "Unexpected preview worker error for %s", filepath.name - ) - if not future.done(): - future.set_exception( - PreviewError(f"unexpected worker error for {filepath.name}") + logger.exception("Preview dispatcher error; continuing") + + async def _dispatch_one(self) -> None: + _priority, _seq, future, args = await self._pending.get() + + if future.cancelled(): + return + + try: + worker = await asyncio.wait_for(self._idle.get(), timeout=PREVIEW_TIMEOUT) + except TimeoutError: + logger.warning( + "Preview worker unavailable (%ds) for %s", + int(PREVIEW_TIMEOUT), + args[0].name, + ) + if not future.done(): + future.set_exception( + PreviewTimeoutError( + args[0].name, + backend=_expected_preview_backend(args[0]), ) - finally: - if replace: - await self._replace_worker(worker) - elif worker.proc.returncode is None: - await self._idle.put(worker) - else: - await self._replace_worker(worker) + ) + return + + filepath = args[0] + replace = False + try: + out, resp = await asyncio.wait_for( + worker.request(*args), + timeout=PREVIEW_TIMEOUT, + ) + if not future.done(): + future.set_result((out, resp)) + except TimeoutError: + replace = True + logger.warning( + "Preview worker pid=%s timed out (%ds) on %s; replacing it", + worker.proc.pid, + int(PREVIEW_TIMEOUT), + filepath.name, + ) + if not future.done(): + future.set_exception( + PreviewTimeoutError( + filepath.name, + backend=_expected_preview_backend(filepath), + ) + ) + except WorkerChecksumError: + replace = True + logger.error( + "Preview checksum mismatch for %s (worker pid=%s); replacing it", + filepath.name, + worker.proc.pid, + ) + if not future.done(): + future.set_exception( + PreviewError(f"worker checksum mismatch for {filepath.name}") + ) + except PreviewError as e: + if not future.done(): + future.set_exception(e) + except ( + WorkerProtocolError, + asyncio.IncompleteReadError, + BrokenPipeError, + ConnectionResetError, + OSError, + ValueError, + msgspec.DecodeError, + ) as e: + replace = True + logger.warning( + "Preview worker pid=%s protocol failure for %s: %s", + worker.proc.pid, + filepath.name, + e, + ) + if not future.done(): + future.set_exception( + PreviewError(f"worker protocol failure for {filepath.name}: {e}") + ) + except Exception: + replace = True + logger.exception("Unexpected preview worker error for %s", filepath.name) + if not future.done(): + future.set_exception( + PreviewError(f"unexpected worker error for {filepath.name}") + ) + finally: + if replace: + await self._replace_worker(worker) + elif worker.proc.returncode is None: + await self._idle.put(worker) + else: + await self._replace_worker(worker) async def start(self) -> None: workers = await asyncio.gather( diff --git a/cista/preview_worker.py b/cista/preview_worker.py index 94c407d..c4f66e7 100644 --- a/cista/preview_worker.py +++ b/cista/preview_worker.py @@ -128,13 +128,20 @@ def _read_request() -> tuple[PreviewRequest, bytes] | None: return req, data +# Raw stdout buffer reserved for the binary protocol once main() redirects +# Python-level stdout to stderr. None means "use sys.stdout.buffer as-is" +# (CLI single-shot mode, where real stdout is wanted). +_protocol_out = None + + def _write_response(resp: PreviewResponse, payload: bytes) -> None: + out = _protocol_out if _protocol_out is not None else sys.stdout.buffer meta_bytes = _enc.encode(resp) packet = struct.pack(" None: if len(sys.argv) > 1: _run_once() return + # The command channel is a binary protocol on fd 1. Anything printed to + # stdout by Python code (e.g. a library emitting a warning via print()) + # would corrupt the protocol, so redirect Python-level stdout to stderr + # (the server log) and keep the raw buffer solely for protocol traffic. + global _protocol_out + _protocol_out = sys.stdout.buffer + sys.stdout = sys.stderr # Eagerly import heavy modules before signalling readiness so the parent # does not hand us a request while we are still initialising. - sys.stdout.buffer.write(b"\x01") - sys.stdout.buffer.flush() + _protocol_out.write(b"\x01") + _protocol_out.flush() _run_loop() diff --git a/pyproject.toml b/pyproject.toml index 4b75858..2cb296d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -161,7 +161,7 @@ ignore = [ "TRY003", # exception-message strictness too noisy on legacy handlers ] isort.known-first-party = ["cista"] -per-file-ignores."tests/*" = ["S", "ANN", "D", "INP", "PLR2004", "ARG001"] +per-file-ignores."tests/*" = ["S", "ANN", "D", "INP", "PLR2004", "ARG001", "SLF001"] per-file-ignores."scripts/*" = ["T20"] [dependency-groups] diff --git a/tests/test_preview_pool.py b/tests/test_preview_pool.py new file mode 100644 index 0000000..2c32836 --- /dev/null +++ b/tests/test_preview_pool.py @@ -0,0 +1,175 @@ +"""Tests for the preview worker pool resilience. + +Regression context: a piped worker stderr that nobody drains used to block +the worker mid-request once the OS pipe buffer filled, and asyncio's +proc.wait() then never resolved even after kill() — wedging one dispatcher +per stuck worker until all preview traffic timed out permanently. +""" + +import asyncio +import sys +import textwrap +import time +from pathlib import Path +from unittest.mock import AsyncMock, Mock + +import pytest + +from cista import preview + +FAKE_WORKER = textwrap.dedent( + """ + import json + import struct + import sys + + import blake3 + + + def read_exact(n): + buf = b"" + while len(buf) < n: + chunk = sys.stdin.buffer.read(n - len(buf)) + if not chunk: + raise EOFError + buf += chunk + return buf + + + sys.stdout.buffer.write(b"\\x01") + sys.stdout.buffer.flush() + while True: + header = sys.stdin.buffer.read(8) + if not header or len(header) < 8: + break + meta_len, payload_len = struct.unpack("