diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-05-12 15:18:27 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-05-12 15:18:27 +0200 |
| commit | 11df42bb77a22fa698a910bdb0173a98bbc1ad6a (patch) | |
| tree | 28e34810822201f7a2b69222e4178c95f7d26475 | |
| parent | 11793ebeedd4d6ed3b5ff3a98b136571e50053e1 (diff) | |
| download | apigrabber-11df42bb77a22fa698a910bdb0173a98bbc1ad6a.tar.gz apigrabber-11df42bb77a22fa698a910bdb0173a98bbc1ad6a.zip | |
move api logic to orm
| -rw-r--r-- | jsonapi.js | 64 | ||||
| -rw-r--r-- | worker.js | 76 |
2 files changed, 16 insertions, 124 deletions
diff --git a/jsonapi.js b/jsonapi.js deleted file mode 100644 index fa3b6ef..0000000 --- a/jsonapi.js +++ /dev/null @@ -1,64 +0,0 @@ -#!/usr/bin/node -/* https://github.com/alex94puchades/superagent-jsonapify/blob/master/common.js */ - -"use strict"; - -var _isUndefined = require('lodash/isUndefined'); -var _isArray = require('lodash/isArray'); -var _map = require('lodash/map'); -var _partial = require('lodash/partial'); -var _each = require('lodash/each'); -var _camelCase = require('lodash/camelCase'); -var _memoize = require('lodash/memoize'); -var _chain = require('lodash/chain'); -var _find = require('lodash/find'); -var _clone = require('lodash/clone'); -var _compact = require('lodash/compact'); - -exports.parse = (obj) => { - let response = obj, - data = obj.data; - if (!data) { - return data; - } else if (_isArray(data)) { - return _map(data, _partial(parseResourceDataObject, response)); - } else { - return parseResourceDataObject(response, data); - } -} - -function parseResourceDataObject(response, data) { - var result = _clone(data); - _each(data.attributes, function(value, name) { - Object.defineProperty(result, _camelCase(name), { value: value, enumerable: true }); - }); - _each(data.relationships, function(value, name) { - if (_isArray(value.data)) { - Object.defineProperty(result, _camelCase(name), { - get: _memoize(function() { - const related = _map(value.data, function(related) { - var resdata = _find(response.included, function(included) { - return included.id === related.id && included.type === related.type; - }); - if(resdata) - return parseResourceDataObject(response, resdata); - }); - return _compact(related); - }), - enumerable: true - }); - } else if (value.data) { - Object.defineProperty(result, _camelCase(name), { - get: _memoize(function() { - var resdata = _find(response.included, function(included) { - return included.id === value.data.id && included.type === value.data.type; - }); - return resdata ? parseResourceDataObject(response, resdata) - : null; - }), - enumerable: true - }); - } - }); - return result; -} @@ -8,7 +8,7 @@ const amqp = require("amqplib"), loggly = require("winston-loggly-bulk"), request = require("request-promise"), sleep = require("sleep-promise"), - jsonapi = require("../orm/jsonapi"); + api = require("../orm/api"); const MADGLORY_TOKEN = process.env.MADGLORY_TOKEN, QUEUE = process.env.QUEUE || "grab", @@ -67,70 +67,26 @@ if (LOGGLY_TOKEN) ch.ack(msg); }, { noAck: false }); - // loop over API data pages, notify of progress + // loop over API data objects async function getAPI(payload, where, notify="global") { - let exhausted = false; - payload.params["page[limit]"] = payload.params["page[limit]"] || 50; - payload.params["page[offset]"] = payload.params["page[offset]"] || 0; - - while (!exhausted) { - let opts = { - uri: "https://api.dc01.gamelockerapp.com/shards/" + payload.region + "/" + where, - headers: { - "X-Title-ID": "semc-vainglory", - "Authorization": MADGLORY_TOKEN - }, - json: true, - gzip: true, - time: true, - forever: true, - strictSSL: true, - resolveWithFullResponse: true - }, failed = false, response; - opts.qs = payload.params; - try { - logger.info("API request", { uri: opts.uri, qs: opts.qs }); - response = await request(opts); - const data = jsonapi.parse(response.body); - if (where == "matches") - // send match structure - await Promise.map(data, (match) => - ch.sendToQueue(PROCESS_QUEUE, new Buffer(JSON.stringify(match)), - { persistent: true, type: "match" })); + await Promise.each(await api.requests(where, + payload.region, payload.params, logger), + async (data, idx, len) => { + if (where == "matches") // send match structure + await ch.sendToQueue(PROCESS_QUEUE, + new Buffer(JSON.stringify(data)), { + persistent: true, type: "match", + headers: idx == len - 1? { notify: notify } : {} + }); + // forward "notify" for the last match on the last page if (where == "samples") // forward to sampler - await Promise.map(data, - (sample) => ch.sendToQueue(SAMPLE_QUEUE, - new Buffer(JSON.stringify(sample.attributes.URL)), - { persistent: true, type: "sample" })); - // tell web about progress - await ch.publish("amq.topic", notify, new Buffer("grab_success")); - } catch (err) { - response = err.response; - if (err.statusCode == 429) { - await sleep(1000); - } else if (err.statusCode == 404) { - exhausted = true; - } else { - try { - logger.error("API error", - { uri: err.options.uri, qs: err.options.qs, error: err.response.body }); - } catch (whatever) { - logger.error("weird API error", err); - } - exhausted = true; - } - failed = true; - } finally { - if (response) - logger.info("API response", - { status: response.statusCode, connection_start: response.timings.connect, connection_end: response.timings.end, ratelimit_remaining: parseInt(response.headers["x-ratelimit-remaining"]) }); + await ch.sendToQueue(SAMPLE_QUEUE, + new Buffer(JSON.stringify(data.attributes.URL)), + { persistent: true, type: "sample" }); } + ); - // next page - if (!failed) - payload.params["page[offset]"] += payload.params["page[limit]"] - } await ch.publish("amq.topic", notify, new Buffer("grab_done")); } |
