Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0c3c3615ce | ||
|
|
07305538dc | ||
|
|
f3b3b5efd9 | ||
|
|
8613d6c25e | ||
|
|
5ed627d9f4 | ||
|
|
f0c3f7a7f9 |
+3
-2
@@ -2,11 +2,12 @@ import asyncio
|
|||||||
from secrets import token_bytes
|
from secrets import token_bytes
|
||||||
|
|
||||||
import msgspec
|
import msgspec
|
||||||
|
from mediapreview.office import is_available_cached
|
||||||
from sanic import Blueprint, json
|
from sanic import Blueprint, json
|
||||||
from sanic.exceptions import BadRequest
|
from sanic.exceptions import BadRequest
|
||||||
from sanic.log import logger
|
from sanic.log import logger
|
||||||
|
|
||||||
from cista import __version__, auth, config, onlyoffice, sharefs, sso, watching
|
from cista import __version__, auth, config, sharefs, sso, watching
|
||||||
from cista.auth import (
|
from cista.auth import (
|
||||||
create_share_token_handler,
|
create_share_token_handler,
|
||||||
create_token_handler,
|
create_token_handler,
|
||||||
@@ -67,7 +68,7 @@ async def watch(req, ws):
|
|||||||
"version": __version__,
|
"version": __version__,
|
||||||
"public": config.config.public,
|
"public": config.config.public,
|
||||||
"paskia": sso.paskia_enabled(),
|
"paskia": sso.paskia_enabled(),
|
||||||
"office_previews": await onlyoffice.is_available_cached(),
|
"office_previews": await is_available_cached(),
|
||||||
},
|
},
|
||||||
"user": user_info,
|
"user": user_info,
|
||||||
}
|
}
|
||||||
|
|||||||
+5
-3
@@ -10,6 +10,8 @@ from wsgiref.handlers import format_date_time
|
|||||||
|
|
||||||
import tracerite
|
import tracerite
|
||||||
from blake3 import blake3
|
from blake3 import blake3
|
||||||
|
from mediapreview.office import close_oo_client, log_reachable_info
|
||||||
|
from mediapreview.pool import shutdown_preview_workers, start_preview_workers
|
||||||
from sanic import Sanic, empty, raw, redirect
|
from sanic import Sanic, empty, raw, redirect
|
||||||
from sanic.exceptions import Forbidden, NotFound
|
from sanic.exceptions import Forbidden, NotFound
|
||||||
from sanic.log import logger
|
from sanic.log import logger
|
||||||
@@ -29,7 +31,6 @@ from cista import (
|
|||||||
watching,
|
watching,
|
||||||
)
|
)
|
||||||
from cista.api import bp
|
from cista.api import bp
|
||||||
from cista.preview import shutdown_preview_workers, start_preview_workers
|
|
||||||
from cista.sanic_logging import (
|
from cista.sanic_logging import (
|
||||||
configure_access_logging,
|
configure_access_logging,
|
||||||
configure_main_logging,
|
configure_main_logging,
|
||||||
@@ -129,6 +130,7 @@ setproctitle("cista-main")
|
|||||||
@app.before_server_start
|
@app.before_server_start
|
||||||
async def main_start(app):
|
async def main_start(app):
|
||||||
config.load_config()
|
config.load_config()
|
||||||
|
onlyoffice.configure()
|
||||||
setproctitle(f"cista {config.config.path.name}")
|
setproctitle(f"cista {config.config.path.name}")
|
||||||
app.ctx.threadexec = ThreadPoolExecutor(
|
app.ctx.threadexec = ThreadPoolExecutor(
|
||||||
max_workers=4, thread_name_prefix="cista-worker"
|
max_workers=4, thread_name_prefix="cista-worker"
|
||||||
@@ -142,7 +144,7 @@ async def main_start(app):
|
|||||||
@app.after_server_start
|
@app.after_server_start
|
||||||
async def main_after_start(app):
|
async def main_after_start(app):
|
||||||
_ = app
|
_ = app
|
||||||
onlyoffice.log_reachable_info()
|
log_reachable_info()
|
||||||
|
|
||||||
|
|
||||||
# Sanic sometimes fails to execute after_server_stop, so we do it before instead (potentially interrupting handlers)
|
# Sanic sometimes fails to execute after_server_stop, so we do it before instead (potentially interrupting handlers)
|
||||||
@@ -150,7 +152,7 @@ async def main_after_start(app):
|
|||||||
async def main_stop(app):
|
async def main_stop(app):
|
||||||
async with asyncio.TaskGroup() as tg:
|
async with asyncio.TaskGroup() as tg:
|
||||||
tg.create_task(asyncio.to_thread(watching.stop, app))
|
tg.create_task(asyncio.to_thread(watching.stop, app))
|
||||||
tg.create_task(onlyoffice.close_oo_client())
|
tg.create_task(close_oo_client())
|
||||||
tg.create_task(shutdown_preview_workers())
|
tg.create_task(shutdown_preview_workers())
|
||||||
tg.create_task(sso.close_client())
|
tg.create_task(sso.close_client())
|
||||||
|
|
||||||
|
|||||||
@@ -1,70 +0,0 @@
|
|||||||
# Patched OnlyOffice Document Server with configurable converter worker count.
|
|
||||||
#
|
|
||||||
# The Community Edition hardcodes the document converter to 1 worker,
|
|
||||||
# which creates a severe bottleneck under concurrent load.
|
|
||||||
# This image patches the open-source license.js to spawn a configurable
|
|
||||||
# number of converter workers (default 8).
|
|
||||||
#
|
|
||||||
# Build:
|
|
||||||
# docker build -t onlyoffice-cista docker/onlyoffice-converter-patch
|
|
||||||
#
|
|
||||||
# Run:
|
|
||||||
# docker run -d -p 8988:80 \
|
|
||||||
# -e WORKERS=16 \
|
|
||||||
# -e JWT_SECRET=your-strong-secret \
|
|
||||||
# --name onlyoffice onlyoffice-cista
|
|
||||||
#
|
|
||||||
# JWT:
|
|
||||||
# Set JWT_SECRET to the same value you pass to Cista as ONLYOFFICE_JWT_SECRET.
|
|
||||||
# OnlyOffice will enable token validation automatically.
|
|
||||||
#
|
|
||||||
# The ONLYOFFICE_VERSION build arg lets you target a specific release.
|
|
||||||
|
|
||||||
ARG ONLYOFFICE_VERSION=9.3.1
|
|
||||||
|
|
||||||
FROM onlyoffice/documentserver:${ONLYOFFICE_VERSION}
|
|
||||||
|
|
||||||
# Prevent interactive apt prompts
|
|
||||||
ENV DEBIAN_FRONTEND=noninteractive
|
|
||||||
|
|
||||||
# Install Node.js, npm, and git so we can run the FileConverter from source.
|
|
||||||
RUN apt-get update -qq && \
|
|
||||||
apt-get install -y -qq --no-install-recommends \
|
|
||||||
nodejs \
|
|
||||||
npm \
|
|
||||||
git \
|
|
||||||
ca-certificates && \
|
|
||||||
rm -rf /var/lib/apt/lists/*
|
|
||||||
|
|
||||||
# Clone the open-source server components (shallow, ~15 MB).
|
|
||||||
# The master branch is used because the Linux/web tags are not published
|
|
||||||
# in the server repo; the license.js file has been stable for years.
|
|
||||||
RUN git clone --depth 1 https://github.com/ONLYOFFICE/server.git /opt/oo-server
|
|
||||||
|
|
||||||
# Patch license.js so the converter worker count is read from an env var
|
|
||||||
# instead of being hardcoded to 1.
|
|
||||||
RUN sed -i \
|
|
||||||
's/count: 1,/count: parseInt(process.env.WORKERS, 10) || 8,/' \
|
|
||||||
/opt/oo-server/Common/sources/license.js
|
|
||||||
|
|
||||||
# Install npm dependencies for the modules the FileConverter touches.
|
|
||||||
# DocService deps are also needed because converter.js pulls in baseConnector.
|
|
||||||
RUN cd /opt/oo-server/Common && npm ci --no-audit --no-fund
|
|
||||||
RUN cd /opt/oo-server/FileConverter && npm ci --no-audit --no-fund
|
|
||||||
RUN cd /opt/oo-server/DocService && npm ci --no-audit --no-fund
|
|
||||||
|
|
||||||
# Back up the compiled pkg binary and replace it with our wrapper.
|
|
||||||
RUN mv /var/www/onlyoffice/documentserver/server/FileConverter/converter \
|
|
||||||
/var/www/onlyoffice/documentserver/server/FileConverter/converter.orig
|
|
||||||
|
|
||||||
COPY converter-wrapper.sh /var/www/onlyoffice/documentserver/server/FileConverter/converter
|
|
||||||
RUN chmod +x /var/www/onlyoffice/documentserver/server/FileConverter/converter
|
|
||||||
|
|
||||||
# Default worker count (override at runtime with -e WORKERS=16).
|
|
||||||
ENV WORKERS=8
|
|
||||||
|
|
||||||
# Use our custom entrypoint to persist the env var to a file that the
|
|
||||||
# non-root converter process (user=ds) can read.
|
|
||||||
COPY entrypoint.sh /app/ds/run-document-server-patched.sh
|
|
||||||
RUN chmod +x /app/ds/run-document-server-patched.sh
|
|
||||||
ENTRYPOINT ["/app/ds/run-document-server-patched.sh"]
|
|
||||||
@@ -1,19 +0,0 @@
|
|||||||
#!/bin/bash
|
|
||||||
# Wrapper that runs the OnlyOffice FileConverter from patched Node.js source.
|
|
||||||
# Replaces the compiled pkg binary shipped with the Community Edition.
|
|
||||||
|
|
||||||
# The env var is not passed through supervisor to the 'ds' user, so we read
|
|
||||||
# it from a file written by the custom entrypoint.
|
|
||||||
if [ -z "${WORKERS}" ] && [ -r /tmp/oo-converter-workers.txt ]; then
|
|
||||||
export WORKERS=$(cat /tmp/oo-converter-workers.txt)
|
|
||||||
fi
|
|
||||||
|
|
||||||
cd /opt/oo-server/FileConverter || exit 1
|
|
||||||
|
|
||||||
export NODE_ENV=production-linux
|
|
||||||
export NODE_CONFIG_DIR=/etc/onlyoffice/documentserver
|
|
||||||
export NODE_DISABLE_COLORS=1
|
|
||||||
export APPLICATION_NAME=onlyoffice
|
|
||||||
export LD_LIBRARY_PATH=/var/www/onlyoffice/documentserver/server/FileConverter/bin
|
|
||||||
|
|
||||||
exec node sources/convertermaster.js "$@"
|
|
||||||
@@ -1,8 +0,0 @@
|
|||||||
#!/bin/bash
|
|
||||||
# Custom entrypoint that persists WORKERS to a file readable by
|
|
||||||
# the non-root user that supervisor uses to run the converter.
|
|
||||||
|
|
||||||
echo "${WORKERS:-8}" > /tmp/oo-converter-workers.txt
|
|
||||||
chmod 644 /tmp/oo-converter-workers.txt
|
|
||||||
|
|
||||||
exec /app/ds/run-document-server.sh "$@"
|
|
||||||
+12
-292
@@ -1,134 +1,27 @@
|
|||||||
"""OnlyOffice Document Server integration for office document preview.
|
"""Cista-specific OnlyOffice setup.
|
||||||
|
|
||||||
Provides server-side conversion of office documents to PNG via the
|
The conversion client itself lives in `mediapreview.office`; this module
|
||||||
OnlyOffice Document Server /ConvertService.ashx API. The resulting PNG
|
only bridges cista's config-derived JWT secret into it and wires the
|
||||||
is passed through pyvips for AVIF compression.
|
`--oosetup` Docker bootstrap to cista's config.
|
||||||
|
|
||||||
Environment requirements:
|
|
||||||
- OnlyOffice Document Server must be running and reachable.
|
|
||||||
- If Document Server runs in Docker, the callback host IP must be
|
|
||||||
reachable from the container (usually the docker bridge IP).
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
|
||||||
import json
|
|
||||||
import os
|
import os
|
||||||
import socket
|
|
||||||
import socketserver
|
|
||||||
import subprocess
|
|
||||||
import threading
|
|
||||||
import urllib.error
|
|
||||||
import urllib.request
|
|
||||||
from functools import partial
|
|
||||||
from http.server import SimpleHTTPRequestHandler
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from time import perf_counter
|
|
||||||
from urllib.parse import quote
|
|
||||||
|
|
||||||
import httpx
|
import mediapreview.office
|
||||||
import jwt
|
|
||||||
from sanic.log import logger
|
|
||||||
|
|
||||||
from cista import config
|
from cista import config
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
# Configuration helpers
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
|
|
||||||
|
def configure() -> None:
|
||||||
_httpx_client: httpx.AsyncClient | None = None
|
"""Point mediapreview's OnlyOffice client at cista's derived JWT secret."""
|
||||||
|
os.environ.setdefault(
|
||||||
|
"ONLYOFFICE_JWT_SECRET", config.derived_secret("onlyoffice", size=16).hex()
|
||||||
def _get_onlyoffice_url() -> str:
|
|
||||||
return os.environ.get("ONLYOFFICE_CISTA_URL", "http://localhost:8988")
|
|
||||||
|
|
||||||
|
|
||||||
def _get_jwt_secret() -> str:
|
|
||||||
return (
|
|
||||||
os.environ.get("ONLYOFFICE_JWT_SECRET")
|
|
||||||
or config.derived_secret("onlyoffice", size=16).hex()
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def _get_callback_host() -> str:
|
|
||||||
"""Return the host IP that OnlyOffice (usually in Docker) can use to reach us."""
|
|
||||||
if host := os.environ.get("ONLYOFFICE_CALLBACK_HOST"):
|
|
||||||
return host
|
|
||||||
# Try to auto-detect docker bridge IP
|
|
||||||
try:
|
|
||||||
result = subprocess.run(
|
|
||||||
["/sbin/ip", "-4", "addr", "show", "docker0"],
|
|
||||||
capture_output=True,
|
|
||||||
text=True,
|
|
||||||
timeout=2,
|
|
||||||
check=False,
|
|
||||||
)
|
|
||||||
for line in result.stdout.splitlines():
|
|
||||||
if "inet " in line:
|
|
||||||
parts = line.strip().split()
|
|
||||||
addr_part = parts[1] # e.g. 172.17.0.1/16
|
|
||||||
return addr_part.split("/")[0]
|
|
||||||
except Exception:
|
|
||||||
logger.debug("Failed to auto-detect docker bridge IP")
|
|
||||||
return "127.0.0.1"
|
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
# Async HTTP client
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
|
|
||||||
|
|
||||||
def get_httpx_client() -> httpx.AsyncClient:
|
|
||||||
"""Return the shared async HTTP client for OnlyOffice requests."""
|
|
||||||
global _httpx_client
|
|
||||||
if _httpx_client is None:
|
|
||||||
_httpx_client = httpx.AsyncClient()
|
|
||||||
return _httpx_client
|
|
||||||
|
|
||||||
|
|
||||||
async def close_oo_client() -> None:
|
|
||||||
"""Close the shared async HTTP client."""
|
|
||||||
global _httpx_client
|
|
||||||
if _httpx_client is not None:
|
|
||||||
await _httpx_client.aclose()
|
|
||||||
_httpx_client = None
|
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
# Availability check
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
|
|
||||||
|
|
||||||
def _probe_status() -> tuple[bool, bool, str | None]:
|
|
||||||
"""Return (ok, responded, detail) for a lightweight reachability probe."""
|
|
||||||
url = _get_onlyoffice_url().rstrip("/") + "/ConvertService.ashx"
|
|
||||||
try:
|
|
||||||
with urllib.request.urlopen(url, timeout=2) as resp: # noqa: S310
|
|
||||||
status = resp.status
|
|
||||||
except urllib.error.HTTPError as e:
|
|
||||||
status = e.code
|
|
||||||
except Exception:
|
|
||||||
return False, False, None
|
|
||||||
|
|
||||||
if status in (200, 405):
|
|
||||||
return True, True, None
|
|
||||||
if status >= 500:
|
|
||||||
return False, True, f"HTTP {status}"
|
|
||||||
return False, True, f"HTTP {status}"
|
|
||||||
|
|
||||||
|
|
||||||
def log_reachable_info() -> None:
|
|
||||||
"""Log info on success, warning on responded probe errors, silent on no-response."""
|
|
||||||
ok, responded, detail = _probe_status()
|
|
||||||
if ok:
|
|
||||||
logger.info("Using OnlyOffice document server at %s", _get_onlyoffice_url())
|
|
||||||
elif responded:
|
|
||||||
suffix = f": {detail}" if detail else ""
|
|
||||||
logger.warning("OnlyOffice probe failed%s", suffix)
|
|
||||||
|
|
||||||
|
|
||||||
def setup_docker(confdir: Path | None = None) -> int:
|
def setup_docker(confdir: Path | None = None) -> int:
|
||||||
"""Build and run the patched OnlyOffice Docker image."""
|
"""Build and run the patched OnlyOffice Docker image (via mediapreview)."""
|
||||||
if confdir is not None:
|
if confdir is not None:
|
||||||
os.environ["CISTA_HOME"] = confdir.as_posix()
|
os.environ["CISTA_HOME"] = confdir.as_posix()
|
||||||
config.init_confdir()
|
config.init_confdir()
|
||||||
@@ -142,178 +35,5 @@ def setup_docker(confdir: Path | None = None) -> int:
|
|||||||
"public": False,
|
"public": False,
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
configure()
|
||||||
secret = config.derived_secret("onlyoffice", size=16).hex()
|
return mediapreview.office.setup_docker()
|
||||||
docker_dir = Path(__file__).parent / "docker"
|
|
||||||
if not docker_dir.is_dir():
|
|
||||||
raise FileNotFoundError(
|
|
||||||
f"Docker files not found at {docker_dir}. Is the package installed correctly?"
|
|
||||||
)
|
|
||||||
|
|
||||||
logger.info("Building OnlyOffice image")
|
|
||||||
build_cmd = ["docker", "build", "-t", "onlyoffice-cista", str(docker_dir)]
|
|
||||||
logger.info("%s", " ".join(build_cmd))
|
|
||||||
result = subprocess.run(build_cmd, check=False, shell=False) # noqa: S603
|
|
||||||
if result.returncode != 0:
|
|
||||||
raise RuntimeError("Failed to build OnlyOffice image")
|
|
||||||
|
|
||||||
logger.info("Starting OnlyOffice container")
|
|
||||||
run_cmd = [
|
|
||||||
"docker",
|
|
||||||
"run",
|
|
||||||
"-d",
|
|
||||||
"-p",
|
|
||||||
"8988:80",
|
|
||||||
"-e",
|
|
||||||
f"JWT_SECRET={secret}",
|
|
||||||
"-e",
|
|
||||||
"WORKERS=8",
|
|
||||||
"--name",
|
|
||||||
"onlyoffice-cista",
|
|
||||||
"--restart",
|
|
||||||
"unless-stopped",
|
|
||||||
"onlyoffice-cista",
|
|
||||||
]
|
|
||||||
logger.info("%s", " ".join(run_cmd))
|
|
||||||
result = subprocess.run(run_cmd, check=False, shell=False) # noqa: S603
|
|
||||||
if result.returncode != 0:
|
|
||||||
raise RuntimeError("Failed to start OnlyOffice container")
|
|
||||||
logger.info("OnlyOffice is running on http://localhost:8988")
|
|
||||||
return 0
|
|
||||||
|
|
||||||
|
|
||||||
async def is_available_async(request_timeout: float = 2.0) -> bool:
|
|
||||||
"""Return True if the configured OnlyOffice Document Server is reachable."""
|
|
||||||
url = _get_onlyoffice_url().rstrip("/") + "/ConvertService.ashx"
|
|
||||||
client = get_httpx_client()
|
|
||||||
try:
|
|
||||||
response = await client.get(url, timeout=request_timeout)
|
|
||||||
except Exception:
|
|
||||||
return False
|
|
||||||
else:
|
|
||||||
return response.status_code in (200, 405)
|
|
||||||
|
|
||||||
|
|
||||||
_oo_available_cache: tuple[bool, float] | None = None
|
|
||||||
OO_AVAILABILITY_CACHE_TTL = 30.0
|
|
||||||
|
|
||||||
|
|
||||||
async def is_available_cached() -> bool:
|
|
||||||
"""Return cached OnlyOffice availability, refreshed every 30 seconds."""
|
|
||||||
global _oo_available_cache
|
|
||||||
now = perf_counter()
|
|
||||||
if _oo_available_cache is not None:
|
|
||||||
result, timestamp = _oo_available_cache
|
|
||||||
if now - timestamp < OO_AVAILABILITY_CACHE_TTL:
|
|
||||||
return result
|
|
||||||
result = await is_available_async()
|
|
||||||
_oo_available_cache = (result, now)
|
|
||||||
return result
|
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
# Temporary HTTP server so OnlyOffice can download the file
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
|
|
||||||
|
|
||||||
class _QuietHandler(SimpleHTTPRequestHandler):
|
|
||||||
def log_message(self, fmt, *args) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
def _get_free_port() -> int:
|
|
||||||
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
|
||||||
s.bind(("0.0.0.0", 0)) # noqa: S104
|
|
||||||
return s.getsockname()[1]
|
|
||||||
|
|
||||||
|
|
||||||
def _serve_file_temporarily(file_path: Path):
|
|
||||||
"""Start a temporary HTTP server for *file_path* and return (url, server)."""
|
|
||||||
directory = str(file_path.parent)
|
|
||||||
filename = file_path.name
|
|
||||||
port = _get_free_port()
|
|
||||||
|
|
||||||
handler = partial(_QuietHandler, directory=directory)
|
|
||||||
httpd = socketserver.TCPServer(("0.0.0.0", port), handler) # noqa: S104
|
|
||||||
thread = threading.Thread(target=httpd.serve_forever, daemon=True)
|
|
||||||
thread.start()
|
|
||||||
|
|
||||||
host = _get_callback_host()
|
|
||||||
url = f"http://{host}:{port}/{quote(filename)}"
|
|
||||||
return url, httpd
|
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
# OnlyOffice conversion client
|
|
||||||
# ---------------------------------------------------------------------------
|
|
||||||
|
|
||||||
|
|
||||||
def _build_jwt_token(payload: dict) -> str | None:
|
|
||||||
secret = _get_jwt_secret()
|
|
||||||
if not secret:
|
|
||||||
return None
|
|
||||||
return jwt.encode(payload, secret, algorithm="HS256")
|
|
||||||
|
|
||||||
|
|
||||||
async def convert_to_png_async(file_path: Path, request_timeout: float = 5.0) -> bytes:
|
|
||||||
"""Convert *file_path* to PNG using OnlyOffice Document Server (async).
|
|
||||||
|
|
||||||
Returns the PNG bytes. Raises RuntimeError on failure.
|
|
||||||
"""
|
|
||||||
oo_url = _get_onlyoffice_url().rstrip("/")
|
|
||||||
convert_url = f"{oo_url}/ConvertService.ashx"
|
|
||||||
client = get_httpx_client()
|
|
||||||
|
|
||||||
# Start temporary HTTP server so OnlyOffice can fetch the file
|
|
||||||
doc_url, httpd = await asyncio.to_thread(_serve_file_temporarily, file_path)
|
|
||||||
try:
|
|
||||||
suffix = file_path.suffix.lstrip(".").lower()
|
|
||||||
payload = {
|
|
||||||
"async": False,
|
|
||||||
"filetype": suffix,
|
|
||||||
"key": f"cista_{(await asyncio.to_thread(file_path.stat)).st_mtime_ns}",
|
|
||||||
"outputtype": "png",
|
|
||||||
"title": file_path.name,
|
|
||||||
"url": doc_url,
|
|
||||||
}
|
|
||||||
|
|
||||||
headers = {"Content-Type": "application/json"}
|
|
||||||
token = _build_jwt_token(payload)
|
|
||||||
if token:
|
|
||||||
# Conversion API expects JWT in request body when token checks are enabled.
|
|
||||||
payload["token"] = token
|
|
||||||
headers["Authorization"] = token
|
|
||||||
|
|
||||||
t_start = perf_counter()
|
|
||||||
response = await client.post(
|
|
||||||
convert_url,
|
|
||||||
content=json.dumps(payload).encode(),
|
|
||||||
headers=headers,
|
|
||||||
timeout=request_timeout,
|
|
||||||
)
|
|
||||||
response.raise_for_status()
|
|
||||||
body = response.content
|
|
||||||
t_end = perf_counter()
|
|
||||||
|
|
||||||
# Parse XML response
|
|
||||||
text = body.decode("utf-8", errors="replace")
|
|
||||||
if "<Error>" in text:
|
|
||||||
code = "unknown"
|
|
||||||
if "<Error>" in text and "</Error>" in text:
|
|
||||||
code = text.split("<Error>")[1].split("</Error>")[0]
|
|
||||||
raise RuntimeError(f"OnlyOffice conversion error: {code}")
|
|
||||||
|
|
||||||
if "<FileUrl>" not in text:
|
|
||||||
raise RuntimeError("OnlyOffice response did not contain FileUrl")
|
|
||||||
|
|
||||||
file_url = text.split("<FileUrl>")[1].split("</FileUrl>")[0]
|
|
||||||
file_url = file_url.replace("&", "&")
|
|
||||||
|
|
||||||
logger.debug("OnlyOffice converted in %.2fs: %s", t_end - t_start, file_url)
|
|
||||||
|
|
||||||
# Download converted PNG
|
|
||||||
png_response = await client.get(file_url, timeout=request_timeout)
|
|
||||||
png_response.raise_for_status()
|
|
||||||
return png_response.content
|
|
||||||
finally:
|
|
||||||
await asyncio.to_thread(httpd.shutdown)
|
|
||||||
|
|||||||
+34
-606
@@ -1,448 +1,42 @@
|
|||||||
|
"""Preview HTTP blueprint: routing, caching and response building.
|
||||||
|
|
||||||
|
All conversion work is delegated to the mediapreview package (worker pool,
|
||||||
|
OnlyOffice integration, classification); this module only wires it into
|
||||||
|
Sanic with auth, etag negotiation and the in-memory response cache.
|
||||||
|
"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import contextlib
|
|
||||||
import mimetypes
|
|
||||||
import struct
|
|
||||||
import sys
|
|
||||||
import threading
|
|
||||||
import urllib.parse
|
import urllib.parse
|
||||||
from collections import OrderedDict
|
from pathlib import PurePosixPath
|
||||||
from dataclasses import dataclass
|
|
||||||
from multiprocessing import cpu_count
|
|
||||||
from pathlib import Path, PurePosixPath
|
|
||||||
from time import perf_counter
|
|
||||||
from urllib.parse import unquote
|
from urllib.parse import unquote
|
||||||
from wsgiref.handlers import format_date_time
|
from wsgiref.handlers import format_date_time
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
import msgspec
|
from mediapreview.formats import OFFICE_PREVIEW_SUFFIXES
|
||||||
from blake3 import blake3
|
from mediapreview.formats import expected_backend as _expected_preview_backend
|
||||||
|
from mediapreview.office import onlyoffice_error_short_text
|
||||||
|
from mediapreview.pool import (
|
||||||
|
PREVIEW_TIMEOUT,
|
||||||
|
PreviewError,
|
||||||
|
PreviewPoolClosedError,
|
||||||
|
PreviewTimeoutError,
|
||||||
|
generate_office_preview,
|
||||||
|
run_preview,
|
||||||
|
)
|
||||||
from sanic import Blueprint, empty, raw, redirect
|
from sanic import Blueprint, empty, raw, redirect
|
||||||
from sanic.exceptions import NotFound
|
from sanic.exceptions import NotFound
|
||||||
from sanic.log import logger
|
from sanic.log import logger
|
||||||
|
|
||||||
from cista import auth, config, onlyoffice, sharefs, watching
|
from cista import auth, config, sharefs, watching
|
||||||
from cista.fileio import fuid
|
from cista.fileio import fuid
|
||||||
from cista.preview_worker import (
|
|
||||||
DOC_PREVIEW_SUFFIXES,
|
|
||||||
OFFICE_PREVIEW_SUFFIXES,
|
|
||||||
PreviewRequest,
|
|
||||||
PreviewResponse,
|
|
||||||
)
|
|
||||||
from cista.util.filename import sanitize
|
from cista.util.filename import sanitize
|
||||||
|
from mediapreview import CachedPreview, PreviewCache, is_previewable_path
|
||||||
|
|
||||||
bp = Blueprint("preview", url_prefix="/preview")
|
bp = Blueprint("preview", url_prefix="/preview")
|
||||||
|
|
||||||
|
|
||||||
@dataclass(slots=True)
|
|
||||||
class CachedPreview:
|
|
||||||
"""Cached preview with headers and body."""
|
|
||||||
|
|
||||||
headers: dict[str, str]
|
|
||||||
body: bytes
|
|
||||||
|
|
||||||
|
|
||||||
class PreviewCache:
|
|
||||||
"""Thread-safe LRU cache for preview responses."""
|
|
||||||
|
|
||||||
def __init__(self, capacity: int = 500):
|
|
||||||
self.capacity = capacity
|
|
||||||
self._cache: OrderedDict[str, CachedPreview] = OrderedDict()
|
|
||||||
self._lock = threading.Lock()
|
|
||||||
|
|
||||||
def get(self, key: str) -> CachedPreview | None:
|
|
||||||
"""Get cached preview, moving it to end (most recently used)."""
|
|
||||||
with self._lock:
|
|
||||||
if key in self._cache:
|
|
||||||
self._cache.move_to_end(key)
|
|
||||||
return self._cache[key]
|
|
||||||
return None
|
|
||||||
|
|
||||||
def set(self, key: str, value: CachedPreview) -> None:
|
|
||||||
"""Cache preview, evicting oldest if at capacity."""
|
|
||||||
with self._lock:
|
|
||||||
if key in self._cache:
|
|
||||||
self._cache.move_to_end(key)
|
|
||||||
else:
|
|
||||||
if len(self._cache) >= self.capacity:
|
|
||||||
self._cache.popitem(last=False)
|
|
||||||
self._cache[key] = value
|
|
||||||
|
|
||||||
def __len__(self) -> int:
|
|
||||||
with self._lock:
|
|
||||||
return len(self._cache)
|
|
||||||
|
|
||||||
|
|
||||||
# Global preview cache instance
|
# Global preview cache instance
|
||||||
_preview_cache = PreviewCache(capacity=500)
|
_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()
|
|
||||||
AVIF_FAST_EFFORT = 0
|
|
||||||
WORKER_CHECKSUM_BYTES = 32
|
|
||||||
WORKER_MAX_JSON_BYTES = 1_000_000
|
|
||||||
|
|
||||||
|
|
||||||
class WorkerChecksumError(Exception):
|
|
||||||
"""Raised when worker response checksum does not match the packet."""
|
|
||||||
|
|
||||||
|
|
||||||
class WorkerProtocolError(Exception):
|
|
||||||
"""Raised when worker response packet is malformed."""
|
|
||||||
|
|
||||||
|
|
||||||
class _PreviewWorker:
|
|
||||||
def __init__(self, proc: asyncio.subprocess.Process):
|
|
||||||
self.proc = proc
|
|
||||||
|
|
||||||
async def request(
|
|
||||||
self,
|
|
||||||
filepath,
|
|
||||||
quality: int,
|
|
||||||
maxsize: int,
|
|
||||||
maxzoom: float,
|
|
||||||
data: bytes | None = None,
|
|
||||||
):
|
|
||||||
if self.proc.returncode is not None:
|
|
||||||
raise WorkerProtocolError("worker already exited")
|
|
||||||
if self.proc.stdin is None or self.proc.stdout is None:
|
|
||||||
raise WorkerProtocolError("worker streams not available")
|
|
||||||
|
|
||||||
meta = msgspec.json.encode(
|
|
||||||
PreviewRequest(
|
|
||||||
path=str(filepath),
|
|
||||||
quality=quality,
|
|
||||||
maxsize=maxsize,
|
|
||||||
maxzoom=maxzoom,
|
|
||||||
)
|
|
||||||
)
|
|
||||||
payload = data or b""
|
|
||||||
packet = struct.pack("<II", len(meta), len(payload)) + meta + payload
|
|
||||||
self.proc.stdin.write(packet)
|
|
||||||
await self.proc.stdin.drain()
|
|
||||||
|
|
||||||
checksum = await self.proc.stdout.readexactly(WORKER_CHECKSUM_BYTES)
|
|
||||||
header = await self.proc.stdout.readexactly(8)
|
|
||||||
json_size, data_size = struct.unpack("<II", header)
|
|
||||||
if json_size > WORKER_MAX_JSON_BYTES:
|
|
||||||
raise WorkerProtocolError(f"worker JSON too large: {json_size}")
|
|
||||||
meta_raw = await self.proc.stdout.readexactly(json_size)
|
|
||||||
payload = await self.proc.stdout.readexactly(data_size)
|
|
||||||
packet = header + meta_raw + payload
|
|
||||||
if blake3(packet).digest() != checksum:
|
|
||||||
raise WorkerChecksumError("worker checksum mismatch")
|
|
||||||
|
|
||||||
resp = msgspec.json.decode(meta_raw, type=PreviewResponse)
|
|
||||||
if not resp.ok:
|
|
||||||
raise PreviewError(
|
|
||||||
resp.error or "preview worker error",
|
|
||||||
stderr=resp.stderr,
|
|
||||||
backend=resp.backend,
|
|
||||||
)
|
|
||||||
return payload or None, resp
|
|
||||||
|
|
||||||
async def kill(self) -> None:
|
|
||||||
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:
|
|
||||||
def __init__(self, size: int):
|
|
||||||
self.size = size
|
|
||||||
self._idle: asyncio.Queue[_PreviewWorker] = asyncio.Queue()
|
|
||||||
self._pending: asyncio.PriorityQueue[tuple[int, int, asyncio.Future, tuple]] = (
|
|
||||||
asyncio.PriorityQueue()
|
|
||||||
)
|
|
||||||
self._workers: set[_PreviewWorker] = set()
|
|
||||||
self._dispatchers: list[asyncio.Task] = []
|
|
||||||
self._seq = 0
|
|
||||||
self._closed = False
|
|
||||||
|
|
||||||
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=None,
|
|
||||||
)
|
|
||||||
_active_procs.add(proc)
|
|
||||||
try:
|
|
||||||
ready = await asyncio.wait_for(proc.stdout.readexactly(1), timeout=30.0)
|
|
||||||
except TimeoutError as err:
|
|
||||||
with contextlib.suppress(ProcessLookupError):
|
|
||||||
proc.kill()
|
|
||||||
with contextlib.suppress(Exception):
|
|
||||||
await proc.wait()
|
|
||||||
raise WorkerProtocolError(
|
|
||||||
"preview worker failed to become ready"
|
|
||||||
" (worker stderr goes to the server log)"
|
|
||||||
) from err
|
|
||||||
except asyncio.IncompleteReadError as 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}")
|
|
||||||
return _PreviewWorker(proc)
|
|
||||||
|
|
||||||
async def _add_worker(self) -> None:
|
|
||||||
worker = await self._spawn_worker()
|
|
||||||
self._workers.add(worker)
|
|
||||||
await self._idle.put(worker)
|
|
||||||
|
|
||||||
async def _replace_worker(self, worker: _PreviewWorker) -> None:
|
|
||||||
self._workers.discard(worker)
|
|
||||||
try:
|
|
||||||
await worker.kill()
|
|
||||||
except Exception:
|
|
||||||
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:
|
|
||||||
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
|
|
||||||
|
|
||||||
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():
|
|
||||||
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]),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
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(
|
|
||||||
*(self._spawn_worker() for _ in range(self.size))
|
|
||||||
)
|
|
||||||
for worker in workers:
|
|
||||||
self._workers.add(worker)
|
|
||||||
await self._idle.put(worker)
|
|
||||||
for _ in range(self.size):
|
|
||||||
self._dispatchers.append(asyncio.create_task(self._dispatch_loop()))
|
|
||||||
|
|
||||||
async def run(
|
|
||||||
self,
|
|
||||||
filepath,
|
|
||||||
quality: int,
|
|
||||||
maxsize: int,
|
|
||||||
maxzoom: float,
|
|
||||||
data: bytes | None = None,
|
|
||||||
):
|
|
||||||
if self._closed:
|
|
||||||
raise PreviewError("preview worker pool closed")
|
|
||||||
loop = asyncio.get_running_loop()
|
|
||||||
future = loop.create_future()
|
|
||||||
self._seq += 1
|
|
||||||
await self._pending.put(
|
|
||||||
(
|
|
||||||
_preview_job_priority(filepath),
|
|
||||||
self._seq,
|
|
||||||
future,
|
|
||||||
(filepath, quality, maxsize, maxzoom, data),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
return await future
|
|
||||||
|
|
||||||
async def close(self) -> None:
|
|
||||||
self._closed = True
|
|
||||||
for task in self._dispatchers:
|
|
||||||
task.cancel()
|
|
||||||
if self._dispatchers:
|
|
||||||
await asyncio.gather(*self._dispatchers, return_exceptions=True)
|
|
||||||
self._dispatchers.clear()
|
|
||||||
workers = list(self._workers)
|
|
||||||
self._workers.clear()
|
|
||||||
while not self._pending.empty():
|
|
||||||
try:
|
|
||||||
_priority, _seq, future, _args = self._pending.get_nowait()
|
|
||||||
except asyncio.QueueEmpty:
|
|
||||||
break
|
|
||||||
if not future.done():
|
|
||||||
future.set_exception(PreviewError("preview worker pool closed"))
|
|
||||||
while not self._idle.empty():
|
|
||||||
try:
|
|
||||||
self._idle.get_nowait()
|
|
||||||
except asyncio.QueueEmpty:
|
|
||||||
break
|
|
||||||
await asyncio.gather(
|
|
||||||
*(worker.kill() for worker in workers), return_exceptions=True
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
async def start_preview_workers() -> None:
|
|
||||||
"""Warm up persistent preview workers during server startup."""
|
|
||||||
global _preview_pool
|
|
||||||
if _preview_pool is not None:
|
|
||||||
return
|
|
||||||
async with _preview_pool_lock:
|
|
||||||
if _preview_pool is not None:
|
|
||||||
return
|
|
||||||
pool = _PreviewWorkerPool(PREVIEW_WORKERS)
|
|
||||||
await pool.start()
|
|
||||||
_preview_pool = pool
|
|
||||||
logger.info("Started %d persistent preview workers", PREVIEW_WORKERS)
|
|
||||||
|
|
||||||
|
|
||||||
async def shutdown_preview_workers() -> None:
|
|
||||||
"""Kill persistent preview workers (called during server shutdown)."""
|
|
||||||
global _preview_pool
|
|
||||||
async with _preview_pool_lock:
|
|
||||||
pool = _preview_pool
|
|
||||||
_preview_pool = None
|
|
||||||
if pool is not None:
|
|
||||||
await pool.close()
|
|
||||||
if not _active_procs:
|
|
||||||
return
|
|
||||||
for proc in list(_active_procs):
|
|
||||||
with contextlib.suppress(ProcessLookupError):
|
|
||||||
proc.kill()
|
|
||||||
await asyncio.gather(
|
|
||||||
*(proc.wait() for proc in list(_active_procs)), return_exceptions=True
|
|
||||||
)
|
|
||||||
_active_procs.clear()
|
|
||||||
|
|
||||||
|
|
||||||
@bp.on_request
|
@bp.on_request
|
||||||
async def verify_preview(request):
|
async def verify_preview(request):
|
||||||
@@ -450,178 +44,6 @@ async def verify_preview(request):
|
|||||||
await auth.verify(request)
|
await auth.verify(request)
|
||||||
|
|
||||||
|
|
||||||
class PreviewTimeoutError(Exception):
|
|
||||||
"""Raised when the preview subprocess exceeds PREVIEW_TIMEOUT."""
|
|
||||||
|
|
||||||
def __init__(self, message: str, *, backend: str | None = None):
|
|
||||||
super().__init__(message)
|
|
||||||
self.backend = backend
|
|
||||||
|
|
||||||
|
|
||||||
class PreviewError(Exception):
|
|
||||||
"""Raised when the preview subprocess exits with a non-zero status."""
|
|
||||||
|
|
||||||
def __init__(
|
|
||||||
self,
|
|
||||||
message: str,
|
|
||||||
*,
|
|
||||||
stderr: str | None = None,
|
|
||||||
backend: str | None = None,
|
|
||||||
):
|
|
||||||
super().__init__(message)
|
|
||||||
self.stderr = stderr
|
|
||||||
self.backend = backend
|
|
||||||
|
|
||||||
|
|
||||||
# Max concurrent OnlyOffice conversion requests. OO has its own queue;
|
|
||||||
# we must not flood it. This is intentionally small.
|
|
||||||
OO_MAX_CONCURRENT = PREVIEW_WORKERS
|
|
||||||
|
|
||||||
|
|
||||||
class OOConversionManager:
|
|
||||||
"""Manages async OnlyOffice conversions with deduplication and concurrency limits."""
|
|
||||||
|
|
||||||
def __init__(self, max_concurrent: int = OO_MAX_CONCURRENT):
|
|
||||||
self._semaphore = asyncio.Semaphore(max_concurrent)
|
|
||||||
self._in_flight: dict[str, asyncio.Future[bytes]] = {}
|
|
||||||
self._tasks: set[asyncio.Task[None]] = set()
|
|
||||||
self._lock = asyncio.Lock()
|
|
||||||
|
|
||||||
async def convert(self, filepath: Path) -> bytes:
|
|
||||||
"""Return PNG bytes for *filepath*, deduplicating concurrent requests."""
|
|
||||||
stat = await asyncio.to_thread(filepath.stat)
|
|
||||||
key = f"{filepath}:{stat.st_mtime_ns}"
|
|
||||||
|
|
||||||
async with self._lock:
|
|
||||||
if key in self._in_flight:
|
|
||||||
future = self._in_flight[key]
|
|
||||||
else:
|
|
||||||
future = asyncio.get_running_loop().create_future()
|
|
||||||
self._in_flight[key] = future
|
|
||||||
task = asyncio.create_task(self._do_convert(filepath, key, future))
|
|
||||||
self._tasks.add(task)
|
|
||||||
task.add_done_callback(self._tasks.discard)
|
|
||||||
|
|
||||||
return await future
|
|
||||||
|
|
||||||
async def _do_convert(
|
|
||||||
self, filepath: Path, key: str, future: asyncio.Future[bytes]
|
|
||||||
) -> None:
|
|
||||||
try:
|
|
||||||
async with self._semaphore:
|
|
||||||
png_bytes = await onlyoffice.convert_to_png_async(
|
|
||||||
filepath, request_timeout=5.0
|
|
||||||
)
|
|
||||||
except Exception as e:
|
|
||||||
if not future.done():
|
|
||||||
future.set_exception(e)
|
|
||||||
async with self._lock:
|
|
||||||
self._in_flight.pop(key, None)
|
|
||||||
else:
|
|
||||||
if not future.done():
|
|
||||||
future.set_result(png_bytes)
|
|
||||||
async with self._lock:
|
|
||||||
self._in_flight.pop(key, None)
|
|
||||||
|
|
||||||
|
|
||||||
_oo_manager: OOConversionManager | None = None
|
|
||||||
|
|
||||||
|
|
||||||
def get_oo_manager() -> OOConversionManager:
|
|
||||||
"""Return the singleton OOConversionManager."""
|
|
||||||
global _oo_manager
|
|
||||||
if _oo_manager is None:
|
|
||||||
_oo_manager = OOConversionManager(max_concurrent=OO_MAX_CONCURRENT)
|
|
||||||
return _oo_manager
|
|
||||||
|
|
||||||
|
|
||||||
async def _generate_office_preview(
|
|
||||||
filepath: Path, quality: int, maxsize: int, maxzoom: float
|
|
||||||
) -> tuple[bytes | None, PreviewResponse | None]:
|
|
||||||
"""Generate a preview for an office file using OnlyOffice + worker AVIF conversion."""
|
|
||||||
manager = get_oo_manager()
|
|
||||||
t_oo_start = perf_counter()
|
|
||||||
png_bytes = await manager.convert(filepath)
|
|
||||||
t_oo_end = perf_counter()
|
|
||||||
|
|
||||||
img, resp = await _run_preview_process(
|
|
||||||
filepath, quality, maxsize, maxzoom, data=png_bytes
|
|
||||||
)
|
|
||||||
|
|
||||||
if resp is not None:
|
|
||||||
resp.backend = "onlyoffice+" + (resp.backend or "pyvips")
|
|
||||||
if resp.timings:
|
|
||||||
resp.timings = [round((t_oo_end - t_oo_start) * 1000, 1), *resp.timings]
|
|
||||||
return img, resp
|
|
||||||
|
|
||||||
|
|
||||||
async def _run_preview_process(
|
|
||||||
filepath, quality: int, maxsize: int, maxzoom: float, data: bytes | None = None
|
|
||||||
) -> tuple[bytes | None, PreviewResponse | None]:
|
|
||||||
"""Run preview request in a persistent worker process."""
|
|
||||||
await start_preview_workers()
|
|
||||||
if _preview_pool is None:
|
|
||||||
raise PreviewError(f"preview worker pool unavailable for {filepath.name}")
|
|
||||||
return await _preview_pool.run(filepath, quality, maxsize, maxzoom, data)
|
|
||||||
|
|
||||||
|
|
||||||
def _onlyoffice_error_short_text(detail: str) -> str:
|
|
||||||
if detail.startswith("OnlyOffice conversion error:"):
|
|
||||||
code = detail.rsplit(":", 1)[-1].strip()
|
|
||||||
return {
|
|
||||||
"-8": "onlyoffice jwt error",
|
|
||||||
"-4": "onlyoffice input error",
|
|
||||||
"-2": "onlyoffice timeout error",
|
|
||||||
"-1": "onlyoffice unknown error",
|
|
||||||
}.get(code, f"onlyoffice {code} error")
|
|
||||||
if "OnlyOffice response did not contain FileUrl" in detail:
|
|
||||||
return "onlyoffice no-fileurl error"
|
|
||||||
return "onlyoffice error"
|
|
||||||
|
|
||||||
|
|
||||||
def _preview_job_priority(path) -> int:
|
|
||||||
"""Return priority for preview job (lower=higher priority).
|
|
||||||
|
|
||||||
Priority order: images (0) < video (1) < PDF (2) < office (3) < unknown (4)
|
|
||||||
"""
|
|
||||||
suffix = path.suffix.lower()
|
|
||||||
if suffix in DOC_PREVIEW_SUFFIXES:
|
|
||||||
return 2
|
|
||||||
if suffix in OFFICE_PREVIEW_SUFFIXES:
|
|
||||||
return 3
|
|
||||||
mime_type, _ = mimetypes.guess_type(path.name)
|
|
||||||
if mime_type and mime_type.startswith("image/"):
|
|
||||||
return 0
|
|
||||||
if mime_type and mime_type.startswith("video/"):
|
|
||||||
return 1
|
|
||||||
return 4
|
|
||||||
|
|
||||||
|
|
||||||
def _expected_preview_backend(path: Path) -> str:
|
|
||||||
"""Best-effort backend label used for timeout/access logging."""
|
|
||||||
suffix = path.suffix.lower()
|
|
||||||
if suffix in OFFICE_PREVIEW_SUFFIXES:
|
|
||||||
return "onlyoffice"
|
|
||||||
if suffix in DOC_PREVIEW_SUFFIXES:
|
|
||||||
return "pdf"
|
|
||||||
mime_type, _ = mimetypes.guess_type(path.name)
|
|
||||||
if mime_type and mime_type.startswith("video/"):
|
|
||||||
return "video"
|
|
||||||
if mime_type and mime_type.startswith("image/"):
|
|
||||||
return "pyvips"
|
|
||||||
return "preview"
|
|
||||||
|
|
||||||
|
|
||||||
def is_previewable_path(path) -> bool:
|
|
||||||
suffix = path.suffix.lower()
|
|
||||||
if suffix in DOC_PREVIEW_SUFFIXES or suffix in OFFICE_PREVIEW_SUFFIXES:
|
|
||||||
return True
|
|
||||||
mime_type, _ = mimetypes.guess_type(path.name)
|
|
||||||
if not mime_type:
|
|
||||||
return False
|
|
||||||
return mime_type.startswith(("image/", "video/"))
|
|
||||||
|
|
||||||
|
|
||||||
@bp.get("/<path:path>")
|
@bp.get("/<path:path>")
|
||||||
async def preview(req, path):
|
async def preview(req, path):
|
||||||
"""Preview a file"""
|
"""Preview a file"""
|
||||||
@@ -664,12 +86,12 @@ async def preview(req, path):
|
|||||||
try:
|
try:
|
||||||
if filepath.suffix.lower() in OFFICE_PREVIEW_SUFFIXES:
|
if filepath.suffix.lower() in OFFICE_PREVIEW_SUFFIXES:
|
||||||
img, preview_resp = await asyncio.wait_for(
|
img, preview_resp = await asyncio.wait_for(
|
||||||
_generate_office_preview(filepath, quality, maxsize, maxzoom),
|
generate_office_preview(filepath, quality, maxsize, maxzoom),
|
||||||
timeout=PREVIEW_TIMEOUT,
|
timeout=PREVIEW_TIMEOUT,
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
img, preview_resp = await asyncio.wait_for(
|
img, preview_resp = await asyncio.wait_for(
|
||||||
_run_preview_process(filepath, quality, maxsize, maxzoom),
|
run_preview(filepath, quality, maxsize, maxzoom),
|
||||||
timeout=PREVIEW_TIMEOUT,
|
timeout=PREVIEW_TIMEOUT,
|
||||||
)
|
)
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
@@ -689,22 +111,28 @@ async def preview(req, path):
|
|||||||
except RuntimeError as e:
|
except RuntimeError as e:
|
||||||
detail = str(e)
|
detail = str(e)
|
||||||
if detail.startswith("OnlyOffice"):
|
if detail.startswith("OnlyOffice"):
|
||||||
req.ctx.log_extra = _onlyoffice_error_short_text(detail)
|
req.ctx.log_extra = onlyoffice_error_short_text(detail)
|
||||||
return empty(503)
|
return empty(503)
|
||||||
raise
|
raise
|
||||||
|
except PreviewPoolClosedError:
|
||||||
|
# Server is shutting down; not an error, just a cancelled preview.
|
||||||
|
req.ctx.log_extra = "preview cancelled"
|
||||||
|
return empty(503)
|
||||||
except PreviewError as e:
|
except PreviewError as e:
|
||||||
if e.backend:
|
|
||||||
req.ctx.log_extra = e.backend
|
|
||||||
detail = str(e)
|
detail = str(e)
|
||||||
if detail == "preview worker error" and e.stderr:
|
if detail == "preview worker error" and e.stderr:
|
||||||
captured = e.stderr.strip()
|
captured = e.stderr.strip()
|
||||||
if captured:
|
if captured:
|
||||||
detail = captured.splitlines()[0]
|
detail = captured.splitlines()[0]
|
||||||
logger.error("%s preview: %s", filepath, detail)
|
# The worker already logged the failure (with traceback where the
|
||||||
|
# error occurred) — annotate the access log instead of re-logging.
|
||||||
|
req.ctx.log_extra = e.backend or detail
|
||||||
return empty(422)
|
return empty(422)
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
# Server shutdown or client disconnect: the connection is being torn
|
||||||
|
# down, so responding is impossible — just annotate the access log.
|
||||||
req.ctx.log_extra = "preview cancelled"
|
req.ctx.log_extra = "preview cancelled"
|
||||||
return empty(503)
|
raise
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Unhandled preview error for %s", filepath)
|
logger.exception("Unhandled preview error for %s", filepath)
|
||||||
return empty(500)
|
return empty(500)
|
||||||
|
|||||||
@@ -1,605 +0,0 @@
|
|||||||
"""Preview generation worker subprocess and synchronous preview engine.
|
|
||||||
|
|
||||||
Two modes are supported:
|
|
||||||
1) Legacy one-shot mode: argv has path/quality/maxsize/maxzoom.
|
|
||||||
2) Long-lived mode: read framed requests from stdin and write framed responses.
|
|
||||||
|
|
||||||
Framed request format (stdin):
|
|
||||||
(uint32 json size)(uint32 data size)(json)(binary data)
|
|
||||||
|
|
||||||
Framed response format (stdout):
|
|
||||||
(blake3(packet))(uint32 json size)(uint32 payload size)(json)(binary payload)
|
|
||||||
where packet = (uint32 json size)(uint32 payload size)(json)(binary payload).
|
|
||||||
"""
|
|
||||||
|
|
||||||
import contextlib
|
|
||||||
import gc
|
|
||||||
import io
|
|
||||||
import logging
|
|
||||||
import mimetypes
|
|
||||||
import os
|
|
||||||
import shlex
|
|
||||||
import struct
|
|
||||||
import subprocess
|
|
||||||
import sys
|
|
||||||
import tempfile
|
|
||||||
from pathlib import Path
|
|
||||||
from time import perf_counter
|
|
||||||
|
|
||||||
import av
|
|
||||||
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"}
|
|
||||||
|
|
||||||
OFFICE_PREVIEW_SUFFIXES = {
|
|
||||||
".doc",
|
|
||||||
".dot",
|
|
||||||
".docx",
|
|
||||||
".docm",
|
|
||||||
".dotx",
|
|
||||||
".dotm",
|
|
||||||
".rtf",
|
|
||||||
".odt",
|
|
||||||
".ott",
|
|
||||||
".txt",
|
|
||||||
".md",
|
|
||||||
".mhtml",
|
|
||||||
".mht",
|
|
||||||
".html",
|
|
||||||
".htm",
|
|
||||||
".xml",
|
|
||||||
".wps",
|
|
||||||
".wri",
|
|
||||||
# Spreadsheets
|
|
||||||
".xls",
|
|
||||||
".xlsx",
|
|
||||||
".xlsm",
|
|
||||||
".xlsb",
|
|
||||||
".xltx",
|
|
||||||
".xltm",
|
|
||||||
".ods",
|
|
||||||
".ots",
|
|
||||||
".csv",
|
|
||||||
# Presentations
|
|
||||||
".ppt",
|
|
||||||
".pptx",
|
|
||||||
".pptm",
|
|
||||||
".pps",
|
|
||||||
".ppsx",
|
|
||||||
".pot",
|
|
||||||
".potx",
|
|
||||||
".odp",
|
|
||||||
".otp",
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
class PreviewRequest(msgspec.Struct, omit_defaults=True):
|
|
||||||
path: str
|
|
||||||
quality: int
|
|
||||||
maxsize: int
|
|
||||||
maxzoom: float
|
|
||||||
|
|
||||||
|
|
||||||
class PreviewResponse(msgspec.Struct, omit_defaults=True):
|
|
||||||
ok: bool
|
|
||||||
mime: str | None = None
|
|
||||||
backend: str | None = None
|
|
||||||
timings: list[float] | None = None
|
|
||||||
error: str | None = None
|
|
||||||
stderr: str | None = None
|
|
||||||
width: int | None = None
|
|
||||||
height: int | None = None
|
|
||||||
|
|
||||||
|
|
||||||
_enc = msgspec.json.Encoder()
|
|
||||||
_dec_req = msgspec.json.Decoder(PreviewRequest)
|
|
||||||
|
|
||||||
|
|
||||||
def _read_exactly(f, n: int) -> bytes:
|
|
||||||
buf = b""
|
|
||||||
while len(buf) < n:
|
|
||||||
chunk = f.read(n - len(buf))
|
|
||||||
if not chunk:
|
|
||||||
raise EOFError
|
|
||||||
buf += chunk
|
|
||||||
return buf
|
|
||||||
|
|
||||||
|
|
||||||
def _read_request() -> tuple[PreviewRequest, bytes] | None:
|
|
||||||
try:
|
|
||||||
header = _read_exactly(sys.stdin.buffer, 8)
|
|
||||||
except EOFError:
|
|
||||||
return None
|
|
||||||
json_size, data_size = struct.unpack("<II", header)
|
|
||||||
meta_raw = _read_exactly(sys.stdin.buffer, json_size)
|
|
||||||
data = b""
|
|
||||||
if data_size:
|
|
||||||
data = _read_exactly(sys.stdin.buffer, data_size)
|
|
||||||
req = _dec_req.decode(meta_raw)
|
|
||||||
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()
|
|
||||||
out.write(checksum)
|
|
||||||
out.write(packet)
|
|
||||||
out.flush()
|
|
||||||
|
|
||||||
|
|
||||||
def dispatch(path, quality, maxsize, maxzoom, data=None):
|
|
||||||
backend = "unknown"
|
|
||||||
try:
|
|
||||||
if data:
|
|
||||||
backend = "pyvips"
|
|
||||||
return process_image_buffer(
|
|
||||||
data, quality=quality, maxsize=maxsize, maxzoom=maxzoom
|
|
||||||
)
|
|
||||||
suffix = path.suffix.lower()
|
|
||||||
if suffix in DOC_PREVIEW_SUFFIXES:
|
|
||||||
backend = "pdf"
|
|
||||||
return process_pdf(path, quality=quality, maxsize=maxsize, maxzoom=maxzoom)
|
|
||||||
mime_type, _ = mimetypes.guess_type(path.name)
|
|
||||||
if mime_type and mime_type.startswith("video/"):
|
|
||||||
backend = "video"
|
|
||||||
return process_video(path, quality=quality, maxsize=maxsize)
|
|
||||||
if mime_type and mime_type.startswith("image/"):
|
|
||||||
backend = "pyvips"
|
|
||||||
return process_image(path, quality=quality, maxsize=maxsize)
|
|
||||||
except ValueError as e:
|
|
||||||
return None, PreviewResponse(ok=False, backend=backend, error=str(e))
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception("Preview dispatch failed for %s", path)
|
|
||||||
return None, PreviewResponse(ok=False, backend=backend, error=str(e))
|
|
||||||
return None, PreviewResponse(ok=False, backend=backend, error="preview unsupported")
|
|
||||||
|
|
||||||
|
|
||||||
def process_image(path, *, maxsize, quality):
|
|
||||||
return process_image_pyvips(path, maxsize=maxsize, quality=quality)
|
|
||||||
|
|
||||||
|
|
||||||
def _get_image_dimensions(path: Path) -> tuple[int, int] | None:
|
|
||||||
"""Probe image dimensions.
|
|
||||||
|
|
||||||
pyvips can read the header of most formats (including HEIC) without
|
|
||||||
fully decoding the image.
|
|
||||||
"""
|
|
||||||
try:
|
|
||||||
img = pyvips.Image.new_from_file(str(path))
|
|
||||||
img = img.autorot()
|
|
||||||
except pyvips.error.Error:
|
|
||||||
return None
|
|
||||||
else:
|
|
||||||
return img.width, img.height
|
|
||||||
|
|
||||||
|
|
||||||
def _image_via_ffmpeg(path: Path, maxsize: int, quality: int) -> bytes:
|
|
||||||
"""Convert any image to AVIF using ffmpeg CLI.
|
|
||||||
|
|
||||||
ffmpeg handles HEIC tile assembly, EXIF rotation, HDR metadata and
|
|
||||||
ICC profile embedding automatically.
|
|
||||||
"""
|
|
||||||
dims = _get_image_dimensions(path)
|
|
||||||
crf = int(63 * (1 - quality / 100) ** 2)
|
|
||||||
with tempfile.NamedTemporaryFile(suffix=".avif", delete=False) as tmp_f:
|
|
||||||
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),
|
|
||||||
"-frames:v",
|
|
||||||
"1",
|
|
||||||
"-c:v",
|
|
||||||
"av1",
|
|
||||||
"-crf",
|
|
||||||
str(crf),
|
|
||||||
"-cpu-used",
|
|
||||||
"8",
|
|
||||||
tmp_path,
|
|
||||||
]
|
|
||||||
if dims is not None:
|
|
||||||
w, h = dims
|
|
||||||
if max(w, h) > maxsize:
|
|
||||||
scale = min(maxsize / w, maxsize / h)
|
|
||||||
new_w = int(w * scale)
|
|
||||||
new_h = int(h * scale)
|
|
||||||
# insert -s <wxh> right after the input file
|
|
||||||
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
|
|
||||||
# stdin, which carries the framed request protocol. An inherited
|
|
||||||
# stdin lets ffmpeg eat protocol bytes and, if the worker is
|
|
||||||
# killed mid-conversion, keeps the orphaned ffmpeg holding the
|
|
||||||
# pipe open so the parent's proc.wait() hangs forever.
|
|
||||||
subprocess.run( # noqa: S603
|
|
||||||
cmd,
|
|
||||||
capture_output=True,
|
|
||||||
check=True,
|
|
||||||
shell=False,
|
|
||||||
stdin=subprocess.DEVNULL,
|
|
||||||
)
|
|
||||||
except subprocess.CalledProcessError as e:
|
|
||||||
shell_cmd = shlex.join(cmd)
|
|
||||||
stderr = (e.stderr or b"").decode(errors="replace").strip()
|
|
||||||
if stderr:
|
|
||||||
raise RuntimeError(
|
|
||||||
f"ffmpeg failed (exit {e.returncode}):\n{shell_cmd}\n{stderr}"
|
|
||||||
) from e
|
|
||||||
raise RuntimeError(
|
|
||||||
f"ffmpeg failed (exit {e.returncode}):\n{shell_cmd}"
|
|
||||||
) from e
|
|
||||||
with Path(tmp_path).open("rb") as f:
|
|
||||||
return f.read()
|
|
||||||
finally:
|
|
||||||
Path(tmp_path).unlink(missing_ok=True)
|
|
||||||
|
|
||||||
|
|
||||||
def process_image_pyvips(path, *, maxsize, quality):
|
|
||||||
t_start = perf_counter()
|
|
||||||
suffix = path.suffix.lower()
|
|
||||||
|
|
||||||
# HEIC/HEIF: ffmpeg handles tile assembly and HDR correctly;
|
|
||||||
# skip pyvips entirely.
|
|
||||||
if suffix in (".heic", ".heif"):
|
|
||||||
heic_dims = _get_image_dimensions(path)
|
|
||||||
width, height = heic_dims or (None, None)
|
|
||||||
ret = _image_via_ffmpeg(path, maxsize, quality)
|
|
||||||
t_end = perf_counter()
|
|
||||||
return ret, PreviewResponse(
|
|
||||||
ok=True,
|
|
||||||
mime="image/avif",
|
|
||||||
backend="ffmpeg",
|
|
||||||
timings=[round((t_end - t_start) * 1000, 1)],
|
|
||||||
width=width,
|
|
||||||
height=height,
|
|
||||||
)
|
|
||||||
|
|
||||||
# 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"}
|
|
||||||
try:
|
|
||||||
img = pyvips.Image.new_from_file(str(path), **load_opts)
|
|
||||||
img = img.autorot()
|
|
||||||
orig_w, orig_h = img.width, img.height
|
|
||||||
scale = min(maxsize / img.width, maxsize / img.height, 1.0)
|
|
||||||
if scale < 1.0:
|
|
||||||
img = img.resize(scale)
|
|
||||||
ret = img.write_to_buffer(
|
|
||||||
".avif",
|
|
||||||
Q=quality,
|
|
||||||
effort=AVIF_FAST_EFFORT,
|
|
||||||
keep="none",
|
|
||||||
)
|
|
||||||
except pyvips.error.Error as e:
|
|
||||||
raise ValueError(f"cannot decode image: {e}") from e
|
|
||||||
backend = "pyvips"
|
|
||||||
t_end = perf_counter()
|
|
||||||
|
|
||||||
return ret, PreviewResponse(
|
|
||||||
ok=True,
|
|
||||||
mime="image/avif",
|
|
||||||
backend=backend,
|
|
||||||
timings=[round((t_end - t_start) * 1000, 1)],
|
|
||||||
width=orig_w,
|
|
||||||
height=orig_h,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def process_image_buffer(data: bytes, *, quality, maxsize, maxzoom):
|
|
||||||
_ = maxzoom
|
|
||||||
t_start = perf_counter()
|
|
||||||
img = pyvips.Image.new_from_buffer(data, "")
|
|
||||||
img = img.autorot()
|
|
||||||
orig_w, orig_h = img.width, img.height
|
|
||||||
scale = min(maxsize / img.width, maxsize / img.height, 1.0)
|
|
||||||
if scale < 1.0:
|
|
||||||
img = img.resize(scale)
|
|
||||||
ret = img.write_to_buffer(
|
|
||||||
".avif",
|
|
||||||
Q=quality,
|
|
||||||
effort=AVIF_FAST_EFFORT,
|
|
||||||
keep="none",
|
|
||||||
)
|
|
||||||
t_end = perf_counter()
|
|
||||||
|
|
||||||
return ret, PreviewResponse(
|
|
||||||
ok=True,
|
|
||||||
mime="image/avif",
|
|
||||||
backend="pyvips",
|
|
||||||
timings=[round((t_end - t_start) * 1000, 1)],
|
|
||||||
width=orig_w,
|
|
||||||
height=orig_h,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def process_pdf(path, *, maxsize, maxzoom, quality, page_number=0):
|
|
||||||
t_load_start = 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, keep="none")
|
|
||||||
backend = "pdf+pyvips"
|
|
||||||
t_save_end = perf_counter()
|
|
||||||
|
|
||||||
return ret, PreviewResponse(
|
|
||||||
ok=True,
|
|
||||||
mime="image/avif",
|
|
||||||
backend=backend,
|
|
||||||
timings=[
|
|
||||||
round((t_load_end - t_load_start) * 1000, 1),
|
|
||||||
round((t_save_end - t_save_start) * 1000, 1),
|
|
||||||
],
|
|
||||||
width=round(w),
|
|
||||||
height=round(h),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def process_video(path, *, maxsize, quality):
|
|
||||||
frame = None
|
|
||||||
imgdata = io.BytesIO()
|
|
||||||
istream = ostream = icc = occ = frame = None
|
|
||||||
t_load_start = perf_counter()
|
|
||||||
# Initialize to avoid "possibly unbound" in static analysis when exceptions occur
|
|
||||||
t_load_end = t_load_start
|
|
||||||
t_save_start = t_load_start
|
|
||||||
t_save_end = t_load_start
|
|
||||||
with (
|
|
||||||
av.open(
|
|
||||||
str(path),
|
|
||||||
options={
|
|
||||||
"analyzeduration": "1000000", # 1 second (in microseconds)
|
|
||||||
"fflags": "fastseek",
|
|
||||||
},
|
|
||||||
) as icontainer,
|
|
||||||
av.open(imgdata, "w", format="avif") as ocontainer,
|
|
||||||
):
|
|
||||||
istream = icontainer.streams.video[0]
|
|
||||||
istream.codec_context.skip_frame = "NONKEY"
|
|
||||||
icontainer.seek((icontainer.duration or 0) // 8)
|
|
||||||
for frame in icontainer.decode(istream):
|
|
||||||
if frame.dts is not None:
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
raise RuntimeError("No frames found in video")
|
|
||||||
|
|
||||||
# Resize frame to thumbnail size
|
|
||||||
# Capture display dimensions before resize (accounting for rotation)
|
|
||||||
disp_w = frame.width
|
|
||||||
disp_h = frame.height
|
|
||||||
if frame.rotation in (90, 270):
|
|
||||||
disp_w, disp_h = disp_h, disp_w
|
|
||||||
if frame.width > maxsize or frame.height > maxsize:
|
|
||||||
scale_factor = min(maxsize / frame.width, maxsize / frame.height)
|
|
||||||
new_width = int(frame.width * scale_factor)
|
|
||||||
new_height = int(frame.height * scale_factor)
|
|
||||||
frame = frame.reformat(width=new_width, height=new_height)
|
|
||||||
|
|
||||||
# Apply EXIF rotation if present
|
|
||||||
if frame.rotation:
|
|
||||||
# frame.rotation indicates clockwise rotation needed to display correctly
|
|
||||||
# np.rot90 rotates counter-clockwise, so we negate k
|
|
||||||
k = (frame.rotation // 90) % 4 # Convert to counter-clockwise rotations
|
|
||||||
if k == 2:
|
|
||||||
# 180° rotation can be done in YUV420p, preserving HDR
|
|
||||||
try:
|
|
||||||
fplanes = frame.to_ndarray()
|
|
||||||
# Split into Y, U, V planes of proper dimensions
|
|
||||||
planes = [
|
|
||||||
fplanes[: frame.height],
|
|
||||||
fplanes[
|
|
||||||
frame.height : frame.height + frame.height // 4
|
|
||||||
].reshape(frame.height // 2, frame.width // 2),
|
|
||||||
fplanes[frame.height + frame.height // 4 :].reshape(
|
|
||||||
frame.height // 2, frame.width // 2
|
|
||||||
),
|
|
||||||
]
|
|
||||||
# Rotate each plane by 180°
|
|
||||||
planes = [np.rot90(p, 2) for p in planes]
|
|
||||||
# Restore PyAV format
|
|
||||||
planes = np.hstack([p.flat for p in planes]).reshape(
|
|
||||||
-1, planes[0].shape[1]
|
|
||||||
)
|
|
||||||
frame = av.VideoFrame.from_ndarray(planes, format=frame.format.name)
|
|
||||||
del planes, fplanes
|
|
||||||
except Exception:
|
|
||||||
logger.exception("Error rotating video frame by 180°")
|
|
||||||
elif k in (1, 3):
|
|
||||||
# 90° or 270° rotation requires RGB conversion (loses HDR)
|
|
||||||
try:
|
|
||||||
rgb = frame.to_ndarray(format="rgb24")
|
|
||||||
rgb = np.rot90(rgb, k)
|
|
||||||
frame = av.VideoFrame.from_ndarray(rgb, format="rgb24")
|
|
||||||
frame = frame.reformat(
|
|
||||||
format="yuv420p"
|
|
||||||
) # Convert back for encoding
|
|
||||||
del rgb
|
|
||||||
except Exception:
|
|
||||||
logger.exception(
|
|
||||||
"Error rotating video frame by %s°", frame.rotation
|
|
||||||
)
|
|
||||||
|
|
||||||
# libsvtav1 rejects full-range JPEG-style YUV pixel formats such as
|
|
||||||
# yuvj420p, so normalize them before opening the encoder.
|
|
||||||
if frame.format.name.startswith("yuvj"):
|
|
||||||
frame = frame.reformat(format="yuv420p")
|
|
||||||
t_load_end = perf_counter()
|
|
||||||
|
|
||||||
t_save_start = perf_counter()
|
|
||||||
crf = str(int(63 * (1 - quality / 100) ** 2)) # Closely matching PIL quality-%
|
|
||||||
ostream = ocontainer.add_stream(
|
|
||||||
"av1",
|
|
||||||
options={
|
|
||||||
"crf": crf,
|
|
||||||
"usage": "realtime",
|
|
||||||
"cpu-used": "8",
|
|
||||||
"threads": "1",
|
|
||||||
},
|
|
||||||
)
|
|
||||||
if not isinstance(ostream, av.VideoStream):
|
|
||||||
raise TypeError("failed to initialize AV1 video stream")
|
|
||||||
ostream.width = frame.width
|
|
||||||
ostream.height = frame.height
|
|
||||||
ostream.pix_fmt = frame.format.name
|
|
||||||
icc = istream.codec_context
|
|
||||||
occ = ostream.codec_context
|
|
||||||
|
|
||||||
# Copy HDR metadata from input video stream
|
|
||||||
occ.color_primaries = icc.color_primaries
|
|
||||||
occ.color_trc = icc.color_trc
|
|
||||||
occ.colorspace = icc.colorspace
|
|
||||||
occ.color_range = icc.color_range
|
|
||||||
|
|
||||||
ocontainer.mux(ostream.encode(frame))
|
|
||||||
ocontainer.mux(ostream.encode(None)) # Flush the stream
|
|
||||||
t_save_end = perf_counter()
|
|
||||||
|
|
||||||
# Capture result before cleanup
|
|
||||||
ret = imgdata.getvalue()
|
|
||||||
resp = PreviewResponse(
|
|
||||||
ok=True,
|
|
||||||
mime="image/avif",
|
|
||||||
backend="video",
|
|
||||||
timings=[
|
|
||||||
round((t_load_end - t_load_start) * 1000, 1),
|
|
||||||
round((t_save_end - t_save_start) * 1000, 1),
|
|
||||||
],
|
|
||||||
width=disp_w,
|
|
||||||
height=disp_h,
|
|
||||||
)
|
|
||||||
del imgdata, istream, ostream, icc, occ, frame
|
|
||||||
gc.collect()
|
|
||||||
return ret, resp
|
|
||||||
|
|
||||||
|
|
||||||
def _run_once() -> None:
|
|
||||||
if len(sys.argv) != 5:
|
|
||||||
sys.stderr.write(f"Usage: {sys.argv[0]} <path> <quality> <maxsize> <maxzoom>\n")
|
|
||||||
sys.exit(1)
|
|
||||||
|
|
||||||
path = Path(sys.argv[1])
|
|
||||||
quality = int(sys.argv[2])
|
|
||||||
maxsize = int(sys.argv[3])
|
|
||||||
maxzoom = float(sys.argv[4])
|
|
||||||
result, _ = dispatch(path, quality, maxsize, maxzoom)
|
|
||||||
if result:
|
|
||||||
sys.stdout.buffer.write(result)
|
|
||||||
sys.stdout.buffer.flush()
|
|
||||||
|
|
||||||
|
|
||||||
def _run_loop() -> None:
|
|
||||||
while True:
|
|
||||||
result = _read_request()
|
|
||||||
if result is None:
|
|
||||||
return
|
|
||||||
req, data = result
|
|
||||||
stderr_capture = io.StringIO()
|
|
||||||
handler = logging.StreamHandler(stderr_capture)
|
|
||||||
root_logger = logging.getLogger()
|
|
||||||
root_logger.addHandler(handler)
|
|
||||||
try:
|
|
||||||
with contextlib.redirect_stderr(stderr_capture):
|
|
||||||
result, resp = dispatch(
|
|
||||||
Path(req.path), req.quality, req.maxsize, req.maxzoom, data
|
|
||||||
)
|
|
||||||
if not resp.ok:
|
|
||||||
captured = stderr_capture.getvalue().strip()
|
|
||||||
if captured:
|
|
||||||
resp = PreviewResponse(
|
|
||||||
ok=False,
|
|
||||||
backend=resp.backend,
|
|
||||||
error=resp.error,
|
|
||||||
stderr=captured,
|
|
||||||
)
|
|
||||||
_write_response(resp, result or b"")
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception("Preview worker error for %s", req.path)
|
|
||||||
captured = stderr_capture.getvalue().strip()
|
|
||||||
_write_response(
|
|
||||||
PreviewResponse(ok=False, error=str(e), stderr=captured or None), b""
|
|
||||||
)
|
|
||||||
finally:
|
|
||||||
root_logger.removeHandler(handler)
|
|
||||||
handler.close()
|
|
||||||
|
|
||||||
|
|
||||||
def main() -> None:
|
|
||||||
# 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()
|
|
||||||
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.
|
|
||||||
_protocol_out.write(b"\x01")
|
|
||||||
_protocol_out.flush()
|
|
||||||
_run_loop()
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
main()
|
|
||||||
@@ -286,6 +286,9 @@ def configure_main_logging() -> None:
|
|||||||
LOGGING_CONFIG_DEFAULTS["formatters"]["generic"] = {
|
LOGGING_CONFIG_DEFAULTS["formatters"]["generic"] = {
|
||||||
"class": "cista.sanic_logging._EmojiFormatter",
|
"class": "cista.sanic_logging._EmojiFormatter",
|
||||||
}
|
}
|
||||||
|
# Silence websockets' built-in "connection closed" INFO messages; we log WS
|
||||||
|
# open/close ourselves in the custom access log instead.
|
||||||
|
logging.getLogger("websockets.server").setLevel(logging.WARNING)
|
||||||
# Also reformat any handlers already attached (covers the initial Sanic() call)
|
# Also reformat any handlers already attached (covers the initial Sanic() call)
|
||||||
for name in ("sanic.root", "sanic.error", "sanic.server", "sanic.websockets"):
|
for name in ("sanic.root", "sanic.error", "sanic.server", "sanic.websockets"):
|
||||||
for handler in logging.getLogger(name).handlers:
|
for handler in logging.getLogger(name).handlers:
|
||||||
|
|||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
services:
|
services:
|
||||||
onlyoffice:
|
onlyoffice:
|
||||||
build:
|
build:
|
||||||
context: ./docker/onlyoffice-converter-patch
|
context: ./mediapreview/mediapreview/docker
|
||||||
args:
|
args:
|
||||||
ONLYOFFICE_VERSION: "9.3.1"
|
ONLYOFFICE_VERSION: "9.3.1"
|
||||||
container_name: onlyoffice
|
container_name: onlyoffice
|
||||||
|
|||||||
@@ -1,70 +0,0 @@
|
|||||||
# Patched OnlyOffice Document Server with configurable converter worker count.
|
|
||||||
#
|
|
||||||
# The Community Edition hardcodes the document converter to 1 worker,
|
|
||||||
# which creates a severe bottleneck under concurrent load.
|
|
||||||
# This image patches the open-source license.js to spawn a configurable
|
|
||||||
# number of converter workers (default 8).
|
|
||||||
#
|
|
||||||
# Build:
|
|
||||||
# docker build -t onlyoffice-cista docker/onlyoffice-converter-patch
|
|
||||||
#
|
|
||||||
# Run:
|
|
||||||
# docker run -d -p 8988:80 \
|
|
||||||
# -e WORKERS=16 \
|
|
||||||
# -e JWT_SECRET=your-strong-secret \
|
|
||||||
# --name onlyoffice onlyoffice-cista
|
|
||||||
#
|
|
||||||
# JWT:
|
|
||||||
# Set JWT_SECRET to the same value you pass to Cista as ONLYOFFICE_JWT_SECRET.
|
|
||||||
# OnlyOffice will enable token validation automatically.
|
|
||||||
#
|
|
||||||
# The ONLYOFFICE_VERSION build arg lets you target a specific release.
|
|
||||||
|
|
||||||
ARG ONLYOFFICE_VERSION=9.3.1
|
|
||||||
|
|
||||||
FROM onlyoffice/documentserver:${ONLYOFFICE_VERSION}
|
|
||||||
|
|
||||||
# Prevent interactive apt prompts
|
|
||||||
ENV DEBIAN_FRONTEND=noninteractive
|
|
||||||
|
|
||||||
# Install Node.js, npm, and git so we can run the FileConverter from source.
|
|
||||||
RUN apt-get update -qq && \
|
|
||||||
apt-get install -y -qq --no-install-recommends \
|
|
||||||
nodejs \
|
|
||||||
npm \
|
|
||||||
git \
|
|
||||||
ca-certificates && \
|
|
||||||
rm -rf /var/lib/apt/lists/*
|
|
||||||
|
|
||||||
# Clone the open-source server components (shallow, ~15 MB).
|
|
||||||
# The master branch is used because the Linux/web tags are not published
|
|
||||||
# in the server repo; the license.js file has been stable for years.
|
|
||||||
RUN git clone --depth 1 https://github.com/ONLYOFFICE/server.git /opt/oo-server
|
|
||||||
|
|
||||||
# Patch license.js so the converter worker count is read from an env var
|
|
||||||
# instead of being hardcoded to 1.
|
|
||||||
RUN sed -i \
|
|
||||||
's/count: 1,/count: parseInt(process.env.WORKERS, 10) || 8,/' \
|
|
||||||
/opt/oo-server/Common/sources/license.js
|
|
||||||
|
|
||||||
# Install npm dependencies for the modules the FileConverter touches.
|
|
||||||
# DocService deps are also needed because converter.js pulls in baseConnector.
|
|
||||||
RUN cd /opt/oo-server/Common && npm ci --no-audit --no-fund
|
|
||||||
RUN cd /opt/oo-server/FileConverter && npm ci --no-audit --no-fund
|
|
||||||
RUN cd /opt/oo-server/DocService && npm ci --no-audit --no-fund
|
|
||||||
|
|
||||||
# Back up the compiled pkg binary and replace it with our wrapper.
|
|
||||||
RUN mv /var/www/onlyoffice/documentserver/server/FileConverter/converter \
|
|
||||||
/var/www/onlyoffice/documentserver/server/FileConverter/converter.orig
|
|
||||||
|
|
||||||
COPY converter-wrapper.sh /var/www/onlyoffice/documentserver/server/FileConverter/converter
|
|
||||||
RUN chmod +x /var/www/onlyoffice/documentserver/server/FileConverter/converter
|
|
||||||
|
|
||||||
# Default worker count (override at runtime with -e WORKERS=16).
|
|
||||||
ENV WORKERS=8
|
|
||||||
|
|
||||||
# Use our custom entrypoint to persist the env var to a file that the
|
|
||||||
# non-root converter process (user=ds) can read.
|
|
||||||
COPY entrypoint.sh /app/ds/run-document-server-patched.sh
|
|
||||||
RUN chmod +x /app/ds/run-document-server-patched.sh
|
|
||||||
ENTRYPOINT ["/app/ds/run-document-server-patched.sh"]
|
|
||||||
@@ -1,19 +0,0 @@
|
|||||||
#!/bin/bash
|
|
||||||
# Wrapper that runs the OnlyOffice FileConverter from patched Node.js source.
|
|
||||||
# Replaces the compiled pkg binary shipped with the Community Edition.
|
|
||||||
|
|
||||||
# The env var is not passed through supervisor to the 'ds' user, so we read
|
|
||||||
# it from a file written by the custom entrypoint.
|
|
||||||
if [ -z "${WORKERS}" ] && [ -r /tmp/oo-converter-workers.txt ]; then
|
|
||||||
export WORKERS=$(cat /tmp/oo-converter-workers.txt)
|
|
||||||
fi
|
|
||||||
|
|
||||||
cd /opt/oo-server/FileConverter || exit 1
|
|
||||||
|
|
||||||
export NODE_ENV=production-linux
|
|
||||||
export NODE_CONFIG_DIR=/etc/onlyoffice/documentserver
|
|
||||||
export NODE_DISABLE_COLORS=1
|
|
||||||
export APPLICATION_NAME=onlyoffice
|
|
||||||
export LD_LIBRARY_PATH=/var/www/onlyoffice/documentserver/server/FileConverter/bin
|
|
||||||
|
|
||||||
exec node sources/convertermaster.js "$@"
|
|
||||||
@@ -1,8 +0,0 @@
|
|||||||
#!/bin/bash
|
|
||||||
# Custom entrypoint that persists WORKERS to a file readable by
|
|
||||||
# the non-root user that supervisor uses to run the converter.
|
|
||||||
|
|
||||||
echo "${WORKERS:-8}" > /tmp/oo-converter-workers.txt
|
|
||||||
chmod 644 /tmp/oo-converter-workers.txt
|
|
||||||
|
|
||||||
exec /app/ds/run-document-server.sh "$@"
|
|
||||||
+2
-1
@@ -33,6 +33,7 @@ dependencies = [
|
|||||||
"html5tagger>=1.3.0",
|
"html5tagger>=1.3.0",
|
||||||
"httpx>=0.28.0",
|
"httpx>=0.28.0",
|
||||||
"inotify>=0.2.12",
|
"inotify>=0.2.12",
|
||||||
|
"mediapreview[standard]",
|
||||||
"msgspec>=0.19.0",
|
"msgspec>=0.19.0",
|
||||||
"natsort>=8.4.0",
|
"natsort>=8.4.0",
|
||||||
"numpy>=2.3.2",
|
"numpy>=2.3.2",
|
||||||
@@ -77,7 +78,7 @@ docs = [
|
|||||||
source = "vcs"
|
source = "vcs"
|
||||||
|
|
||||||
[tool.hatch.build]
|
[tool.hatch.build]
|
||||||
artifacts = ["cista/frontend-build", "cista/docker"]
|
artifacts = ["cista/frontend-build"]
|
||||||
targets.sdist.hooks.custom.path = "scripts/fastapi-vue/buildhook.py"
|
targets.sdist.hooks.custom.path = "scripts/fastapi-vue/buildhook.py"
|
||||||
targets.sdist.include = [
|
targets.sdist.include = [
|
||||||
"/cista",
|
"/cista",
|
||||||
|
|||||||
@@ -1,175 +0,0 @@
|
|||||||
"""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