Improved threading and configuration options
This commit is contained in:
parent
d9ceeea556
commit
f21b3e7354
3 changed files with 101 additions and 44 deletions
19
bumper.py
19
bumper.py
|
|
@ -2,15 +2,18 @@
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import bumper
|
import bumper
|
||||||
|
import sys, socket
|
||||||
|
|
||||||
logging.basicConfig(level=logging.INFO,
|
logging.basicConfig(level=logging.INFO,
|
||||||
format='%(levelname)-8s %(message)s')
|
format='%(asctime)s %(levelname)-8s %(message)s')
|
||||||
|
|
||||||
bumper.ConfServer()
|
conf_address = (socket.gethostbyname(socket.gethostname()), 8007)
|
||||||
bumper.XMPPServer()
|
xmpp_address = (socket.gethostbyname(socket.gethostname()), 5223)
|
||||||
|
|
||||||
try:
|
# start conf server (async)
|
||||||
while True:
|
conf_server = bumper.ConfServer(conf_address, ssl=False, async=True)
|
||||||
pass
|
|
||||||
except KeyboardInterrupt:
|
# start xmpp server (sync)
|
||||||
logging.info('keyboard interrupt')
|
xmpp_server = bumper.XMPPServer(xmpp_address)
|
||||||
|
|
||||||
|
conf_server.disconnect()
|
||||||
|
|
|
||||||
|
|
@ -3,7 +3,8 @@
|
||||||
from http.server import HTTPServer
|
from http.server import HTTPServer
|
||||||
from http.server import BaseHTTPRequestHandler
|
from http.server import BaseHTTPRequestHandler
|
||||||
from http import HTTPStatus
|
from http import HTTPStatus
|
||||||
import socket, logging, _thread, json
|
from threading import Thread
|
||||||
|
import socket, logging, ssl, json, sys
|
||||||
|
|
||||||
|
|
||||||
class RequestHandler(BaseHTTPRequestHandler):
|
class RequestHandler(BaseHTTPRequestHandler):
|
||||||
|
|
@ -15,10 +16,18 @@ class RequestHandler(BaseHTTPRequestHandler):
|
||||||
logging.debug("Headers: " + str(self.headers))
|
logging.debug("Headers: " + str(self.headers))
|
||||||
request_body = post_data.decode('utf-8')
|
request_body = post_data.decode('utf-8')
|
||||||
logging.debug("Request: " + request_body)
|
logging.debug("Request: " + request_body)
|
||||||
if request_body.find('EcoMsgNew') > -1:
|
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()))
|
body = '{{"result":"ok","ip":"{}","port":5223}}'.format(socket.gethostbyname(socket.gethostname()))
|
||||||
else:
|
elif service == 'EcoUpdate':
|
||||||
body = '{"result":"ok","ip":"47.88.66.164","port":8005}'
|
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)
|
logging.debug("Response: " + body)
|
||||||
body = body.encode()
|
body = body.encode()
|
||||||
self.send_response(HTTPStatus.OK)
|
self.send_response(HTTPStatus.OK)
|
||||||
|
|
@ -27,17 +36,53 @@ class RequestHandler(BaseHTTPRequestHandler):
|
||||||
self.send_header('Content-Length', len(body))
|
self.send_header('Content-Length', len(body))
|
||||||
self.end_headers()
|
self.end_headers()
|
||||||
self.wfile.write(body)
|
self.wfile.write(body)
|
||||||
return
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error(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():
|
class ConfServer():
|
||||||
def __init__(self):
|
def __init__(self, address, ssl=False, async=True):
|
||||||
try:
|
try:
|
||||||
server_address = (socket.gethostbyname(socket.gethostname()), 8007)
|
self.async = async
|
||||||
httpd = HTTPServer(server_address, RequestHandler)
|
self.server = HTTPServerThread(address)
|
||||||
logging.info("ConfServer: running on http://{}:{}".format(server_address[0], server_address[1]))
|
if ssl:
|
||||||
_thread.start_new_thread(httpd.serve_forever, ())
|
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:
|
except Exception as e:
|
||||||
logging.error(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')
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,6 @@
|
||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
|
|
||||||
import sys, socket, _thread, re, time, logging, uuid
|
import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET
|
||||||
import xml.etree.ElementTree as ET
|
|
||||||
|
|
||||||
|
|
||||||
class XMPPServer():
|
class XMPPServer():
|
||||||
|
|
@ -9,35 +8,46 @@ class XMPPServer():
|
||||||
bot_id = 'bumpy'
|
bot_id = 'bumpy'
|
||||||
client_id = None
|
client_id = None
|
||||||
clients = []
|
clients = []
|
||||||
|
exit_flag = False
|
||||||
|
|
||||||
def __init__(self):
|
def __init__(self, address):
|
||||||
try:
|
try:
|
||||||
# Initialize bot server
|
# Initialize bot server
|
||||||
server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||||
server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||||
server_address = (socket.gethostbyname(socket.gethostname()), 5223)
|
self.socket.bind(address)
|
||||||
server.bind(server_address)
|
self.socket.listen(1)
|
||||||
server.listen(1)
|
logging.info('XMPPServer: listening on {}:{}'.format(address[0], address[1]))
|
||||||
logging.info('XMPPServer: listening on {}:{}'.format(server_address[0], server_address[1]))
|
while not self.exit_flag:
|
||||||
while True:
|
|
||||||
logging.info('XMPPServer: awaiting connection')
|
logging.info('XMPPServer: awaiting connection')
|
||||||
connection, client_address = server.accept()
|
connection, client_address = self.socket.accept()
|
||||||
# disconnect any clients with this ip
|
# disconnect any clients with this ip
|
||||||
for client in self.clients:
|
for client in self.clients:
|
||||||
if client.address == client_address[0]:
|
if client.address == client_address[0]:
|
||||||
client.disconnect()
|
client.disconnect()
|
||||||
_thread.start_new_thread(Client,(connection, client_address))
|
thread_id = uuid.uuid4()
|
||||||
except KeyboardInterrupt:
|
client = Client(thread_id, connection, client_address)
|
||||||
logging.info('keyboard interrupt')
|
client.start()
|
||||||
server.shutdown(2)
|
self.clients.append(client)
|
||||||
server.close()
|
self.socket.close()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
server.shutdown(2)
|
logging.error('e: ' + e)
|
||||||
server.close()
|
except KeyboardInterrupt:
|
||||||
logging.error(e)
|
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
|
IDLE = 0
|
||||||
CONNECT = 1
|
CONNECT = 1
|
||||||
INIT = 2
|
INIT = 2
|
||||||
|
|
@ -48,14 +58,13 @@ class Client():
|
||||||
BOT = 1
|
BOT = 1
|
||||||
CONTROLLER = 2
|
CONTROLLER = 2
|
||||||
|
|
||||||
def __init__(self, connection, client_address):
|
def __init__(self, thread_id, connection, client_address):
|
||||||
self.id = uuid.uuid4()
|
threading.Thread.__init__(self)
|
||||||
|
self.id = thread_id
|
||||||
self.type = self.UNKNOWN
|
self.type = self.UNKNOWN
|
||||||
self.state = self.IDLE
|
self.state = self.IDLE
|
||||||
self.connection = connection
|
self.connection = connection
|
||||||
self.address = client_address[0]
|
self.address = client_address[0]
|
||||||
XMPPServer.clients.append(self)
|
|
||||||
self._main()
|
|
||||||
|
|
||||||
def send(self, command):
|
def send(self, command):
|
||||||
logging.debug('to {}: {}'.format(self.address, command))
|
logging.debug('to {}: {}'.format(self.address, command))
|
||||||
|
|
@ -102,7 +111,7 @@ class Client():
|
||||||
logging.debug('sending result: ' + data.decode('utf-8'))
|
logging.debug('sending result: ' + data.decode('utf-8'))
|
||||||
client.send(data.decode('utf-8'))
|
client.send(data.decode('utf-8'))
|
||||||
|
|
||||||
def _main(self):
|
def run(self):
|
||||||
try:
|
try:
|
||||||
logging.info('client connected: {}'.format(self.address))
|
logging.info('client connected: {}'.format(self.address))
|
||||||
self._set_state('CONNECT')
|
self._set_state('CONNECT')
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue