import { EventEmitter } from 'events'; import type { Job, JobProcessor, WorkerOptions, ClientOptions } from './types'; import { type WorkerHooks } from './hooks'; /** * Typed event map for Worker events. * Use with worker.on() for type-safe event handling. */ export interface WorkerEvents { /** Emitted when worker is ready to process jobs */ ready: []; /** Emitted when a job starts processing */ active: [job: Job & { data: T; }, workerId: number]; /** Emitted when a job completes successfully */ completed: [job: Job & { data: T; }, result: R, workerId: number]; /** Emitted when a job fails */ failed: [job: Job & { data: T; }, error: Error, workerId: number]; /** Emitted when worker starts stopping */ stopping: []; /** Emitted when worker has fully stopped */ stopped: []; /** Emitted when all jobs are drained */ drained: []; /** Emitted on worker errors (connection, timeout, etc.) */ error: [error: unknown]; /** Emitted during reconnection attempts */ reconnecting: [info: { attempt: number; delay: number; }]; /** Emitted after successful reconnection */ reconnected: []; } /** * Typed EventEmitter for Worker with proper event signatures. */ export interface TypedWorkerEmitter { on>(event: K, listener: (...args: WorkerEvents[K]) => void): this; once>(event: K, listener: (...args: WorkerEvents[K]) => void): this; emit>(event: K, ...args: WorkerEvents[K]): boolean; off>(event: K, listener: (...args: WorkerEvents[K]) => void): this; removeListener>(event: K, listener: (...args: WorkerEvents[K]) => void): this; removeAllListeners>(event?: K): this; } export interface BullMQWorkerOptions extends Omit, WorkerOptions { /** Auto-start worker (BullMQ-compatible, default: true) */ autorun?: boolean; /** Enable debug logging (default: false) */ debug?: boolean; /** Graceful shutdown timeout in ms (default: 30000). Use 0 for infinite wait. */ closeTimeout?: number; /** Worker hooks for observability */ workerHooks?: WorkerHooks; } type WorkerState = 'idle' | 'starting' | 'running' | 'stopping' | 'stopped'; /** * FlashQ Worker (BullMQ-compatible) * * @example * ```typescript * // BullMQ-style: auto-starts by default * const worker = new Worker('emails', async (job) => { * await sendEmail(job.data.to); * return { sent: true }; * }); * * // With options * const worker = new Worker('tasks', processor, { * concurrency: 10, * autorun: false, // disable auto-start * }); * await worker.start(); * * // Type-safe event handling * worker.on('completed', (job, result, workerId) => { * console.log(`Job ${job.id} completed by worker ${workerId}:`, result); * }); * * // Graceful shutdown * process.on('SIGTERM', () => worker.close()); * ``` */ export declare class Worker extends EventEmitter implements TypedWorkerEmitter { private clients; private clientOptions; private queues; private processor; private options; private state; private processing; private jobsProcessed; private workers; private startPromise; private stopPromise; private abortController; private workerHooks?; private logger; constructor(queues: string | string[], processor: JobProcessor, options?: BullMQWorkerOptions); /** * Start processing jobs */ start(): Promise; private doStart; /** * Close the worker (BullMQ-compatible alias for stop) * @param force If true, don't wait for jobs to complete (default: false) */ close(force?: boolean): Promise; /** * Stop processing jobs (graceful shutdown) * @param force If true, don't wait for jobs to complete (default: false) */ stop(force?: boolean): Promise; private doStop; /** * Wait for all currently processing jobs to complete * @param timeout Max wait time in ms (default: uses closeTimeout option) * @returns true if all jobs completed, false if timeout */ waitForJobs(timeout?: number): Promise; /** * Check if worker is running */ isRunning(): boolean; /** * Get current worker state */ getState(): WorkerState; /** * Get number of jobs currently being processed */ getProcessingCount(): number; /** * Get total number of jobs processed by this worker */ getJobsProcessed(): number; /** * Batch worker loop - pulls and processes jobs in batches for maximum throughput */ private batchWorkerLoop; /** * Process a batch of jobs - always completes even during shutdown */ private processJobBatch; private processJob; private sleep; /** * Update progress for the current job * (Use this within your processor function) * @throws Error if worker is not running */ updateProgress(jobId: number, progress: number, message?: string): Promise; } export default Worker; //# sourceMappingURL=worker.d.ts.map