summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-24 16:48:08 +0100
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-24 16:48:08 +0100
commit3e96a0355ef2bda1928d7cd34250b62c4a9af49e (patch)
tree556cd19c08542e088514af1fd7a1335a2b247a7e
parent90c0cd497930b37cbb1b0e2c3be34759a5bd6549 (diff)
downloadbridge-3e96a0355ef2bda1928d7cd34250b62c4a9af49e.tar.gz
bridge-3e96a0355ef2bda1928d7cd34250b62c4a9af49e.zip
refactor - send to player specific channels
-rw-r--r--api.js424
-rw-r--r--index.html60
2 files changed, 239 insertions, 245 deletions
diff --git a/api.js b/api.js
index 8df8f9d..d6525dc 100644
--- a/api.js
+++ b/api.js
@@ -34,10 +34,11 @@ if (APITOKEN == undefined) throw "Need a valid API token!";
http.listen(8080);
/* API helper */
-/* searches for player name across all shards */
-async function findPlayer(name) {
+/* search for player name across all shards */
+async function api_playerByAttr(attr, val) {
var regions = ["na", "eu", "sg"],
- finds = [];
+ finds = [],
+ filter = "filter[" + attr + "]";
for (let region of regions) {
var options = {
@@ -47,7 +48,7 @@ async function findPlayer(name) {
"Authorization": APITOKEN
},
qs: {
- "filter[playerNames]": name
+ filter: val
},
json: true,
gzip: true
@@ -57,7 +58,8 @@ async function findPlayer(name) {
finds.push({
"region": res.data[0].attributes.shardId,
"id": res.data[0].id,
- "last_update": res.data[0].attributes.createdAt
+ "last_update": res.data[0].attributes.createdAt,
+ "source": "api"
});
} catch (err) {
// TODO
@@ -72,141 +74,191 @@ async function findPlayer(name) {
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);
+}
-/* routes */
-app.get("/api/player/:name", async (req, res) => {
- var name = req.params.name;
- var raw = await pool_raw.connect(),
- web = await pool_web.connect();
+/* 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.code != "40001") throw err;
+ }
+ /* transaction end */
+ } while (!commit_success);
+ return res;
+}
- /* search for the player in our db */
- var player = await web.query(`
- SELECT api_id, shard_id, last_match_created_date
- FROM player WHERE name=$1
- `, [name]);
- var player_id, player_region, grab_start;
+/* 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();
- /* found */
if (player.rows.length > 0) {
- console.log("player '" + name + "' was found in db");
- player_id = player.rows[0].api_id;
- player_region = player.rows[0].shard_id;
- grab_start = player.rows[0].last_match_created_date;
- // TODO db does not save time zone offset @stormcaller remove this
- if (grab_start != undefined)
- grab_start.setMinutes(grab_start.getMinutes() -
- (new Date().getTimezoneOffset()));
+ 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"
+ };
}
- /* not found */
- if (player.rows.length == 0) {
- console.log("player '" + name + "' not found in db");
- /* search in all regions */
- player = await findPlayer(name);
- if (player == undefined) {
- console.log("player '" + name + "' not found in API");
- /* give up */
- res.sendStatus(404);
- return;
- }
+ return undefined;
+}
- player_id = player.id;
- player_region = player.region;
- }
- if (grab_start == undefined) {
- grab_start = new Date("2017-01-01T00:00:00Z");
- }
+async function db_playerByName(name) {
+ return await db_playerByAttr("name", name);
+}
+async function db_playerById(id) {
+ return await db_playerByAttr("api_id", id);
+}
- var timedelta_minutes = ((new Date()) - grab_start) / 1000 / 60;
- // TODO make it configurable
- if (timedelta_minutes < 30) {
- console.log("player '" + name + "' update skipped");
- res.sendStatus(304);
- return;
- }
+/* returns a player by name from db or API */
+async function playerByName(name) {
+ var player = await db_playerByName(name);
+ if (player != undefined) return player;
- /* request update job */
- var jobid = -1;
+ console.log("player '" + name + "' not found in db");
+ player = await api_playerByName(name);
+ if (player != undefined) return player;
- // createdAt-start <= x <= createdAt-end
- grab_start.setSeconds(grab_start.getSeconds() + 1);
+ 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;
- payload = {
- "region": player_region,
- "params": {
- "filter[playerIds]": player_id,
- "filter[playerNames]": name, // TODO remove in 2.0 - backwards compat
- "filter[createdAt-start]": grab_start.toISOString(),
- "filter[gameMode]": "casual,ranked"
- }
- };
+ 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;
+}
+
+
+/* update request helpers */
+/* upsert a job */
+async function upsertGrabjob(payload) {
+ var raw = await pool_raw.connect(),
+ job;
+
+ // find and prioritize existing jobs
+ job = await db_serialized(raw, `
+ UPDATE jobs SET priority=0
+ WHERE
+ (
+ (type='grab' AND payload=$1) OR
+ (type='process' AND payload->>'playername'=$1->'params'->>'filter[playerNames]') OR
+ (type='compile' AND payload->>'type'='player' AND payload->>'id'=$1->'params'->>'filter[playerIds]')
+ ) AND status<>'finished' AND status<>'failed'
+ RETURNING id
+ `, [payload]);
+ // TODO job dependency information format on jobs is shit
+ // (2.0)
- /* transaction begin */
- try {
- await raw.query("BEGIN");
- await raw.query("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE");
- var job = await raw.query(`
- UPDATE jobs SET priority=0
- WHERE
- (
- (type='grab' AND payload=$1) OR
- (type='process' AND payload->>'playername'=$1->'params'->>'filter[playerNames]') OR
- (type='compile' AND payload->>'type'='player' AND payload->>'id'=$1->'params'->>'filter[playerIds]')
- ) AND status<>'finished' AND status<>'failed'
+ if (job.rows.length == 0) {
+ // 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]);
- await raw.query("COMMIT");
- } catch (err) {
- if (err.code == "40001") {
- // serialization error - expected
- res.sendStatus(202); // try again
- return;
- } else {
- throw err;
- }
+ // wake apigrabber up
+ await raw.query(`NOTIFY grab_open`);
+ console.log("job requested: '%j'", payload);
}
- /* transaction end */
- if (job.rows.length == 0) {
- /* transaction begin */
- try {
- await raw.query("BEGIN");
- await raw.query("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE");
- job = await raw.query(`
- INSERT INTO jobs(type, payload, priority)
- VALUES('grab', $1, 0)
- RETURNING id
- `, [payload]);
- await raw.query("COMMIT");
- } catch (err) {
- if (err.code == "40001") {
- // serialization error - expected
- res.sendStatus(202); // try again
- return;
- } else {
- throw err;
- }
- }
- /* transaction end */
+ raw.release();
+ return job.rows[0];
+}
- // wake apigrabber up
- await raw.query(`NOTIFY grab_open`, []);
- console.log("player '" + name + "' new job requested");
- }
+async function playerRequestUpdate(name, id) {
+ if (name == undefined && id == undefined) // fail hard and die
+ throw "playerRequestUpdate needs either name or id";
- console.log("player '" + name + "' updating after " + grab_start.toISOString());
+ var player;
+ if (id != undefined) // prefer id over name
+ player = await playerById(id);
+ else
+ player = await playerByName(name);
+ if (player == undefined)
+ return undefined;
- jobid = job.rows[0].id;
+ console.log("updating player '%j'", player);
- /* clean up */
- raw.release();
- web.release();
+ // 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
+ */
- res.json({
- "job_id": jobid,
- "player_id": player_id,
- "player_region": player_region
- });
+ // 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"
+ }
+ };
+ job = await upsertGrabjob(payload);
+
+ 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 */
@@ -214,117 +266,65 @@ app.get("/", async (req, res) => {
res.sendFile(__dirname + "/index.html");
});
+
/* notifications from database */
async function listen() {
- var client = new pg.Client(db_config_raw),
- last_broadcast_ids = {};
-
- /* save the current state of the job queue */
+ var client = new pg.Client(db_config_raw);
await client.connect();
- last_broadcast_ids.grab = (await client.query(`
- SELECT id FROM jobs
- WHERE type='process'
- ORDER BY id DESC LIMIT 1
- `)).rows[0].id;
- last_broadcast_ids.process = (await client.query(`
- SELECT id FROM jobs
- WHERE type='process'
- ORDER BY id DESC LIMIT 1
- `)).rows[0].id;
- last_broadcast_ids.compile = (await client.query(`
- SELECT id FROM jobs
- WHERE type='compile'
- ORDER BY id DESC LIMIT 1
- `)).rows[0].id;
/* job status change notification listener */
client.on('notification', async (msg) => {
- var raw = await pool_raw.connect();
- // TODO DRY
+ var raw = await pool_raw.connect(),
+ jobs;
+
+ // find all interesting jobs, delete them & forward their notification
if (msg.channel == "grab_failed") {
- // get all jobs between the last time and the notification
- // TODO remove playername in 2.0
- var grab_jobs = await raw.query(`
- SELECT
- MAX(id) AS id,
- payload->'params'->>'filter[playerIds]' AS playerid,
- payload->'params'->>'filter[playerNames]' AS playername
- FROM jobs
- WHERE type='grab' AND status='failed' AND id>$1
- GROUP BY playerid, playername
- `, [last_broadcast_ids.grab]);
- grab_jobs.rows.sort((a, b) => { return a.id < b.id; });
- for (let grab_job of grab_jobs.rows) {
- // send a notification for each player name
- io.emit("player grab failed", {
- "name": grab_job.playername,
- "id": grab_job.playerid
- });
- }
- if (grab_jobs.rows.length > 0) {
- last_broadcast_ids.grab = grab_jobs.rows[0].id;
- }
+ jobs = await raw.query(`
+ DELETE FROM jobs WHERE
+ type='grab' AND status='failed' AND payload->>'error'='Not Found'
+ RETURNING
+ payload->'params'->>'filter[playerIds]' AS player_id,
+ payload->'params'->>'filter[playerNames]' AS player_name
+ `);
}
if (msg.channel == "process_finished") {
- // get all jobs between the last time and the notification
- var process_jobs = await raw.query(`
- SELECT
- MAX(id) AS id, payload->>'playername' AS name FROM jobs
- WHERE type='process' AND status='finished' AND id>$1
- GROUP BY name
- `, [last_broadcast_ids.process]);
- process_jobs.rows.sort((a, b) => { return a.id < b.id; });
- for (let process_job of process_jobs.rows) {
- // send a notification for each player name
- io.emit("player processed", {
- "id": null,
- "name": process_job.name
- });
- }
- if (process_jobs.rows.length > 0) {
- last_broadcast_ids.process = process_jobs.rows[0].id;
- }
+ jobs = await raw.query(`
+ DELETE FROM jobs WHERE
+ type='process' AND status='finished'
+ RETURNING payload->>'playername' AS player_name
+ `);
}
if (msg.channel == "compile_finished") {
- // get all jobs between the last time and the notification
- var compile_jobs = await raw.query(`
- SELECT
- MAX(id) AS id, payload->>'id' AS playerid FROM jobs
- WHERE type='compile' AND payload->>'type'='player'
- AND status='finished' AND id>$1
- GROUP BY playerid
- `, [last_broadcast_ids.compile]);
- compile_jobs.rows.sort((a, b) => { return a.id < b.id; });
- for (let compile_job of compile_jobs.rows) {
- // send a notification for each player name
- io.emit("player compiled", {
- "id": compile_job.playerid,
- "name": null
- });
- }
- if (compile_jobs.rows.length > 0) {
- last_broadcast_ids.compile = compile_jobs.rows[0].id;
- }
+ jobs = await raw.query(`
+ DELETE FROM jobs WHERE
+ type='compile' AND payload->>'type'='player' AND status='finished'
+ RETURNING payload->>'id' AS player_id
+ `);
}
- io.emit("job update", msg.channel);
+
raw.release();
+ if (jobs == undefined) return; // nothing to do
+
+ // forward notification to all playername / playerid channels
+ for (let job of jobs.rows) {
+ var name = job.player_name;
+ var id = job.player_id;
+ var player;
+ if (name == undefined && id == undefined) throw "notification needs either name or ID";
+ if (name == undefined)
+ player = await db_playerById(job.player_id);
+ else
+ player = await db_playerByName(job.player_name);
+ if (player == undefined) throw "player had a job, but doesn't exist";
+ console.log("sending '%s' notification for player '%s'", msg.channel, player.name);
+ io.emit(player.name, msg.channel);
+ io.emit(player.id, msg.channel);
+ }
});
- client.query("LISTEN grab_open");
- client.query("LISTEN process_open");
- client.query("LISTEN compile_open");
- client.query("LISTEN analyze_open");
- client.query("LISTEN grab_running");
- client.query("LISTEN process_running");
- client.query("LISTEN compile_running");
- client.query("LISTEN analyze_running");
- client.query("LISTEN grab_finished");
client.query("LISTEN process_finished");
client.query("LISTEN compile_finished");
client.query("LISTEN analyze_finished");
client.query("LISTEN grab_failed");
- client.query("LISTEN process_failed");
- client.query("LISTEN compile_failed");
- client.query("LISTEN analyze_failed");
// keep open forever
}
diff --git a/index.html b/index.html
index 7ab1345..b97d287 100644
--- a/index.html
+++ b/index.html
@@ -13,13 +13,39 @@
<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();
$("#update-form").submit(function(e) {
- $.get("/api/player/" + $("input:first").val()).done(function() {
+ var name = $("input:first").val();
+ $.get("/api/player/name/" + name).done(function() {
$("#updates").append(
$("<li>")
.text("(found)")
.css("color", "black")
);
+ /* subscribe to socket notifications */
+ socket.on(name, function(msg) {
+ if (msg == "grab_failed") {
+ $("#updates").append(
+ $("<li>")
+ .text(name + ": no new data")
+ .css("color", "red")
+ );
+ }
+ if (msg == "process_finished") {
+ $("#updates").append(
+ $("<li>")
+ .text(name + ": processed")
+ .css("color", "green")
+ );
+ }
+ if (msg == "compile_finished") {
+ $("#updates").append(
+ $("<li>")
+ .text(name + ": compiled")
+ .css("color", "orange")
+ );
+ }
+ });
}).fail(function() {
$("#updates").append(
$("<li>")
@@ -29,38 +55,6 @@
});
e.preventDefault();
});
- var socket = io();
- socket.on("job update", function(msg) {
- $("#updates").append(
- $("<li>")
- .text(msg)
- .css("color", "grey")
- );
- });
- socket.on("player grab failed", function(msg) {
- $("#updates").append(
- $("<li>")
- .text("no new data for player " + msg.id +
- " (" + msg.name + ")")
- .css("color", "red")
- );
- });
- socket.on("player processed", function(msg) {
- $("#updates").append(
- $("<li>")
- .text("processed player " + msg.id +
- " (" + msg.name + ")")
- .css("color", "green")
- );
- });
- socket.on("player compiled", function(msg) {
- $("#updates").append(
- $("<li>")
- .text("compiled player " + msg.id +
- " (" + msg.name + ")")
- .css("color", "orange")
- );
- });
</script>
</body>
</html>