-
Notifications
You must be signed in to change notification settings - Fork 3
/
Copy pathpush_tcp.py
228 lines (200 loc) · 6.95 KB
/
push_tcp.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
#!/usr/bin/env python
"""
Mark Nottingham's push-based asynchronous TCP library (with Dugsong's pyevent)
*** Create Server
server = push_tcp.create_server(host, port, conn_handler)
*** Create Client
host = 'www.example.com'
port = '80'
push_tcp.create_client(host, port, conn_handler, error_handler)
*** Handler
def conn_handler(tcp_conn):
print "connected to %s:%s" % tcp_conn.host, tcp_conn.port
return read_cb, close_cb, pause_cb
def error_handler(host, port, reason):
print "can't connect to %s:%s: %s" % (host, port, reason)
"""
import sys
import socket
import errno
import event
class _TCPConnection(object):
"""
Base class for a TCP connection
"""
write_bufsize = 8
read_bufsize = 1024 * 8
def __init__(self, sock, host, port, connect_error_handler=None):
self.socket = sock
self.host = host
self.port = port
self.connect_error_handler = connect_error_handler
self.read_cb = None
self.close_cb = None
self._close_cb_called = False
self.pause_cb = None
self.tcp_connected = True
self._paused = False
self._closing = False
self.write_buffer = []
self._revent = event.read(sock, self.handle_read)
self._wevent = event.write(sock, self.handle_write)
def handle_read(self):
"""
The connection has data read for reading;
call read_cb if appropriate
"""
try:
data = self.socket.recv(self.read_bufsize)
except socket.error, why:
if why[0] in [errno.EBADF, errno.ECONNRESET, errno.EPIPE, errno.TIMEOUT]:
self.conn_closed()
return
elif why[0] in [errno.ECONNREFUSED, errno.ENETUNREACH] and self.connect_error_handler:
self.tcp_connected = False
self.connect_error_handler(why[0])
return
else:
raise
if data == "":
self.conn_closed()
else:
self.read_cb(data)
if self.read_cb and self.tcp_connected and not self._paused:
return self._revent
def handle_write(self):
"""
The connection is ready for writing;
write any buffered data
"""
if len(self._write_buffer) > 0:
data = "".join(self._write_buffer)
try:
sent = self.socket.send(data)
except socket.error, why:
if why[0] in [errno.EBADF, errno.ECONNRESET, errno.EPIPE, errno.ETIMEOUT]:
self.conn_closed()
return
elif why[0] in [errno.ECONNREFUSED, errno.ENETUNREACH] and self.connect_error_handler:
self.tcp_connected = False
self.connect_error_handler(why[0])
return
else:
raise
if sent < len(data):
self._write_buffer = [data[sent:]]
else:
self._write_buffer = []
if self.pause_cb and len(self._write_buffer) < self.write_bufsize:
self.pause_cb(False)
if self._closing:
self.close()
if self.tcp_connected and (len(self._write_buffer) > 0 or self._closing):
return self._wevent
def conn_closed(self):
"""
The connection has been closed by the other side.
Do local cleanup and then call close_cb
"""
self.tcp_connected = False
if self._close_cb_called:
return
elif self.close_cb:
self._close_cb_called = True
self.close_cb()
else:
schedule(1, self.conn_closed)
def write(self, data):
"""
Write data to the connection
"""
self._write_buffer.append(data)
if self.pause_cb and len(self._write_buffer) > self.write_bufsize:
self.pause_cb(True)
if not self._wevent.pending():
self._wevent.add()
def pause(self, paused):
"""
Temporarily stop/start reading from the connection
and pushing it to the app
"""
if paused:
if self._revent.pending():
self._revent.delete()
else:
if not self._revent.pending():
self._revent.add()
self._paused = paused
def close(self):
"""
Flush buffered data (if any) and close the connection
"""
self.pause(True)
if len(self._write_buffer) > 0:
self._closing = True
else:
self.socket.close()
self.tcp_connected = False
class create_server(object):
"""
An asynchronous TCP server
"""
def __init__(self, host, port, conn_handler):
self.host = host
self.port = port
self.conn_handler = conn_handler
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setblocking(0)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADD, 1)
sock.bind((host, port))
sock.listen(socket.SOMAXCONN)
event.event(self.handle_accept, handle=sock, evtype=event.EV_READ|event.EV_PERSIST).add()
def handle_accept(self, *args):
conn, addr = args[1].accept()
tcp_conn = _TCPConnection(conn, self.host, self.port, self.handle_error)
tcp_conn.read_cb, tcp_conn.close_cb, tcp_conn.pause_cb = self.conn_handler(tcp_conn)
def handle_error(self, err=None):
raise AssertionError, "this (%s) should never happen for a server" % err
class create_client(object):
"""
An asynchronous TCP client
"""
def __init__(self, host, port, conn_handler, connect_error_handler, timeout=None):
self.host = host
self.port = port
self.conn_handler = conn_handler
self.connect_error_handler = connect_error_handler
self._timeout_ev = None
self._conn_handled = False
self._error_sent = False
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setblocking(0)
event.write(sock, self.handle_connect, sock).add()
try:
err = sock.connect_ex((host, port))
except socket.error, why:
self.handle_error(why)
return
if err != errno.EINPROGRESS:
self.handle_error(err)
if timeout:
to_err = errno.ETIMEDOUT
self._timeout_ev = schedule(timeout, self.handle_error, to_err)
def handle_connect(self, sock=None):
if self._timeout_ev:
self._timeout_ev.delete()
tcp_conn = _TCPConnection(sock, self.host, self.port, self.handle_error)
tcp_conn.read_cb, tcp_conn.close_cb, tcp_conn.pause_cb = self.conn_handler(tcp_conn)
def handle_write(self):
pass
def handle_error(self, err=None):
if self._timeout_ev:
self._timeout_ev.delete()
if not self._error_sent:
self._error_sent = True
if err == None:
t, err, tb = sys.exc_info()
self.connect_error_handler(self.host, self.port, err)
schedule = event.timeout
run = event.dispatch
stop = event.abort