From 360e0ee7a6290357bd7d52e15685bf8cba1ab55a Mon Sep 17 00:00:00 2001 From: Anthony GAEREMYNCK <1@anthony.sh> Date: Thu, 30 Jul 2026 09:21:42 +0200 Subject: [PATCH] feat: add worker process with bounded concurrency and per-job timeout Co-Authored-By: Claude Sonnet 5 --- src/worker.js | 78 +++++++++++++++++++++++++++++++++++ test/worker.test.js | 99 +++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 177 insertions(+) create mode 100644 src/worker.js create mode 100644 test/worker.test.js diff --git a/src/worker.js b/src/worker.js new file mode 100644 index 0000000..d5c4cea --- /dev/null +++ b/src/worker.js @@ -0,0 +1,78 @@ +import { loadConfig } from './config.js'; +import { getPool } from './db.js'; +import { ensureStorageDirs, uploadPath, outputPath, deleteIfExists } from './storage.js'; +import { outputMimeType } from './mime.js'; +import { resolve as resolveConverter } from './converters/registry.js'; +import { registerImageConverters } from './converters/image.js'; +import { registerImageToPdfConverter } from './converters/imageToPdf.js'; +import { registerDocumentConverters } from './converters/document.js'; +import { findPendingJobs, markProcessing, markDone, markFailed } from './jobs/jobRepository.js'; + +const JOB_TIMEOUT_MS = 60000; + +function withTimeout(promise, ms) { + let timeoutId; + const timeout = new Promise((_, reject) => { + timeoutId = setTimeout(() => reject(new Error(`Conversion timed out after ${ms}ms`)), ms); + }); + return Promise.race([promise, timeout]).finally(() => clearTimeout(timeoutId)); +} + +async function processJob(pool, config, job) { + await markProcessing(pool, job.id); + + const inputFilePath = uploadPath(config, job.id, job.sourceFormat); + const outputFilePath = outputPath(config, job.id, job.targetFormat); + + try { + const entry = resolveConverter(job.sourceFormat, job.targetFormat); + if (!entry) { + throw new Error(`No converter registered for ${job.sourceFormat} -> ${job.targetFormat}`); + } + + await withTimeout(entry.convert(inputFilePath, outputFilePath), JOB_TIMEOUT_MS); + + await markDone(pool, job.id, { + outputPath: `${job.id}.${job.targetFormat}`, + outputMimeType: outputMimeType(job.targetFormat), + }); + } catch (error) { + await deleteIfExists(outputFilePath); + await markFailed(pool, job.id, { + errorMessage: 'Conversion failed, please try again.', + errorLog: error.stack ?? String(error), + }); + } +} + +export async function processPendingJobs(pool, config) { + const jobs = await findPendingJobs(pool, config.workerConcurrency); + await Promise.all(jobs.map((job) => processJob(pool, config, job))); + return jobs.length; +} + +export function startWorker(pool, config) { + const interval = setInterval(() => { + processPendingJobs(pool, config).catch((error) => { + console.error('Error while processing pending jobs:', error); + }); + }, config.workerPollIntervalMs); + return () => clearInterval(interval); +} + +async function main() { + const config = loadConfig(); + await ensureStorageDirs(config); + const pool = getPool(config); + + registerImageConverters(); + registerImageToPdfConverter(); + registerDocumentConverters(); + + startWorker(pool, config); + console.log(`Worker started, polling every ${config.workerPollIntervalMs}ms`); +} + +if (import.meta.url === `file://${process.argv[1]}`) { + main(); +} diff --git a/test/worker.test.js b/test/worker.test.js new file mode 100644 index 0000000..c928217 --- /dev/null +++ b/test/worker.test.js @@ -0,0 +1,99 @@ +import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest'; +import fs from 'node:fs/promises'; +import path from 'node:path'; +import os from 'node:os'; +import { getPool, closePool } from '../src/db.js'; +import { loadConfig } from '../src/config.js'; +import { ensureStorageDirs, uploadPath, outputPath } from '../src/storage.js'; +import { createJob, getJobById, getJobErrorLog } from '../src/jobs/jobRepository.js'; +import { registerImageConverters } from '../src/converters/image.js'; +import { processPendingJobs } from '../src/worker.js'; + +let pool; +let config; + +beforeAll(async () => { + registerImageConverters(); + config = { ...loadConfig(), storageDir: await fs.mkdtemp(path.join(os.tmpdir(), 'converter-worker-')) }; + await ensureStorageDirs(config); + pool = getPool(config); +}); + +afterAll(async () => { + await closePool(); + await fs.rm(config.storageDir, { recursive: true, force: true }); +}); + +beforeEach(async () => { + await pool.query('DELETE FROM conversion_jobs'); +}); + +async function createPendingImageJob(id, sourceFormat, targetFormat) { + const fixturePath = path.join(import.meta.dirname, 'fixtures', 'sample.png'); + const inputFilePath = uploadPath(config, id, sourceFormat); + await fs.copyFile(fixturePath, inputFilePath); + + await createJob(pool, { + id, + family: 'image', + sourceFormat, + targetFormat, + originalFilename: `photo.${sourceFormat}`, + inputPath: `${id}.${sourceFormat}`, + inputMimeType: 'image/png', + expiresAt: new Date(Date.now() + 3600 * 1000), + }); +} + +describe('processPendingJobs', () => { + it('converts a pending image job to done', async () => { + const id = '99999999-9999-4999-8999-999999999999'; + await createPendingImageJob(id, 'png', 'webp'); + + const processedCount = await processPendingJobs(pool, config); + + expect(processedCount).toBe(1); + const job = await getJobById(pool, id); + expect(job.status).toBe('done'); + expect(job.outputPath).toBe(`${id}.webp`); + expect(job.outputMimeType).toBe('image/webp'); + + const stat = await fs.stat(outputPath(config, id, 'webp')); + expect(stat.size).toBeGreaterThan(0); + }); + + it('marks a job failed with a safe message and a detailed log when the converter throws', async () => { + const id = 'aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa'; + await createJob(pool, { + id, + family: 'image', + sourceFormat: 'png', + targetFormat: 'webp', + originalFilename: 'missing.png', + inputPath: `${id}.png`, + inputMimeType: 'image/png', + expiresAt: new Date(Date.now() + 3600 * 1000), + }); + // Note: input file is intentionally never written, so sharp will throw ENOENT. + + await processPendingJobs(pool, config); + + const job = await getJobById(pool, id); + expect(job.status).toBe('failed'); + expect(job.errorMessage).toBe('Conversion failed, please try again.'); + + const errorLog = await getJobErrorLog(pool, id); + expect(errorLog).toMatch(/input file is missing/i); + }); + + it('only picks up as many jobs as workerConcurrency allows', async () => { + await createPendingImageJob('bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb', 'png', 'webp'); + await createPendingImageJob('cccccccc-cccc-4ccc-8ccc-cccccccccccc', 'png', 'webp'); + await createPendingImageJob('dddddddd-dddd-4ddd-8ddd-dddddddddddd', 'png', 'webp'); + + const limitedConfig = { ...config, workerConcurrency: 2 }; + const processedCount = await processPendingJobs(pool, limitedConfig); + + expect(processedCount).toBe(2); + }); +});