|
| 1 | +#!/usr/bin/env python3.6 |
| 2 | +import os |
| 3 | +import asyncio |
| 4 | +from time import time |
| 5 | + |
| 6 | +import chevron |
| 7 | +import uvloop |
| 8 | +from aiohttp import web, ClientError, ClientSession |
| 9 | +from aiohttp_session import SimpleCookieStorage, get_session |
| 10 | +from aiohttp_session import setup as session_setup |
| 11 | +from arq import Actor, BaseWorker, RedisSettings, concurrent |
| 12 | + |
| 13 | +R_OUTPUT = 'output' |
| 14 | + |
| 15 | +asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) |
| 16 | + |
| 17 | + |
| 18 | +class Downloader(Actor): |
| 19 | + re_enqueue_jobs = True |
| 20 | + |
| 21 | + async def startup(self): |
| 22 | + self.session = ClientSession(loop=self.loop) |
| 23 | + |
| 24 | + @concurrent |
| 25 | + async def download_content(self, url, count): |
| 26 | + total_size = 0 |
| 27 | + errors = [] |
| 28 | + start = time() |
| 29 | + for _ in range(count): |
| 30 | + try: |
| 31 | + async with self.session.get(url) as r: |
| 32 | + content = await r.read() |
| 33 | + total_size += len(content) |
| 34 | + if r.status != 200: |
| 35 | + errors.append(f'{r.status} length: {len(content)}') |
| 36 | + except ClientError as e: |
| 37 | + errors.append(f'{e.__class__.__name__}: {e}') |
| 38 | + output = f'{time() - start:0.2f}s, {count} downloads, total size: {total_size}' |
| 39 | + if errors: |
| 40 | + output += ', errors: ' + ', '.join(errors) |
| 41 | + async with self.redis_pool.get() as redis: |
| 42 | + await redis.rpush(R_OUTPUT, output.encode()) |
| 43 | + return total_size |
| 44 | + |
| 45 | + async def shutdown(self): |
| 46 | + self.session.close() |
| 47 | + |
| 48 | + |
| 49 | +html_template = """ |
| 50 | +<h1>arq demo</h1> |
| 51 | +
|
| 52 | +{{#message}} |
| 53 | +<div>{{ message }}</div> |
| 54 | +{{/message}} |
| 55 | +
|
| 56 | +<form method="post" action="/start-job/"> |
| 57 | + <p> |
| 58 | + <label for="url">Url to download</label> |
| 59 | + <input type="url" name="url" id="url" value="https://httpbin.org/get" required/> |
| 60 | + </p> |
| 61 | + <p> |
| 62 | + <label for="count">Download count</label> |
| 63 | + <input type="number" step="1" name="count" id="count" value="10" required/> |
| 64 | + </p> |
| 65 | + <p> |
| 66 | + <input type="submit" value="Download"/> |
| 67 | + </p> |
| 68 | +</form> |
| 69 | +
|
| 70 | +<h2>Results:</h2> |
| 71 | +{{#results}} |
| 72 | +<p>{{ . }}</p> |
| 73 | +{{/results}} |
| 74 | +""" |
| 75 | + |
| 76 | + |
| 77 | +async def index(request): |
| 78 | + async with await request.app['downloader'].get_redis_conn() as redis: |
| 79 | + data = await redis.lrange(R_OUTPUT, 0, -1) |
| 80 | + results = [r.decode() for r in data] |
| 81 | + |
| 82 | + session = await get_session(request) |
| 83 | + html = chevron.render(html_template, {'message': session.get('message'), 'results': results}) |
| 84 | + session.invalidate() |
| 85 | + return web.Response(text=html, content_type='text/html') |
| 86 | + |
| 87 | + |
| 88 | +async def start_job(request): |
| 89 | + data = await request.post() |
| 90 | + session = await get_session(request) |
| 91 | + try: |
| 92 | + url = data['url'] |
| 93 | + count = int(data['count']) |
| 94 | + except (KeyError, ValueError) as e: |
| 95 | + session['message'] = f'Invalid input, {e.__class__.__name__}: {e}' |
| 96 | + else: |
| 97 | + await request.app['downloader'].download_content(url, count) |
| 98 | + session['message'] = f'Downloading "{url}" ' + (f'{count} times.' if count > 1 else 'once.') |
| 99 | + raise web.HTTPFound(location='/') |
| 100 | + |
| 101 | + |
| 102 | +redis_settings = RedisSettings(host=os.getenv('REDIS_HOST', 'localhost')) |
| 103 | + |
| 104 | + |
| 105 | +async def shutdown(app): |
| 106 | + await app['downloader'].close() |
| 107 | + |
| 108 | + |
| 109 | +def create_app(): |
| 110 | + app = web.Application() |
| 111 | + app.router.add_get('/', index) |
| 112 | + app.router.add_post('/start-job/', start_job) |
| 113 | + app['downloader'] = Downloader(redis_settings=redis_settings) |
| 114 | + app.on_shutdown.append(shutdown) |
| 115 | + session_setup(app, SimpleCookieStorage()) |
| 116 | + return app |
| 117 | + |
| 118 | + |
| 119 | +class Worker(BaseWorker): |
| 120 | + # used by `arq app.py` command |
| 121 | + shadows = [Downloader] |
| 122 | + # set to small value so we can play with timeouts |
| 123 | + timeout_seconds = 10 |
| 124 | + |
| 125 | + def __init__(self, *args, **kwargs): |
| 126 | + kwargs['redis_settings'] = redis_settings |
| 127 | + super().__init__(*args, **kwargs) |
| 128 | + |
| 129 | + |
| 130 | +if __name__ == '__main__': |
| 131 | + # when called directly run the webserver |
| 132 | + app = create_app() |
| 133 | + web.run_app(app, port=8000) |
0 commit comments