Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathreactor.py
More file actions
Latest commit
163 lines (139 loc) · 4.7 KB
/
Copy pathreactor.py
File metadata and controls
163 lines (139 loc) · 4.7 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
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
importsocket
importselect
try:
importssl
assertssl
exceptImportError:
ssl=False
try:
from . importapi, msg
from .. importeditor
from ..common.exc_fmtimportstr_e, pp_e
from ..common.handlersimporttcp_server
assertmsgandtcp_server
except (ImportError, ValueError):
fromfloo.common.exc_fmtimportstr_e, pp_e
fromfloo.common.handlersimporttcp_server
fromfloo.commonimportapi, msg
fromflooimporteditor
reactor=None
class_Reactor(object):
''' Low level event driver '''
def__init__(self):
self._protos= []
self._handlers= []
self.on_stop=None
defconnect(self, factory, host, port, secure, conn=None):
proto=factory.build_protocol(host, port, secure)
self._protos.append(proto)
proto.connect(conn)
self._handlers.append(factory)
deflisten(self, factory, host='127.0.0.1', port=0):
listener_factory=tcp_server.TCPServerHandler(factory, self)
proto=listener_factory.build_protocol(host, port)
factory.listener_factory=listener_factory
self._protos.append(proto)
self._handlers.append(listener_factory)
returnproto.sockname()
defstop_handler(self, handler):
try:
handler.proto.stop()
exceptExceptionase:
msg.warn('Error stopping connection: ', str_e(e))
try:
self._handlers.remove(handler)
exceptException:
pass
try:
self._protos.remove(handler.proto)
exceptException:
pass
ifhasattr(handler, 'listener_factory'):
returnhandler.listener_factory.stop()
ifnotself._handlersandnotself._protos:
msg.log('All handlers stopped. Stopping reactor.')
self.stop()
defstop(self):
for_conninself._protos:
_conn.stop()
self._protos= []
self._handlers= []
msg.log('Reactor shut down.')
editor.status_message('Disconnected.')
ifself.on_stop:
self.on_stop()
defis_ready(self):
ifnotself._handlers:
returnFalse
forfinself._handlers:
ifnotf.is_ready():
returnFalse
returnTrue
def_reconnect(self, fd, *fd_sets):
forfd_setinfd_sets:
try:
fd_set.remove(fd)
exceptValueError:
pass
fd.reconnect()
@api.send_errors
deftick(self, timeout=0):
forfactoryinself._handlers:
factory.tick()
self.select(timeout)
editor.call_timeouts()
defblock(self):
whileself._protosorself._handlers:
self.tick(.05)
defselect(self, timeout=0):
ifnotself._protos:
return
readable= []
writeable= []
errorable= []
fd_map= {}
forfdinself._protos:
fileno=fd.fileno()
ifnotfileno:
continue
fd.fd_set(readable, writeable, errorable)
fd_map[fileno] =fd
ifnotreadableandnotwriteable:
return
try:
_in, _out, _except=select.select(readable, writeable, errorable, timeout)
except (select.error, socket.error, Exception) ase:
# TODO: with multiple FDs, must call select with just one until we find the error :(
forfilenoinreadable:
try:
select.select([fileno], [], [], 0)
except (select.error, socket.error, Exception) ase:
fd_map[fileno].reconnect()
msg.error('Error in select(): ', fileno, str_e(e))
return
forfilenoin_except:
fd=fd_map[fileno]
self._reconnect(fd, _in, _out)
forfilenoin_out:
fd=fd_map[fileno]
try:
fd.write()
exceptssl.SSLErrorase:
ife.args[0] !=ssl.SSL_ERROR_WANT_WRITE:
raise
exceptExceptionase:
msg.error('Couldn\'t write to socket: ', str_e(e))
msg.debug('Couldn\'t write to socket: ', pp_e(e))
returnself._reconnect(fd, _in)
forfilenoin_in:
fd=fd_map[fileno]
try:
fd.read()
exceptssl.SSLErrorase:
ife.args[0] !=ssl.SSL_ERROR_WANT_READ:
raise
exceptExceptionase:
msg.error('Couldn\'t read from socket: ', str_e(e))
msg.debug('Couldn\'t read from socket: ', pp_e(e))
fd.reconnect()
reactor=_Reactor()