Merge pull request #5358 from MugeSo/rsmq

Add RSMQ
This commit is contained in:
Masahiro Wakame
2015-08-17 22:28:16 +09:00
4 changed files with 191 additions and 0 deletions
+21
View File
@@ -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();
+52
View File
@@ -0,0 +1,52 @@
// Type definitions for rsmq-worker 0.3.5
// Project: http://smrchy.github.io/rsmq/rsmq-worker/
// Definitions by: Qubo <https://github.com/MugeSo>
// Definitions: https://github.com/borisyankov/DefinitelyTyped
/// <reference path='../rsmq/rsmq.d.ts'/>
declare module "rsmq-worker" {
import redis = require('redis');
import events = require('events');
interface CallbackT<R> {
(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<string>): RSMQWorker;
send(message: string, cb: CallbackT<string>): RSMQWorker;
del(id: string, cb?: CallbackT<void>): RSMQWorker;
changeInterval(interval: number|number[]): RSMQWorker;
}
var worker: RSMQWorkerStatic;
export = worker;
}
+24
View File
@@ -0,0 +1,24 @@
/// <reference path="rsmq.d.ts" />
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) => {});
+94
View File
@@ -0,0 +1,94 @@
// Type definitions for rsmq 0.3.16
// Project: http://smrchy.github.io/rsmq/
// Definitions by: Qubo <https://github.com/MugeSo>
// Definitions: https://github.com/borisyankov/DefinitelyTyped
/// <reference path='../redis/redis.d.ts'/>
declare module RedisSMQ {
interface CallbackT<R> {
(e?:Error, res?:R): void;
}
interface Client {
createQueue(options:QueueOptions, cb:CallbackT<number>): void;
changeMessageVisibility(options:VisibilityOptions, cb:CallbackT<number>): void;
deleteMessage(options:MessageIdentifier, cb:CallbackT<number>): void;
deleteQueue(options:QueueIdentifier, cb:CallbackT<number>): void;
getQueueAttributes(options:QueueIdentifier, cb:CallbackT<QueueAttributes>): void;
listQueues(cb:CallbackT<string[]>): void;
receiveMessage(options:ReceiveOptions, cb:CallbackT<Message>): void;
sendMessage(options:NewMessage, cb:CallbackT<string>): void;
setQueueAttributes(options:QueueOptions, cb:CallbackT<QueueAttributes>): 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;
}