Newer
Older
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
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
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
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()
## process events for read/write operations only
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
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
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 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')