-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathindex.js
54 lines (48 loc) · 1.57 KB
/
index.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
const { PubSub } = require("@google-cloud/pubsub");
const logger = require("lib-logger");
async function subscribe(topic_name, subscription_name, subscriber, options) {
const pubsub = new PubSub();
try {
logger.info(`Creating subscription ${subscription_name} on topic ${topic_name}`)
await pubsub.topic(topic_name).createSubscription(subscription_name);
logger.info(`Subscription ${subscription_name} on topic ${topic_name} created successfully`)
} catch(err) {
if (err.code === 6) {
logger.info(`subscription ${subscription_name} already exist, skipping creating it...`)
} else {
throw err;
}
}
// Default subscriber options
// Documentation here https://cloud.google.com/pubsub/docs/pull#config
const defaultOptions = {
flowControl: {
maxMessages: 10,
},
};
const subscription = pubsub.subscription(subscription_name, options || defaultOptions);
logger.info(`Listening...`);
subscription.on(`message`, async function processMessage(message) {
const { data } = message;
try {
let payload = {};
try {
payload = JSON.parse(data);
} catch (e) {
logger.error(`Invalid message received: ${data.toString("utf8")}`);
message.ack();
return;
}
await subscriber(payload, { message, topic_name });
message.ack();
logger.info(payload);
} catch (e) {
logger.error(`error processing: ${data.toString("utf8")} - ${e.message}`);
message.nack();
}
});
}
module.exports = {
subscribe
}
process.on('unhandledRejection', err => { throw err })