diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-28 18:29:39 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-03-28 18:29:39 +0200 |
| commit | 151d1ff466c1d2904eafe6f0d383f89978114cb0 (patch) | |
| tree | fdb2a0420f6be6f8cd16373fa591ed9f9a7e48c3 | |
| parent | 3c7b19d77476d774762b4a6c97bff18bcd0b249a (diff) | |
| download | apigrabber-151d1ff466c1d2904eafe6f0d383f89978114cb0.tar.gz apigrabber-151d1ff466c1d2904eafe6f0d383f89978114cb0.zip | |
rewrite in NodeJS
| -rw-r--r-- | .gitignore | 116 | ||||
| -rw-r--r-- | .gitmodules | 3 | ||||
| -rw-r--r-- | Dockerfile | 14 | ||||
| -rw-r--r-- | cli.py | 40 | ||||
| -rw-r--r-- | crawler.py | 90 | ||||
| m--------- | joblib | 0 | ||||
| -rw-r--r-- | package.json | 11 | ||||
| -rw-r--r-- | requirements.txt | 6 | ||||
| -rw-r--r-- | worker.js | 62 | ||||
| -rw-r--r-- | worker.py | 55 |
10 files changed, 124 insertions, 273 deletions
@@ -1,91 +1,59 @@ -# Byte-compiled / optimized / DLL files -__pycache__/ -*.py[cod] -*$py.class +# Logs +logs +*.log +npm-debug.log* +yarn-debug.log* +yarn-error.log* -# C extensions -*.so +# Runtime data +pids +*.pid +*.seed +*.pid.lock -# Distribution / packaging -.Python -env/ -build/ -develop-eggs/ -dist/ -downloads/ -eggs/ -.eggs/ -lib/ -lib64/ -parts/ -sdist/ -var/ -wheels/ -*.egg-info/ -.installed.cfg -*.egg +# Directory for instrumented libs generated by jscoverage/JSCover +lib-cov -# PyInstaller -# Usually these files are written by a python script from a template -# before PyInstaller builds the exe, so as to inject date/other infos into it. -*.manifest -*.spec +# Coverage directory used by tools like istanbul +coverage -# Installer logs -pip-log.txt -pip-delete-this-directory.txt +# nyc test coverage +.nyc_output -# Unit test / coverage reports -htmlcov/ -.tox/ -.coverage -.coverage.* -.cache -nosetests.xml -coverage.xml -*,cover -.hypothesis/ +# Grunt intermediate storage (http://gruntjs.com/creating-plugins#storing-task-files) +.grunt -# Translations -*.mo -*.pot +# Bower dependency directory (https://bower.io/) +bower_components -# Django stuff: -*.log -local_settings.py +# node-waf configuration +.lock-wscript -# Flask stuff: -instance/ -.webassets-cache +# Compiled binary addons (http://nodejs.org/api/addons.html) +build/Release -# Scrapy stuff: -.scrapy +# Dependency directories +node_modules/ +jspm_packages/ -# Sphinx documentation -docs/_build/ +# Typescript v1 declaration files +typings/ -# PyBuilder -target/ +# Optional npm cache directory +.npm -# Jupyter Notebook -.ipynb_checkpoints +# Optional eslint cache +.eslintcache -# pyenv -.python-version +# Optional REPL history +.node_repl_history -# celery beat schedule file -celerybeat-schedule +# Output of 'npm pack' +*.tgz -# dotenv -.env +# Yarn Integrity file +.yarn-integrity -# virtualenv -.venv/ -venv/ -ENV/ - -# Spyder project settings -.spyderproject +# dotenv environment variables file +.env -# Rope project settings -.ropeproject diff --git a/.gitmodules b/.gitmodules index 52948d0..e69de29 100644 --- a/.gitmodules +++ b/.gitmodules @@ -1,3 +0,0 @@ -[submodule "joblib"] - path = joblib - url = https://github.com/vainglorygame/joblib @@ -1,5 +1,9 @@ -FROM python:3.6-alpine -ADD requirements.txt /code/requirements.txt -WORKDIR /code -RUN pip install -r requirements.txt -CMD ["python", "worker.py"] +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", "worker.js"] @@ -1,40 +0,0 @@ -#!/usr/bin/python3 -import os -import argparse - -import joblib.joblib - - -RABBIT = { - "host": os.environ.get("RABBITMQ_HOST"), - "port": os.environ.get("RABBITMQ_PORT"), - "credentials": os.environ.get("RABBITMQ_CREDS") -} - - -def main(args): - queue = joblib.joblib.JobQueue() - queue.connect(**RABBIT) - - if args.player: - payload = { - "region": args.region, - "params": { - "filter[createdAt-start]": "2017-02-01T00:00:00Z", - "filter[playerNames]": args.player - } - } - queue.request("grab", payload) - -parser = argparse.ArgumentParser(description="Request a Vainsocial update.") -parser.add_argument("-n", "--player", - help="Player name", - type=str) -parser.add_argument("-r", "--region", - help="Specify a region", - type=str, - choices=["na", "eu", "sg"], - default="na") - -if __name__ == "__main__": - main(parser.parse_args()) diff --git a/crawler.py b/crawler.py deleted file mode 100644 index 84e4692..0000000 --- a/crawler.py +++ /dev/null @@ -1,90 +0,0 @@ -#!/usr/bin/python - -import json -import time -import logging -import requests - -APIURL = "https://api.dc01.gamelockerapp.com/" - - -class ApiError(Exception): - pass - -class Crawler(object): - def __init__(self, token): - """Sets constants.""" - self._apiurl = APIURL - self._token = token - self._pagelimit = 50 - - def _req(self, path, params): - """Sends an API request and returns the response dict. - - :param path: URL path. - :type path: str - :param params: Request parameters. - :type params: dict - :return: API response. - :rtype: dict - """ - headers = { - "Authorization": "Bearer " + self._token, - "X-TITLE-ID": "semc-vainglory", - "Accept": "application/vnd.api+json", - "Accept-Encoding": "gzip" - } - retries = 5 - while True: - try: - response = requests.get(self._apiurl + path, - headers=headers, - params=params) - if response.status_code == 429: - logging.warning("rate limited, retrying") - else: - if response.status_code > 500: - logging.error("API server error %s", - response.status_code) - raise ApiError(response.status_code) - else: - return response.json() - except (json.decoder.JSONDecodeError) as err: - # API bug? - logging.error("API error '%s', retrying", err) - retries -= 1 - if retries == 0: - logging.error("Giving up") - raise ApiError(str(err)) - - time.sleep(5) - - def matches(self, params, region="na"): - """Queries the API for matches and their related data. - - :param region: (optional) Region where the matches were played. - Defaults to "na" (North America). - :type region: str - :param params: Additional filters. - :type params: dict - """ - params["page[offset]"] = 0 - params["page[limit]"] = self._pagelimit - while True: - res = self._req("shards/" + region + "/matches", - params) - - if "errors" in res: - if res["errors"][0].get("title") == "Not Found" \ - and params["page[offset]"] > 0: - # a query returned exactly 50 matches - # which is expected, so don't fail. - return - raise ApiError(res["errors"]) - - yield res - - if len(res["data"]) < self._pagelimit: - # asked for 50, got less -> exhausted - break - params["page[offset]"] += params["page[limit]"] diff --git a/joblib b/joblib deleted file mode 160000 -Subproject 7ced841aa44a7dd7bd2d3a68ceb989ef488ce44 diff --git a/package.json b/package.json new file mode 100644 index 0000000..401ffb8 --- /dev/null +++ b/package.json @@ -0,0 +1,11 @@ +{ + "name": "apigrabber", + "version": "2.0.0", + "description": "", + "main": "worker.js", + "scripts": { + "test": "echo \"Error: no test specified\" && exit 1" + }, + "author": "schneefux", + "license": "UNLICENSED" +} diff --git a/requirements.txt b/requirements.txt deleted file mode 100644 index b01174d..0000000 --- a/requirements.txt +++ /dev/null @@ -1,6 +0,0 @@ -appdirs==1.4.3 -packaging==16.8 -pika==0.10.0 -pyparsing==2.2.0 -requests==2.13.0 -six==1.10.0 diff --git a/worker.js b/worker.js new file mode 100644 index 0000000..8dc1645 --- /dev/null +++ b/worker.js @@ -0,0 +1,62 @@ +#!/usr/bin/node +/* jshint esnext:true */ + +var amqp = require("amqplib"), + request = require("request-promise"), + sleep = require("sleep-promise"); + +var MADGLORY_TOKEN = process.env.MADGLORY_TOKEN, + RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost"; +if (MADGLORY_TOKEN == undefined) throw "Need an API token"; + +(async () => { + var rabbit = await amqp.connect(RABBITMQ_URI), + ch = await rabbit.createChannel(); + + await ch.assertQueue("grab", {durable: true}); + await ch.assertQueue("process", {durable: true}); + await ch.prefetch(1); + + ch.consume("grab", async (msg) => { + let exhausted = false; + payload = JSON.parse(msg.content); + 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 + "/matches", + headers: { + "X-Title-ID": "semc-vainglory", + "Authorization": MADGLORY_TOKEN + }, + json: true, + gzip: true + }; + opts.qs = payload.params; + try { + console.log("API request: %j", opts.qs); + res = await request(opts); + console.log("got a few matches"); + await ch.sendToQueue("process", new Buffer(JSON.stringify(res)), { persistent: true }); + } catch (err) { + if (err.statusCode == 429) { + await sleep(1000); + } else if (err.statusCode == 404) { + // TODO stop early if len(matches) < pagelen + exhausted = true; + } else { + console.error(err); + exhausted = true; + } + console.log(err.statusCode); + } + + // next page + payload.params["page[offset]"] += payload.params["page[limit]"] + } + + console.log("done"); + ch.ack(msg); + }, { noAck: false }); +})(); diff --git a/worker.py b/worker.py deleted file mode 100644 index d4c5fab..0000000 --- a/worker.py +++ /dev/null @@ -1,55 +0,0 @@ -#!/usr/bin/python - -import os -import logging - -import crawler -import joblib.joblib - - -RABBIT = { - "host": os.environ.get("RABBITMQ_HOST"), - "port": os.environ.get("RABBITMQ_PORT"), - "credentials": os.environ.get("RABBITMQ_CREDS") -} - -APITOKEN = os.environ["MADGLORY_TOKEN"] - - -class Apigrabber(joblib.joblib.Worker): - def __init__(self, apitoken): - super().__init__("grab") - self._apitoken = apitoken - - def work(self, payload): - """Finish a job.""" - api = crawler.Crawler(self._apitoken) - logging.info("running on %s with parameters '%s'", - payload["region"], payload["params"]) - try: - for data in api.matches(region=payload["region"], - params=payload["params"]): - items = data["data"] + data["included"] - for item in items: - self.request("process", - payload={ - "id": item["id"], - "type": item["type"], - "data": item - }) - except crawler.ApiError as error: - logging.warning("API returned error '%s'", error.args[0]) - raise joblib.JobFailed(error.args[0], - False) # not critical - - -def startup(): - worker = Apigrabber(APITOKEN) - worker.connect(**RABBIT) - worker.setup() - worker.run() - - -if __name__ == "__main__": - logging.basicConfig(level=logging.INFO) - startup() |
