Compare commits

...
2 Commits
Author SHA1 Message Date
LeoVasanko 7b1c6f6772 Fix preview pool permanently wedging after worker stderr pipe fills
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.
2026-08-11 04:57:20 +00:00
LeoVasanko c025e7af95 Fix deprecation warning from fitz being renamed to pymupdf: update import. 2026-08-11 02:23:48 +00:00
4 changed files with 361 additions and 139 deletions
+79 -46
View File
@@ -77,6 +77,9 @@ _preview_cache = PreviewCache(capacity=500)
PREVIEW_TIMEOUT = 10.0 # seconds until preview subprocess is killed PREVIEW_TIMEOUT = 10.0 # seconds until preview subprocess is killed
PREVIEW_WORKERS = max(2, min(8, cpu_count())) 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() _active_procs: set[asyncio.subprocess.Process] = set()
_preview_pool = None _preview_pool = None
_preview_pool_lock = asyncio.Lock() _preview_pool_lock = asyncio.Lock()
@@ -144,13 +147,24 @@ class _PreviewWorker:
return payload or None, resp return payload or None, resp
async def kill(self) -> None: async def kill(self) -> None:
try:
if self.proc.returncode is None: if self.proc.returncode is None:
# Safe to hard-kill: the worker is stateless per request, and its # Safe to hard-kill: the worker is stateless per request.
# subprocesses (ffmpeg) use stdin=DEVNULL so they never hold the # proc.wait() must not be awaited unaided: if a pipe
# worker's pipes open — proc.wait() cannot hang on pipe EOF. # 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): with contextlib.suppress(ProcessLookupError):
self.proc.kill() self.proc.kill()
await self.proc.wait() 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) _active_procs.discard(self.proc)
@@ -166,22 +180,19 @@ class _PreviewWorkerPool:
self._seq = 0 self._seq = 0
self._closed = False 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: 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( proc = await asyncio.create_subprocess_exec(
sys.executable, sys.executable,
"-m", "-m",
"cista.preview_worker", "cista.preview_worker",
stdin=asyncio.subprocess.PIPE, stdin=asyncio.subprocess.PIPE,
stdout=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE, stderr=None,
) )
_active_procs.add(proc) _active_procs.add(proc)
try: try:
@@ -191,21 +202,14 @@ class _PreviewWorkerPool:
proc.kill() proc.kill()
with contextlib.suppress(Exception): with contextlib.suppress(Exception):
await proc.wait() await proc.wait()
stderr = await self._read_startup_stderr(proc)
if stderr:
raise WorkerProtocolError( raise WorkerProtocolError(
"preview worker failed to become ready: " + stderr.splitlines()[-1] "preview worker failed to become ready"
" (worker stderr goes to the server log)"
) from err ) from err
raise WorkerProtocolError("preview worker failed to become ready") from err
except asyncio.IncompleteReadError as 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( raise WorkerProtocolError(
"preview worker exited before signalling readiness" "preview worker exited before signalling readiness"
" (worker stderr goes to the server log)"
) from err ) from err
if ready != b"\x01": if ready != b"\x01":
raise WorkerProtocolError(f"preview worker ready signal invalid: {ready!r}") raise WorkerProtocolError(f"preview worker ready signal invalid: {ready!r}")
@@ -218,28 +222,48 @@ class _PreviewWorkerPool:
async def _replace_worker(self, worker: _PreviewWorker) -> None: async def _replace_worker(self, worker: _PreviewWorker) -> None:
self._workers.discard(worker) self._workers.discard(worker)
try:
await worker.kill() await worker.kill()
if self._closed: except Exception:
return 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: try:
await self._add_worker() await self._add_worker()
except Exception: except Exception:
logger.exception("Failed to replace preview worker") logger.exception(
"Failed to replace preview worker (pool %d/%d); retrying in %ds",
async def _dispatch_loop(self) -> None: len(self._workers),
while True: self.size,
try: int(delay),
_priority, _seq, future, args = await self._pending.get() )
except asyncio.CancelledError: await asyncio.sleep(delay)
delay = min(delay * 2, WORKER_RESPAWN_DELAY_MAX)
else:
return return
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:
await self._dispatch_one()
except asyncio.CancelledError:
return
except Exception:
logger.exception("Preview dispatcher error; continuing")
async def _dispatch_one(self) -> None:
_priority, _seq, future, args = await self._pending.get()
if future.cancelled(): if future.cancelled():
continue return
try: try:
worker = await asyncio.wait_for( worker = await asyncio.wait_for(self._idle.get(), timeout=PREVIEW_TIMEOUT)
self._idle.get(), timeout=PREVIEW_TIMEOUT
)
except TimeoutError: except TimeoutError:
logger.warning( logger.warning(
"Preview worker unavailable (%ds) for %s", "Preview worker unavailable (%ds) for %s",
@@ -253,7 +277,7 @@ class _PreviewWorkerPool:
backend=_expected_preview_backend(args[0]), backend=_expected_preview_backend(args[0]),
) )
) )
continue return
filepath = args[0] filepath = args[0]
replace = False replace = False
@@ -266,6 +290,12 @@ class _PreviewWorkerPool:
future.set_result((out, resp)) future.set_result((out, resp))
except TimeoutError: except TimeoutError:
replace = True 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(): if not future.done():
future.set_exception( future.set_exception(
PreviewTimeoutError( PreviewTimeoutError(
@@ -275,7 +305,11 @@ class _PreviewWorkerPool:
) )
except WorkerChecksumError: except WorkerChecksumError:
replace = True replace = True
logger.error("Preview checksum mismatch for %s", filepath.name) logger.error(
"Preview checksum mismatch for %s (worker pid=%s); replacing it",
filepath.name,
worker.proc.pid,
)
if not future.done(): if not future.done():
future.set_exception( future.set_exception(
PreviewError(f"worker checksum mismatch for {filepath.name}") PreviewError(f"worker checksum mismatch for {filepath.name}")
@@ -290,23 +324,22 @@ class _PreviewWorkerPool:
ConnectionResetError, ConnectionResetError,
OSError, OSError,
ValueError, ValueError,
msgspec.json.DecodeError, msgspec.DecodeError,
) as e: ) as e:
replace = True replace = True
logger.warning( logger.warning(
"Preview worker protocol failure for %s: %s", filepath.name, e "Preview worker pid=%s protocol failure for %s: %s",
worker.proc.pid,
filepath.name,
e,
) )
if not future.done(): if not future.done():
future.set_exception( future.set_exception(
PreviewError( PreviewError(f"worker protocol failure for {filepath.name}: {e}")
f"worker protocol failure for {filepath.name}: {e}"
)
) )
except Exception: except Exception:
replace = True replace = True
logger.exception( logger.exception("Unexpected preview worker error for %s", filepath.name)
"Unexpected preview worker error for %s", filepath.name
)
if not future.done(): if not future.done():
future.set_exception( future.set_exception(
PreviewError(f"unexpected worker error for {filepath.name}") PreviewError(f"unexpected worker error for {filepath.name}")
+22 -8
View File
@@ -26,9 +26,9 @@ from pathlib import Path
from time import perf_counter from time import perf_counter
import av import av
import fitz # PyMuPDF
import msgspec import msgspec
import numpy as np import numpy as np
import pymupdf
import pyvips import pyvips
from blake3 import blake3 from blake3 import blake3
@@ -128,13 +128,20 @@ def _read_request() -> tuple[PreviewRequest, bytes] | None:
return req, data 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: 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) meta_bytes = _enc.encode(resp)
packet = struct.pack("<II", len(meta_bytes), len(payload)) + meta_bytes + payload packet = struct.pack("<II", len(meta_bytes), len(payload)) + meta_bytes + payload
checksum = blake3(packet).digest() checksum = blake3(packet).digest()
sys.stdout.buffer.write(checksum) out.write(checksum)
sys.stdout.buffer.write(packet) out.write(packet)
sys.stdout.buffer.flush() out.flush()
def dispatch(path, quality, maxsize, maxzoom, data=None): def dispatch(path, quality, maxsize, maxzoom, data=None):
@@ -329,11 +336,11 @@ def process_image_buffer(data: bytes, *, quality, maxsize, maxzoom):
def process_pdf(path, *, maxsize, maxzoom, quality, page_number=0): def process_pdf(path, *, maxsize, maxzoom, quality, page_number=0):
t_load_start = perf_counter() t_load_start = perf_counter()
pdf = fitz.open(path) with pymupdf.open(path) as pdf:
page = pdf.load_page(page_number) page = pdf.load_page(page_number)
w, h = page.rect[2:4] w, h = page.rect[2:4]
zoom = min(maxsize / w, maxsize / h, maxzoom) zoom = min(maxsize / w, maxsize / h, maxzoom)
mat = fitz.Matrix(zoom, zoom) mat = pymupdf.Matrix(zoom, zoom)
pix = page.get_pixmap(matrix=mat) pix = page.get_pixmap(matrix=mat)
t_load_end = perf_counter() t_load_end = perf_counter()
@@ -557,10 +564,17 @@ def main() -> None:
if len(sys.argv) > 1: if len(sys.argv) > 1:
_run_once() _run_once()
return 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 # Eagerly import heavy modules before signalling readiness so the parent
# does not hand us a request while we are still initialising. # does not hand us a request while we are still initialising.
sys.stdout.buffer.write(b"\x01") _protocol_out.write(b"\x01")
sys.stdout.buffer.flush() _protocol_out.flush()
_run_loop() _run_loop()
+1 -1
View File
@@ -161,7 +161,7 @@ ignore = [
"TRY003", # exception-message strictness too noisy on legacy handlers "TRY003", # exception-message strictness too noisy on legacy handlers
] ]
isort.known-first-party = ["cista"] 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"] per-file-ignores."scripts/*" = ["T20"]
[dependency-groups] [dependency-groups]
+175
View File
@@ -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("<II", header)
meta = read_exact(meta_len)
read_exact(payload_len)
req = json.loads(meta)
if req["path"].endswith(".block"):
# Simulate a worker stuck on an undrained stderr pipe:
# flood stderr past the OS pipe buffer, then never respond.
import os
import time
try:
os.write(2, b"x" * 10_000_000)
except OSError:
pass
while True:
time.sleep(3600)
resp = json.dumps({"ok": True, "mime": "image/avif", "backend": "fake"}).encode()
payload = b"FAKEIMG"
packet = struct.pack("<II", len(resp), len(payload)) + resp + payload
sys.stdout.buffer.write(blake3.blake3(packet).digest())
sys.stdout.buffer.write(packet)
sys.stdout.buffer.flush()
"""
)
@pytest.mark.asyncio
async def test_pool_recovers_from_wedged_worker(monkeypatch, tmp_path):
"""A worker wedged on an undrained stderr pipe must not kill the pool.
Recreates the old production setup (stderr=PIPE, never drained) and
verifies the request times out, the stuck worker's kill() cannot hang
the dispatcher, and the pool serves the next request normally.
"""
monkeypatch.setattr(preview, "PREVIEW_TIMEOUT", 1.0)
monkeypatch.setattr(preview, "WORKER_KILL_GRACE", 0.5)
monkeypatch.setattr(preview, "WORKER_RESPAWN_DELAY", 0.05)
script = tmp_path / "fake_worker.py"
script.write_text(FAKE_WORKER)
async def fake_spawn(self):
proc = await asyncio.create_subprocess_exec(
sys.executable,
str(script),
stdin=asyncio.subprocess.PIPE,
stdout=asyncio.subprocess.PIPE,
# Deliberately piped-and-undrained, recreating the old
# production setup that wedges a worker on stderr writes.
stderr=asyncio.subprocess.PIPE,
)
preview._active_procs.add(proc)
await asyncio.wait_for(proc.stdout.readexactly(1), timeout=10)
return preview._PreviewWorker(proc)
monkeypatch.setattr(preview._PreviewWorkerPool, "_spawn_worker", fake_spawn)
pool = preview._PreviewWorkerPool(1)
await pool.start()
try:
with pytest.raises(preview.PreviewTimeoutError):
await pool.run(Path("wedged.block"), 60, 512, 2.0)
out, resp = await asyncio.wait_for(
pool.run(Path("ok.jpg"), 60, 512, 2.0), timeout=10
)
assert out == b"FAKEIMG"
assert resp.ok
assert all(not task.done() for task in pool._dispatchers)
finally:
await pool.close()
@pytest.mark.asyncio
async def test_worker_kill_grace_when_wait_hangs(monkeypatch):
"""kill() must return even if asyncio never resolves proc.wait()."""
monkeypatch.setattr(preview, "WORKER_KILL_GRACE", 0.1)
proc = Mock()
proc.returncode = None
proc.pid = 1234
never = asyncio.Future()
async def wait():
await never
proc.wait = wait
worker = preview._PreviewWorker(proc)
preview._active_procs.add(proc)
start = time.monotonic()
await worker.kill()
assert time.monotonic() - start < 2
assert proc not in preview._active_procs
never.cancel()
@pytest.mark.asyncio
async def test_replace_worker_retries_failed_spawn(monkeypatch):
"""A failed replacement spawn must be retried, not silently dropped."""
monkeypatch.setattr(preview, "WORKER_RESPAWN_DELAY", 0.01)
pool = preview._PreviewWorkerPool(1)
old_worker = Mock()
old_worker.proc = Mock(pid=4321)
old_worker.kill = AsyncMock()
attempts = 0
async def add_worker():
nonlocal attempts
attempts += 1
if attempts < 3:
raise OSError("too many open files")
pool._add_worker = add_worker
await pool._replace_worker(old_worker)
assert attempts == 3
@pytest.mark.asyncio
async def test_dispatch_loop_survives_body_errors():
"""Exceptions escaping a dispatch cycle must not kill the dispatcher."""
pool = preview._PreviewWorkerPool(1)
calls = 0
async def dispatch_one():
nonlocal calls
calls += 1
if calls == 1:
raise RuntimeError("boom")
raise asyncio.CancelledError
pool._dispatch_one = dispatch_one
await pool._dispatch_loop()
assert calls == 2