-
Notifications
You must be signed in to change notification settings - Fork 0
/
asyncio_subscriber.py
76 lines (59 loc) · 2.23 KB
/
asyncio_subscriber.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
74
75
import struct
import asyncio
class SubscriberProtocol(asyncio.Protocol):
hdrfmt = struct.Struct('>I')
def __init__(self, loop):
self.loop = loop
def connection_made(self, transport):
print('connection_made')
self.transport = transport
self.reset()
def connection_lost(self, exc):
print('connection_lost')
self.loop.stop()
def data_received(self, data):
self.accum_buffer.extend(data)
while len(self.accum_buffer) >= self.packet_size:
pkt = memoryview(self.accum_buffer)[:self.packet_size]
self.accum_buffer = self.accum_buffer[self.packet_size:]
if self.wait_hdr:
self.wait_hdr = False
self.packet_size = self.hdrfmt.unpack(pkt)[0]
else:
self.wait_hdr = True
self.packet_size = self.hdrfmt.size
self.handle_msg(pkt)
def reset(self):
self.wait_hdr = True
self.packet_size = self.hdrfmt.size
self.accum_buffer = bytearray()
def handle_msg(self, msg):
idx = struct.unpack('I', msg[:4])[0]
print('{0}: {1}'.format(idx, len(msg)))
@asyncio.coroutine
def reconnect(loop):
coro = loop.create_connection(lambda: SubscriberProtocol(loop), '127.0.0.1', 10000)
transport, protocol = yield from coro
def handle_msg(msg):
idx = struct.unpack('I', msg[:4])[0]
print('{0}: {1}'.format(idx, len(msg)))
@asyncio.coroutine
def subscribe_stuff():
reader, writer = yield from asyncio.open_connection('127.0.0.1', 10000)
hdrfmt = struct.Struct('>I')
while True:
try:
hdr = yield from reader.readexactly(hdrfmt.size)
packet_size = hdrfmt.unpack(hdr)[0]
payload = yield from reader.readexactly(packet_size)
except asyncio.IncompleteReadError as e:
print('partial read {} / {}'.format(len(e.partial), e.expected))
break
handle_msg(payload)
loop = asyncio.get_event_loop()
#coro = loop.create_connection(lambda: SubscriberProtocol(loop), '127.0.0.1', 10000)
#transport, protocol = loop.run_until_complete(coro)
asyncio.async(reconnect(loop))
loop.run_forever()
loop.close()
#loop.run_until_complete(subscribe_stuff())