forked from Tribler/dispersy
-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathendpoint.py
376 lines (289 loc) · 13.9 KB
/
endpoint.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
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
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
import errno
import logging
import socket
import sys
import threading
from abc import ABCMeta, abstractmethod
from itertools import product
from select import select
from time import time
from twisted.internet import reactor
from .candidate import Candidate
if sys.platform == 'win32':
SOCKET_BLOCK_ERRORCODE = 10035 # WSAEWOULDBLOCK
else:
SOCKET_BLOCK_ERRORCODE = errno.EWOULDBLOCK
TUNNEL_PREFIX = "ffffffff".decode("HEX")
TUNNEL_PREFIX_LENGHT = 4
class Endpoint(object):
__metaclass__ = ABCMeta
def __init__(self):
self._logger = logging.getLogger(self.__class__.__name__)
self._dispersy = None
@abstractmethod
def get_address(self):
pass
@abstractmethod
def send(self, candidates, packets):
pass
@abstractmethod
def send_packet(self, candidate, packet):
pass
def open(self, dispersy):
self._dispersy = dispersy
return True
def close(self, timeout=0.0):
assert self._dispersy, "Should not be called before open(...)"
assert isinstance(timeout, float), type(timeout)
return True
def log_packet(self, sock_addr, packet, outbound=True):
try:
community = self._dispersy.get_community(packet[2:22], load=False, auto_load=False)
# find associated conversion
conversion = community.get_conversion_for_packet(packet)
name = conversion.decode_meta_message(packet).name
except:
name = "???"
self._logger.debug("%30s %s %15s:%-5d %4d bytes", name, '->' if outbound else '<-',
sock_addr[0], sock_addr[1], len(packet))
if outbound:
self._dispersy.statistics.dict_inc(u"endpoint_send", name)
else:
self._dispersy.statistics.dict_inc(u"endpoint_recv", name)
class NullEndpoint(Endpoint):
"""
NullEndpoint will ignore not send or receive anything.
This Endpoint can be used during unit tests that should not communicate with other peers.
"""
def __init__(self, address=("0.0.0.0", 42)):
super(NullEndpoint, self).__init__()
self._address = address
def get_address(self):
return self._address
def send(self, candidates, packets):
if any(len(packet) > 2 ** 16 - 60 for packet in packets):
raise RuntimeError("UDP does not support %d byte packets" % max(len(packet) for packet in packets))
self._dispersy.statistics.total_up += sum(len(packet) for packet in packets) * len(candidates)
def listen_to(self, prefix, handler):
pass
def send_packet(self, candidate, packet):
if len(packet) > 2 ** 16 - 60:
raise RuntimeError("UDP does not support %d byte packets" % len(packet))
self._dispersy.statistics.total_up += len(packet)
class StandaloneEndpoint(Endpoint):
def __init__(self, port, ip="0.0.0.0"):
super(StandaloneEndpoint, self).__init__()
self._port = port
self._ip = ip
self._running = False
self._add_task = lambda task, delay = 0.0, id = "": None
self._sendqueue_lock = threading.RLock()
self._sendqueue = []
# _THREAD and _THREAD are set during open(...)
self._thread = None
self._socket = None
self.packet_handlers = {}
def listen_to(self, prefix, handler):
self.packet_handlers[prefix] = handler
def stop_listen_to(self, prefix):
del self.packet_handlers[prefix]
def get_address(self):
assert self._dispersy, "Should not be called before open(...)"
return self._socket.getsockname()
def open(self, dispersy):
super(StandaloneEndpoint, self).open(dispersy)
for _ in xrange(10000):
try:
self._logger.debug("Listening at %d", self._port)
self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self._socket.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 870400)
self._socket.bind((self._ip, self._port))
self._socket.setblocking(0)
self._port = self._socket.getsockname()[1]
except socket.error:
self._port += 1
continue
break
self._running = True
self._thread = threading.Thread(name="StandaloneEndpoint", target=self._loop)
self._thread.daemon = True
self._thread.start()
return True
def close(self, timeout=10.0):
self._running = False
result = True
if timeout > 0.0:
self._thread.join(timeout)
if self._thread.is_alive():
self._logger.error("the endpoint thread is still running (after waiting %f seconds)", timeout)
result = False
else:
if self._thread.is_alive():
self._logger.debug("the endpoint thread is still running (use timeout > 0.0 to ensure the thread stops)")
result = False
try:
self._socket.close()
except socket.error as exception:
self._logger.exception("%s", exception)
result = False
return super(StandaloneEndpoint, self).close(timeout) and result
def _loop(self):
assert self._dispersy, "Should not be called before open(...)"
recvfrom = self._socket.recvfrom
socket_list = [self._socket.fileno()]
prev_sendqueue = 0
while self._running:
# This is a tricky, if we are running on the DAS4 whenever a socket is ready for writing all processes of
# this node will try to write. Therefore, we have to limit the frequency of trying to write a bit.
if self._sendqueue and (time() - prev_sendqueue) > 0.1:
read_list, write_list, _ = select(socket_list, socket_list, [], 0.1)
else:
read_list, write_list, _ = select(socket_list, [], [], 0.1)
# Furthermore, if we are allowed to send, process sendqueue immediately
if write_list:
self._process_sendqueue()
prev_sendqueue = time()
if read_list:
packets = []
try:
while True:
(data, sock_addr) = recvfrom(65535)
if data:
packets.append((sock_addr, data))
else:
break
except socket.error as e:
if e.errno != errno.EAGAIN:
self._dispersy.statistics.dict_inc(u"endpoint_recv", u"socket-error-'%s'" % repr(e))
finally:
if packets:
self._logger.debug('%d came in, %d bytes in total', len(packets), sum(len(packet) for _, packet in packets))
self.data_came_in(packets)
def data_came_in(self, packets, cache=True):
assert self._dispersy, "Should not be called before open(...)"
assert isinstance(packets, (list, tuple)), type(packets)
normal_packets = []
for packet in packets:
prefix = next((p for p in self.packet_handlers if
packet[1].startswith(p)), None)
if prefix:
sock_addr, data = packet
self.packet_handlers[prefix](sock_addr, data[len(prefix):])
else:
normal_packets.append(packet)
if normal_packets:
self._dispersy.statistics.total_down += sum(len(data) for _, data in normal_packets)
if self._logger.isEnabledFor(logging.DEBUG):
for sock_addr, data in normal_packets:
self.log_packet(sock_addr, data, outbound=False)
# The endpoint runs on it's own thread, so we can't do a callLater here
reactor.callFromThread(self.dispersythread_data_came_in, normal_packets, time(), cache)
def dispersythread_data_came_in(self, packets, timestamp, cache=True):
assert self._dispersy, "Should not be called before open(...)"
def strip_if_tunnel(packets):
for sock_addr, data in packets:
if data.startswith(TUNNEL_PREFIX):
yield True, sock_addr, data[TUNNEL_PREFIX_LENGHT:]
else:
yield False, sock_addr, data
self._dispersy.on_incoming_packets([(Candidate(sock_addr, tunnel), data)
for tunnel, sock_addr, data
in strip_if_tunnel(packets)],
cache,
timestamp,
u"standalone_ep")
def send(self, candidates, packets, prefix=None):
assert self._dispersy, "Should not be called before open(...)"
assert isinstance(candidates, (tuple, list, set)), type(candidates)
assert all(isinstance(candidate, Candidate) for candidate in candidates), [type(candidate) for candidate in candidates]
assert isinstance(packets, (tuple, list, set)), type(packets)
assert all(isinstance(packet, str) for packet in packets), [type(packet) for packet in packets]
assert all(len(packet) > 0 for packet in packets), [len(packet) for packet in packets]
prefix = prefix or ''
packets = [prefix + packet for packet in packets]
if any(len(packet) > 2 ** 16 - 60 for packet in packets):
raise RuntimeError("UDP does not support %d byte packets" % max(len(packet) for packet in packets))
send_packet = False
for candidate, packet in product(candidates, packets):
if self.send_packet(candidate, packet):
send_packet = True
return send_packet
def send_packet(self, candidate, packet, prefix=None):
assert self._dispersy, "Should not be called before open(...)"
assert isinstance(candidate, Candidate), type(candidate)
assert isinstance(packet, str), type(packet)
assert len(packet) > 0
packet = (prefix or '') + packet
if len(packet) > 2 ** 16 - 60:
raise RuntimeError("UDP does not support %d byte packets" % len(packet))
self._dispersy.statistics.total_up += len(packet)
self._dispersy.statistics.total_send += 1
data = TUNNEL_PREFIX + packet if candidate.tunnel else packet
try:
self._socket.sendto(data, candidate.sock_addr)
if self._logger.isEnabledFor(logging.DEBUG):
self.log_packet(candidate.sock_addr, packet)
except socket.error:
with self._sendqueue_lock:
did_have_senqueue = bool(self._sendqueue)
self._sendqueue.append((time(), candidate.sock_addr, data))
# If we did not have a sendqueue, then we need to call process_sendqueue in order send these messages
if not did_have_senqueue:
self._process_sendqueue()
return True
def _process_sendqueue(self):
assert self._dispersy, "Should not be called before start(...)"
with self._sendqueue_lock:
if self._sendqueue:
index = 0
NUM_PACKETS = min(max(50, len(self._sendqueue) / 10), len(self._sendqueue))
self._logger.debug("%d left in sendqueue, trying to send %d packets",
len(self._sendqueue), NUM_PACKETS)
allowed_timestamp = time() - 300
for i in xrange(NUM_PACKETS):
queued_at, sock_addr, data = self._sendqueue[i]
if queued_at > allowed_timestamp:
try:
self._socket.sendto(data, sock_addr)
index += 1
if self._logger.isEnabledFor(logging.DEBUG):
self.log_packet(sock_addr, data)
except socket.error as e:
if e[0] != SOCKET_BLOCK_ERRORCODE:
self._logger.warning("could not send %d to %s (%d in sendqueue)",
len(data), sock_addr, len(self._sendqueue))
self._dispersy.statistics.dict_inc(u"endpoint_send", u"socket-error")
break
else:
self._dispersy.statistics.dict_inc(u"endpoint_send", u"packet-expired")
index += 1
self._sendqueue = self._sendqueue[index:]
if self._sendqueue:
# And schedule a new attempt
self._add_task(self._process_sendqueue, 0.1, "process_sendqueue")
self._logger.debug("%d left in sendqueue", len(self._sendqueue))
self._dispersy.statistics.cur_sendqueue = len(self._sendqueue)
class ManualEnpoint(StandaloneEndpoint):
def __init__(self, *args, **kwargs):
StandaloneEndpoint.__init__(self, *args, **kwargs)
self.receive_lock = threading.RLock()
self.received_packets = []
def data_came_in(self, packets):
self._logger.debug('added %d packets to receivequeue, %d packets are queued in total', len(packets), len(packets) + len(self.received_packets))
with self.receive_lock:
self.received_packets.extend(packets)
def clear_receive_queue(self):
with self.receive_lock:
packets = self.received_packets
self.received_packets = []
if packets:
self._logger.debug('returning %d packets, %d bytes in total',
len(packets), sum(len(packet) for _, packet in packets))
return packets
def process_receive_queue(self):
packets = self.clear_receive_queue()
self.process_packets(packets)
return packets
def process_packets(self, packets, cache=True):
self._logger.debug('processing %d packets', len(packets))
StandaloneEndpoint.data_came_in(self, packets, cache=cache)