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