diff --git a/admin/app/jobs/check_service_updates_job.ts b/admin/app/jobs/check_service_updates_job.ts index 6fb7335..58be73c 100644 --- a/admin/app/jobs/check_service_updates_job.ts +++ b/admin/app/jobs/check_service_updates_job.ts @@ -95,7 +95,7 @@ export class CheckServiceUpdatesJob { } static async scheduleNightly() { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) await queue.upsertJobScheduler( @@ -114,7 +114,7 @@ export class CheckServiceUpdatesJob { } static async dispatch() { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const job = await queue.add( diff --git a/admin/app/jobs/check_update_job.ts b/admin/app/jobs/check_update_job.ts index 046d9c0..aaac08d 100644 --- a/admin/app/jobs/check_update_job.ts +++ b/admin/app/jobs/check_update_job.ts @@ -42,7 +42,7 @@ export class CheckUpdateJob { } static async scheduleNightly() { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) await queue.upsertJobScheduler( @@ -61,7 +61,7 @@ export class CheckUpdateJob { } static async dispatch() { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const job = await queue.add(this.key, {}, { diff --git a/admin/app/jobs/download_model_job.ts b/admin/app/jobs/download_model_job.ts index f189021..2ba0080 100644 --- a/admin/app/jobs/download_model_job.ts +++ b/admin/app/jobs/download_model_job.ts @@ -34,7 +34,7 @@ export class DownloadModelJob { /** Signal cancellation via Redis so the worker process can pick it up on its next poll tick */ static async signalCancel(jobId: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const client = await queue.client await client.set(this.cancelKey(jobId), '1', 'EX', 300) // 5 min TTL @@ -66,7 +66,7 @@ export class DownloadModelJob { DownloadModelJob.abortControllers.set(job.id!, abortController) // Get Redis client for checking cancel signals from the API process - const queueService = new QueueService() + const queueService = QueueService.getInstance() const cancelRedis = await queueService.getQueue(DownloadModelJob.queue).client // Track whether cancellation was explicitly requested by the user. Only user-initiated @@ -154,14 +154,14 @@ export class DownloadModelJob { } static async getByModelName(modelName: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(modelName) return await queue.getJob(jobId) } static async dispatch(params: DownloadModelJobParams) { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(params.modelName) diff --git a/admin/app/jobs/embed_file_job.ts b/admin/app/jobs/embed_file_job.ts index c7566fa..771cedf 100644 --- a/admin/app/jobs/embed_file_job.ts +++ b/admin/app/jobs/embed_file_job.ts @@ -184,7 +184,7 @@ export class EmbedFileJob { } static async listActiveJobs(): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobs = await queue.getJobs(['waiting', 'active', 'delayed']) @@ -198,14 +198,14 @@ export class EmbedFileJob { } static async getByFilePath(filePath: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(filePath) return await queue.getJob(jobId) } static async dispatch(params: EmbedFileJobParams) { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) // Continuation batches (batchOffset > 0) must NOT reuse the deterministic @@ -267,7 +267,7 @@ export class EmbedFileJob { } static async listFailedJobs(): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) // Jobs that have failed at least once are in 'delayed' (retrying) or terminal 'failed' state. // We identify them by job.data.status === 'failed' set in the catch block of handle(). @@ -286,7 +286,7 @@ export class EmbedFileJob { } static async cleanupFailedJobs(): Promise<{ cleaned: number; filesDeleted: number }> { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const allJobs = await queue.getJobs(['waiting', 'delayed', 'failed']) const failedJobs = allJobs.filter((job) => (job.data as any).status === 'failed') diff --git a/admin/app/jobs/run_benchmark_job.ts b/admin/app/jobs/run_benchmark_job.ts index 0ae41e8..962e663 100644 --- a/admin/app/jobs/run_benchmark_job.ts +++ b/admin/app/jobs/run_benchmark_job.ts @@ -53,7 +53,7 @@ export class RunBenchmarkJob { } static async dispatch(params: RunBenchmarkJobParams) { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) try { @@ -89,7 +89,7 @@ export class RunBenchmarkJob { } static async getJob(benchmarkId: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) return await queue.getJob(benchmarkId) } diff --git a/admin/app/jobs/run_download_job.ts b/admin/app/jobs/run_download_job.ts index 12b3532..5fad2b6 100644 --- a/admin/app/jobs/run_download_job.ts +++ b/admin/app/jobs/run_download_job.ts @@ -31,7 +31,7 @@ export class RunDownloadJob { /** Signal cancellation via Redis so the worker process can pick it up */ static async signalCancel(jobId: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const client = await queue.client await client.set(this.cancelKey(jobId), '1', 'EX', 300) // 5 min TTL @@ -46,7 +46,7 @@ export class RunDownloadJob { RunDownloadJob.abortControllers.set(job.id!, abortController) // Get Redis client for checking cancel signals from the API process - const queueService = new QueueService() + const queueService = QueueService.getInstance() const cancelRedis = await queueService.getQueue(RunDownloadJob.queue).client let lastKnownProgress: Pick = { @@ -199,7 +199,7 @@ export class RunDownloadJob { } static async getByUrl(url: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(url) return await queue.getJob(jobId) @@ -229,7 +229,7 @@ export class RunDownloadJob { } static async dispatch(params: RunDownloadJobParams) { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(params.url) diff --git a/admin/app/jobs/run_extract_pmtiles_job.ts b/admin/app/jobs/run_extract_pmtiles_job.ts index 73c4eed..de7049f 100644 --- a/admin/app/jobs/run_extract_pmtiles_job.ts +++ b/admin/app/jobs/run_extract_pmtiles_job.ts @@ -49,7 +49,7 @@ export class RunExtractPmtilesJob { } static async signalCancel(jobId: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const client = await queue.client await client.set(this.cancelKey(jobId), '1', 'EX', 300) @@ -77,7 +77,7 @@ export class RunExtractPmtilesJob { `maxzoom=${maxzoom ?? 'source-max'} out=${outputFilepath}` ) - const queueService = new QueueService() + const queueService = QueueService.getInstance() const cancelRedis = await queueService.getQueue(RunExtractPmtilesJob.queue).client let userCancelled = false @@ -249,13 +249,13 @@ export class RunExtractPmtilesJob { } static async getById(jobId: string): Promise { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) return await queue.getJob(jobId) } static async dispatch(params: RunExtractPmtilesJobParams) { - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue(this.queue) const jobId = this.getJobId(params.sourceUrl, params.regionFilepath, params.maxzoom) diff --git a/admin/app/services/queue_service.ts b/admin/app/services/queue_service.ts index fa3a050..bad976b 100644 --- a/admin/app/services/queue_service.ts +++ b/admin/app/services/queue_service.ts @@ -1,9 +1,25 @@ import { Queue } from 'bullmq' import queueConfig from '#config/queue' +// Process-wide singleton. Each `Queue` opens two ioredis connections (one for +// commands, one blocking). Instantiating a fresh QueueService per dispatch / +// status lookup leaks both, and under sustained job churn (e.g. multi-batch ZIM +// ingestion enqueueing a continuation every few seconds) it saturates Redis's +// maxclients within hours. export class QueueService { private queues: Map = new Map() + private static _instance: QueueService | null = null + + private constructor() {} + + static getInstance(): QueueService { + if (!QueueService._instance) { + QueueService._instance = new QueueService() + } + return QueueService._instance + } + getQueue(name: string): Queue { if (!this.queues.has(name)) { const queue = new Queue(name, { @@ -18,5 +34,6 @@ export class QueueService { for (const queue of this.queues.values()) { await queue.close() } + this.queues.clear() } } diff --git a/admin/app/services/zim_service.ts b/admin/app/services/zim_service.ts index 1cc9e97..538d59f 100644 --- a/admin/app/services/zim_service.ts +++ b/admin/app/services/zim_service.ts @@ -314,7 +314,7 @@ export class ZimService { if (restart) { // Check if there are any remaining ZIM download jobs before restarting const { QueueService } = await import('./queue_service.js') - const queueService = new QueueService() + const queueService = QueueService.getInstance() const queue = queueService.getQueue('downloads') // Get all active and waiting jobs