-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathtelegram_dispatcher.js
More file actions
80 lines (66 loc) · 2.63 KB
/
Copy pathtelegram_dispatcher.js
File metadata and controls
80 lines (66 loc) · 2.63 KB
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
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
const signalNotifier = require('./dispatching/signals')
const signalHelper = require('./dispatching/signal-helper')
const Consumer = require('sqs-consumer')
const AWS = require('aws-sdk')
var moment = require('moment')
var database = require('./database')
database.connect()
var tradingAlertController = require('./controllers/tradingAlertsController')
var userController = require('./controllers/usersController')
userController.refreshCache().then(() => console.log('Users cache refreshed'))
AWS.config.update({
region: 'us-east-1',
accessKeyId: process.env.AWS_KEY,
secretAccessKey: process.env.AWS_SECRET
})
console.log('Starting telegram dispatching service 📬');
var aws_queue_url = process.env.LOCAL_ENV
? `${process.env.AWS_SQS_LOCAL_QUEUE_URL}`
: `${process.env.AWS_SQS_QUEUE_URL}`;
const app = Consumer.create({
queueUrl: aws_queue_url,
handleMessage: (message, done) => {
var signalValidity = signalHelper.checkValidity(message)
if (signalValidity.isValid) {
signalNotifier.notify(signalValidity.decoded_message_body).then((result) => {
var ta = {
signalId: result.signal_id,
awsSQSId: message.MessageId,
rejections: result.rejections,
reasons: result.reasons,
sent_at: result.sent_at
}
tradingAlertController.addTradingAlert(ta).then(() => { console.log(`[Notified] Message ${message.MessageId}`) })
console.log(`[Notified] Message ${message.MessageId}`)
}).catch((reason) => {
console.log(reason)
console.log(`[Not notified] Message ${message.MessageId}`)
})
} else {
var sent_at = moment()
if (moment(signalValidity.decoded_message_body.sent).isValid()) {
sent_at = moment(signalValidity.decoded_message_body.sent)
}
tradingAlertController.addTradingAlert({ signalId: signalValidity.decoded_message_body.id, reasons: signalValidity.reasons.split(','), awsSQSId: message.MessageId, sent_at: sent_at }).then(() => {
console.log(`[Invalid][Sent at ${signalValidity.decoded_message_body.sent}] SQS message ${message.MessageId} for signal ${signalValidity.decoded_message_body.id} ${signalValidity.reasons}`)
}).catch((err) => console.log(err))
}
done()
},
sqs: new AWS.SQS()
})
app.on('message_received', (msg) => {
app.handleMessage(msg, function (err) {
if (err) console.log(err);
})
})
app.on('message_processed', (msg) => {
console.log(`[Processed] Message ${msg.MessageId}`);
})
app.on('processing_error', (err, signal) => {
console.log(err.message);
})
app.on('error', (err) => {
console.log(err.message);
})
signalHelper.init().then(() => app.start())