summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-28 18:32:03 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-28 18:32:03 +0200
commitb48d68e5a2129512b62c13c73f6312d684256bab (patch)
tree96ddb71b108cf8b9b381d2e6c9d80324e9d50e2c
parentd830686ced729dc37b05293353d838e697dcf60a (diff)
downloadbridge-b48d68e5a2129512b62c13c73f6312d684256bab.tar.gz
bridge-b48d68e5a2129512b62c13c73f6312d684256bab.zip
rewrite for AMQP
-rw-r--r--Dockerfile9
-rw-r--r--api.js395
-rw-r--r--index.html6
-rw-r--r--package.json19
-rw-r--r--server.js212
5 files changed, 233 insertions, 408 deletions
diff --git a/Dockerfile b/Dockerfile
new file mode 100644
index 0000000..1289cc4
--- /dev/null
+++ b/Dockerfile
@@ -0,0 +1,9 @@
+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", "server.js"]
diff --git a/api.js b/api.js
deleted file mode 100644
index 1ef7020..0000000
--- a/api.js
+++ /dev/null
@@ -1,395 +0,0 @@
-#!/usr/bin/env node
-/* jshint esnext: true */
-
-var request = require("request-promise");
-
-var pg = require("pg").native;
-var Pool = pg.Pool;
-var db_config_raw = {
- user: process.env.POSTGRESQL_SOURCE_USER || "vainraw",
- password: process.env.POSTGRESQL_SOURCE_PASSWORD || "vainraw",
- host: process.env.POSTGRESQL_SOURCE_HOST || "localhost",
- database: process.env.POSTGRESQL_SOURCE_DB || "vainsocial-raw",
- port: process.env.POSTGRESQL_SOURCE_PORT || 5433,
- min: 2,
- max: 10
-};
-var db_config_web = {
- user: process.env.POSTGRESQL_DEST_USER || "vainweb",
- password: process.env.POSTGRESQL_DEST_PASSWORD || "vainweb",
- host: process.env.POSTGRESQL_DEST_HOST || "localhost",
- database: process.env.POSTGRESQL_DEST_DB || "vainsocial-web",
- port: process.env.POSTGRESQL_DEST_PORT || 5432,
- min: 4,
- max: 10
-};
-var pool_raw = new Pool(db_config_raw);
-var pool_web = new Pool(db_config_web);
-
-var app = require("express")();
-var http = require("http").Server(app);
-var io = require("socket.io")(http);
-
-var APITOKEN = process.env.VAINSOCIAL_APITOKEN;
-if (APITOKEN == undefined) throw "Need a valid API token!";
-
-http.listen(8080);
-
-
-function sleep(ms) {
- return new Promise(resolve => {
- setTimeout(resolve, ms)
- });
-}
-
-/* API helper */
-/* search for player name across all shards */
-async function api_playerByAttr(attr, val) {
- var regions = ["na", "eu", "sg", "sa", "ea"],
- finds = [],
- filter = "filter[" + attr + "]";
-
- for (let region of regions) {
- var options = {
- uri: "https://api.dc01.gamelockerapp.com/shards/" + region + "/players",
- headers: {
- "X-Title-Id": "semc-vainglory",
- "Authorization": APITOKEN
- },
- qs: {},
- json: true,
- gzip: true
- };
- options.qs[filter] = val;
- retry = true;
- while (retry) {
- try {
- res = await request(options);
- finds.push({
- "region": res.data[0].attributes.shardId,
- "id": res.data[0].id,
- "name": res.data[0].attributes.name,
- "last_update": res.data[0].attributes.createdAt,
- "source": "api"
- });
- retry = false;
- } catch (err) {
- if (err.statusCode == 429) {
- await sleep(100);
- retry = true;
- } else if (err.statusCode == 404) {
- retry = false;
- } else {
- console.error(err);
- retry = false;
- }
- // TODO
- }
- }
- }
-
- if (finds.length == 0)
- return undefined;
-
- // due to an API bug, many players are also present in NA
- // TODO: get history for all regions in case of region transfer
- finds.sort((a, b) => { return a.last_update < b.last_update; });
- return finds[0];
-}
-async function api_playerByName(name) {
- return await api_playerByAttr("playerNames", name);
-}
-async function api_playerById(id) {
- return await api_playerByAttr("playerIds", id);
-}
-
-
-/* DB helper */
-/* retry until no serialization error */
-async function db_serialized(con, query, data) {
- var commit_success = false,
- res;
- do {
- /* transaction begin */
- try {
- await con.query("BEGIN");
- await con.query("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE");
- res = await con.query(query, data);
- await con.query("COMMIT");
- commit_success = true;
- } catch (err) {
- await con.query("ROLLBACK");
- // serialization error - expected, else rethrow
- console.error(err);
- if (err.sqlState != "40001") throw err;
- // important! with pg, it is err.code,
- // with pg-native, it is err.sqlState
- console.log("serialization error, retrying");
- }
- /* transaction end */
- } while (!commit_success);
- return res;
-}
-
-/* find player by id or by name in db */
-async function db_playerByAttr(attr, val) {
- var web = await pool_web.connect(),
- player = await web.query(`
- SELECT name, api_id, shard_id, last_match_created_date
- FROM player WHERE ` + attr + `=$1
- `, [val]);
- web.release();
-
- if (player.rows.length > 0) {
- return {
- "name": player.rows[0].name,
- "id": player.rows[0].api_id,
- "region": player.rows[0].shard_id,
- "last_update": player.rows[0].last_match_created_date,
- "source": "db"
- };
- }
- return undefined;
-}
-
-async function db_playerByName(name) {
- return await db_playerByAttr("name", name);
-}
-async function db_playerById(id) {
- return await db_playerByAttr("api_id", id);
-}
-
-/* returns a player by name from db or API */
-async function playerByName(name) {
- var player = await db_playerByName(name);
- if (player != undefined) return player;
-
- console.log("player '" + name + "' not found in db");
- player = await api_playerByName(name);
- if (player != undefined) return player;
-
- console.log("player '" + name + "' not found in API");
- return undefined;
-}
-/* returns a player by id from db or API */
-async function playerById(id) {
- var player = await db_playerById(id);
- if (player != undefined) return player;
-
- console.log("player with id '" + name + "' not found in db");
- player = await api_playerById(id);
- if (player != undefined) return player;
-
- console.log("player with id '" + name + "' not found in API");
- return undefined;
-}
-/* returns true if a player has pending jobs */
-async function anyJobsRunningFor(con, name, id) {
- result = await con.query(`
- SELECT COUNT(*)>0 AS jobs_running FROM jobs
- WHERE
- (
- (type='grab' AND payload->'params'->>'filter[playerNames]'=$1) OR
- (type='process' AND payload->>'playername'=$1) OR
- (type='compile' AND payload->>'type'='player' AND payload->>'id'=$2)
- ) AND status<>'finished' AND status<>'failed'
- `, [name, id]); // TODO improve dependency tracking
- console.log("jobs running for %s: %s", name, result.rows[0].jobs_running);
- return result.rows[0].jobs_running;
-}
-
-
-/* update request helpers */
-/* upsert a job */
-async function upsertGrabjob(payload) {
- var raw = await pool_raw.connect(),
- job, jobs_running;
-
- jobs_running = await anyJobsRunningFor(raw,
- payload.params["filter[playerNames]"],
- payload.params["filter[playerIds]"])
-
- if (!jobs_running) {
- // this job is currently not running, insert it
- job = await db_serialized(raw, `
- INSERT INTO jobs(type, payload, priority)
- VALUES('grab', $1, 0)
- RETURNING id
- `, [payload]);
- // wake apigrabber up
- await raw.query(`NOTIFY grab_open`);
- console.log("job requested: '%j'", payload);
- }
-
- raw.release();
-}
-
-async function playerRequestUpdate(name, id) {
- if (name == undefined && id == undefined) // fail hard and die
- throw "playerRequestUpdate needs either name or id";
-
- var player;
- if (id != undefined) // prefer id over name
- player = await playerById(id);
- else
- player = await playerByName(name);
- if (player == undefined)
- return undefined;
-
- console.log("updating player '%j'", player);
-
- // if last_update is from our db, use it as a start, else get the whole history
- if (player.source == "api" || player.last_update == undefined)
- player.last_update = new Date(value=0); // forever ago
- /* comment out on shutter's machine
- if (player.source == "db")
- player.last_update.setMinutes(player.last_update.getMinutes() - new Date().getTimezoneOffset()); // TODO workaround for my broken db schema
- */
-
- // add 1s, because createdAt-start <= x <= createdAt-end
- // so without the +1s, we'd always get the last_match_created_date match back
- player.last_update.setSeconds(player.last_update.getSeconds() + 1);
-
- var timedelta_minutes = ((new Date()) - player.last_update) / 1000 / 60;
- if (timedelta_minutes < 30) {
- console.log("player '" + player.name + "' update skipped");
- return;
- }
-
- var payload = {
- "region": player.region,
- "params": {
- "filter[playerIds]": player.id,
- "filter[playerNames]": player.name, // TODO remove in 2.0 - backwards compat
- "filter[createdAt-start]": player.last_update.toISOString(),
- "filter[gameMode]": "casual,ranked"
- }
- };
- upsertGrabjob(payload); // TODO do something with response?
-
- return player;
-}
-async function playerRequestUpdateByName(name) {
- return await playerRequestUpdate(name, undefined);
-}
-async function playerRequestUpdateById(id) {
- return await playerRequestUpdate(undefined, id);
-}
-
-/* routes */
-/* request a grab job */
-app.get("/api/player/name/:name", async (req, res) => {
- player = await playerRequestUpdateByName(req.params.name);
- if (player == undefined) res.sendStatus(404);
- else res.json(player);
-});
-app.get("/api/player/id/:id", async (req, res) => {
- player = await playerRequestUpdateById(req.params.id);
- if (player == undefined) res.sendStatus(404);
- else res.json(player);
-});
-
-/* internal monitoring */
-app.get("/", async (req, res) => {
- res.sendFile(__dirname + "/index.html");
-});
-
-
-/* notifications from database */
-async function listen() {
- var client = new pg.Client(db_config_raw);
- await client.connect();
-
- /* job status change notification listener */
- client.on('notification', async (msg) => {
- var raw = await pool_raw.connect(),
- jobs;
-
- await raw.query("BEGIN"); // TODO catch error & rollback
- // find all interesting jobs, delete them & forward their notification
- if (msg.channel == "grab_failed") {
- jobs = await raw.query(`
- WITH grab_failed AS (
- DELETE FROM jobs WHERE
- type='grab' AND status='failed' AND payload->'error'->>'title'='Not Found'
- RETURNING
- payload->'params'->>'filter[playerIds]' AS player_id,
- payload->'params'->>'filter[playerNames]' AS player_name
- )
- SELECT DISTINCT * FROM grab_failed
- `);
- }
- if (msg.channel == "process_finished") {
- jobs = await raw.query(`
- WITH process_finished AS (
- DELETE FROM jobs WHERE
- type='process' AND status='finished'
- RETURNING payload->>'playername' AS player_name
- )
- SELECT DISTINCT * FROM process_finished
- `);
- }
- if (msg.channel == "compile_finished") {
- jobs = await raw.query(`
- WITH compile_finished AS (
- DELETE FROM jobs WHERE
- type='compile' AND payload->>'type'='player' AND status='finished'
- RETURNING payload->>'id' AS player_id
- )
- SELECT DISTINCT * FROM compile_finished
- `);
- }
-
- if (jobs == undefined) {
- console.log("notification was about no jobs, exiting");
- return; // nothing to do
- }
-
- console.log("forwarding notifications for %s", msg.channel);
- // forward notification to all playername / playerid channels
- for (let job of jobs.rows) {
- var name = job.player_name;
- var id = job.player_id;
- if (name == undefined && id == undefined) throw "notification needs either name or ID";
-
- // attempt to fill gaps (TODO will not be needed in 2.0)
- if (name == undefined && id != undefined) {
- var player = await db_playerById(id);
- if (player != undefined)
- name = player.name;
- else console.log("id %s: warning! player had a job, but doesn't exist in db yet", id);
- }
- if (id == undefined && name != undefined) {
- var player = await db_playerByName(name);
- if (player != undefined)
- id = player.id;
- else console.log("name %s: warning! player had a job, but doesn't exist in db yet", name);
- }
- console.log("sending '%s' notification for player '%s' ('%s')", msg.channel, name, id);
-
- if (name != undefined) io.emit(name, msg.channel);
- if (id != undefined) io.emit(id, msg.channel);
-
- // don't give up on player not being found in db (= playerByAttr returns undefined),
- // be optimistic and try to notify what we can
-
- if (name != undefined && id != undefined) {
- // re - above: when we get grab_failed/compile_finished, the player will be COMMITted already for sure, so we don't miss 'done'
- if (msg.channel == "grab_failed" || (msg.channel == "compile_finished" && !await anyJobsRunningFor(raw, name, id))) {
- io.emit(name, "done");
- io.emit(id, "done");
- }
- }
- }
-
- await raw.query("COMMIT");
- raw.release();
- });
- client.query("LISTEN process_finished");
- client.query("LISTEN compile_finished");
- client.query("LISTEN analyze_finished");
- client.query("LISTEN grab_failed");
- // keep open forever
-}
-
-listen();
diff --git a/index.html b/index.html
index 6ccdbb1..350cf8a 100644
--- a/index.html
+++ b/index.html
@@ -10,10 +10,10 @@
</form>
<ul id="updates"></ul>
- <script src="/socket.io/socket.io.js"></script>
+ <!-- <script src="/socket.io/socket.io.js"></script> -->
<script src="https://code.jquery.com/jquery-3.1.1.min.js"></script>
<script>
- var socket = io();
+ //var socket = io();
$("#update-form").submit(function(e) {
var name = $("input:first").val();
$.get("/api/player/name/" + name).done(function() {
@@ -23,6 +23,7 @@
.css("color", "black")
);
/* subscribe to socket notifications */
+ /*
socket.on(name, function(msg) {
if (msg == "grab_failed") {
$("#updates").append(
@@ -53,6 +54,7 @@
);
}
});
+ */
}).fail(function() {
$("#updates").append(
$("<li>")
diff --git a/package.json b/package.json
index bd514d2..eca65e0 100644
--- a/package.json
+++ b/package.json
@@ -1,19 +1,16 @@
{
- "name": "vainsocial-updater",
+ "name": "updater",
"version": "2.0.0",
- "description": "Vainsocial frontend to backend communication service",
- "main": "api.js",
+ "description": "",
+ "main": "server.js",
"dependencies": {
- "express": "^4.15.2",
- "pg": "^6.1.4",
- "request": "^2.81.0",
- "request-promise": "^4.2.0",
- "socketio": "^1.0.0"
+ "pg-native": "^1.10.0"
},
"devDependencies": {},
"scripts": {
- "test": "echo \"Error: no test specified\" && exit 1"
+ "test": "echo \"Error: no test specified\" && exit 1",
+ "start": "node server.js"
},
- "author": "",
- "license": ""
+ "author": "schneefux",
+ "license": "UNLICENSED"
}
diff --git a/server.js b/server.js
new file mode 100644
index 0000000..1f5b126
--- /dev/null
+++ b/server.js
@@ -0,0 +1,212 @@
+#!/usr/bin/env node
+/* jshint esnext: true */
+
+var amqp = require("amqplib"),
+ request = require("request-promise"),
+ app = require("express")(),
+ http = require("http").Server(app),
+ sleep = require("sleep-promise");
+
+var MADGLORY_TOKEN = process.env.MADGLORY_TOKEN,
+ RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
+ REGIONS = ["na", "eu", "sg", "sa", "ea"];
+if (MADGLORY_TOKEN == undefined) throw "Need an API token";
+
+(async () => {
+ var rabbit = await amqp.connect(RABBITMQ_URI),
+ ch = await rabbit.createChannel();
+
+ http.listen(8880);
+
+
+ /* API helper */
+ /* search for player name across all shards */
+ async function api_playerByAttr(attr, val) {
+ let finds = [],
+ filter = "filter[" + attr + "]";
+
+ for (let region of REGIONS) {
+ let options = {
+ uri: "https://api.dc01.gamelockerapp.com/shards/" + region + "/players",
+ headers: {
+ "X-Title-Id": "semc-vainglory",
+ "Authorization": MADGLORY_TOKEN
+ },
+ qs: {},
+ json: true,
+ gzip: true
+ };
+ options.qs[filter] = val;
+ retry = true;
+ while (retry) {
+ try {
+ res = await request(options);
+ finds.push({
+ "region": res.data[0].attributes.shardId,
+ "id": res.data[0].id,
+ "name": res.data[0].attributes.name,
+ "last_update": res.data[0].attributes.createdAt,
+ "source": "api"
+ });
+ retry = false;
+ } catch (err) {
+ if (err.statusCode == 429) {
+ await sleep(100);
+ retry = true;
+ } else if (err.statusCode == 404) {
+ retry = false;
+ } else {
+ console.error(err);
+ retry = false;
+ }
+ // TODO
+ }
+ }
+ }
+
+ if (finds.length == 0)
+ return undefined;
+
+ // due to an API bug, many players are also present in NA
+ // TODO: get history for all regions in case of region transfer
+ finds.sort((a, b) => { return a.last_update < b.last_update; });
+ return finds[0];
+ }
+ async function api_playerByName(name) {
+ return await api_playerByAttr("playerNames", name);
+ }
+ async function api_playerById(id) {
+ return await api_playerByAttr("playerIds", id);
+ }
+
+
+ /* DB helper */
+ /* find player by id or by name in db */
+ async function db_playerByAttr(attr, val) {
+ var web = await pool_web.connect(),
+ player = await web.query(`
+ SELECT name, api_id, shard_id, last_match_created_date
+ FROM player WHERE ` + attr + `=$1
+ `, [val]);
+ web.release();
+
+ if (player.rows.length > 0) {
+ return {
+ "name": player.rows[0].name,
+ "id": player.rows[0].api_id,
+ "region": player.rows[0].shard_id,
+ "last_update": player.rows[0].last_match_created_date,
+ "source": "db"
+ };
+ }
+ return undefined;
+ }
+
+ async function db_playerByName(name) {
+ return await db_playerByAttr("name", name);
+ }
+ async function db_playerById(id) {
+ return await db_playerByAttr("api_id", id);
+ }
+
+ /* returns a player by name from db or API */
+ async function playerByName(name) {
+ /*
+ var player = await db_playerByName(name);
+ if (player != undefined) return player;
+ */
+
+ console.log("player '" + name + "' not found in db");
+ player = await api_playerByName(name);
+ if (player != undefined) return player;
+
+ console.log("player '" + name + "' not found in API");
+ return undefined;
+ }
+ /* returns a player by id from db or API */
+ async function playerById(id) {
+ /*
+ var player = await db_playerById(id);
+ if (player != undefined) return player;
+ */
+
+ console.log("player with id '" + name + "' not found in db");
+ player = await api_playerById(id);
+ if (player != undefined) return player;
+
+ console.log("player with id '" + name + "' not found in API");
+ return undefined;
+ }
+ /* returns true if a player has pending jobs */
+ async function anyJobsRunningFor(con, name, id) {
+ return false; // TODO
+ }
+
+ async function playerRequestUpdate(name, id) {
+ if (name == undefined && id == undefined) // fail hard and die
+ throw "playerRequestUpdate needs either name or id";
+
+ let player;
+ if (id != undefined) // prefer id over name
+ player = await playerById(id);
+ else
+ player = await playerByName(name);
+ if (player == undefined)
+ return undefined; // 404
+
+ console.log("updating player '%j'", player);
+
+ // if last_update is from our db, use it as a start, else get the whole history
+ if (player.source == "api" || player.last_update == undefined)
+ player.last_update = new Date(value=0); // forever ago
+
+ // add 1s, because createdAt-start <= x <= createdAt-end
+ // so without the +1s, we'd always get the last_match_created_date match back
+ player.last_update.setSeconds(player.last_update.getSeconds() + 1);
+
+ var timedelta_minutes = ((new Date()) - player.last_update) / 1000 / 60;
+ if (timedelta_minutes < 30) {
+ console.log("player '" + player.name + "' update skipped");
+ return player;
+ }
+
+ var payload = {
+ "region": player.region,
+ "params": {
+ "filter[playerIds]": player.id,
+ "filter[createdAt-start]": player.last_update.toISOString(),
+ "filter[gameMode]": "casual,ranked"
+ }
+ };
+ await ch.sendToQueue("grab", new Buffer(JSON.stringify(payload)), { persistent: true });
+
+ return player;
+ }
+ async function playerRequestUpdateByName(name) {
+ return await playerRequestUpdate(name, undefined);
+ }
+ async function playerRequestUpdateById(id) {
+ return await playerRequestUpdate(undefined, id);
+ }
+
+ /* routes */
+ /* request a grab job */
+ app.get("/api/player/name/:name", async (req, res) => {
+ player = await playerRequestUpdateByName(req.params.name);
+ if (player == undefined) res.sendStatus(404);
+ else res.json(player);
+ });
+ app.get("/api/player/id/:id", async (req, res) => {
+ player = await playerRequestUpdateById(req.params.id);
+ if (player == undefined) res.sendStatus(404);
+ else res.json(player);
+ });
+
+ /* internal monitoring */
+ app.get("/", async (req, res) => {
+ res.sendFile(__dirname + "/index.html");
+ });
+
+
+ //io.emit(name, "done");
+})();