/** * FlashQ Client - High-performance job queue client. * * This is the main facade that exposes all operations. * Methods are organized into logical modules under ./methods/ */ import { FlashQConnection } from './connection'; import { type ClientHooks } from '../hooks'; import type { Job, JobState, JobWithState, PushOptions, QueueInfo, QueueStats, Metrics, CronJob, CronOptions, JobLogEntry, FlowChild, FlowResult, FlowOptions, BatchPushResult } from './types'; export type { BatchPushResult } from './types'; /** FlashQ Client - High-performance job queue client with auto-connect. */ export declare class FlashQ extends FlashQConnection { /** Hooks for observability */ get hooks(): ClientHooks | undefined; /** Push a job to a queue */ push(queueName: string, data: T, opts?: PushOptions): Promise; /** Add a job to a queue (alias for push) */ add(queueName: string, data: T, opts?: PushOptions): Promise; /** Push multiple jobs in a single batch */ pushBatch(queueName: string, jobList: Array<{ data: T; } & PushOptions>): Promise; /** Add multiple jobs (alias for pushBatch) */ addBulk(queueName: string, jobList: Array<{ data: T; } & PushOptions>): Promise; /** Push multiple jobs with partial failure handling */ pushBatchSafe(queueName: string, jobList: Array<{ data: T; } & PushOptions>): Promise; /** Pull a job from a queue (blocking with timeout) */ pull(queueName: string, timeout?: number): Promise<(Job & { data: T; }) | null>; /** Pull multiple jobs from a queue */ pullBatch(queueName: string, count: number, timeout?: number): Promise>; /** Acknowledge a job as completed */ ack(jobId: number, result?: unknown): Promise; /** Acknowledge multiple jobs at once */ ackBatch(jobIds: number[]): Promise; /** Fail a job (will retry or move to DLQ) */ fail(jobId: number, error?: string): Promise; /** Get a job with its current state */ getJob(jobId: number): Promise; /** Get job state only */ getState(jobId: number): Promise; /** Get job result */ getResult(jobId: number): Promise; /** Wait for a job to complete (finished() promise pattern) */ finished(jobId: number, timeout?: number): Promise; /** Get a job by its custom ID */ getJobByCustomId(customId: string): Promise; /** Get multiple jobs by their IDs */ getJobsBatch(jobIds: number[]): Promise; /** Cancel a pending job */ cancel(jobId: number): Promise; /** Update job progress */ progress(jobId: number, value: number, message?: string): Promise; /** Get job progress */ getProgress(jobId: number): Promise<{ progress: number; message?: string; }>; /** Add a log entry to a job */ log(jobId: number, message: string, level?: 'info' | 'warn' | 'error'): Promise; /** Get log entries for a job */ getLogs(jobId: number): Promise; /** Send a heartbeat for a long-running job */ heartbeat(jobId: number): Promise; /** Send partial result for streaming jobs (LLM tokens, chunks, etc.) */ partial(jobId: number, data: unknown, index?: number): Promise; /** Pause a queue */ pause(queueName: string): Promise; /** Resume a paused queue */ resume(queueName: string): Promise; /** Check if a queue is paused */ isPaused(queueName: string): Promise; /** Set rate limit for a queue (jobs per second) */ setRateLimit(queueName: string, limit: number): Promise; /** Clear rate limit for a queue */ clearRateLimit(queueName: string): Promise; /** Set concurrency limit for a queue */ setConcurrency(queueName: string, limit: number): Promise; /** Clear concurrency limit for a queue */ clearConcurrency(queueName: string): Promise; /** List all queues */ listQueues(): Promise; /** Get jobs from the dead letter queue */ getDlq(queueName: string, count?: number): Promise; /** Retry jobs from the dead letter queue */ retryDlq(queueName: string, jobId?: number): Promise; /** Purge all jobs from the dead letter queue */ purgeDlq(queueName: string): Promise; /** Add a cron job for scheduled recurring tasks */ addCron(name: string, options: CronOptions): Promise; /** Delete a cron job */ deleteCron(name: string): Promise; /** List all cron jobs */ listCrons(): Promise; /** Get queue statistics */ stats(): Promise; /** Get detailed metrics */ metrics(): Promise; /** Push a flow (parent job with children) */ pushFlow(queueName: string, parentData: T, children: FlowChild[], options?: FlowOptions): Promise; /** Get children job IDs for a parent job */ getChildren(jobId: number): Promise; /** Get jobs filtered by queue and/or state with pagination */ getJobs(options?: { queue?: string; state?: JobState; limit?: number; offset?: number; }): Promise<{ jobs: JobWithState[]; total: number; }>; /** Get job counts by state for a queue */ getJobCounts(queueName: string): Promise<{ waiting: number; active: number; delayed: number; completed: number; failed: number; }>; /** Get total count of jobs in a queue (waiting + delayed) */ count(queueName: string): Promise; /** Clean jobs older than grace period by state */ clean(queueName: string, grace: number, state: 'waiting' | 'delayed' | 'completed' | 'failed', limit?: number): Promise; /** Drain all waiting jobs from a queue */ drain(queueName: string): Promise; /** Remove ALL data for a queue */ obliterate(queueName: string): Promise; /** Change job priority */ changePriority(jobId: number, priority: number): Promise; /** Move job from processing back to delayed */ moveToDelayed(jobId: number, delay: number): Promise; /** Promote delayed job to waiting immediately */ promote(jobId: number): Promise; /** Update job data */ update(jobId: number, data: T): Promise; /** Discard job - move directly to DLQ */ discard(jobId: number): Promise; /** Subscribe to real-time events via SSE */ subscribe(queueName?: string): import('../events').EventSubscriber; /** Subscribe to real-time events via WebSocket */ subscribeWs(queueName?: string): import('../events').EventSubscriber; } export default FlashQ; //# sourceMappingURL=index.d.ts.map