summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--worker.js23
1 files changed, 16 insertions, 7 deletions
diff --git a/worker.js b/worker.js
index 183b552..c6d947c 100644
--- a/worker.js
+++ b/worker.js
@@ -12,7 +12,10 @@ var amqp = require("amqplib"),
var RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
BATCHSIZE = parseInt(process.env.BATCHSIZE) || 2 * 50 * (1 + 5), // matches + players
- IDLE_TIMEOUT = parseFloat(process.env.IDLE_TIMEOUT) || 1000; // ms
+ IDLE_TIMEOUT = parseFloat(process.env.IDLE_TIMEOUT) || 1000, // ms
+ PREMIUM_FEATURES = process.env.PREMIUM_FEATURES || false; // calculate on demand for non-premium users
+
+console.log("features for premium users activated", PREMIUM_FEATURES);
(async () => {
let seq, model, rabbit, ch;
@@ -61,6 +64,10 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI,
.map((role) => role_db_map[role.name] = role.id)
]);
+ // TODO expire this cache after some time
+ let premium_users = (await model.Gamer.findAll()).map((gamer) =>
+ gamer.name);
+
// as long as the queue is filled, msg are not ACKed
// server sends as long as there are less than `prefetch` unACKed
await ch.prefetch(BATCHSIZE);
@@ -331,12 +338,14 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI,
type: "participant"
})
));
- Promise.all(player_records.map(async (p) =>
- await ch.sendToQueue("crunch", new Buffer(p.api_id), {
- persistent: true,
- type: "player"
- })
- ));
+ Promise.all(player_records.map(async (p) => {
+ if (PREMIUM_FEATURES == false || premium_users.indexOf(p.name) != -1) {
+ await ch.sendToQueue("crunch", new Buffer(p.api_id), {
+ persistent: true,
+ type: "player"
+ })
+ }
+ }));
}
// Split participant API data into participant and participant_stats