Compare commits

..
9 Commits
Author SHA1 Message Date
LeoVasanko 953ec628a0 Add -nostdin to ffmpeg preview conversions to suppress keyboard prompts 2026-08-11 05:32:27 +00:00
LeoVasanko e678c8c267 Put the failing ffmpeg command on its own line in error messages 2026-08-11 05:31:30 +00:00
LeoVasanko 69d58f99e3 Drop noisy ffmpeg fallback for non-HEIC images, quiet ffmpeg output
A corrupt TIFF in production produced a wall of ffmpeg error output:
pyvips could not decode it, the generic ffmpeg fallback was tried, and
ffmpeg's TIFF decoder failed just the same — with banner, configuration
and stream-mapping spam included.

- Non-HEIC images are now decoded by pyvips only; a pyvips failure
  raises a clean one-line ValueError ("cannot decode image: ...", a 422
  like any other undecodable file) instead of invoking ffmpeg. The
  ffmpeg path is kept for HEIC/HEIF, where pyvips genuinely falls short
  (tile assembly, HDR metadata).
- ffmpeg runs with -hide_banner -loglevel error -nostats: error output
  is still shown on failure, without the version/configuration/progress
  noise. The -s insertion index no longer depends on fixed positions.
2026-08-11 05:29:32 +00:00
LeoVasanko 36764885ed Silence pyvips deprecation and INFO spam in preview worker
- AVIF saves: replace deprecated strip=True with keep="none" (libvips
  8.15+; production already runs a version that deprecates strip).
- Set the pyvips logger to WARNING in the worker: its INFO messages
  ("threadpool completed ...") are pure spam on every operation.
2026-08-11 05:11:44 +00:00
LeoVasanko 5a82560cf2 Format preview worker logs like the main process, tagged with worker pid
Worker stderr is now inherited by the parent, so its log lines land in the
server log — but they arrived with the default logging format and a noisy
"preview-worker config=..." line at every spawn.

- Extract the emoji level-prefix formatting from cista.sanic_logging into
  cista.util.logformat, which has no Sanic dependency (the worker must not
  import Sanic: import-time prints could corrupt the stdout protocol).
- Worker configures its stderr handler with the same emoji prefixes plus a
  worker[pid] tag, and the config-loaded info message is removed.
