diff options
| -rw-r--r-- | .gitignore | 118 | ||||
| -rw-r--r-- | .gitmodules | 3 | ||||
| -rw-r--r-- | Dockerfile | 15 | ||||
| m--------- | joblib | 0 | ||||
| -rw-r--r-- | model.js | 123 | ||||
| -rw-r--r-- | package.json | 18 | ||||
| -rw-r--r-- | queries/match.sql | 15 | ||||
| -rw-r--r-- | queries/participant.sql | 40 | ||||
| -rw-r--r-- | queries/player.sql | 34 | ||||
| -rw-r--r-- | queries/roster.sql | 20 | ||||
| -rw-r--r-- | requirements.txt | 6 | ||||
| -rw-r--r-- | worker.js | 124 | ||||
| -rw-r--r-- | worker.py | 86 |
13 files changed, 316 insertions, 286 deletions
@@ -1,93 +1,59 @@ -# Byte-compiled / optimized / DLL files -__pycache__/ -*.py[cod] -*$py.class +# Logs +logs +*.log +npm-debug.log* +yarn-debug.log* +yarn-error.log* -# C extensions -*.so +# Runtime data +pids +*.pid +*.seed +*.pid.lock -# Distribution / packaging -.Python -env/ -build/ -develop-eggs/ -dist/ -downloads/ -eggs/ -.eggs/ -lib/ -lib64/ -parts/ -sdist/ -var/ -wheels/ -*.egg-info/ -.installed.cfg -*.egg +# Directory for instrumented libs generated by jscoverage/JSCover +lib-cov -# PyInstaller -# Usually these files are written by a python script from a template -# before PyInstaller builds the exe, so as to inject date/other infos into it. -*.manifest -*.spec +# Coverage directory used by tools like istanbul +coverage -# Installer logs -pip-log.txt -pip-delete-this-directory.txt +# nyc test coverage +.nyc_output -# Unit test / coverage reports -htmlcov/ -.tox/ -.coverage -.coverage.* -.cache -nosetests.xml -coverage.xml -*,cover -.hypothesis/ +# Grunt intermediate storage (http://gruntjs.com/creating-plugins#storing-task-files) +.grunt -# Translations -*.mo -*.pot +# Bower dependency directory (https://bower.io/) +bower_components -# Django stuff: -*.log -!/logs -!run*.log -local_settings.py +# node-waf configuration +.lock-wscript -# Flask stuff: -instance/ -.webassets-cache +# Compiled binary addons (http://nodejs.org/api/addons.html) +build/Release -# Scrapy stuff: -.scrapy +# Dependency directories +node_modules/ +jspm_packages/ -# Sphinx documentation -docs/_build/ +# Typescript v1 declaration files +typings/ -# PyBuilder -target/ +# Optional npm cache directory +.npm -# Jupyter Notebook -.ipynb_checkpoints +# Optional eslint cache +.eslintcache -# pyenv -.python-version +# Optional REPL history +.node_repl_history -# celery beat schedule file -celerybeat-schedule +# Output of 'npm pack' +*.tgz -# dotenv -.env +# Yarn Integrity file +.yarn-integrity -# virtualenv -.venv/ -venv/ -ENV/ - -# Spyder project settings -.spyderproject +# dotenv environment variables file +.env -# Rope project settings -.ropeproject diff --git a/.gitmodules b/.gitmodules index 52948d0..e69de29 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,3 +0,0 @@ -[submodule "joblib"] - path = joblib - url = https://github.com/vainglorygame/joblib @@ -1,6 +1,9 @@ -FROM python:3.6-alpine -RUN apk add --no-cache postgresql-dev gcc python3-dev musl-dev -ADD requirements.txt /code/requirements.txt -WORKDIR /code -RUN pip install -r requirements.txt -CMD ["python", "worker.py"] +FROM node:7.7-alpine + +RUN mkdir -p /usr/src/app +WORKDIR /usr/src/app + +COPY . /usr/src/app +RUN npm install && npm cache clean + +CMD ["node", "worker.js"] diff --git a/joblib b/joblib deleted file mode 160000 -Subproject 7ced841aa44a7dd7bd2d3a68ceb989ef488ce44 diff --git a/model.js b/model.js new file mode 100644 index 0000000..fa7c517 --- /dev/null +++ b/model.js @@ -0,0 +1,123 @@ +/* jshint esnext:true */ +'use strict'; + +var JsonField = require("sequelize-json"); + +module.exports = (seq, Seq) => { + let Player = seq.define("player", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + name: Seq.STRING, + shardId: Seq.STRING, + /* stats */ + level: Seq.INTEGER, + lifetimeGold: Seq.DECIMAL, + lossStreak: Seq.INTEGER, + played: Seq.INTEGER, + played_ranked: Seq.INTEGER, + winStreak: Seq.INTEGER, + wins: Seq.INTEGER, + xp: Seq.INTEGER + }, { + freezeTableName: true, + underscored: true + }); + let Team = seq.define("team", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + name: Seq.STRING, + shardId: Seq.STRING + /* stats */ + }, { + freezeTableName: true, + underscored: true + }); + let Asset = seq.define("asset", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + URL: Seq.STRING, + contentType: Seq.STRING, + createdAt: Seq.DATE, + description: Seq.TEXT, + filename: Seq.STRING, + name: Seq.STRING + }, { + freezeTableName: true, + underscored: true, + createdAt: false + }); + let Match = seq.define("match", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + createdAt: Seq.DATE, + duration: Seq.INTEGER, + gameMode: Seq.STRING, + patchVersion: Seq.STRING, + shardId: Seq.STRING, + /* stats */ + endGameReason: Seq.STRING, + queue: Seq.STRING + }, { + freezeTableName: true, + underscored: true, + createdAt: false + }); + let Roster = seq.define("roster", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + /* stats */ + acesEarned: Seq.INTEGER, + gold: Seq.INTEGER, + heroKills: Seq.INTEGER, + krakenCaptures: Seq.INTEGER, + side: Seq.STRING, + turretKills: Seq.INTEGER, + turretsRemaining: Seq.INTEGER, + }, { + freezeTableName: true, + underscored: true + }); + let Participant = seq.define("participant", { + id: { type: Seq.STRING, unique: true, primaryKey: true }, + /* attributes */ + actor: Seq.STRING, + /* stats */ + assists: Seq.INTEGER, + crystalMineCaptures: Seq.INTEGER, + deaths: Seq.INTEGER, + farm: Seq.DECIMAL, + firstAfkTime: Seq.INTEGER, + goldMineCaptures: Seq.INTEGER, + itemGrants: JsonField(seq, "Participant", "itemGrants"), + itemSells: JsonField(seq, "Participant", "itemSells"), + itemUses: JsonField(seq, "Participant", "itemUses"), + items: JsonField(seq, "Participant", "items"), + jungleKills: Seq.INTEGER, + karmaLevel: Seq.INTEGER, + kills: Seq.INTEGER, + krakenCaptures: Seq.INTEGER, + level: Seq.INTEGER, + minionKills: Seq.INTEGER, + nonJungleMinionKills: Seq.INTEGER, + skillTier: Seq.INTEGER, + skinKey: Seq.STRING, + turretCaptures: Seq.INTEGER, + wentAfk: Seq.BOOLEAN, + winner: Seq.BOOLEAN, + }, { + freezeTableName: true, + underscored: true + }); + + Match.hasMany(Roster, {as: "Rosters"}); + Match.hasMany(Asset, {as: "Assets"}); + Asset.belongsTo(Match); + Roster.belongsTo(Match); + Roster.hasMany(Participant, {as: "Participants"}); + Roster.hasOne(Team, {as: "Team"}); + Team.belongsTo(Roster); + Participant.belongsTo(Roster); + Participant.belongsTo(Player); + + return {Match, Roster, Team, Participant, Player, Asset}; +}; diff --git a/package.json b/package.json new file mode 100644 index 0000000..ff3adc7 --- /dev/null +++ b/package.json @@ -0,0 +1,18 @@ +{ + "name": "processor", + "version": "2.0.0", + "description": "", + "main": "model.js", + "dependencies": { + "amqplib": "^0.5.1", + "mysql": "^2.13.0", + "sequelize": "^3.30.4", + "superagent-jsonapify": "^1.4.4" + }, + "devDependencies": {}, + "scripts": { + "test": "echo \"Error: no test specified\" && exit 1" + }, + "author": "schneefux", + "license": "UNLICENSED" +} diff --git a/queries/match.sql b/queries/match.sql deleted file mode 100644 index fc8ab7e..0000000 --- a/queries/match.sql +++ /dev/null @@ -1,15 +0,0 @@ -SELECT - -match.id AS "api_id", -(match.data->'data'->'attributes'->>'createdAt')::timestamp as "created_at", -(match.data->'data'->'attributes'->>'duration')::int AS "duration", -match.data->'data'->'attributes'->>'gameMode' AS "game_mode", -COALESCE(NULLIF(match.data->'data'->'attributes'->>'patch_version', ''), '2.2') AS "patch_version", -match.data->'data'->'attributes'->>'shardId' AS "shard_id", -match.data->'data'->'attributes'->'stats'->>'endGameReason' AS "end_game_reason", -match.data->'data'->'attributes'->'stats'->>'queue' AS "queue", - -match.data->'relations'->0->'data'->>'id' as "roster_1", -match.data->'relations'->1->'data'->>'id' as "roster_2" - -FROM match WHERE id=$1 diff --git a/queries/participant.sql b/queries/participant.sql deleted file mode 100644 index 8e86e3d..0000000 --- a/queries/participant.sql +++ /dev/null @@ -1,40 +0,0 @@ -WITH rosters AS (SELECT - (match.data->'data'->'attributes'->>'createdAt')::timestamp AS matchdate, - JSONB_ARRAY_ELEMENTS(match.data->'relations') AS roster -FROM match WHERE id=$1), -participants AS (SELECT - matchdate, - rosters.roster->'data'->>'id' AS rosterid, - JSONB_ARRAY_ELEMENTS(rosters.roster->'relations') AS participant -FROM rosters) -SELECT - -participant->'data'->>'id' AS "api_id", -rosterid AS "roster_api_id", -participant->'relations'->0->'data'->>'id' AS "player_api_id", -matchdate AS "created_at", - -participant->'data'->'attributes'->>'actor' AS "hero", -(participant->'data'->'attributes'->'stats'->>'assists')::int AS "assists", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'crystalMineCaptures', ''), '0')::int AS "crystal_mine_captures", -(participant->'data'->'attributes'->'stats'->>'deaths')::int AS "deaths", -(participant->'data'->'attributes'->'stats'->>'farm')::float AS "farm", -(participant->'data'->'attributes'->'stats'->>'firstAfkTime')::float AS "first_afk_time", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'goldMineCaptures', ''), '0')::int AS "gold_mine_captures", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'jungleKills', ''), '0')::int AS "jungle_kills", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'karmaLevel', ''), '0')::int AS "karma_level", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'kills', ''), '0')::int AS "kills", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'krakenCaptures', ''), '0')::int AS "kraken_captures", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'level', ''), '0')::int AS "level", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'minionKills', ''), '0')::int AS "minion_kills", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'skillTier', ''), '0')::int AS "skill_tier", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'skinKey', ''), '0') AS "skin_key", -COALESCE(NULLIF(participant->'data'->'attributes'->'stats'->>'turretCaptures', ''), '0')::int AS "turret_kills", -(participant->'data'->'attributes'->'stats'->>'wentAfk')::bool AS "went_afk", -(participant->'data'->'attributes'->'stats'->>'winner')::bool AS "winner", -participant->'data'->'attributes'->'stats'->'itemGrants' AS "item_grants", -participant->'data'->'attributes'->'stats'->'itemSells' AS "item_sells", -participant->'data'->'attributes'->'stats'->'itemUses' AS "item_uses", -participant->'data'->'attributes'->'stats'->'items' AS "items" - -FROM participants diff --git a/queries/player.sql b/queries/player.sql deleted file mode 100644 index d718391..0000000 --- a/queries/player.sql +++ /dev/null @@ -1,34 +0,0 @@ -WITH rosters AS (SELECT - match.data->'data'->'attributes'->>'shardId' AS shardid, - (match.data->'data'->'attributes'->>'createdAt')::timestamp AS matchdate, - JSONB_ARRAY_ELEMENTS(match.data->'relations') AS roster -FROM match WHERE id=$1), -participants AS (SELECT - shardid, - matchdate, - JSONB_ARRAY_ELEMENTS(rosters.roster->'relations') AS participant -FROM rosters), -players AS (SELECT - participants.participant->'data'->'attributes'->'stats'->>'skillTier' AS skill_tier, - shardid, - matchdate, - JSONB_ARRAY_ELEMENTS(participants.participant->'relations') AS player -FROM participants) - -SELECT - -players.player->'data'->>'id' AS "api_id", -shardid AS "shard_id", -players.player->'data'->'attributes'->>'name' AS "name", -(players.player->'data'->'attributes'->'stats'->>'level')::int AS "level", -(players.player->'data'->'attributes'->'stats'->>'xp')::float::int AS "xp", -(players.player->'data'->'attributes'->'stats'->>'played')::int AS "played", -(players.player->'data'->'attributes'->'stats'->>'played_ranked')::int AS "played_ranked", -(players.player->'data'->'attributes'->'stats'->>'played')::int - (players.player->'data'->'attributes'->'stats'->>'played_ranked')::int AS "played_casual", -(players.player->'data'->'attributes'->'stats'->>'wins')::int AS "wins", -(players.player->'data'->'attributes'->'stats'->>'lifetimeGold')::float AS "lifetime_gold", -matchdate AS "last_match_created_date", -skill_tier AS "skill_tier", -0 AS "streak" - -FROM players diff --git a/queries/roster.sql b/queries/roster.sql deleted file mode 100644 index f6500fb..0000000 --- a/queries/roster.sql +++ /dev/null @@ -1,20 +0,0 @@ -WITH rosters AS (SELECT - id AS matchid, - JSONB_ARRAY_ELEMENTS(match.data->'relations') AS roster -FROM match WHERE id=$1) -SELECT -roster->'data'->>'id' AS "api_id", -matchid AS "match_api_id", -(roster->'data'->'attributes'->'stats'->>'acesEarned')::int AS "aces_earned", -(roster->'data'->'attributes'->'stats'->>'gold')::int AS "gold", -(roster->'data'->'attributes'->'stats'->>'heroKills')::int AS "hero_kills", -(roster->'data'->'attributes'->'stats'->>'krakenCaptures')::int AS "kraken_captures", -roster->'data'->'attributes'->'stats'->>'side' AS "side", -roster->'data'->'attributes'->'stats'->>'side' AS "team_color", -(roster->'data'->'attributes'->'stats'->>'turretKills')::int AS "turret_kills", -(roster->'data'->'attributes'->'stats'->>'turretsRemaining')::int AS "turrets_remaining", - -roster->'relations'->0->'data'->>'id' AS "participant_1", -roster->'relations'->1->'data'->>'id' AS "participant_2", -roster->'relations'->2->'data'->>'id' AS "participant_3" -FROM rosters diff --git a/requirements.txt b/requirements.txt deleted file mode 100644 index b549b13..0000000 --- a/requirements.txt +++ /dev/null @@ -1,6 +0,0 @@ -appdirs==1.4.3 -packaging==16.8 -pika==0.10.0 -psycopg2==2.7.1 -pyparsing==2.2.0 -six==1.10.0 diff --git a/worker.js b/worker.js new file mode 100644 index 0000000..31894da --- /dev/null +++ b/worker.js @@ -0,0 +1,124 @@ +#!/usr/bin/node +/* jshint esnext:true */ +'use strict'; + +var amqp = require("amqplib"), + Seq = require("sequelize"), + Bluebird = require("bluebird"), + jsonapi = Bluebird.promisifyAll(require("superagent-jsonapify/common")); + +var MADGLORY_TOKEN = process.env.MADGLORY_TOKEN, + RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost", + DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite"; +if (MADGLORY_TOKEN == undefined) throw "Need an API token"; + +(async () => { + let seq = new Seq(DATABASE_URI), + model = require("./model")(seq, Seq), + rabbit = await amqp.connect(RABBITMQ_URI), + ch = await rabbit.createChannel(); + + /* recreate for debugging + await seq.query("SET FOREIGN_KEY_CHECKS=0"); + await seq.sync({force: true}); + */ + await seq.sync(); + + await ch.assertQueue("process", {durable: true}); + await ch.prefetch(1); + + ch.consume("process", async (msg) => { + let data = await jsonapi.parse(msg.content); + + // TODO commit less often if possible, avoid deadlocks + let transaction = await seq.transaction({ autocommit: false }); + + data.data.forEach((match_data) => { + function flatten(obj) { + let attrs = obj.attributes || {}, + stats = attrs.stats || {}, + o = Object.assign({}, obj, attrs, stats); + delete o.type; + delete o.attributes; + delete o.stats; + delete o.relationships; + return o; + } + let match = JSON.parse(JSON.stringify(match_data)); // deep clone + + /* bring jsonapi response into our db structure-like shape */ + match.rosters = match.rosters.map((roster) => { + roster.participants = roster.participants.map((participant) => { + participant.player = flatten(participant.player); + return flatten(participant); + }); + return flatten(roster); + }); + match = flatten(match); + + /* upsert everything */ + model.Match.upsert(match, { + include: [ + { + model: model.Roster, + as: "Rosters", + include: [ + { + model: model.Participant, + as: "Participants", + include: [ + model.Player + ] + }, + { + model: model.Team, + as: "Team" + } + ] + }, + { + model: model.Asset, + as: "Assets" + } + ] + }); + + match.rosters.forEach((roster) => { + model.Roster.upsert(roster, { + include: [ + { + model: model.Participant, + as: "Participants", + include: [ + model.Player + ] + }, + { + model: model.Team, + as: "Team" + } + ] + }); + + roster.participants.forEach((participant) => { + model.Participant.upsert(participant, { + include: [ + model.Player + ] + }); + + model.Player.upsert(participant.player); + }); + + if (roster.team != null) model.Team.upsert(roster.team); + }); + + match.assets.forEach((asset) => { + model.Asset.upsert(asset); + }); + }); + + await transaction.commit(); // TODO rollback on err + ch.ack(msg); + }, { noAck: false }); +})(); diff --git a/worker.py b/worker.py deleted file mode 100644 index 9a00f06..0000000 --- a/worker.py +++ /dev/null @@ -1,86 +0,0 @@ -#!/usr/bin/python3 - -import os -import logging -import psycopg2 -import psycopg2.extras -import psycopg2.extensions - -import joblib.joblib - - -RABBIT = { - "host": os.environ.get("RABBITMQ_HOST"), - "port": os.environ.get("RABBITMQ_PORT"), - "credentials": os.environ.get("RABBITMQ_CREDS") -} - -DB = { - "host": os.environ.get("POSTGRESQL_HOST") or "vaindock_postgres_web", - "port": os.environ.get("POSTGRESQL_PORT") or 5432, - "user": os.environ.get("POSTGRESQL_USER") or "vainweb", - "password": os.environ.get("POSTGRESQL_PASSWORD") or "vainweb", - "dbname": os.environ.get("POSTGRESQL_DB") or "vainsocial-web" -} - - -class Processor(joblib.joblib.Worker): - def __init__(self): - super().__init__(jobtype="process") - - def connect(self, rabbit, db): - """Connect to database.""" - logging.warning("connecting to database") - super().connect(**rabbit) - self._con = psycopg2.connect(**db) - - def commit(self, failed): - self._con.commit() - - def work(self, payload): - """Finish a job.""" - obj_id = payload["id"] - obj_type = payload["type"] - o = payload["data"]["attributes"] - r = payload["data"].get("relationships") - - c = self._con.cursor( - cursor_factory=psycopg2.extras.DictCursor) - - d = {} - if obj_type == "match": - d["api_id"] = obj_id - d["created_at"] = o["createdAt"] - d["duration"] = o["duration"] - d["game_mode"] = o["gameMode"] - d["patch_version"] = o.get("patchVersion") or "2.2" - d["shard_id"] = o["shardId"] - d["end_game_reason"] = o["stats"]["endGameReason"] - d["queue"] = o["stats"]["queue"] - d["roster_1"] = r["rosters"]["data"][0]["id"] - d["roster_2"] = r["rosters"]["data"][1]["id"] - - if d == {}: - raise Exception - # TODO - - # TODO player upsert - # TODO request compile & preload - l = [(c, v) for c, v in d.items()] - columns = ",".join([t[0] for t in l]) - values = tuple([t[1] for t in l]) - c.execute("INSERT INTO " + obj_type + "(%s) VALUES %s", - ([psycopg2.extensions.AsIs(columns)] + [values])) - c.close() - - -def startup(): - worker = Processor() - worker.connect(RABBIT, DB) - worker.setup() - worker.run() - -logging.basicConfig(level=logging.DEBUG) - -if __name__ == "__main__": - startup() |
