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.
This commit is contained in:
+154
-121
@@ -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(
|
||||
|
||||
+30
-16
@@ -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("<II", len(meta_bytes), len(payload)) + meta_bytes + payload
|
||||
checksum = blake3(packet).digest()
|
||||
sys.stdout.buffer.write(checksum)
|
||||
sys.stdout.buffer.write(packet)
|
||||
sys.stdout.buffer.flush()
|
||||
out.write(checksum)
|
||||
out.write(packet)
|
||||
out.flush()
|
||||
|
||||
|
||||
def dispatch(path, quality, maxsize, maxzoom, data=None):
|
||||
@@ -329,18 +336,18 @@ def process_image_buffer(data: bytes, *, quality, maxsize, maxzoom):
|
||||
|
||||
def process_pdf(path, *, maxsize, maxzoom, quality, page_number=0):
|
||||
t_load_start = perf_counter()
|
||||
pdf = pymupdf.open(path)
|
||||
page = pdf.load_page(page_number)
|
||||
w, h = page.rect[2:4]
|
||||
zoom = min(maxsize / w, maxsize / h, maxzoom)
|
||||
mat = pymupdf.Matrix(zoom, zoom)
|
||||
pix = page.get_pixmap(matrix=mat)
|
||||
t_load_end = perf_counter()
|
||||
with pymupdf.open(path) as pdf:
|
||||
page = pdf.load_page(page_number)
|
||||
w, h = page.rect[2:4]
|
||||
zoom = min(maxsize / w, maxsize / h, maxzoom)
|
||||
mat = pymupdf.Matrix(zoom, zoom)
|
||||
pix = page.get_pixmap(matrix=mat)
|
||||
t_load_end = perf_counter()
|
||||
|
||||
t_save_start = perf_counter()
|
||||
img = pyvips.Image.new_from_memory(
|
||||
pix.samples_mv, pix.width, pix.height, pix.n, "uchar"
|
||||
)
|
||||
t_save_start = perf_counter()
|
||||
img = pyvips.Image.new_from_memory(
|
||||
pix.samples_mv, pix.width, pix.height, pix.n, "uchar"
|
||||
)
|
||||
ret = img.write_to_buffer(".avif", Q=quality, effort=AVIF_FAST_EFFORT, strip=True)
|
||||
backend = "pdf+pyvips"
|
||||
t_save_end = perf_counter()
|
||||
@@ -557,10 +564,17 @@ def main() -> 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()
|
||||
|
||||
|
||||
|
||||
+1
-1
@@ -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]
|
||||
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user