diff options
Diffstat (limited to 'worker.js')
| -rw-r--r-- | worker.js | 40 |
1 files changed, 24 insertions, 16 deletions
@@ -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 |
