From f21b3e7354649f884486b5711b628f58e33554f3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Torbj=C3=B6rn=20Axelsson?= Date: Sat, 16 Dec 2017 08:43:08 -0800 Subject: [PATCH] Improved threading and configuration options --- bumper.py | 19 +++++++------ bumper/confserver.py | 67 ++++++++++++++++++++++++++++++++++++-------- bumper/xmppserver.py | 59 +++++++++++++++++++++----------------- 3 files changed, 101 insertions(+), 44 deletions(-) diff --git a/bumper.py b/bumper.py index cdbd414..8dbf6dd 100644 --- a/bumper.py +++ b/bumper.py @@ -2,15 +2,18 @@ import logging import bumper +import sys, socket logging.basicConfig(level=logging.INFO, - format='%(levelname)-8s %(message)s') + format='%(asctime)s %(levelname)-8s %(message)s') -bumper.ConfServer() -bumper.XMPPServer() +conf_address = (socket.gethostbyname(socket.gethostname()), 8007) +xmpp_address = (socket.gethostbyname(socket.gethostname()), 5223) -try: - while True: - pass -except KeyboardInterrupt: - logging.info('keyboard interrupt') +# start conf server (async) +conf_server = bumper.ConfServer(conf_address, ssl=False, async=True) + +# start xmpp server (sync) +xmpp_server = bumper.XMPPServer(xmpp_address) + +conf_server.disconnect() diff --git a/bumper/confserver.py b/bumper/confserver.py index 0f6bebc..1fa4357 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -3,7 +3,8 @@ from http.server import HTTPServer from http.server import BaseHTTPRequestHandler from http import HTTPStatus -import socket, logging, _thread, json +from threading import Thread +import socket, logging, ssl, json, sys class RequestHandler(BaseHTTPRequestHandler): @@ -15,10 +16,18 @@ class RequestHandler(BaseHTTPRequestHandler): logging.debug("Headers: " + str(self.headers)) request_body = post_data.decode('utf-8') logging.debug("Request: " + request_body) - if request_body.find('EcoMsgNew') > -1: - body = '{{"result":"ok","ip":"{}","port":5223}}'.format(socket.gethostbyname(socket.gethostname())) - else: - body = '{"result":"ok","ip":"47.88.66.164","port":8005}' + json_body = json.loads(request_body) + todo = json_body['todo'] + if todo == 'FindBest': + service = json_body['service'] + if service == 'EcoMsgNew': + body = '{{"result":"ok","ip":"{}","port":5223}}'.format(socket.gethostbyname(socket.gethostname())) + elif service == 'EcoUpdate': + body = '{"result":"ok","ip":"47.88.66.164","port":8005}' + elif todo == 'loginByItToken': + body = "{{'todo': 'result', 'result': 'ok', 'userId': '{}', 'resource': '{}', 'token': '{}'}}".format(json_body['userId'], json_body['resource'], json_body['token']) + elif todo == 'GetDeviceList': + body = "{'todo': 'result', 'result': 'ok', 'devices': [{'did': '{}', 'name': '{}', 'class': '{}', 'resource': 'atom', 'nick': None, 'company': 'eco'}]}" logging.debug("Response: " + body) body = body.encode() self.send_response(HTTPStatus.OK) @@ -27,17 +36,53 @@ class RequestHandler(BaseHTTPRequestHandler): self.send_header('Content-Length', len(body)) self.end_headers() self.wfile.write(body) - return except Exception as e: logging.error(e) +class HTTPServerThread(HTTPServer, Thread): + def __init__(self, server_address): + Thread.__init__(self) + self.server_address = server_address + self.handler = RequestHandler + self.exit_flag = False + def handle_error(self, request, client_address): + self.close_request(request) + def run(self): + try: + HTTPServer.__init__(self, self.server_address, self.handler) + logging.info('ConfServer: listening on {}:{}'.format(self.server_address[0], self.server_address[1])) + while not self.exit_flag: + self.handle_request() + except Exception as e: + logging.error(e) + def disconnect(self): + self.exit_flag = True + # make a connection to + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.connect(self.server_address) + s.close() + + class ConfServer(): - def __init__(self): + def __init__(self, address, ssl=False, async=True): try: - server_address = (socket.gethostbyname(socket.gethostname()), 8007) - httpd = HTTPServer(server_address, RequestHandler) - logging.info("ConfServer: running on http://{}:{}".format(server_address[0], server_address[1])) - _thread.start_new_thread(httpd.serve_forever, ()) + self.async = async + self.server = HTTPServerThread(address) + if ssl: + self.server.socket = ssl.wrap_socket(self.server.socket, keyfile='./certs/key.pem', certfile='./certs/cert.pem', server_side=True) + if self.async: + self.server.start() + else: + try: + self.server.run() + except KeyboardInterrupt: + self.disconnect() except Exception as e: logging.error(e) + def disconnect(self): + logging.info('ConfServer: shutting down...') + self.server.disconnect() + if(self.async): + self.server.join() + logging.info('ConfServer: bye') diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index 3ce3dc4..583b481 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -1,7 +1,6 @@ #!/usr/bin/env python3 -import sys, socket, _thread, re, time, logging, uuid -import xml.etree.ElementTree as ET +import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET class XMPPServer(): @@ -9,35 +8,46 @@ class XMPPServer(): bot_id = 'bumpy' client_id = None clients = [] + exit_flag = False - def __init__(self): + def __init__(self, address): try: # Initialize bot server - server = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - server_address = (socket.gethostbyname(socket.gethostname()), 5223) - server.bind(server_address) - server.listen(1) - logging.info('XMPPServer: listening on {}:{}'.format(server_address[0], server_address[1])) - while True: + self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self.socket.bind(address) + self.socket.listen(1) + logging.info('XMPPServer: listening on {}:{}'.format(address[0], address[1])) + while not self.exit_flag: logging.info('XMPPServer: awaiting connection') - connection, client_address = server.accept() + connection, client_address = self.socket.accept() # disconnect any clients with this ip for client in self.clients: if client.address == client_address[0]: client.disconnect() - _thread.start_new_thread(Client,(connection, client_address)) - except KeyboardInterrupt: - logging.info('keyboard interrupt') - server.shutdown(2) - server.close() + thread_id = uuid.uuid4() + client = Client(thread_id, connection, client_address) + client.start() + self.clients.append(client) + self.socket.close() except Exception as e: - server.shutdown(2) - server.close() - logging.error(e) + logging.error('e: ' + e) + except KeyboardInterrupt: + logging.debug('XMPPServer: Keyboard interrupt') + finally: + self.disconnect() + logging.info('XMPPServer: bye') + + def disconnect(self): + logging.info('XMPPServer: waiting for all client threads to exit') + for client in self.clients: + client.disconnect() + client.join() + self.exit_flag = True + logging.info('XMPPServer: shutting down...') -class Client(): +class Client(threading.Thread): IDLE = 0 CONNECT = 1 INIT = 2 @@ -48,14 +58,13 @@ class Client(): BOT = 1 CONTROLLER = 2 - def __init__(self, connection, client_address): - self.id = uuid.uuid4() + def __init__(self, thread_id, connection, client_address): + threading.Thread.__init__(self) + self.id = thread_id self.type = self.UNKNOWN self.state = self.IDLE self.connection = connection self.address = client_address[0] - XMPPServer.clients.append(self) - self._main() def send(self, command): logging.debug('to {}: {}'.format(self.address, command)) @@ -102,7 +111,7 @@ class Client(): logging.debug('sending result: ' + data.decode('utf-8')) client.send(data.decode('utf-8')) - def _main(self): + def run(self): try: logging.info('client connected: {}'.format(self.address)) self._set_state('CONNECT')