Merge async and xmpp work #34

Merged
bmartin5692 merged 20 commits from dev_broken-XMPP into master 2019-05-23 15:22:58 +02:00
10 changed files with 1313 additions and 614 deletions

1
.gitignore vendored
View file

@ -6,3 +6,4 @@ __pycache__
nosetests.xml nosetests.xml
tests/report tests/report
tests/tmp.db tests/tmp.db
logs/

View file

@ -15,6 +15,7 @@ nose = "*"
coverage = "*" coverage = "*"
mock = "*" mock = "*"
pylint = "*" pylint = "*"
pbr = "*"
[pipenv] [pipenv]
allow_prereleases = true allow_prereleases = true

180
Pipfile.lock generated
View file

@ -1,7 +1,7 @@
{ {
"_meta": { "_meta": {
"hash": { "hash": {
"sha256": "f6a2679f2c348e81b120bfa9d06f544ed25c1fcea44d94c419e6af514579fdee" "sha256": "6e5b4ccc879dfdf82f53f7966eac54af90f9b9a2ffd1da6f1ce126d0ca125fc9"
}, },
"pipfile-spec": 6, "pipfile-spec": 6,
"requires": {}, "requires": {},
@ -244,90 +244,68 @@
], ],
"version": "==7.0" "version": "==7.0"
}, },
"colorama": {
"hashes": [
"sha256:05eed71e2e327246ad6b38c540c4a3117230b19679b875190486ddd2d721422d",
"sha256:f8ac84de7840f5b9c4e3347b3c1eaa50f7e49c2b07596221daec5edaabbd7c48"
],
"markers": "sys_platform == 'win32'",
"version": "==0.4.1"
},
"coverage": { "coverage": {
"hashes": [ "hashes": [
"sha256:029c69deaeeeae1b15bc6c59f0ffa28aa8473721c614a23f2c2976dec245cd12", "sha256:0402b1822d513d0231589494bceddb067d20581f5083598c451b56c684b0e5d6",
"sha256:02abbbebc6e9d5abe13cd28b5e963dedb6ffb51c146c916d17b18f141acd9947", "sha256:0644e28e8aea9d9d563607ee8b7071b07dd57a4a3de11f8684cd33c51c0d1b93",
"sha256:1bbfe5b82a3921d285e999c6d256c1e16b31c554c29da62d326f86c173d30337", "sha256:0874a283686803884ec0665018881130604956dbaa344f2539c46d82cbe29eda",
"sha256:210c02f923df33a8d0e461c86fdcbbb17228ff4f6d92609fc06370a98d283c2d", "sha256:0988c3837df4bc371189bb3425d5232cf150055452034c232dda9cbe04f9c38e",
"sha256:2d0807ba935f540d20b49d5bf1c0237b90ce81e133402feda906e540003f2f7a", "sha256:20bc3205b3100956bb72293fabb97f0ed972c81fed10b3251c90c70dcb0599ab",
"sha256:35d7a013874a7c927ce997350d314144ffc5465faf787bb4e46e6c4f381ef562", "sha256:2cc9142a3367e74eb6b19d58c53ebb1dfd7336b91cdcc91a6a2888bf8c7af984",
"sha256:3636f9d0dcb01aed4180ef2e57a4e34bb4cac3ecd203c2a23db8526d86ab2fb4", "sha256:3ae9a0a59b058ce0761c3bd2c2d66ecb2ee2b8ac592620184370577f7a546fb3",
"sha256:42f4be770af2455a75e4640f033a82c62f3fb0d7a074123266e143269d7010ef", "sha256:3b2e30b835df58cb973f478d09f3d82e90c98c8e5059acc245a8e4607e023801",
"sha256:48440b25ba6cda72d4c638f3a9efa827b5b87b489c96ab5f4ff597d976413156", "sha256:401e9b04894eb1498c639c6623ee78a646990ce5f095248e2440968aafd6e90e",
"sha256:4dac8dfd1acf6a3ac657475dfdc66c621f291b1b7422a939cc33c13ac5356473", "sha256:41ec5812d5decdaa72708be3018e7443e90def4b5a71294236a4df192cf9eab9",
"sha256:4e8474771c69c2991d5eab65764289a7dd450bbea050bc0ebb42b678d8222b42", "sha256:475769b638a055e75b3d3219e054fe2a023c0b077ff15bff6c95aba7e93e6cac",
"sha256:551f10ddfeff56a1325e5a34eff304c5892aa981fd810babb98bfee77ee2fb17", "sha256:61424f4e2e82c4129a4ba71e10ebacb32a9ecd6f80de2cd05bdead6ba75ed736",
"sha256:5b104982f1809c1577912519eb249f17d9d7e66304ad026666cb60a5ef73309c", "sha256:811969904d4dd0bee7d958898be8d9d75cef672d9b7e7db819dfeac3d20d2d0c",
"sha256:5c62aef73dfc87bfcca32cee149a1a7a602bc74bac72223236b0023543511c88", "sha256:86224bb99abfd672bf2f9fcecad5e8d7a3fa94f7f71513f2210460a0350307cd",
"sha256:633151f8d1ad9467b9f7e90854a7f46ed8f2919e8bc7d98d737833e8938fc081", "sha256:9a238a20a3af00665f8381f7e53e9c606f9bb652d2423f6b822f6cb790d887e8",
"sha256:772207b9e2d5bf3f9d283b88915723e4e92d9a62c83f44ec92b9bd0cd685541b", "sha256:a23b3fbc14d4e6182ecebfd22f3729beef0636d151d94764a1c28330d185e4e5",
"sha256:7d5e02f647cd727afc2659ec14d4d1cc0508c47e6cfb07aea33d7aa9ca94d288", "sha256:ac162b4ebe51b7a2b7f5e462c4402802633eb81e77c94f8a7c1ed8a556e72c75",
"sha256:a9798a4111abb0f94584000ba2a2c74841f2cfe5f9254709756367aabbae0541", "sha256:b6187378726c84365bf297b5dcdae8789b6a5823b200bea23797777e5a63be09",
"sha256:b38ea741ab9e35bfa7015c93c93bbd6a1623428f97a67083fc8ebd366238b91f", "sha256:bcd5723d905ed4a825f17410a53535f880b6d7548ae3d89078db7b1ceefcd853",
"sha256:b6a5478c904236543c0347db8a05fac6fc0bd574c870e7970faa88e1d9890044", "sha256:c48a4f9c5fb385269bb7fbaf9c1326a94863b65ec7f5c96b2ea56b252f01ad08",
"sha256:c6248bfc1de36a3844685a2e10ba17c18119ba6252547f921062a323fb31bff1", "sha256:cd40199d6f1c29c85b170d25589be9a97edff8ee7e62be180a2a137823896030",
"sha256:c705ab445936457359b1424ef25ccc0098b0491b26064677c39f1d14a539f056", "sha256:d1bc331a7d069485ac1d8c25a0ea1f6aab6cb2a87146fb652222481c1bddc9ff",
"sha256:d95a363d663ceee647291131dbd213af258df24f41350246842481ec3709bd33", "sha256:d7e0cdc249aa0f94aa2e531b03999ddaf03a10b4fa090a894712d4c8066abd89",
"sha256:e27265eb80cdc5dab55a40ef6f890e04ecc618649ad3da5265f128b141f93f78", "sha256:e9ee8fcd8e067fcc5d7276d46e07e863102b70a52545ef4254df1ff0893ce75f",
"sha256:ebc276c9cb5d917bd2ae959f84ffc279acafa9c9b50b0fa436ebb70bbe2166ea", "sha256:eb313c23d983b7810504f42104e8dcd1c7ccdda8fbaab82aab92ab79fea19345",
"sha256:f4d229866d030863d0fe3bf297d6d11e6133ca15bbb41ed2534a8b9a3d6bd061", "sha256:f9cfd478654b509941b85ed70f870f5e3c74678f566bec12fd26545e5340ba47",
"sha256:f95675bd88b51474d4fe5165f3266f419ce754ffadfb97f10323931fa9ac95e5", "sha256:fae1fa144034d021a52cb9ea200eb8dedf91869c6df8202ad5d149b41ed91cc8"
"sha256:f95bc54fb6d61b9f9ff09c4ae8ff6a3f5edc937cda3ca36fc937302a7c152bf1",
"sha256:fd0f6be53de40683584e5331c341e65a679dbe5ec489a0697cec7c2ef1a48cda"
], ],
"index": "pypi", "index": "pypi",
"version": "==5.0a4" "version": "==5.0a5"
}, },
"isort": { "isort": {
"hashes": [ "hashes": [
"sha256:18c796c2cd35eb1a1d3f012a214a542790a1aed95e29768bdcb9f2197eccbd0b", "sha256:c40744b6bc5162bbb39c1257fe298b7a393861d50978b565f3ccd9cb9de0182a",
"sha256:96151fca2c6e736503981896495d344781b60d18bfda78dc11b290c6125ebdb6" "sha256:f57abacd059dc3bd666258d1efb0377510a89777fda3e3274e3c01f7c03ae22d"
], ],
"version": "==4.3.15" "version": "==4.3.20"
}, },
"lazy-object-proxy": { "lazy-object-proxy": {
"hashes": [ "hashes": [
"sha256:0ce34342b419bd8f018e6666bfef729aec3edf62345a53b537a4dcc115746a33", "sha256:159a745e61422217881c4de71f9eafd9d703b93af95618635849fe469a283661",
"sha256:1b668120716eb7ee21d8a38815e5eb3bb8211117d9a90b0f8e21722c0758cc39", "sha256:23f63c0821cc96a23332e45dfaa83266feff8adc72b9bcaef86c202af765244f",
"sha256:209615b0fe4624d79e50220ce3310ca1a9445fd8e6d3572a896e7f9146bbf019", "sha256:3b11be575475db2e8a6e11215f5aa95b9ec14de658628776e10d96fa0b4dac13",
"sha256:27bf62cb2b1a2068d443ff7097ee33393f8483b570b475db8ebf7e1cba64f088", "sha256:3f447aff8bc61ca8b42b73304f6a44fa0d915487de144652816f950a3f1ab821",
"sha256:27ea6fd1c02dcc78172a82fc37fcc0992a94e4cecf53cb6d73f11749825bd98b", "sha256:4ba73f6089cd9b9478bc0a4fa807b47dbdb8fad1d8f31a0f0a5dbf26a4527a71",
"sha256:2c1b21b44ac9beb0fc848d3993924147ba45c4ebc24be19825e57aabbe74a99e", "sha256:4f53eadd9932055eac465bd3ca1bd610e4d7141e1278012bd1f28646aebc1d0e",
"sha256:2df72ab12046a3496a92476020a1a0abf78b2a7db9ff4dc2036b8dd980203ae6", "sha256:64483bd7154580158ea90de5b8e5e6fc29a16a9b4db24f10193f0c1ae3f9d1ea",
"sha256:320ffd3de9699d3892048baee45ebfbbf9388a7d65d832d7e580243ade426d2b", "sha256:6f72d42b0d04bfee2397aa1862262654b56922c20a9bb66bb76b6f0e5e4f9229",
"sha256:50e3b9a464d5d08cc5227413db0d1c4707b6172e4d4d915c1c70e4de0bbff1f5", "sha256:7c7f1ec07b227bdc561299fa2328e85000f90179a2f44ea30579d38e037cb3d4",
"sha256:5276db7ff62bb7b52f77f1f51ed58850e315154249aceb42e7f4c611f0f847ff", "sha256:7c8b1ba1e15c10b13cad4171cfa77f5bb5ec2580abc5a353907780805ebe158e",
"sha256:61a6cf00dcb1a7f0c773ed4acc509cb636af2d6337a08f362413c76b2b47a8dd", "sha256:8559b94b823f85342e10d3d9ca4ba5478168e1ac5658a8a2f18c991ba9c52c20",
"sha256:6ae6c4cb59f199d8827c5a07546b2ab7e85d262acaccaacd49b62f53f7c456f7", "sha256:a262c7dfb046f00e12a2bdd1bafaed2408114a89ac414b0af8755c696eb3fc16",
"sha256:7661d401d60d8bf15bb5da39e4dd72f5d764c5aff5a86ef52a042506e3e970ff", "sha256:acce4e3267610c4fdb6632b3886fe3f2f7dd641158a843cf6b6a68e4ce81477b",
"sha256:7bd527f36a605c914efca5d3d014170b2cb184723e423d26b1fb2fd9108e264d", "sha256:be089bb6b83fac7f29d357b2dc4cf2b8eb8d98fe9d9ff89f9ea6012970a853c7",
"sha256:7cb54db3535c8686ea12e9535eb087d32421184eacc6939ef15ef50f83a5e7e2", "sha256:bfab710d859c779f273cc48fb86af38d6e9210f38287df0069a63e40b45a2f5c",
"sha256:7f3a2d740291f7f2c111d86a1c4851b70fb000a6c8883a59660d95ad57b9df35", "sha256:c10d29019927301d524a22ced72706380de7cfc50f767217485a912b4c8bd82a",
"sha256:81304b7d8e9c824d058087dcb89144842c8e0dea6d281c031f59f0acf66963d4", "sha256:dd6e2b598849b3d7aee2295ac765a578879830fb8966f70be8cd472e6069932e",
"sha256:933947e8b4fbe617a51528b09851685138b49d511af0b6c0da2539115d6d4514", "sha256:e408f1eacc0a68fed0c08da45f31d0ebb38079f043328dce69ff133b95c29dc1"
"sha256:94223d7f060301b3a8c09c9b3bc3294b56b2188e7d8179c762a1cda72c979252",
"sha256:ab3ca49afcb47058393b0122428358d2fbe0408cf99f1b58b295cfeb4ed39109",
"sha256:bd6292f565ca46dee4e737ebcc20742e3b5be2b01556dafe169f6c65d088875f",
"sha256:c0e2945ebf5b6eb32848e9bd6f850f558722caf7ae9427c9fff8e0ea7b185a2f",
"sha256:cb924aa3e4a3fb644d0c463cad5bc2572649a6a3f68a7f8e4fbe44aaa6d77e4c",
"sha256:d0fc7a286feac9077ec52a927fc9fe8fe2fabab95426722be4c953c9a8bede92",
"sha256:ddc34786490a6e4ec0a855d401034cbd1242ef186c20d79d2166d6a4bd449577",
"sha256:e34b155e36fa9da7e1b7c738ed7767fc9491a62ec6af70fe9da4a057759edc2d",
"sha256:e5b9e8f6bda48460b7b143c3821b21b452cb3a835e6bbd5dd33aa0c8d3f5137d",
"sha256:e81ebf6c5ee9684be8f2c87563880f93eedd56dd2b6146d8a725b50b7e5adb0f",
"sha256:eb91be369f945f10d3a49f5f9be8b3d0b93a4c2be8f8a5b83b0571b8123e0a7a",
"sha256:f460d1ceb0e4a5dcb2a652db0904224f367c9b3c1470d5a7683c0480e582468b"
], ],
"version": "==1.3.1" "version": "==1.4.1"
}, },
"mccabe": { "mccabe": {
"hashes": [ "hashes": [
@ -338,11 +316,11 @@
}, },
"mock": { "mock": {
"hashes": [ "hashes": [
"sha256:5ce3c71c5545b472da17b72268978914d0252980348636840bd34a00b5cc96c1", "sha256:83657d894c90d5681d62155c82bda9c1187827525880eda8ff5df4ec813437c3",
"sha256:b158b6df76edd239b8208d481dc46b6afd45a846b7812ff0ce58971cf5bc8bba" "sha256:d157e52d4e5b938c550f39eb2fd15610db062441a9c2747d3dbfa9298211d0f8"
], ],
"index": "pypi", "index": "pypi",
"version": "==2.0.0" "version": "==3.0.5"
}, },
"nose": { "nose": {
"hashes": [ "hashes": [
@ -355,10 +333,11 @@
}, },
"pbr": { "pbr": {
"hashes": [ "hashes": [
"sha256:8257baf496c8522437e8a6cfe0f15e00aedc6c0e0e7c9d55eeeeab31e0853843", "sha256:6901995b9b686cb90cceba67a0f6d4d14ae003cd59bc12beb61549bdfbe3bc89",
"sha256:8c361cc353d988e4f5b998555c88098b9d5964c2e11acf7b0d21925a66bb5824" "sha256:d950c64aeea5456bbd147468382a5bb77fe692c13c9f00f0219814ce5b642755"
], ],
"version": "==5.1.3" "index": "pypi",
"version": "==5.2.0"
}, },
"pylint": { "pylint": {
"hashes": [ "hashes": [
@ -384,33 +363,32 @@
}, },
"typed-ast": { "typed-ast": {
"hashes": [ "hashes": [
"sha256:035a54ede6ce1380599b2ce57844c6554666522e376bd111eb940fbc7c3dad23", "sha256:132eae51d6ef3ff4a8c47c393a4ef5ebf0d1aecc96880eb5d6c8ceab7017cc9b",
"sha256:037c35f2741ce3a9ac0d55abfcd119133cbd821fffa4461397718287092d9d15", "sha256:18141c1484ab8784006c839be8b985cfc82a2e9725837b0ecfa0203f71c4e39d",
"sha256:049feae7e9f180b64efacbdc36b3af64a00393a47be22fa9cb6794e68d4e73d3", "sha256:2baf617f5bbbfe73fd8846463f5aeafc912b5ee247f410700245d68525ec584a",
"sha256:19228f7940beafc1ba21a6e8e070e0b0bfd1457902a3a81709762b8b9039b88d", "sha256:3d90063f2cbbe39177e9b4d888e45777012652d6110156845b828908c51ae462",
"sha256:2ea681e91e3550a30c2265d2916f40a5f5d89b59469a20f3bad7d07adee0f7a6", "sha256:4304b2218b842d610aa1a1d87e1dc9559597969acc62ce717ee4dfeaa44d7eee",
"sha256:3a6b0a78af298d82323660df5497bcea0f0a4a25a0b003afd0ce5af049bd1f60", "sha256:4983ede548ffc3541bae49a82675996497348e55bafd1554dc4e4a5d6eda541a",
"sha256:5385da8f3b801014504df0852bf83524599df890387a3c2b17b7caa3d78b1773", "sha256:5315f4509c1476718a4825f45a203b82d7fdf2a6f5f0c8f166435975b1c9f7d4",
"sha256:606d8afa07eef77280c2bf84335e24390055b478392e1975f96286d99d0cb424", "sha256:6cdfb1b49d5345f7c2b90d638822d16ba62dc82f7616e9b4caa10b72f3f16649",
"sha256:69245b5b23bbf7fb242c9f8f08493e9ecd7711f063259aefffaeb90595d62287", "sha256:7b325f12635598c604690efd7a0197d0b94b7d7778498e76e0710cd582fd1c7a",
"sha256:6f6d839ab09830d59b7fa8fb6917023d8cb5498ee1f1dbd82d37db78eb76bc99", "sha256:8d3b0e3b8626615826f9a626548057c5275a9733512b137984a68ba1598d3d2f",
"sha256:730888475f5ac0e37c1de4bd05eeb799fdb742697867f524dc8a4cd74bcecc23", "sha256:8f8631160c79f53081bd23446525db0bc4c5616f78d04021e6e434b286493fd7",
"sha256:9819b5162ffc121b9e334923c685b0d0826154e41dfe70b2ede2ce29034c71d8", "sha256:912de10965f3dc89da23936f1cc4ed60764f712e5fa603a09dd904f88c996760",
"sha256:9e60ef9426efab601dd9aa120e4ff560f4461cf8442e9c0a2b92548d52800699", "sha256:b010c07b975fe853c65d7bbe9d4ac62f1c69086750a574f6292597763781ba18",
"sha256:af5fbdde0690c7da68e841d7fc2632345d570768ea7406a9434446d7b33b0ee1", "sha256:c908c10505904c48081a5415a1e295d8403e353e0c14c42b6d67f8f97fae6616",
"sha256:b64efdbdf3bbb1377562c179f167f3bf301251411eb5ac77dec6b7d32bcda463", "sha256:c94dd3807c0c0610f7c76f078119f4ea48235a953512752b9175f9f98f5ae2bd",
"sha256:bac5f444c118aeb456fac1b0b5d14c6a71ea2a42069b09c176f75e9bd4c186f6", "sha256:ce65dee7594a84c466e79d7fb7d3303e7295d16a83c22c7c4037071b059e2c21",
"sha256:bda9068aafb73859491e13b99b682bd299c1b5fd50644d697533775828a28ee0", "sha256:eaa9cfcb221a8a4c2889be6f93da141ac777eb8819f077e1d09fb12d00a09a93",
"sha256:d659517ca116e6750101a1326107d3479028c5191f0ecee3c7203c50f5b915b0", "sha256:f3376bc31bad66d46d44b4e6522c5c21976bf9bca4ef5987bb2bf727f4506cbb",
"sha256:eddd3fb1f3e0f82e5915a899285a39ee34ce18fd25d89582bc89fc9fb16cd2c6" "sha256:f9202fa138544e13a4ec1a6792c35834250a85958fde1251b6a22e07d1260ae7"
], ],
"markers": "implementation_name == 'cpython'", "markers": "implementation_name == 'cpython'",
"version": "==1.3.1" "version": "==1.3.5"
}, },
"wrapt": { "wrapt": {
"hashes": [ "hashes": [
"sha256:4aea003270831cceb8a90ff27c4031da6ead7ec1886023b80ce0dfe0adf61533", "sha256:4aea003270831cceb8a90ff27c4031da6ead7ec1886023b80ce0dfe0adf61533"
"sha256:71ad0a3729a2f29eb25ff0dd6e66fc9768c749ff0c3f8899a6ac35c059f502c1"
], ],
"version": "==1.11.1" "version": "==1.11.1"
} }

View file

@ -1,25 +1,21 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
from .confserver import ConfServer from bumper.confserver import ConfServer
from .mqttserver import MQTTServer from bumper.mqttserver import MQTTServer, MQTTHelperBot
from .mqttserver import MQTTHelperBot from bumper.xmppserver import XMPPServer
from .xmppserver import XMPPServer
import asyncio import asyncio
import contextvars
import json import json
import time import time
from datetime import datetime, timedelta from datetime import datetime, timedelta
import platform import platform
import os import os, sys
import logging import logging
from logging.handlers import RotatingFileHandler
from base64 import b64decode, b64encode from base64 import b64decode, b64encode
from tinydb import TinyDB, Query from tinydb import TinyDB, Query
import json
from tinydb.storages import MemoryStorage 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" ca_cert = "./certs/CA/cacert.pem"
server_cert = "./certs/cert.pem" server_cert = "./certs/cert.pem"
server_key = "./certs/key.pem" server_key = "./certs/key.pem"
@ -29,21 +25,49 @@ token_validity_seconds = 3600 # 1 hour
db = None db = None
# Logs # 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") 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") 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 # Override the logging level
# confserverlog.setLevel(logging.INFO) # confserverlog.setLevel(logging.INFO)
mqttserverlog = logging.getLogger("mqttserver") 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 # Override the logging level
# mqttserverlog.setLevel(logging.INFO) # mqttserverlog.setLevel(logging.INFO)
helperbotlog = logging.getLogger("helperbot") 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 # Override the logging level
# helperbotlog.setLevel(logging.INFO) # helperbotlog.setLevel(logging.INFO)
xmppserverlog = logging.getLogger("xmppserver") 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 # Override the logging level
# xmppserverlog.setLevel(logging.INFO) # xmppserverlog.setLevel(logging.INFO)
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
def get_milli_time(timetoconvert): def get_milli_time(timetoconvert):
return int(round(timetoconvert * 1000)) return int(round(timetoconvert * 1000))
@ -57,12 +81,14 @@ def db_file():
def os_db_path(): def os_db_path():
if platform.system() == "Windows": 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") return os.path.join(os.getenv("APPDATA"), "bumper.db")
else: else:
os.makedirs(os.path.expanduser("~/.config"), exist_ok=True) #Ensure db_path directory exists or create
return os.path.expanduser("~/.config/bumper.db") return os.path.expanduser("~/.config/bumper.db")
def db_get(): def db_get():
try:
# Will create the database if it doesn't exist # Will create the database if it doesn't exist
db = TinyDB(db_file()) db = TinyDB(db_file())
@ -75,6 +101,13 @@ def db_get():
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): class BumperUser(object):
def __init__(self, userid=""): def __init__(self, userid=""):
self.userid = userid self.userid = userid
@ -632,7 +665,8 @@ def bot_add(sn, did, devclass, resource, company):
newbot.company = company newbot.company = company
bot = bot_get(did) bot = bot_get(did)
if not bot: 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( bumperlog.info(
"Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did) "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did)
) )

View file

@ -8,7 +8,6 @@ import bumper
import time import time
from datetime import datetime, timedelta from datetime import datetime, timedelta
import asyncio import asyncio
import contextvars
from aiohttp import web from aiohttp import web
import uuid import uuid
import xml.etree.ElementTree as ET import xml.etree.ElementTree as ET
@ -32,9 +31,7 @@ class aiohttp_filter(logging.Filter):
confserverlog = logging.getLogger("confserver") confserverlog = logging.getLogger("confserver")
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) #Add logging filter above to aiohttp.access
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter())
class EcoVacs_Login: class EcoVacs_Login:
accessToken = "" accessToken = ""
@ -61,42 +58,6 @@ class ConfServer:
self.run_async = False self.run_async = False
self.app = None 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): def confserver_app(self):
self.app = web.Application() self.app = web.Application()
@ -197,6 +158,7 @@ class ConfServer:
async def start_server(self): async def start_server(self):
try: try:
confserverlog.info("Starting ConfServer at {}:{}".format(self.address[0], self.address[1]))
runner = web.AppRunner(self.app) runner = web.AppRunner(self.app)
await runner.setup() await runner.setup()
@ -352,6 +314,7 @@ class ConfServer:
def check_token(self, apptype, countrycode, user, token): def check_token(self, apptype, countrycode, user, token):
try:
if bumper.check_token(user["userid"], token): if bumper.check_token(user["userid"], token):
if "global_" in apptype: #EcoVacs Home if "global_" in apptype: #EcoVacs Home
@ -392,16 +355,27 @@ class ConfServer:
} }
return web.json_response(body) return web.json_response(body)
except Exception as e:
confserverlog.exception("{}".format(e))
def generate_token(self, user): def generate_token(self, user):
try:
tmpaccesstoken = uuid.uuid4().hex tmpaccesstoken = uuid.uuid4().hex
bumper.user_add_token(user["userid"], tmpaccesstoken) bumper.user_add_token(user["userid"], tmpaccesstoken)
return tmpaccesstoken return tmpaccesstoken
except Exception as e:
confserverlog.exception("{}".format(e))
def generate_authcode(self, user, countrycode, token): def generate_authcode(self, user, countrycode, token):
try:
tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex) tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex)
bumper.user_add_authcode(user["userid"], token, tmpauthcode) bumper.user_add_authcode(user["userid"], token, tmpauthcode)
return tmpauthcode return tmpauthcode
except Exception as e:
confserverlog.exception("{}".format(e))
def _auth_any(self, devid, apptype, country, request): def _auth_any(self, devid, apptype, country, request):
try: try:
user_devid = devid user_devid = devid
@ -816,7 +790,6 @@ class ConfServer:
confserverlog.exception("{}".format(e)) confserverlog.exception("{}".format(e))
async def handle_getProductIotMap(self, request): async def handle_getProductIotMap(self, request):
user_devid = request.match_info.get("devid", "")
try: try:
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
@ -910,13 +883,23 @@ class ConfServer:
if todo == "FindBest": if todo == "FindBest":
service = postbody["service"] service = postbody["service"]
if service == "EcoMsgNew": if service == "EcoMsgNew":
srvip = socket.gethostbyname(socket.gethostname())
srvport = 5223
confserverlog.info(
"Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport)
)
body = { body = {
"result": "ok", "result": "ok",
"ip": socket.gethostbyname(socket.gethostname()), "ip": srvip,
"port": 5223, "port": srvport,
} }
elif service == "EcoUpdate": 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": elif todo == "loginByItToken":
if "userId" in postbody: if "userId" in postbody:
@ -996,6 +979,7 @@ class ConfServer:
for bot in bots: for bot in bots:
if bot["class"] != "": if bot["class"] != "":
b = bumper.bot_toEcoVacsHome_JSON(bot) b = bumper.bot_toEcoVacsHome_JSON(bot)
if not b is None: #Happens if the bot isn't on the EcoVacs Home list
botlist.append(json.loads(b)) botlist.append(json.loads(b))
body = { body = {
@ -1036,9 +1020,12 @@ class ConfServer:
if todo == "FindBest": if todo == "FindBest":
service = postbody["service"] service = postbody["service"]
if service == "EcoMsgNew": if service == "EcoMsgNew":
srvip = socket.gethostbyname(socket.gethostname()) 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 = json.dumps(msgserver)
msgserver = msgserver.replace( msgserver = msgserver.replace(
" ", "" " ", ""
@ -1095,11 +1082,16 @@ class ConfServer:
if did != "": if did != "":
confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did) bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
body = "" body = ""
retcmd = await self.helperbot.send_command(json_body, randomid) retcmd = await self.helperbot.send_command(json_body, randomid)
confserverlog.debug(
"Send Bot - {}".format(json_body)
)
confserverlog.debug(
"Bot Response - {}".format(body)
)
logs = [] logs = []
logsroot = ET.fromstring(retcmd["resp"]) logsroot = ET.fromstring(retcmd["resp"])
if logsroot.attrib["ret"] == "ok": if logsroot.attrib["ret"] == "ok":
@ -1148,13 +1140,15 @@ class ConfServer:
did = json_body["toId"] did = json_body["toId"]
if did != "": if did != "":
confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did) bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
retcmd = await self.helperbot.send_command(json_body, randomid) retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd body = retcmd
confserverlog.debug( 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) return web.json_response(body)
else: else:
@ -1181,7 +1175,7 @@ class ConfServer:
confserverlog.exception("{}".format(e)) confserverlog.exception("{}".format(e))
async def handle_dim_devmanager(self, request): async def handle_dim_devmanager(self, request): #Used in EcoVacs Home App
try: try:
json_body = json.loads(await request.text()) json_body = json.loads(await request.text())
@ -1191,13 +1185,15 @@ class ConfServer:
did = json_body["toId"] did = json_body["toId"]
if did != "": if did != "":
confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did) bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
retcmd = await self.helperbot.send_command(json_body, randomid) retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd body = retcmd
confserverlog.debug( 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) return web.json_response(body)
else: else:
@ -1223,13 +1219,13 @@ class ConfServer:
except Exception as e: except Exception as e:
confserverlog.exception("{}".format(e)) confserverlog.exception("{}".format(e))
def disconnect(self): async def disconnect(self):
try: try:
confserverlog.info("shutting down") confserverlog.info("shutting down")
if self.run_async: if self.run_async:
self.confthread.join() self.confthread.join()
else: else:
self.app.shutdown() await self.app.shutdown()
except Exception as e: except Exception as e:
confserverlog.exception("{}".format(e)) confserverlog.exception("{}".format(e))

View file

@ -8,13 +8,13 @@ from hbmqtt.broker import Broker
from hbmqtt.client import MQTTClient from hbmqtt.client import MQTTClient
from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2 from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2
import pkg_resources import pkg_resources
import contextvars
import time import time
from threading import Thread from threading import Thread
import ssl import ssl
import bumper import bumper
import json import json
from datetime import datetime, timedelta from datetime import datetime, timedelta
import bumper
helperbotlog = logging.getLogger("helperbot") helperbotlog = logging.getLogger("helperbot")
mqttserverlog = logging.getLogger("mqttserver") mqttserverlog = logging.getLogger("mqttserver")
@ -40,39 +40,16 @@ class MQTTHelperBot:
): ):
self.address = address self.address = address
self.client_id = "helper1@bumper/helper1" self.client_id = "helper1@bumper/helper1"
self.command_responses = contextvars.ContextVar("command_responses", default=[]) self.command_responses = []
self.helperthread = None 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): async def start_helper_bot(self):
try: try:
self.Client = MQTTClient(
client_id=self.client_id, config={"check_hostname": False}
)
await self.Client.connect( await self.Client.connect(
"mqtts://{}:{}/".format(self.address[0], self.address[1]), "mqtts://{}:{}/".format(self.address[0], self.address[1]),
cafile=bumper.ca_cert, cafile=bumper.ca_cert,
@ -81,8 +58,10 @@ class MQTTHelperBot:
[ [
("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0), ("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0),
("iot/p2p/+", QOS_0), ("iot/p2p/+", QOS_0),
("iot/atr/+", QOS_0),
] ]
) )
asyncio.create_task(self.get_msg())
except Exception as e: except Exception as e:
helperbotlog.exception("{}".format(e)) helperbotlog.exception("{}".format(e))
@ -92,29 +71,33 @@ class MQTTHelperBot:
while True: while True:
message = await self.Client.deliver_message() 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": 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(), "time": time.time(),
"topic": message.topic, "topic": message.topic,
"payload": str(message.data.decode("utf-8")), "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 # Cleanup "expired messages" > 60 seconds from time
for msg in cresp: for msg in self.command_responses:
expire_time = ( expire_time = (
datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10) datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10)
).timestamp() ).timestamp()
if time.time() > expire_time: if time.time() > expire_time:
# helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time)) helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
cresp.remove(msg) self.command_responses.remove(msg)
self.command_responses.set(cresp)
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
except Exception as e: except Exception as e:
helperbotlog.exception("{}".format(e)) helperbotlog.exception("{}".format(e))
@ -126,9 +109,8 @@ class MQTTHelperBot:
while time.time() < t_end: while time.time() < t_end:
await asyncio.sleep(0.1) await asyncio.sleep(0.1)
responses = self.command_responses.get() if len(self.command_responses) > 0:
if len(responses) > 0: for msg in self.command_responses:
for msg in responses:
topic = str(msg["topic"]).split("/") topic = str(msg["topic"]).split("/")
if topic[6] == "helper1" and topic[10] == requestid: if topic[6] == "helper1" and topic[10] == requestid:
# helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload'])) # helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
@ -137,9 +119,7 @@ class MQTTHelperBot:
else: else:
resppayload = str(msg["payload"]) resppayload = str(msg["payload"])
resp = {"id": requestid, "ret": "ok", "resp": resppayload} resp = {"id": requestid, "ret": "ok", "resp": resppayload}
cresp = self.command_responses.get() self.command_responses.remove(msg)
cresp.remove(msg)
self.command_responses.set(cresp)
return resp return resp
return {"id": requestid, "errno": 500, "ret": "fail", "debug": "wait for response timed out"} return {"id": requestid, "errno": 500, "ret": "fail", "debug": "wait for response timed out"}
@ -173,6 +153,7 @@ class MQTTHelperBot:
except Exception as e: except Exception as e:
helperbotlog.exception("{}".format(e)) helperbotlog.exception("{}".format(e))
return {}
class MQTTServer: class MQTTServer:
@ -180,6 +161,7 @@ class MQTTServer:
async def broker_coro(self): async def broker_coro(self):
try: try:
mqttserverlog.info("Starting MQTT Server at {}:{}".format(self.address[0], self.address[1]))
broker = hbmqtt.broker.Broker(config=self.default_config) broker = hbmqtt.broker.Broker(config=self.default_config)
await broker.start() await broker.start()
@ -238,32 +220,6 @@ class MQTTServer:
except Exception as e: except Exception as e:
mqttserverlog.exception("{}".format(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: class BumperMQTTServer_Plugin:
def __init__(self, context): def __init__(self, context):
@ -300,11 +256,9 @@ class BumperMQTTServer_Plugin:
client_id = session.client_id client_id = session.client_id
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
# If this isn't a fake user (fuid) then add as a bot if not ( # if ecouser or bumper aren't in details it is a bot
if not ( "ecouser" in didsplit[1]
str(didsplit[0]).startswith("fuid") or "bumper" in didsplit[1]):
or str(didsplit[0]).startswith("helper")
):
tmpbotdetail = str(didsplit[1]).split("/") tmpbotdetail = str(didsplit[1]).split("/")
bumper.bot_add( bumper.bot_add(
username, username,

769
bumper/xmpp_old_client.py Normal file
View file

@ -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(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.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(
'<iq id="{}" to="{}@{}/{}" from="rl.ecorobot.net" type="result"/>'.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("<query", '<query xmlns="com:ctl"')
if client.type == self.BOT:
if client.uid.lower() in ctl_to.lower():
xmppserverlog.info(
"Sending ctl to bot: {}".format(rxmlstring)
)
client.send(rxmlstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
def _handle_ping(self, xml, data):
try:
if xml.get("to").find("@") == -1: # No to address
# Ping to server - respond
pingresp = '<iq type="result" id="{}" from="{}" />'.format(
xml.get("id"), xml.get("to")
)
# xmppserverlog.debug("Server Ping resp: {}".format(pingresp))
self.send(pingresp)
else:
pingto = xml.get("to")
pingfrom = self.bumper_jid
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("<query", '<query xmlns="com:ctl"')
for client in XMPPServer.clients:
if (
client.bumper_jid != self.bumper_jid
and client.state == client.READY
):
if pingto.lower() in client.bumper_jid.lower():
pingstring = '<iq type="result" id="{}" from="{}" to="{}" />'.format(
xml.get("id"), pingfrom, pingto
)
xmppserverlog.debug(
"ping from {} to {}".format(pingfrom, pingto)
)
client.send(pingstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
def _handle_result(self, xml, data):
try:
ctl_to = xml.get("to")
xml.attrib["from"] = self.bumper_jid
if (
"errno='103' error='permission denied," in data
): # No permissions, usually if bot was last on Ecovac network
if self.type == self.BOT:
xquery = xml.getchildren()
ctl = xquery[0].getchildren()
ctlerr = ctl[0].attrib["error"]
adminuser = ctlerr.replace("permission denied, please contact ", "")
adminuser = adminuser.replace(" ", "")
if not (
adminuser.startswith("fuid_") or bumper.use_auth
): # if not fuid_ then its ecovacs OR ignore bumper auth
# TODO: Implement auth later, should this user have access to bot?
# Add user jid to bot
newuser = ctl_to.split("/")[0]
adduser = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="AddUser" id="0000" jid="{}" /></query></iq>'.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 = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="SetAC" id="1111" jid="{}"><acs><ac name="userman" allow="1"/><ac name="setting" allow="1"/><ac name="clean" allow="1"/></acs></ctl></query></iq>'.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(
'<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetUserInfo" id="4444" /><UserInfos/></query></iq>'.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("<query", '<query xmlns="com:ctl"')
if self.type == self.BOT:
if ctl_to == "de.ecorobot.net": # Send to all clients
xmppserverlog.debug(
"Sending to all clients because of de: {}".format(
rxmlstring
)
)
for client in XMPPServer.clients:
client.send(rxmlstring)
if xml.get("to").find("@") == -1: # No to address
ctl_to = xml.get("to")
else:
ctl_to = "{}@ecouser.net".format(ctl_to.split("@")[0])
for client in XMPPServer.clients:
if (
client.bumper_jid != self.bumper_jid
and client.state == client.READY
):
if not "@" in ctl_to: # No user@, send to all clients?
# TODO: Revisit later, this may be wrong
client.send(rxmlstring)
elif (
client.uid.lower() in ctl_to.lower()
): # If client matches TO=
xmppserverlog.debug(
"Sending from {} to client {}: {}".format(
self.uid, client.uid, rxmlstring
)
)
client.send(rxmlstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
def _handle_connect(self, data, xml=None):
try:
if self.state == self.CONNECT:
if xml == None:
# Client first connecting, send our features
if data.decode("utf-8").find("jabber:client") > -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(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id
)
)
# with STARTTLS
# self.send('<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns:tls="http://www.ietf.org/rfc/rfc2595.txt" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(XMPPServer.server_id))
time.sleep(0.25)
# send authentication support for iq-auth (fallback) and SASL
self.send(
'<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/><mechanisms xmlns="urn:ietf:params:xml:ns:xmpp-sasl"><mechanism>PLAIN</mechanism></mechanisms></stream:features>'
)
# self.send('<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/></stream:features>')
else:
self.send("</stream>")
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(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id
)
)
time.sleep(0.25)
# session
self.send(
'<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>'
)
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(
'<iq type="result" id="{}"><query xmlns="jabber:iq:auth"><username/><password/></query></iq>'.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('<iq type="result" id="{}"/>'.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('<iq type="result" id="{}"/>'.format(xml.get("id")))
else:
# Failed auth
self.send(
'<iq type="error" id="{}"><error code="401" type="auth"><not-authorized xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.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 xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # 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 xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Success
else:
# Failed to authenticate
self.send(
'<response xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # 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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.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 = '<iq type="result" id="{}" />'.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(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
)
# If it is a BOT, send extras
if self.type == self.BOT:
# get device info
self.send(
'<iq type="set" id="14" to="{}" from="{}"><query xmlns="com:ctl"><ctl td="GetDeviceInfo"/></query></iq>'.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(
'<presence to="{}"> dummy </presence>'.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(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
)
except Exception as e:
xmppserverlog.exception("{}".format(e))
def _parse_data(self, data):
if data.decode("utf-8").startswith(
"<?xml"
): # Strip <?xml and add artificial root
newdata = (
re.sub(r"(<\?xml[^>]+\?>)", r"<root>", data.decode("utf-8")) + "</root>"
)
else:
newdata = "<root>{}</root>".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 "<stream:stream " in newdata:
if self.state == self.CONNECT or self.state == self.INIT:
self._handle_connect(newdata.encode("utf-8"))
else:
if not (newdata == "" or newdata == " "):
xmppserverlog.error(
"xml parse error - {} - {}".format(newdata, e)
)
elif "not well-formed (invalid token)" in e.msg:
# If a lone </stream:stream> - client is signalling end of session/disconnect
if not "</stream:stream>" in newdata:
xmppserverlog.error("xml parse error - {} - {}".format(newdata, e))
else:
self.send("</stream:stream>") # Close stream
else:
if "<stream:stream" in newdata: # Handle start stream and connect
if self.state == self.CONNECT or self.state == self.INIT:
xmppserverlog.debug(
"Handling connect data - {}".format(newdata)
)
self._handle_connect(newdata.encode("utf-8"))
else:
if not "</stream:stream>" in newdata:
xmppserverlog.error(
"xml parse error - {} - {}".format(newdata, e)
)
else:
self.send("</stream:stream>") # 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))

View file

@ -4,8 +4,8 @@ from threading import Thread
import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET
import base64 import base64
import ssl import ssl
import contextvars
import bumper import bumper
import asyncio, functools
xmppserverlog = logging.getLogger("xmppserver") xmppserverlog = logging.getLogger("xmppserver")
@ -20,87 +20,39 @@ class XMPPServer:
# Initialize bot server # Initialize bot server
self.address = address self.address = address
def run(self, run_async=False): async def async_server(self):
if run_async: xmppserverlog.info(
xmppserverlog.debug("Starting XMPPServer Thread: 1") "Starting XMPP Server at {}:{}".format(self.address[0], self.address[1])
self.xmppthread = Thread(name="XMPPServer_Thread", target=self.run_server) )
self.xmppthread.setDaemon(True) server = await asyncio.start_server(
self.xmppthread.start() self.accept_client, self.address[0], self.address[1]
)
else: await server.serve_forever()
# self.clients = {} # task -> (reader, writer)
def accept_client(self, client_reader, client_writer):
try: try:
self.run_server() aclient = XMPPAsyncClient(client_reader, client_writer)
except KeyboardInterrupt: task = asyncio.Task(aclient.handle_async_client())
self.disconnect() aclient._async_task = task
self.clients.append(aclient)
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
)
self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
def client_done(aclient, task):
try: try:
self.socket.bind(self.address) self.clients.remove(aclient)
self.socket.listen(5) 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))
xmppserverlog.debug( clientaddr = client_writer.get_extra_info("peername")
"listening on {}:{}".format(self.address[0], self.address[1]) xmppserverlog.debug("New Connection from {}:{}".format(clientaddr[0],clientaddr[1]))
) task.add_done_callback(functools.partial(client_done, aclient))
while not self.exit_flag:
connection, client_address = self.socket.accept()
# 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)
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.error("{}".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()
def disconnect(self): def disconnect(self):
try: try:
@ -112,43 +64,9 @@ class XMPPServer:
xmppserverlog.debug("shutting down") xmppserverlog.debug("shutting down")
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.error("{}".format(e))
def remove_client_byip(self, ip): class XMPPAsyncClient:
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):
IDLE = 0 IDLE = 0
CONNECT = 1 CONNECT = 1
INIT = 2 INIT = 2
@ -158,52 +76,70 @@ class Client(threading.Thread):
UNKNOWN = 0 UNKNOWN = 0
BOT = 1 BOT = 1
CONTROLLER = 2 CONTROLLER = 2
_async_task = None
def __init__(self, thread_id, connection, client_address): def __init__(self, client_reader, client_writer):
threading.Thread.__init__(self)
self.id = thread_id
self.name = "XMPP_Client_{}".format(client_address[0])
self.type = self.UNKNOWN self.type = self.UNKNOWN
self.state = self.IDLE self.state = self.IDLE
self.connection = connection self.address = client_writer.get_extra_info("peername")
self.address = client_address[0] self.client_reader = client_reader
self.client_writer = client_writer
self.clientresource = "" self.clientresource = ""
self.devclass = "" self.devclass = ""
self.bumper_jid = "" self.bumper_jid = ""
self.uid = "" 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 self.log_incoming_data = True # Set to true to log sends
xmppserverlog.debug( xmppserverlog.debug("new client with ip {}".format(self.address))
"new client thread init for 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: try:
if not self.connection._closed: # if not self.connection._closed:
if self.log_sent_message: if self.log_sent_message:
xmppserverlog.debug("send {} - {}".format(self.address, command)) xmppserverlog.debug("send to ({}:{} | {}) - {}".format(self.address[0], self.address[1], self.bumper_jid, command))
self.connection.send(command.encode())
self.client_writer.write(command.encode())
await self.client_writer.drain()
except BrokenPipeError as e: except BrokenPipeError as e:
xmppserverlog.debug("{}".format(e)) #xmppserverlog.debug("{}".format(e))
self._set_state("DISCONNECT") await self._set_state("DISCONNECT")
except ConnectionResetError as e: except ConnectionResetError as e:
xmppserverlog.debug("{}".format(e)) #xmppserverlog.debug("{}".format(e))
self._set_state("DISCONNECT") await self._set_state("DISCONNECT")
except ConnectionAbortedError as e: except ConnectionAbortedError as e:
xmppserverlog.debug("{}".format(e)) #xmppserverlog.debug("{}".format(e))
self._set_state("DISCONNECT") await self._set_state("DISCONNECT")
except OSError as e: except OSError as e:
xmppserverlog.debug("{}".format(e)) xmppserverlog.error("{}".format(e))
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _disconnect(self): async def _disconnect(self):
try: try:
bot = bumper.bot_get(self.uid) bot = bumper.bot_get(self.uid)
@ -214,23 +150,23 @@ class Client(threading.Thread):
if client: if client:
bumper.client_set_xmpp(client["resource"], False) bumper.client_set_xmpp(client["resource"], False)
self.connection.close() self.client_writer.close()
except Exception as e: 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: try:
if tag[0] == "{": if tag[0] == "{":
_, _, tag = tag[1:].partition("}") _, _, tag = tag[1:].partition("}")
return tag return tag
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.error("{}".format(e))
def _set_state(self, state): async def _set_state(self, state):
try: try:
new_state = getattr(Client, state) new_state = getattr(XMPPAsyncClient, state)
if self.state > new_state: if self.state > new_state:
raise Exception( raise Exception(
"{} illegal state change {}->{}".format( "{} 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 self.state = new_state
if new_state == 5: if new_state == 5:
self._disconnect() await self._disconnect()
except Exception as e: 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: try:
if "roster" in data: if "roster" in data:
# Return not-implemented for roster # Return not-implemented for roster
self.send( await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format( '<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id") xml.get("id")
) )
) )
return return
if "disco#items" in data:
# Return not-implemented for disco#items
await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
))
return
if "disco#info" in data:
# Return not-implemented for disco#info
await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
)
)
return
if xml.get("type") == "set": if xml.get("type") == "set":
if ( if (
"com:sf" in data and xml.get("to") == "rl.ecorobot.net" "com:sf" in data and xml.get("to") == "rl.ecorobot.net"
): # Android bind? Not sure what this does yet. ): # Android bind? Not sure what this does yet.
self.send( await self.send(
'<iq id="{}" to="{}@{}/{}" from="rl.ecorobot.net" type="result"/>'.format( '<iq id="{}" to="{}@{}/{}" from="rl.ecorobot.net" type="result"/>'.format(
xml.get("id"), xml.get("id"),
self.uid, self.uid,
@ -273,7 +227,7 @@ class Client(threading.Thread):
) )
) )
if xml[0][0]: if len(xml[0]) > 0:
ctl = xml[0][0] ctl = xml[0][0]
if ctl.get("admin") and self.type == self.BOT: if ctl.get("admin") and self.type == self.BOT:
xmppserverlog.debug( xmppserverlog.debug(
@ -299,15 +253,15 @@ class Client(threading.Thread):
if client.type == self.BOT: if client.type == self.BOT:
if client.uid.lower() in ctl_to.lower(): if client.uid.lower() in ctl_to.lower():
xmppserverlog.info( xmppserverlog.debug(
"Sending ctl to bot: {}".format(rxmlstring) "Sending ctl to bot: {}".format(rxmlstring)
) )
client.send(rxmlstring) await client.send(rxmlstring)
except Exception as e: 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: try:
if xml.get("to").find("@") == -1: # No to address if xml.get("to").find("@") == -1: # No to address
# Ping to server - respond # Ping to server - respond
@ -315,7 +269,7 @@ class Client(threading.Thread):
xml.get("id"), xml.get("to") xml.get("id"), xml.get("to")
) )
# xmppserverlog.debug("Server Ping resp: {}".format(pingresp)) # xmppserverlog.debug("Server Ping resp: {}".format(pingresp))
self.send(pingresp) await self.send(pingresp)
else: else:
pingto = xml.get("to") pingto = xml.get("to")
@ -326,8 +280,9 @@ class Client(threading.Thread):
# clean up string to remove namespaces added by ET # clean up string to remove namespaces added by ET
pingstring = pingstring.replace("xmlns:ns0=", "xmlns=") pingstring = pingstring.replace("xmlns:ns0=", "xmlns=")
pingstring = pingstring.replace("ns0:", "") pingstring = pingstring.replace("ns0:", "")
pingstring = pingstring.replace('iq xmlns="com:ctl"', "iq") pingstring = pingstring.replace('iq xmlns="urn:xmpp:ping"', "iq")
pingstring = pingstring.replace("<query", '<query xmlns="com:ctl"') pingstring = pingstring.replace("<ping", '<ping xmlns="urn:xmpp:ping"')
for client in XMPPServer.clients: for client in XMPPServer.clients:
if ( if (
@ -335,18 +290,27 @@ class Client(threading.Thread):
and client.state == client.READY and client.state == client.READY
): ):
if pingto.lower() in client.bumper_jid.lower(): if pingto.lower() in client.bumper_jid.lower():
pingstring = '<iq type="result" id="{}" from="{}" to="{}" />'.format( #pingstring = '<iq type="result" id="{}" from="{}" to="{}" />'.format(
xml.get("id"), pingfrom, pingto # xml.get("id"), pingfrom, pingto
) #)
xmppserverlog.debug( #xmppserverlog.debug(
"ping from {} to {}".format(pingfrom, pingto) # "ping from {} to {} with {}".format(pingfrom, pingto, pingstring)
) #)
client.send(pingstring)
await client.send(pingstring)
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_result(self, xml, data):
async def schedule_ping(self, time):
if not self.state == 5: #disconnected
pingstring = "<iq from='{}' to='{}' id='s2c1' type='get'><ping xmlns='urn:xmpp:ping'/></iq>".format(XMPPServer.server_id, self.bumper_jid)
await self.send(pingstring)
await asyncio.sleep(time)
asyncio.Task(self.schedule_ping(time))
async def _handle_result(self, xml, data):
try: try:
ctl_to = xml.get("to") ctl_to = xml.get("to")
xml.attrib["from"] = self.bumper_jid xml.attrib["from"] = self.bumper_jid
@ -360,7 +324,7 @@ class Client(threading.Thread):
adminuser = ctlerr.replace("permission denied, please contact ", "") adminuser = ctlerr.replace("permission denied, please contact ", "")
adminuser = adminuser.replace(" ", "") adminuser = adminuser.replace(" ", "")
if not ( if not (
adminuser.startswith("fuid_") or bumper.use_auth adminuser.startswith("fuid_") or adminuser.startswith("fusername_") or bumper.use_auth
): # if not fuid_ then its ecovacs OR ignore bumper auth ): # if not fuid_ then its ecovacs OR ignore bumper auth
# TODO: Implement auth later, should this user have access to bot? # TODO: Implement auth later, should this user have access to bot?
@ -370,17 +334,17 @@ class Client(threading.Thread):
uuid.uuid4(), adminuser, self.bumper_jid, newuser uuid.uuid4(), adminuser, self.bumper_jid, newuser
) )
xmppserverlog.debug("Add User: {}".format(adduser)) xmppserverlog.debug("Add User: {}".format(adduser))
self.send(adduser) await self.send(adduser)
# Add user ACs - Manage users, settings, and clean (full access) # Add user ACs - Manage users, settings, and clean (full access)
adduseracs = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="SetAC" id="1111" jid="{}"><acs><ac name="userman" allow="1"/><ac name="setting" allow="1"/><ac name="clean" allow="1"/></acs></ctl></query></iq>'.format( adduseracs = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="SetAC" id="1111" jid="{}"><acs><ac name="userman" allow="1"/><ac name="setting" allow="1"/><ac name="clean" allow="1"/></acs></ctl></query></iq>'.format(
uuid.uuid4(), adminuser, self.bumper_jid, newuser uuid.uuid4(), adminuser, self.bumper_jid, newuser
) )
xmppserverlog.debug("Add User ACs: {}".format(adduseracs)) xmppserverlog.debug("Add User ACs: {}".format(adduseracs))
self.send(adduseracs) await self.send(adduseracs)
# GetUserInfo - Just to confirm it set correctly # GetUserInfo - Just to confirm it set correctly
self.send( await self.send(
'<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetUserInfo" id="4444" /><UserInfos/></query></iq>'.format( '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetUserInfo" id="4444" /><UserInfos/></query></iq>'.format(
uuid.uuid4(), adminuser, self.bumper_jid uuid.uuid4(), adminuser, self.bumper_jid
) )
@ -401,7 +365,7 @@ class Client(threading.Thread):
) )
) )
for client in XMPPServer.clients: for client in XMPPServer.clients:
client.send(rxmlstring) await client.send(rxmlstring)
if xml.get("to").find("@") == -1: # No to address if xml.get("to").find("@") == -1: # No to address
ctl_to = xml.get("to") ctl_to = xml.get("to")
@ -415,7 +379,7 @@ class Client(threading.Thread):
): ):
if not "@" in ctl_to: # No user@, send to all clients? if not "@" in ctl_to: # No user@, send to all clients?
# TODO: Revisit later, this may be wrong # TODO: Revisit later, this may be wrong
client.send(rxmlstring) await client.send(rxmlstring)
elif ( elif (
client.uid.lower() in ctl_to.lower() client.uid.lower() in ctl_to.lower()
@ -425,12 +389,12 @@ class Client(threading.Thread):
self.uid, client.uid, rxmlstring self.uid, client.uid, rxmlstring
) )
) )
client.send(rxmlstring) await client.send(rxmlstring)
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_connect(self, data, xml=None): async def _handle_connect(self, data, xml=None):
try: try:
if self.state == self.CONNECT: if self.state == self.CONNECT:
@ -443,30 +407,32 @@ class Client(threading.Thread):
self.devclass = data.decode("utf-8")[sc + 4 : ec] self.devclass = data.decode("utf-8")[sc + 4 : ec]
# ack jabbr:client # ack jabbr:client
# no STARTTLS # no STARTTLS
self.send( await self.send(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format( '<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id XMPPServer.server_id
) )
) )
# with STARTTLS # with STARTTLS
# self.send('<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns:tls="http://www.ietf.org/rfc/rfc2595.txt" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(XMPPServer.server_id)) # await self.send('<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns:tls="http://www.ietf.org/rfc/rfc2595.txt" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(XMPPServer.server_id))
time.sleep(0.25)
await asyncio.sleep(0.25)
#time.sleep(0.25)
# send authentication support for iq-auth (fallback) and SASL # send authentication support for iq-auth (fallback) and SASL
self.send( await self.send(
'<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/><mechanisms xmlns="urn:ietf:params:xml:ns:xmpp-sasl"><mechanism>PLAIN</mechanism></mechanisms></stream:features>' '<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/><mechanisms xmlns="urn:ietf:params:xml:ns:xmpp-sasl"><mechanism>PLAIN</mechanism></mechanisms></stream:features>'
) )
# self.send('<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/></stream:features>') # await self.send('<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/></stream:features>')
else: else:
self.send("</stream>") await self.send("</stream>")
else: else:
if "jabber:iq:auth" in xml.tag: # Handle iq-auth if "jabber:iq:auth" in xml.tag: # Handle iq-auth
self._handle_iq_auth(xml) await self._handle_iq_auth(xml)
elif ( elif (
"urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag "urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag
): # Handle SASL Auth ): # Handle SASL Auth
self._handle_sasl_auth(xml) await self._handle_sasl_auth(xml)
else: else:
xmppserverlog.error("Couldn't handle: {}".format(xml)) xmppserverlog.error("Couldn't handle: {}".format(xml))
@ -475,33 +441,34 @@ class Client(threading.Thread):
# Client getting session after authentication # Client getting session after authentication
if data.decode("utf-8").find("jabber:client") > -1: if data.decode("utf-8").find("jabber:client") > -1:
# ack jabbr:client # ack jabbr:client
self.send( await self.send(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format( '<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id XMPPServer.server_id
) )
) )
time.sleep(0.25) await asyncio.sleep(0.25)
#time.sleep(0.25)
# session # session
self.send( await self.send(
'<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>' '<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>'
) )
else: # Handle init bind else: # Handle init bind
if len(xml): if len(xml):
child = self._tag_strip_uri(xml[0].tag) child = await self._tag_strip_uri(xml[0].tag)
else: else:
child = None child = None
if xml.tag == "iq": if xml.tag == "iq":
if child == "bind": if child == "bind":
self._handle_bind(xml) await self._handle_bind(xml)
else: else:
xmppserverlog.error("Couldn't handle: {}".format(xml)) xmppserverlog.error("Couldn't handle: {}".format(xml))
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_iq_auth(self, data): async def _handle_iq_auth(self, data):
try: try:
xml = ET.fromstring(data.decode("utf-8")) xml = ET.fromstring(data.decode("utf-8"))
ctl = xml[0][0] ctl = xml[0][0]
@ -512,7 +479,7 @@ class Client(threading.Thread):
and "auth}username" in ctl.tag and "auth}username" in ctl.tag
and self.type == self.UNKNOWN and self.type == self.UNKNOWN
): ):
self.send( await self.send(
'<iq type="result" id="{}"><query xmlns="jabber:iq:auth"><username/><password/></query></iq>'.format( '<iq type="result" id="{}"><query xmlns="jabber:iq:auth"><username/><password/></query></iq>'.format(
xml.get("id") xml.get("id")
) )
@ -527,6 +494,7 @@ class Client(threading.Thread):
xmlauth = xml[0].getchildren() xmlauth = xml[0].getchildren()
# uid = "" # uid = ""
password = "" password = ""
authcode = ""
resource = "" resource = ""
for aitem in xmlauth: for aitem in xmlauth:
if "username" in aitem.tag: if "username" in aitem.tag:
@ -540,17 +508,15 @@ class Client(threading.Thread):
self.clientresource = aitem.text self.clientresource = aitem.text
resource = self.clientresource resource = self.clientresource
if not self.uid.startswith("fuid"): if self.devclass: # if there is a devclass it is a bot
# Need sample data to see details here
bumper.bot_add("", self.uid, "", resource, "eco-legacy") 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 # Client authenticated, move to next state
self._set_state("INIT") await self._set_state("INIT")
# Successful auth # Successful auth
self.send('<iq type="result" id="{}"/>'.format(xml.get("id"))) await self.send('<iq type="result" id="{}"/>'.format(xml.get("id")))
else: else:
auth = False auth = False
@ -564,14 +530,16 @@ class Client(threading.Thread):
xmppserverlog.debug("client authenticated {}".format(self.uid)) xmppserverlog.debug("client authenticated {}".format(self.uid))
# Client authenticated, move to next state # Client authenticated, move to next state
self._set_state("INIT") await self._set_state("INIT")
# Successful auth # Successful auth
self.send('<iq type="result" id="{}"/>'.format(xml.get("id"))) await self.send(
'<iq type="result" id="{}"/>'.format(xml.get("id"))
)
else: else:
# Failed auth # Failed auth
self.send( await self.send(
'<iq type="error" id="{}"><error code="401" type="auth"><not-authorized xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format( '<iq type="error" id="{}"><error code="401" type="auth"><not-authorized xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id") xml.get("id")
) )
@ -594,12 +562,13 @@ class Client(threading.Thread):
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_sasl_auth(self, xml): async def _handle_sasl_auth(self, xml):
try: try:
saslauth = base64.b64decode(xml.text).decode("utf-8").split("/") saslauth = base64.b64decode(xml.text).decode("utf-8").split("/")
username = saslauth[0] username = saslauth[0]
username = saslauth[0].split("\x00")[1] username = saslauth[0].split("\x00")[1]
authcode = ""
self.uid = username self.uid = username
if len(saslauth) > 1: if len(saslauth) > 1:
resource = saslauth[1] resource = saslauth[1]
@ -611,18 +580,17 @@ class Client(threading.Thread):
if len(saslauth) > 2: if len(saslauth) > 2:
authcode = saslauth[2] authcode = saslauth[2]
if not self.uid.startswith("fuid"): if self.devclass: # if there is a devclass it is a bot
# Need sample data to see details here
bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy") bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy")
self.type = self.BOT self.type = self.BOT
xmppserverlog.info("bot authenticated {}".format(self.uid)) xmppserverlog.debug("bot authenticated {}".format(self.uid))
# Send response # Send response
self.send( await self.send(
'<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>' '<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Success ) # Success
# Client authenticated, move to next state # Client authenticated, move to next state
self._set_state("INIT") await self._set_state("INIT")
else: else:
auth = False auth = False
@ -637,23 +605,23 @@ class Client(threading.Thread):
xmppserverlog.debug("client authenticated {}".format(self.uid)) xmppserverlog.debug("client authenticated {}".format(self.uid))
# Client authenticated, move to next state # Client authenticated, move to next state
self._set_state("INIT") await self._set_state("INIT")
# Send response # Send response
self.send( await self.send(
'<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>' '<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Success ) # Success
else: else:
# Failed to authenticate # Failed to authenticate
self.send( await self.send(
'<response xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>' '<response xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Fail ) # Fail
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_bind(self, xml): async def _handle_bind(self, xml):
try: try:
bot = bumper.bot_get(self.uid) bot = bumper.bot_get(self.uid)
@ -671,7 +639,7 @@ class Client(threading.Thread):
self.bumper_jid = "{}@{}.ecorobot.net/atom".format( self.bumper_jid = "{}@{}.ecorobot.net/atom".format(
self.uid, self.devclass 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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format( res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(
xml.get("id"), self.bumper_jid xml.get("id"), self.bumper_jid
) )
@ -681,83 +649,84 @@ class Client(threading.Thread):
self.bumper_jid = "{}@{}/{}".format( self.bumper_jid = "{}@{}/{}".format(
self.uid, XMPPServer.server_id, self.clientresource self.uid, XMPPServer.server_id, self.clientresource
) )
xmppserverlog.debug( xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid))
"new client {} using resource {}".format(
self.uid, self.clientresource
)
)
res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format( res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(
xml.get("id"), self.bumper_jid xml.get("id"), self.bumper_jid
) )
else: else:
self.name = "XMPP_Client_{}_{}".format(self.uid, self.address) self.name = "XMPP_Client_{}_{}".format(self.uid, self.address)
self.bumper_jid = "{}@{}".format(self.uid, XMPPServer.server_id) 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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format( res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(
xml.get("id"), self.bumper_jid xml.get("id"), self.bumper_jid
) )
self._set_state("BIND") await self._set_state("BIND")
self.send(res) await self.send(res)
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_session(self, xml): async def _handle_session(self, xml):
try: try:
res = '<iq type="result" id="{}" />'.format(xml.get("id")) res = '<iq type="result" id="{}" />'.format(xml.get("id"))
self._set_state("READY") await self._set_state("READY")
self.send(res) await self.send(res)
asyncio.Task(self.schedule_ping(30))
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_presence(self, xml): async def _handle_presence(self, xml):
try: try:
if len(xml) and xml[0].tag == "status": if len(xml) and xml[0].tag == "status":
xmppserverlog.debug( 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 # Most likely a bot, possibly hello world in text
# Send dummy return # Send dummy return
self.send( await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid) '<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
) )
# If it is a BOT, send extras # If it is a BOT, send extras
if self.type == self.BOT: if self.type == self.BOT:
# get device info # get device info
self.send( await self.send(
'<iq type="set" id="14" to="{}" from="{}"><query xmlns="com:ctl"><ctl td="GetDeviceInfo"/></query></iq>'.format( '<iq type="set" id="14" to="{}" from="{}"><query xmlns="com:ctl"><ctl td="GetDeviceInfo"/></query></iq>'.format(
self.bumper_jid, XMPPServer.server_id self.bumper_jid, XMPPServer.server_id
) )
) )
else: else:
xmppserverlog.debug( 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": if xml.get("type") == "available":
xmppserverlog.debug( xmppserverlog.debug(
"client presence available - {} ".format( "client presence available - {} ".format(
ET.tostring(xml, encoding="utf-8") ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
) )
# Send dummy return # Send dummy return
self.send( await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid) '<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
) )
elif xml.get("type") == "unavailable": elif xml.get("type") == "unavailable":
xmppserverlog.debug( xmppserverlog.debug(
"client presence unavailable (DISCONNECT) - {} ".format( "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: else:
# Sometimes the android app sends these # Sometimes the android app sends these
xmppserverlog.debug( xmppserverlog.debug(
@ -766,14 +735,14 @@ class Client(threading.Thread):
) )
) )
# Send dummy return # Send dummy return
self.send( await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid) '<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
) )
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _parse_data(self, data): async def _parse_data(self, data):
if data.decode("utf-8").startswith( if data.decode("utf-8").startswith(
"<?xml" "<?xml"
@ -794,8 +763,8 @@ class Client(threading.Thread):
if item.tag == "iq": if item.tag == "iq":
if self.log_incoming_data: if self.log_incoming_data:
xmppserverlog.debug( xmppserverlog.debug(
"from {} - {}".format( "from ({}:{} | {}) - {}".format(
self.address, self.address[0],self.address[1],self.bumper_jid,
str( str(
ET.tostring(item, encoding="utf-8").decode( ET.tostring(item, encoding="utf-8").decode(
"utf-8" "utf-8"
@ -803,16 +772,16 @@ class Client(threading.Thread):
).replace("ns0:", ""), ).replace("ns0:", ""),
) )
) )
self._handle_iq(item, newdata) await self._handle_iq(item, newdata)
item.clear() item.clear()
elif "auth" in item.tag: elif "auth" in item.tag:
if "urn:ietf:params:xml:ns:xmpp-sasl" in item.tag: # SASL Auth if "urn:ietf:params:xml:ns:xmpp-sasl" in item.tag: # SASL Auth
self._handle_sasl_auth(item) await self._handle_sasl_auth(item)
item.clear() item.clear()
elif "presence" in item.tag: elif "presence" in item.tag:
self._handle_presence(item) await self._handle_presence(item)
item.clear() item.clear()
else: else:
@ -834,7 +803,7 @@ class Client(threading.Thread):
# Happens wth connect stream often # Happens wth connect stream often
if "<stream:stream " in newdata: if "<stream:stream " in newdata:
if self.state == self.CONNECT or self.state == self.INIT: if self.state == self.CONNECT or self.state == self.INIT:
self._handle_connect(newdata.encode("utf-8")) await self._handle_connect(newdata.encode("utf-8"))
else: else:
if not (newdata == "" or newdata == " "): if not (newdata == "" or newdata == " "):
xmppserverlog.error( xmppserverlog.error(
@ -846,7 +815,7 @@ class Client(threading.Thread):
if not "</stream:stream>" in newdata: if not "</stream:stream>" in newdata:
xmppserverlog.error("xml parse error - {} - {}".format(newdata, e)) xmppserverlog.error("xml parse error - {} - {}".format(newdata, e))
else: else:
self.send("</stream:stream>") # Close stream await self.send("</stream:stream>") # Close stream
else: else:
if "<stream:stream" in newdata: # Handle start stream and connect if "<stream:stream" in newdata: # Handle start stream and connect
@ -854,67 +823,48 @@ class Client(threading.Thread):
xmppserverlog.debug( xmppserverlog.debug(
"Handling connect data - {}".format(newdata) "Handling connect data - {}".format(newdata)
) )
self._handle_connect(newdata.encode("utf-8")) await self._handle_connect(newdata.encode("utf-8"))
else: else:
if not "</stream:stream>" in newdata: if not "</stream:stream>" in newdata:
xmppserverlog.error( xmppserverlog.error(
"xml parse error - {} - {}".format(newdata, e) "xml parse error - {} - {}".format(newdata, e)
) )
else: else:
self.send("</stream:stream>") # Close stream await self.send("</stream:stream>") # Close stream
self._set_state("DISCONNECT") await self._set_state("DISCONNECT")
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(e)) xmppserverlog.exception("{}".format(e))
def _handle_iq(self, xml, data): async def _handle_iq(self, xml, data):
try: try:
if len(xml): if len(xml):
child = self._tag_strip_uri(xml[0].tag) child = await self._tag_strip_uri(xml[0].tag)
else: else:
child = None child = None
if xml.tag == "iq": if xml.tag == "iq":
if child == "bind": if child == "bind":
self._handle_bind(xml) await self._handle_bind(xml)
elif child == "session": elif child == "session":
self._handle_session(xml) await self._handle_session(xml)
elif child == "ping": elif child == "ping":
self._handle_ping(xml, data) await self._handle_ping(xml, data)
elif child == "query": elif child == "query":
if self.type == self.BOT: if self.type == self.BOT:
self._handle_result(xml, data) await self._handle_result(xml, data)
else: else:
self._handle_ctl(xml, data) await self._handle_ctl(xml, data)
elif xml.get("type") == "result": elif xml.get("type") == "result":
if self.type == self.BOT: if self.type == self.BOT:
self._handle_result(xml, data) await self._handle_result(xml, data)
else: else:
self._handle_result(xml, data) await self._handle_result(xml, data)
elif xml.get("type") == "set": elif xml.get("type") == "set":
if self.type == self.BOT: if self.type == self.BOT:
self._handle_result(xml, data) await self._handle_result(xml, data)
else: else:
self._handle_result(xml, data) await self._handle_result(xml, data)
except Exception as e: except Exception as e:
xmppserverlog.exception("{}".format(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))

View file

@ -5,10 +5,19 @@ import bumper
import sys, socket import sys, socket
import time import time
import platform 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 args = sys.argv
listen_host = ""
if len(args) > 0: if len(args) > 0:
if "--debug" in args: if "--debug" in args:
@ -16,23 +25,22 @@ def main():
level=logging.DEBUG, level=logging.DEBUG,
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s", 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: else:
logging.basicConfig( logging.basicConfig(
level=logging.INFO, level=logging.INFO,
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s", 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 "--listen" in args:
if (len(args) - 1) >= (listen_host + 1): listen_host = args[args.index("--listen") + 1]
listen_host = args[listen_host+1]
else: if listen_host == "":
if platform.system() == "Darwin": # If a Mac, use 0.0.0.0 for listening if platform.system() == "Darwin": # If a Mac, use 0.0.0.0 for listening
listen_host = "0.0.0.0" listen_host = "0.0.0.0"
else: else:
listen_host = socket.gethostbyname(socket.gethostname()) listen_host = socket.gethostbyname(socket.gethostname())
#listen_host = "localhost" # Try this if the above doesn't work
conf_address_443 = (listen_host, 443) conf_address_443 = (listen_host, 443)
conf_address_8007 = (listen_host, 8007) conf_address_8007 = (listen_host, 8007)
@ -51,38 +59,39 @@ def main():
conf_address_8007, usessl=False, helperbot=mqtt_helperbot conf_address_8007, usessl=False, helperbot=mqtt_helperbot
) )
# add user # Start web servers
# users = bumper.bumper_users_var.get() conf_server.confserver_app()
# user1 = bumper.BumperUser('user1') task_conf_server = asyncio.create_task(conf_server.start_server())
# user1.add_device('devid') bumper.bumperlog.debug("task_conf_server added")
# user1.add_bot('bot_did') await task_conf_server
# users.append(user1)
# bumper.bumper_users_var.set(users)
# start xmpp server on port 5223 (sync) conf_server_2.confserver_app()
xmpp_server.run(run_async=True) # Start in new thread 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) # Start MQTT Server
mqtt_server.run(run_async=True) # Start in new thread 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
task_mqtt_helperbot = asyncio.create_task(mqtt_helperbot.start_helper_bot())
bumper.bumperlog.debug("task_mqtt_helperbot added")
await task_mqtt_helperbot
# start mqtt_helperbot (async) # Start XMPP Server
mqtt_helperbot.run(run_async=True) # Start in new thread task_xmpp_server = asyncio.create_task(xmpp_server.async_server())
bumper.bumperlog.debug("task_xmpp_server added")
# start conf server on port 443 (async) - Used for most https calls await task_xmpp_server
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
while True: while True:
try: try:
time.sleep(30) await asyncio.sleep(30)
bumper.revoke_expired_tokens() bumper.revoke_expired_tokens()
disconnected_clients = bumper.get_disconnected_xmpp_clients() #disconnected_clients = bumper.get_disconnected_xmpp_clients()
for client in disconnected_clients: #for client in disconnected_clients:
xmpp_server.remove_client_byuid(client["userid"]) # xmpp_server.remove_client_byuid(client["userid"])
except KeyboardInterrupt: except KeyboardInterrupt:
bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt") bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt")
@ -91,4 +100,4 @@ def main():
if __name__ == "__main__": if __name__ == "__main__":
main() asyncio.run(main())

View file

@ -18,7 +18,14 @@ def async_return(result):
return f return f
def test_disconnect(): 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(): def test_base():
if os.path.exists("tests/tmp.db"): if os.path.exists("tests/tmp.db"):