more grace
more graceful shutdown and error handling
This commit is contained in:
parent
4e46d7e810
commit
6c429b93e6
5 changed files with 86 additions and 50 deletions
|
|
@ -13,12 +13,14 @@ from tinydb import TinyDB, Query
|
|||
from tinydb.storages import MemoryStorage
|
||||
import socket
|
||||
|
||||
|
||||
def strtobool(strbool):
|
||||
if str(strbool).lower() in ['true', '1', 't', 'y', 'on','yes']:
|
||||
if str(strbool).lower() in ["true", "1", "t", "y", "on", "yes"]:
|
||||
return True
|
||||
else:
|
||||
return False
|
||||
|
||||
|
||||
# os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG
|
||||
|
||||
bumper_dir = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir))
|
||||
|
|
@ -118,6 +120,7 @@ xmppserverlog.addHandler(xmpp_rotate)
|
|||
|
||||
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
||||
|
||||
|
||||
async def start():
|
||||
|
||||
try:
|
||||
|
|
@ -173,7 +176,7 @@ async def start():
|
|||
# Start maintenance
|
||||
while not shutting_down:
|
||||
asyncio.create_task(maintenance())
|
||||
await asyncio.sleep(30)
|
||||
await asyncio.sleep(5)
|
||||
|
||||
|
||||
async def maintenance():
|
||||
|
|
@ -183,11 +186,20 @@ async def maintenance():
|
|||
async def shutdown():
|
||||
try:
|
||||
bumperlog.info("Shutting down")
|
||||
await mqtt_server.broker.shutdown()
|
||||
xmpp_server.server.close()
|
||||
await xmpp_server.server.wait_closed()
|
||||
|
||||
await conf_server.stop_server()
|
||||
await conf_server_2.stop_server()
|
||||
if mqtt_server.broker.transitions.state == "started":
|
||||
await mqtt_server.broker.shutdown()
|
||||
elif mqtt_server.broker.transitions.state == "starting":
|
||||
while mqtt_server.broker.transitions.state == "starting":
|
||||
await asyncio.sleep(0.1)
|
||||
await mqtt_server.broker.shutdown()
|
||||
await mqtt_helperbot.Client.disconnect()
|
||||
if xmpp_server.server:
|
||||
if xmpp_server.server._serving:
|
||||
xmpp_server.server.close()
|
||||
await xmpp_server.server.wait_closed()
|
||||
global shutting_down
|
||||
shutting_down = True
|
||||
|
||||
|
|
|
|||
|
|
@ -64,6 +64,7 @@ class ConfServer:
|
|||
self.run_async = False
|
||||
self.app = None
|
||||
self.site = None
|
||||
self.runner = None
|
||||
|
||||
def confserver_app(self):
|
||||
self.app = web.Application(loop=asyncio.get_event_loop())
|
||||
|
|
@ -168,40 +169,40 @@ class ConfServer:
|
|||
confserverlog.info(
|
||||
"Starting ConfServer at {}:{}".format(self.address[0], self.address[1])
|
||||
)
|
||||
runner = web.AppRunner(self.app)
|
||||
await runner.setup()
|
||||
self.runner = web.AppRunner(self.app)
|
||||
await self.runner.setup()
|
||||
|
||||
if self.usessl:
|
||||
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
|
||||
ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key)
|
||||
self.site = web.TCPSite(
|
||||
runner,
|
||||
self.runner,
|
||||
host=self.address[0],
|
||||
port=self.address[1],
|
||||
ssl_context=ssl_ctx,
|
||||
)
|
||||
|
||||
else:
|
||||
self.site = web.TCPSite(runner, host=self.address[0], port=self.address[1])
|
||||
self.site = web.TCPSite(
|
||||
self.runner, host=self.address[0], port=self.address[1]
|
||||
)
|
||||
|
||||
await self.site.start()
|
||||
|
||||
except PermissionError as e:
|
||||
if "bind" in e.strerror:
|
||||
confserverlog.exception(
|
||||
"Error binding confserver, exiting. Try using a different hostname or IP - {}".format(
|
||||
e
|
||||
)
|
||||
)
|
||||
exit(1)
|
||||
confserverlog.error(e.strerror)
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
exit(1)
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
|
||||
async def stop_server(self):
|
||||
try:
|
||||
await self.site.stop()
|
||||
await self.runner.shutdown()
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
|
|
|||
|
|
@ -58,8 +58,20 @@ class MQTTHelperBot:
|
|||
("iot/atr/+", QOS_0),
|
||||
]
|
||||
)
|
||||
|
||||
asyncio.create_task(self.get_msg())
|
||||
|
||||
except ConnectionRefusedError as e:
|
||||
helperbotlog.Error(e)
|
||||
pass
|
||||
|
||||
except asyncio.CancelledError as e:
|
||||
pass
|
||||
|
||||
except hbmqtt.client.ConnectException as e:
|
||||
helperbotlog.Error(e)
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
|
||||
|
|
@ -192,25 +204,23 @@ class MQTTServer:
|
|||
broker = None
|
||||
|
||||
async def broker_coro(self):
|
||||
try:
|
||||
|
||||
mqttserverlog.info(
|
||||
"Starting MQTT Server at {}:{}".format(self.address[0], self.address[1])
|
||||
)
|
||||
self.broker = hbmqtt.broker.Broker(config=self.default_config)
|
||||
|
||||
try:
|
||||
await self.broker.start()
|
||||
|
||||
except PermissionError as e:
|
||||
if "bind" in e.strerror:
|
||||
mqttserverlog.exception(
|
||||
"Error binding mqttserver, exiting. Try using a different hostname or IP - {}".format(
|
||||
e
|
||||
)
|
||||
)
|
||||
exit(1)
|
||||
except hbmqtt.broker.BrokerException as e:
|
||||
mqttserverlog.exception(e)
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
mqttserverlog.exception("{}".format(e))
|
||||
exit(1)
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
|
||||
def __init__(self, address):
|
||||
try:
|
||||
|
|
@ -299,9 +309,7 @@ class BumperMQTTServer_Plugin:
|
|||
)
|
||||
|
||||
mqttserverlog.info(
|
||||
"bot authenticated SN: {} DID: {}".format(
|
||||
username, didsplit[0]
|
||||
)
|
||||
"bot authenticated SN: {} DID: {}".format(username, didsplit[0])
|
||||
)
|
||||
authenticated = True
|
||||
|
||||
|
|
@ -322,9 +330,7 @@ class BumperMQTTServer_Plugin:
|
|||
|
||||
if auth:
|
||||
bumper.client_add(userid, realm, resource)
|
||||
mqttserverlog.info(
|
||||
"client authenticated {}".format(userid)
|
||||
)
|
||||
mqttserverlog.info("client authenticated {}".format(userid))
|
||||
authenticated = True
|
||||
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ class XMPPServer:
|
|||
self.xmpp_protocol = lambda: XMPPServer_Protocol()
|
||||
|
||||
async def start_async_server(self):
|
||||
try:
|
||||
xmppserverlog.info(
|
||||
"Starting XMPP Server at {}:{}".format(self.address[0], self.address[1])
|
||||
)
|
||||
|
|
@ -36,6 +37,18 @@ class XMPPServer:
|
|||
|
||||
self.server_coro = loop.create_task(self.server.serve_forever())
|
||||
|
||||
except PermissionError as e:
|
||||
xmppserverlog.error(e.strerror)
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
pass
|
||||
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
xmppserverlog.exception("{}".format(e))
|
||||
asyncio.create_task(bumper.shutdown())
|
||||
|
||||
def disconnect(self):
|
||||
try:
|
||||
xmppserverlog.debug("waiting for all client threads to exit")
|
||||
|
|
|
|||
|
|
@ -7,11 +7,11 @@ import asyncio
|
|||
if __name__ == "__main__":
|
||||
try:
|
||||
parser = argparse.ArgumentParser()
|
||||
parser.add_argument("--listen", type=str, default=None, help="listen address")
|
||||
parser.add_argument("--listen", type=str, default=None, help="start serving on address")
|
||||
parser.add_argument(
|
||||
"--announce", type=str, default=None, help="announce address (for bot)"
|
||||
"--announce", type=str, default=None, help="announce address to bots on checkin"
|
||||
)
|
||||
parser.add_argument("--debug", action="store_true")
|
||||
parser.add_argument("--debug", action="store_true", help="enable debug logs")
|
||||
args = parser.parse_args()
|
||||
|
||||
if args.debug:
|
||||
|
|
@ -29,5 +29,9 @@ if __name__ == "__main__":
|
|||
bumper.bumperlog.info("Keyboard Interrupt!")
|
||||
pass
|
||||
|
||||
except Exception as e:
|
||||
bumper.bumperlog.Exception(e)
|
||||
pass
|
||||
|
||||
finally:
|
||||
asyncio.run(bumper.shutdown())
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue