-
Notifications
You must be signed in to change notification settings - Fork 3
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
8 changed files
with
137 additions
and
12 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,61 @@ | ||
const { ServiceBroker } = require("moleculer"); | ||
const QueueMixin = require("../../index"); | ||
|
||
let broker = new ServiceBroker({ | ||
logger: console, | ||
transporter: "TCP", | ||
}); | ||
|
||
const queueMixin = QueueMixin({ | ||
connection: "amqp://localhost", | ||
asyncActions: true, // Enable auto generate .async version for actions | ||
localPublisher: false, // Enable/Disable call this.actions.callAsync to call remote async | ||
}); | ||
|
||
broker.createService({ | ||
name: "consumer", | ||
version: 1, | ||
|
||
mixins: [ | ||
queueMixin, | ||
], | ||
|
||
settings: { | ||
amqp: { | ||
connection: "amqp://localhost", // You can also override setting from service setting | ||
}, | ||
}, | ||
|
||
actions: { | ||
hello: { | ||
queue: { // Enable queue for this action | ||
// Options for AMQP queue | ||
channel: { | ||
assert: { | ||
durable: true, | ||
}, | ||
prefetch: 0, | ||
}, | ||
consume: { | ||
noAck: false, | ||
}, | ||
}, | ||
params: { | ||
name: "string|convert:true|empty:false", | ||
}, | ||
async handler(ctx) { | ||
this.logger.info(`[CONSUMER] PID: ${process.pid} Received job with name=${ctx.params.name}`); | ||
return new Promise((resolve) => { | ||
setTimeout(() => { | ||
this.logger.info(`[CONSUMER] PID: ${process.pid} Processed job with name=${ctx.params.name}`); | ||
return resolve(`hello ${ctx.params.name}`); | ||
}, 1000); | ||
}); | ||
}, | ||
}, | ||
}, | ||
}); | ||
|
||
broker.start().then(() => { | ||
broker.repl(); | ||
}); |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,54 @@ | ||
const { ServiceBroker } = require("moleculer"); | ||
const QueueMixin = require("../../index"); | ||
|
||
let broker = new ServiceBroker({ | ||
logger: console, | ||
transporter: "TCP", | ||
}); | ||
|
||
const queueMixin = QueueMixin({ | ||
connection: "amqp://localhost", | ||
asyncActions: true, // Enable auto generate .async version for actions | ||
localPublisher: true, // Enable/Disable call this.actions.callAsync to call remote async | ||
}); | ||
|
||
broker.createService({ | ||
name: "publisher", | ||
version: 1, | ||
|
||
mixins: [ | ||
queueMixin, | ||
], | ||
|
||
settings: { | ||
amqp: { | ||
connection: "amqp://localhost", // You can also override setting from service setting | ||
}, | ||
}, | ||
|
||
async started() { | ||
// await broker.waitForServices({ name: "consumer", version: 1 }); | ||
|
||
let name = 1; | ||
setInterval(async () => { | ||
const response = await this.actions.callAsync({ | ||
// remote async action name | ||
action: "v1.consumer.hello", | ||
// `params` is the real param will be passed to original action | ||
params: { | ||
name, | ||
}, | ||
// `options` is the real options will be passed to original action | ||
options: { | ||
timeout: 2000, | ||
}, | ||
}); | ||
this.logger.info(`[PUBLISHER] PID: ${process.pid} Called job with name=${name} response=${JSON.stringify(response)}`); | ||
name++; | ||
}, 2000); | ||
} | ||
}); | ||
|
||
broker.start().then(() => { | ||
broker.repl(); | ||
}); |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters