summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-05-12 15:18:27 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-05-12 15:18:27 +0200
commit11df42bb77a22fa698a910bdb0173a98bbc1ad6a (patch)
tree28e34810822201f7a2b69222e4178c95f7d26475
parent11793ebeedd4d6ed3b5ff3a98b136571e50053e1 (diff)
downloadapigrabber-11df42bb77a22fa698a910bdb0173a98bbc1ad6a.tar.gz
apigrabber-11df42bb77a22fa698a910bdb0173a98bbc1ad6a.zip
move api logic to orm
-rw-r--r--jsonapi.js64
-rw-r--r--worker.js76
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;
-}
diff --git a/worker.js b/worker.js
index 6955532..937925a 100644
--- a/worker.js
+++ b/worker.js
@@ -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"));
}