summaryrefslogtreecommitdiff
path: root/worker.py
diff options
context:
space:
mode:
Diffstat (limited to 'worker.py')
-rw-r--r--worker.py86
1 files changed, 0 insertions, 86 deletions
diff --git a/worker.py b/worker.py
deleted file mode 100644
index 9a00f06..0000000
--- a/worker.py
+++ /dev/null
@@ -1,86 +0,0 @@
-#!/usr/bin/python3
-
-import os
-import logging
-import psycopg2
-import psycopg2.extras
-import psycopg2.extensions
-
-import joblib.joblib
-
-
-RABBIT = {
- "host": os.environ.get("RABBITMQ_HOST"),
- "port": os.environ.get("RABBITMQ_PORT"),
- "credentials": os.environ.get("RABBITMQ_CREDS")
-}
-
-DB = {
- "host": os.environ.get("POSTGRESQL_HOST") or "vaindock_postgres_web",
- "port": os.environ.get("POSTGRESQL_PORT") or 5432,
- "user": os.environ.get("POSTGRESQL_USER") or "vainweb",
- "password": os.environ.get("POSTGRESQL_PASSWORD") or "vainweb",
- "dbname": os.environ.get("POSTGRESQL_DB") or "vainsocial-web"
-}
-
-
-class Processor(joblib.joblib.Worker):
- def __init__(self):
- super().__init__(jobtype="process")
-
- def connect(self, rabbit, db):
- """Connect to database."""
- logging.warning("connecting to database")
- super().connect(**rabbit)
- self._con = psycopg2.connect(**db)
-
- def commit(self, failed):
- self._con.commit()
-
- def work(self, payload):
- """Finish a job."""
- obj_id = payload["id"]
- obj_type = payload["type"]
- o = payload["data"]["attributes"]
- r = payload["data"].get("relationships")
-
- c = self._con.cursor(
- cursor_factory=psycopg2.extras.DictCursor)
-
- d = {}
- if obj_type == "match":
- d["api_id"] = obj_id
- d["created_at"] = o["createdAt"]
- d["duration"] = o["duration"]
- d["game_mode"] = o["gameMode"]
- d["patch_version"] = o.get("patchVersion") or "2.2"
- d["shard_id"] = o["shardId"]
- d["end_game_reason"] = o["stats"]["endGameReason"]
- d["queue"] = o["stats"]["queue"]
- d["roster_1"] = r["rosters"]["data"][0]["id"]
- d["roster_2"] = r["rosters"]["data"][1]["id"]
-
- if d == {}:
- raise Exception
- # TODO
-
- # TODO player upsert
- # TODO request compile & preload
- l = [(c, v) for c, v in d.items()]
- columns = ",".join([t[0] for t in l])
- values = tuple([t[1] for t in l])
- c.execute("INSERT INTO " + obj_type + "(%s) VALUES %s",
- ([psycopg2.extensions.AsIs(columns)] + [values]))
- c.close()
-
-
-def startup():
- worker = Processor()
- worker.connect(RABBIT, DB)
- worker.setup()
- worker.run()
-
-logging.basicConfig(level=logging.DEBUG)
-
-if __name__ == "__main__":
- startup()