diff --git a/.gitignore b/.gitignore index 3595135..ebef324 100644 --- a/.gitignore +++ b/.gitignore @@ -5,4 +5,5 @@ __pycache__ .noseids nosetests.xml tests/report -tests/tmp.db \ No newline at end of file +tests/tmp.db +logs/ \ No newline at end of file diff --git a/Pipfile b/Pipfile index 6435d3c..b7df90c 100644 --- a/Pipfile +++ b/Pipfile @@ -15,6 +15,7 @@ nose = "*" coverage = "*" mock = "*" pylint = "*" +pbr = "*" [pipenv] allow_prereleases = true diff --git a/Pipfile.lock b/Pipfile.lock index d3afc2f..a824b67 100644 --- a/Pipfile.lock +++ b/Pipfile.lock @@ -1,7 +1,7 @@ { "_meta": { "hash": { - "sha256": "f6a2679f2c348e81b120bfa9d06f544ed25c1fcea44d94c419e6af514579fdee" + "sha256": "6e5b4ccc879dfdf82f53f7966eac54af90f9b9a2ffd1da6f1ce126d0ca125fc9" }, "pipfile-spec": 6, "requires": {}, @@ -244,90 +244,68 @@ ], "version": "==7.0" }, - "colorama": { - "hashes": [ - "sha256:05eed71e2e327246ad6b38c540c4a3117230b19679b875190486ddd2d721422d", - "sha256:f8ac84de7840f5b9c4e3347b3c1eaa50f7e49c2b07596221daec5edaabbd7c48" - ], - "markers": "sys_platform == 'win32'", - "version": "==0.4.1" - }, "coverage": { "hashes": [ - "sha256:029c69deaeeeae1b15bc6c59f0ffa28aa8473721c614a23f2c2976dec245cd12", - "sha256:02abbbebc6e9d5abe13cd28b5e963dedb6ffb51c146c916d17b18f141acd9947", - "sha256:1bbfe5b82a3921d285e999c6d256c1e16b31c554c29da62d326f86c173d30337", - "sha256:210c02f923df33a8d0e461c86fdcbbb17228ff4f6d92609fc06370a98d283c2d", - "sha256:2d0807ba935f540d20b49d5bf1c0237b90ce81e133402feda906e540003f2f7a", - "sha256:35d7a013874a7c927ce997350d314144ffc5465faf787bb4e46e6c4f381ef562", - "sha256:3636f9d0dcb01aed4180ef2e57a4e34bb4cac3ecd203c2a23db8526d86ab2fb4", - "sha256:42f4be770af2455a75e4640f033a82c62f3fb0d7a074123266e143269d7010ef", - "sha256:48440b25ba6cda72d4c638f3a9efa827b5b87b489c96ab5f4ff597d976413156", - "sha256:4dac8dfd1acf6a3ac657475dfdc66c621f291b1b7422a939cc33c13ac5356473", - "sha256:4e8474771c69c2991d5eab65764289a7dd450bbea050bc0ebb42b678d8222b42", - "sha256:551f10ddfeff56a1325e5a34eff304c5892aa981fd810babb98bfee77ee2fb17", - "sha256:5b104982f1809c1577912519eb249f17d9d7e66304ad026666cb60a5ef73309c", - "sha256:5c62aef73dfc87bfcca32cee149a1a7a602bc74bac72223236b0023543511c88", - "sha256:633151f8d1ad9467b9f7e90854a7f46ed8f2919e8bc7d98d737833e8938fc081", - "sha256:772207b9e2d5bf3f9d283b88915723e4e92d9a62c83f44ec92b9bd0cd685541b", - "sha256:7d5e02f647cd727afc2659ec14d4d1cc0508c47e6cfb07aea33d7aa9ca94d288", - "sha256:a9798a4111abb0f94584000ba2a2c74841f2cfe5f9254709756367aabbae0541", - "sha256:b38ea741ab9e35bfa7015c93c93bbd6a1623428f97a67083fc8ebd366238b91f", - "sha256:b6a5478c904236543c0347db8a05fac6fc0bd574c870e7970faa88e1d9890044", - "sha256:c6248bfc1de36a3844685a2e10ba17c18119ba6252547f921062a323fb31bff1", - "sha256:c705ab445936457359b1424ef25ccc0098b0491b26064677c39f1d14a539f056", - "sha256:d95a363d663ceee647291131dbd213af258df24f41350246842481ec3709bd33", - "sha256:e27265eb80cdc5dab55a40ef6f890e04ecc618649ad3da5265f128b141f93f78", - "sha256:ebc276c9cb5d917bd2ae959f84ffc279acafa9c9b50b0fa436ebb70bbe2166ea", - "sha256:f4d229866d030863d0fe3bf297d6d11e6133ca15bbb41ed2534a8b9a3d6bd061", - "sha256:f95675bd88b51474d4fe5165f3266f419ce754ffadfb97f10323931fa9ac95e5", - "sha256:f95bc54fb6d61b9f9ff09c4ae8ff6a3f5edc937cda3ca36fc937302a7c152bf1", - "sha256:fd0f6be53de40683584e5331c341e65a679dbe5ec489a0697cec7c2ef1a48cda" + "sha256:0402b1822d513d0231589494bceddb067d20581f5083598c451b56c684b0e5d6", + "sha256:0644e28e8aea9d9d563607ee8b7071b07dd57a4a3de11f8684cd33c51c0d1b93", + "sha256:0874a283686803884ec0665018881130604956dbaa344f2539c46d82cbe29eda", + "sha256:0988c3837df4bc371189bb3425d5232cf150055452034c232dda9cbe04f9c38e", + "sha256:20bc3205b3100956bb72293fabb97f0ed972c81fed10b3251c90c70dcb0599ab", + "sha256:2cc9142a3367e74eb6b19d58c53ebb1dfd7336b91cdcc91a6a2888bf8c7af984", + "sha256:3ae9a0a59b058ce0761c3bd2c2d66ecb2ee2b8ac592620184370577f7a546fb3", + "sha256:3b2e30b835df58cb973f478d09f3d82e90c98c8e5059acc245a8e4607e023801", + "sha256:401e9b04894eb1498c639c6623ee78a646990ce5f095248e2440968aafd6e90e", + "sha256:41ec5812d5decdaa72708be3018e7443e90def4b5a71294236a4df192cf9eab9", + "sha256:475769b638a055e75b3d3219e054fe2a023c0b077ff15bff6c95aba7e93e6cac", + "sha256:61424f4e2e82c4129a4ba71e10ebacb32a9ecd6f80de2cd05bdead6ba75ed736", + "sha256:811969904d4dd0bee7d958898be8d9d75cef672d9b7e7db819dfeac3d20d2d0c", + "sha256:86224bb99abfd672bf2f9fcecad5e8d7a3fa94f7f71513f2210460a0350307cd", + "sha256:9a238a20a3af00665f8381f7e53e9c606f9bb652d2423f6b822f6cb790d887e8", + "sha256:a23b3fbc14d4e6182ecebfd22f3729beef0636d151d94764a1c28330d185e4e5", + "sha256:ac162b4ebe51b7a2b7f5e462c4402802633eb81e77c94f8a7c1ed8a556e72c75", + "sha256:b6187378726c84365bf297b5dcdae8789b6a5823b200bea23797777e5a63be09", + "sha256:bcd5723d905ed4a825f17410a53535f880b6d7548ae3d89078db7b1ceefcd853", + "sha256:c48a4f9c5fb385269bb7fbaf9c1326a94863b65ec7f5c96b2ea56b252f01ad08", + "sha256:cd40199d6f1c29c85b170d25589be9a97edff8ee7e62be180a2a137823896030", + "sha256:d1bc331a7d069485ac1d8c25a0ea1f6aab6cb2a87146fb652222481c1bddc9ff", + "sha256:d7e0cdc249aa0f94aa2e531b03999ddaf03a10b4fa090a894712d4c8066abd89", + "sha256:e9ee8fcd8e067fcc5d7276d46e07e863102b70a52545ef4254df1ff0893ce75f", + "sha256:eb313c23d983b7810504f42104e8dcd1c7ccdda8fbaab82aab92ab79fea19345", + "sha256:f9cfd478654b509941b85ed70f870f5e3c74678f566bec12fd26545e5340ba47", + "sha256:fae1fa144034d021a52cb9ea200eb8dedf91869c6df8202ad5d149b41ed91cc8" ], "index": "pypi", - "version": "==5.0a4" + "version": "==5.0a5" }, "isort": { "hashes": [ - "sha256:18c796c2cd35eb1a1d3f012a214a542790a1aed95e29768bdcb9f2197eccbd0b", - "sha256:96151fca2c6e736503981896495d344781b60d18bfda78dc11b290c6125ebdb6" + "sha256:c40744b6bc5162bbb39c1257fe298b7a393861d50978b565f3ccd9cb9de0182a", + "sha256:f57abacd059dc3bd666258d1efb0377510a89777fda3e3274e3c01f7c03ae22d" ], - "version": "==4.3.15" + "version": "==4.3.20" }, "lazy-object-proxy": { "hashes": [ - "sha256:0ce34342b419bd8f018e6666bfef729aec3edf62345a53b537a4dcc115746a33", - "sha256:1b668120716eb7ee21d8a38815e5eb3bb8211117d9a90b0f8e21722c0758cc39", - "sha256:209615b0fe4624d79e50220ce3310ca1a9445fd8e6d3572a896e7f9146bbf019", - "sha256:27bf62cb2b1a2068d443ff7097ee33393f8483b570b475db8ebf7e1cba64f088", - "sha256:27ea6fd1c02dcc78172a82fc37fcc0992a94e4cecf53cb6d73f11749825bd98b", - "sha256:2c1b21b44ac9beb0fc848d3993924147ba45c4ebc24be19825e57aabbe74a99e", - "sha256:2df72ab12046a3496a92476020a1a0abf78b2a7db9ff4dc2036b8dd980203ae6", - "sha256:320ffd3de9699d3892048baee45ebfbbf9388a7d65d832d7e580243ade426d2b", - "sha256:50e3b9a464d5d08cc5227413db0d1c4707b6172e4d4d915c1c70e4de0bbff1f5", - "sha256:5276db7ff62bb7b52f77f1f51ed58850e315154249aceb42e7f4c611f0f847ff", - "sha256:61a6cf00dcb1a7f0c773ed4acc509cb636af2d6337a08f362413c76b2b47a8dd", - "sha256:6ae6c4cb59f199d8827c5a07546b2ab7e85d262acaccaacd49b62f53f7c456f7", - "sha256:7661d401d60d8bf15bb5da39e4dd72f5d764c5aff5a86ef52a042506e3e970ff", - "sha256:7bd527f36a605c914efca5d3d014170b2cb184723e423d26b1fb2fd9108e264d", - "sha256:7cb54db3535c8686ea12e9535eb087d32421184eacc6939ef15ef50f83a5e7e2", - "sha256:7f3a2d740291f7f2c111d86a1c4851b70fb000a6c8883a59660d95ad57b9df35", - "sha256:81304b7d8e9c824d058087dcb89144842c8e0dea6d281c031f59f0acf66963d4", - "sha256:933947e8b4fbe617a51528b09851685138b49d511af0b6c0da2539115d6d4514", - "sha256:94223d7f060301b3a8c09c9b3bc3294b56b2188e7d8179c762a1cda72c979252", - "sha256:ab3ca49afcb47058393b0122428358d2fbe0408cf99f1b58b295cfeb4ed39109", - "sha256:bd6292f565ca46dee4e737ebcc20742e3b5be2b01556dafe169f6c65d088875f", - "sha256:c0e2945ebf5b6eb32848e9bd6f850f558722caf7ae9427c9fff8e0ea7b185a2f", - "sha256:cb924aa3e4a3fb644d0c463cad5bc2572649a6a3f68a7f8e4fbe44aaa6d77e4c", - "sha256:d0fc7a286feac9077ec52a927fc9fe8fe2fabab95426722be4c953c9a8bede92", - "sha256:ddc34786490a6e4ec0a855d401034cbd1242ef186c20d79d2166d6a4bd449577", - "sha256:e34b155e36fa9da7e1b7c738ed7767fc9491a62ec6af70fe9da4a057759edc2d", - "sha256:e5b9e8f6bda48460b7b143c3821b21b452cb3a835e6bbd5dd33aa0c8d3f5137d", - "sha256:e81ebf6c5ee9684be8f2c87563880f93eedd56dd2b6146d8a725b50b7e5adb0f", - "sha256:eb91be369f945f10d3a49f5f9be8b3d0b93a4c2be8f8a5b83b0571b8123e0a7a", - "sha256:f460d1ceb0e4a5dcb2a652db0904224f367c9b3c1470d5a7683c0480e582468b" + "sha256:159a745e61422217881c4de71f9eafd9d703b93af95618635849fe469a283661", + "sha256:23f63c0821cc96a23332e45dfaa83266feff8adc72b9bcaef86c202af765244f", + "sha256:3b11be575475db2e8a6e11215f5aa95b9ec14de658628776e10d96fa0b4dac13", + "sha256:3f447aff8bc61ca8b42b73304f6a44fa0d915487de144652816f950a3f1ab821", + "sha256:4ba73f6089cd9b9478bc0a4fa807b47dbdb8fad1d8f31a0f0a5dbf26a4527a71", + "sha256:4f53eadd9932055eac465bd3ca1bd610e4d7141e1278012bd1f28646aebc1d0e", + "sha256:64483bd7154580158ea90de5b8e5e6fc29a16a9b4db24f10193f0c1ae3f9d1ea", + "sha256:6f72d42b0d04bfee2397aa1862262654b56922c20a9bb66bb76b6f0e5e4f9229", + "sha256:7c7f1ec07b227bdc561299fa2328e85000f90179a2f44ea30579d38e037cb3d4", + "sha256:7c8b1ba1e15c10b13cad4171cfa77f5bb5ec2580abc5a353907780805ebe158e", + "sha256:8559b94b823f85342e10d3d9ca4ba5478168e1ac5658a8a2f18c991ba9c52c20", + "sha256:a262c7dfb046f00e12a2bdd1bafaed2408114a89ac414b0af8755c696eb3fc16", + "sha256:acce4e3267610c4fdb6632b3886fe3f2f7dd641158a843cf6b6a68e4ce81477b", + "sha256:be089bb6b83fac7f29d357b2dc4cf2b8eb8d98fe9d9ff89f9ea6012970a853c7", + "sha256:bfab710d859c779f273cc48fb86af38d6e9210f38287df0069a63e40b45a2f5c", + "sha256:c10d29019927301d524a22ced72706380de7cfc50f767217485a912b4c8bd82a", + "sha256:dd6e2b598849b3d7aee2295ac765a578879830fb8966f70be8cd472e6069932e", + "sha256:e408f1eacc0a68fed0c08da45f31d0ebb38079f043328dce69ff133b95c29dc1" ], - "version": "==1.3.1" + "version": "==1.4.1" }, "mccabe": { "hashes": [ @@ -338,11 +316,11 @@ }, "mock": { "hashes": [ - "sha256:5ce3c71c5545b472da17b72268978914d0252980348636840bd34a00b5cc96c1", - "sha256:b158b6df76edd239b8208d481dc46b6afd45a846b7812ff0ce58971cf5bc8bba" + "sha256:83657d894c90d5681d62155c82bda9c1187827525880eda8ff5df4ec813437c3", + "sha256:d157e52d4e5b938c550f39eb2fd15610db062441a9c2747d3dbfa9298211d0f8" ], "index": "pypi", - "version": "==2.0.0" + "version": "==3.0.5" }, "nose": { "hashes": [ @@ -355,10 +333,11 @@ }, "pbr": { "hashes": [ - "sha256:8257baf496c8522437e8a6cfe0f15e00aedc6c0e0e7c9d55eeeeab31e0853843", - "sha256:8c361cc353d988e4f5b998555c88098b9d5964c2e11acf7b0d21925a66bb5824" + "sha256:6901995b9b686cb90cceba67a0f6d4d14ae003cd59bc12beb61549bdfbe3bc89", + "sha256:d950c64aeea5456bbd147468382a5bb77fe692c13c9f00f0219814ce5b642755" ], - "version": "==5.1.3" + "index": "pypi", + "version": "==5.2.0" }, "pylint": { "hashes": [ @@ -384,33 +363,32 @@ }, "typed-ast": { "hashes": [ - "sha256:035a54ede6ce1380599b2ce57844c6554666522e376bd111eb940fbc7c3dad23", - "sha256:037c35f2741ce3a9ac0d55abfcd119133cbd821fffa4461397718287092d9d15", - "sha256:049feae7e9f180b64efacbdc36b3af64a00393a47be22fa9cb6794e68d4e73d3", - "sha256:19228f7940beafc1ba21a6e8e070e0b0bfd1457902a3a81709762b8b9039b88d", - "sha256:2ea681e91e3550a30c2265d2916f40a5f5d89b59469a20f3bad7d07adee0f7a6", - "sha256:3a6b0a78af298d82323660df5497bcea0f0a4a25a0b003afd0ce5af049bd1f60", - "sha256:5385da8f3b801014504df0852bf83524599df890387a3c2b17b7caa3d78b1773", - "sha256:606d8afa07eef77280c2bf84335e24390055b478392e1975f96286d99d0cb424", - "sha256:69245b5b23bbf7fb242c9f8f08493e9ecd7711f063259aefffaeb90595d62287", - "sha256:6f6d839ab09830d59b7fa8fb6917023d8cb5498ee1f1dbd82d37db78eb76bc99", - "sha256:730888475f5ac0e37c1de4bd05eeb799fdb742697867f524dc8a4cd74bcecc23", - "sha256:9819b5162ffc121b9e334923c685b0d0826154e41dfe70b2ede2ce29034c71d8", - "sha256:9e60ef9426efab601dd9aa120e4ff560f4461cf8442e9c0a2b92548d52800699", - "sha256:af5fbdde0690c7da68e841d7fc2632345d570768ea7406a9434446d7b33b0ee1", - "sha256:b64efdbdf3bbb1377562c179f167f3bf301251411eb5ac77dec6b7d32bcda463", - "sha256:bac5f444c118aeb456fac1b0b5d14c6a71ea2a42069b09c176f75e9bd4c186f6", - "sha256:bda9068aafb73859491e13b99b682bd299c1b5fd50644d697533775828a28ee0", - "sha256:d659517ca116e6750101a1326107d3479028c5191f0ecee3c7203c50f5b915b0", - "sha256:eddd3fb1f3e0f82e5915a899285a39ee34ce18fd25d89582bc89fc9fb16cd2c6" + "sha256:132eae51d6ef3ff4a8c47c393a4ef5ebf0d1aecc96880eb5d6c8ceab7017cc9b", + "sha256:18141c1484ab8784006c839be8b985cfc82a2e9725837b0ecfa0203f71c4e39d", + "sha256:2baf617f5bbbfe73fd8846463f5aeafc912b5ee247f410700245d68525ec584a", + "sha256:3d90063f2cbbe39177e9b4d888e45777012652d6110156845b828908c51ae462", + "sha256:4304b2218b842d610aa1a1d87e1dc9559597969acc62ce717ee4dfeaa44d7eee", + "sha256:4983ede548ffc3541bae49a82675996497348e55bafd1554dc4e4a5d6eda541a", + "sha256:5315f4509c1476718a4825f45a203b82d7fdf2a6f5f0c8f166435975b1c9f7d4", + "sha256:6cdfb1b49d5345f7c2b90d638822d16ba62dc82f7616e9b4caa10b72f3f16649", + "sha256:7b325f12635598c604690efd7a0197d0b94b7d7778498e76e0710cd582fd1c7a", + "sha256:8d3b0e3b8626615826f9a626548057c5275a9733512b137984a68ba1598d3d2f", + "sha256:8f8631160c79f53081bd23446525db0bc4c5616f78d04021e6e434b286493fd7", + "sha256:912de10965f3dc89da23936f1cc4ed60764f712e5fa603a09dd904f88c996760", + "sha256:b010c07b975fe853c65d7bbe9d4ac62f1c69086750a574f6292597763781ba18", + "sha256:c908c10505904c48081a5415a1e295d8403e353e0c14c42b6d67f8f97fae6616", + "sha256:c94dd3807c0c0610f7c76f078119f4ea48235a953512752b9175f9f98f5ae2bd", + "sha256:ce65dee7594a84c466e79d7fb7d3303e7295d16a83c22c7c4037071b059e2c21", + "sha256:eaa9cfcb221a8a4c2889be6f93da141ac777eb8819f077e1d09fb12d00a09a93", + "sha256:f3376bc31bad66d46d44b4e6522c5c21976bf9bca4ef5987bb2bf727f4506cbb", + "sha256:f9202fa138544e13a4ec1a6792c35834250a85958fde1251b6a22e07d1260ae7" ], "markers": "implementation_name == 'cpython'", - "version": "==1.3.1" + "version": "==1.3.5" }, "wrapt": { "hashes": [ - "sha256:4aea003270831cceb8a90ff27c4031da6ead7ec1886023b80ce0dfe0adf61533", - "sha256:71ad0a3729a2f29eb25ff0dd6e66fc9768c749ff0c3f8899a6ac35c059f502c1" + "sha256:4aea003270831cceb8a90ff27c4031da6ead7ec1886023b80ce0dfe0adf61533" ], "version": "==1.11.1" } diff --git a/bumper/__init__.py b/bumper/__init__.py index 8501819..e2ffaaa 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -1,25 +1,21 @@ #!/usr/bin/env python3 -from .confserver import ConfServer -from .mqttserver import MQTTServer -from .mqttserver import MQTTHelperBot -from .xmppserver import XMPPServer +from bumper.confserver import ConfServer +from bumper.mqttserver import MQTTServer, MQTTHelperBot +from bumper.xmppserver import XMPPServer import asyncio -import contextvars import json import time from datetime import datetime, timedelta import platform -import os +import os, sys import logging +from logging.handlers import RotatingFileHandler from base64 import b64decode, b64encode from tinydb import TinyDB, Query +import json from tinydb.storages import MemoryStorage -bumper_users_var = contextvars.ContextVar("bumper_users", default=[]) -bumper_clients_var = contextvars.ContextVar("bumper_clients", default=[]) -bumper_bots_var = contextvars.ContextVar("bumper_bots", default=[]) - ca_cert = "./certs/CA/cacert.pem" server_cert = "./certs/cert.pem" server_key = "./certs/key.pem" @@ -29,21 +25,49 @@ token_validity_seconds = 3600 # 1 hour db = None # Logs +os.makedirs("logs", exist_ok=True) #Ensure logs directory exists or create +# Set format for all logs +logformat = logging.Formatter("[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s") + bumperlog = logging.getLogger("bumper") +bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5) +bumper_rotate.setFormatter(logformat) +bumperlog.addHandler(bumper_rotate) +# Override the logging level +# bumperlog.setLevel(logging.INFO) + confserverlog = logging.getLogger("confserver") +conf_rotate = RotatingFileHandler("logs/confserver.log", maxBytes=5000000, backupCount=5) +conf_rotate.setFormatter(logformat) +confserverlog.addHandler(conf_rotate) # Override the logging level # confserverlog.setLevel(logging.INFO) + mqttserverlog = logging.getLogger("mqttserver") +mqtt_rotate = RotatingFileHandler("logs/mqttserver.log", maxBytes=5000000, backupCount=5) +mqtt_rotate.setFormatter(logformat) +mqttserverlog.addHandler(mqtt_rotate) # Override the logging level # mqttserverlog.setLevel(logging.INFO) + helperbotlog = logging.getLogger("helperbot") +helperbot_rotate = RotatingFileHandler("logs/helperbot.log", maxBytes=5000000, backupCount=5) +helperbot_rotate.setFormatter(logformat) +helperbotlog.addHandler(helperbot_rotate) # Override the logging level # helperbotlog.setLevel(logging.INFO) + xmppserverlog = logging.getLogger("xmppserver") +xmpp_rotate = RotatingFileHandler("logs/xmppserver.log", maxBytes=5000000, backupCount=5) +xmpp_rotate.setFormatter(logformat) +xmppserverlog.addHandler(xmpp_rotate) # Override the logging level # xmppserverlog.setLevel(logging.INFO) +logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger + + def get_milli_time(timetoconvert): return int(round(timetoconvert * 1000)) @@ -56,23 +80,32 @@ def db_file(): def os_db_path(): - if platform.system() == "Windows": - return os.path.join(os.getenv("APPDATA"), "bumper.db") + if platform.system() == "Windows": + os.makedirs(os.getenv("APPDATA"), exist_ok=True) #Ensure db_path directory exists or create + return os.path.join(os.getenv("APPDATA"), "bumper.db") else: + os.makedirs(os.path.expanduser("~/.config"), exist_ok=True) #Ensure db_path directory exists or create return os.path.expanduser("~/.config/bumper.db") - def db_get(): - # Will create the database if it doesn't exist - db = TinyDB(db_file()) + try: + # Will create the database if it doesn't exist + db = TinyDB(db_file()) - # Will create the tables if they don't exist - db.table("users", cache_size=0) - db.table("clients", cache_size=0) - db.table("bots", cache_size=0) - db.table("tokens", cache_size=0) + # Will create the tables if they don't exist + db.table("users", cache_size=0) + db.table("clients", cache_size=0) + db.table("bots", cache_size=0) + db.table("tokens", cache_size=0) - return db + return db + + + except json.decoder.JSONDecodeError as jerr: + bumperlog.error("JsonErr: {} - Doc: {}".format(jerr.msg, jerr.doc)) + + except Exception as ex: + bumperlog.error(ex) class BumperUser(object): @@ -632,11 +665,12 @@ def bot_add(sn, did, devclass, resource, company): newbot.company = company bot = bot_get(did) - if not bot: - bumperlog.info( - "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did) - ) - bot_full_upsert(newbot.asdict()) + if not bot: # Not existing bot in database + if not devclass == "" or "@" not in sn or "tmp" not in sn: # try to prevent bad additions to the bot list + bumperlog.info( + "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did) + ) + bot_full_upsert(newbot.asdict()) def bot_remove(did): diff --git a/bumper/confserver.py b/bumper/confserver.py index 7af5fa6..d18cf6f 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -8,7 +8,6 @@ import bumper import time from datetime import datetime, timedelta import asyncio -import contextvars from aiohttp import web import uuid import xml.etree.ElementTree as ET @@ -32,9 +31,7 @@ class aiohttp_filter(logging.Filter): confserverlog = logging.getLogger("confserver") - -logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger -logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) +logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) #Add logging filter above to aiohttp.access class EcoVacs_Login: accessToken = "" @@ -61,42 +58,6 @@ class ConfServer: self.run_async = False self.app = None - def run(self, run_async=False): - try: - if run_async: - self.run_async = True - confserverlog.debug("Starting ConfServer Thread: 1") - self.confthread = Thread( - name="ConfServer_{}_Thread".format(self.address[1]), - target=self.run_server, - ) - self.confthread.setDaemon(True) - self.confthread.start() - - else: - try: - self.run_server() - except KeyboardInterrupt: - self.disconnect() - - except Exception as e: - confserverlog.exception("{}".format(e)) - - def run_server(self): - logging.info("Starting ConfServer at {}".format(self.address)) - print("Starting ConfServer at {}".format(self.address)) - try: - loop = asyncio.get_event_loop() - except: - loop = asyncio.new_event_loop() - - try: - self.confserver_app() - loop.run_until_complete(self.start_server()) - loop.run_forever() - except Exception as e: - confserverlog.exception("{}".format(e)) - def confserver_app(self): self.app = web.Application() @@ -197,6 +158,7 @@ class ConfServer: async def start_server(self): try: + confserverlog.info("Starting ConfServer at {}:{}".format(self.address[0], self.address[1])) runner = web.AppRunner(self.app) await runner.setup() @@ -352,55 +314,67 @@ class ConfServer: def check_token(self, apptype, countrycode, user, token): - if bumper.check_token(user["userid"], token): - - if "global_" in apptype: #EcoVacs Home - login_details = EcoVacsHome_Login() - login_details.ucUid = "fuid_{}".format(user["userid"]) - login_details.loginName = "fusername_{}".format(user["userid"]) - login_details.mobile = None + try: + if bumper.check_token(user["userid"], token): + + if "global_" in apptype: #EcoVacs Home + login_details = EcoVacsHome_Login() + login_details.ucUid = "fuid_{}".format(user["userid"]) + login_details.loginName = "fusername_{}".format(user["userid"]) + login_details.mobile = None + else: + login_details = EcoVacs_Login() + + login_details.accessToken = token + login_details.uid = "fuid_{}".format(user["userid"]) + login_details.username = "fusername_{}".format(user["userid"]) + login_details.country = countrycode + login_details.email = "null@null.com" + + body = { + "code": bumper.RETURN_API_SUCCESS, + "data": json.loads(login_details.toJSON()), + #{ + # "accessToken": self.generate_token(tmpuser), # Generate a token + # "country": countrycode, + # "email": "null@null.com", + # "uid": "fuid_{}".format(tmpuser["userid"]), + # "username": "fusername_{}".format(tmpuser["userid"]), + #}, + "msg": "操作成功", + "time": bumper.get_milli_time(datetime.utcnow().timestamp()), + } + return web.json_response(body) + else: - login_details = EcoVacs_Login() - - login_details.accessToken = token - login_details.uid = "fuid_{}".format(user["userid"]) - login_details.username = "fusername_{}".format(user["userid"]) - login_details.country = countrycode - login_details.email = "null@null.com" - - body = { - "code": bumper.RETURN_API_SUCCESS, - "data": json.loads(login_details.toJSON()), - #{ - # "accessToken": self.generate_token(tmpuser), # Generate a token - # "country": countrycode, - # "email": "null@null.com", - # "uid": "fuid_{}".format(tmpuser["userid"]), - # "username": "fusername_{}".format(tmpuser["userid"]), - #}, - "msg": "操作成功", - "time": bumper.get_milli_time(datetime.utcnow().timestamp()), - } - return web.json_response(body) - - else: - body = { - "code": bumper.ERR_TOKEN_INVALID, - "data": None, - "msg": "当前密码错误", - "time": bumper.get_milli_time(datetime.utcnow().timestamp()), - } - return web.json_response(body) + body = { + "code": bumper.ERR_TOKEN_INVALID, + "data": None, + "msg": "当前密码错误", + "time": bumper.get_milli_time(datetime.utcnow().timestamp()), + } + return web.json_response(body) + + except Exception as e: + confserverlog.exception("{}".format(e)) def generate_token(self, user): - tmpaccesstoken = uuid.uuid4().hex - bumper.user_add_token(user["userid"], tmpaccesstoken) - return tmpaccesstoken + try: + tmpaccesstoken = uuid.uuid4().hex + bumper.user_add_token(user["userid"], tmpaccesstoken) + return tmpaccesstoken + + except Exception as e: + confserverlog.exception("{}".format(e)) def generate_authcode(self, user, countrycode, token): - tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex) - bumper.user_add_authcode(user["userid"], token, tmpauthcode) - return tmpauthcode + try: + tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex) + bumper.user_add_authcode(user["userid"], token, tmpauthcode) + return tmpauthcode + + except Exception as e: + confserverlog.exception("{}".format(e)) def _auth_any(self, devid, apptype, country, request): try: @@ -816,7 +790,6 @@ class ConfServer: confserverlog.exception("{}".format(e)) async def handle_getProductIotMap(self, request): - user_devid = request.match_info.get("devid", "") try: body = { "code": bumper.RETURN_API_SUCCESS, @@ -910,13 +883,23 @@ class ConfServer: if todo == "FindBest": service = postbody["service"] if service == "EcoMsgNew": + srvip = socket.gethostbyname(socket.gethostname()) + srvport = 5223 + confserverlog.info( + "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport) + ) body = { "result": "ok", - "ip": socket.gethostbyname(socket.gethostname()), - "port": 5223, + "ip": srvip, + "port": srvport, } elif service == "EcoUpdate": - body = {"result": "ok", "ip": "47.88.66.164", "port": 8005} + srvip = "47.88.66.164" #EcoVacs Server + srvport = 8005 + confserverlog.info( + "Reporting FindBest-EcoUpdate Server to Bot as: {}:{}".format(srvip, srvport) + ) + body = {"result": "ok", "ip": srvip, "port": srvport} elif todo == "loginByItToken": if "userId" in postbody: @@ -996,7 +979,8 @@ class ConfServer: for bot in bots: if bot["class"] != "": b = bumper.bot_toEcoVacsHome_JSON(bot) - botlist.append(json.loads(b)) + if not b is None: #Happens if the bot isn't on the EcoVacs Home list + botlist.append(json.loads(b)) body = { "code": 0, @@ -1036,9 +1020,12 @@ class ConfServer: if todo == "FindBest": service = postbody["service"] if service == "EcoMsgNew": - srvip = socket.gethostbyname(socket.gethostname()) - msgserver = {"ip": srvip, "port": 5223, "result": "ok"} + srvport = 5223 + confserverlog.info( + "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport) + ) + msgserver = {"ip": srvip, "port": srvport, "result": "ok"} msgserver = json.dumps(msgserver) msgserver = msgserver.replace( " ", "" @@ -1095,11 +1082,16 @@ class ConfServer: if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) bot = bumper.bot_get(did) if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: body = "" retcmd = await self.helperbot.send_command(json_body, randomid) + confserverlog.debug( + "Send Bot - {}".format(json_body) + ) + confserverlog.debug( + "Bot Response - {}".format(body) + ) logs = [] logsroot = ET.fromstring(retcmd["resp"]) if logsroot.attrib["ret"] == "ok": @@ -1148,13 +1140,15 @@ class ConfServer: did = json_body["toId"] if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) bot = bumper.bot_get(did) if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: retcmd = await self.helperbot.send_command(json_body, randomid) body = retcmd confserverlog.debug( - "\r\n POST: {} \r\n Response: {}".format(json_body, body) + "Send Bot - {}".format(json_body) + ) + confserverlog.debug( + "Bot Response - {}".format(body) ) return web.json_response(body) else: @@ -1181,7 +1175,7 @@ class ConfServer: confserverlog.exception("{}".format(e)) - async def handle_dim_devmanager(self, request): + async def handle_dim_devmanager(self, request): #Used in EcoVacs Home App try: json_body = json.loads(await request.text()) @@ -1191,13 +1185,15 @@ class ConfServer: did = json_body["toId"] if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) bot = bumper.bot_get(did) if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: retcmd = await self.helperbot.send_command(json_body, randomid) body = retcmd confserverlog.debug( - "\r\n POST: {} \r\n Response: {}".format(json_body, body) + "Send Bot - {}".format(json_body) + ) + confserverlog.debug( + "Bot Response - {}".format(body) ) return web.json_response(body) else: @@ -1223,13 +1219,13 @@ class ConfServer: except Exception as e: confserverlog.exception("{}".format(e)) - def disconnect(self): + async def disconnect(self): try: confserverlog.info("shutting down") if self.run_async: self.confthread.join() else: - self.app.shutdown() + await self.app.shutdown() except Exception as e: confserverlog.exception("{}".format(e)) diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index 3749f0c..a9c3ca1 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -8,13 +8,13 @@ from hbmqtt.broker import Broker from hbmqtt.client import MQTTClient from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2 import pkg_resources -import contextvars import time from threading import Thread import ssl import bumper import json from datetime import datetime, timedelta +import bumper helperbotlog = logging.getLogger("helperbot") mqttserverlog = logging.getLogger("mqttserver") @@ -40,39 +40,16 @@ class MQTTHelperBot: ): self.address = address self.client_id = "helper1@bumper/helper1" - self.command_responses = contextvars.ContextVar("command_responses", default=[]) + self.command_responses = [] self.helperthread = None - def run(self, run_async=False): - if run_async: - hloop = asyncio.new_event_loop() - helperbotlog.debug("Starting MQTT HelperBot Thread: 1") - self.helperthread = Thread( - name="MQTTHelperBot_Thread", target=self.run_helperbot, args=(hloop,) - ) - self.helperthread.setDaemon(True) - self.helperthread.start() - - else: - self.run_helperbot(asyncio.get_event_loop()) - - def run_helperbot(self, loop): - logging.info("Starting MQTT HelperBot") - print("Starting MQTT HelperBot") - try: - asyncio.set_event_loop(loop) - self.Client = MQTTClient( - client_id=self.client_id, config={"check_hostname": False} - ) - loop.run_until_complete(self.start_helper_bot()) - loop.run_until_complete(self.get_msg()) - loop.run_forever() - except Exception as e: - helperbotlog.exception("{}".format(e)) - async def start_helper_bot(self): try: + self.Client = MQTTClient( + client_id=self.client_id, config={"check_hostname": False} + ) + await self.Client.connect( "mqtts://{}:{}/".format(self.address[0], self.address[1]), cafile=bumper.ca_cert, @@ -81,8 +58,10 @@ class MQTTHelperBot: [ ("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0), ("iot/p2p/+", QOS_0), + ("iot/atr/+", QOS_0), ] ) + asyncio.create_task(self.get_msg()) except Exception as e: helperbotlog.exception("{}".format(e)) @@ -92,29 +71,33 @@ class MQTTHelperBot: while True: message = await self.Client.deliver_message() - # helperbotlog.debug("HelperBot MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8")))) - cresp = self.command_responses.get() - if str(message.topic).split("/")[6] == "helper1": - cresp.append( + #Response to command + helperbotlog.debug("Received Response - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8")))) + self.command_responses.append( { "time": time.time(), "topic": message.topic, "payload": str(message.data.decode("utf-8")), } ) + elif str(message.topic).split("/")[3] == "helper1": + #Helperbot sending command + helperbotlog.debug("Send Command - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8")))) + elif str(message.topic).split("/")[1] == "atr": + #Broadcast message received on atr + helperbotlog.debug("Received Broadcast - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8")))) + else: + helperbotlog.debug("Received Message - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8")))) # Cleanup "expired messages" > 60 seconds from time - for msg in cresp: + for msg in self.command_responses: expire_time = ( datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10) ).timestamp() if time.time() > expire_time: - # helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time)) - cresp.remove(msg) - - self.command_responses.set(cresp) - # helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp)) + helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time)) + self.command_responses.remove(msg) except Exception as e: helperbotlog.exception("{}".format(e)) @@ -125,10 +108,9 @@ class MQTTHelperBot: t_end = (datetime.now() + timedelta(seconds=10)).timestamp() while time.time() < t_end: - await asyncio.sleep(0.1) - responses = self.command_responses.get() - if len(responses) > 0: - for msg in responses: + await asyncio.sleep(0.1) + if len(self.command_responses) > 0: + for msg in self.command_responses: topic = str(msg["topic"]).split("/") if topic[6] == "helper1" and topic[10] == requestid: # helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload'])) @@ -136,10 +118,8 @@ class MQTTHelperBot: resppayload = json.loads(msg["payload"]) else: resppayload = str(msg["payload"]) - resp = {"id": requestid, "ret": "ok", "resp": resppayload} - cresp = self.command_responses.get() - cresp.remove(msg) - self.command_responses.set(cresp) + resp = {"id": requestid, "ret": "ok", "resp": resppayload} + self.command_responses.remove(msg) return resp return {"id": requestid, "errno": 500, "ret": "fail", "debug": "wait for response timed out"} @@ -173,6 +153,7 @@ class MQTTHelperBot: except Exception as e: helperbotlog.exception("{}".format(e)) + return {} class MQTTServer: @@ -180,6 +161,7 @@ class MQTTServer: async def broker_coro(self): try: + mqttserverlog.info("Starting MQTT Server at {}:{}".format(self.address[0], self.address[1])) broker = hbmqtt.broker.Broker(config=self.default_config) await broker.start() @@ -238,32 +220,6 @@ class MQTTServer: except Exception as e: mqttserverlog.exception("{}".format(e)) - def run(self, run_async=False): - if run_async: - sloop = asyncio.new_event_loop() - mqttserverlog.debug("Starting MQTTServer Thread: 1") - self.mqttserverthread = Thread( - name="MQTTServer_Thread", target=self.run_server, args=(sloop,) - ) - self.mqttserverthread.setDaemon(True) - self.mqttserverthread.start() - - else: - self.run_server(asyncio.get_event_loop()) - - def run_server(self, loop): - - logging.info("Starting MQTT Server at {}".format(self.address)) - print("Starting MQTT Server at {}".format(self.address)) - try: - asyncio.set_event_loop(loop) - loop.run_until_complete(self.broker_coro()) - # loop.run_until_complete(self.active_bot_listing()) - loop.run_forever() - - except Exception as e: - mqttserverlog.exception("{}".format(e)) - class BumperMQTTServer_Plugin: def __init__(self, context): @@ -300,11 +256,9 @@ class BumperMQTTServer_Plugin: client_id = session.client_id didsplit = str(client_id).split("@") - # If this isn't a fake user (fuid) then add as a bot - if not ( - str(didsplit[0]).startswith("fuid") - or str(didsplit[0]).startswith("helper") - ): + if not ( # if ecouser or bumper aren't in details it is a bot + "ecouser" in didsplit[1] + or "bumper" in didsplit[1]): tmpbotdetail = str(didsplit[1]).split("/") bumper.bot_add( username, diff --git a/bumper/xmpp_old_client.py b/bumper/xmpp_old_client.py new file mode 100644 index 0000000..ea5b24b --- /dev/null +++ b/bumper/xmpp_old_client.py @@ -0,0 +1,769 @@ +class Client(threading.Thread): + IDLE = 0 + CONNECT = 1 + INIT = 2 + BIND = 3 + READY = 4 + DISCONNECT = 5 + UNKNOWN = 0 + BOT = 1 + CONTROLLER = 2 + + def __init__(self, thread_id, connection, client_address): + threading.Thread.__init__(self) + self.id = thread_id + self.name = "XMPP_Client_{}".format(client_address[0]) + self.type = self.UNKNOWN + self.state = self.IDLE + self.connection = connection + self.address = client_address[0] + self.clientresource = "" + self.devclass = "" + self.bumper_jid = "" + self.uid = "" + self.log_sent_message = False # Set to true to log sends + self.log_incoming_data = True # Set to true to log sends + + xmppserverlog.debug( + "new client thread init for client with ip {}".format(self.address) + ) + + def send(self, command): + try: + if not self.connection._closed: + if self.log_sent_message: + xmppserverlog.debug("send {} - {}".format(self.address, command)) + self.connection.send(command.encode()) + + except BrokenPipeError as e: + xmppserverlog.debug("{}".format(e)) + self._set_state("DISCONNECT") + + except ConnectionResetError as e: + xmppserverlog.debug("{}".format(e)) + self._set_state("DISCONNECT") + + except ConnectionAbortedError as e: + xmppserverlog.debug("{}".format(e)) + self._set_state("DISCONNECT") + + except OSError as e: + xmppserverlog.debug("{}".format(e)) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _disconnect(self): + try: + + bot = bumper.bot_get(self.uid) + if bot: + bumper.bot_set_xmpp(bot["did"], False) + + client = bumper.client_get(self.clientresource) + if client: + bumper.client_set_xmpp(client["resource"], False) + + self.connection.close() + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _tag_strip_uri(self, tag): + try: + if tag[0] == "{": + _, _, tag = tag[1:].partition("}") + return tag + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _set_state(self, state): + try: + new_state = getattr(Client, state) + if self.state > new_state: + raise Exception( + "{} illegal state change {}->{}".format( + self.address, self.state, new_state + ) + ) + + xmppserverlog.debug("{} state: {}".format(self.address, state)) + + self.state = new_state + + if new_state == 5: + self._disconnect() + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_ctl(self, xml, data): + try: + + if "roster" in data: + # Return not-implemented for roster + self.send( + ''.format( + xml.get("id") + ) + ) + return + + if xml.get("type") == "set": + if ( + "com:sf" in data and xml.get("to") == "rl.ecorobot.net" + ): # Android bind? Not sure what this does yet. + self.send( + ''.format( + xml.get("id"), + self.uid, + XMPPServer.server_id, + self.clientresource, + ) + ) + + if xml[0][0]: + ctl = xml[0][0] + if ctl.get("admin") and self.type == self.BOT: + xmppserverlog.debug( + "admin username received from bot: {}".format(ctl.get("admin")) + ) + XMPPServer.client_id = ctl.get("admin") + return + + # forward + for client in XMPPServer.clients: + if ( + client.bumper_jid != self.bumper_jid + and client.state == client.READY + ): + ctl_to = xml.get("to") + xml.attrib["from"] = "{}".format(self.bumper_jid) + rxmlstring = ET.tostring(xml).decode("utf-8") + # clean up string to remove namespaces added by ET + rxmlstring = rxmlstring.replace("xmlns:ns0=", "xmlns=") + rxmlstring = rxmlstring.replace("ns0:", "") + rxmlstring = rxmlstring.replace('iq xmlns="com:ctl"', "iq") + rxmlstring = rxmlstring.replace("'.format( + uuid.uuid4(), adminuser, self.bumper_jid, newuser + ) + xmppserverlog.debug("Add User: {}".format(adduser)) + self.send(adduser) + + # Add user ACs - Manage users, settings, and clean (full access) + adduseracs = ''.format( + uuid.uuid4(), adminuser, self.bumper_jid, newuser + ) + xmppserverlog.debug("Add User ACs: {}".format(adduseracs)) + self.send(adduseracs) + + # GetUserInfo - Just to confirm it set correctly + self.send( + ''.format( + uuid.uuid4(), adminuser, self.bumper_jid + ) + ) + + else: + rxmlstring = ET.tostring(xml).decode("utf-8") + # clean up string to remove namespaces added by ET + rxmlstring = rxmlstring.replace("xmlns:ns0=", "xmlns=") + rxmlstring = rxmlstring.replace("ns0:", "") + rxmlstring = rxmlstring.replace('iq xmlns="com:ctl"', "iq") + rxmlstring = rxmlstring.replace(" -1: + sc = data.decode("utf-8").find("to=") + ec = data.decode("utf-8").find(".ecorobot.net") + if ec > -1: + self.devclass = data.decode("utf-8")[sc + 4 : ec] + # ack jabbr:client + # no STARTTLS + self.send( + ''.format( + XMPPServer.server_id + ) + ) + # with STARTTLS + # self.send(''.format(XMPPServer.server_id)) + time.sleep(0.25) + # send authentication support for iq-auth (fallback) and SASL + self.send( + 'PLAIN' + ) + # self.send('') + + else: + self.send("") + + else: + if "jabber:iq:auth" in xml.tag: # Handle iq-auth + self._handle_iq_auth(xml) + elif ( + "urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag + ): # Handle SASL Auth + self._handle_sasl_auth(xml) + else: + xmppserverlog.error("Couldn't handle: {}".format(xml)) + + elif self.state == self.INIT: + if xml == None: + # Client getting session after authentication + if data.decode("utf-8").find("jabber:client") > -1: + # ack jabbr:client + self.send( + ''.format( + XMPPServer.server_id + ) + ) + time.sleep(0.25) + # session + self.send( + '' + ) + + else: # Handle init bind + if len(xml): + child = self._tag_strip_uri(xml[0].tag) + else: + child = None + + if xml.tag == "iq": + if child == "bind": + self._handle_bind(xml) + else: + xmppserverlog.error("Couldn't handle: {}".format(xml)) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_iq_auth(self, data): + try: + xml = ET.fromstring(data.decode("utf-8")) + ctl = xml[0][0] + xmppserverlog.info("IQ AUTH XML: {}".format(xml)) + # Received username and auth tag, send username/password requirement + if ( + xml.get("type") == "get" + and "auth}username" in ctl.tag + and self.type == self.UNKNOWN + ): + self.send( + ''.format( + xml.get("id") + ) + ) + + # Received username, password, resource - Handle auth here and return pass or fail + if ( + xml.get("type") == "set" + and "auth}username" in ctl.tag + and self.type == self.UNKNOWN + ): + xmlauth = xml[0].getchildren() + # uid = "" + password = "" + resource = "" + for aitem in xmlauth: + if "username" in aitem.tag: + self.uid = aitem.text + + elif "password" in aitem.tag: + password = aitem.text.split("/")[2] + authcode = password + + elif "resource" in aitem.tag: + self.clientresource = aitem.text + resource = self.clientresource + + if not self.uid.startswith("fuid"): + + # Need sample data to see details here + bumper.bot_add("", self.uid, "", resource, "eco-legacy") + xmppserverlog.info("bot authenticated {}".format(self.uid)) + + # Client authenticated, move to next state + self._set_state("INIT") + + # Successful auth + self.send(''.format(xml.get("id"))) + + else: + auth = False + if bumper.check_authcode(self.uid, authcode): + auth = True + elif bumper.use_auth == False: + auth = True + + if auth: + bumper.client_add(self.uid, "bumper", self.clientresource) + xmppserverlog.debug("client authenticated {}".format(self.uid)) + + # Client authenticated, move to next state + self._set_state("INIT") + + # Successful auth + self.send(''.format(xml.get("id"))) + + else: + # Failed auth + self.send( + ''.format( + xml.get("id") + ) + ) + + except ET.ParseError as e: + if "no element found" in e.msg: + xmppserverlog.debug( + "xml parse error - {} - {}".format(data.decode("utf-8"), e) + ) + elif "not well-formed (invalid token)" in e.msg: + xmppserverlog.debug( + "xml parse error - {} - {}".format(data.decode("utf-8"), e) + ) + else: + xmppserverlog.debug( + "xml parse error - {} - {}".format(data.decode("utf-8"), e) + ) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_sasl_auth(self, xml): + try: + + saslauth = base64.b64decode(xml.text).decode("utf-8").split("/") + username = saslauth[0] + username = saslauth[0].split("\x00")[1] + self.uid = username + if len(saslauth) > 1: + resource = saslauth[1] + self.clientresource = resource + elif len(saslauth[0].split("\x00")) > 2: + resource = saslauth[0].split("\x00")[2] + self.clientresource = resource + + if len(saslauth) > 2: + authcode = saslauth[2] + + if not self.uid.startswith("fuid"): + # Need sample data to see details here + bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy") + self.type = self.BOT + xmppserverlog.info("bot authenticated {}".format(self.uid)) + # Send response + self.send( + '' + ) # Success + + # Client authenticated, move to next state + self._set_state("INIT") + + else: + auth = False + if bumper.check_authcode(self.uid, authcode): + auth = True + elif bumper.use_auth == False: + auth = True + + if auth: + self.type = self.CONTROLLER + bumper.client_add(self.uid, "bumper", self.clientresource) + xmppserverlog.debug("client authenticated {}".format(self.uid)) + + # Client authenticated, move to next state + self._set_state("INIT") + + # Send response + self.send( + '' + ) # Success + + else: + # Failed to authenticate + self.send( + '' + ) # Fail + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_bind(self, xml): + try: + + bot = bumper.bot_get(self.uid) + if bot: + bumper.bot_set_xmpp(bot["did"], True) + + client = bumper.client_get(self.clientresource) + if client: + bumper.client_set_xmpp(client["resource"], True) + + clientbindxml = xml.getchildren() + clientresourcexml = clientbindxml[0].getchildren() + if self.devclass: # its a bot + self.name = "XMPP_Client_{}_{}".format(self.uid, self.devclass) + self.bumper_jid = "{}@{}.ecorobot.net/atom".format( + self.uid, self.devclass + ) + xmppserverlog.debug("new bot {}".format(self.uid)) + res = '{}'.format( + xml.get("id"), self.bumper_jid + ) + elif len(clientresourcexml) > 0: + self.clientresource = clientresourcexml[0].text + self.name = "XMPP_Client_{}".format(self.clientresource) + self.bumper_jid = "{}@{}/{}".format( + self.uid, XMPPServer.server_id, self.clientresource + ) + xmppserverlog.debug( + "new client {} using resource {}".format( + self.uid, self.clientresource + ) + ) + res = '{}'.format( + xml.get("id"), self.bumper_jid + ) + else: + self.name = "XMPP_Client_{}_{}".format(self.uid, self.address) + self.bumper_jid = "{}@{}".format(self.uid, XMPPServer.server_id) + xmppserverlog.debug("new client {}".format(self.uid)) + res = '{}'.format( + xml.get("id"), self.bumper_jid + ) + + self._set_state("BIND") + self.send(res) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_session(self, xml): + try: + res = ''.format(xml.get("id")) + self._set_state("READY") + self.send(res) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_presence(self, xml): + try: + + if len(xml) and xml[0].tag == "status": + xmppserverlog.debug( + "bot presence {} ".format(ET.tostring(xml, encoding="utf-8")) + ) + # Most likely a bot, possibly hello world in text + + # Send dummy return + self.send( + ' dummy '.format(self.bumper_jid) + ) + + # If it is a BOT, send extras + if self.type == self.BOT: + # get device info + self.send( + ''.format( + self.bumper_jid, XMPPServer.server_id + ) + ) + + else: + xmppserverlog.debug( + "client presence - {} ".format(ET.tostring(xml, encoding="utf-8")) + ) + + if xml.get("type") == "available": + xmppserverlog.debug( + "client presence available - {} ".format( + ET.tostring(xml, encoding="utf-8") + ) + ) + # Send dummy return + self.send( + ' dummy '.format(self.bumper_jid) + ) + elif xml.get("type") == "unavailable": + xmppserverlog.debug( + "client presence unavailable (DISCONNECT) - {} ".format( + ET.tostring(xml, encoding="utf-8") + ) + ) + + self._set_state("DISCONNECT") + else: + # Sometimes the android app sends these + xmppserverlog.debug( + "client presence (UNKNOWN) - {} ".format( + ET.tostring(xml, encoding="utf-8") + ) + ) + # Send dummy return + self.send( + ' dummy '.format(self.bumper_jid) + ) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _parse_data(self, data): + + if data.decode("utf-8").startswith( + "]+\?>)", r"", data.decode("utf-8")) + "" + ) + + else: + newdata = "{}".format( + data.decode("utf-8") + ) # Add artificial root + + try: + root = ET.fromstring(newdata) + for item in root.iter(): + if item.tag != "root": + if item.tag == "iq": + if self.log_incoming_data: + xmppserverlog.debug( + "from {} - {}".format( + self.address, + str( + ET.tostring(item, encoding="utf-8").decode( + "utf-8" + ) + ).replace("ns0:", ""), + ) + ) + self._handle_iq(item, newdata) + item.clear() + + elif "auth" in item.tag: + if "urn:ietf:params:xml:ns:xmpp-sasl" in item.tag: # SASL Auth + self._handle_sasl_auth(item) + item.clear() + + elif "presence" in item.tag: + self._handle_presence(item) + item.clear() + + else: + if self.log_incoming_data: + xmppserverlog.debug( + "Unparsed Item - {}".format( + str( + ET.tostring(item, encoding="utf-8").decode( + "utf-8" + ) + ).replace("ns0:", "") + ) + ) + + except ET.ParseError as e: + if ( + "no element found" in e.msg + ): # Element not closed or not all bytes received + # Happens wth connect stream often + if " - client is signalling end of session/disconnect + if not "" in newdata: + xmppserverlog.error("xml parse error - {} - {}".format(newdata, e)) + else: + self.send("") # Close stream + + else: + if "" in newdata: + xmppserverlog.error( + "xml parse error - {} - {}".format(newdata, e) + ) + else: + self.send("") # Close stream + self._set_state("DISCONNECT") + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def _handle_iq(self, xml, data): + try: + if len(xml): + child = self._tag_strip_uri(xml[0].tag) + else: + child = None + + if xml.tag == "iq": + if child == "bind": + self._handle_bind(xml) + elif child == "session": + self._handle_session(xml) + elif child == "ping": + self._handle_ping(xml, data) + elif child == "query": + if self.type == self.BOT: + self._handle_result(xml, data) + else: + self._handle_ctl(xml, data) + elif xml.get("type") == "result": + if self.type == self.BOT: + self._handle_result(xml, data) + else: + self._handle_result(xml, data) + elif xml.get("type") == "set": + if self.type == self.BOT: + self._handle_result(xml, data) + else: + self._handle_result(xml, data) + + except Exception as e: + xmppserverlog.exception("{}".format(e)) + + def run(self): + # xmppserverlog.info('client connected - {}'.format(self.address)) + await self._set_state("CONNECT") + while not self.state == self.DISCONNECT and not self.connection._closed: + data = b"" + time.sleep(0.1) + if not self.connection._closed: + try: + data = self.connection.recv(4096) + if data != b"": + self._parse_data(data) + except ConnectionResetError as e: + xmppserverlog.debug("{}".format(e)) + except OSError as e: + xmppserverlog.debug("{}".format(e)) + except Exception as e: + xmppserverlog.exception("{}".format(e)) diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index e044790..7d00fc5 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -4,8 +4,8 @@ from threading import Thread import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET import base64 import ssl -import contextvars import bumper +import asyncio, functools xmppserverlog = logging.getLogger("xmppserver") @@ -20,87 +20,39 @@ class XMPPServer: # Initialize bot server self.address = address - def run(self, run_async=False): - if run_async: - xmppserverlog.debug("Starting XMPPServer Thread: 1") - self.xmppthread = Thread(name="XMPPServer_Thread", target=self.run_server) - self.xmppthread.setDaemon(True) - self.xmppthread.start() - - else: - try: - self.run_server() - except KeyboardInterrupt: - self.disconnect() - - def run_server(self): - logging.info("Starting XMPP Server at {}".format(self.address)) - print("Starting XMPP Server at {}".format(self.address)) - - # xmppserverlog.setLevel(logging.DEBUG) - - # Set SSL Context - self.ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) - self.ssl_ctx.load_cert_chain( - certfile=bumper.server_cert, keyfile=bumper.server_key + async def async_server(self): + xmppserverlog.info( + "Starting XMPP Server at {}:{}".format(self.address[0], self.address[1]) + ) + server = await asyncio.start_server( + self.accept_client, self.address[0], self.address[1] ) - self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + await server.serve_forever() + # self.clients = {} # task -> (reader, writer) + + def accept_client(self, client_reader, client_writer): try: - self.socket.bind(self.address) - self.socket.listen(5) + aclient = XMPPAsyncClient(client_reader, client_writer) + task = asyncio.Task(aclient.handle_async_client()) + aclient._async_task = task + self.clients.append(aclient) - xmppserverlog.debug( - "listening on {}:{}".format(self.address[0], self.address[1]) - ) - while not self.exit_flag: - connection, client_address = self.socket.accept() + def client_done(aclient, task): + try: + self.clients.remove(aclient) + xmppserverlog.debug("End Connection for ({}:{} | {})".format(client_writer.get_extra_info("peername")[0], client_writer.get_extra_info("peername")[1], aclient.bumper_jid)) + client_writer.close() + except Exception as e: + xmppserverlog.error("{}".format(e)) - # disconnect any clients with this ip - for client in self.clients: - if client.address == client_address[0]: - xmppserverlog.debug( - "disconnecting existing client {} with resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.remove_client_byip(client.address) - - xmppserverlog.debug( - "starting new client with ip {}".format(client_address[0]) - ) - thread_id = uuid.uuid4() - client = Client(thread_id, connection, client_address) - client.setDaemon(True) - client.start() - self.clients.append(client) - - except PermissionError as e: - if "bind" in e.strerror: - xmppserverlog.exception( - "Error binding XMPPServer, exiting. Try using a different hostname or IP - {}".format( - e - ) - ) - exit(1) + clientaddr = client_writer.get_extra_info("peername") + xmppserverlog.debug("New Connection from {}:{}".format(clientaddr[0],clientaddr[1])) + task.add_done_callback(functools.partial(client_done, aclient)) except Exception as e: - xmppserverlog.exception("{}".format(e)) - exit(1) - - except KeyboardInterrupt as e: - xmppserverlog.exception("{}".format(e)) - - finally: - connection.shutdown(socket.SHUT_RDWR) - connection.close() - self.disconnect() - xmppserverlog.info("disconnecting") - - self.socket.close() + xmppserverlog.error("{}".format(e)) def disconnect(self): try: @@ -112,43 +64,9 @@ class XMPPServer: xmppserverlog.debug("shutting down") except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) - def remove_client_byip(self, ip): - for client in self.clients: - if client.address == ip: - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - - def remove_client_byresource(self, resource): - for client in self.clients: - if str(client.clientresource).lower() == str(resource).lower(): - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - - def remove_client_byuid(self, uid): - for client in self.clients: - if str(client.uid).lower() == str(uid).lower(): - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - - -class Client(threading.Thread): +class XMPPAsyncClient: IDLE = 0 CONNECT = 1 INIT = 2 @@ -158,52 +76,70 @@ class Client(threading.Thread): UNKNOWN = 0 BOT = 1 CONTROLLER = 2 + _async_task = None - def __init__(self, thread_id, connection, client_address): - threading.Thread.__init__(self) - self.id = thread_id - self.name = "XMPP_Client_{}".format(client_address[0]) + def __init__(self, client_reader, client_writer): self.type = self.UNKNOWN self.state = self.IDLE - self.connection = connection - self.address = client_address[0] + self.address = client_writer.get_extra_info("peername") + self.client_reader = client_reader + self.client_writer = client_writer self.clientresource = "" self.devclass = "" self.bumper_jid = "" self.uid = "" - self.log_sent_message = False # Set to true to log sends + self.log_sent_message = True # Set to true to log sends self.log_incoming_data = True # Set to true to log sends - xmppserverlog.debug( - "new client thread init for client with ip {}".format(self.address) - ) + xmppserverlog.debug("new client with ip {}".format(self.address)) - def send(self, command): + async def handle_async_client(self): + # xmppserverlog.info('client connected - {}'.format(self.address)) + await self._set_state("CONNECT") + while True: + await asyncio.sleep(0.05) + if not self.state == self.DISCONNECT: + data = await self.client_reader.read(4096) + if data is None or data == b'': + xmppserverlog.debug("Received no data") + # exit loop and disconnect + return + else: + await self._parse_data(data) + else: + break + + # exit loop and disconnect + return + + async def send(self, command): try: - if not self.connection._closed: - if self.log_sent_message: - xmppserverlog.debug("send {} - {}".format(self.address, command)) - self.connection.send(command.encode()) + # if not self.connection._closed: + if self.log_sent_message: + xmppserverlog.debug("send to ({}:{} | {}) - {}".format(self.address[0], self.address[1], self.bumper_jid, command)) + + self.client_writer.write(command.encode()) + await self.client_writer.drain() except BrokenPipeError as e: - xmppserverlog.debug("{}".format(e)) - self._set_state("DISCONNECT") + #xmppserverlog.debug("{}".format(e)) + await self._set_state("DISCONNECT") except ConnectionResetError as e: - xmppserverlog.debug("{}".format(e)) - self._set_state("DISCONNECT") + #xmppserverlog.debug("{}".format(e)) + await self._set_state("DISCONNECT") except ConnectionAbortedError as e: - xmppserverlog.debug("{}".format(e)) - self._set_state("DISCONNECT") + #xmppserverlog.debug("{}".format(e)) + await self._set_state("DISCONNECT") except OSError as e: - xmppserverlog.debug("{}".format(e)) + xmppserverlog.error("{}".format(e)) except Exception as e: xmppserverlog.exception("{}".format(e)) - def _disconnect(self): + async def _disconnect(self): try: bot = bumper.bot_get(self.uid) @@ -214,23 +150,23 @@ class Client(threading.Thread): if client: bumper.client_set_xmpp(client["resource"], False) - self.connection.close() + self.client_writer.close() except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) - def _tag_strip_uri(self, tag): + async def _tag_strip_uri(self, tag): try: if tag[0] == "{": _, _, tag = tag[1:].partition("}") return tag except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) - def _set_state(self, state): + async def _set_state(self, state): try: - new_state = getattr(Client, state) + new_state = getattr(XMPPAsyncClient, state) if self.state > new_state: raise Exception( "{} illegal state change {}->{}".format( @@ -238,33 +174,51 @@ class Client(threading.Thread): ) ) - xmppserverlog.debug("{} state: {}".format(self.address, state)) + xmppserverlog.debug("({}:{} | {}) state: {}".format(self.address[0],self.address[1],self.bumper_jid, state)) self.state = new_state if new_state == 5: - self._disconnect() + await self._disconnect() except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) - def _handle_ctl(self, xml, data): + async def _handle_ctl(self, xml, data): try: if "roster" in data: # Return not-implemented for roster - self.send( + await self.send( ''.format( xml.get("id") ) ) return + + if "disco#items" in data: + # Return not-implemented for disco#items + await self.send( + ''.format( + xml.get("id") + )) + return + + if "disco#info" in data: + # Return not-implemented for disco#info + await self.send( + ''.format( + xml.get("id") + ) + ) + return + if xml.get("type") == "set": if ( "com:sf" in data and xml.get("to") == "rl.ecorobot.net" ): # Android bind? Not sure what this does yet. - self.send( + await self.send( ''.format( xml.get("id"), self.uid, @@ -273,7 +227,7 @@ class Client(threading.Thread): ) ) - if xml[0][0]: + if len(xml[0]) > 0: ctl = xml[0][0] if ctl.get("admin") and self.type == self.BOT: xmppserverlog.debug( @@ -299,15 +253,15 @@ class Client(threading.Thread): if client.type == self.BOT: if client.uid.lower() in ctl_to.lower(): - xmppserverlog.info( + xmppserverlog.debug( "Sending ctl to bot: {}".format(rxmlstring) ) - client.send(rxmlstring) + await client.send(rxmlstring) except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) - def _handle_ping(self, xml, data): + async def _handle_ping(self, xml, data): try: if xml.get("to").find("@") == -1: # No to address # Ping to server - respond @@ -315,38 +269,48 @@ class Client(threading.Thread): xml.get("id"), xml.get("to") ) # xmppserverlog.debug("Server Ping resp: {}".format(pingresp)) - self.send(pingresp) + await self.send(pingresp) else: pingto = xml.get("to") pingfrom = self.bumper_jid - xml.attrib["from"] = pingfrom + xml.attrib["from"] = pingfrom pingstring = ET.tostring(xml).decode("utf-8") # clean up string to remove namespaces added by ET pingstring = pingstring.replace("xmlns:ns0=", "xmlns=") pingstring = pingstring.replace("ns0:", "") - pingstring = pingstring.replace('iq xmlns="com:ctl"', "iq") - pingstring = pingstring.replace("'.format( uuid.uuid4(), adminuser, self.bumper_jid ) @@ -401,7 +365,7 @@ class Client(threading.Thread): ) ) for client in XMPPServer.clients: - client.send(rxmlstring) + await client.send(rxmlstring) if xml.get("to").find("@") == -1: # No to address ctl_to = xml.get("to") @@ -415,7 +379,7 @@ class Client(threading.Thread): ): if not "@" in ctl_to: # No user@, send to all clients? # TODO: Revisit later, this may be wrong - client.send(rxmlstring) + await client.send(rxmlstring) elif ( client.uid.lower() in ctl_to.lower() @@ -425,12 +389,12 @@ class Client(threading.Thread): self.uid, client.uid, rxmlstring ) ) - client.send(rxmlstring) + await client.send(rxmlstring) except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_connect(self, data, xml=None): + async def _handle_connect(self, data, xml=None): try: if self.state == self.CONNECT: @@ -443,30 +407,32 @@ class Client(threading.Thread): self.devclass = data.decode("utf-8")[sc + 4 : ec] # ack jabbr:client # no STARTTLS - self.send( + await self.send( ''.format( XMPPServer.server_id ) ) # with STARTTLS - # self.send(''.format(XMPPServer.server_id)) - time.sleep(0.25) + # await self.send(''.format(XMPPServer.server_id)) + + await asyncio.sleep(0.25) + #time.sleep(0.25) # send authentication support for iq-auth (fallback) and SASL - self.send( + await self.send( 'PLAIN' ) - # self.send('') + # await self.send('') else: - self.send("") + await self.send("") else: if "jabber:iq:auth" in xml.tag: # Handle iq-auth - self._handle_iq_auth(xml) + await self._handle_iq_auth(xml) elif ( "urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag ): # Handle SASL Auth - self._handle_sasl_auth(xml) + await self._handle_sasl_auth(xml) else: xmppserverlog.error("Couldn't handle: {}".format(xml)) @@ -475,33 +441,34 @@ class Client(threading.Thread): # Client getting session after authentication if data.decode("utf-8").find("jabber:client") > -1: # ack jabbr:client - self.send( + await self.send( ''.format( XMPPServer.server_id ) ) - time.sleep(0.25) + await asyncio.sleep(0.25) + #time.sleep(0.25) # session - self.send( + await self.send( '' ) else: # Handle init bind if len(xml): - child = self._tag_strip_uri(xml[0].tag) + child = await self._tag_strip_uri(xml[0].tag) else: child = None if xml.tag == "iq": if child == "bind": - self._handle_bind(xml) + await self._handle_bind(xml) else: xmppserverlog.error("Couldn't handle: {}".format(xml)) except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_iq_auth(self, data): + async def _handle_iq_auth(self, data): try: xml = ET.fromstring(data.decode("utf-8")) ctl = xml[0][0] @@ -512,7 +479,7 @@ class Client(threading.Thread): and "auth}username" in ctl.tag and self.type == self.UNKNOWN ): - self.send( + await self.send( ''.format( xml.get("id") ) @@ -527,6 +494,7 @@ class Client(threading.Thread): xmlauth = xml[0].getchildren() # uid = "" password = "" + authcode = "" resource = "" for aitem in xmlauth: if "username" in aitem.tag: @@ -540,17 +508,15 @@ class Client(threading.Thread): self.clientresource = aitem.text resource = self.clientresource - if not self.uid.startswith("fuid"): - - # Need sample data to see details here + if self.devclass: # if there is a devclass it is a bot bumper.bot_add("", self.uid, "", resource, "eco-legacy") - xmppserverlog.info("bot authenticated {}".format(self.uid)) + xmppserverlog.debug("bot authenticated {}".format(self.uid)) # Client authenticated, move to next state - self._set_state("INIT") + await self._set_state("INIT") # Successful auth - self.send(''.format(xml.get("id"))) + await self.send(''.format(xml.get("id"))) else: auth = False @@ -564,14 +530,16 @@ class Client(threading.Thread): xmppserverlog.debug("client authenticated {}".format(self.uid)) # Client authenticated, move to next state - self._set_state("INIT") + await self._set_state("INIT") # Successful auth - self.send(''.format(xml.get("id"))) + await self.send( + ''.format(xml.get("id")) + ) else: # Failed auth - self.send( + await self.send( ''.format( xml.get("id") ) @@ -594,12 +562,13 @@ class Client(threading.Thread): except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_sasl_auth(self, xml): + async def _handle_sasl_auth(self, xml): try: saslauth = base64.b64decode(xml.text).decode("utf-8").split("/") username = saslauth[0] username = saslauth[0].split("\x00")[1] + authcode = "" self.uid = username if len(saslauth) > 1: resource = saslauth[1] @@ -610,19 +579,18 @@ class Client(threading.Thread): if len(saslauth) > 2: authcode = saslauth[2] - - if not self.uid.startswith("fuid"): - # Need sample data to see details here + + if self.devclass: # if there is a devclass it is a bot bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy") self.type = self.BOT - xmppserverlog.info("bot authenticated {}".format(self.uid)) + xmppserverlog.debug("bot authenticated {}".format(self.uid)) # Send response - self.send( + await self.send( '' ) # Success # Client authenticated, move to next state - self._set_state("INIT") + await self._set_state("INIT") else: auth = False @@ -637,23 +605,23 @@ class Client(threading.Thread): xmppserverlog.debug("client authenticated {}".format(self.uid)) # Client authenticated, move to next state - self._set_state("INIT") + await self._set_state("INIT") # Send response - self.send( + await self.send( '' ) # Success else: # Failed to authenticate - self.send( + await self.send( '' ) # Fail except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_bind(self, xml): + async def _handle_bind(self, xml): try: bot = bumper.bot_get(self.uid) @@ -671,7 +639,7 @@ class Client(threading.Thread): self.bumper_jid = "{}@{}.ecorobot.net/atom".format( self.uid, self.devclass ) - xmppserverlog.debug("new bot {}".format(self.uid)) + xmppserverlog.debug("new bot ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) @@ -681,83 +649,84 @@ class Client(threading.Thread): self.bumper_jid = "{}@{}/{}".format( self.uid, XMPPServer.server_id, self.clientresource ) - xmppserverlog.debug( - "new client {} using resource {}".format( - self.uid, self.clientresource - ) - ) + xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) else: self.name = "XMPP_Client_{}_{}".format(self.uid, self.address) self.bumper_jid = "{}@{}".format(self.uid, XMPPServer.server_id) - xmppserverlog.debug("new client {}".format(self.uid)) + xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) - self._set_state("BIND") - self.send(res) + await self._set_state("BIND") + await self.send(res) except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_session(self, xml): + async def _handle_session(self, xml): try: res = ''.format(xml.get("id")) - self._set_state("READY") - self.send(res) + await self._set_state("READY") + await self.send(res) + asyncio.Task(self.schedule_ping(30)) + except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_presence(self, xml): + async def _handle_presence(self, xml): try: if len(xml) and xml[0].tag == "status": xmppserverlog.debug( - "bot presence {} ".format(ET.tostring(xml, encoding="utf-8")) + "bot presence {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) # Most likely a bot, possibly hello world in text # Send dummy return - self.send( + await self.send( ' dummy '.format(self.bumper_jid) ) + + # If it is a BOT, send extras if self.type == self.BOT: # get device info - self.send( + await self.send( ''.format( self.bumper_jid, XMPPServer.server_id ) ) + + else: xmppserverlog.debug( - "client presence - {} ".format(ET.tostring(xml, encoding="utf-8")) + "client presence - {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) - + if xml.get("type") == "available": xmppserverlog.debug( "client presence available - {} ".format( - ET.tostring(xml, encoding="utf-8") - ) + ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) + # Send dummy return - self.send( + await self.send( ' dummy '.format(self.bumper_jid) ) elif xml.get("type") == "unavailable": xmppserverlog.debug( "client presence unavailable (DISCONNECT) - {} ".format( - ET.tostring(xml, encoding="utf-8") - ) - ) + ET.tostring(xml, encoding="utf-8").decode("utf-8")) + ) - self._set_state("DISCONNECT") + await self._set_state("DISCONNECT") else: # Sometimes the android app sends these xmppserverlog.debug( @@ -766,14 +735,14 @@ class Client(threading.Thread): ) ) # Send dummy return - self.send( + await self.send( ' dummy '.format(self.bumper_jid) ) except Exception as e: xmppserverlog.exception("{}".format(e)) - def _parse_data(self, data): + async def _parse_data(self, data): if data.decode("utf-8").startswith( "" in newdata: xmppserverlog.error("xml parse error - {} - {}".format(newdata, e)) else: - self.send("") # Close stream + await self.send("") # Close stream else: if "" in newdata: xmppserverlog.error( "xml parse error - {} - {}".format(newdata, e) ) else: - self.send("") # Close stream - self._set_state("DISCONNECT") + await self.send("") # Close stream + await self._set_state("DISCONNECT") except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_iq(self, xml, data): + async def _handle_iq(self, xml, data): try: if len(xml): - child = self._tag_strip_uri(xml[0].tag) + child = await self._tag_strip_uri(xml[0].tag) else: child = None if xml.tag == "iq": if child == "bind": - self._handle_bind(xml) + await self._handle_bind(xml) elif child == "session": - self._handle_session(xml) + await self._handle_session(xml) elif child == "ping": - self._handle_ping(xml, data) + await self._handle_ping(xml, data) elif child == "query": if self.type == self.BOT: - self._handle_result(xml, data) + await self._handle_result(xml, data) else: - self._handle_ctl(xml, data) + await self._handle_ctl(xml, data) elif xml.get("type") == "result": if self.type == self.BOT: - self._handle_result(xml, data) + await self._handle_result(xml, data) else: - self._handle_result(xml, data) + await self._handle_result(xml, data) elif xml.get("type") == "set": if self.type == self.BOT: - self._handle_result(xml, data) + await self._handle_result(xml, data) else: - self._handle_result(xml, data) + await self._handle_result(xml, data) except Exception as e: xmppserverlog.exception("{}".format(e)) - - def run(self): - # xmppserverlog.info('client connected - {}'.format(self.address)) - self._set_state("CONNECT") - while not self.state == self.DISCONNECT and not self.connection._closed: - data = b"" - time.sleep(0.1) - if not self.connection._closed: - try: - data = self.connection.recv(4096) - if data != b"": - self._parse_data(data) - except ConnectionResetError as e: - xmppserverlog.debug("{}".format(e)) - except OSError as e: - xmppserverlog.debug("{}".format(e)) - except Exception as e: - xmppserverlog.exception("{}".format(e)) - diff --git a/start_bumper.py b/start_bumper.py index 1975e80..a01cd81 100644 --- a/start_bumper.py +++ b/start_bumper.py @@ -5,10 +5,19 @@ import bumper import sys, socket import time import platform +import os +#os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG +import asyncio -def main(): +async def main(): + try: + loop = asyncio.get_event_loop() + except: + loop = asyncio.new_event_loop() + args = sys.argv + listen_host = "" if len(args) > 0: if "--debug" in args: @@ -16,23 +25,22 @@ def main(): level=logging.DEBUG, format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s", ) + loop.set_debug(True) # Set asyncio loop to debug + #logging.getLogger("asyncio").setLevel(logging.DEBUG) # Show debug asyncio logs (disabled in init, uncomment for debugging asyncio) else: logging.basicConfig( level=logging.INFO, format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s", ) - # format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s") - listen_host = args.index("--listen") - if (len(args) - 1) >= (listen_host + 1): - listen_host = args[listen_host+1] - else: + if "--listen" in args: + listen_host = args[args.index("--listen") + 1] + + if listen_host == "": if platform.system() == "Darwin": # If a Mac, use 0.0.0.0 for listening listen_host = "0.0.0.0" else: - listen_host = socket.gethostbyname(socket.gethostname()) - #listen_host = "localhost" # Try this if the above doesn't work - + listen_host = socket.gethostbyname(socket.gethostname()) conf_address_443 = (listen_host, 443) conf_address_8007 = (listen_host, 8007) @@ -51,38 +59,39 @@ def main(): conf_address_8007, usessl=False, helperbot=mqtt_helperbot ) - # add user - # users = bumper.bumper_users_var.get() - # user1 = bumper.BumperUser('user1') - # user1.add_device('devid') - # user1.add_bot('bot_did') - # users.append(user1) - # bumper.bumper_users_var.set(users) + # Start web servers + conf_server.confserver_app() + task_conf_server = asyncio.create_task(conf_server.start_server()) + bumper.bumperlog.debug("task_conf_server added") + await task_conf_server - # start xmpp server on port 5223 (sync) - xmpp_server.run(run_async=True) # Start in new thread + conf_server_2.confserver_app() + task_conf_server2 = asyncio.create_task(conf_server_2.start_server()) + bumper.bumperlog.debug("task_conf_server2 added") + await task_conf_server2 - # start mqtt server on port 8883 (async) - mqtt_server.run(run_async=True) # Start in new thread + # Start MQTT Server + task_mqtt_server = asyncio.create_task(mqtt_server.broker_coro()) + bumper.bumperlog.debug("task_mqtt_server added") + await task_mqtt_server - time.sleep(1.5) # Wait for broker startup - - # start mqtt_helperbot (async) - mqtt_helperbot.run(run_async=True) # Start in new thread - - # start conf server on port 443 (async) - Used for most https calls - conf_server.run(run_async=True) # Start in new thread - - # start conf server on port 8007 (async) - Used for a load balancer request - conf_server_2.run(run_async=True) # Start in new thread + # Start MQTT Helperbot + task_mqtt_helperbot = asyncio.create_task(mqtt_helperbot.start_helper_bot()) + bumper.bumperlog.debug("task_mqtt_helperbot added") + await task_mqtt_helperbot + + # Start XMPP Server + task_xmpp_server = asyncio.create_task(xmpp_server.async_server()) + bumper.bumperlog.debug("task_xmpp_server added") + await task_xmpp_server while True: try: - time.sleep(30) + await asyncio.sleep(30) bumper.revoke_expired_tokens() - disconnected_clients = bumper.get_disconnected_xmpp_clients() - for client in disconnected_clients: - xmpp_server.remove_client_byuid(client["userid"]) + #disconnected_clients = bumper.get_disconnected_xmpp_clients() + #for client in disconnected_clients: + # xmpp_server.remove_client_byuid(client["userid"]) except KeyboardInterrupt: bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt") @@ -91,4 +100,4 @@ def main(): if __name__ == "__main__": - main() + asyncio.run(main()) diff --git a/tests/test_confserver.py b/tests/test_confserver.py index 0de11dc..b7a4849 100644 --- a/tests/test_confserver.py +++ b/tests/test_confserver.py @@ -18,7 +18,14 @@ def async_return(result): return f def test_disconnect(): - confserver.disconnect() + + async def test_disconnect_async(): + await confserver.disconnect() + + loop = asyncio.get_event_loop() + # Test + loop.run_until_complete(test_disconnect_async()) + def test_base(): if os.path.exists("tests/tmp.db"):