2026-08-11 05:05:07 +00:00
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
LeoVasanko 3bad311e35 Tell Paskia SSO not to renew session on WebSocket connections where we cannot pass back the refreshed cookie. 2026-08-11 02:08:01 +00:00
LeoVasanko fdc4fe0a3e Forward client user-agent to SSO backend on validation refreshes. Matches function of existing proxy_auth_request (copies all headers) and proxy_auth_websocket (copies user-agent when present). 2026-08-10 14:08:17 +00:00
8 changed files with 455 additions and 189 deletions
+3 -1
View File
@@ -40,7 +40,9 @@ async def watch(req, ws):
if sso.paskia_enabled():
# SSO auth: call validation to get user info (don't enforce auth in public mode)
try:
await sso.validate_sso_request(req)
# WebSocket cannot forward Set-Cookie, so ask the auth backend not to
# renew the session here; renewal happens on the HTTP side instead.
await sso.validate_sso_request(req, renew=False)
except Exception as e:
logger.debug("watch SSO validation failed: %s", e)
if sso_user := getattr(req.ctx, "sso_user", None):
+154 -121
View File
@@ -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(
+71 -34
View File
@@ -17,6 +17,7 @@ import gc
import io
import logging
import mimetypes
import os
import shlex
import struct
import subprocess
@@ -26,16 +27,26 @@ from pathlib import Path
from time import perf_counter
import av
import fitz # PyMuPDF
import msgspec
import numpy as np
import pymupdf
import pyvips
from blake3 import blake3
from cista import config
from cista.util.logformat import format_level_prefix
logger = logging.getLogger(__name__)
class _WorkerLogFormatter(logging.Formatter):
"""Emoji level prefix like the main process, tagged with the worker pid."""
def format(self, record: logging.LogRecord) -> str:
prefix = format_level_prefix(record.levelno)
return f"{prefix}worker[{os.getpid()}]: {record.getMessage()}"
AVIF_FAST_EFFORT = 0
DOC_PREVIEW_SUFFIXES = {".pdf", ".xps", ".epub", ".mobi"}
@@ -128,13 +139,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):
@@ -195,6 +213,13 @@ def _image_via_ffmpeg(path: Path, maxsize: int, quality: int) -> bytes:
tmp_path = tmp_f.name
cmd = [
"ffmpeg",
# Keep error messages, drop the banner/config/stream-mapping spam.
"-hide_banner",
"-loglevel",
"error",
"-nostats",
# No interactive keyboard prompts ("Press [q] to stop ...").
"-nostdin",
"-y",
"-i",
str(path),
@@ -215,8 +240,9 @@ def _image_via_ffmpeg(path: Path, maxsize: int, quality: int) -> bytes:
new_w = int(w * scale)
new_h = int(h * scale)
# insert -s <wxh> right after the input file
cmd.insert(4, "-s")
cmd.insert(5, f"{new_w}x{new_h}")
input_index = cmd.index(str(path)) + 1
cmd.insert(input_index, "-s")
cmd.insert(input_index + 1, f"{new_w}x{new_h}")
try:
try:
# stdin=DEVNULL is critical: ffmpeg must not inherit the worker's
@@ -236,10 +262,10 @@ def _image_via_ffmpeg(path: Path, maxsize: int, quality: int) -> bytes:
stderr = (e.stderr or b"").decode(errors="replace").strip()
if stderr:
raise RuntimeError(
f"ffmpeg failed (exit {e.returncode}): {shell_cmd}\n{stderr}"
f"ffmpeg failed (exit {e.returncode}):\n{shell_cmd}\n{stderr}"
) from e
raise RuntimeError(
f"ffmpeg failed (exit {e.returncode}): {shell_cmd}"
f"ffmpeg failed (exit {e.returncode}):\n{shell_cmd}"
) from e
with Path(tmp_path).open("rb") as f:
return f.read()
@@ -267,9 +293,10 @@ def process_image_pyvips(path, *, maxsize, quality):
height=height,
)
# Other image formats: pyvips first, ffmpeg fallback.
# Other image formats: pyvips only. ffmpeg is not a useful fallback
# here — when pyvips cannot decode a file, ffmpeg's image decoders
# cannot either, and their failure output is far noisier.
load_opts = {"access": "sequential"}
orig_w = orig_h = None
try:
img = pyvips.Image.new_from_file(str(path), **load_opts)
img = img.autorot()
@@ -281,13 +308,11 @@ def process_image_pyvips(path, *, maxsize, quality):
".avif",
Q=quality,
effort=AVIF_FAST_EFFORT,
strip=True,
keep="none",
)
backend = "pyvips"
except pyvips.error.Error:
orig_w, orig_h = None, None
ret = _image_via_ffmpeg(path, maxsize, quality)
backend = "ffmpeg"
except pyvips.error.Error as e:
raise ValueError(f"cannot decode image: {e}") from e
backend = "pyvips"
t_end = perf_counter()
return ret, PreviewResponse(
@@ -313,7 +338,7 @@ def process_image_buffer(data: bytes, *, quality, maxsize, maxzoom):
".avif",
Q=quality,
effort=AVIF_FAST_EFFORT,
strip=True,
keep="none",
)
t_end = perf_counter()
@@ -329,19 +354,19 @@ 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 = fitz.open(path)
page = pdf.load_page(page_number)
w, h = page.rect[2:4]
zoom = min(maxsize / w, maxsize / h, maxzoom)
mat = fitz.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"
)
ret = img.write_to_buffer(".avif", Q=quality, effort=AVIF_FAST_EFFORT, strip=True)
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, keep="none")
backend = "pdf+pyvips"
t_save_end = perf_counter()
@@ -547,20 +572,32 @@ def _run_loop() -> None:
def main() -> None:
# Configure all log output to stderr before any imports that may emit logs.
logging.basicConfig(stream=sys.stderr, level=logging.INFO)
# Configure all log output to stderr before any imports that may emit
# logs. stderr is inherited by the parent, so this lands in the server
# log, formatted like the main process and tagged with the worker pid.
handler = logging.StreamHandler(sys.stderr)
handler.setFormatter(_WorkerLogFormatter())
logging.basicConfig(level=logging.INFO, handlers=[handler])
# pyvips is chatty at INFO ("threadpool completed ..." per operation).
logging.getLogger("pyvips").setLevel(logging.WARNING)
try:
config.load_config()
logger.info("preview-worker config=%s", config.conffile)
except Exception:
logger.exception("preview-worker failed to load config at startup")
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()
+3 -31
View File
@@ -3,11 +3,13 @@
import logging
import os
import sys
import unicodedata
from ipaddress import IPv6Address
from sanic.log import LOGGING_CONFIG_DEFAULTS
from cista.util.logformat import EmojiFormatter as _EmojiFormatter
from cista.util.logformat import display_width as _display_width
logger = logging.getLogger("cista.access")
@@ -132,14 +134,6 @@ def format_duration_ms(duration_ms: float) -> str:
return f"{hours}h{minutes}m"
def _display_width(text: str) -> int:
return sum(
1 + (unicodedata.east_asian_width(c) in "FW")
for c in text
if unicodedata.category(c) != "Mn"
)
def _format_left(label: str) -> str:
return label[:19].ljust(19)
@@ -279,28 +273,6 @@ def configure_access_logging() -> None:
logger.propagate = False
_LEVEL_EMOJI = {
logging.DEBUG: "🔍",
logging.INFO: "", # noqa: RUF001
logging.WARNING: "⚠️",
logging.ERROR: "🛑",
logging.CRITICAL: "🛑",
}
def _format_level_prefix(levelno: int) -> str:
emoji = _LEVEL_EMOJI.get(levelno, "▪️")
prefix = f"{emoji} "
return prefix + (" " * max(0, 3 - _display_width(prefix)))
class _EmojiFormatter(logging.Formatter):
"""Compact formatter: emoji + message, no timestamp/level text/logger name."""
def format(self, record: logging.LogRecord) -> str:
return _format_level_prefix(record.levelno) + record.getMessage()
def configure_main_logging() -> None:
"""Replace Sanic's verbose 'Main yyyy-mm-dd INFO:' prefix with emoji-only format.
+11 -1
View File
@@ -62,12 +62,18 @@ async def close_client():
_client = None
async def validate_sso_request(request, *, perm: str = "cista:login") -> dict | None:
async def validate_sso_request(
request, *, perm: str = "cista:login", renew: bool = True
) -> dict | None:
"""Validate an SSO request against the auth backend.
Args:
request: The Sanic request object
perm: Permission to validate (default: cista:login, privileged also cista:admin)
renew: Whether to allow the auth backend to renew the session cookie.
Use ``False`` for WebSocket validation where Set-Cookie cannot be
forwarded to the client; this makes the request read-only and avoids
resetting the backend renewal timeout.
Returns:
User info dict if valid, None if validation fails with auth required response
@@ -88,12 +94,16 @@ async def validate_sso_request(request, *, perm: str = "cista:login") -> dict |
headers["cookie"] = request.headers["cookie"]
if "authorization" in request.headers:
headers["authorization"] = request.headers["authorization"]
if "user-agent" in request.headers:
headers["user-agent"] = request.headers["user-agent"]
headers["accept"] = "application/json"
headers["x-forwarded-for"] = request.client_ip
headers["x-forwarded-host"] = request.host
headers["x-forwarded-proto"] = request.scheme
url = f"{PASKIA_BACKEND_URL}/auth/api/validate?perm={perm}"
if not renew:
url += "&renew=0"
try:
response = await client.post(
+37
View File
@@ -0,0 +1,37 @@
"""Shared log formatting helpers with no Sanic dependency.
Used by the main process (cista.sanic_logging) and by the preview worker
subprocess, which must not import Sanic.
"""
import logging
import unicodedata
LEVEL_EMOJI = {
logging.DEBUG: "🔍",
logging.INFO: "", # noqa: RUF001
logging.WARNING: "⚠️",
logging.ERROR: "🛑",
logging.CRITICAL: "🛑",
}
def display_width(text: str) -> int:
return sum(
1 + (unicodedata.east_asian_width(c) in "FW")
for c in text
if unicodedata.category(c) != "Mn"
)
def format_level_prefix(levelno: int) -> str:
emoji = LEVEL_EMOJI.get(levelno, "▪️")
prefix = f"{emoji} "
return prefix + (" " * max(0, 3 - display_width(prefix)))
class EmojiFormatter(logging.Formatter):
"""Compact formatter: emoji + message, no timestamp/level text/logger name."""
def format(self, record: logging.LogRecord) -> str:
return format_level_prefix(record.levelno) + record.getMessage()
+1 -1
View File
@@ -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]
+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