Newer
Older
time.sleep(retry_interval)
if self._cancelled:
if not connected:
continue
## tunnel data between sockets
# init socket poller
mask_ro = select.EPOLLIN ^ select.EPOLLERR ^ select.EPOLLHUP
mask_rw = mask_ro ^ select.EPOLLOUT
poller = select.epoll()
poller.register(self._sock_server, mask_ro)
poller.register(self._sock_client, mask_ro)
# forward data until a connection is closed
buf_server = ''
buf_client = ''
connected = True
while not self._cancelled and connected:
# wait for events on both sockets
poll = poller.poll()
for fd, events in poll:
# events are for server socket
if fd == self._sock_server.fileno():
# read available
if events & select.EPOLLIN:
# read incoming data
try:
read = self._sock_server.recv(4096)
except socket.error:
connected = False
if not len(read):
empty_buf += 1
else:
empty_buf = 0
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
buf_server += read
# set the other socket to notify us when it's
# available for writing
if len(buf_server):
poller.modify(self._sock_client, mask_rw)
# write available
if events & select.EPOLLOUT:
# try to send the whole buffer
try:
sent = self._sock_server.send(buf_client)
except socket.error:
connected = False
# drop sent data from the buffer
buf_client = buf_client[sent:]
# if the buffer becomes empty, stop write polling
if not len(buf_client):
poller.modify(self._sock_server, mask_ro)
# events for client socket
else:
# read available
if events & select.EPOLLIN:
# read incoming data
try:
read = self._sock_client.recv(4096)
except socket.error:
connected = False
if not len(read):
empty_buf += 1
else:
empty_buf = 0
buf_client += read
# set the other socket to notify us when it's
# available for writing
if len(buf_client):
poller.modify(self._sock_server, mask_rw)
# write available
if events & select.EPOLLOUT:
# try to send the whole buffer
try:
sent = self._sock_client.send(buf_server)
except socket.error:
connected = False
# drop sent data from the buffer
buf_server = buf_server[sent:]
# if the buffer becomes empty, stop write polling
if not len(buf_server):
poller.modify(self._sock_client, mask_ro)
if empty_buf >= 10:
connected = False
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
if connected is False:
self._log.append('disconnected')
try:
self._sock_server.close()
except: pass
try:
self._sock_client.close()
except: pass
except Exception as err:
traceback.print_exc()
raise err
def tunnel_get_server_port(self):
'''
'''
self._event_server_bound.wait()
return self._server_port
def tunnel_connect(self, endpoint):
'''
'''
with self._mutex:
if not self._event_client_conf.is_set():
self._endpoint = endpoint
self._event_client_conf.set()
else:
raise TCPTunnelJobError('remote endpoint already set')
def tunnel_listen(self, local):
'''
'''
with self._mutex:
if not self._event_server_conf.is_set():
self._server_local = local is True
self._event_server_conf.set()
else:
raise TCPTunnelJobError('server parameters already set')