Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
041090cce9 |
@@ -145,6 +145,8 @@ def convert_to_png(file_path: Path, timeout: float = 30.0) -> bytes:
|
|||||||
headers = {"Content-Type": "application/json"}
|
headers = {"Content-Type": "application/json"}
|
||||||
token = _build_jwt_token(payload)
|
token = _build_jwt_token(payload)
|
||||||
if token:
|
if token:
|
||||||
|
# Conversion API expects JWT in request body when token checks are enabled.
|
||||||
|
payload["token"] = token
|
||||||
headers["Authorization"] = token
|
headers["Authorization"] = token
|
||||||
|
|
||||||
req = urllib.request.Request( # noqa: S310
|
req = urllib.request.Request( # noqa: S310
|
||||||
|
|||||||
+113
-43
@@ -90,7 +90,7 @@ class PreviewCache:
|
|||||||
# Global preview cache instance
|
# Global preview cache instance
|
||||||
_preview_cache = PreviewCache(capacity=500)
|
_preview_cache = PreviewCache(capacity=500)
|
||||||
|
|
||||||
PREVIEW_TIMEOUT = 3.0 # seconds until preview subprocess is killed
|
PREVIEW_TIMEOUT = 10.0 # seconds until preview subprocess is killed
|
||||||
PREVIEW_WORKERS = max(2, min(8, cpu_count()))
|
PREVIEW_WORKERS = max(2, min(8, cpu_count()))
|
||||||
_active_procs: set[asyncio.subprocess.Process] = set()
|
_active_procs: set[asyncio.subprocess.Process] = set()
|
||||||
_preview_pool = None
|
_preview_pool = None
|
||||||
@@ -164,7 +164,12 @@ class _PreviewWorkerPool:
|
|||||||
def __init__(self, size: int):
|
def __init__(self, size: int):
|
||||||
self.size = size
|
self.size = size
|
||||||
self._idle: asyncio.Queue[_PreviewWorker] = asyncio.Queue()
|
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._workers: set[_PreviewWorker] = set()
|
||||||
|
self._dispatchers: list[asyncio.Task] = []
|
||||||
|
self._seq = 0
|
||||||
self._closed = False
|
self._closed = False
|
||||||
|
|
||||||
async def _spawn_worker(self) -> _PreviewWorker:
|
async def _spawn_worker(self) -> _PreviewWorker:
|
||||||
@@ -195,61 +200,108 @@ class _PreviewWorkerPool:
|
|||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("Failed to replace preview worker")
|
logger.exception("Failed to replace preview worker")
|
||||||
|
|
||||||
|
async def _dispatch_loop(self) -> None:
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
_priority, _seq, future, args = await self._pending.get()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
return
|
||||||
|
|
||||||
|
if future.cancelled():
|
||||||
|
continue
|
||||||
|
|
||||||
|
worker = await self._idle.get()
|
||||||
|
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 timeout (%ds) for %s", int(PREVIEW_TIMEOUT), filepath.name
|
||||||
|
)
|
||||||
|
if not future.done():
|
||||||
|
future.set_exception(PreviewTimeoutError(filepath.name))
|
||||||
|
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}"
|
||||||
|
)
|
||||||
|
)
|
||||||
|
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:
|
async def start(self) -> None:
|
||||||
for _ in range(self.size):
|
for _ in range(self.size):
|
||||||
await self._add_worker()
|
await self._add_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):
|
async def run(self, filepath, quality: int, maxsize: int, maxzoom: float):
|
||||||
if self._closed:
|
if self._closed:
|
||||||
raise PreviewError("preview worker pool closed")
|
raise PreviewError("preview worker pool closed")
|
||||||
worker = await self._idle.get()
|
loop = asyncio.get_running_loop()
|
||||||
replace = False
|
future = loop.create_future()
|
||||||
try:
|
self._seq += 1
|
||||||
out, resp = await asyncio.wait_for(
|
await self._pending.put(
|
||||||
worker.request(filepath, quality, maxsize, maxzoom),
|
(
|
||||||
timeout=PREVIEW_TIMEOUT,
|
_preview_job_priority(filepath),
|
||||||
|
self._seq,
|
||||||
|
future,
|
||||||
|
(filepath, quality, maxsize, maxzoom),
|
||||||
)
|
)
|
||||||
return out, resp
|
)
|
||||||
except TimeoutError:
|
return await future
|
||||||
replace = True
|
|
||||||
logger.warning(
|
|
||||||
"Preview timeout (%ds) for %s", int(PREVIEW_TIMEOUT), filepath.name
|
|
||||||
)
|
|
||||||
raise PreviewTimeoutError(filepath.name) from None
|
|
||||||
except WorkerChecksumError as e:
|
|
||||||
replace = True
|
|
||||||
logger.error("Preview checksum mismatch for %s", filepath.name)
|
|
||||||
raise PreviewError(f"worker checksum mismatch for {filepath.name}") from e
|
|
||||||
except PreviewError:
|
|
||||||
raise
|
|
||||||
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
|
|
||||||
)
|
|
||||||
raise PreviewError(
|
|
||||||
f"worker protocol failure for {filepath.name}: {e}"
|
|
||||||
) from e
|
|
||||||
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 close(self) -> None:
|
async def close(self) -> None:
|
||||||
self._closed = True
|
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)
|
workers = list(self._workers)
|
||||||
self._workers.clear()
|
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():
|
while not self._idle.empty():
|
||||||
try:
|
try:
|
||||||
self._idle.get_nowait()
|
self._idle.get_nowait()
|
||||||
@@ -372,6 +424,24 @@ OFFICE_PREVIEW_SUFFIXES = {
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
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 is_previewable_path(path) -> bool:
|
def is_previewable_path(path) -> bool:
|
||||||
suffix = path.suffix.lower()
|
suffix = path.suffix.lower()
|
||||||
if suffix in DOC_PREVIEW_SUFFIXES or suffix in OFFICE_PREVIEW_SUFFIXES:
|
if suffix in DOC_PREVIEW_SUFFIXES or suffix in OFFICE_PREVIEW_SUFFIXES:
|
||||||
|
|||||||
Reference in New Issue
Block a user