Rev 742 | Only display areas with differences | Regard whitespace | Details | Blame | Last modification | View Log | SVN | RSS feed
| Rev 742 | Rev 743 | ||
|---|---|---|---|
| 1 | import socket |
1 | import socket |
| 2 | from threading import Thread |
2 | from threading import Thread |
| 3 | from threading import Lock |
3 | from threading import Lock |
| 4 | import logging as log |
4 | import logging as log |
| 5 | import time |
5 | import time |
| 6 | 6 | ||
| 7 | class TCPServerThread(Thread): |
7 | class TCPServerThread(Thread): |
| 8 | 8 | ||
| 9 | ownerNotifier = None |
9 | ownerNotifier = None |
| 10 | sock = None |
10 | sock = None |
| 11 | terminated = False |
11 | terminated = False |
| 12 | 12 | ||
| 13 | clientSocks = [] |
13 | clientSocks = [] |
| 14 | clientSocksLock = None |
14 | clientSocksLock = None |
| 15 | 15 | ||
| 16 | def __init__ (self, sock, ownerNotifier): |
16 | def __init__ (self, sock, ownerNotifier): |
| 17 | Thread.__init__(self) |
17 | Thread.__init__(self) |
| 18 | self.sock = sock |
18 | self.sock = sock |
| 19 | self.ownerNotifier = ownerNotifier |
19 | self.ownerNotifier = ownerNotifier |
| - | 20 | self.clientSocks = [] # may not be gc'd between instances of object |
|
| 20 | self.clientSocksLock = Lock() |
21 | self.clientSocksLock = Lock() |
| 21 | 22 | ||
| 22 | def cleanup(self): |
23 | def cleanup(self): |
| 23 | self.sock.close() |
24 | self.sock.close() |
| 24 | self.clientSocksLock.acquire() |
25 | self.clientSocksLock.acquire() |
| 25 | for clsock, addr in self.clientSocks: |
26 | for clsock, addr in self.clientSocks: |
| 26 | clsock.close() |
27 | clsock.close() |
| 27 | self.clientSocksLock.release() |
28 | self.clientSocksLock.release() |
| 28 | 29 | ||
| 29 | if self.ownerNotifier is not None: |
30 | if self.ownerNotifier is not None: |
| 30 | self.ownerNotifier.notify(self.ownerNotifier.TERMINATE) |
31 | self.ownerNotifier.notify(self.ownerNotifier.TERMINATE) |
| 31 | 32 | ||
| 32 | def run(self): |
33 | def run(self): |
| 33 | try: |
34 | try: |
| 34 | log.debug('TCPServerThread.run') |
35 | log.debug('TCPServerThread.run') |
| 35 | while not self.terminated: |
36 | while not self.terminated: |
| 36 | try: |
37 | try: |
| 37 | conn, addr = self.sock.accept() |
38 | conn, addr = self.sock.accept() |
| 38 | print 'Client connected from ' + str(addr) |
39 | print 'Client connected from ' + str(addr) |
| 39 | conn.setblocking(0) |
40 | conn.setblocking(0) |
| 40 | self.clientSocksLock.acquire() |
41 | self.clientSocksLock.acquire() |
| 41 | self.clientSocks.append((conn, addr)) |
42 | self.clientSocks.append((conn, addr)) |
| 42 | self.clientSocksLock.release() |
43 | self.clientSocksLock.release() |
| 43 | except socket.timeout: |
44 | except socket.timeout: |
| 44 | pass |
45 | pass |
| 45 | self.cleanup() |
46 | self.cleanup() |
| 46 | except: # sys.excepthook does not work in threads |
47 | except: # sys.excepthook does not work in threads |
| 47 | import traceback |
48 | import traceback |
| 48 | traceback.print_exc() |
49 | traceback.print_exc() |
| 49 | self.terminated = True |
50 | self.terminated = True |
| 50 | self.cleanup() |
51 | self.cleanup() |
| 51 | 52 | ||
| 52 | def terminate(self): |
53 | def terminate(self): |
| 53 | log.debug('TCPServerThread.terminate') |
54 | log.debug('TCPServerThread.terminate') |
| 54 | self.terminated = True |
55 | self.terminated = True |
| 55 | 56 | ||
| 56 | class TCPServer(): |
57 | class TCPServer(): |
| 57 | 58 | ||
| 58 | serverThread = None |
59 | serverThread = None |
| 59 | ownerNotifier = None |
60 | ownerNotifier = None |
| 60 | port = 50007 |
61 | port = 50007 |
| 61 | host = '' |
62 | host = '' |
| 62 | 63 | ||
| 63 | def __init__(self, port, ownerNotifier = None): |
64 | def __init__(self, port, ownerNotifier = None): |
| 64 | self.port = port |
65 | self.port = port |
| 65 | self.ownerNotifier = ownerNotifier |
66 | self.ownerNotifier = ownerNotifier |
| 66 | 67 | ||
| 67 | def start(self): |
68 | def start(self): |
| 68 | if self.serverThread is not None: |
69 | if self.serverThread is not None: |
| 69 | log.debug('TCPServer is already running') |
70 | log.debug('TCPServer is already running') |
| 70 | return |
71 | return |
| 71 | try: |
72 | try: |
| 72 | s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) |
73 | s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) |
| 73 | except Exception, e: |
74 | except Exception, e: |
| 74 | print 'TCP Server start failed:', e |
75 | print 'TCP Server start failed:', e |
| 75 | return False |
76 | return False |
| 76 | try: |
77 | try: |
| 77 | s.bind((self.host, self.port)) |
78 | s.bind((self.host, self.port)) |
| 78 | s.listen(1) |
79 | s.listen(1) |
| 79 | s.settimeout(0.5) |
80 | s.settimeout(0.5) |
| 80 | except Exception, e: |
81 | except Exception, e: |
| 81 | print 'TCP Server start failed:', e |
82 | print 'TCP Server start failed:', e |
| 82 | s.close() |
83 | s.close() |
| 83 | return False |
84 | return False |
| 84 | try: |
85 | try: |
| 85 | self.serverThread = TCPServerThread(s, self.ownerNotifier) |
86 | self.serverThread = TCPServerThread(s, self.ownerNotifier) |
| 86 | self.serverThread.start() |
87 | self.serverThread.start() |
| 87 | except Exception, e: |
88 | except Exception, e: |
| 88 | print 'TCP Server start failed:', e |
89 | print 'TCP Server start failed:', e |
| 89 | if self.serverThread is not None and self.serverThread.running(): |
90 | if self.serverThread is not None and self.serverThread.running(): |
| 90 | self.serverThread.terminate() |
91 | self.serverThread.terminate() |
| 91 | s.close() |
92 | s.close() |
| 92 | return False |
93 | return False |
| 93 | return True |
94 | return True |
| 94 | 95 | ||
| 95 | def stop(self): |
96 | def stop(self): |
| 96 | if self.running(): |
97 | if self.running(): |
| 97 | self.serverThread.terminate() |
98 | self.serverThread.terminate() |
| 98 | else: |
99 | else: |
| 99 | log.debug('TCPServer not started, or already stopped') |
100 | log.debug('TCPServer not started, or already stopped') |
| 100 | 101 | ||
| 101 | def running(self): |
102 | def running(self): |
| 102 | if self.serverThread is not None: |
103 | if self.serverThread is not None: |
| 103 | return not self.serverThread.terminated |
104 | return not self.serverThread.terminated |
| 104 | else: |
105 | else: |
| 105 | return False |
106 | return False |
| - | 107 | ||
| - | 108 | def getClients(self): |
|
| - | 109 | return self.serverThread.clientSocks |
|
| 106 | 110 | ||
| 107 | def readAll(self): |
111 | def readAll(self): |
| 108 | cldata = [] |
112 | cldata = [] |
| 109 | self.serverThread.clientSocksLock.acquire() |
113 | self.serverThread.clientSocksLock.acquire() |
| 110 | for clsock,addr in self.serverThread.clientSocks: |
114 | for clsock,addr in self.serverThread.clientSocks: |
| 111 | try: |
115 | try: |
| 112 | data = clsock |
116 | data = clsock.recv(1024) |
| 113 | except Exception, e: |
117 | except Exception, e: |
| 114 | log.debug('TCPServer client socket error: ' + str(e)) |
118 | log.debug('TCPServer client ' + str(addr) + ' socket error: ' + str(e)) |
| 115 | self.serverThread.clientSocksLock.release() |
119 | self.serverThread.clientSocksLock.release() |
| 116 | return None |
120 | return None |
| 117 | cldata.append((addr, data)) |
121 | cldata.append((addr, data)) |
| 118 | self.serverThread.clientSocksLock.release() |
122 | self.serverThread.clientSocksLock.release() |
| 119 | return cldata |
123 | return cldata |
| 120 | 124 | ||
| 121 | def writeAll(self, data): |
125 | def writeAll(self, data): |
| 122 | self.serverThread.clientSocksLock.acquire() |
126 | self.serverThread.clientSocksLock.acquire() |
| 123 | for clsock,addr in self.serverThread.clientSocks: |
127 | for clsock,addr in self.serverThread.clientSocks: |
| 124 | try: |
128 | try: |
| 125 | clsock.send(data) |
129 | clsock.send(data) |
| 126 | except Exception, e: |
130 | except Exception, e: |
| 127 | log.debug('TCPServer client socket error: ' + str(e)) |
131 | log.debug('TCPServer client ' + str(addr) + ' socket error: ' + str(e)) |
| 128 | self.serverThread. |
132 | self.serverThread.clientSocks.remove((clsock,addr)) |
| 129 | return None |
- | |
| 130 | self.serverThread.clientSocksLock.release() |
133 | self.serverThread.clientSocksLock.release() |
| 131 | - | ||
| - | 134 | return False |
|
| - | 135 | self.serverThread.clientSocksLock.release() |
|
| - | 136 | return True |
|
| 132 | 137 | ||