diff --git a/types/bull/bull-tests.tsx b/types/bull/bull-tests.tsx index fcf339b3c3..2ff4579931 100644 --- a/types/bull/bull-tests.tsx +++ b/types/bull/bull-tests.tsx @@ -4,68 +4,69 @@ import * as Queue from "bull" -var videoQueue = Queue( 'video transcoding', 6379, '127.0.0.1' ); -var audioQueue = Queue( 'audio transcoding', 6379, '127.0.0.1' ); -var imageQueue = Queue( 'image transcoding', 6379, '127.0.0.1' ); +var videoQueue = new Queue('video transcoding', 'redis://127.0.0.1:6379'); +var audioQueue = new Queue('audio transcoding', {redis: {port: 6379, host: '127.0.0.1'}}); // Specify Redis connection using object +var imageQueue = new Queue('image transcoding'); +var pdfQueue = new Queue('pdf transcoding'); -videoQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { +videoQueue.process(function(job, done){ // job.data contains the custom data passed when the job was created // job.jobId contains id of this job. // transcode video asynchronously and report progress - job.progress( 42 ); + job.progress(42); // call done when finished done(); // or give a error if error - done( Error( 'error transcoding' ) ); + done(new Error('error transcoding')); // or pass it a result - done( null, { framerate: 29.5 /* etc... */ } ); + done(null, { framerate: 29.5 /* etc... */ }); // If the job throws an unhandled exception it is also handled correctly - throw (Error( 'some unexpected error' )); -} ); + throw new Error('some unexpected error'); +}); -audioQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { +audioQueue.process(function(job, done){ // transcode audio asynchronously and report progress - job.progress( 42 ); + job.progress(42); // call done when finished done(); // or give a error if error - done( Error( 'error transcoding' ) ); + done(new Error('error transcoding')); // or pass it a result - done( null, { samplerate: 48000 /* etc... */ } ); + done(null, { samplerate: 48000 /* etc... */ }); // If the job throws an unhandled exception it is also handled correctly - throw (Error( 'some unexpected error' )); -} ); + throw new Error('some unexpected error'); +}); -imageQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { +imageQueue.process(function(job, done){ // transcode image asynchronously and report progress - job.progress( 42 ); + job.progress(42); // call done when finished done(); // or give a error if error - done( Error( 'error transcoding' ) ); + done(new Error('error transcoding')); // or pass it a result - done( null, { width: 1280, height: 720 /* etc... */ } ); + done(null, { width: 1280, height: 720 /* etc... */ }); // If the job throws an unhandled exception it is also handled correctly - throw (Error( 'some unexpected error' )); -} ); + throw new Error('some unexpected error'); +}); -videoQueue.add( { video: 'http://example.com/video1.mov' } ); -audioQueue.add( { audio: 'http://example.com/audio1.mp3' } ); -imageQueue.add( { image: 'http://example.com/image1.tiff' } ); +videoQueue.add({video: 'http://example.com/video1.mov'}); +audioQueue.add({audio: 'http://example.com/audio1.mp3'}); +imageQueue.add({image: 'http://example.com/image1.tiff'}); ////////////////////////////////////////////////////////////////////////////////// @@ -74,42 +75,43 @@ imageQueue.add( { image: 'http://example.com/image1.tiff' } ); // ////////////////////////////////////////////////////////////////////////////////// -const fetchVideo = ( url: string ): Promise => { return null } -const transcodeVideo = ( data: any ): Promise => { return null } -interface VideoJob extends Queue.Job { - data: {url: string} -} +pdfQueue.process(function(job){ + // Processors can also return promises instead of using the done callback + return Promise.resolve(); +}); -videoQueue.process( ( job: VideoJob ) => { // don't forget to remove the done callback! - // Simply return a promise - fetchVideo( job.data.url ).then( transcodeVideo ); - - // Handles promise rejection - Promise.reject( new Error( 'error transcoding' ) ); - - // Passes the value the promise is resolved with to the "completed" event - Promise.resolve( { framerate: 29.5 /* etc... */ } ); - - // same as - Promise.reject( new Error( 'some unexpected error' ) ); - - // If the job throws an unhandled exception it is also handled correctly - throw new Error( 'some unexpected error' ); -} ); - - -var addVideo1Job = videoQueue.add( { video: 'http://example.com/video1.mov' } ); - -addVideo1Job.then((video1Job) => { +videoQueue.add({ video: 'http://example.com/video1.mov' }, { jobId: 1 }) +.then((video1Job) => { // When job has successfully be placed in the queue the job is returned // then wait for completion return video1Job.finished(); }) .then(() => { - // video1Job completed successfully + // completed successfully }) .catch((err) => { // error }); + + +////////////////////////////////////////////////////////////////////////////////// +// +// Typed Event Handlers +// +////////////////////////////////////////////////////////////////////////////////// + + +pdfQueue +.on('ready', () => undefined) +.on('error', (err: Error) => undefined) +.on('active', (job: Queue.Job, jobPromise: Queue.JobPromise) => jobPromise.cancel()) +.on('active', (job: Queue.Job) => undefined) +.on('stalled', (job: Queue.Job) => undefined) +.on('progress', (job: Queue.Job) => undefined) +.on('completed', (job: Queue.Job) => undefined) +.on('failed', (job: Queue.Job) => undefined) +.on('paused', () => undefined) +.on('resumed', () => undefined) +.on('cleaned', (jobs: Queue.Job[], status: Queue.JobStatus) => undefined) diff --git a/types/bull/index.d.ts b/types/bull/index.d.ts index 3c89655101..3c49109903 100644 --- a/types/bull/index.d.ts +++ b/types/bull/index.d.ts @@ -1,6 +1,8 @@ -// Type definitions for bull 2.1.2 +// Type definitions for bull v.3.0.0-rc.2 // Project: https://github.com/OptimalBits/bull -// Definitions by: Bruno Grieder , Cameron Crothers +// Definitions by: Bruno Grieder +// Cameron Crothers +// Marshall Cottrell // Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped /// @@ -14,22 +16,72 @@ declare module "bull" { * It creates a new Queue that is persisted in Redis. * Everytime the same queue is instantiated it tries to process all the old jobs that may exist from a previous unfinished session. */ - function Bull(queueName: string, redisPort: number, redisHost: string, redisOpt?: Redis.ClientOpts): Bull.Queue; + var Bull: { + (queueName: string, opts?: Bull.QueueOptions): Bull.Queue; + (queueName: string, url?: string): Bull.Queue; + new (queueName: string, opts?: Bull.QueueOptions): Bull.Queue; + new (queueName: string, url?: string): Bull.Queue; + } namespace Bull { - export interface DoneCallback { - (error?: Error, value?: any): void + export interface QueueOptions { + + /** + * Options passed directly to the `ioredis` constructor + */ + redis?: Redis.ClientOpts; + + /** + * Prefix to use for all redis keys + */ + prefix?: string; + + settings?: AdvancedSettings; } + export interface AdvancedSettings { + + /** + * Key expiration time for job locks + */ + lockDuration?: number; + + /** + * How often check for stalled jobs (use 0 for never checking) + */ + stalledInterval?: number; + + /** + * Max amount of times a stalled job will be re-processed + */ + maxStalledCount?: number; + + /** + * Poll interval for delayed jobs and added jobs + */ + guardInterval?: number; + + /** + * Delay before processing next job in case of internal error + */ + retryProcessDelay?: number; + } + + export interface DoneCallback { + (error?: Error | null, value?: any): void + } + + export type JobId = number | string; + export interface Job { - jobId: string + id: JobId; /** - * The custom data passed when the job was created - */ - data: Object; + * The custom data passed when the job was created + */ + data: any; /** * Report progress on a job @@ -37,13 +89,13 @@ declare module "bull" { progress(value: any): Promise; /** - * Removes a Job from the queue from all the lists where it may be included. + * Removes a job from the queue and from any lists it may be included in. * @returns {Promise} A promise that resolves when the job is removed. */ remove(): Promise; /** - * Rerun a Job that has failed. + * Re-run a job that has failed. * @returns {Promise} A promise that resolves when the job is scheduled for retry. */ retry(): Promise; @@ -56,12 +108,14 @@ declare module "bull" { finished(): Promise; } - export interface Backoff { + export type JobStatus = 'completed' | 'waiting' | 'active' | 'delayed' | 'failed'; + + export interface BackoffOptions { /** * Backoff type, which can be either `fixed` or `exponential` */ - type: string + type: 'fixed' | 'exponential'; /** * Backoff delay, in milliseconds @@ -69,22 +123,52 @@ declare module "bull" { delay: number; } - export interface AddOptions { + export interface RepeatOptions { + + /** + * Cron pattern specifying when the job should execute + */ + cron: string; + + /** + * Timezone + */ + tz?: string; + + /** + * End date when the repeat job should stop repeating + */ + endDate?: Date | string | number; + } + + export interface JobOptions { + + /** + * Optional priority value. ranges from 1 (highest priority) to MAX_INT (lowest priority). + * Note that using priorities has a slight impact on performance, so do not use it if not required + */ + priority?: number; + /** * An amount of miliseconds to wait until this job can be processed. - * Note that for accurate delays, both server and clients should have their clocks synchronized + * Note that for accurate delays, both server and clients should have their clocks synchronized. [optional] */ delay?: number; /** - * A number of attempts to retry if the job fails [optional] + * The total number of attempts to try the job until it completes */ attempts?: number; + /** + * Repeat job according to a cron specification + */ + repeat?: RepeatOptions; + /** * Backoff setting for automatic retries if the job fails */ - backoff?: number | Backoff + backoff?: number | BackoffOptions; /** * A boolean which, if true, adds the job to the right @@ -96,6 +180,35 @@ declare module "bull" { * The number of milliseconds after which the job should be fail with a timeout error */ timeout?: number; + + /** + * Override the job ID - by default, the job ID is a unique + * integer, but you can use this setting to override it. + * If you use this option, it is up to you to ensure the + * jobId is unique. If you attempt to add a job with an id that + * already exists, it will not be added. + */ + jobId?: JobId; + + /** + * A boolean which, if true, removes the job when it successfully completes. + * Default behavior is to keep the job in the completed set. + */ + removeOnComplete?: boolean; + + /** + * A boolean which, if true, removes the job when it fails after all attempts + * Default behavior is to keep the job in the completed set. + */ + removeOnFail?: boolean; + } + + export interface JobCounts { + wait: number; + active: number; + completed: number; + failed: number; + delayed: number; } export interface Queue { @@ -112,7 +225,7 @@ declare module "bull" { * results, as a second argument to the "completed" event. * * concurrency: Bull will then call you handler in parallel respecting this max number. - */ + */ process(concurrency: number, callback: (job: Job, done: DoneCallback) => void): void; /** @@ -125,7 +238,7 @@ declare module "bull" { * or with a result as second argument as second argument (e.g.: done(null, result);) when the job is successful. * Errors will be passed as a second argument to the "failed" event; * results, as a second argument to the "completed" event. - */ + */ process(callback: (job: Job, done: DoneCallback) => void): void; /** @@ -139,7 +252,7 @@ declare module "bull" { * If it is resolved, its value will be the "completed" event's second argument. * * concurrency: Bull will then call you handler in parallel respecting this max number. - */ + */ process(concurrency: number, callback: (job: Job) => void): Promise; /** @@ -151,17 +264,15 @@ declare module "bull" { * A promise must be returned to signal job completion. * If the promise is rejected, the error will be passed as a second argument to the "failed" event. * If it is resolved, its value will be the "completed" event's second argument. - */ + */ process(callback: (job: Job) => void): Promise; - // process(callback: (job: Job, done?: DoneCallback) => void): Promise; - /** * Creates a new job and adds it to the queue. * If the queue is empty the job will be executed directly, * otherwise it will be placed in the queue and executed as soon as possible. */ - add(data: Object, opts?: AddOptions): Promise; + add(data: Object, opts?: JobOptions): Promise; /** * Returns a promise that resolves when the queue is paused. @@ -204,113 +315,118 @@ declare module "bull" { * Returns a promise that will return the job instance associated with the jobId parameter. * If the specified job cannot be located, the promise callback parameter will be set to null. */ - getJob(jobId: string): Promise; + getJob(jobId: JobId): Promise; + + /** + * Returns a promise that resolves with the job counts for the given queue + */ + getJobCounts(): Promise; /** * Tells the queue remove all jobs created outside of a grace period in milliseconds. * You can clean the jobs with the following states: completed, waiting, active, delayed, and failed. */ - clean(gracePeriod: number, jobsState?: string): Promise; + clean(grace: number, status?: JobStatus, limit?: number): Promise; /** * Listens to queue events - * 'ready', 'error', 'activ', 'progress', 'completed', 'failed', 'paused', 'resumed', 'cleaned' */ - on(eventName: string, callback: EventCallback): void; + on(event: string, callback: (...args: any[]) => void): this; + + /** + * Redis is connected and the queue is ready to accept jobs + */ + on(event: 'ready', callback: EventCallback): this; + + /** + * An error occured + */ + on(event: 'error', callback: ErrorEventCallback): this; + + /** + * A job has started. You can use `jobPromise.cancel()` to abort it + */ + on(event: 'active', callback: ActiveEventCallback): this; + + /** + * A job has been marked as stalled. + * This is useful for debugging job workers that crash or pause the event loop. + */ + on(event: 'stalled', callback: StalledEventCallback): this; + + /** + * A job's progress was updated + */ + on(event: 'progress', callback: ProgressEventCallback): this; + + /** + * A job successfully completed with a `result` + */ + on(event: 'completed', callback: CompletedEventCallback): this; + + /** + * A job failed with `err` as the reason + */ + on(event: 'failed', callback: FailedEventCallback): this; + + /** + * The queue has been paused + */ + on(event: 'paused', callback: EventCallback): this; + + /** + * The queue has been resumed + */ + on(event: 'resumed', callback: EventCallback): this; + + /** + * Old jobs have been cleaned from the queue. + * `jobs` is an array of jobs that were removed, and `type` is the type of those jobs. + * + * @see Queue#clean() for details + */ + on(event: 'cleaned', callback: CleanedEventCallback): this; } - interface EventCallback { - (...args: any[]): void - } - - interface ReadyEventCallback extends EventCallback { + export interface EventCallback { (): void; } - interface ErrorEventCallback extends EventCallback { + export interface ErrorEventCallback { (error: Error): void; } - interface JobPromise { + export interface JobPromise { /** * Abort this job */ cancel(): void } - interface ActiveEventCallback extends EventCallback { - (job: Job, jobPromise: JobPromise): void; + export interface ActiveEventCallback { + (job: Job, jobPromise?: JobPromise): void; } - interface ProgressEventCallback extends EventCallback { + export interface StalledEventCallback { + (job: Job): void; + } + + export interface ProgressEventCallback { (job: Job, progress: any): void; } - interface CompletedEventCallback extends EventCallback { - (job: Job, result: Object): void; + export interface CompletedEventCallback { + (job: Job, result: any): void; } - interface FailedEventCallback extends EventCallback { + export interface FailedEventCallback { (job: Job, error: Error): void; } - interface PausedEventCallback extends EventCallback { - (): void; - } - - interface ResumedEventCallback extends EventCallback { - (job?: Job): void; - } - - /** - * @see clean() for details - */ - interface CleanedEventCallback extends EventCallback { - (jobs: Job[], type: string): void; + export interface CleanedEventCallback { + (jobs: Job[], status: JobStatus): void; } } export = Bull; -} - -declare module "bull/lib/priority-queue" { - - import * as Bull from "bull"; - import * as Redis from "redis"; - - /** - * This is the Queue constructor of priority queue. - * - * It works same a normal queue, with same function and parameters. - * The only difference is that the Queue#add() allow an options opts.priority - * that could take ["low", "normal", "medium", "hight", "critical"]. If no options provider, "normal" will be taken. - * - * The priority queue will process more often highter priority jobs than lower. - */ - function PQueue(queueName: string, redisPort: number, redisHost: string, redisOpt?: Redis.ClientOpts): PQueue.PriorityQueue; - - namespace PQueue { - - export interface AddOptions extends Bull.AddOptions { - - /** - * "low", "normal", "medium", "high", "critical" - */ - priority?: string; - } - - - export interface PriorityQueue extends Bull.Queue { - - /** - * Creates a new job and adds it to the queue. - * If the queue is empty the job will be executed directly, - * otherwise it will be placed in the queue and executed as soon as possible. - */ - add(data: Object, opts?: PQueue.AddOptions): Promise; - - } - } - - export = PQueue; -} +} \ No newline at end of file diff --git a/types/bull/v2/bull-tests.tsx b/types/bull/v2/bull-tests.tsx new file mode 100644 index 0000000000..fcf339b3c3 --- /dev/null +++ b/types/bull/v2/bull-tests.tsx @@ -0,0 +1,115 @@ +/** + * Created by Bruno Grieder + */ + +import * as Queue from "bull" + +var videoQueue = Queue( 'video transcoding', 6379, '127.0.0.1' ); +var audioQueue = Queue( 'audio transcoding', 6379, '127.0.0.1' ); +var imageQueue = Queue( 'image transcoding', 6379, '127.0.0.1' ); + +videoQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { + + // job.data contains the custom data passed when the job was created + // job.jobId contains id of this job. + + // transcode video asynchronously and report progress + job.progress( 42 ); + + // call done when finished + done(); + + // or give a error if error + done( Error( 'error transcoding' ) ); + + // or pass it a result + done( null, { framerate: 29.5 /* etc... */ } ); + + // If the job throws an unhandled exception it is also handled correctly + throw (Error( 'some unexpected error' )); +} ); + +audioQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { + // transcode audio asynchronously and report progress + job.progress( 42 ); + + // call done when finished + done(); + + // or give a error if error + done( Error( 'error transcoding' ) ); + + // or pass it a result + done( null, { samplerate: 48000 /* etc... */ } ); + + // If the job throws an unhandled exception it is also handled correctly + throw (Error( 'some unexpected error' )); +} ); + +imageQueue.process( ( job: Queue.Job, done: Queue.DoneCallback ) => { + // transcode image asynchronously and report progress + job.progress( 42 ); + + // call done when finished + done(); + + // or give a error if error + done( Error( 'error transcoding' ) ); + + // or pass it a result + done( null, { width: 1280, height: 720 /* etc... */ } ); + + // If the job throws an unhandled exception it is also handled correctly + throw (Error( 'some unexpected error' )); +} ); + +videoQueue.add( { video: 'http://example.com/video1.mov' } ); +audioQueue.add( { audio: 'http://example.com/audio1.mp3' } ); +imageQueue.add( { image: 'http://example.com/image1.tiff' } ); + + +////////////////////////////////////////////////////////////////////////////////// +// +// Using Promises +// +////////////////////////////////////////////////////////////////////////////////// + +const fetchVideo = ( url: string ): Promise => { return null } +const transcodeVideo = ( data: any ): Promise => { return null } + +interface VideoJob extends Queue.Job { + data: {url: string} +} + + +videoQueue.process( ( job: VideoJob ) => { // don't forget to remove the done callback! + // Simply return a promise + fetchVideo( job.data.url ).then( transcodeVideo ); + + // Handles promise rejection + Promise.reject( new Error( 'error transcoding' ) ); + + // Passes the value the promise is resolved with to the "completed" event + Promise.resolve( { framerate: 29.5 /* etc... */ } ); + + // same as + Promise.reject( new Error( 'some unexpected error' ) ); + + // If the job throws an unhandled exception it is also handled correctly + throw new Error( 'some unexpected error' ); +} ); + + +var addVideo1Job = videoQueue.add( { video: 'http://example.com/video1.mov' } ); + +addVideo1Job.then((video1Job) => { + // When job has successfully be placed in the queue the job is returned + // then wait for completion + return video1Job.finished(); +}) +.then(() => { + // video1Job completed successfully +}) +.catch((err) => { + // error +}); diff --git a/types/bull/v2/index.d.ts b/types/bull/v2/index.d.ts new file mode 100644 index 0000000000..3c89655101 --- /dev/null +++ b/types/bull/v2/index.d.ts @@ -0,0 +1,316 @@ +// Type definitions for bull 2.1.2 +// Project: https://github.com/OptimalBits/bull +// Definitions by: Bruno Grieder , Cameron Crothers +// Definitions: https://github.com/DefinitelyTyped/DefinitelyTyped + +/// + +declare module "bull" { + + import * as Redis from "redis"; + + /** + * This is the Queue constructor. + * It creates a new Queue that is persisted in Redis. + * Everytime the same queue is instantiated it tries to process all the old jobs that may exist from a previous unfinished session. + */ + function Bull(queueName: string, redisPort: number, redisHost: string, redisOpt?: Redis.ClientOpts): Bull.Queue; + + namespace Bull { + + export interface DoneCallback { + (error?: Error, value?: any): void + } + + export interface Job { + + jobId: string + + /** + * The custom data passed when the job was created + */ + data: Object; + + /** + * Report progress on a job + */ + progress(value: any): Promise; + + /** + * Removes a Job from the queue from all the lists where it may be included. + * @returns {Promise} A promise that resolves when the job is removed. + */ + remove(): Promise; + + /** + * Rerun a Job that has failed. + * @returns {Promise} A promise that resolves when the job is scheduled for retry. + */ + retry(): Promise; + + /** + * Returns a promise the resolves when the job has been finished. + * TODO: Add a watchdog to check if the job has finished periodically. + * since pubsub does not give any guarantees. + */ + finished(): Promise; + } + + export interface Backoff { + + /** + * Backoff type, which can be either `fixed` or `exponential` + */ + type: string + + /** + * Backoff delay, in milliseconds + */ + delay: number; + } + + export interface AddOptions { + /** + * An amount of miliseconds to wait until this job can be processed. + * Note that for accurate delays, both server and clients should have their clocks synchronized + */ + delay?: number; + + /** + * A number of attempts to retry if the job fails [optional] + */ + attempts?: number; + + /** + * Backoff setting for automatic retries if the job fails + */ + backoff?: number | Backoff + + /** + * A boolean which, if true, adds the job to the right + * of the queue instead of the left (default false) + */ + lifo?: boolean; + + /** + * The number of milliseconds after which the job should be fail with a timeout error + */ + timeout?: number; + } + + export interface Queue { + + /** + * Defines a processing function for the jobs placed into a given Queue. + * + * The callback is called everytime a job is placed in the queue. + * It is passed an instance of the job as first argument. + * + * The done callback can be called with an Error instance, to signal that the job did not complete successfully, + * or with a result as second argument as second argument (e.g.: done(null, result);) when the job is successful. + * Errors will be passed as a second argument to the "failed" event; + * results, as a second argument to the "completed" event. + * + * concurrency: Bull will then call you handler in parallel respecting this max number. + */ + process(concurrency: number, callback: (job: Job, done: DoneCallback) => void): void; + + /** + * Defines a processing function for the jobs placed into a given Queue. + * + * The callback is called everytime a job is placed in the queue. + * It is passed an instance of the job as first argument. + * + * The done callback can be called with an Error instance, to signal that the job did not complete successfully, + * or with a result as second argument as second argument (e.g.: done(null, result);) when the job is successful. + * Errors will be passed as a second argument to the "failed" event; + * results, as a second argument to the "completed" event. + */ + process(callback: (job: Job, done: DoneCallback) => void): void; + + /** + * Defines a processing function for the jobs placed into a given Queue. + * + * The callback is called everytime a job is placed in the queue. + * It is passed an instance of the job as first argument. + * + * A promise must be returned to signal job completion. + * If the promise is rejected, the error will be passed as a second argument to the "failed" event. + * If it is resolved, its value will be the "completed" event's second argument. + * + * concurrency: Bull will then call you handler in parallel respecting this max number. + */ + process(concurrency: number, callback: (job: Job) => void): Promise; + + /** + * Defines a processing function for the jobs placed into a given Queue. + * + * The callback is called everytime a job is placed in the queue. + * It is passed an instance of the job as first argument. + * + * A promise must be returned to signal job completion. + * If the promise is rejected, the error will be passed as a second argument to the "failed" event. + * If it is resolved, its value will be the "completed" event's second argument. + */ + process(callback: (job: Job) => void): Promise; + + // process(callback: (job: Job, done?: DoneCallback) => void): Promise; + + /** + * Creates a new job and adds it to the queue. + * If the queue is empty the job will be executed directly, + * otherwise it will be placed in the queue and executed as soon as possible. + */ + add(data: Object, opts?: AddOptions): Promise; + + /** + * Returns a promise that resolves when the queue is paused. + * The pause is global, meaning that all workers in all queue instances for a given queue will be paused. + * A paused queue will not process new jobs until resumed, + * but current jobs being processed will continue until they are finalized. + * + * Pausing a queue that is already paused does nothing. + */ + pause(): Promise; + + /** + * Returns a promise that resolves when the queue is resumed after being paused. + * The resume is global, meaning that all workers in all queue instances for a given queue will be resumed. + * + * Resuming a queue that is not paused does nothing. + */ + resume(): Promise; + + /** + * Returns a promise that returns the number of jobs in the queue, waiting or paused. + * Since there may be other processes adding or processing jobs, this value may be true only for a very small amount of time. + */ + count(): Promise; + + /** + * Empties a queue deleting all the input lists and associated jobs. + */ + empty(): Promise; + + /** + * Closes the underlying redis client. Use this to perform a graceful shutdown. + * + * `close` can be called from anywhere, with one caveat: + * if called from within a job handler the queue won't close until after the job has been processed + */ + close(): Promise; + + /** + * Returns a promise that will return the job instance associated with the jobId parameter. + * If the specified job cannot be located, the promise callback parameter will be set to null. + */ + getJob(jobId: string): Promise; + + /** + * Tells the queue remove all jobs created outside of a grace period in milliseconds. + * You can clean the jobs with the following states: completed, waiting, active, delayed, and failed. + */ + clean(gracePeriod: number, jobsState?: string): Promise; + + /** + * Listens to queue events + * 'ready', 'error', 'activ', 'progress', 'completed', 'failed', 'paused', 'resumed', 'cleaned' + */ + on(eventName: string, callback: EventCallback): void; + } + + interface EventCallback { + (...args: any[]): void + } + + interface ReadyEventCallback extends EventCallback { + (): void; + } + + interface ErrorEventCallback extends EventCallback { + (error: Error): void; + } + + interface JobPromise { + /** + * Abort this job + */ + cancel(): void + } + + interface ActiveEventCallback extends EventCallback { + (job: Job, jobPromise: JobPromise): void; + } + + interface ProgressEventCallback extends EventCallback { + (job: Job, progress: any): void; + } + + interface CompletedEventCallback extends EventCallback { + (job: Job, result: Object): void; + } + + interface FailedEventCallback extends EventCallback { + (job: Job, error: Error): void; + } + + interface PausedEventCallback extends EventCallback { + (): void; + } + + interface ResumedEventCallback extends EventCallback { + (job?: Job): void; + } + + /** + * @see clean() for details + */ + interface CleanedEventCallback extends EventCallback { + (jobs: Job[], type: string): void; + } + } + + export = Bull; +} + +declare module "bull/lib/priority-queue" { + + import * as Bull from "bull"; + import * as Redis from "redis"; + + /** + * This is the Queue constructor of priority queue. + * + * It works same a normal queue, with same function and parameters. + * The only difference is that the Queue#add() allow an options opts.priority + * that could take ["low", "normal", "medium", "hight", "critical"]. If no options provider, "normal" will be taken. + * + * The priority queue will process more often highter priority jobs than lower. + */ + function PQueue(queueName: string, redisPort: number, redisHost: string, redisOpt?: Redis.ClientOpts): PQueue.PriorityQueue; + + namespace PQueue { + + export interface AddOptions extends Bull.AddOptions { + + /** + * "low", "normal", "medium", "high", "critical" + */ + priority?: string; + } + + + export interface PriorityQueue extends Bull.Queue { + + /** + * Creates a new job and adds it to the queue. + * If the queue is empty the job will be executed directly, + * otherwise it will be placed in the queue and executed as soon as possible. + */ + add(data: Object, opts?: PQueue.AddOptions): Promise; + + } + } + + export = PQueue; +} diff --git a/types/bull/v2/tsconfig.json b/types/bull/v2/tsconfig.json new file mode 100644 index 0000000000..2948048b1e --- /dev/null +++ b/types/bull/v2/tsconfig.json @@ -0,0 +1,26 @@ +{ + "compilerOptions": { + "module": "commonjs", + "lib": [ + "es6" + ], + "noImplicitAny": true, + "noImplicitThis": true, + "strictNullChecks": false, + "baseUrl": "../", + "typeRoots": [ + "../" + ], + "types": [], + "paths": { + "bull": [ "bull/v2" ], + "bull/*": [ "bull/v2/*" ] + }, + "noEmit": true, + "forceConsistentCasingInFileNames": true + }, + "files": [ + "index.d.ts", + "bull-tests.tsx" + ] +} \ No newline at end of file