Repository navigation
Expand file tree
/
Copy pathformat.py
More file actions
148 lines (107 loc) · 3.86 KB
/
Copy pathformat.py
File metadata and controls
148 lines (107 loc) · 3.86 KB
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
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
"""
Message formats for IPC.
"""
from __future__ import absolute_import
from __future__ import division
from __future__ import unicode_literals
import pickle
import pickletools
import json
class PickleMessageParser(object):
"""
Message parser for the pickle stream format.
Args:
buffer_max_len (int): The maximum number of bytes buffered while
parsing a stream of incoming messages. Defaults to 32768.
"""
MAX_LENGTH = 32768
HEADER_MAX_LEN = 24
def __init__(self, buffer_max_len=MAX_LENGTH):
self._buffer_max_len = buffer_max_len
self._buffer = b''
def push(self, data):
"""
Push data onto the message parser buffer.
Args:
data (bytes): Data as received from the network. Partial messages
are allowed.
Raises:
RuntimeError: If the buffer is full.
"""
if len(self._buffer) + len(data) > self._buffer_max_len:
raise RuntimeError('Buffer length exceeded')
self._buffer += data
def messages(self):
"""
Iterate over all available messages.
Yields:
object: The next decoded message.
"""
frame_start = 0
while self._buffer.find(b'.', frame_start, frame_start + self.HEADER_MAX_LEN) > -1:
header = self._buffer[frame_start:frame_start + self.HEADER_MAX_LEN]
header_len = list(pickletools.genops(header))[-1][2]+1
doc_len = pickle.loads(header[:header_len])
if not isinstance(doc_len, int) or doc_len < 0:
raise ValueError('Document length must be a positive integer')
doc_start = frame_start + header_len
doc_end = doc_start + doc_len
if doc_end > len(self._buffer):
break
yield pickle.loads(self._buffer[doc_start:doc_end])
frame_start += header_len + doc_len
self._buffer = self._buffer[frame_start:]
class PickleMessageBuilder(object):
"""
Message builder for the pickle stream format.
"""
def __init__(self, protocol=2):
self.protocol = protocol
def message(self, msg):
data = pickle.dumps(msg, protocol=self.protocol)
header = pickle.dumps(len(data), protocol=self.protocol)
return header + data
class JsonMessageParser(object):
"""
Message parser for the JSON lines stream format.
Args:
buffer_max_len (int): The maximum number of bytes buffered while
parsing a stream of incoming messages. Defaults to 32768.
"""
MAX_LENGTH = 32768
def __init__(self, buffer_max_len=MAX_LENGTH):
self._buffer_max_len = buffer_max_len
self._buffer = b''
def push(self, data):
"""
Push data onto the message parser buffer.
Args:
data (bytes): Data as received from the network. Partial messages
are allowed.
Raises:
RuntimeError: If the buffer is full.
"""
if len(self._buffer) + len(data) > self._buffer_max_len:
raise RuntimeError('Buffer length exceeded')
self._buffer += data
def messages(self):
"""
Iterate over all available messages.
Yields:
object: The next decoded message.
"""
frame_start = 0
while self._buffer.find(b'\n') > -1:
frame_len = self._buffer.index(b'\n') + 1
frame_end = frame_start + frame_len
if frame_end > len(self._buffer):
break
yield json.loads(self._buffer[frame_start:frame_end].decode('utf-8'))
frame_start += frame_len
self._buffer = self._buffer[frame_start:]
class JsonMessageBuilder(object):
"""
Message builder for the JSON lines stream format.
"""
def message(self, msg):
return json.dumps(msg).encode('utf-8') + b'\n'