penpot/media-processor/test/queue.test.ts
Andrey Antukh aeedb96260
Add media-processor service for image and font processing (#10767)
*  Add media-processor service for image and font processing

Externalizes ImageMagick and FontForge subprocess invocations into a
separate Node.js HTTP service (media-processor/). Backend dispatches
via feature flag :use-remote-media-processing.

Key changes:
- media-processor module (TypeScript, Express 5, Sharp, FontForge/woff)
  - POST /api/image/info, /api/image/thumbnail, /api/font/generate
  - Resource limits: 128MP rejection, prlimit (512MB + 30s CPU)
  - Streaming multipart via SequenceInputStream
- app.media split into validation (leaf), local (shell impls), remote (HTTP)
- Schema enforcement: :upload and :input schemas in validation namespace
- Configurable timeout (PENPOT_MEDIA_PROCESSING_SERVICE_TIMEOUT)
- 78 tests across 4 files (image, font, middleware, config)
- FontForge path escaping for command injection prevention
- Parallel font variant conversions with Promise.all

AI-assisted-by: mimo-v2.5-pro

* 🐳 Revert docker-compose changes from media-processor commit

Remove docker-compose.yaml modifications that were part of the media-processor
service commit. The media-processor service definition, flags, and environment
variables are reverted to their previous state.

AI-assisted-by: qwen3.7-plus

* ⬆️ Update dependencies

* 🐛 Fix PR review issues in media-processor

- Font path bug: sfntToWoff and woff2ToSfnt now copy input to temp dir
  when input is a file path, ensuring output lands in expected location
- Error preservation: execCommand preserves killed/signal/code properties
  from child process errors for OOM detection
- Content-Length: service-multipart-request calculates and includes
  Content-Length header for streaming multipart requests

AI-assisted-by: qwen3.7-plus

* 🐛 Fix code review issues in media-processor

- Rename PENPOT_MEDIA_PROCESSOR_SECRET_KEY to PENPOT_MEDIA_PROCESSOR_SHARED_KEY
  in devenv to match backend config key
- Fix timeout middleware to destroy request AFTER response finishes,
  preventing truncated 504 responses
- Fix quality=0 parsing to preserve explicit zero (was silently overridden to 85)
- Replace require('fs') with proper ES module import in upload-storage.ts
- Refactor font conversion temp-dir boilerplate into withTempInput helper
- Document FontForge escaping limitations (single quotes only)
- Fix misleading comment in image.ts about sharp metadata decoding

AI-assisted-by: qwen3.7-plus

* 🐛 Fix code review issues in media-processor (round 2)

- Fix queue middleware to skip next() when response already ended,
  preventing orphaned work after timeout
- Fix hybrid storage to use disk when Content-Length is absent (chunked
  transfer), preventing unbounded memory allocation
- Add source image format validation in generateThumbnail to reject
  unsupported formats (TIFF, BMP, etc.) with 400 instead of 500
- Remove dead code in convertFont for unreachable woff→woff path
- Remove unused isEnabled() method from LokiLogTransport
- Fix sfntToWoff to use correct extension (.ttf/.otf) based on source type
- Extract queue middleware to separate file for testability
- Add comprehensive tests for queue middleware and upload storage

AI-assisted-by: qwen3.7-plus

* 🐛 Fix code review issues in media-processor (round 3)

- Fix disk-backed upload cleanup after successful requests by adding
  cleanup middleware that removes temp files on response finish/close
- Wrap sharp metadata/decoding errors as 400 validation errors instead
  of 500 internal errors
- Only apply flatten() for JPEG output to preserve alpha channel in
  PNG and WebP outputs

AI-assisted-by: qwen3.7-plus

*  Add comprehensive tests for media-processor

Phase 1 - Cleanup verification:
- Add cleanup middleware unit tests (6 tests)
- Add HTTP upload cleanup integration tests (5 tests)

Phase 2 - Error handling & alpha preservation:
- Add sharp error wrapping tests (4 tests)
- Add HTTP malformed image tests (2 tests)
- Add alpha preservation tests (3 tests)

Phase 3 - Edge cases:
- Add upload storage edge case tests (3 tests)
- Add queue middleware edge case tests (4 tests)

Phase 4 - Backend mock verification:
- Fix backend mocks to include :mtype field in image info responses
- Verify all error codes match actual service behavior

Total: 27 new tests added (160 tests passing)

AI-assisted-by: qwen3.7-plus

* 🐛 Fix code review issues in media-processor (round 4)

- Add Zod validation constraints for config values (int, positive, min)
- Fix auth middleware to compare Buffer byte lengths instead of string lengths
- Validate requested output dimensions in generateThumbnail (crop mode)
- Change queue middleware to release slot via callback in finally block
- Add comprehensive tests for all fixes

AI-assisted-by: qwen3.7-plus

* 🐛 Close HTTP response streams in backend media remote

- Wrap stream consumption in try/finally with .close() calls
- Add tests to verify stream closure for info, font-convert, and thumbnail

AI-assisted-by: qwen3.7-plus

* 🐛 Fix queue slot leak on upload failures

Make releaseQueue idempotent and attach fallback listener to release
slot when response finishes. This covers Multer errors that bypass
the route handler's finally block, preventing permanent queue stall.

AI-assisted-by: qwen3.7-plus

* 🐛 Cancel processing on timeout

Create AbortController in timeout middleware and abort signal when
timeout fires. Pass signal to Sharp and FontForge to cancel ongoing
processing and release resources when request is cancelled.

AI-assisted-by: qwen3.7-plus

* 🐛 Fix code review issues in media-processor (round 6)

- Error handler: check headersSent before writing response to prevent
  ERR_HTTP_HEADERS_SENT when timeout already sent 504
- Timeout config: increase default requestTimeout from 60s to 180s to
  match font processing timeout (120s) and backend request timeout
- Image processing: check abort signal before starting Sharp operations
  to cancel processing when timeout fires
- Queue lifecycle: remove res.on('close', release) fallback to hold
  queue slot until processing completes, preventing concurrency limit
  violation when client disconnects

AI-assisted-by: qwen3.7-plus

* 🐛 Close HTTP response stream in download-image

Wrap response body in with-open to ensure stream is closed after
writing to temp file, preventing HTTP connection leaks on repeated
URL imports.

AI-assisted-by: qwen3.7-plus

* 🐛 Close HTTP response stream on validation errors in download-image

Move with-open to wrap the entire validation and processing block,
ensuring the response body stream is closed even when validation fails
(non-2xx status, missing size, invalid media type). This prevents
HTTP connection leaks on repeated failed downloads.

Add test to verify stream closure on validation errors.

AI-assisted-by: qwen3.7-plus

* 🐛 Pass abort signal to Sharp toBuffer for timeout cancellation

Wrap Sharp's toBuffer() with Promise.race to check abort signal during
processing. This ensures large thumbnails stop processing when the
request times out, preventing wasted CPU/memory and queue capacity.

Add test to verify abort during toBuffer operation.

AI-assisted-by: qwen3.7-plus

* 🐛 Hold queue slot until Sharp completes and handle client disconnect

- Remove Promise.race from generateThumbnail — Sharp processing now
  completes fully before queue slot is released, preventing concurrency
  limit violations under timeout conditions
- Remove res.on("finish", release) fallback from queue middleware —
  error handler now explicitly calls releaseQueue in all error paths
- Add res.on("close") handler in timeout middleware to abort signal
  when client disconnects, ensuring processing stops early
- Add tests for client disconnect handling and queue slot lifecycle

AI-assisted-by: qwen3.7-plus

* 🐛 Address round 9 review findings

- Document Sharp 0.35.3 cancellation limitation in image.ts
- Add integration test for timeout cleanup with large images
- Fix font tools (sfntToWoff, woffToSfnt, woff2ToSfnt) to throw
  ProcessingError on resource limit kills instead of returning null
- Validate font signatures for same-format conversions to prevent
  arbitrary files from being persisted as valid fonts
- Fix concurrent mkdtemp race in upload-storage by using shared
  initialization promise

AI-assisted-by: qwen3.7-plus

* 🐛 Address round 10 review findings

- Add tmpdir assertion in font.ts to prevent path injection
- Preserve original error in queue middleware catch handler
- Change auth middleware response type from "internal" to "authorization"
- Add cleanup flag to prevent double cleanup in cleanup middleware
- Move quality clamping into parseQuality function for consistency
- Add integration tests for quality parameter clamping at route level
- Update existing tests to match new auth response type

AI-assisted-by: qwen3.7-plus

* 🐛 Address round 11 review findings

- Extract releaseSlot helper in error-handler to reduce duplication
- Remove redundant try/catch in font.ts withTempDir cleanup
- Improve font path validation error message for clarity
- Move path validation before try/catch to prevent swallowing
- Add debug logging for cleanup failures in cleanup middleware
- Inline TransportTargetSpec type alias in logger.ts
- Extract logging middleware to separate file for consistency
- Remove duplicate MIME validation in image thumbnail route
- Add test for font path validation (outside tmpdir rejection)
- Add tests for error handler queue release across all branches

AI-assisted-by: qwen3.7-plus

* 🐛 Remove Content-Length header from multipart requests

The JDK's HttpClient rejects Content-Length as a restricted header,
causing IllegalArgumentException when sending multipart requests to the
media-processor. Remove the explicit Content-Length header and let the
JDK use chunked transfer encoding. The media-processor will use disk
storage for all multipart requests (safe default behavior).

Remove unused size computations (file-size, header-bytes, footer-bytes,
total-size) that were only used for Content-Length.

Update test to verify Content-Length is not present in request headers.

AI-assisted-by: qwen3.7-plus

* 🐛 Fix pino ESM bundling for media-processor

Mark pino and its transports (pino-pretty, pino-loki) as external to
avoid bundling issues with worker thread modules that reference
__dirname (not available in ES modules).

AI-assisted-by: qwen3.7-plus
2026-08-05 09:41:48 +02:00

378 lines
11 KiB
TypeScript

import { describe, it, expect, vi, beforeEach } from "vitest";
import { createQueueMiddleware } from "../src/middleware/queue.js";
import type { Request, Response, NextFunction } from "express";
import { EventEmitter } from "node:events";
function mockRes() {
const res = new EventEmitter() as any;
res.status = vi.fn().mockReturnThis();
res.json = vi.fn().mockReturnThis();
res.send = vi.fn().mockReturnThis();
res.headersSent = false;
res.writableEnded = false;
return res as Response;
}
function mockReq() {
return {} as Request;
}
describe("queueMiddleware", () => {
it("calls next() when queue has capacity", async () => {
const middleware = createQueueMiddleware(1);
const req = mockReq();
const res = mockRes();
const next = vi.fn();
middleware(req, res, next);
expect(next).toHaveBeenCalled();
});
it("skips next() when res.writableEnded is true (timeout already sent)", async () => {
const middleware = createQueueMiddleware(1);
const req = mockReq();
const res = mockRes();
(res as any).writableEnded = true;
const next = vi.fn();
middleware(req, res, next);
// next() should NOT be called because response already ended
expect(next).not.toHaveBeenCalled();
});
it("queues requests when concurrency limit reached", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
// First request takes the slot
middleware(req1, res1, next1);
expect(next1).toHaveBeenCalled();
// Second request should queue
middleware(req2, res2, next2);
expect(next2).not.toHaveBeenCalled();
// Release first request's queue slot (simulating processing completion)
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
// Now second request should proceed
await new Promise((resolve) => setTimeout(resolve, 10));
expect(next2).toHaveBeenCalled();
});
it("resolves promise when releaseQueue is called", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// Release first request's queue slot
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
expect(next2).toHaveBeenCalled();
});
it("processes requests sequentially with concurrency 1", async () => {
const middleware = createQueueMiddleware(1);
const order: number[] = [];
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn(() => order.push(1));
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn(() => order.push(2));
const req3 = mockReq();
const res3 = mockRes();
const next3 = vi.fn(() => order.push(3));
middleware(req1, res1, next1);
middleware(req2, res2, next2);
middleware(req3, res3, next3);
// Only first should be called immediately
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
expect(next3).not.toHaveBeenCalled();
// Release first request's queue slot
const releaseQueue1 = (res1 as any).locals?.releaseQueue;
expect(releaseQueue1).toBeDefined();
releaseQueue1();
await new Promise((resolve) => setTimeout(resolve, 10));
// Now second should be called
expect(next2).toHaveBeenCalled();
expect(next3).not.toHaveBeenCalled();
// Release second request's queue slot
const releaseQueue2 = (res2 as any).locals?.releaseQueue;
expect(releaseQueue2).toBeDefined();
releaseQueue2();
await new Promise((resolve) => setTimeout(resolve, 10));
// Now third should be called
expect(next3).toHaveBeenCalled();
// Verify sequential order
expect(order).toEqual([1, 2, 3]);
});
it("processes requests in parallel with concurrency 10", async () => {
const middleware = createQueueMiddleware(10);
const calls: number[] = [];
// Create 5 requests (less than concurrency limit)
const requests = Array.from({ length: 5 }, (_, i) => {
const req = mockReq();
const res = mockRes();
const next = vi.fn(() => calls.push(i));
return { req, res, next };
});
// All should be called immediately
requests.forEach(({ req, res, next }) => {
middleware(req, res, next);
});
// All 5 should be called immediately since concurrency is 10
expect(calls.length).toBe(5);
expect(calls).toEqual([0, 1, 2, 3, 4]);
});
it("handles request errors gracefully", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// Simulate error by releasing queue slot (as would happen in finally block)
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Second request should still proceed even after first "error"
expect(next2).toHaveBeenCalled();
});
it("doesn't block on slow requests within concurrency limit", async () => {
const middleware = createQueueMiddleware(2);
const calls: number[] = [];
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn(() => calls.push(1));
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn(() => calls.push(2));
const req3 = mockReq();
const res3 = mockRes();
const next3 = vi.fn(() => calls.push(3));
// Start first two requests (concurrency is 2)
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// Both should be called immediately
expect(next1).toHaveBeenCalled();
expect(next2).toHaveBeenCalled();
expect(next3).not.toHaveBeenCalled();
// Third request should wait
middleware(req3, res3, next3);
expect(next3).not.toHaveBeenCalled();
// Release first request's queue slot
const releaseQueue1 = (res1 as any).locals?.releaseQueue;
expect(releaseQueue1).toBeDefined();
releaseQueue1();
await new Promise((resolve) => setTimeout(resolve, 10));
// Now third should proceed
expect(next3).toHaveBeenCalled();
});
it("releases slot when error handler calls releaseQueue (covers Multer error path)", async () => {
// This test verifies that the error handler releases the queue slot
// by calling releaseQueue from res.locals. This covers the Multer error
// case where the route handler never runs.
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// First request is processing, second is queued
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
// Simulate error handler calling releaseQueue (e.g., Multer error)
const releaseQueue = (res1 as any).locals.releaseQueue;
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Second request should proceed because slot was released
expect(next2).toHaveBeenCalled();
});
it("releases slot when releaseQueue callback is called", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// First request is processing, second is queued
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
// Simulate processing completing by calling releaseQueue
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Now second request should proceed
expect(next2).toHaveBeenCalled();
});
it("releases slot on processing error (via finally block)", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// First request is processing, second is queued
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
// Simulate processing error and release in finally block
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Second request should proceed even after error
expect(next2).toHaveBeenCalled();
});
it("releaseQueue is idempotent (can be called multiple times)", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// First request is processing, second is queued
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
// Call releaseQueue multiple times
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
releaseQueue(); // Should not throw or cause issues
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Second request should proceed
expect(next2).toHaveBeenCalled();
});
it("holds queue slot when client disconnects (close event)", async () => {
const middleware = createQueueMiddleware(1);
const req1 = mockReq();
const res1 = mockRes();
const next1 = vi.fn();
const req2 = mockReq();
const res2 = mockRes();
const next2 = vi.fn();
middleware(req1, res1, next1);
middleware(req2, res2, next2);
// First request is processing, second is queued
expect(next1).toHaveBeenCalled();
expect(next2).not.toHaveBeenCalled();
// Simulate client disconnect (close event)
res1.emit("close");
await new Promise((resolve) => setTimeout(resolve, 10));
// Second request should NOT proceed because slot is still held
expect(next2).not.toHaveBeenCalled();
// Now release the slot (simulating processing completion)
const releaseQueue = (res1 as any).locals?.releaseQueue;
expect(releaseQueue).toBeDefined();
releaseQueue();
await new Promise((resolve) => setTimeout(resolve, 10));
// Now second request should proceed
expect(next2).toHaveBeenCalled();
});
});