Rev 751 | Show entire file | Regard whitespace | Details | Blame | Last modification | View Log | SVN | RSS feed
| Rev 751 | Rev 752 | ||
|---|---|---|---|
| Line 4... | Line 4... | ||
| 4 | # or monitoring apps |
4 | # or monitoring apps |
| 5 | # |
5 | # |
| 6 | ########################################################### |
6 | ########################################################### |
| 7 | 7 | ||
| 8 | from TCPServerBase import TCPServerBase |
8 | from TCPServerBase import TCPServerBase |
| - | 9 | ||
| - | 10 | from threading import Thread |
|
| - | 11 | from threading import Lock |
|
| - | 12 | import logging as log |
|
| - | 13 | import time |
|
| - | 14 | from TCPServer import TCPServer |
|
| - | 15 | ||
| - | 16 | class TCPServer1Thread(Thread): |
|
| - | 17 | ownerNotifier = None |
|
| - | 18 | terminated = False |
|
| - | 19 | serverListener = None |
|
| - | 20 | tcpServer = None |
|
| - | 21 | port = None |
|
| - | 22 | ||
| - | 23 | def __init__ (self, ownerNotifier, serverListener, port): |
|
| - | 24 | Thread.__init__(self) |
|
| - | 25 | self.ownerNotifier = ownerNotifier |
|
| - | 26 | self.serverListener = serverListener |
|
| - | 27 | self.port = port |
|
| - | 28 | ||
| - | 29 | def cleanup(self): |
|
| - | 30 | if self.tcpServer.running(): |
|
| - | 31 | self.tcpServer.stop() |
|
| - | 32 | self.ownerNotifier.notify(self.ownerNotifier.TERMINATE) |
|
| - | 33 | ||
| - | 34 | def run(self): |
|
| - | 35 | try: |
|
| - | 36 | log.debug('TCPServer1Thread.run') |
|
| - | 37 | self.tcpServer = TCPServer(self.port, self.ownerNotifier) |
|
| - | 38 | self.tcpServer.start() |
|
| - | 39 | ||
| - | 40 | while not self.terminated: |
|
| - | 41 | while not self.serverListener.empty(): |
|
| - | 42 | data = self.serverListener.read() + '\n' |
|
| - | 43 | self.tcpServer.writeAll(data) |
|
| - | 44 | time.sleep(0.01) |
|
| - | 45 | self.cleanup() |
|
| - | 46 | except: # sys.excepthook does not work in threads |
|
| - | 47 | import traceback |
|
| - | 48 | traceback.print_exc() |
|
| - | 49 | self.terminated = True |
|
| - | 50 | self.cleanup() |
|
| - | 51 | ||
| - | 52 | def terminate(self): |
|
| - | 53 | log.debug('TCPServer1Thread.terminate') |
|
| - | 54 | self.terminated = True |
|
| 9 | 55 | ||
| 10 | class TCPServer1(TCPServerBase): |
56 | class TCPServer1(TCPServerBase): |
| - | 57 | ||
| - | 58 | DEFAULT_CONFIG = {'ethif':'any', |
|
| - | 59 | 'port':1200, |
|
| - | 60 | 'compat_1': True, |
|
| - | 61 | 'canif': None, |
|
| - | 62 | } |
|
| - | 63 | ||
| - | 64 | pktHandler = None |
|
| - | 65 | serverThread = None |
|
| - | 66 | ownerNotifier = None |
|
| - | 67 | ||
| - | 68 | pktBuffer = [] |
|
| - | 69 | pktBufferLock = None |
|
| - | 70 | ||
| - | 71 | class ServerListener(): |
|
| - | 72 | ||
| - | 73 | parent = None |
|
| - | 74 | ||
| - | 75 | def __init__(self, parent): |
|
| - | 76 | self.parent = parent |
|
| - | 77 | ||
| - | 78 | def empty(self): |
|
| - | 79 | if len(self.parent.pktBuffer) > 0: |
|
| - | 80 | return False |
|
| - | 81 | return True |
|
| - | 82 | ||
| - | 83 | def full(self): |
|
| - | 84 | return False |
|
| - | 85 | ||
| - | 86 | def lock(self): |
|
| - | 87 | self.parent.pktBufferLock.acquire() |
|
| - | 88 | ||
| - | 89 | def unlock(self): |
|
| - | 90 | self.parent.pktBufferLock.release() |
|
| - | 91 | ||
| - | 92 | def write(self, data): |
|
| - | 93 | self.lock() |
|
| - | 94 | self.parent.pktBuffer.append(data) |
|
| - | 95 | self.unlock() |
|
| - | 96 | ||
| - | 97 | def read(self): |
|
| - | 98 | self.lock() |
|
| - | 99 | data = self.parent.pktBuffer.pop(0) |
|
| - | 100 | self.unlock() |
|
| - | 101 | return data |
|
| - | 102 | ||
| - | 103 | ||
| - | 104 | def __init__(self, pktHandler, cfg = None): |
|
| - | 105 | if cfg is None: |
|
| - | 106 | self.config = self.DEFAULT_CONFIG |
|
| - | 107 | else: |
|
| - | 108 | self.config = cfg |
|
| - | 109 | self.pktBufferLock = Lock() |
|
| - | 110 | self.pktHandler = pktHandler |
|
| - | 111 | log.debug('Creating tcpd server, config: ' + str(self.config)) |
|
| - | 112 | ||
| - | 113 | def setPktHandler(self, pktHandler): |
|
| - | 114 | self.pktHandler = pktHandler |
|
| - | 115 | ||
| - | 116 | def setIfNotifier(self, ifNotifier): |
|
| - | 117 | self.ownerNotifier = ifNotifier |
|
| - | 118 | ||
| - | 119 | def start(self): |
|
| - | 120 | ethif = self.config['ethif'] |
|
| - | 121 | port = self.config['port'] |
|
| - | 122 | compat_1 = self.config['compat_1'] |
|
| - | 123 | canif = self.config['canif'] |
|
| - | 124 | ||
| - | 125 | serverListener = self.ServerListener(self) |
|
| - | 126 | self.pktHandler.addServerListener(serverListener) |
|
| 11 |
|
127 | |
| - | 128 | self.serverThread = TCPServer1Thread(self.ownerNotifier, serverListener, port) |
|
| - | 129 | # self.serverThread = TCPServer(port, self.ownerNotifier) |
|
| - | 130 | self.serverThread.start() |
|
| - | 131 | return True |
|
| - | 132 | ||
| - | 133 | def stop(self): |
|
| - | 134 | if self.running(): |
|
| - | 135 | self.serverThread.terminate() |
|
| - | 136 | else: |
|
| - | 137 | log.debug('TCP interface not started, or already stopped') |
|
| - | 138 | ||
| - | 139 | def running(self): |
|
| - | 140 | if self.serverThread is not None: |
|
| - | 141 | return not self.serverThread.terminated |
|
| - | 142 | else: |
|
| - | 143 | return False |
|