diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-24 16:48:08 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-24 16:48:08 +0100 |
| commit | 3e96a0355ef2bda1928d7cd34250b62c4a9af49e (patch) | |
| tree | 556cd19c08542e088514af1fd7a1335a2b247a7e | |
| parent | 90c0cd497930b37cbb1b0e2c3be34759a5bd6549 (diff) | |
| download | bridge-3e96a0355ef2bda1928d7cd34250b62c4a9af49e.tar.gz bridge-3e96a0355ef2bda1928d7cd34250b62c4a9af49e.zip | |
refactor - send to player specific channels
| -rw-r--r-- | api.js | 424 | ||||
| -rw-r--r-- | index.html | 60 |
2 files changed, 239 insertions, 245 deletions
@@ -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 } @@ -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> |
