summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--package.json3
-rw-r--r--worker.js40
2 files changed, 26 insertions, 17 deletions
diff --git a/package.json b/package.json
index 54e49ae..ba632fc 100644
--- a/package.json
+++ b/package.json
@@ -8,7 +8,8 @@
"mysql": "^2.13.0",
"sequelize": "^3.30.4",
"sleep-promise": "^2.0.0",
- "winston": "^2.3.1"
+ "winston": "^2.3.1",
+ "winston-loggly-bulk": "^1.4.2"
},
"devDependencies": {},
"scripts": {
diff --git a/worker.js b/worker.js
index 77c243c..78f7f37 100644
--- a/worker.js
+++ b/worker.js
@@ -4,6 +4,7 @@
const amqp = require("amqplib"),
winston = require("winston"),
+ loggly = require("winston-loggly-bulk"),
Seq = require("sequelize"),
item_name_map = require("../orm/items"),
hero_name_map = require("../orm/heroes"),
@@ -12,6 +13,7 @@ const amqp = require("amqplib"),
const RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
+ LOGGLY_TOKEN = process.env.LOGGLY_TOKEN,
// matches + players, 5 players with 50 matches as default
BATCHSIZE = parseInt(process.env.BATCHSIZE) || 5 * (50 + 1),
// maximum number of elements to be inserted in one statement
@@ -22,13 +24,20 @@ const RABBITMQ_URI = process.env.RABBITMQ_URI,
const logger = new (winston.Logger)({
transports: [
new (winston.transports.Console)({
- timestamp: () => Date.now(),
- formatter: (options) => winston.config.colorize(options.level,
-`${new Date(options.timestamp()).toISOString()} ${options.level.toUpperCase()} ${(options.message? options.message:"")} ${(options.meta && Object.keys(options.meta).length? JSON.stringify(options.meta):"")}`)
+ timestamp: true,
+ colorize: true
})
]
});
+// loggly integration
+if (LOGGLY_TOKEN)
+ logger.add(winston.transports.Loggly, {
+ inputToken: LOGGLY_TOKEN,
+ subdomain: "kvahuja",
+ tags: ["backend", "processor"],
+ json: true
+ });
// helpers
const camelCaseRegExp = new RegExp(/([a-z])([A-Z]+)/g);
@@ -172,8 +181,8 @@ function flatten(obj) {
profiler.done("buffer filled");
profiler = undefined;
- logger.info("processing batch of %s players and %s matches",
- player_data.size, match_data.size);
+ logger.info("processing batch",
+ { players: player_data.size, matches: match_data.size });
// clean up to allow processor to accept while we wait for db
clearTimeout(idle_timer);
@@ -208,8 +217,8 @@ function flatten(obj) {
await Promise.each(player_objects, async (p) => {
const player = flatten(p);
if (processed_players.has(player.api_id)) {
- logger.warn("got player '%s' in additional region '%s'",
- player.name, player.shard_id);
+ logger.warn("got player in additional region",
+ { name: player.name, region: player.shard_id });
// see below, this is handling region changes
// when player objects end up in the same batch
const duplicate = player_records_direct.find(
@@ -219,8 +228,8 @@ function flatten(obj) {
player_records_direct.splice(
player_records_direct.indexOf(duplicate), 1);
} else {
- logger.warn("ignoring player '%s' from additional region '%s'",
- player.name, player.shard_id);
+ logger.warn("ignoring player from additional region",
+ { name: player.name, region: player.shard_id });
return;
}
} else {
@@ -229,8 +238,8 @@ function flatten(obj) {
// player objects that arrive here came from a search
// with search, updater can't update last_update
player.last_update = seq.fn("NOW");
- logger.info("processing player '%s' in '%s'",
- player.name, player.shard_id);
+ logger.info("processing player",
+ { name: player.name, region: player.shard_id });
// check whether there is a player in db
// that has a more recent `created_at`
// this is only the case with region changes
@@ -243,8 +252,8 @@ function flatten(obj) {
}
}});
if (count > 0) {
- logger.warn("ignoring player '%s' who seems to have switched from region '%s'",
- player.name, player.shard_id);
+ logger.warn("ignoring player who seems to have switched from region",
+ { name: player.name, region: player.shard_id });
return;
} else {
player_records_direct.push(player);
@@ -431,7 +440,7 @@ function flatten(obj) {
});
await seq.query("SET unique_checks=1");
- logger.info("acking batch with %s messages", msgs.size);
+ logger.info("acking batch", { size: msgs.size });
for (let msg of msgs)
await ch.ack(msg);
} catch (err) {
@@ -439,8 +448,7 @@ function flatten(obj) {
// it *must not* fail due to broken schema or missing dependency
// TODO: eliminate such cases earlier in the chain
// and immediately NACK those broken matches, requeueing only the rest
- logger.error("SQL error: %s, %j, %s",
- err.name, err.errors, err.parent.sql);
+ logger.error("SQL error", err);
for (let msg of msgs)
await ch.nack(msg, false, true);
return; // give up