diff --git a/rsmq-worker/rsmq-worker-tests.ts b/rsmq-worker/rsmq-worker-tests.ts new file mode 100644 index 000000000..4568302dd --- /dev/null +++ b/rsmq-worker/rsmq-worker-tests.ts @@ -0,0 +1,21 @@ + +import RSMQWorker = require('rsmq-worker'); + +var worker = new RSMQWorker("my-queue"); +worker.changeInterval(1); +worker.changeInterval([0, 1, 5, 10]); + +worker.on('message', (message: RedisSMQ.Message, next: Function) => { + console.log('received message: "' + message.message + "'"); + next(); +}); +worker.start(); + +worker.send('message1'); +worker.send('message2', 1, (e: Error, id:string) => { + worker.del(id); +}); +worker.send('message3', () => {}); + +worker.stop(); + diff --git a/rsmq-worker/rsmq-worker.d.ts b/rsmq-worker/rsmq-worker.d.ts new file mode 100644 index 000000000..8783d914f --- /dev/null +++ b/rsmq-worker/rsmq-worker.d.ts @@ -0,0 +1,52 @@ +// Type definitions for rsmq-worker 0.3.5 +// Project: http://smrchy.github.io/rsmq/rsmq-worker/ +// Definitions by: Qubo +// Definitions: https://github.com/borisyankov/DefinitelyTyped + +/// + +declare module "rsmq-worker" { + import redis = require('redis'); + import events = require('events'); + + interface CallbackT { + (e?:Error, res?:R): void; + } + + interface RSMQWorkerStatic { + new(queuename: string, options?: WorkerOptions): RSMQWorker; + } + + interface WorkerOptions { + interval?: number; + maxReceiveCount?: number; + invisibletime?: number; + defaultDelay?: number; + autostart?: boolean; + timeout: number; + customExceedCheck?: CustomExceedCheckCallback; + rsmq?: RedisSMQ.Client; + redis?: redis.RedisClient; + redisPrefix?: string; + host?: string; + port?: number; + options?: redis.ClientOpts; + } + + interface CustomExceedCheckCallback { + (message: RedisSMQ.Message): boolean; + } + + + interface RSMQWorker extends events.EventEmitter { + start(): RSMQWorker; + stop(): RSMQWorker; + send(message: string, delay?: number, cb?: CallbackT): RSMQWorker; + send(message: string, cb: CallbackT): RSMQWorker; + del(id: string, cb?: CallbackT): RSMQWorker; + changeInterval(interval: number|number[]): RSMQWorker; + } + + var worker: RSMQWorkerStatic; + export = worker; +} diff --git a/rsmq/rsmq-tests.ts b/rsmq/rsmq-tests.ts new file mode 100644 index 000000000..c7a89490a --- /dev/null +++ b/rsmq/rsmq-tests.ts @@ -0,0 +1,24 @@ +/// + +var RedisSMQ = require("rsmq"); +var rsmq = new RedisSMQ( {host: "127.0.0.1", port: 6379, ns: "rsmq"} ); +var rsmq2 = new RedisSMQ({client: rsmq.redis, ns: "rsmq2"}); + +rsmq.createQueue({qname: "my-queue"}, (e: Error, success: number) => { + if (e) { + console.error(e); + } + + console.info('"my-queue" has created'); +}); + +rsmq.sendMessage({qname: "my-queue", message: "first message"}, (e: Error, id: string) => { + rsmq.changeMessageVisibility({qname: "my-queue", id: id, vt: 100}, (e: Error, success: number) => {}); +}); +rsmq.getQueueAttributes({qname: 'my-queue'}, (e: Error, attr: RedisSMQ.QueueAttributes) =>{} ); +rsmq.setQueueAttributes({qname: "my-queue", vt: 20}, (e: Error, attr: RedisSMQ.QueueAttributes) => {}); +rsmq.receiveMessage({qname: "my-queue"}, (e: Error, message: RedisSMQ.Message) => { + rsmq.deleteMessage(message, (e: Error, success: number) => {}); +}); +rsmq.listQueues((e: Error, list: string[]) => {}); +rsmq.deleteQueue({qname: "my-queue"}, (e: Error, res: number) => {}); \ No newline at end of file diff --git a/rsmq/rsmq.d.ts b/rsmq/rsmq.d.ts new file mode 100644 index 000000000..34a5adbd4 --- /dev/null +++ b/rsmq/rsmq.d.ts @@ -0,0 +1,94 @@ +// Type definitions for rsmq 0.3.16 +// Project: http://smrchy.github.io/rsmq/ +// Definitions by: Qubo +// Definitions: https://github.com/borisyankov/DefinitelyTyped + +/// + +declare module RedisSMQ { + interface CallbackT { + (e?:Error, res?:R): void; + } + + interface Client { + createQueue(options:QueueOptions, cb:CallbackT): void; + changeMessageVisibility(options:VisibilityOptions, cb:CallbackT): void; + deleteMessage(options:MessageIdentifier, cb:CallbackT): void; + deleteQueue(options:QueueIdentifier, cb:CallbackT): void; + getQueueAttributes(options:QueueIdentifier, cb:CallbackT): void; + listQueues(cb:CallbackT): void; + receiveMessage(options:ReceiveOptions, cb:CallbackT): void; + sendMessage(options:NewMessage, cb:CallbackT): void; + setQueueAttributes(options:QueueOptions, cb:CallbackT): void; + quit(): void; + } + + interface QueueIdentifier { + qname: string; + } + + interface QueueOptions extends QueueIdentifier { + vt?: number; + delay?: number; + maxsize?: number; + } + + interface MessageIdentifier extends QueueIdentifier { + id: string; + } + + interface VisibilityOptions extends MessageIdentifier { + vt: number; + } + + export interface QueueAttributes extends QueueIdentifier { + vt: number; + delay: number; + maxsize: number; + totalrecv: number; + totalsent: number; + created: number; + modified: number; + msgs: number; + hiddenmsgs: number; + } + + interface ReceiveOptions extends QueueIdentifier { + vt?: number; + } + + export interface Message extends MessageIdentifier { + message: string; + sent: number; + fr: number; + rc: number; + } + + interface NewMessage extends QueueIdentifier { + message: string; + delay?: number; + } +} + +declare module 'rsmq' { + import redis = require('redis'); + + interface RedisSMQStatic { + new (options:ClientOptions): Client; + } + + interface Client extends RedisSMQ.Client{ + redis: redis.RedisClient; + } + + interface ClientOptions { + host?: string; + port?: number; + options?: redis.ClientOpts; + client?: redis.RedisClient; + ns?: string; + } + + var rsmq: RedisSMQStatic; + export = rsmq; +}