diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-02-27 18:40:33 +0100 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-02-27 18:40:33 +0100 |
| commit | 6ac4d07004c94cf73134580e7d104e514b17bf82 (patch) | |
| tree | 7b21bce58b1a63beaad0e1e4aeb3f1f0e84ce21b /api.py | |
| download | compiler-6ac4d07004c94cf73134580e7d104e514b17bf82.tar.gz compiler-6ac4d07004c94cf73134580e7d104e514b17bf82.zip | |
allow writing web-web queries
Diffstat (limited to 'api.py')
| -rw-r--r-- | api.py | 102 |
1 files changed, 102 insertions, 0 deletions
@@ -0,0 +1,102 @@ +#!/usr/bin/python3 + +import asyncio +import os +import glob +import logging +import json +import asyncpg + +import joblib.joblib + +queue_db = { + "host": os.environ.get("POSTGRESQL_SOURCE_HOST") or "vaindock_postgres_raw", + "port": os.environ.get("POSTGRESQL_SOURCE_PORT") or 5532, + "user": os.environ.get("POSTGRESQL_SOURCE_USER") or "vainraw", + "password": os.environ.get("POSTGRESQL_SOURCE_PASSWORD") or "vainraw", + "database": os.environ.get("POSTGRESQL_SOURCE_DB") or "vainsocial-raw" +} + +db_config = { + "host": os.environ.get("POSTGRESQL_DEST_HOST") or "vaindock_postgres_web", + "port": os.environ.get("POSTGRESQL_DEST_PORT") or 5432, + "user": os.environ.get("POSTGRESQL_DEST_USER") or "vainweb", + "password": os.environ.get("POSTGRESQL_DEST_PASSWORD") or "vainweb", + "database": os.environ.get("POSTGRESQL_DEST_DB") or "vainsocial-web" +} + + +class Worker(object): + def __init__(self): + self._queue = None + self._pool = None + self._queries = {} + + async def connect(self, dbconf, queuedb): + """Connect to database.""" + logging.info("connecting to database") + self._queue = joblib.joblib.JobQueue() + await self._queue.connect(**queuedb) + await self._queue.setup() + self._pool = await asyncpg.create_pool(**dbconf) + + async def setup(self): + """Initialize the database.""" + scriptroot = os.path.realpath( + os.path.join(os.getcwd(), os.path.dirname(__file__))) + for path in glob.glob(scriptroot + "/queries/*.sql"): + # utf-8-sig is used by pgadmin, doesn't hurt to specify + # file names: web target table + table = os.path.splitext(os.path.basename(path))[0] + with open(path, "r", encoding="utf-8-sig") as file: + self._queries[table] = file.read() + logging.info("loaded query '%s'", table) + + async def _execute_job(self, jobid, payload): + """Finish a job.""" + object_id = payload["id"] + table = payload["type"] + if table not in self._queries: + logging.debug("%s: nothing to do", jobid) + return + async with self._pool.acquire() as con: + logging.debug("%s: compiling '%s' from '%s'", + jobid, object_id, table) + async with con.transaction(): + await con.execute(self._queries[table], + object_id) + + async def _work(self): + """Fetch a job and run it.""" + jobid, payload, _ = await self._queue.acquire(jobtype="compile") + if jobid is None: + raise LookupError("no jobs available") + logging.debug("%s: starting job", jobid) + await self._execute_job(jobid, payload) + await self._queue.finish(jobid) + logging.debug("%s: finished job", jobid) + + async def run(self): + """Start jobs forever.""" + while True: + try: + await self._work() + except LookupError: + logging.info("nothing to do, idling") + await asyncio.sleep(10) + + async def start(self, number=1): + """Start jobs in background.""" + for _ in range(number): + asyncio.ensure_future(self.run()) + +async def startup(): + worker = Worker() + await worker.connect(db_config, queue_db) + await worker.setup() + await worker.start(10) + +logging.basicConfig(level=logging.DEBUG) +loop = asyncio.get_event_loop() +loop.run_until_complete(startup()) +loop.run_forever() |
