From ac8c4512fdd6ff15c0b6b41f6a8434fb18200c58 Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Fri, 27 Feb 2026 23:45:22 +0100 Subject: [PATCH] [FIX] Migrated to Tensamin Transport Protocol --- Cargo.lock | 267 +++++--- Cargo.toml | 8 +- src/data/communication.rs | 403 ----------- src/data/mod.rs | 1 - src/main.rs | 6 +- src/server/api.rs | 234 ++++--- src/server/mod.rs | 1 - src/server/omikron_connection.rs | 1068 +++++++----------------------- src/server/server.rs | 14 +- src/server/socket.rs | 129 ---- src/util/logger.rs | 87 ++- 11 files changed, 626 insertions(+), 1592 deletions(-) delete mode 100644 src/data/communication.rs delete mode 100644 src/data/mod.rs delete mode 100755 src/server/socket.rs diff --git a/Cargo.lock b/Cargo.lock index 7469a5e..40c314f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -21,6 +21,8 @@ dependencies = [ "color-eyre", "dashmap", "dotenv", + "epsilon-core", + "epsilon-native", "futures", "futures-util", "hex", @@ -32,16 +34,17 @@ dependencies = [ "json", "once_cell", "pnet", + "quinn", "rand 0.8.5", "rand_core 0.6.4", "reqwest", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-pemfile 2.2.0", "sha1", "sha2", "sqlx", - "strum", - "strum_macros", + "strum 0.27.2", + "strum_macros 0.27.2", "sysinfo", "tokio", "tokio-rustls", @@ -99,9 +102,9 @@ dependencies = [ [[package]] name = "actix-http" -version = "3.11.2" +version = "3.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7926860314cbe2fb5d1f13731e387ab43bd32bca224e82e6e2db85de0a3dba49" +checksum = "f860ee6746d0c5b682147b2f7f8ef036d4f92fe518251a3a35ffa3650eafdf0e" dependencies = [ "actix-codec", "actix-rt", @@ -149,9 +152,9 @@ dependencies = [ [[package]] name = "actix-router" -version = "0.5.3" +version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "13d324164c51f63867b57e73ba5936ea151b8a41a1d23d1031eeb9f70d0236f8" +checksum = "14f8c75c51892f18d9c46150c5ac7beb81c95f78c8b83a634d49f4ca32551fe7" dependencies = [ "bytestring", "cfg-if", @@ -231,9 +234,9 @@ dependencies = [ [[package]] name = "actix-web" -version = "4.12.1" +version = "4.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1654a77ba142e37f049637a3e5685f864514af11fcbc51cb51eb6596afe5b8d6" +checksum = "ff87453bc3b56e9b2b23c1cc0b1be8797184accf51d2abe0f8a33ec275d316bf" dependencies = [ "actix-codec", "actix-http", @@ -405,9 +408,9 @@ dependencies = [ [[package]] name = "anyhow" -version = "1.0.101" +version = "1.0.102" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5f0e0fee31ef5ed1ba1316088939cea399010ed7731dba877ed44aeb407a75ea" +checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" [[package]] name = "arbitrary" @@ -443,9 +446,9 @@ dependencies = [ [[package]] name = "async-executor" -version = "1.13.3" +version = "1.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "497c00e0fd83a72a79a39fcbd8e3e2f055d6f6c7e025f3b3d91f4f8e76527fb8" +checksum = "c96bf972d85afc50bf5ab8fe2d54d1586b4e0b46c97c50a0c9e71e2f7bcd812a" dependencies = [ "async-task", "concurrent-queue", @@ -503,7 +506,7 @@ dependencies = [ "futures-lite 2.6.1", "parking", "polling 3.11.0", - "rustix 1.1.3", + "rustix 1.1.4", "slab", "windows-sys 0.61.2", ] @@ -649,9 +652,9 @@ checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" [[package]] name = "aws-lc-rs" -version = "1.15.4" +version = "1.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7b7b6141e96a8c160799cc2d5adecd5cbbe5054cb8c7c4af53da0f83bb7ad256" +checksum = "d9a7b350e3bb1767102698302bc37256cbd48422809984b98d292c40e2579aa9" dependencies = [ "aws-lc-sys", "zeroize", @@ -817,9 +820,9 @@ dependencies = [ [[package]] name = "bumpalo" -version = "3.19.1" +version = "3.20.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5dd9dc738b7a8311c7ade152424974d8115f2cdad61e8dab8dac9f2362298510" +checksum = "5d20789868f4b01b2f2caec9f5c4e0213b41e3e5702a50157d699ae31ced2fcb" [[package]] name = "byteorder" @@ -1113,9 +1116,9 @@ checksum = "d7a1e2f27636f116493b8b860f5546edb47c8d8f8ea73e1d2a20be88e28d1fea" [[package]] name = "deflate64" -version = "0.1.10" +version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26bf8fc351c5ed29b5c2f0cbbac1b209b74f60ecd62e675a998df72c49af5204" +checksum = "807800ff3288b621186fe0a8f3392c4652068257302709c24efd918c3dffcdc2" [[package]] name = "der" @@ -1130,9 +1133,9 @@ dependencies = [ [[package]] name = "deranged" -version = "0.5.6" +version = "0.5.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc3dc5ad92c2e2d1c193bbbbdf2ea477cb81331de4f3103f267ca18368b988c4" +checksum = "7cd812cc2bc1d69d4764bd80df88b4317eaef9e773c75226407d9bc0876b211c" dependencies = [ "powerfmt", ] @@ -1241,6 +1244,32 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "epsilon-core" +version = "0.1.0" +source = "git+https://github.com/Tensamin/Epsilon.git#dc17684bc9308c79559e0d4648c6e730cc567c89" +dependencies = [ + "byteorder", + "rand 0.8.5", + "strum 0.28.0", + "strum_macros 0.28.0", +] + +[[package]] +name = "epsilon-native" +version = "0.1.0" +source = "git+https://github.com/Tensamin/Epsilon.git#dc17684bc9308c79559e0d4648c6e730cc567c89" +dependencies = [ + "anyhow", + "async-trait", + "bytes", + "epsilon-core", + "quinn", + "rustls 0.23.37", + "thiserror 2.0.18", + "tokio", +] + [[package]] name = "equivalent" version = "1.0.2" @@ -1305,6 +1334,18 @@ dependencies = [ "once_cell", ] +[[package]] +name = "fastbloom" +version = "0.14.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4e7f34442dbe69c60fe8eaf58a8cafff81a1f278816d8ab4db255b3bef4ac3c4" +dependencies = [ + "getrandom 0.3.4", + "libm", + "rand 0.9.2", + "siphasher", +] + [[package]] name = "fastrand" version = "1.9.0" @@ -1398,9 +1439,9 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" [[package]] name = "futures" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "65bc07b1a8bc7c85c5f2e110c476c7389b4554ba72af57d8445ea63a576b0876" +checksum = "8b147ee9d1f6d097cef9ce628cd2ee62288d963e16fb287bd9286455b241382d" dependencies = [ "futures-channel", "futures-core", @@ -1413,9 +1454,9 @@ dependencies = [ [[package]] name = "futures-channel" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2dff15bf788c671c1934e366d07e30c1814a8ef514e1af724a602e8a2fbe1b10" +checksum = "07bbe89c50d7a535e539b8c17bc0b49bdb77747034daa8087407d655f3f7cc1d" dependencies = [ "futures-core", "futures-sink", @@ -1423,15 +1464,15 @@ dependencies = [ [[package]] name = "futures-core" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "05f29059c0c2090612e8d742178b0580d2dc940c837851ad723096f87af6663e" +checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d" [[package]] name = "futures-executor" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e28d1d997f585e54aebc3f97d39e72338912123a67330d723fdbb564d646c9f" +checksum = "baf29c38818342a3b26b5b923639e7b1f4a61fc5e76102d4b1981c6dc7a7579d" dependencies = [ "futures-core", "futures-task", @@ -1451,9 +1492,9 @@ dependencies = [ [[package]] name = "futures-io" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e5c1b78ca4aae1ac06c48a526a655760685149f0d465d21f37abfe57ce075c6" +checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718" [[package]] name = "futures-lite" @@ -1485,9 +1526,9 @@ dependencies = [ [[package]] name = "futures-macro" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "162ee34ebcb7c64a8abebc059ce0fee27c2262618d7b60ed8faf72fef13c3650" +checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" dependencies = [ "proc-macro2", "quote", @@ -1496,21 +1537,21 @@ dependencies = [ [[package]] name = "futures-sink" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e575fab7d1e0dcb8d0c7bcf9a63ee213816ab51902e6d244a95819acacf1d4f7" +checksum = "c39754e157331b013978ec91992bde1ac089843443c49cbc7f46150b0fad0893" [[package]] name = "futures-task" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f90f7dce0722e95104fcb095585910c0977252f286e354b5e3bd38902cd99988" +checksum = "037711b3d59c33004d3856fbdc83b99d4ff37a24768fa1be9ce3538a1cde4393" [[package]] name = "futures-util" -version = "0.3.31" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9fa08315bb612088cc391249efdc3bc77536f16c91f6cf495e6fbe85b20a4a81" +checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ "futures-channel", "futures-core", @@ -1520,7 +1561,6 @@ dependencies = [ "futures-task", "memchr", "pin-project-lite", - "pin-utils", "slab", ] @@ -1903,7 +1943,7 @@ dependencies = [ "hyper", "hyper-util", "log", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-native-certs", "rustls-pki-types", "tokio", @@ -2163,9 +2203,9 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.85" +version = "0.3.91" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8c942ebf8e95485ca0d52d97da7c5a2c387d0e7f0ba4c35e93bfcaee045955b3" +checksum = "b49715b7073f385ba4bc528e5747d02e66cb39c6146efb66b781f131f0fb399c" dependencies = [ "once_cell", "wasm-bindgen", @@ -2233,7 +2273,7 @@ checksum = "3d0b95e02c851351f877147b7deea7b1afb1df71b63aa5f8270716e0c5720616" dependencies = [ "bitflags 2.11.0", "libc", - "redox_syscall 0.7.1", + "redox_syscall 0.7.3", ] [[package]] @@ -2254,9 +2294,9 @@ checksum = "ef53942eb7bf7ff43a617b3e2c1c4a5ecf5944a7c1bc12d7ee39bbb15e5c1519" [[package]] name = "linux-raw-sys" -version = "0.11.0" +version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df1d3c3b53da64cf5760482273a98e575c651a67eec7f77df96b5b642de8f039" +checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" [[package]] name = "litemap" @@ -2367,9 +2407,9 @@ dependencies = [ [[package]] name = "native-tls" -version = "0.2.16" +version = "0.2.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d5d26952a508f321b4d3d2e80e78fc2603eaefcdf0c30783867f19586518bdc" +checksum = "465500e14ea162429d264d44189adc38b199b62b1c21eea9f69e4b73cb03bbf2" dependencies = [ "libc", "log", @@ -2535,9 +2575,9 @@ dependencies = [ [[package]] name = "owo-colors" -version = "4.2.3" +version = "4.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9c6901729fa79e91a0913333229e9ca5dc725089d1c363b2f4b4760709dc4a52" +checksum = "d211803b9b6b570f68772237e415a029d5a50c65d382910b879fb19d3271f94d" [[package]] name = "parking" @@ -2595,9 +2635,9 @@ checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" [[package]] name = "pin-project-lite" -version = "0.2.16" +version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3b3cff922bd51709b605d9ead9aa71031d81447142d828eb4a6eba76fe619f9b" +checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" [[package]] name = "pin-utils" @@ -2760,7 +2800,7 @@ dependencies = [ "concurrent-queue", "hermit-abi 0.5.2", "pin-project-lite", - "rustix 1.1.3", + "rustix 1.1.4", "windows-sys 0.61.2", ] @@ -2846,7 +2886,7 @@ dependencies = [ "quinn-proto", "quinn-udp", "rustc-hash", - "rustls 0.23.36", + "rustls 0.23.37", "socket2 0.6.2", "thiserror 2.0.18", "tokio", @@ -2862,13 +2902,15 @@ checksum = "f1906b49b0c3bc04b5fe5d86a77925ae6524a19b816ae38ce1e426255f1d8a31" dependencies = [ "aws-lc-rs", "bytes", + "fastbloom", "getrandom 0.3.4", "lru-slab", "rand 0.9.2", "ring", "rustc-hash", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-pki-types", + "rustls-platform-verifier", "slab", "thiserror 2.0.18", "tinyvec", @@ -2981,9 +3023,9 @@ dependencies = [ [[package]] name = "redox_syscall" -version = "0.7.1" +version = "0.7.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "35985aa610addc02e24fc232012c86fd11f14111180f902b67e2d5331f8ebf2b" +checksum = "6ce70a74e890531977d37e532c34d45e9055d2409ed08ddba14529471ed0be16" dependencies = [ "bitflags 2.11.0", ] @@ -3019,9 +3061,9 @@ checksum = "cab834c73d247e67f4fae452806d17d3c7501756d98c8808d7c9c7aa7d18f973" [[package]] name = "regex-syntax" -version = "0.8.9" +version = "0.8.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a96887878f22d7bad8a3b6dc5b7440e0ada9a245242924394987b21cf2210a4c" +checksum = "dc897dd8d9e8bd1ed8cdad82b5966c3e0ecae09fb1907d58efaa013543185d0a" [[package]] name = "reqwest" @@ -3046,7 +3088,7 @@ dependencies = [ "percent-encoding", "pin-project-lite", "quinn", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-pki-types", "rustls-platform-verifier", "sync_wrapper", @@ -3132,14 +3174,14 @@ dependencies = [ [[package]] name = "rustix" -version = "1.1.3" +version = "1.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "146c9e247ccc180c1f61615433868c99f3de3ae256a30a43b49f67c2d9171f34" +checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" dependencies = [ "bitflags 2.11.0", "errno", "libc", - "linux-raw-sys 0.11.0", + "linux-raw-sys 0.12.1", "windows-sys 0.61.2", ] @@ -3157,13 +3199,14 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.36" +version = "0.23.37" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c665f33d38cea657d9614f766881e4d510e0eda4239891eea56b4cadcf01801b" +checksum = "758025cb5fccfd3bc2fd74708fd4682be41d99e5dff73c377c0646c6012c73a4" dependencies = [ "aws-lc-rs", "log", "once_cell", + "ring", "rustls-pki-types", "rustls-webpki 0.103.9", "subtle", @@ -3221,7 +3264,7 @@ dependencies = [ "jni", "log", "once_cell", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-native-certs", "rustls-platform-verifier-android", "rustls-webpki 0.103.9", @@ -3307,9 +3350,9 @@ dependencies = [ [[package]] name = "security-framework" -version = "3.6.0" +version = "3.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d17b898a6d6948c3a8ee4372c17cb384f90d2e6e912ef00895b14fd7ab54ec38" +checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d" dependencies = [ "bitflags 2.11.0", "core-foundation 0.10.1", @@ -3320,9 +3363,9 @@ dependencies = [ [[package]] name = "security-framework-sys" -version = "2.16.0" +version = "2.17.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "321c8673b092a9a42605034a9879d73cb79101ed5fd117bc9a597b89b4e9e61a" +checksum = "6ce2691df843ecc5d231c0b14ece2acc3efb62c0a398c7e1d875f3983ce020e3" dependencies = [ "core-foundation-sys", "libc", @@ -3472,6 +3515,12 @@ version = "0.3.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e320a6c5ad31d271ad523dcf3ad13e2767ad8b1cb8f047f75a8aeaf8da139da2" +[[package]] +name = "siphasher" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b2aa850e253778c88a04c3d7323b043aeda9d3e30d5971937c1855769763678e" + [[package]] name = "slab" version = "0.4.12" @@ -3747,6 +3796,12 @@ version = "0.27.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" +[[package]] +name = "strum" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9628de9b8791db39ceda2b119bbe13134770b56c138ec1d3af810d045c04f9bd" + [[package]] name = "strum_macros" version = "0.27.2" @@ -3759,6 +3814,18 @@ dependencies = [ "syn", ] +[[package]] +name = "strum_macros" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ab85eea0270ee17587ed4156089e10b9e6880ee688791d45a905f5b1ca36f664" +dependencies = [ + "heck", + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "subtle" version = "2.6.1" @@ -3767,9 +3834,9 @@ checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" [[package]] name = "syn" -version = "2.0.115" +version = "2.0.117" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6e614ed320ac28113fa64972c4262d5dbc89deacdfd00c34a3e4cea073243c12" +checksum = "e665b8803e7b1d2a727f4023456bbbbe74da67099c585258af0ad9c5013b9b99" dependencies = [ "proc-macro2", "quote", @@ -3852,14 +3919,14 @@ checksum = "df7f62577c25e07834649fc3b39fafdc597c0a3527dc1c60129201ccfcbaa50c" [[package]] name = "tempfile" -version = "3.25.0" +version = "3.26.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0136791f7c95b1f6dd99f9cc786b91bb81c3800b639b3478e561ddb7be95e5f1" +checksum = "82a72c767771b47409d2345987fda8628641887d5466101319899796367354a0" dependencies = [ "fastrand 2.3.0", "getrandom 0.4.1", "once_cell", - "rustix 1.1.3", + "rustix 1.1.4", "windows-sys 0.61.2", ] @@ -4023,7 +4090,7 @@ version = "0.26.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ - "rustls 0.23.36", + "rustls 0.23.37", "tokio", ] @@ -4094,9 +4161,9 @@ dependencies = [ [[package]] name = "toml_parser" -version = "1.0.8+spec-1.1.0" +version = "1.0.9+spec-1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0742ff5ff03ea7e67c8ae6c93cac239e0d9784833362da3f9a9c1da8dfefcbdc" +checksum = "702d4415e08923e7e1ef96cd5727c0dfed80b4d2fa25db9647fe5eb6f7c5a4c4" dependencies = [ "winnow", ] @@ -4226,7 +4293,7 @@ dependencies = [ "log", "native-tls", "rand 0.9.2", - "rustls 0.23.36", + "rustls 0.23.37", "rustls-pki-types", "sha1", "thiserror 2.0.18", @@ -4248,9 +4315,9 @@ checksum = "5c1cb5db39152898a79168971543b1cb5020dff7fe43c8dc468b0885f5e29df5" [[package]] name = "unicode-ident" -version = "1.0.23" +version = "1.0.24" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "537dd038a89878be9b64dd4bd1b260315c1bb94f4d784956b81e27a088d9a09e" +checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" [[package]] name = "unicode-normalization" @@ -4417,9 +4484,9 @@ checksum = "b8dad83b4f25e74f184f64c43b150b91efe7647395b42289f38e50566d82855b" [[package]] name = "wasm-bindgen" -version = "0.2.108" +version = "0.2.114" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "64024a30ec1e37399cf85a7ffefebdb72205ca1c972291c51512360d90bd8566" +checksum = "6532f9a5c1ece3798cb1c2cfdba640b9b3ba884f5db45973a6f442510a87d38e" dependencies = [ "cfg-if", "once_cell", @@ -4430,9 +4497,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-futures" -version = "0.4.58" +version = "0.4.64" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "70a6e77fd0ae8029c9ea0063f87c46fde723e7d887703d74ad2616d792e51e6f" +checksum = "e9c5522b3a28661442748e09d40924dfb9ca614b21c00d3fd135720e48b67db8" dependencies = [ "cfg-if", "futures-util", @@ -4444,9 +4511,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro" -version = "0.2.108" +version = "0.2.114" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "008b239d9c740232e71bd39e8ef6429d27097518b6b30bdf9086833bd5b6d608" +checksum = "18a2d50fcf105fb33bb15f00e7a77b772945a2ee45dcf454961fd843e74c18e6" dependencies = [ "quote", "wasm-bindgen-macro-support", @@ -4454,9 +4521,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.108" +version = "0.2.114" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5256bae2d58f54820e6490f9839c49780dff84c65aeab9e772f15d5f0e913a55" +checksum = "03ce4caeaac547cdf713d280eda22a730824dd11e6b8c3ca9e42247b25c631e3" dependencies = [ "bumpalo", "proc-macro2", @@ -4467,9 +4534,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-shared" -version = "0.2.108" +version = "0.2.114" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f01b580c9ac74c8d8f0c0e4afb04eeef2acf145458e52c03845ee9cd23e3d12" +checksum = "75a326b8c223ee17883a4251907455a2431acc2791c98c26279376490c378c16" dependencies = [ "unicode-ident", ] @@ -4510,9 +4577,9 @@ dependencies = [ [[package]] name = "web-sys" -version = "0.3.85" +version = "0.3.91" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "312e32e551d92129218ea9a2452120f4aabc03529ef03e4d0d82fb2780608598" +checksum = "854ba17bb104abfb26ba36da9729addc7ce7f06f5c0f90f3c391f8461cca21f9" dependencies = [ "js-sys", "wasm-bindgen", @@ -5179,18 +5246,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.39" +version = "0.8.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "db6d35d663eadb6c932438e763b262fe1a70987f9ae936e60158176d710cae4a" +checksum = "a789c6e490b576db9f7e6b6d661bcc9799f7c0ac8352f56ea20193b2681532e5" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.39" +version = "0.8.40" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4122cd3169e94605190e77839c9a40d40ed048d305bfdc146e7df40ab0f3e517" +checksum = "f65c489a7071a749c849713807783f70672b28094011623e200cb86dcb835953" dependencies = [ "proc-macro2", "quote", @@ -5300,9 +5367,9 @@ dependencies = [ [[package]] name = "zlib-rs" -version = "0.6.0" +version = "0.6.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a7948af682ccbc3342b6e9420e8c51c1fe5d7bf7756002b4a3c6cabfe96a7e3c" +checksum = "c745c48e1007337ed136dc99df34128b9faa6ed542d80a1c673cf55a6d7236c8" [[package]] name = "zmij" diff --git a/Cargo.toml b/Cargo.toml index 27766b7..8432d26 100755 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,6 +4,9 @@ version = "0.1.0" edition = "2024" [dependencies] +epsilon-core = { git = "https://github.com/Tensamin/Epsilon.git", package = "epsilon-core" } +epsilon-native = { git = "https://github.com/Tensamin/Epsilon.git", package = "epsilon-native" } + actix = "0.13.5" actix-rt = "2.11.0" actix-web = { version = "4.12.1", features = ["rustls-0_23"] } @@ -44,7 +47,7 @@ async-tungstenite = { version = "0.32.0", features = [ "verbose-logging", "webpki-roots", ] } -axum = { version = "0.8.8", features = [ "ws", "http2" ] } +axum = { version = "0.8.8", features = ["ws", "http2"] } base64 = "0.22.1" bytes = "1.11.1" color-eyre = "0.6.5" @@ -64,7 +67,7 @@ pnet = "0.35.0" rand = "0.8" rand_core = { version = "0.6", features = ["getrandom", "std"] } reqwest = "0.13.2" -rustls = "0.23.36" +rustls = "0.23.37" rustls-pemfile = "2.2.0" sha1 = "0.10.6" sha2 = "0.10.9" @@ -82,3 +85,4 @@ uuid = { version = "1.19.0", features = ["v4"] } walkdir = "2.5.0" x448 = "0.6.0" zip = "6.0.0" +quinn = "0.11.9" diff --git a/src/data/communication.rs b/src/data/communication.rs deleted file mode 100644 index 99229f5..0000000 --- a/src/data/communication.rs +++ /dev/null @@ -1,403 +0,0 @@ -use json::number::Number; -use json::{Array, JsonValue, object, parse}; -use std::collections::HashMap; -use std::time::{SystemTime, UNIX_EPOCH}; -use strum::IntoEnumIterator; -use strum_macros::EnumIter; -use uuid::Uuid; - -#[derive(Eq, Hash, PartialEq, EnumIter, Clone, Debug)] -#[allow(non_camel_case_types, dead_code)] -pub enum DataTypes { - error_type, - accepted_ids, - uuid, - register_id, - - link, - - settings, - settings_name, - chat_partner_id, - chat_partner_name, - iota_id, - user_id, - user_ids, - iota_ids, - user_state, - user_states, - user_pings, - call_state, - screen_share, - private_key_hash, - accepted, - accepted_profiles, - denied_profiles, - content, - messages, - notifications, - send_time, - get_time, - get_variant, - shared_secret_own, - shared_secret_other, - shared_secret_sign, - shared_secret, - call_id, - call_token, - untill, - enabled, - start_date, - end_date, - receiver_id, - sender_id, - signature, - signed, - message, - message_state, - last_ping, - ping_iota, - ping_clients, - matches, - omikron, - offset, - amount, - position, - name, - path, - codec, - function, - payload, - result, - interactables, - want_to_watch, - watcher, - created_at, - username, - display, - avatar, - about, - status, - public_key, - sub_level, - sub_end, - community_address, - challenge, - community_title, - communities, - rho_connections, - user, - online_status, - omikron_id, - omikron_connections, - reset_token, - new_token, -} - -impl DataTypes { - pub fn parse(p0: String) -> DataTypes { - for datatype in DataTypes::iter() { - if datatype.to_string().to_lowercase().replace('_', "") - == p0.to_lowercase().replace('_', "") - { - return datatype; - } - } - DataTypes::error_type - } - pub fn to_string(&self) -> String { - return format!("{:?}", self); - } -} - -#[derive(PartialEq, Clone, EnumIter, Debug)] -#[allow(non_camel_case_types, dead_code)] -pub enum CommunicationType { - error, - error_anonymous, - error_internal, - error_invalid_data, - error_invalid_user_id, - error_invalid_omikron_id, - error_not_found, - error_not_authenticated, - error_no_iota, - error_invalid_challenge, - error_invalid_secret, - error_invalid_private_key, - error_invalid_public_key, - error_no_user_id, - error_no_call_id, - error_invalid_call_id, - success, - - shorten_link, - - settings_save, - settings_load, - settings_list, - message, - message_state, - message_send, - message_live, - message_other_iota, - message_chunk, - messages_get, - - push_notification, - read_notification, - get_notifications, - - change_confirm, - confirm_receive, - confirm_read, - get_chats, - get_states, - add_community, - remove_community, - get_communities, - challenge, - challenge_response, - register, - register_response, - identification, - identification_response, - register_iota, - register_iota_success, - ping, - pong, - add_conversation, - send_chat, - client_changed, - client_connected, - client_disconnected, - client_closed, - public_key, - private_key, - webrtc_sdp, - webrtc_ice, - start_stream, - end_stream, - watch_stream, - call_token, - call_invite, - call_disconnect_user, - call_timeout_user, - call_set_anonymous_joining, - end_call, - function, - update, - create_user, - rho_update, - - user_connected, - user_disconnected, - iota_connected, - iota_disconnected, - sync_client_iota_status, - - get_user_data, - get_iota_data, - iota_user_data, - - change_user_data, - change_iota_data, - - get_register, - complete_register_user, - complete_register_iota, - delete_user, - delete_iota, - - start_register, - complete_register, -} -impl CommunicationType { - pub fn parse(p0: String) -> CommunicationType { - for datatype in CommunicationType::iter() { - if datatype.to_string().to_lowercase().replace('_', "") - == p0.to_lowercase().replace('_', "") - { - return datatype; - } - } - CommunicationType::error - } - pub fn to_string(&self) -> String { - return format!("{:?}", self); - } -} - -#[derive(Debug, Clone)] -pub struct CommunicationValue { - id: Uuid, - comm_type: CommunicationType, - sender: i64, - receiver: i64, - data: HashMap, -} - -#[allow(dead_code)] -impl CommunicationValue { - pub fn new(comm_type: CommunicationType) -> Self { - Self { - id: Uuid::new_v4(), - comm_type, - sender: 0, - receiver: 0, - data: HashMap::new(), - } - } - pub fn with_id(mut self, p0: Uuid) -> Self { - self.id = p0; - self - } - pub fn get_id(&self) -> Uuid { - self.id.clone() - } - pub fn with_sender(mut self, sender: i64) -> Self { - self.sender = sender; - self - } - pub fn get_sender(&self) -> i64 { - self.sender.clone() - } - pub fn with_receiver(mut self, receiver: i64) -> Self { - self.receiver = receiver; - self - } - pub fn get_receiver(&self) -> i64 { - self.receiver.clone() - } - pub fn add_data_num(mut self, key: DataTypes, value: Number) -> Self { - self.data.insert(key, JsonValue::Number(value)); - self - } - pub fn add_data_str(mut self, key: DataTypes, value: String) -> Self { - self.data.insert(key, JsonValue::String(value)); - self - } - pub fn add_data(mut self, key: DataTypes, value: JsonValue) -> Self { - self.data.insert(key, value); - self - } - pub fn add_array(mut self, key: DataTypes, value: Array) -> Self { - self.data.insert(key, JsonValue::Array(value)); - self - } - pub fn get_data(&self, key: DataTypes) -> Option<&JsonValue> { - self.data.get(&key) - } - - pub fn get_type(&self) -> CommunicationType { - self.comm_type.clone() - } - pub fn is_type(&self, p0: CommunicationType) -> bool { - self.comm_type == p0 - } - pub fn to_json(&self) -> JsonValue { - let mut jdata = object! {}; - for (k, v) in &self.data { - jdata[&format!("{:?}", k)] = JsonValue::from(v.clone()); - } - if self.sender > 0 && self.receiver > 0 { - object! { - id: self.id.to_string(), - type: format!("{:?}", self.comm_type), - sender: self.sender, - receiver: self.receiver, - data: jdata - } - } else if self.sender > 0 { - object! { - id: self.id.to_string(), - type: format!("{:?}", self.comm_type), - sender: self.sender, - data: jdata - } - } else if self.receiver > 0 { - object! { - id: self.id.to_string(), - type: format!("{:?}", self.comm_type), - receiver: self.receiver, - data: jdata - } - } else { - object! { - id: self.id.to_string(), - type: format!("{:?}", self.comm_type), - data: jdata - } - } - } - - pub fn from_json(json_str: &str) -> Self { - if let Ok(parsed) = parse(json_str) { - let comm_type = CommunicationType::parse(parsed["type"].to_string()); - let mut sender: i64 = 0; - if parsed.has_key("sender") { - sender = parsed["sender"].as_i64().unwrap_or(0); - } - let mut receiver: i64 = 0; - if parsed.has_key("receiver") { - receiver = parsed["receiver"].as_i64().unwrap_or(0); - } - - let uuid = - Uuid::parse_str(parsed["id"].as_str().unwrap_or("")).unwrap_or(Uuid::new_v4()); - let mut data = HashMap::new(); - if parsed["data"].is_object() { - for (k, v) in parsed["data"].entries() { - data.insert(DataTypes::parse(k.to_string()), v.clone()); - } - } - - Self { - id: uuid, - comm_type, - sender, - receiver, - data, - } - } else { - Self { - id: Uuid::new_v4(), - comm_type: CommunicationType::error, - sender: 0, - receiver: 0, - data: HashMap::new(), - } - } - } - pub fn forward_to_other_iota(original: &mut CommunicationValue) -> CommunicationValue { - let receiver = original - .get_data(DataTypes::receiver_id) - .unwrap_or(&JsonValue::Number(Number::from(0))) - .as_i64() - .unwrap_or(0); - - let now_ms = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_millis() as i64; - - let sender = original.get_sender(); - CommunicationValue::new(CommunicationType::message_other_iota) - .with_id(original.get_id()) - .with_receiver(receiver) - .add_data( - DataTypes::receiver_id, - JsonValue::Number(Number::from(receiver)), - ) - .with_sender(sender) - .add_data(DataTypes::send_time, JsonValue::String(now_ms.to_string())) - .add_data( - DataTypes::sender_id, - JsonValue::Number(Number::from(sender)), - ) - .add_data( - DataTypes::content, - JsonValue::String(original.get_data(DataTypes::content).unwrap().to_string()), - ) - } -} diff --git a/src/data/mod.rs b/src/data/mod.rs deleted file mode 100644 index 73c2657..0000000 --- a/src/data/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod communication; diff --git a/src/main.rs b/src/main.rs index 88637e6..043700c 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,8 +1,8 @@ -mod data; mod server; mod sql; mod util; +use crate::server::omikron_connection; use crate::sql::sql::initialize_db; use crate::sql::sql::print_users; use crate::util::crypto_helper::load_public_key; @@ -43,7 +43,7 @@ async fn main() { } else { log!(" Users"); } - let _ = server::server::start(9187).await; - + let _ = server::server::start(9188).await; + let _ = omikron_connection::OmikronServer::start(9187).await; tokio::signal::ctrl_c().await.unwrap(); } diff --git a/src/server/api.rs b/src/server/api.rs index f826e3e..7f3c728 100644 --- a/src/server/api.rs +++ b/src/server/api.rs @@ -1,4 +1,3 @@ -use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; use crate::get_public_key; use crate::server::omikron_manager::get_random_omikron; use crate::sql::sql; @@ -11,7 +10,6 @@ use actix_web::HttpResponse; use actix_web::http::{StatusCode, header}; use base64::Engine as _; use json::JsonValue; -use json::number::Number; pub async fn handle(path: &str, body_string: Option) -> HttpResponse { if path == "OPTIONS" { @@ -35,6 +33,9 @@ pub async fn handle(path: &str, body_string: Option) -> HttpResponse { }; let (status, body_text) = match path_parts.as_slice() { + // ================================================== + // DOWNLOAD IOTA FRONTEND + // ================================================== ["api", "download", "iota_frontend"] => { let file_path = "downloads/[iota_frontend].zip"; @@ -50,87 +51,107 @@ pub async fn handle(path: &str, body_string: Option) -> HttpResponse { .body(file_bytes); } Err(_) => { + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); return HttpResponse::NotFound() .insert_header(("Access-Control-Allow-Origin", "*")) - .body("File not found"); + .body(res.dump()); } } } + + // ================================================== + // GET RANDOM OMIKRON + // ================================================== ["api", "get", "omikron"] => { if let Ok(omikron_conn) = get_random_omikron().await { - if let Ok((public_key, ip_address)) = - sql::get_omikron_by_id(omikron_conn.get_omikron_id().await).await - { - ( - StatusCode::OK, - format!( - "{{\"id\": {}, \"public_key\": \"{}\", \"ip_address\": \"{}\"}}", - omikron_conn.get_omikron_id().await, - public_key, - ip_address - ), - ) + let id = omikron_conn.get_omikron_id().await; + + if let Ok((public_key, ip_address)) = sql::get_omikron_by_id(id).await { + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["id"] = id.into(); + res["public_key"] = public_key.into(); + res["ip_address"] = ip_address.into(); + + (StatusCode::OK, res.dump()) } else { - ( - StatusCode::INTERNAL_SERVER_ERROR, - "selected an invalid omikron".to_string(), - ) + let mut res = JsonValue::new_object(); + res["status"] = "error".into(); + (StatusCode::INTERNAL_SERVER_ERROR, res.dump()) } } else { - ( - StatusCode::NOT_FOUND, - "couldn't find online omikron".to_string(), - ) + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::NOT_FOUND, res.dump()) } } + + // ================================================== + // GET OMIKRON BY ID + // ================================================== ["api", "get", "omikron", id] => { let id = id.parse::().unwrap_or(0); + if id == 0 { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_bad_request".into(); + (StatusCode::BAD_REQUEST, res.dump()) } else if let Ok((public_key, ip_address)) = get_omikron_by_id(id).await { - ( - StatusCode::OK, - format!( - "{{\"id\": {}, \"public_key\": \"{}\", \"ip_address\": \"{}\"}}", - id, public_key, ip_address - ), - ) + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["id"] = id.into(); + res["public_key"] = public_key.into(); + res["ip_address"] = ip_address.into(); + (StatusCode::OK, res.dump()) } else if let Some(omikron_id) = get_iota_primary_omikron_connection(id) { if let Ok((public_key, ip_address)) = get_omikron_by_id(omikron_id).await { - ( - StatusCode::OK, - format!( - "{{\"id\": {}, \"public_key\": \"{}\", \"ip_address\": \"{}\"}}", - omikron_id, public_key, ip_address - ), - ) + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["id"] = omikron_id.into(); + res["public_key"] = public_key.into(); + res["ip_address"] = ip_address.into(); + (StatusCode::OK, res.dump()) } else { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::NOT_FOUND, res.dump()) } } else if let Ok((_, iota_id, _, _, _, _, _, _, _, _, _, _)) = get_by_user_id(id).await { if let Some(omikron_id) = get_iota_primary_omikron_connection(iota_id) { if let Ok((public_key, ip_address)) = get_omikron_by_id(omikron_id).await { - ( - StatusCode::OK, - format!( - "{{\"id\": {}, \"public_key\": \"{}\", \"ip_address\": \"{}\"}}", - omikron_id, public_key, ip_address - ), - ) + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["id"] = omikron_id.into(); + res["public_key"] = public_key.into(); + res["ip_address"] = ip_address.into(); + (StatusCode::OK, res.dump()) } else { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::NOT_FOUND, res.dump()) } } else { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::NOT_FOUND, res.dump()) } } else { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::NOT_FOUND, res.dump()) } } + + // ================================================== + // GET ID BY USERNAME + // ================================================== ["api", "get", "id", username] => { if username.is_empty() { - not_found() + let mut res = JsonValue::new_object(); + res["status"] = "error_bad_request".into(); + (StatusCode::BAD_REQUEST, res.dump()) } else if let Ok(( id, iota_id, @@ -146,37 +167,49 @@ pub async fn handle(path: &str, body_string: Option) -> HttpResponse { _, )) = sql::get_by_username(username).await { - let cv = CommunicationValue::new(CommunicationType::success) - .add_data_str(DataTypes::username, username) - .add_data_str(DataTypes::public_key, public_key) - .add_data(DataTypes::user_id, JsonValue::Number(Number::from(id))) - .add_data(DataTypes::iota_id, JsonValue::Number(Number::from(iota_id))) - .add_data( - DataTypes::sub_level, - JsonValue::Number(Number::from(sub_level)), - ) - .add_data(DataTypes::sub_end, JsonValue::Number(Number::from(sub_end))); - (StatusCode::OK, cv.to_json().to_string()) + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["username"] = username.into(); + res["public_key"] = public_key.into(); + res["user_id"] = id.into(); + res["iota_id"] = iota_id.into(); + res["sub_level"] = sub_level.into(); + res["sub_end"] = sub_end.into(); + + (StatusCode::OK, res.dump()) } else { - ( - StatusCode::OK, - CommunicationValue::new(CommunicationType::error_not_found) - .to_json() - .to_string(), - ) + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::OK, res.dump()) } } - ["api", "get", "public_key"] => (StatusCode::OK, public_key_to_base64(&get_public_key())), + + // ================================================== + // GET SERVER PUBLIC KEY + // ================================================== + ["api", "get", "public_key"] => { + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["public_key"] = public_key_to_base64(&get_public_key()).into(); + (StatusCode::OK, res.dump()) + } + + // ================================================== + // GET USER BY ID + // ================================================== ["api", "get", "user", id] => { let id: i64 = id.parse().unwrap_or(0); + if id == 0 { - bad_request() + let mut res = JsonValue::new_object(); + res["status"] = "error_bad_request".into(); + (StatusCode::BAD_REQUEST, res.dump()) } else if let Ok(( id, iota_id, username, display, - status, + status_msg, about, avatar, sub_level, @@ -186,47 +219,46 @@ pub async fn handle(path: &str, body_string: Option) -> HttpResponse { _, )) = sql::get_by_user_id(id).await { - let mut cv = CommunicationValue::new(CommunicationType::success) - .add_data_str(DataTypes::username, username) - .add_data_str(DataTypes::public_key, public_key) - .add_data(DataTypes::user_id, JsonValue::Number(Number::from(id))) - .add_data(DataTypes::iota_id, JsonValue::Number(Number::from(iota_id))) - .add_data( - DataTypes::sub_level, - JsonValue::Number(Number::from(sub_level)), - ) - .add_data(DataTypes::sub_end, JsonValue::Number(Number::from(sub_end))); + let mut res = JsonValue::new_object(); + res["status"] = "success".into(); + res["username"] = username.into(); + res["public_key"] = public_key.into(); + res["user_id"] = id.into(); + res["iota_id"] = iota_id.into(); + res["sub_level"] = sub_level.into(); + res["sub_end"] = sub_end.into(); + if let Some(display) = display { - cv = cv.add_data_str(DataTypes::display, display); + res["display"] = display.into(); } - if let Some(status) = status { - cv = cv.add_data_str(DataTypes::status, status); + if let Some(status_msg) = status_msg { + res["status_message"] = status_msg.into(); } if let Some(about) = about { - cv = cv.add_data_str(DataTypes::about, about); + res["about"] = about.into(); } if let Some(avatar) = avatar { - cv = cv.add_data_str( - DataTypes::avatar, - base64::engine::general_purpose::STANDARD.encode(avatar), - ); + res["avatar"] = base64::engine::general_purpose::STANDARD + .encode(avatar) + .into(); } - (StatusCode::OK, cv.to_json().to_string()) + + (StatusCode::OK, res.dump()) } else { - ( - StatusCode::OK, - CommunicationValue::new(CommunicationType::error_not_found) - .to_json() - .to_string(), - ) + let mut res = JsonValue::new_object(); + res["status"] = "error_not_found".into(); + (StatusCode::OK, res.dump()) } } - _ => ( - StatusCode::INTERNAL_SERVER_ERROR, - CommunicationValue::new(CommunicationType::error) - .to_json() - .to_string(), - ), + + // ================================================== + // DEFAULT + // ================================================== + _ => { + let mut res = JsonValue::new_object(); + res["status"] = "error".into(); + (StatusCode::INTERNAL_SERVER_ERROR, res.dump()) + } }; let body_bytes = body_text.into_bytes(); diff --git a/src/server/mod.rs b/src/server/mod.rs index 72ff9f4..aab6ab6 100644 --- a/src/server/mod.rs +++ b/src/server/mod.rs @@ -3,4 +3,3 @@ pub mod omikron_connection; pub mod omikron_manager; pub mod server; pub mod short_link; -pub mod socket; diff --git a/src/server/omikron_connection.rs b/src/server/omikron_connection.rs index 08a2ead..0ae49b7 100644 --- a/src/server/omikron_connection.rs +++ b/src/server/omikron_connection.rs @@ -1,43 +1,98 @@ -use crate::data::communication::{CommunicationType, CommunicationValue, DataTypes}; -use crate::server::omikron_manager; -use crate::server::short_link::add_short_link; -use crate::server::socket::{WsSendMessage, WsSession}; -use crate::sql::connection_status::UserStatus; -use crate::sql::sql::{self, get_by_user_id, get_by_username, get_iota_by_id, get_omikron_by_id}; -use crate::sql::user_online_tracker::{self}; -use crate::util::crypto_helper::encrypt; -use crate::util::logger::PrintType; -use crate::{get_private_key, get_public_key, log_cv_in, log_cv_out, log_in, log_out}; -use actix::Addr; +use crate::server::{omikron_manager, short_link::add_short_link}; +use crate::sql::{sql, sql::get_omikron_by_id, user_online_tracker}; +use crate::util::file_util::load_file_buf; +use crate::util::{crypto_helper::encrypt, logger::PrintType}; +use crate::{get_private_key, get_public_key, log_cv_in, log_cv_out, log_in}; use base64::{Engine as _, engine::general_purpose::STANDARD}; use dashmap::DashMap; -use json::JsonValue; -use json::number::Number; -use rand::Rng; -use rand::distributions::Alphanumeric; +use epsilon_core::{CommunicationType, CommunicationValue, DataTypes, DataValue}; +use epsilon_native::{Receiver, Sender, host}; +use quinn::ServerConfig; +use rand::{Rng, distributions::Alphanumeric}; +use rustls::pki_types::PrivateKeyDer; use std::sync::Arc; use tokio::sync::RwLock; -use uuid::Uuid; use x448::PublicKey; +pub struct OmikronServer; + +impl OmikronServer { + pub async fn start(port: u16) -> Result<(), Box> { + let tls_cfg = OmikronServer::load_tls().expect("TLS config failed"); + let server_crypto = quinn::crypto::rustls::QuicServerConfig::try_from(tls_cfg) + .expect("Failed to convert to QuicServerConfig"); + let mut host = host(port, ServerConfig::with_crypto(Arc::new(server_crypto))).await?; + tokio::spawn(async move { + while let Some((sender, receiver)) = host.next().await { + let connection = OmikronConnection::new(sender); + tokio::spawn(Self::connection_loop(connection, receiver)); + } + }); + + Ok(()) + } + fn load_tls() -> Option { + let mut cert_file_buf = load_file_buf("certs", "cert.pem").ok()?; + let mut key_file_buf = load_file_buf("certs", "cert.key").ok()?; + + let cert_chain = rustls_pemfile::certs(&mut cert_file_buf) + .collect::, _>>() + .ok()?; + + let mut keys: Vec = rustls_pemfile::pkcs8_private_keys(&mut key_file_buf) + .map(|k| k.map(Into::into)) + .collect::, _>>() + .ok()?; + + if keys.is_empty() { + let mut key_file_buf = load_file_buf("certs", "cert.key").ok()?; + keys = rustls_pemfile::rsa_private_keys(&mut key_file_buf) + .map(|k| k.map(Into::into)) + .collect::, _>>() + .ok()?; + } + + if keys.is_empty() { + return None; + } + + let cfg = rustls::ServerConfig::builder() + .with_no_client_auth() + .with_single_cert(cert_chain, keys.remove(0)) + .ok()?; + + Some(cfg) + } + + async fn connection_loop(conn: Arc, receiver: Receiver) { + while let Ok(cv) = receiver.receive().await { + log_cv_in!(PrintType::Omikron, cv); + conn.clone().handle_message(cv).await; + } + + conn.handle_close().await; + } +} + pub struct OmikronConnection { - pub ws_addr: Arc>>, - pub omikron_id: Arc>, - pub pub_key: Arc>>>, + sender: Arc>>, + omikron_id: Arc>, + pub_key: Arc>>>, + identified: Arc>, challenged: Arc>, challenge: Arc>, - pub ping: Arc>, - waiting_tasks: DashMap< - Uuid, - Box, CommunicationValue) -> bool + Send + Sync>, - >, + + ping: Arc>, + + waiting_tasks: + DashMap, CommunicationValue) -> bool + Send + Sync>>, } impl OmikronConnection { - pub fn new(ws_addr: Addr) -> Arc { + pub fn new(sender: Sender) -> Arc { Arc::new(Self { - ws_addr: Arc::new(RwLock::new(ws_addr)), + sender: Arc::new(RwLock::new(Some(sender))), omikron_id: Arc::new(RwLock::new(0)), pub_key: Arc::new(RwLock::new(None)), identified: Arc::new(RwLock::new(false)), @@ -47,836 +102,220 @@ impl OmikronConnection { waiting_tasks: DashMap::new(), }) } - pub async fn send_message(&self, cv: &CommunicationValue) { - let text = cv.to_json().to_string(); + async fn send(&self, cv: &CommunicationValue) { log_cv_out!(PrintType::Omikron, cv); - self.ws_addr.read().await.do_send(WsSendMessage(text)); + if let Some(sender) = self.sender.read().await.as_ref() { + if sender.send(cv).await.is_err() { + self.handle_close().await; + } + } + } + + async fn send_error_response(&self, id: u32, comm_type: CommunicationType) { + let response = CommunicationValue::new(comm_type).with_id(id); + self.send(&response).await; } pub async fn get_omikron_id(&self) -> i64 { - let omikron_id = { - let guard = self.omikron_id.read().await; - guard.clone() - }; - omikron_id + *self.omikron_id.read().await } - pub async fn is_identified(&self) -> bool { + + async fn is_identified(&self) -> bool { *self.identified.read().await && *self.challenged.read().await } - pub async fn handle_message(self: Arc, message: String) { - let cv = CommunicationValue::from_json(&message); - log_cv_in!(PrintType::Omikron, cv); + pub async fn handle_message(self: Arc, cv: CommunicationValue) { if cv.is_type(CommunicationType::ping) { self.handle_ping(cv).await; return; } + if let Some((_, task)) = self.waiting_tasks.remove(&cv.get_id()) { let _ = task(self.clone(), cv.clone()); return; } - // ────────────────────────────── - // Identification - // ────────────────────────────── if !self.is_identified().await { - let identified = *self.identified.read().await; - let challenged = *self.challenged.read().await; + self.handle_identification(cv).await; + return; + } - if !identified && cv.is_type(CommunicationType::identification) { - let omikron_id = cv - .get_data(DataTypes::omikron) - .and_then(|v| v.as_i64()) - .unwrap_or(0); + self.handle_authenticated(cv).await; + } - let (public_key, _) = match get_omikron_by_id(omikron_id).await { - Ok(v) => v, - Err(e) => { - self.send_message( + async fn handle_identification(self: &Arc, cv: CommunicationValue) { + let identified = *self.identified.read().await; + let challenged = *self.challenged.read().await; + + if !identified && cv.is_type(CommunicationType::identification) { + let omikron_id = cv.get_data(DataTypes::omikron).as_number().unwrap_or(0); + + let (public_key, _) = match get_omikron_by_id(omikron_id).await { + Ok(v) => v, + Err(_) => { + let _ = self + .send( &CommunicationValue::new(CommunicationType::error_not_authenticated) - .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), + .with_id(cv.get_id()), ) .await; - return; - } - }; + return; + } + }; - let pub_key_bytes = match STANDARD.decode(&public_key) { - Ok(b) => b, - Err(_) => { - self.send_error_response( - &cv.get_id(), - CommunicationType::error_invalid_omikron_id, + let pub_key_bytes = match STANDARD.decode(&public_key) { + Ok(b) => b, + Err(_) => { + let _ = self + .send( + &CommunicationValue::new(CommunicationType::error_invalid_omikron_id) + .with_id(cv.get_id()), ) .await; - return; - } - }; + return; + } + }; - let omikron_pub_key = match PublicKey::from_bytes(&pub_key_bytes) { - Some(k) => k, - _ => { - self.send_error_response( - &cv.get_id(), - CommunicationType::error_invalid_public_key, + let omikron_pub_key = match PublicKey::from_bytes(&pub_key_bytes) { + Some(k) => k, + _ => { + let _ = self + .send( + &CommunicationValue::new(CommunicationType::error_invalid_public_key) + .with_id(cv.get_id()), ) .await; - return; - } - }; - - let challenge: String = rand::thread_rng() - .sample_iter(&Alphanumeric) - .take(32) - .map(char::from) - .collect(); - - *self.omikron_id.write().await = omikron_id; - *self.challenge.write().await = challenge.clone(); - *self.pub_key.write().await = Some(pub_key_bytes); - *self.identified.write().await = true; - - let encrypted = - encrypt(get_private_key(), omikron_pub_key, &challenge).unwrap_or_default(); - - let response = CommunicationValue::new(CommunicationType::challenge) - .with_id(cv.get_id()) - .add_data_str( - DataTypes::public_key, - STANDARD.encode(get_public_key().as_bytes()), - ) - .add_data_str(DataTypes::challenge, encrypted); - - self.send_message(&response).await; - return; - } - - if identified && !challenged && cv.is_type(CommunicationType::challenge_response) { - let client_response = cv - .get_data(DataTypes::challenge) - .and_then(|v| v.as_str()) - .unwrap_or(""); - - if client_response == *self.challenge.read().await { - *self.challenged.write().await = true; - omikron_manager::add_omikron(self.clone()).await; - self.send_message( - &CommunicationValue::new(CommunicationType::identification_response) - .with_id(cv.get_id()) - .add_data(DataTypes::accepted, JsonValue::Boolean(true)), - ) - .await; - log_in!(PrintType::Omega, "Omikron Connected"); - } else { - self.send_error_response( - &cv.get_id(), - CommunicationType::error_invalid_challenge, - ) - .await; - self.close().await; + return; } - return; - } + }; - self.send_error_response(&cv.get_id(), CommunicationType::error_not_authenticated) - .await; - self.close().await; - } - let omikron_id = self.get_omikron_id().await; + let challenge: String = rand::thread_rng() + .sample_iter(&Alphanumeric) + .take(32) + .map(char::from) + .collect(); - if cv.is_type(CommunicationType::shorten_link) { - if let Some(link) = cv.get_data(DataTypes::link) { - if let Some(link) = link.as_str() { - if let Ok(short_link) = add_short_link(link).await { - let response = CommunicationValue::new(CommunicationType::shorten_link) - .with_id(cv.get_id()) - .add_data(DataTypes::link, JsonValue::String(short_link)); - self.send_message(&response).await; - return; - } - } - } - } + *self.omikron_id.write().await = omikron_id; + *self.challenge.write().await = challenge.clone(); + *self.pub_key.write().await = Some(pub_key_bytes); + *self.identified.write().await = true; - // ONLINE STATUS TRACKING - if cv.is_type(CommunicationType::user_connected) { - log_in!(PrintType::Omega, "User connected"); - if let Some(user_id) = cv.get_data(DataTypes::user_id).and_then(|v| v.as_i64()) { - user_online_tracker::track_user_status( - user_id, - UserStatus::user_online, - omikron_id, - ); - } - return; - } - if cv.is_type(CommunicationType::user_disconnected) { - log_in!(PrintType::Omega, "User disconnected"); - if let Some(user_id) = cv.get_data(DataTypes::user_id).and_then(|v| v.as_i64()) { - if let Some(status) = user_online_tracker::get_user_status(user_id) { - user_online_tracker::track_user_status( - user_id, - UserStatus::user_offline, - status.omikron_id, - ); - } - } - return; - } + let encrypted = + encrypt(get_private_key(), omikron_pub_key, &challenge).unwrap_or_default(); - if cv.is_type(CommunicationType::iota_connected) { - log_in!(PrintType::Omega, "IOTA connected"); - if let Some(iota_id) = cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()) { - user_online_tracker::track_iota_connection(iota_id, omikron_id, true); - let mut user_ids = JsonValue::Array(Vec::new()); - if let Ok(users) = sql::get_users_by_iota_id(iota_id).await { - for user in users { - let _ = user_ids.push(JsonValue::from(user.0)); - user_online_tracker::track_user_status( - user.0, - UserStatus::user_offline, - omikron_id, - ); - } - } else { - log_in!(PrintType::General, "SQL error, when loading users for IOTA"); - } - self.send_message( - &CommunicationValue::new(CommunicationType::iota_user_data) - .with_id(cv.get_id()) - .add_data(DataTypes::user_ids, user_ids), - ) - .await; - } else { - log_in!(PrintType::General, "No IOTA ID found"); - } - return; - } - if cv.is_type(CommunicationType::iota_disconnected) { - log_in!(PrintType::Omega, "IOTA disconnected"); - - if let Some(v) = cv.get_data(DataTypes::iota_id) { - if let Some(iota_id) = v.as_i64() { - let iota_offline = - user_online_tracker::untrack_iota_connection(iota_id, omikron_id); - if iota_offline { - if let Ok(users) = sql::get_users_by_iota_id(iota_id).await { - let user_ids: Vec = users.iter().map(|u| u.0).collect(); - user_online_tracker::untrack_many_users(&user_ids); - } - } - } - } - - return; - } - - if cv.is_type(CommunicationType::sync_client_iota_status) { - if let Some(json::JsonValue::Array(user_ids)) = - cv.get_data(DataTypes::user_ids).cloned() - { - for user_id_json in user_ids { - if let Some(user_id) = user_id_json.as_i64() { - user_online_tracker::track_user_status( - user_id, - UserStatus::user_online, - omikron_id, - ); - } - } - } - if let Some(json::JsonValue::Array(iota_ids)) = - cv.get_data(DataTypes::iota_ids).cloned() - { - for iota_id_json in iota_ids { - if let Some(iota_id) = iota_id_json.as_i64() { - user_online_tracker::track_iota_connection(iota_id, omikron_id, true); - } - } - } - return; - } - - // DATA RETURN - if cv.is_type(CommunicationType::get_user_data) { - if let Some(user_id) = cv.get_data(DataTypes::user_id).cloned() { - if let Some(user_id) = user_id.as_i64() { - if let Ok(( - id, - iota_id, - username, - display, - status, - about, - avatar, - sub_level, - sub_end, - public_key, - _, - _, - )) = get_by_user_id(user_id).await - { - let mut response = - CommunicationValue::new(CommunicationType::get_user_data) - .with_id(cv.get_id()) - .add_data_str(DataTypes::username, username.clone()) - .add_data_str(DataTypes::public_key, public_key) - .add_data(DataTypes::user_id, JsonValue::Number(Number::from(id))) - .add_data( - DataTypes::iota_id, - JsonValue::Number(Number::from(iota_id)), - ) - .add_data( - DataTypes::sub_level, - JsonValue::Number(Number::from(sub_level)), - ) - .add_data( - DataTypes::sub_end, - JsonValue::Number(Number::from(sub_end)), - ); - - if let Some(display) = display { - if display.is_empty() { - response = response.add_data_str(DataTypes::display, username); - } else { - response = response.add_data_str(DataTypes::display, display); - } - } - if let Some(status) = status { - if !status.is_empty() { - response = response.add_data_str(DataTypes::status, status); - } - } - if let Some(about) = about { - if !about.is_empty() { - response = response.add_data_str(DataTypes::about, about); - } - } - if let Some(avatar) = avatar { - response = - response.add_data_str(DataTypes::avatar, STANDARD.encode(avatar)); - } - - let user_status = user_online_tracker::get_user_status(id); - let iota_connections = - user_online_tracker::get_iota_omikron_connections(iota_id) - .unwrap_or_default(); - response = response.add_data( - DataTypes::omikron_connections, - JsonValue::Array( - iota_connections - .into_iter() - .map(|id| JsonValue::Number(Number::from(id))) - .collect(), - ), - ); - - if let Some(user_status) = user_status { - response = response.add_data( - DataTypes::online_status, - JsonValue::String(user_status.connection_type.to_string()), - ); - response = response.add_data( - DataTypes::omikron_id, - JsonValue::Number(Number::from(user_status.omikron_id)), - ); - } else { - response = response.add_data( - DataTypes::online_status, - JsonValue::String(UserStatus::iota_offline.to_string()), - ); - } - - self.send_message(&response).await; - return; - } - } - } - if let Some(username) = cv.get_data(DataTypes::username).cloned() { - if let Some(username) = username.as_str() { - if let Ok(( - id, - iota_id, - username, - display, - status, - about, - avatar, - sub_level, - sub_end, - public_key, - _, - _, - )) = get_by_username(username).await - { - let mut response = - CommunicationValue::new(CommunicationType::get_user_data) - .with_id(cv.get_id()) - .add_data_str(DataTypes::username, username.clone()) - .add_data_str(DataTypes::public_key, public_key) - .add_data(DataTypes::user_id, JsonValue::Number(Number::from(id))) - .add_data( - DataTypes::iota_id, - JsonValue::Number(Number::from(iota_id)), - ) - .add_data( - DataTypes::sub_level, - JsonValue::Number(Number::from(sub_level)), - ) - .add_data( - DataTypes::sub_end, - JsonValue::Number(Number::from(sub_end)), - ); - - if let Some(display) = display { - if display.is_empty() { - response = response.add_data_str(DataTypes::display, username); - } else { - response = response.add_data_str(DataTypes::display, display); - } - } - if let Some(status) = status { - if !status.is_empty() { - response = response.add_data_str(DataTypes::status, status); - } - } - if let Some(about) = about { - if !about.is_empty() { - response = response.add_data_str(DataTypes::about, about); - } - } - if let Some(avatar) = avatar { - response = - response.add_data_str(DataTypes::avatar, STANDARD.encode(avatar)); - } - - let user_status = user_online_tracker::get_user_status(id); - let iota_connections = - user_online_tracker::get_iota_omikron_connections(iota_id) - .unwrap_or_default(); - - if let Some(user_status) = user_status { - response = response.add_data( - DataTypes::online_status, - JsonValue::String(user_status.connection_type.to_string()), - ); - response = response.add_data( - DataTypes::omikron_id, - JsonValue::Number(Number::from(user_status.omikron_id)), - ); - } else { - response = response.add_data( - DataTypes::online_status, - JsonValue::String(UserStatus::iota_offline.to_string()), - ); - } - response = response.add_data( - DataTypes::omikron_connections, - JsonValue::Array( - iota_connections - .iter() - .map(|&id| JsonValue::Number(Number::from(id))) - .collect(), - ), - ); - - self.send_message(&response).await; - return; - } - } - } - - let response = - CommunicationValue::new(CommunicationType::error_not_found).with_id(cv.get_id()); - self.send_message(&response).await; - - return; - } - if cv.is_type(CommunicationType::get_iota_data) { - if let Some(iota_id) = cv.get_data(DataTypes::iota_id) { - if let Some(iota_id) = iota_id.as_i64() { - if let Ok((iota_id, public_key)) = get_iota_by_id(iota_id).await { - let mut response = - CommunicationValue::new(CommunicationType::get_iota_data) - .with_id(cv.get_id()) - .add_data_str(DataTypes::public_key, public_key) - .add_data(DataTypes::iota_id, JsonValue::from(iota_id)); - - let iota_connections = - user_online_tracker::get_iota_omikron_connections(iota_id) - .unwrap_or_default(); - - response = response.add_data( - DataTypes::omikron_connections, - JsonValue::Array( - iota_connections - .iter() - .map(|&id| JsonValue::Number(Number::from(id))) - .collect(), - ), - ); - self.send_message(&response).await; - return; - } - } - } - if let Some(user_id) = cv.get_data(DataTypes::user_id).cloned() { - if let Some(user_id) = user_id.as_i64() { - if let Ok((_, iota_id, _, _, _, _, _, _, _, _, _, _)) = - get_by_user_id(user_id).await - { - if let Ok((iota_id, public_key)) = get_iota_by_id(iota_id).await { - let mut response = - CommunicationValue::new(CommunicationType::get_iota_data) - .with_id(cv.get_id()) - .add_data_str(DataTypes::public_key, public_key) - .add_data( - DataTypes::user_id, - JsonValue::Number(Number::from(user_id)), - ) - .add_data( - DataTypes::iota_id, - JsonValue::Number(Number::from(iota_id)), - ); - - let iota_connections = - user_online_tracker::get_iota_omikron_connections(iota_id) - .unwrap_or_default(); - - response = response.add_data( - DataTypes::omikron_connections, - JsonValue::Array( - iota_connections - .iter() - .map(|&id| JsonValue::Number(Number::from(id))) - .collect(), - ), - ); - - self.send_message(&response).await; - return; - } - } - } - } - if let Some(username) = cv.get_data(DataTypes::username).cloned() { - if let Some(username) = username.as_str() { - if let Ok((user_id, iota_id, _, _, _, _, _, _, _, _, _, _)) = - get_by_username(username).await - { - if let Ok((iota_id, public_key)) = get_iota_by_id(iota_id).await { - let mut response = - CommunicationValue::new(CommunicationType::get_iota_data) - .with_id(cv.get_id()) - .add_data_str(DataTypes::public_key, public_key) - .add_data( - DataTypes::user_id, - JsonValue::Number(Number::from(user_id)), - ) - .add_data_str(DataTypes::username, username.to_string()) - .add_data( - DataTypes::iota_id, - JsonValue::Number(Number::from(iota_id)), - ); - - let iota_connections = - user_online_tracker::get_iota_omikron_connections(iota_id) - .unwrap_or_default(); - - response = response.add_data( - DataTypes::omikron_connections, - JsonValue::Array( - iota_connections - .iter() - .map(|&id| JsonValue::Number(Number::from(id))) - .collect(), - ), - ); - self.send_message(&response).await; - return; - } - } - } - } - let response = - CommunicationValue::new(CommunicationType::error_not_found).with_id(cv.get_id()); - self.send_message(&response).await; - - return; - } - - // REGISTERING - if cv.is_type(CommunicationType::get_register) { - let register_id = sql::get_register_id().await; - let response = CommunicationValue::new(CommunicationType::get_register) + let response = CommunicationValue::new(CommunicationType::challenge) .with_id(cv.get_id()) .add_data( - DataTypes::user_id, - JsonValue::Number(Number::from(register_id)), - ); - self.send_message(&response).await; + DataTypes::public_key, + DataValue::Str(STANDARD.encode(get_public_key().as_bytes())), + ) + .add_data(DataTypes::challenge, DataValue::Str(encrypted)); + + let _ = self.send(&response).await; return; } - if cv.is_type(CommunicationType::complete_register_iota) { - let iota_id_opt = cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()); - if let Some(public_key) = cv.get_data(DataTypes::public_key).and_then(|v| v.as_str()) { - if let Some(iota_id) = iota_id_opt { - match sql::register_complete_iota(iota_id, public_key.to_string()).await { - Ok(_) => { - let response = CommunicationValue::new(CommunicationType::success) - .with_id(cv.get_id()); - self.send_message(&response).await; - } - Err(e) => { - self.send_message( - &CommunicationValue::new(CommunicationType::error) - .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), - ) - .await; - } - } - } else { - // New logic to create iota and return id - match sql::create_new_iota(public_key.to_string()).await { - Ok(new_iota_id) => { - let response = - CommunicationValue::new(CommunicationType::complete_register_iota) - .with_id(cv.get_id()) - .add_data(DataTypes::iota_id, JsonValue::from(new_iota_id)); - self.send_message(&response).await; - } - Err(e) => { - self.send_message( - &CommunicationValue::new(CommunicationType::error) - .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), - ) - .await; - } - } - } + if identified && !challenged && cv.is_type(CommunicationType::challenge_response) { + let client_response = cv.get_data(DataTypes::challenge).as_str().unwrap_or(""); + + if client_response == *self.challenge.read().await { + *self.challenged.write().await = true; + + omikron_manager::add_omikron(self.clone()).await; + + let _ = self + .send( + &CommunicationValue::new(CommunicationType::identification_response) + .with_id(cv.get_id()) + .add_data(DataTypes::accepted, DataValue::BoolTrue), + ) + .await; + + log_in!(PrintType::Omega, "Omikron Connected"); } else { - self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) + let _ = self + .send( + &CommunicationValue::new(CommunicationType::error_invalid_challenge) + .with_id(cv.get_id()), + ) .await; } - return; } - if cv.is_type(CommunicationType::complete_register_user) { - if let ( - Some(user_id), - Some(username), - Some(public_key), - Some(iota_id), - Some(reset_token), - ) = ( - cv.get_data(DataTypes::user_id).and_then(|v| v.as_i64()), - cv.get_data(DataTypes::username).and_then(|v| v.as_str()), - cv.get_data(DataTypes::public_key).and_then(|v| v.as_str()), - cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()), - cv.get_data(DataTypes::reset_token).and_then(|v| v.as_str()), - ) { - match sql::register_complete_user( - user_id, - username.to_string(), - public_key.to_string(), - iota_id, - reset_token.to_string(), - ) - .await - { - Ok(_) => { - let response = CommunicationValue::new(CommunicationType::success) - .with_id(cv.get_id()); - self.send_message(&response).await; - } - Err(e) => { - self.send_message( - &CommunicationValue::new(CommunicationType::error) + } + + async fn handle_authenticated(self: &Arc, cv: CommunicationValue) { + if cv.is_type(CommunicationType::shorten_link) { + if let Some(link) = cv.get_data(DataTypes::link).as_str() { + if let Ok(short) = add_short_link(link).await { + let _ = self + .send( + &CommunicationValue::new(CommunicationType::shorten_link) .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), + .add_data(DataTypes::link, DataValue::Str(short)), ) .await; - } } - } else { - self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) - .await; } return; } - // CHANGING DATA - if cv.is_type(CommunicationType::change_user_data) { - let user_id = cv.get_sender(); - let mut success = true; - let mut error_message = String::new(); - - if let Some(username) = cv.get_data(DataTypes::username).and_then(|v| v.as_str()) { - if let Err(e) = sql::change_username(user_id, username.to_string()).await { - success = false; - error_message = e.to_string(); - } - } - if let Some(display) = cv.get_data(DataTypes::display).and_then(|v| v.as_str()) { - if let Err(e) = sql::change_display_name(user_id, display.to_string()).await { - success = false; - error_message = e.to_string(); - } - } - if let Some(avatar) = cv.get_data(DataTypes::avatar).and_then(|v| v.as_str()) { - if let Err(e) = sql::change_avatar(user_id, avatar.to_string()).await { - success = false; - error_message = e.to_string(); - } - } - if let Some(about) = cv.get_data(DataTypes::about).and_then(|v| v.as_str()) { - if let Err(e) = sql::change_about(user_id, about.to_string()).await { - success = false; - error_message = e.to_string(); - } - } - if let Some(status) = cv.get_data(DataTypes::status).and_then(|v| v.as_str()) { - if let Err(e) = sql::change_status(user_id, status.to_string()).await { - success = false; - error_message = e.to_string(); - } - } - if let Some(public_key) = cv.get_data(DataTypes::public_key).and_then(|v| v.as_str()) { - if let Some(private_key_hash) = cv - .get_data(DataTypes::private_key_hash) - .and_then(|v| v.as_str()) - { - if let Err(e) = sql::change_keys( - user_id, - public_key.to_string(), - private_key_hash.to_string(), - ) - .await - { - success = false; - error_message = e.to_string(); - } - } else { - success = false; - error_message = - "private_key_hash is required when changing public_key".to_string(); - } - } - - if success { - let response = - CommunicationValue::new(CommunicationType::success).with_id(cv.get_id()); - self.send_message(&response).await; - } else { - self.send_message( - &CommunicationValue::new(CommunicationType::error) + if cv.is_type(CommunicationType::get_register) { + let register_id = sql::get_register_id().await; + let _ = self + .send( + &CommunicationValue::new(CommunicationType::get_register) .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, error_message), + .add_data(DataTypes::user_id, DataValue::Number(register_id as i64)), ) .await; - } return; } - if cv.is_type(CommunicationType::change_iota_data) { - if let (user_id, Some(iota_id), Some(reset_token), Some(new_token)) = ( - cv.get_sender(), - cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()), - cv.get_data(DataTypes::reset_token).and_then(|v| v.as_str()), - cv.get_data(DataTypes::new_token).and_then(|v| v.as_str()), - ) { - match sql::get_by_user_id(user_id).await { - Ok(user) => { - let current_token = user.11; - if current_token == reset_token { - let mut success = true; - let mut error_message = String::new(); - if let Err(e) = sql::change_iota_id(user_id, iota_id).await { - success = false; - error_message = e.to_string(); - } - if success { - if let Err(e) = - sql::change_token(user_id, new_token.to_string()).await - { - success = false; - error_message = e.to_string(); - } - } - if success { - let response = CommunicationValue::new(CommunicationType::success) - .with_id(cv.get_id()); - self.send_message(&response).await; - } else { - self.send_message( - &CommunicationValue::new(CommunicationType::error) - .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, error_message), - ) - .await; - } - } else { - self.send_error_response( - &cv.get_id(), - CommunicationType::error_invalid_challenge, - ) - .await; - } - } - Err(_) => { - self.send_error_response(&cv.get_id(), CommunicationType::error_not_found) - .await; - } - } - } else { - self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) - .await; - } - return; - } if cv.is_type(CommunicationType::delete_user) { let user_id = cv.get_sender(); - match sql::delete_user(user_id).await { + match sql::delete_user(user_id as i64).await { Ok(_) => { - let response = - CommunicationValue::new(CommunicationType::success).with_id(cv.get_id()); - self.send_message(&response).await; + let _ = self + .send( + &CommunicationValue::new(CommunicationType::success) + .with_id(cv.get_id()), + ) + .await; } Err(e) => { - self.send_message( - &CommunicationValue::new(CommunicationType::error) - .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), - ) - .await; + let _ = self + .send( + &CommunicationValue::new(CommunicationType::error) + .with_id(cv.get_id()) + .add_data(DataTypes::error_type, DataValue::Str(e.to_string())), + ) + .await; } } return; } + if cv.is_type(CommunicationType::delete_iota) { - if let Some(iota_id) = cv.get_data(DataTypes::iota_id).and_then(|v| v.as_i64()) { - match sql::delete_iota(iota_id).await { + if let DataValue::Number(iota_id) = cv.get_data(DataTypes::iota_id) { + match sql::delete_iota(*iota_id).await { Ok(_) => { let response = CommunicationValue::new(CommunicationType::success) .with_id(cv.get_id()); - self.send_message(&response).await; + self.send(&response).await; } Err(e) => { - self.send_message( + self.send( &CommunicationValue::new(CommunicationType::error) .with_id(cv.get_id()) - .add_data_str(DataTypes::error_type, e.to_string()), + .add_data(DataTypes::error_type, DataValue::Str(e.to_string())), ) .await; } } } else { - self.send_error_response(&cv.get_id(), CommunicationType::error_invalid_data) + self.send_error_response(cv.get_id(), CommunicationType::error_invalid_data) .await; } return; @@ -885,85 +324,68 @@ impl OmikronConnection { // NOTIFICATIONS if cv.is_type(CommunicationType::get_notifications) { let user_id = cv.get_sender(); - if let Ok(notifications) = sql::get_notifications(user_id).await { + if let Ok(notifications) = sql::get_notifications(user_id as i64).await { let mut json_array = Vec::new(); for (sender, amount) in notifications { - let mut obj = JsonValue::new_object(); - let _ = obj.insert("sender", JsonValue::from(sender)); - let _ = obj.insert("amount", JsonValue::from(amount)); - json_array.push(obj); + let mut obj = Vec::new(); + let _ = obj.push((DataTypes::sender_id, DataValue::Number(sender))); + let _ = obj.push((DataTypes::amount, DataValue::Number(amount))); + json_array.push(DataValue::Container(obj)); } let response = CommunicationValue::new(CommunicationType::get_notifications) .with_id(cv.get_id()) - .add_array(DataTypes::notifications, json_array); - self.send_message(&response).await; + .add_data(DataTypes::notifications, DataValue::Array(json_array)); + self.send(&response).await; } } if cv.is_type(CommunicationType::read_notification) { if let (user_id, Some(other_id)) = ( cv.get_sender(), - cv.get_data(DataTypes::sender_id) - .unwrap_or(&JsonValue::Null) - .as_i64(), + cv.get_data(DataTypes::sender_id).as_number(), ) { - if let Ok(_) = sql::read_notification(user_id, other_id).await { + if let Ok(_) = sql::read_notification(user_id as i64, other_id).await { let response = CommunicationValue::new(CommunicationType::read_notification) .with_id(cv.get_id()); - self.send_message(&response).await; + self.send(&response).await; } } } if cv.is_type(CommunicationType::push_notification) { if let (user_id, Some(other_id)) = ( cv.get_sender(), - cv.get_data(DataTypes::sender_id) - .unwrap_or(&JsonValue::Null) - .as_i64(), + cv.get_data(DataTypes::sender_id).as_number(), ) { - if let Ok(_) = sql::add_notification(user_id, other_id).await { + if let Ok(_) = sql::add_notification(user_id as i64, other_id).await { let response = CommunicationValue::new(CommunicationType::push_notification) .with_id(cv.get_id()); - self.send_message(&response).await; + self.send(&response).await; } } } } - async fn send_error_response(&self, message_id: &Uuid, error_type: CommunicationType) { - let error = CommunicationValue::new(error_type).with_id(*message_id); - self.send_message(&error).await; - } - pub async fn close(&self) { - if self.is_identified().await { - let omikron_id = self.get_omikron_id().await; - if omikron_id != 0 { - log_in!(PrintType::Omega, "Omikron Disconnected"); - omikron_manager::remove_omikron(omikron_id).await; - user_online_tracker::untrack_omikron(omikron_id).await; - } - } - } - pub async fn handle_close(self: Arc) { - if self.is_identified().await { - let omikron_id = self.get_omikron_id().await; - if omikron_id != 0 { - log_in!(PrintType::Omega, "Omikron Disconnected"); - omikron_manager::remove_omikron(omikron_id).await; - user_online_tracker::untrack_omikron(omikron_id).await; - } - } - } - async fn handle_ping(&self, cv: CommunicationValue) { - if let Some(last_ping) = cv.get_data(DataTypes::last_ping) { - if let Ok(ping_val) = last_ping.to_string().parse::() { - let mut ping_guard = self.ping.write().await; - *ping_guard = ping_val; + if let DataValue::Number(last_ping) = cv.get_data(DataTypes::last_ping) { + if let Ok(val) = last_ping.to_string().parse::() { + *self.ping.write().await = val; } } - let response = CommunicationValue::new(CommunicationType::pong).with_id(cv.get_id()); + let _ = self + .send(&CommunicationValue::new(CommunicationType::pong).with_id(cv.get_id())) + .await; + } - self.send_message(&response).await; + pub async fn handle_close(&self) { + if self.is_identified().await { + let omikron_id = self.get_omikron_id().await; + + if omikron_id != 0 { + log_in!(PrintType::Omega, "Omikron Disconnected"); + + omikron_manager::remove_omikron(omikron_id).await; + user_online_tracker::untrack_omikron(omikron_id).await; + } + } } } diff --git a/src/server/server.rs b/src/server/server.rs index 5167444..c1b0e10 100644 --- a/src/server/server.rs +++ b/src/server/server.rs @@ -1,11 +1,10 @@ use crate::{ log, - server::{api, short_link::get_short_link, socket}, + server::{api, short_link::get_short_link}, util::file_util::get_directory, }; use actix_web::{App, HttpRequest, HttpResponse, HttpServer, Responder, http::header, web}; -use actix_web_actors::ws; use rustls::ServerConfig; use rustls::pki_types::{CertificateDer, PrivateKeyDer}; @@ -42,7 +41,6 @@ pub async fn start(port: u16) -> anyhow::Result<()> { App::new() .route("/api/{path:.*}", web::to(api_handler)) .route("/direct/{path:.*}", web::to(direct_handler)) - .route("/ws/{path:.*}", web::get().to(ws_handler)) }) .bind_rustls_0_23(addr, config)? .run() @@ -64,16 +62,6 @@ async fn direct_handler(req: HttpRequest) -> impl Responder { .finish() } } -async fn ws_handler( - req: HttpRequest, - stream: web::Payload, - path: web::Path, -) -> Result { - let path = path.into_inner(); - println!("WS handler reached: {}", path); - - ws::start(socket::WsSession::new(path), &req, stream) -} async fn api_handler(req: HttpRequest, body: web::Bytes) -> HttpResponse { let path = req.uri().path().to_string(); diff --git a/src/server/socket.rs b/src/server/socket.rs deleted file mode 100755 index 7d04b26..0000000 --- a/src/server/socket.rs +++ /dev/null @@ -1,129 +0,0 @@ -use std::sync::Arc; -use std::time::{Duration, Instant}; - -use actix::{Actor, ActorContext, AsyncContext, StreamHandler}; -use actix_web_actors::ws; - -use crate::data::communication::{CommunicationType, CommunicationValue}; -use crate::log; -use crate::server::omikron_connection::OmikronConnection; - -const IDLE_TIMEOUT: Duration = Duration::from_secs(30); - -use actix::Message; - -#[derive(Message)] -#[rtype(result = "()")] -pub struct WsSendMessage(pub String); - -impl actix::Handler for WsSession { - type Result = (); - - fn handle(&mut self, msg: WsSendMessage, ctx: &mut Self::Context) { - ctx.text(msg.0); - } -} - -pub struct WsSession { - path: String, - last_heartbeat: Instant, - omikron: Option>, -} - -impl WsSession { - pub fn new(path: String) -> Self { - Self { - path, - last_heartbeat: Instant::now(), - omikron: None, - } - } - - fn start_heartbeat(&self, ctx: &mut ws::WebsocketContext) { - ctx.run_interval(Duration::from_secs(5), |act, ctx| { - if Instant::now().duration_since(act.last_heartbeat) > IDLE_TIMEOUT { - log!("[ws_handler] Heartbeat failed. Disconnecting."); - ctx.close(None); - ctx.stop(); - return; - } - - let ping = CommunicationValue::new(CommunicationType::ping) - .to_json() - .to_string(); - - ctx.text(ping); - }); - } -} - -impl Actor for WsSession { - type Context = ws::WebsocketContext; - - fn started(&mut self, ctx: &mut Self::Context) { - self.start_heartbeat(ctx); - - if self.path == "omikron" { - let addr = ctx.address(); - let connection = OmikronConnection::new(addr); - self.omikron = Some(connection); - } - } - - fn stopped(&mut self, _: &mut Self::Context) { - log!( - "[ws] WebSocket handling task for path: {} is finished.", - self.path - ); - - if let Some(conn) = &self.omikron { - let conn = conn.clone(); - actix_rt::spawn(async move { - conn.handle_close().await; - }); - } - } -} -impl StreamHandler> for WsSession { - fn handle(&mut self, msg: Result, ctx: &mut Self::Context) { - match msg { - Ok(ws::Message::Text(text)) => { - self.last_heartbeat = Instant::now(); - - if let Some(conn) = &self.omikron { - let conn_clone = conn.clone(); - actix_rt::spawn(async move { - conn_clone.handle_message(text.to_string()).await; - }); - } - } - - Ok(ws::Message::Ping(msg)) => { - self.last_heartbeat = Instant::now(); - ctx.pong(&msg); - } - - Ok(ws::Message::Pong(_)) => { - self.last_heartbeat = Instant::now(); - log!("[ws_handler] Received Pong. Connection alive."); - } - - Ok(ws::Message::Close(reason)) => { - log!("[ws_handler] Received Close. Disconnecting."); - ctx.close(reason); - ctx.stop(); - } - - Ok(ws::Message::Binary(_)) => { - log!("[ws_handler] Binary message ignored."); - } - - Err(e) => { - log!("[ERROR] WS Error: {}. Closing.", e); - ctx.stop(); - } - - _ => {} - } - } -} diff --git a/src/util/logger.rs b/src/util/logger.rs index 80bc7c6..e861e69 100644 --- a/src/util/logger.rs +++ b/src/util/logger.rs @@ -1,4 +1,5 @@ use std::{ + collections::HashMap, fs::{self, OpenOptions}, io::Write, path::Path, @@ -8,10 +9,9 @@ use std::{ }; use ansi_term::Color; +use epsilon_core::{CommunicationValue, DataTypes, DataValue}; use json::JsonValue; -use crate::data::communication::CommunicationValue; - static LOGGER: OnceLock> = OnceLock::new(); #[derive(Clone, Copy)] @@ -224,7 +224,7 @@ pub fn log_cv_internal( let formatted = format_cv(cv); log_internal( - Some(cv.get_sender()), + Some(cv.get_sender() as i64), print_type.unwrap_or(PrintType::General), prefix, false, @@ -249,24 +249,79 @@ pub fn format_cv(cv: &CommunicationValue) -> String { let comm_type = cv.get_type().to_string(); parts.push(format!("{}", comm_type)); - let mut data_parts = Vec::new(); - if let JsonValue::Object(data) = &cv.clone().to_json()["data"] { - for (key, value) in data.iter() { - let val_string = match value { - JsonValue::String(s) => s.clone(), - _ => value.dump(), - }; + let data: &HashMap = cv.get_data_container(); - data_parts.push(format!("{} {}", key, val_string)); - } - } + let formated_data = + format_data_container(data.iter().map(|(k, v)| (k.clone(), v.clone())).collect()); - if !data_parts.is_empty() { - parts.push(format!("{}", data_parts.join(", "))); - } + parts.push(format!("{}", formated_data)); parts.join(": ") } + +fn format_data_container(data: Vec<(DataTypes, DataValue)>) -> String { + let parts: Vec = data + .into_iter() + .map(|(key, value)| { + let key_str = key.to_string(); + + match value { + DataValue::Str(s) => format!("{}=\"{}\"", key_str, s), + + DataValue::Container(inner) => { + let inner_formatted = format_data_container(inner); + format!("{}={{ {} }}", key_str, inner_formatted) + } + + DataValue::Array(arr) => { + let arr_formatted = format_array(arr); + format!("{}=[{}]", key_str, arr_formatted) + } + + DataValue::Bool(b) => format!("{}={}", key_str, b), + + DataValue::BoolTrue => format!("{}=true", key_str), + DataValue::BoolFalse => format!("{}=false", key_str), + + DataValue::Number(num) => format!("{}={}", key_str, num), + + _ => "".to_string(), + } + }) + .collect(); + + parts.join(", ") +} + +fn format_array(arr: Vec) -> String { + let parts: Vec = arr + .into_iter() + .map(|value| match value { + DataValue::Str(s) => format!("\"{}\"", s), + + DataValue::Container(inner) => { + let inner_formatted = format_data_container(inner); + format!("{{ {} }}", inner_formatted) + } + + DataValue::Array(inner_arr) => { + let formatted = format_array(inner_arr); + format!("[{}]", formatted) + } + + DataValue::Bool(b) => b.to_string(), + + DataValue::BoolTrue => "true".to_string(), + DataValue::BoolFalse => "false".to_string(), + + DataValue::Number(num) => num.to_string(), + + _ => String::new(), + }) + .collect(); + + parts.join(", ") +} #[macro_export] macro_rules! log_cv { ($kind:expr, $cv:expr) => {