import { Queue } from '../core/queue.ts'; import type { JobStatus, JobMeta, QueueMessage, BaseJobOptions, WithPriority, WithDelay } from '../interfaces/job.ts'; import type { DatabaseAdapter } from '../interfaces/database.ts'; import type { QueueOptions } from '../interfaces/plugin.ts'; // Driver-specific job request interface export interface DbJobRequest extends BaseJobOptions, WithPriority, WithDelay { /** Job payload */ payload: TPayload; // DB adapters may or may not support delay/priority - we allow them for flexibility // The specific DatabaseAdapter implementation determines actual support } export class DbQueue> extends Queue> { constructor( private db: DatabaseAdapter, options: QueueOptions ) { super(options); } get adapter(): DatabaseAdapter { return this.db; } protected async pushMessage(payload: unknown, meta: JobMeta): Promise { return await this.db.insertJob(payload, meta); } protected async reserve(timeout: number): Promise { const record = await this.db.reserveJob(timeout); if (!record) { return null; } return { id: record.id, payload: record.payload, meta: record.meta, }; } protected async completeJob(message: QueueMessage): Promise { await this.db.completeJob(message.id); } protected async failJob(message: QueueMessage, error: unknown): Promise { const errorMessage = error instanceof Error ? error.message : String(error); await this.db.failJob(message.id, errorMessage); } async status(id: string): Promise { const status = await this.db.getJobStatus(id); return status || 'done'; } }