Details | Last modification | View Log | SVN | RSS feed
| Rev | Author | Line No. | Line |
|---|---|---|---|
| 736 | olof | 1 | ########################################################### |
| 2 | # |
||
| 3 | # Implements the TCP server interface to administrative |
||
| 4 | # or monitoring apps |
||
| 5 | # |
||
| 6 | ########################################################### |
||
| 7 | |||
| 751 | olof | 8 | from TCPServerBase import TCPServerBase |
| 9 | |||
| 752 | olof | 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 |
||
| 55 | |||
| 736 | olof | 56 | class TCPServer1(TCPServerBase): |
| 752 | olof | 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) |
||
| 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 |