This repository has been archived by the owner on Mar 6, 2024. It is now read-only.
/
worker.py
73 lines (58 loc) · 1.97 KB
/
worker.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
"""Proxy for a remote worker process."""
import asyncio
import logging
from collections import deque
_logger = logging.getLogger('distex.Worker')
class Worker(asyncio.Protocol):
"""
Worker that submits tasks to and gets results from a
local or remote processor.
"""
def __init__(self, serializer):
self.serializer = serializer
self.loop = asyncio.get_event_loop_policy().get_event_loop()
self.transport = None
self.peername = ''
self.disconnected = None
self.futures = deque()
self.tasks = deque()
def __repr__(self):
return f'<Worker {self.peername}>'
def run_task(self, task):
"""Send the task to the processor and return Future for the result."""
self.serializer.write_request(self.transport.write, task)
future = self.loop.create_future()
self.futures.append(future)
self.tasks.append(task)
return future
def stop(self):
"""Close connection to the processor."""
if self.transport:
self.transport.close()
self.transport = None
# protocol callbacks:
def connection_made(self, transport):
self.transport = transport
hp = transport.get_extra_info('peername')
if hp:
host, port = hp
self.peername = f'{host}:{port}'
else:
self.peername = 'Unix socket'
_logger.info(f'Connection from {self.peername}')
def connection_lost(self, exc):
if exc:
self.disconnected(self)
_logger.error(f'Connection lost from {self.peername}: {exc}')
self.transport = None
def data_received(self, data):
self.serializer.add_data(data)
for resp in self.serializer.get_responses():
self.futures.popleft().set_result(resp)
self.tasks.popleft()
def eof_received(self):
pass
def pause_writing(self):
pass
def resume_writing(self):
pass