summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-28 18:30:06 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-28 18:30:06 +0200
commitdfa0872dd7940b6d2759d8901274f459edf1b46e (patch)
tree761d559790e496cf750908fa9018671124435236
parent57ac9c2b03f9e46de6ea4a16eb3954868d1427c9 (diff)
downloadprocessor-dfa0872dd7940b6d2759d8901274f459edf1b46e.tar.gz
processor-dfa0872dd7940b6d2759d8901274f459edf1b46e.zip
rewrite in NodeJS
-rw-r--r--.gitignore118
-rw-r--r--.gitmodules3
-rw-r--r--Dockerfile15
m---------joblib0
-rw-r--r--model.js123
-rw-r--r--package.json18
-rw-r--r--queries/match.sql15
-rw-r--r--queries/participant.sql40
-rw-r--r--queries/player.sql34
-rw-r--r--queries/roster.sql20
-rw-r--r--requirements.txt6
-rw-r--r--worker.js124
-rw-r--r--worker.py86
13 files changed, 316 insertions, 286 deletions
diff --git a/.gitignore b/.gitignore
index a3401aa..00cbbdf 100644
--- a/.gitignore
+++ b/.gitignore
@@ -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
diff --git a/Dockerfile b/Dockerfile
index 6344e0f..9c3debd 100644
--- a/Dockerfile
+++ b/Dockerfile
@@ -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()