diff --git a/Cargo.lock b/Cargo.lock index e2550e3..2625b80 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -235,7 +235,7 @@ checksum = "965c2d33e53cb6b267e148a4cb0760bc01f4904c1cd4bb4002a085bb016d1490" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "synstructure", ] @@ -247,7 +247,7 @@ checksum = "7b18050c2cd6fe86c3a76584ef5e0baf286d038cda203eb6223df2cc413565f7" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -384,7 +384,7 @@ checksum = "3b43422f69d8ff38f95f1b2bb76517c91589a924d1559a0e935d7c8ce0274c11" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -458,7 +458,7 @@ checksum = "e539d3fca749fcee5236ab05e93a52867dd549cc157c8cb7f99595f3cedffdb5" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -520,9 +520,9 @@ checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" [[package]] name = "autocfg" -version = "1.4.0" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ace50bade8e6234aa140d9a2f552bbee1db4d353f69b8217bc503490fc1a9f26" +checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" [[package]] name = "axum" @@ -725,7 +725,7 @@ dependencies = [ "proc-macro-crate 3.3.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1189,7 +1189,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "13b588ba4ac1a99f7f2964d24b3d896ddc6bf847ee3855dbd4366f058cfcd331" dependencies = [ "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1199,7 +1199,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a2785755761f3ddc1492979ce1e48d2c00d09311c39e4466429188f3dd6501" dependencies = [ "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1234,7 +1234,7 @@ checksum = "f46882e17999c6cc590af592290432be3bce0428cb0d5f8b6715e4dc7b383eb3" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1258,7 +1258,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1269,7 +1269,7 @@ checksum = "fc34b93ccb385b40dc71c6fceac4b2ad23662c7eeb248cf10d529b7e055b6ead" dependencies = [ "darling_core", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1323,7 +1323,7 @@ dependencies = [ "proc-macro2", "quote", "rustc_version", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1343,7 +1343,7 @@ checksum = "bda628edc44c4bb645fbe0f758797143e4e07926f7ebf4e9bdfbd3d2ce621df3" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "unicode-xid", ] @@ -1416,7 +1416,7 @@ checksum = "97369cbbc041bc366949bc74d34658d6cda5621039731c6310521892a3a20ae0" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1439,7 +1439,7 @@ checksum = "788160fb30de9cdd857af31c6a2675904b16ece8fc2737b2c7127ba368c9d0f4" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1515,9 +1515,9 @@ dependencies = [ [[package]] name = "embed-resource" -version = "3.0.3" +version = "3.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8fe7d068ca6b3a5782ca5ec9afc244acd99dd441e4686a83b1c3973aba1d489" +checksum = "0963f530273dc3022ab2bdc3fcd6d488e850256f2284a82b7413cb9481ee85dd" dependencies = [ "cc", "memchr", @@ -1566,7 +1566,7 @@ checksum = "67c78a4d8fdf9953a5c9d458f9efe940fd97a0cab0941c075a813ac594733827" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1610,12 +1610,12 @@ dependencies = [ [[package]] name = "errno" -version = "0.3.12" +version = "0.3.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cea14ef9355e3beab063703aa9dab15afd25f0667c341310c1e5274bb1d0da18" +checksum = "778e2ac28f6c47af28e4907f13ffd1e1ddbd400980a9abd7c8df189bf578a5ad" dependencies = [ "libc", - "windows-sys 0.59.0", + "windows-sys 0.60.2", ] [[package]] @@ -1729,7 +1729,7 @@ checksum = "1a5c6c585bc94aaf2c7b51dd4c2ba22680844aba4c687be581871a6f518c5742" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -1832,7 +1832,7 @@ checksum = "162ee34ebcb7c64a8abebc059ce0fee27c2262618d7b60ed8faf72fef13c3650" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -2114,7 +2114,7 @@ dependencies = [ "proc-macro-error", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -2228,7 +2228,7 @@ dependencies = [ "proc-macro-error", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -2402,7 +2402,7 @@ dependencies = [ "tokio", "tokio-rustls", "tower-service", - "webpki-roots 1.0.0", + "webpki-roots 1.0.1", ] [[package]] @@ -2736,7 +2736,7 @@ checksum = "03343451ff899767262ec32146f6d559dd759fdadf42ff0e227c7c48f72594b4" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -2858,9 +2858,9 @@ dependencies = [ [[package]] name = "libc" -version = "0.2.173" +version = "0.2.174" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d8cfeafaffdbc32176b64fb251369d52ea9f0a8fbc6f8759edffef7b525d64bb" +checksum = "1171693293099992e19cddea4e8b849964e9846f4acee11b3948bcc337be8776" [[package]] name = "libloading" @@ -2943,14 +2943,9 @@ version = "0.1.0" dependencies = [ "anyhow", "chrono", - "const-str", - "futures", "log", - "matchbox_socket", - "rand 0.9.1", - "rand_chacha 0.9.0", - "reqwest", - "rmp-serde", + "manhunt-logic", + "manhunt-transport", "serde", "serde_json", "specta", @@ -2965,7 +2960,20 @@ dependencies = [ "tauri-plugin-store", "tauri-specta", "tokio", - "tokio-util", + "uuid", +] + +[[package]] +name = "manhunt-logic" +version = "0.1.0" +dependencies = [ + "anyhow", + "chrono", + "rand 0.9.1", + "rand_chacha 0.9.0", + "serde", + "specta", + "tokio", "uuid", ] @@ -2984,6 +2992,26 @@ dependencies = [ "tokio", ] +[[package]] +name = "manhunt-transport" +version = "0.1.0" +dependencies = [ + "anyhow", + "const-str", + "futures", + "log", + "manhunt-logic", + "matchbox_protocol", + "matchbox_socket", + "rand 0.9.1", + "reqwest", + "rmp-serde", + "serde", + "tokio", + "tokio-util", + "uuid", +] + [[package]] name = "markup5ever" version = "0.11.0" @@ -3292,23 +3320,24 @@ dependencies = [ [[package]] name = "num_enum" -version = "0.7.3" +version = "0.7.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4e613fc340b2220f734a8595782c551f1250e969d87d3be1ae0579e8d4065179" +checksum = "a973b4e44ce6cad84ce69d797acf9a044532e4184c4f267913d1b546a0727b7a" dependencies = [ "num_enum_derive", + "rustversion", ] [[package]] name = "num_enum_derive" -version = "0.7.3" +version = "0.7.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "af1844ef2428cc3e1cb900be36181049ef3d3193c63e43026cfe202983b27a56" +checksum = "77e878c846a8abae00dd069496dbe8751b16ac1c3d6bd2a7283a938e8228f90d" dependencies = [ "proc-macro-crate 3.3.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -3827,7 +3856,7 @@ dependencies = [ "phf_shared 0.11.3", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4142,9 +4171,9 @@ dependencies = [ [[package]] name = "quinn-udp" -version = "0.5.12" +version = "0.5.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ee4e529991f949c5e25755532370b8af5d114acae52326361d68d47af64aa842" +checksum = "fcebb1209ee276352ef14ff8732e24cc2b02bbac986cd74a4c81bcb2f9881970" dependencies = [ "cfg_aliases", "libc", @@ -4165,9 +4194,9 @@ dependencies = [ [[package]] name = "r-efi" -version = "5.2.0" +version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" +checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" [[package]] name = "radium" @@ -4342,7 +4371,7 @@ checksum = "1165225c21bff1f3bbce98f5a1f889949bc902d3575308cc7b0de30b4f6d27c7" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4424,7 +4453,7 @@ dependencies = [ "wasm-bindgen-futures", "wasm-streams", "web-sys", - "webpki-roots 1.0.0", + "webpki-roots 1.0.1", ] [[package]] @@ -4735,7 +4764,7 @@ dependencies = [ "proc-macro2", "quote", "serde_derive_internals", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4866,7 +4895,7 @@ checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4877,7 +4906,7 @@ checksum = "18d26a20a969b9e3fdf2fc2d9f21eda6c40e2de84c9408bb5d3b05d499aae711" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4910,7 +4939,7 @@ checksum = "175ee3e80ae9982737ca543e96133087cbd9a485eecc3bc4de9c1a37b47ea59c" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -4962,7 +4991,7 @@ dependencies = [ "darling", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5154,6 +5183,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ab7f01e9310a820edd31c80fde3cae445295adde21a3f9416517d7d65015b971" dependencies = [ "chrono", + "ctor", "paste", "specta-macros", "thiserror 1.0.69", @@ -5169,7 +5199,7 @@ dependencies = [ "Inflector", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5304,9 +5334,9 @@ dependencies = [ [[package]] name = "syn" -version = "2.0.103" +version = "2.0.104" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e4307e30089d6fd6aff212f2da3a1f9e32f3223b1f010fb09b7c95f90f3ca1e8" +checksum = "17b6f705963418cdb9927482fa304bc562ece2fdd4f616084c50b7023b435a40" dependencies = [ "proc-macro2", "quote", @@ -5330,7 +5360,7 @@ checksum = "728a70f3dbaf5bab7f0c4b1ac8d7ae5ea60a4b5549c8a5914361c99147a709d2" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5414,7 +5444,7 @@ checksum = "f4e16beb8b2ac17db28eab8bca40e62dbfbb34c0fcdc6d9826b11b7b5d047dfd" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5521,7 +5551,7 @@ dependencies = [ "serde", "serde_json", "sha2", - "syn 2.0.103", + "syn 2.0.104", "tauri-utils", "thiserror 2.0.12", "time", @@ -5539,7 +5569,7 @@ dependencies = [ "heck 0.5.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "tauri-codegen", "tauri-utils", ] @@ -5603,9 +5633,9 @@ dependencies = [ [[package]] name = "tauri-plugin-geolocation" -version = "2.2.4" +version = "2.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f70978d3dbca4d900f8708dad1e0cff32c52de578101998d0f30dc7e12104e85" +checksum = "3dc06e23cb12d17533aca6001252f7f0fdbd295868194748821d0f0e6fbfbb59" dependencies = [ "log", "serde", @@ -5617,9 +5647,9 @@ dependencies = [ [[package]] name = "tauri-plugin-log" -version = "2.4.0" +version = "2.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8d2b582d860eb214f28323f4ce4f2797ae3b78f197e27b11677f976f9f52aedb" +checksum = "1063890b541877c4367d0ee9b1758360d86398862fbaafdd99aac02a55e28bf8" dependencies = [ "android_logger", "byte-unit", @@ -5639,9 +5669,9 @@ dependencies = [ [[package]] name = "tauri-plugin-notification" -version = "2.2.2" +version = "2.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c474c7cc524385e682ccc1e149e13913a66fd8586ac4c2319cf01b78f070d309" +checksum = "d1c87f171cdb35c3aa8f17e8dfd84c1b9f68eb4086ec16a5a1b9f13b5541c574" dependencies = [ "log", "notify-rust", @@ -5658,9 +5688,9 @@ dependencies = [ [[package]] name = "tauri-plugin-opener" -version = "2.2.7" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "66644b71a31ec1a8a52c4a16575edd28cf763c87cf4a7da24c884122b5c77097" +checksum = "2c8983f50326d34437142a6d560b5c3426e91324297519b6eeb32ed0a1d1e0f2" dependencies = [ "dunce", "glob", @@ -5680,9 +5710,9 @@ dependencies = [ [[package]] name = "tauri-plugin-store" -version = "2.2.0" +version = "2.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1c0c08fae6995909f5e9a0da6038273b750221319f2c0f3b526d6de1cde21505" +checksum = "ada7e7aeea472dec9b8d09d25301e59fe3e8330dc11dbcf903d6388126cb3722" dependencies = [ "dunce", "serde", @@ -5768,7 +5798,7 @@ dependencies = [ "heck 0.5.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5888,7 +5918,7 @@ checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5899,7 +5929,7 @@ checksum = "7f7cf42b4507d8ea322120659672cf1b9dbb93f8f2d4ecfd6e51350ff5b17a1d" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -5987,7 +6017,7 @@ checksum = "6e06d43f1345a3bcd39f6a56dbb7dcab2ba47e68e8ac134855e7e2bdbaf8cab8" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -6161,13 +6191,13 @@ dependencies = [ [[package]] name = "tracing-attributes" -version = "0.1.29" +version = "0.1.30" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1b1ffbcf9c6f6b99d386e7444eb608ba646ae452a36b39737deb9663b610f662" +checksum = "81383ab64e72a7a8b8e13130c49e3dab29def6d0c7d76a03087b3cf71c5c6903" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -6512,7 +6542,7 @@ dependencies = [ "log", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "wasm-bindgen-shared", ] @@ -6547,7 +6577,7 @@ checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "wasm-bindgen-backend", "wasm-bindgen-shared", ] @@ -6659,9 +6689,9 @@ dependencies = [ [[package]] name = "webpki-roots" -version = "1.0.0" +version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2853738d1cc4f2da3a225c18ec6c3721abb31961096e9dbf5ab35fa88b19cfdb" +checksum = "8782dd5a41a24eed3a4f40b606249b3e236ca61adf1f25ea4d45c73de122b502" dependencies = [ "rustls-pki-types", ] @@ -6897,7 +6927,7 @@ checksum = "1d228f15bba3b9d56dde8bddbee66fa24545bd17b48d5128ccf4a8742b18e431" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7011,7 +7041,7 @@ checksum = "a47fddd13af08290e67f4acabf4b459f647552718f683a7b415d290ac744a836" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7022,7 +7052,7 @@ checksum = "bd9211b69f8dcdfa817bfd14bf1c97c9188afa36f4750130fcdf3f400eca9fa8" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7504,7 +7534,7 @@ checksum = "38da3c9736e16c5d3c8c597a9aaa5d1fa565d0532ae05e27c24aa62fb32c0ab6" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "synstructure", ] @@ -7551,7 +7581,7 @@ dependencies = [ "proc-macro-crate 3.3.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "zbus_names", "zvariant", "zvariant_utils", @@ -7571,22 +7601,22 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.25" +version = "0.8.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a1702d9583232ddb9174e01bb7c15a2ab8fb1bc6f227aa1233858c351a3ba0cb" +checksum = "1039dd0d3c310cf05de012d8a39ff557cb0d23087fd44cad61df08fc31907a2f" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.25" +version = "0.8.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "28a6e20d751156648aa063f3800b706ee209a32c0b4d9f24be3d980b01be55ef" +checksum = "9ecf5b4cc5364572d7f4c329661bcc82724222973f2cab6f050a4e5c22f75181" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7606,7 +7636,7 @@ checksum = "d71e5d6e06ab090c67b5e44993ec16b72dcbaabc526db883a360057678b48502" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "synstructure", ] @@ -7627,7 +7657,7 @@ checksum = "ce36e65b0d2999d2aafac989fb249189a141aee1f53c612c1f37d72631959f69" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7660,7 +7690,7 @@ checksum = "5b96237efa0c878c64bd89c436f661be4e46b2f3eff1ebb976f7ef2321d2f58f" dependencies = [ "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", ] [[package]] @@ -7687,7 +7717,7 @@ dependencies = [ "proc-macro-crate 3.3.0", "proc-macro2", "quote", - "syn 2.0.103", + "syn 2.0.104", "zvariant_utils", ] @@ -7701,6 +7731,6 @@ dependencies = [ "quote", "serde", "static_assertions", - "syn 2.0.103", + "syn 2.0.104", "winnow 0.7.11", ] diff --git a/Cargo.toml b/Cargo.toml index e6d80cc..33a1228 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["backend", "manhunt-signaling"] +members = ["backend", "manhunt-logic", "manhunt-signaling", "manhunt-transport"] resolver = "3" [profile.release] diff --git a/TODO.md b/TODO.md index d5bc337..6c7a466 100644 --- a/TODO.md +++ b/TODO.md @@ -19,7 +19,7 @@ - [x] Meta : Recipes for type binding generation - [x] Signaling: All of it - [x] Backend : Better transport error handling -- [ ] Backend : Abstract lobby? Separate crate? +- [x] Backend : Abstract lobby? Separate crate? - [x] Transport : Handle transport cancellation better - [x] Backend : Add checks for when the `powerup_locations` field is an empty array in settings - [ ] Backend : More tests diff --git a/backend/Cargo.toml b/backend/Cargo.toml index a040c75..8957b31 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -22,24 +22,18 @@ tauri = { version = "2", features = [] } tauri-plugin-opener = "2" serde = { version = "1", features = ["derive"] } serde_json = "1" -chrono = { version = "0.4", features = ["serde", "now"] } tokio = { version = "1.45", features = ["sync", "macros", "time", "fs"] } -rand = { version = "0.9", features = ["thread_rng"] } tauri-plugin-geolocation = "2" -rand_chacha = "0.9.0" -futures = "0.3.31" -matchbox_socket = "0.12.0" -uuid = { version = "1.17.0", features = ["serde", "v4"] } -rmp-serde = "1.3.0" tauri-plugin-store = "2.2.0" -specta = { version = "=2.0.0-rc.22", features = ["chrono", "uuid"] } +specta = { version = "=2.0.0-rc.22", features = ["chrono", "uuid", "export"] } tauri-specta = { version = "=2.0.0-rc.21", features = ["derive", "typescript"] } specta-typescript = "0.0.9" tauri-plugin-log = "2" tauri-plugin-notification = "2" log = "0.4.27" -tokio-util = "0.7.15" anyhow = "1.0.98" -reqwest = { version = "0.12.20", default-features = false, features = ["charset", "http2", "rustls-tls", "system-proxy"] } -const-str = "0.6.2" tauri-plugin-dialog = "2" +manhunt-logic = { version = "0.1.0", path = "../manhunt-logic" } +manhunt-transport = { version = "0.1.0", path = "../manhunt-transport" } +uuid = { version = "1.17.0", features = ["serde"] } +chrono = { version = "0.4.41", features = ["serde"] } diff --git a/backend/src/game/transport.rs b/backend/src/game/transport.rs deleted file mode 100644 index ea89f4e..0000000 --- a/backend/src/game/transport.rs +++ /dev/null @@ -1,10 +0,0 @@ -use super::events::GameEvent; - -pub trait Transport { - /// Receive an event - async fn receive_messages(&self) -> impl Iterator; - /// Send an event - async fn send_message(&self, msg: GameEvent); - /// Disconnect from the transport - fn disconnect(&self) {} -} diff --git a/backend/src/history.rs b/backend/src/history.rs index f8c20b0..d1e45cb 100644 --- a/backend/src/history.rs +++ b/backend/src/history.rs @@ -1,14 +1,15 @@ +use anyhow::Context; use serde::{Deserialize, Serialize}; -use std::{collections::HashMap, sync::Arc}; +use std::{collections::HashMap, result::Result as StdResult, sync::Arc}; use tauri::{AppHandle, Runtime}; use tauri_plugin_store::{Store, StoreExt}; use uuid::Uuid; -use crate::{ - game::{GameHistory, UtcDT}, - prelude::*, - profile::PlayerProfile, -}; +use manhunt_logic::{GameHistory, PlayerProfile}; + +use crate::UtcDT; + +type Result = StdResult; #[derive(Debug, Clone, Serialize, Deserialize, specta::Type)] pub struct AppGameHistory { diff --git a/backend/src/lib.rs b/backend/src/lib.rs index e006856..77175e8 100644 --- a/backend/src/lib.rs +++ b/backend/src/lib.rs @@ -1,37 +1,52 @@ -mod game; mod history; -mod lobby; mod location; -mod profile; -mod server; -mod transport; +mod profiles; -use std::{collections::HashMap, sync::Arc, time::Duration}; +use std::{collections::HashMap, marker::PhantomData, sync::Arc, time::Duration}; -use game::{Game as BaseGame, GameSettings}; -use history::AppGameHistory; -use lobby::{Lobby, LobbyState, StartGameInfo}; +use anyhow::Context; use location::TauriLocation; use log::{error, info, warn, LevelFilter}; -use profile::PlayerProfile; +use manhunt_logic::{ + Game as BaseGame, GameSettings, GameUiState, Lobby as BaseLobby, LobbyState, PlayerProfile, + StartGameInfo, StateUpdateSender, +}; +use manhunt_transport::{generate_join_code, room_exists, MatchboxTransport}; use serde::{Deserialize, Serialize}; use tauri::{AppHandle, Manager, State}; use tauri_plugin_dialog::{DialogExt, MessageDialogKind}; use tauri_specta::{collect_commands, collect_events, ErrorHandlingMode, Event}; use tokio::sync::RwLock; -use transport::MatchboxTransport; use uuid::Uuid; -mod prelude { - pub use anyhow::{anyhow, bail, Context, Error as AnyhowError}; - pub use std::result::Result as StdResult; +type UtcDT = chrono::DateTime; - pub type Result = StdResult; +/// The state of the game has changed +#[derive(Serialize, Deserialize, Clone, Default, Debug, specta::Type, tauri_specta::Event)] +struct GameStateUpdate; + +/// The state of the lobby has changed +#[derive(Serialize, Deserialize, Clone, Default, Debug, specta::Type, tauri_specta::Event)] +struct LobbyStateUpdate; + +struct TauriStateUpdateSender(AppHandle, PhantomData); + +impl TauriStateUpdateSender { + fn new(app: &AppHandle) -> Self { + Self(app.clone(), PhantomData) + } } -use prelude::*; +impl StateUpdateSender for TauriStateUpdateSender { + fn send_update(&self) { + if let Err(why) = E::default().emit(&self.0) { + error!("Error sending Game state update to UI: {why:?}"); + } + } +} -type Game = BaseGame; +type Game = BaseGame>; +type Lobby = BaseLobby>; enum AppState { Setup, @@ -52,83 +67,66 @@ enum AppScreen { type AppStateHandle = RwLock; -fn generate_join_code() -> String { - // 5 character sequence of A-Z - (0..5) - .map(|_| (b'A' + rand::random_range(0..26)) as char) - .collect::() -} - const GAME_TICK_RATE: Duration = Duration::from_secs(1); /// The app is changing screens, contains the screen it's switching to #[derive(Serialize, Deserialize, Clone, Debug, specta::Type, tauri_specta::Event)] struct ChangeScreen(AppScreen); -/// The state of the game has updated in some way, you're expected to call [get_game_state] when -/// receiving this -#[derive(Serialize, Deserialize, Clone, Debug, specta::Type, tauri_specta::Event)] -struct GameStateUpdate; - -struct TauriStateUpdateSender(AppHandle); - -impl StateUpdateSender for TauriStateUpdateSender { - fn send_update(&self) { - if let Err(why) = GameStateUpdate.emit(&self.0) { - error!("Error sending Game state update to UI: {why:?}"); - } - } +fn error_dialog(app: &AppHandle, msg: &str) { + app.dialog() + .message(msg) + .kind(MessageDialogKind::Error) + .show(|_| {}); } impl AppState { - pub async fn start_game(&mut self, app: AppHandle, my_id: Uuid, start: StartGameInfo) { + pub async fn start_game(&mut self, app: AppHandle, start: StartGameInfo) { if let AppState::Lobby(lobby) = self { let transport = lobby.clone_transport(); let profiles = lobby.clone_profiles().await; let location = TauriLocation::new(app.clone()); - let state_updates = TauriStateUpdateSender(app.clone()); + let state_updates = TauriStateUpdateSender::new(&app); let game = Arc::new(Game::new( - my_id, GAME_TICK_RATE, - start.initial_caught_state, - start.settings, + start, transport, location, state_updates, )); *self = AppState::Game(game.clone(), profiles.clone()); + Self::game_loop(app.clone(), game, profiles); Self::emit_screen_change(&app, AppScreen::Game); - tokio::spawn(async move { - let res = game.main_loop().await; - let app2 = app.clone(); - let state_handle = app.state::(); - let mut state = state_handle.write().await; - match res { - Ok(Some(history)) => { - let history = AppGameHistory::new(history, profiles); - if let Err(why) = history.save_history(&app2) { - error!("Failed to save game history: {why:?}"); - app2.dialog() - .message("Failed to save the history of this game") - .kind(MessageDialogKind::Error) - .show(|_| {}); - } - state.quit_to_menu(app2); - } - Ok(None) => { - info!("User quit game"); - } - Err(why) => { - error!("Game Error: {why:?}"); - app2.dialog() - .message(format!("Connection Error: {why}")) - .kind(MessageDialogKind::Error) - .show(|_| {}); - state.quit_to_menu(app2); + } + } + + fn game_loop(app: AppHandle, game: Arc, profiles: HashMap) { + tokio::spawn(async move { + let res = game.main_loop().await; + let state_handle = app.state::(); + let mut state = state_handle.write().await; + match res { + Ok(Some(history)) => { + let history = AppGameHistory::new(history, profiles); + if let Err(why) = history.save_history(&app) { + error!("Failed to save game history: {why:?}"); + error_dialog(&app, "Failed to save the history of this game"); } + state.quit_to_menu(app.clone()).await; } - }); - } + Ok(None) => { + info!("User quit game"); + } + Err(why) => { + error!("Game Error: {why:?}"); + app.dialog() + .message(format!("Connection Error: {why}")) + .kind(MessageDialogKind::Error) + .show(|_| {}); + state.quit_to_menu(app.clone()).await; + } + } + }); } pub fn get_menu(&self) -> Result<&PlayerProfile> { @@ -185,7 +183,7 @@ impl AppState { pub fn complete_setup(&mut self, app: &AppHandle, profile: PlayerProfile) -> Result { if let AppState::Setup = self { - profile.write_to_store(app); + write_profile_to_store(app, profile.clone()); *self = AppState::Menu(profile); Self::emit_screen_change(app, AppScreen::Menu); Ok(()) @@ -207,7 +205,30 @@ impl AppState { } } - pub fn start_lobby( + fn lobby_loop(app: AppHandle, lobby: Arc) { + tokio::spawn(async move { + let res = lobby.main_loop().await; + let app_game = app.clone(); + let state_handle = app.state::(); + let mut state = state_handle.write().await; + match res { + Ok(Some(start)) => { + info!("Starting game as"); + state.start_game(app_game, start).await; + } + Ok(None) => { + info!("User quit lobby"); + } + Err(why) => { + error!("Lobby Error: {why}"); + error_dialog(&app_game, &format!("Error joining the lobby: {why}")); + state.quit_to_menu(app_game).await; + } + } + }); + } + + pub async fn start_lobby( &mut self, join_code: Option, app: AppHandle, @@ -216,44 +237,26 @@ impl AppState { if let AppState::Menu(profile) = self { let host = join_code.is_none(); let room_code = join_code.unwrap_or_else(generate_join_code); - let lobby = Arc::new(Lobby::new( - &room_code, - host, - profile.clone(), - settings, - app.clone(), - )); - *self = AppState::Lobby(lobby.clone()); - let app2 = app.clone(); - tokio::spawn(async move { - let res = lobby.open().await; - let app_game = app2.clone(); - let state_handle = app2.state::(); - let mut state = state_handle.write().await; - match res { - Ok(Some((my_id, start))) => { - info!("Starting game as {my_id}"); - state.start_game(app_game, my_id, start).await; - } - Ok(None) => { - info!("User quit lobby"); - } - Err(why) => { - error!("Lobby Error: {why}"); - app_game - .dialog() - .message(format!("Error joining the lobby: {why}")) - .kind(MessageDialogKind::Error) - .show(|_| {}); - state.quit_to_menu(app_game); - } + let state_updates = TauriStateUpdateSender::::new(&app); + let lobby = + Lobby::new(&room_code, host, profile.clone(), settings, state_updates).await; + match lobby { + Ok(lobby) => { + *self = AppState::Lobby(lobby.clone()); + Self::lobby_loop(app.clone(), lobby); + Self::emit_screen_change(&app, AppScreen::Lobby); } - }); - Self::emit_screen_change(&app, AppScreen::Lobby); + Err(why) => { + error_dialog( + &app, + &format!("Couldn't connect you to the lobby\n\n{why:?}"), + ); + } + } } } - pub fn quit_to_menu(&mut self, app: AppHandle) { + pub async fn quit_to_menu(&mut self, app: AppHandle) { let profile = match self { AppState::Setup => None, AppState::Menu(_) => { @@ -261,14 +264,14 @@ impl AppState { return; } AppState::Lobby(lobby) => { - lobby.quit_lobby(); - Some(lobby.self_profile.clone()) + lobby.quit_lobby().await; + read_profile_from_store(&app) } AppState::Game(game, _) => { - game.quit_game(); - PlayerProfile::load_from_store(&app) + game.quit_game().await; + read_profile_from_store(&app) } - AppState::Replay(_) => PlayerProfile::load_from_store(&app), + AppState::Replay(_) => read_profile_from_store(&app), }; let screen = if let Some(profile) = profile { *self = AppState::Menu(profile); @@ -284,7 +287,10 @@ impl AppState { use std::result::Result as StdResult; -use crate::game::{GameUiState, StateUpdateSender, UtcDT}; +use crate::{ + history::AppGameHistory, + profiles::{read_profile_from_store, write_profile_to_store}, +}; type Result = StdResult; @@ -309,7 +315,7 @@ async fn get_current_screen(state: State<'_, AppStateHandle>) -> Result) -> Result { let mut state = state.write().await; - state.quit_to_menu(app); + state.quit_to_menu(app).await; Ok(()) } @@ -358,9 +364,7 @@ async fn replay_game(id: UtcDT, app: AppHandle, state: State<'_, AppStateHandle> /// (Screen: Menu) Check if a room code is valid to join, use this before starting a game /// for faster error checking. async fn check_room_code(code: &str) -> Result { - server::room_exists(code) - .await - .map_err(|err| err.to_string()) + room_exists(code).await.map_err(|err| err.to_string()) } #[tauri::command] @@ -371,7 +375,7 @@ async fn update_profile( app: AppHandle, state: State<'_, AppStateHandle>, ) -> Result { - new_profile.write_to_store(&app); + write_profile_to_store(&app, new_profile.clone()); let mut state = state.write().await; let profile = state.get_menu_mut()?; *profile = new_profile; @@ -389,7 +393,7 @@ async fn start_lobby( state: State<'_, AppStateHandle>, ) -> Result { let mut state = state.write().await; - state.start_lobby(join_code, app, settings); + state.start_lobby(join_code, app, settings).await; Ok(()) } @@ -523,7 +527,7 @@ pub fn mk_specta() -> tauri_specta::Builder { .events(collect_events![ ChangeScreen, GameStateUpdate, - lobby::LobbyStateUpdate + LobbyStateUpdate ]) } @@ -551,7 +555,7 @@ pub fn run() { let handle = app.handle().clone(); tauri::async_runtime::spawn(async move { - if let Some(profile) = PlayerProfile::load_from_store(&handle) { + if let Some(profile) = read_profile_from_store(&handle) { let state_handle = handle.state::(); let mut state = state_handle.write().await; *state = AppState::Menu(profile); diff --git a/backend/src/lobby.rs b/backend/src/lobby.rs deleted file mode 100644 index fc6c96a..0000000 --- a/backend/src/lobby.rs +++ /dev/null @@ -1,244 +0,0 @@ -use std::{collections::HashMap, sync::Arc}; - -use log::{error, warn}; -use serde::{Deserialize, Serialize}; -use tauri::AppHandle; -use tauri_specta::Event; -use tokio::sync::Mutex; -use uuid::Uuid; - -use crate::{ - game::GameSettings, - prelude::*, - profile::PlayerProfile, - server, - transport::{MatchboxTransport, TransportMessage}, -}; - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct StartGameInfo { - pub settings: GameSettings, - pub initial_caught_state: HashMap, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub enum LobbyMessage { - /// Message sent on a new peer, to sync profiles - PlayerSync(PlayerProfile), - /// Message sent on a new peer from the host, to sync game settings - HostPush(GameSettings), - /// Host signals starting the game - StartGame(StartGameInfo), - /// A player has switched teams - PlayerSwitch(bool), -} - -#[derive(Clone, Serialize, Deserialize, specta::Type)] -pub struct LobbyState { - profiles: HashMap, - join_code: String, - /// True represents seeker, false hider - teams: HashMap, - self_id: Option, - self_seeker: bool, - is_host: bool, - settings: GameSettings, -} - -pub struct Lobby { - is_host: bool, - join_code: String, - pub self_profile: PlayerProfile, - state: Mutex, - transport: Arc, - app: AppHandle, -} - -/// The lobby state has updated in some way, you're expected to call [get_lobby_state] after -/// receiving this -#[derive(Serialize, Deserialize, Clone, Debug, specta::Type, tauri_specta::Event)] -pub struct LobbyStateUpdate; - -impl Lobby { - pub fn new( - join_code: &str, - host: bool, - profile: PlayerProfile, - settings: GameSettings, - app: AppHandle, - ) -> Self { - Self { - app, - transport: Arc::new(MatchboxTransport::new(join_code, host)), - is_host: host, - self_profile: profile, - join_code: join_code.to_string(), - state: Mutex::new(LobbyState { - teams: HashMap::with_capacity(5), - join_code: join_code.to_string(), - profiles: HashMap::with_capacity(5), - self_seeker: false, - self_id: None, - is_host: host, - settings, - }), - } - } - - fn emit_state_update(&self) { - if let Err(why) = LobbyStateUpdate.emit(&self.app) { - error!("Error emitting Lobby state update: {why:?}"); - } - } - - pub fn clone_transport(&self) -> Arc { - self.transport.clone() - } - - pub async fn clone_state(&self) -> LobbyState { - self.state.lock().await.clone() - } - - pub async fn clone_profiles(&self) -> HashMap { - let state = self.state.lock().await; - state.profiles.clone() - } - - /// Set self as seeker or hider - pub async fn switch_teams(&self, seeker: bool) { - let mut state = self.state.lock().await; - state.self_seeker = seeker; - if let Some(id) = state.self_id { - if let Some(state_seeker) = state.teams.get_mut(&id) { - *state_seeker = seeker; - } - } - drop(state); - self.transport - .send_transport_message(None, LobbyMessage::PlayerSwitch(seeker).into()) - .await; - self.emit_state_update(); - } - - /// (Host) Update game settings - pub async fn update_settings(&self, new_settings: GameSettings) { - if self.is_host { - let mut state = self.state.lock().await; - state.settings = new_settings.clone(); - drop(state); - let msg = LobbyMessage::HostPush(new_settings); - self.send_transport_message(None, msg).await; - self.emit_state_update(); - } - } - - async fn send_transport_message(&self, id: Option, msg: LobbyMessage) { - self.transport.send_transport_message(id, msg.into()).await - } - - async fn signaling_mark_started(&self) -> Result { - server::mark_room_started(&self.join_code).await - } - - /// (Host) Start the game - pub async fn start_game(&self) { - if self.is_host { - let state = self.state.lock().await; - let start_game_info = StartGameInfo { - settings: state.settings.clone(), - initial_caught_state: state.teams.clone(), - }; - drop(state); - let msg = LobbyMessage::StartGame(start_game_info); - self.send_transport_message(None, msg).await; - if let Err(why) = self.signaling_mark_started().await { - warn!("Failed to tell signalling server that the match started: {why:?}"); - } - self.emit_state_update(); - } - } - - pub fn quit_lobby(&self) { - self.transport.cancel(); - } - - pub async fn open(&self) -> Result> { - let transport_inner = self.transport.clone(); - tokio::spawn(async move { transport_inner.transport_loop().await }); - - let res = 'lobby: loop { - self.emit_state_update(); - - let msgs = self.transport.recv_transport_messages().await; - - for (peer, msg) in msgs { - match msg { - TransportMessage::IdAssigned(id) => { - let mut state = self.state.lock().await; - state.self_id = Some(id); - let seeker = state.self_seeker; - state.teams.insert(id, seeker); - state.profiles.insert(id, self.self_profile.clone()); - } - TransportMessage::Disconnected => { - break 'lobby Ok(None); - } - TransportMessage::Error(why) => { - break 'lobby Err(anyhow!("Transport error: {why}")); - } - TransportMessage::Game(game_event) => { - eprintln!("Peer {peer:?} sent a GameEvent: {game_event:?}"); - } - TransportMessage::Lobby(lobby_message) => match *lobby_message { - LobbyMessage::PlayerSync(player_profile) => { - let mut state = self.state.lock().await; - state.profiles.insert(peer, player_profile); - } - LobbyMessage::HostPush(game_settings) => { - let mut state = self.state.lock().await; - state.settings = game_settings; - } - LobbyMessage::StartGame(start_game_info) => { - let id = self - .state - .lock() - .await - .self_id - .expect("Error getting self ID"); - break 'lobby Ok(Some((id, start_game_info))); - } - LobbyMessage::PlayerSwitch(seeker) => { - let mut state = self.state.lock().await; - state.teams.insert(peer, seeker); - } - }, - TransportMessage::PeerConnect => { - let msg = LobbyMessage::PlayerSync(self.self_profile.clone()); - let mut state = self.state.lock().await; - state.teams.insert(peer, false); - drop(state); - self.send_transport_message(Some(peer), msg).await; - if self.is_host { - let state = self.state.lock().await; - let msg = LobbyMessage::HostPush(state.settings.clone()); - drop(state); - self.send_transport_message(Some(peer), msg).await; - } - } - TransportMessage::PeerDisconnect => { - let mut state = self.state.lock().await; - state.profiles.remove(&peer); - state.teams.remove(&peer); - } - TransportMessage::Seq(_) => {} - } - } - }; - - if res.is_err() { - self.transport.cancel(); - } - - res - } -} diff --git a/backend/src/location.rs b/backend/src/location.rs index d8c73ee..201d04f 100644 --- a/backend/src/location.rs +++ b/backend/src/location.rs @@ -1,7 +1,7 @@ use tauri::AppHandle; use tauri_plugin_geolocation::{GeolocationExt, PositionOptions}; -use crate::game::{Location, LocationService}; +use manhunt_logic::{Location, LocationService}; pub struct TauriLocation(AppHandle); diff --git a/backend/src/profile.rs b/backend/src/profile.rs deleted file mode 100644 index 074dc4c..0000000 --- a/backend/src/profile.rs +++ /dev/null @@ -1,36 +0,0 @@ -use serde::{Deserialize, Serialize}; -use tauri::AppHandle; -use tauri_plugin_store::StoreExt; - -#[derive(Clone, Default, Debug, Serialize, Deserialize, specta::Type)] -pub struct PlayerProfile { - display_name: String, - pfp_base64: Option, -} - -const STORE_NAME: &str = "profile.json"; - -impl PlayerProfile { - // pub fn has_pfp(&self) -> bool { - // self.pfp_base64.is_some() - // } - - pub fn load_from_store(app: &AppHandle) -> Option { - let store = app.store(STORE_NAME).expect("Couldn't Create Store"); - - let profile = store - .get("profile") - .and_then(|v| serde_json::from_value::(v).ok()); - - store.close_resource(); - - profile - } - - pub fn write_to_store(&self, app: &AppHandle) { - let store = app.store(STORE_NAME).expect("Couldn't create store"); - - let value = serde_json::to_value(self.clone()).expect("Failed to serialize"); - store.set("profile", value); - } -} diff --git a/backend/src/profiles.rs b/backend/src/profiles.rs new file mode 100644 index 0000000..a3bc2a1 --- /dev/null +++ b/backend/src/profiles.rs @@ -0,0 +1,24 @@ +use manhunt_logic::PlayerProfile; +use tauri::AppHandle; +use tauri_plugin_store::StoreExt; + +const STORE_NAME: &str = "profile"; + +pub fn read_profile_from_store(app: &AppHandle) -> Option { + let store = app.store(STORE_NAME).expect("Couldn't Create Store"); + + let profile = store + .get("profile") + .and_then(|v| serde_json::from_value::(v).ok()); + + store.close_resource(); + + profile +} + +pub fn write_profile_to_store(app: &AppHandle, profile: PlayerProfile) { + let store = app.store(STORE_NAME).expect("Couldn't create store"); + + let value = serde_json::to_value(profile).expect("Failed to serialize"); + store.set("profile", value); +} diff --git a/backend/src/transport.rs b/backend/src/transport.rs deleted file mode 100644 index c84b3f5..0000000 --- a/backend/src/transport.rs +++ /dev/null @@ -1,378 +0,0 @@ -use std::{ - collections::{HashMap, HashSet}, - time::Duration, -}; - -use anyhow::Context; -use futures::FutureExt; -use log::error; -use matchbox_socket::{Error as SocketError, PeerId, PeerState, WebRtcSocket}; -use serde::{Deserialize, Serialize}; -use tokio::sync::{Mutex, RwLock}; -use tokio_util::sync::CancellationToken; -use uuid::Uuid; - -use crate::{ - game::{GameEvent, Transport}, - lobby::LobbyMessage, - prelude::*, - server, -}; - -#[derive(Serialize, Deserialize, Debug, Clone)] -pub struct TransportChunk { - id: u64, - current: usize, - total: usize, - data: Vec, -} - -#[derive(Debug, Serialize, Deserialize, Clone)] -pub enum TransportMessage { - /// The transport has received a peer id - IdAssigned(Uuid), - /// Message related to the actual game - /// Boxed for space reasons - Game(Box), - /// Message related to the pre-game lobby - Lobby(Box), - /// Internal message when peer connects - PeerConnect, - /// Internal message when peer disconnects - PeerDisconnect, - /// Event sent when the transport gets disconnected, used to help consumers know when to stop - /// consuming messages - Disconnected, - /// Event when the transport encounters a critical error and needs to disconnect - Error(String), - /// Internal message for packet chunking - Seq(TransportChunk), -} - -// Max packet size according to: https://github.com/johanhelsing/matchbox/issues/272 -const MAX_PACKET_SIZE: usize = 65535; - -// Align packets with a bit of extra space for [TransportMessage::Seq] header -const PACKET_ALIGNMENT: usize = MAX_PACKET_SIZE - 128; - -impl TransportMessage { - pub fn serialize(&self) -> Vec { - rmp_serde::to_vec(self).expect("Failed to encode") - } - - pub fn deserialize(data: &[u8]) -> Result { - rmp_serde::from_slice(data).context("While deserializing message") - } - - pub fn from_packets(packets: impl Iterator>) -> Result { - let full_data = packets.flatten().collect::>(); - Self::deserialize(&full_data).context("While decoding a multi-part message") - } - - pub fn to_packets(&self) -> Vec> { - let bytes = self.serialize(); - if bytes.len() > MAX_PACKET_SIZE { - let id = rand::random_range(0..u64::MAX); - let packets_needed = bytes.len().div_ceil(PACKET_ALIGNMENT); - let rem = bytes.len() % PACKET_ALIGNMENT; - let bytes = bytes.into_boxed_slice(); - (0..packets_needed) - .map(|idx| { - let start = PACKET_ALIGNMENT * idx; - let end = if idx == packets_needed - 1 { - start + rem - } else { - PACKET_ALIGNMENT * (idx + 1) - }; - let data = bytes[start..end].to_vec(); - let chunk = TransportChunk { - id, - current: idx, - total: packets_needed, - data, - }; - TransportMessage::Seq(chunk).serialize() - }) - .collect() - } else { - vec![bytes] - } - } -} - -impl From for TransportMessage { - fn from(v: GameEvent) -> Self { - Self::Game(Box::new(v)) - } -} - -impl From for TransportMessage { - fn from(v: LobbyMessage) -> Self { - Self::Lobby(Box::new(v)) - } -} - -type OutgoingMsgPair = (Option, TransportMessage); -type OutgoingQueueSender = tokio::sync::mpsc::Sender; -type OutgoingQueueReceiver = tokio::sync::mpsc::Receiver; - -type IncomingMsgPair = (Uuid, TransportMessage); -type IncomingQueueSender = tokio::sync::mpsc::Sender; -type IncomingQueueReceiver = tokio::sync::mpsc::Receiver; - -pub struct MatchboxTransport { - ws_url: String, - incoming: (IncomingQueueSender, Mutex), - outgoing: (OutgoingQueueSender, Mutex), - my_id: RwLock>, - cancel_token: CancellationToken, -} - -impl MatchboxTransport { - pub fn new(join_code: &str, is_host: bool) -> Self { - let (itx, irx) = tokio::sync::mpsc::channel(15); - let (otx, orx) = tokio::sync::mpsc::channel(15); - - Self { - ws_url: server::room_url(join_code, is_host), - incoming: (itx, Mutex::new(irx)), - outgoing: (otx, Mutex::new(orx)), - my_id: RwLock::new(None), - cancel_token: CancellationToken::new(), - } - } - - pub async fn send_transport_message(&self, peer: Option, msg: TransportMessage) { - self.outgoing - .0 - .send((peer, msg)) - .await - .expect("Failed to add to outgoing queue"); - } - - pub async fn recv_transport_messages(&self) -> Vec { - let mut incoming_rx = self.incoming.1.lock().await; - let mut buffer = Vec::with_capacity(60); - incoming_rx.recv_many(&mut buffer, 60).await; - buffer - } - - async fn push_incoming(&self, id: Uuid, msg: TransportMessage) { - self.incoming - .0 - .send((id, msg)) - .await - .expect("Failed to push to incoming queue"); - } - - async fn handle_send( - &self, - socket: &mut WebRtcSocket, - all_peers: &HashSet, - messages: &mut Vec, - ) { - if let Some(my_id) = *self.my_id.read().await { - for (_, msg) in messages.iter().filter(|(id, _)| id.is_none()) { - self.push_incoming(my_id, msg.clone()).await; - } - } - - let packets = messages.drain(..).flat_map(|(id, msg)| { - msg.to_packets() - .into_iter() - .map(move |packet| (id, packet.into_boxed_slice())) - }); - - for (peer, packet) in packets { - if let Some(peer) = peer { - let channel = socket.channel_mut(0); - channel.send(packet, PeerId(peer)); - } else { - let channel = socket.channel_mut(0); - - for peer in all_peers.iter() { - // TODO: Any way around having to clone here? - let data = packet.clone(); - channel.send(data, *peer); - } - } - } - } - - pub fn cancel(&self) { - self.cancel_token.cancel(); - } - - pub async fn transport_loop(&self) { - let (mut socket, loop_fut) = WebRtcSocket::new_reliable(&self.ws_url); - - let loop_fut = async { - let msg = match loop_fut.await { - Ok(_) => TransportMessage::Disconnected, - Err(e) => { - let msg = match e { - SocketError::ConnectionFailed(e) => { - format!("Failed to connect to server: {e}") - } - SocketError::Disconnected(e) => { - format!("Disconnected from server, network error or kick: {e}") - } - }; - TransportMessage::Error(msg) - } - }; - self.push_incoming(self.my_id.read().await.unwrap_or_default(), msg) - .await; - } - .fuse(); - - tokio::pin!(loop_fut); - - let mut all_peers = HashSet::::with_capacity(20); - let mut my_id = None; - - let mut timer = tokio::time::interval(Duration::from_millis(100)); - - let mut partial_packets = - HashMap::>>)>::with_capacity(3); - - loop { - for (peer, state) in socket.update_peers() { - let msg = match state { - PeerState::Connected => { - all_peers.insert(peer); - TransportMessage::PeerConnect - } - PeerState::Disconnected => { - all_peers.remove(&peer); - TransportMessage::PeerDisconnect - } - }; - self.push_incoming(peer.0, msg).await; - } - - let messages = socket.channel_mut(0).receive(); - - let mut messages = messages - .into_iter() - .filter_map(|(id, data)| { - let msg = TransportMessage::deserialize(&data).ok(); - - if let Some(TransportMessage::Seq(TransportChunk { - id: multipart_id, - current, - total, - data, - })) = msg - { - if let Some((_, map)) = partial_packets.get_mut(&multipart_id) { - map.insert(current, Some(data)); - } else { - let mut map = HashMap::from_iter((0..total).map(|idx| (idx, None))); - map.insert(current, Some(data)); - partial_packets.insert(multipart_id, (id.0, map)); - } - None - } else { - msg.map(|msg| (id.0, msg)) - } - }) - .collect::>(); - - let complete_messages = partial_packets - .keys() - .copied() - .filter(|id| { - partial_packets - .get(id) - .is_some_and(|(_, v)| v.values().all(Option::is_some)) - }) - .collect::>(); - - for id in complete_messages { - let (peer, packet_map) = partial_packets.remove(&id).unwrap(); - - let res = TransportMessage::from_packets( - packet_map - .into_values() - .map(|v| v.unwrap().into_boxed_slice()), - ); - - match res { - Ok(msg) => messages.push((peer, msg)), - Err(why) => error!("Error receiving message: {why:?}"), - } - } - - let push_iter = self - .incoming - .0 - .reserve_many(messages.len()) - .await - .expect("Couldn't reserve space"); - - for (sender, msg) in push_iter.zip(messages.into_iter()) { - sender.send(msg); - } - - if my_id.is_none() { - if let Some(new_id) = socket.id() { - my_id = Some(new_id.0); - *self.my_id.write().await = Some(new_id.0); - self.push_incoming(new_id.0, TransportMessage::IdAssigned(new_id.0)) - .await; - } - } - - let mut outgoing_rx = self.outgoing.1.lock().await; - - let mut buffer = Vec::with_capacity(30); - - tokio::select! { - _ = self.cancel_token.cancelled() => { - // Break if cancelled externally - break; - } - - _ = &mut loop_fut => { - // Break if disconnected - break; - } - - // Rerun every tick - _ = timer.tick() => {} - - _ = outgoing_rx.recv_many(&mut buffer, 30) => { - // Handle sending new messages - self.handle_send(&mut socket, &all_peers, &mut buffer).await; - } - } - } - - drop(socket); - loop_fut.await - } -} - -impl Transport for MatchboxTransport { - async fn receive_messages(&self) -> impl Iterator { - self.recv_transport_messages() - .await - .into_iter() - .filter_map(|(id, msg)| match msg { - TransportMessage::Game(game_event) => Some(*game_event), - TransportMessage::PeerDisconnect => Some(GameEvent::DroppedPlayer(id)), - TransportMessage::Disconnected => Some(GameEvent::TransportDisconnect), - TransportMessage::Error(err) => Some(GameEvent::TransportError(err)), - _ => None, - }) - } - - async fn send_message(&self, msg: GameEvent) { - self.send_transport_message(None, msg.into()).await; - } - - fn disconnect(&self) { - self.cancel(); - } -} diff --git a/frontend/package-lock.json b/frontend/package-lock.json index 92ae411..80fbd66 100644 --- a/frontend/package-lock.json +++ b/frontend/package-lock.json @@ -1080,9 +1080,9 @@ } }, "node_modules/@rolldown/pluginutils": { - "version": "1.0.0-beta.11", - "resolved": "https://registry.npmjs.org/@rolldown/pluginutils/-/pluginutils-1.0.0-beta.11.tgz", - "integrity": "sha512-L/gAA/hyCSuzTF1ftlzUSI/IKr2POHsv1Dd78GfqkR83KMNuswWD61JxGV2L7nRwBBBSDr6R1gCkdTmoN7W4ag==", + "version": "1.0.0-beta.19", + "resolved": "https://registry.npmjs.org/@rolldown/pluginutils/-/pluginutils-1.0.0-beta.19.tgz", + "integrity": "sha512-3FL3mnMbPu0muGOCaKAhhFEYmqv9eTfPSJRJmANrCwtgK8VuxpsZDGK+m0LYAGoyO8+0j5uRe4PeyPDK1yA/hA==", "dev": true, "license": "MIT" }, @@ -1511,17 +1511,17 @@ } }, "node_modules/@typescript-eslint/eslint-plugin": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.34.1.tgz", - "integrity": "sha512-STXcN6ebF6li4PxwNeFnqF8/2BNDvBupf2OPx2yWNzr6mKNGF7q49VM00Pz5FaomJyqvbXpY6PhO+T9w139YEQ==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-8.35.0.tgz", + "integrity": "sha512-ijItUYaiWuce0N1SoSMrEd0b6b6lYkYt99pqCPfybd+HKVXtEvYhICfLdwp42MhiI5mp0oq7PKEL+g1cNiz/Eg==", "dev": true, "license": "MIT", "dependencies": { "@eslint-community/regexpp": "^4.10.0", - "@typescript-eslint/scope-manager": "8.34.1", - "@typescript-eslint/type-utils": "8.34.1", - "@typescript-eslint/utils": "8.34.1", - "@typescript-eslint/visitor-keys": "8.34.1", + "@typescript-eslint/scope-manager": "8.35.0", + "@typescript-eslint/type-utils": "8.35.0", + "@typescript-eslint/utils": "8.35.0", + "@typescript-eslint/visitor-keys": "8.35.0", "graphemer": "^1.4.0", "ignore": "^7.0.0", "natural-compare": "^1.4.0", @@ -1535,7 +1535,7 @@ "url": "https://opencollective.com/typescript-eslint" }, "peerDependencies": { - "@typescript-eslint/parser": "^8.34.1", + "@typescript-eslint/parser": "^8.35.0", "eslint": "^8.57.0 || ^9.0.0", "typescript": ">=4.8.4 <5.9.0" } @@ -1551,16 +1551,16 @@ } }, "node_modules/@typescript-eslint/parser": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/parser/-/parser-8.34.1.tgz", - "integrity": "sha512-4O3idHxhyzjClSMJ0a29AcoK0+YwnEqzI6oz3vlRf3xw0zbzt15MzXwItOlnr5nIth6zlY2RENLsOPvhyrKAQA==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/parser/-/parser-8.35.0.tgz", + "integrity": "sha512-6sMvZePQrnZH2/cJkwRpkT7DxoAWh+g6+GFRK6bV3YQo7ogi3SX5rgF6099r5Q53Ma5qeT7LGmOmuIutF4t3lA==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/scope-manager": "8.34.1", - "@typescript-eslint/types": "8.34.1", - "@typescript-eslint/typescript-estree": "8.34.1", - "@typescript-eslint/visitor-keys": "8.34.1", + "@typescript-eslint/scope-manager": "8.35.0", + "@typescript-eslint/types": "8.35.0", + "@typescript-eslint/typescript-estree": "8.35.0", + "@typescript-eslint/visitor-keys": "8.35.0", "debug": "^4.3.4" }, "engines": { @@ -1576,14 +1576,14 @@ } }, "node_modules/@typescript-eslint/project-service": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/project-service/-/project-service-8.34.1.tgz", - "integrity": "sha512-nuHlOmFZfuRwLJKDGQOVc0xnQrAmuq1Mj/ISou5044y1ajGNp2BNliIqp7F2LPQ5sForz8lempMFCovfeS1XoA==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/project-service/-/project-service-8.35.0.tgz", + "integrity": "sha512-41xatqRwWZuhUMF/aZm2fcUsOFKNcG28xqRSS6ZVr9BVJtGExosLAm5A1OxTjRMagx8nJqva+P5zNIGt8RIgbQ==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/tsconfig-utils": "^8.34.1", - "@typescript-eslint/types": "^8.34.1", + "@typescript-eslint/tsconfig-utils": "^8.35.0", + "@typescript-eslint/types": "^8.35.0", "debug": "^4.3.4" }, "engines": { @@ -1598,14 +1598,14 @@ } }, "node_modules/@typescript-eslint/scope-manager": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/scope-manager/-/scope-manager-8.34.1.tgz", - "integrity": "sha512-beu6o6QY4hJAgL1E8RaXNC071G4Kso2MGmJskCFQhRhg8VOH/FDbC8soP8NHN7e/Hdphwp8G8cE6OBzC8o41ZA==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/scope-manager/-/scope-manager-8.35.0.tgz", + "integrity": "sha512-+AgL5+mcoLxl1vGjwNfiWq5fLDZM1TmTPYs2UkyHfFhgERxBbqHlNjRzhThJqz+ktBqTChRYY6zwbMwy0591AA==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/types": "8.34.1", - "@typescript-eslint/visitor-keys": "8.34.1" + "@typescript-eslint/types": "8.35.0", + "@typescript-eslint/visitor-keys": "8.35.0" }, "engines": { "node": "^18.18.0 || ^20.9.0 || >=21.1.0" @@ -1616,9 +1616,9 @@ } }, "node_modules/@typescript-eslint/tsconfig-utils": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.34.1.tgz", - "integrity": "sha512-K4Sjdo4/xF9NEeA2khOb7Y5nY6NSXBnod87uniVYW9kHP+hNlDV8trUSFeynA2uxWam4gIWgWoygPrv9VMWrYg==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/tsconfig-utils/-/tsconfig-utils-8.35.0.tgz", + "integrity": "sha512-04k/7247kZzFraweuEirmvUj+W3bJLI9fX6fbo1Qm2YykuBvEhRTPl8tcxlYO8kZZW+HIXfkZNoasVb8EV4jpA==", "dev": true, "license": "MIT", "engines": { @@ -1633,14 +1633,14 @@ } }, "node_modules/@typescript-eslint/type-utils": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/type-utils/-/type-utils-8.34.1.tgz", - "integrity": "sha512-Tv7tCCr6e5m8hP4+xFugcrwTOucB8lshffJ6zf1mF1TbU67R+ntCc6DzLNKM+s/uzDyv8gLq7tufaAhIBYeV8g==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/type-utils/-/type-utils-8.35.0.tgz", + "integrity": "sha512-ceNNttjfmSEoM9PW87bWLDEIaLAyR+E6BoYJQ5PfaDau37UGca9Nyq3lBk8Bw2ad0AKvYabz6wxc7DMTO2jnNA==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/typescript-estree": "8.34.1", - "@typescript-eslint/utils": "8.34.1", + "@typescript-eslint/typescript-estree": "8.35.0", + "@typescript-eslint/utils": "8.35.0", "debug": "^4.3.4", "ts-api-utils": "^2.1.0" }, @@ -1657,9 +1657,9 @@ } }, "node_modules/@typescript-eslint/types": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/types/-/types-8.34.1.tgz", - "integrity": "sha512-rjLVbmE7HR18kDsjNIZQHxmv9RZwlgzavryL5Lnj2ujIRTeXlKtILHgRNmQ3j4daw7zd+mQgy+uyt6Zo6I0IGA==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/types/-/types-8.35.0.tgz", + "integrity": "sha512-0mYH3emanku0vHw2aRLNGqe7EXh9WHEhi7kZzscrMDf6IIRUQ5Jk4wp1QrledE/36KtdZrVfKnE32eZCf/vaVQ==", "dev": true, "license": "MIT", "engines": { @@ -1671,16 +1671,16 @@ } }, "node_modules/@typescript-eslint/typescript-estree": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/typescript-estree/-/typescript-estree-8.34.1.tgz", - "integrity": "sha512-rjCNqqYPuMUF5ODD+hWBNmOitjBWghkGKJg6hiCHzUvXRy6rK22Jd3rwbP2Xi+R7oYVvIKhokHVhH41BxPV5mA==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/typescript-estree/-/typescript-estree-8.35.0.tgz", + "integrity": "sha512-F+BhnaBemgu1Qf8oHrxyw14wq6vbL8xwWKKMwTMwYIRmFFY/1n/9T/jpbobZL8vp7QyEUcC6xGrnAO4ua8Kp7w==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/project-service": "8.34.1", - "@typescript-eslint/tsconfig-utils": "8.34.1", - "@typescript-eslint/types": "8.34.1", - "@typescript-eslint/visitor-keys": "8.34.1", + "@typescript-eslint/project-service": "8.35.0", + "@typescript-eslint/tsconfig-utils": "8.35.0", + "@typescript-eslint/types": "8.35.0", + "@typescript-eslint/visitor-keys": "8.35.0", "debug": "^4.3.4", "fast-glob": "^3.3.2", "is-glob": "^4.0.3", @@ -1739,16 +1739,16 @@ } }, "node_modules/@typescript-eslint/utils": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/utils/-/utils-8.34.1.tgz", - "integrity": "sha512-mqOwUdZ3KjtGk7xJJnLbHxTuWVn3GO2WZZuM+Slhkun4+qthLdXx32C8xIXbO1kfCECb3jIs3eoxK3eryk7aoQ==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/utils/-/utils-8.35.0.tgz", + "integrity": "sha512-nqoMu7WWM7ki5tPgLVsmPM8CkqtoPUG6xXGeefM5t4x3XumOEKMoUZPdi+7F+/EotukN4R9OWdmDxN80fqoZeg==", "dev": true, "license": "MIT", "dependencies": { "@eslint-community/eslint-utils": "^4.7.0", - "@typescript-eslint/scope-manager": "8.34.1", - "@typescript-eslint/types": "8.34.1", - "@typescript-eslint/typescript-estree": "8.34.1" + "@typescript-eslint/scope-manager": "8.35.0", + "@typescript-eslint/types": "8.35.0", + "@typescript-eslint/typescript-estree": "8.35.0" }, "engines": { "node": "^18.18.0 || ^20.9.0 || >=21.1.0" @@ -1763,13 +1763,13 @@ } }, "node_modules/@typescript-eslint/visitor-keys": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/@typescript-eslint/visitor-keys/-/visitor-keys-8.34.1.tgz", - "integrity": "sha512-xoh5rJ+tgsRKoXnkBPFRLZ7rjKM0AfVbC68UZ/ECXoDbfggb9RbEySN359acY1vS3qZ0jVTVWzbtfapwm5ztxw==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/@typescript-eslint/visitor-keys/-/visitor-keys-8.35.0.tgz", + "integrity": "sha512-zTh2+1Y8ZpmeQaQVIc/ZZxsx8UzgKJyNg1PTvjzC7WMhPSVS8bfDX34k1SrwOf016qd5RU3az2UxUNue3IfQ5g==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/types": "8.34.1", + "@typescript-eslint/types": "8.35.0", "eslint-visitor-keys": "^4.2.1" }, "engines": { @@ -1781,16 +1781,16 @@ } }, "node_modules/@vitejs/plugin-react": { - "version": "4.5.2", - "resolved": "https://registry.npmjs.org/@vitejs/plugin-react/-/plugin-react-4.5.2.tgz", - "integrity": "sha512-QNVT3/Lxx99nMQWJWF7K4N6apUEuT0KlZA3mx/mVaoGj3smm/8rc8ezz15J1pcbcjDK0V15rpHetVfya08r76Q==", + "version": "4.6.0", + "resolved": "https://registry.npmjs.org/@vitejs/plugin-react/-/plugin-react-4.6.0.tgz", + "integrity": "sha512-5Kgff+m8e2PB+9j51eGHEpn5kUzRKH2Ry0qGoe8ItJg7pqnkPrYPkDQZGgGmTa0EGarHrkjLvOdU3b1fzI8otQ==", "dev": true, "license": "MIT", "dependencies": { "@babel/core": "^7.27.4", "@babel/plugin-transform-react-jsx-self": "^7.27.1", "@babel/plugin-transform-react-jsx-source": "^7.27.1", - "@rolldown/pluginutils": "1.0.0-beta.11", + "@rolldown/pluginutils": "1.0.0-beta.19", "@types/babel__core": "^7.20.5", "react-refresh": "^0.17.0" }, @@ -2399,9 +2399,9 @@ } }, "node_modules/electron-to-chromium": { - "version": "1.5.171", - "resolved": "https://registry.npmjs.org/electron-to-chromium/-/electron-to-chromium-1.5.171.tgz", - "integrity": "sha512-scWpzXEJEMrGJa4Y6m/tVotb0WuvNmasv3wWVzUAeCgKU0ToFOhUW6Z+xWnRQANMYGxN4ngJXIThgBJOqzVPCQ==", + "version": "1.5.172", + "resolved": "https://registry.npmjs.org/electron-to-chromium/-/electron-to-chromium-1.5.172.tgz", + "integrity": "sha512-fnKW9dGgmBfsebbYognQSv0CGGLFH1a5iV9EDYTBwmAQn+whbzHbLFlC+3XbHc8xaNtpO0etm8LOcRXs1qMRkQ==", "dev": true, "license": "ISC" }, @@ -4269,9 +4269,9 @@ } }, "node_modules/prettier": { - "version": "3.5.3", - "resolved": "https://registry.npmjs.org/prettier/-/prettier-3.5.3.tgz", - "integrity": "sha512-QQtaxnoDJeAkDvDKWCLiwIXkTgRhwYDEQCghU9Z6q03iyek/rxRh/2lC3HB7P8sWT2xC/y5JDctPLBIGzHKbhw==", + "version": "3.6.0", + "resolved": "https://registry.npmjs.org/prettier/-/prettier-3.6.0.tgz", + "integrity": "sha512-ujSB9uXHJKzM/2GBuE0hBOUgC77CN3Bnpqa+g80bkv3T3A93wL/xlzDATHhnhkzifz/UE2SNOvmbTz5hSkDlHw==", "dev": true, "license": "MIT", "bin": { @@ -5082,15 +5082,15 @@ } }, "node_modules/typescript-eslint": { - "version": "8.34.1", - "resolved": "https://registry.npmjs.org/typescript-eslint/-/typescript-eslint-8.34.1.tgz", - "integrity": "sha512-XjS+b6Vg9oT1BaIUfkW3M3LvqZE++rbzAMEHuccCfO/YkP43ha6w3jTEMilQxMF92nVOYCcdjv1ZUhAa1D/0ow==", + "version": "8.35.0", + "resolved": "https://registry.npmjs.org/typescript-eslint/-/typescript-eslint-8.35.0.tgz", + "integrity": "sha512-uEnz70b7kBz6eg/j0Czy6K5NivaYopgxRjsnAJ2Fx5oTLo3wefTHIbL7AkQr1+7tJCRVpTs/wiM8JR/11Loq9A==", "dev": true, "license": "MIT", "dependencies": { - "@typescript-eslint/eslint-plugin": "8.34.1", - "@typescript-eslint/parser": "8.34.1", - "@typescript-eslint/utils": "8.34.1" + "@typescript-eslint/eslint-plugin": "8.35.0", + "@typescript-eslint/parser": "8.35.0", + "@typescript-eslint/utils": "8.35.0" }, "engines": { "node": "^18.18.0 || ^20.9.0 || >=21.1.0" diff --git a/manhunt-logic/Cargo.toml b/manhunt-logic/Cargo.toml new file mode 100644 index 0000000..54ba696 --- /dev/null +++ b/manhunt-logic/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "manhunt-logic" +version = "0.1.0" +edition = "2024" + +[dependencies] +anyhow = "1.0.98" +chrono = { version = "0.4.41", features = ["serde", "now"] } +rand = { version = "0.9.1", features = ["thread_rng"] } +rand_chacha = "0.9.0" +serde = { version = "1.0.219", features = ["derive"] } +specta = { version = "=2.0.0-rc.22", features = ["uuid", "chrono", "derive"] } +tokio = { version = "1.45.1", features = ["macros", "rt", "sync", "time"] } +uuid = { version = "1.17.0", features = ["serde", "v4"] } diff --git a/backend/src/game/mod.rs b/manhunt-logic/src/game.rs similarity index 80% rename from backend/src/game/mod.rs rename to manhunt-logic/src/game.rs index f3462b0..8bcab77 100644 --- a/backend/src/game/mod.rs +++ b/manhunt-logic/src/game.rs @@ -1,28 +1,25 @@ +use anyhow::bail; use chrono::{DateTime, Utc}; -pub use events::GameEvent; -use powerups::PowerUpType; -pub use settings::GameSettings; -use std::{collections::HashMap, sync::Arc, time::Duration}; +use std::{sync::Arc, time::Duration}; use uuid::Uuid; use tokio::{sync::RwLock, time::MissedTickBehavior}; -mod events; -mod location; -mod powerups; -mod settings; -mod state; -mod transport; +use crate::StartGameInfo; +use crate::{prelude::*, transport::TransportMessage}; -use crate::prelude::*; - -pub use location::{Location, LocationService}; -pub use state::{GameHistory, GameState, GameUiState}; -pub use transport::Transport; +use crate::{ + game_events::GameEvent, + game_state::{GameHistory, GameState, GameUiState}, + location::LocationService, + powerups::PowerUpType, + settings::GameSettings, + transport::Transport, +}; pub type Id = Uuid; -/// Convenence alias for UTC DT +/// Convenience alias for UTC DT pub type UtcDT = DateTime; pub trait StateUpdateSender { @@ -42,15 +39,17 @@ pub struct Game { impl Game { pub fn new( - my_id: Id, interval: Duration, - initial_caught_state: HashMap, - settings: GameSettings, + start_info: StartGameInfo, transport: Arc, location: L, state_update_sender: S, ) -> Self { - let state = GameState::new(settings, my_id, initial_caught_state); + let state = GameState::new( + start_info.settings, + transport.self_id(), + start_info.initial_caught_state, + ); Self { transport, @@ -61,6 +60,10 @@ impl Game { } } + async fn send_event(&self, event: GameEvent) { + self.transport.send_message(event.into()).await; + } + pub async fn mark_caught(&self) { let mut state = self.state.write().await; let id = state.id; @@ -69,9 +72,7 @@ impl Game { // TODO: Maybe reroll for new powerups (specifically seeker ones) instead of just erasing it state.use_powerup(); - self.transport - .send_message(GameEvent::PlayerCaught(state.id)) - .await; + self.send_event(GameEvent::PlayerCaught(state.id)).await; } pub async fn clone_settings(&self) -> GameSettings { @@ -85,9 +86,7 @@ impl Game { pub async fn get_powerup(&self) { let mut state = self.state.write().await; state.get_powerup(); - self.transport - .send_message(GameEvent::PowerupDespawn(state.id)) - .await; + self.send_event(GameEvent::PowerupDespawn(state.id)).await; } pub async fn use_powerup(&self) { @@ -98,9 +97,7 @@ impl Game { PowerUpType::PingSeeker => {} PowerUpType::PingAllSeekers => { for seeker in state.iter_seekers() { - self.transport - .send_message(GameEvent::ForcePing(seeker, None)) - .await; + self.send_event(GameEvent::ForcePing(seeker, None)).await; } } PowerUpType::ForcePingOther => { @@ -108,16 +105,14 @@ impl Game { let target = state.random_other_hider().or_else(|| state.random_seeker()); if let Some(target) = target { - self.transport - .send_message(GameEvent::ForcePing(target, None)) - .await; + self.send_event(GameEvent::ForcePing(target, None)).await; } } } } } - async fn consume_event(&self, state: &mut GameState, event: GameEvent) -> Result { + async fn consume_event(&self, state: &mut GameState, event: GameEvent) { if !state.game_ended() { state.event_history.push((Utc::now(), event.clone())); } @@ -126,7 +121,7 @@ impl Game { GameEvent::Ping(player_ping) => state.add_ping(player_ping), GameEvent::ForcePing(target, display) => { if target != state.id { - return Ok(false); + return; } let ping = if let Some(display) = display { @@ -137,7 +132,7 @@ impl Game { if let Some(ping) = ping { state.add_ping(ping.clone()); - self.transport.send_message(GameEvent::Ping(ping)).await; + self.send_event(GameEvent::Ping(ping)).await; } } GameEvent::PowerupDespawn(_) => state.despawn_powerup(), @@ -145,23 +140,36 @@ impl Game { state.mark_caught(player); state.remove_ping(player); } - GameEvent::DroppedPlayer(id) => { - state.remove_player(id); - } - GameEvent::TransportDisconnect => { - return Ok(true); - } - GameEvent::TransportError(err) => { - bail!("Transport error: {err}"); - } GameEvent::PostGameSync(id, history) => { state.insert_player_location_history(id, history); } } self.state_update_sender.send_update(); + } - Ok(false) + async fn consume_message( + &self, + state: &mut GameState, + _id: Option, + msg: TransportMessage, + ) -> Result { + match msg { + TransportMessage::Game(event) => { + self.consume_event(state, *event).await; + Ok(false) + } + TransportMessage::PeerDisconnect(id) => { + state.remove_player(id); + Ok(false) + } + TransportMessage::Disconnected => { + // Expected disconnect, exit + Ok(true) + } + TransportMessage::Error(err) => bail!("Transport error: {err}"), + _ => Ok(false), + } } /// Perform a tick for a specific moment in time @@ -172,7 +180,7 @@ impl Game { if state.check_end_game() { // If we're at the point where the game is over, send out our location history let msg = GameEvent::PostGameSync(state.id, state.location_history.clone()); - self.transport.send_message(msg).await; + self.send_event(msg).await; send_update = true; } @@ -208,15 +216,15 @@ impl Game { // We have a powerup that lets us ping a seeker as us, use it. if let Some(seeker) = state.random_seeker() { state.use_powerup(); - self.transport - .send_message(GameEvent::ForcePing(seeker, Some(state.id))) + self.send_event(GameEvent::ForcePing(seeker, Some(state.id))) .await; state.start_pings(now); } } else { // No powerup, normal ping if let Some(ping) = state.create_self_ping() { - self.transport.send_message(GameEvent::Ping(ping)).await; + self.send_event(GameEvent::Ping(ping.clone())).await; + state.add_ping(ping); state.start_pings(now); } } @@ -248,8 +256,8 @@ impl Game { self.tick(&mut state, now).await; } - pub fn quit_game(&self) { - self.transport.disconnect(); + pub async fn quit_game(&self) { + self.transport.disconnect().await; } /// Main loop of the game, handles ticking and receiving messages from [Transport]. @@ -262,13 +270,13 @@ impl Game { tokio::select! { biased; - events = self.transport.receive_messages() => { + messages = self.transport.receive_messages() => { let mut state = self.state.write().await; - for event in events { - match self.consume_event(&mut state, event).await { + for (id, msg) in messages { + match self.consume_message(&mut state, id, msg).await { Ok(should_break) => { if should_break { - break 'game Ok(None); + break 'game Ok(None); } } Err(why) => { break 'game Err(why); } @@ -288,7 +296,7 @@ impl Game { } }; - self.transport.disconnect(); + self.transport.disconnect().await; res } @@ -296,53 +304,16 @@ impl Game { #[cfg(test)] mod tests { - use std::sync::Arc; + use std::{collections::HashMap, sync::Arc}; - use crate::game::{location::Location, settings::PingStartCondition}; + use crate::{ + location::Location, + settings::PingStartCondition, + tests::{DummySender, MockLocation, MockTransport}, + }; use super::*; - use tokio::{sync::Mutex, task::yield_now, test}; - - type GameEventRx = tokio::sync::mpsc::Receiver; - type GameEventTx = tokio::sync::mpsc::Sender; - - struct MockTransport { - rx: Mutex, - txs: Vec, - } - - impl Transport for MockTransport { - async fn receive_messages(&self) -> impl Iterator { - let mut rx = self.rx.lock().await; - let mut buf = Vec::with_capacity(20); - rx.recv_many(&mut buf, 20).await; - buf.into_iter() - } - - async fn send_message(&self, msg: GameEvent) { - for tx in self.txs.iter() { - tx.send(msg.clone()).await.expect("Failed to send msg"); - } - } - } - - struct MockLocation; - - impl LocationService for MockLocation { - fn get_loc(&self) -> Option { - Some(location::Location { - lat: 0.0, - long: 0.0, - heading: None, - }) - } - } - - struct DummySender; - - impl StateUpdateSender for DummySender { - fn send_update(&self) {} - } + use tokio::{task::yield_now, test}; type TestGame = Game; @@ -357,36 +328,24 @@ mod tests { impl MockMatch { pub fn new(settings: GameSettings, players: u32, seekers: u32) -> Self { - let uuids = (0..players) - .map(|_| uuid::Uuid::new_v4()) - .collect::>(); - - let channels = (0..players) - .map(|_| tokio::sync::mpsc::channel(10)) - .collect::>(); + let (uuids, transports) = MockTransport::create_mesh(players); let initial_caught_state = (0..players) .map(|id| (uuids[id as usize], id < seekers)) .collect::>(); - let txs = channels - .iter() - .map(|(tx, _)| tx.clone()) - .collect::>(); - let games = channels + let games = transports .into_iter() .enumerate() - .map(|(id, (_, rx))| { - let transport = MockTransport { - rx: Mutex::new(rx), - txs: txs.clone(), - }; + .map(|(id, transport)| { let location = MockLocation; + let start_info = StartGameInfo { + initial_caught_state: initial_caught_state.clone(), + settings: settings.clone(), + }; let game = TestGame::new( - uuids[id], INTERVAL, - initial_caught_state.clone(), - settings.clone(), + start_info, Arc::new(transport), location, DummySender, @@ -495,7 +454,11 @@ mod tests { }) .await; - // Game over, See TODO in main_loop for more assertions + // Extra tick for post-game syncing + mat.tick().await; + + mat.assert_all_states(|s| assert!(s.game_ended(), "Game {} has not ended", s.id)) + .await; } #[test] @@ -681,4 +644,34 @@ mod tests { }) .await; } + + #[test] + async fn test_player_dropped() { + let settings = mk_settings(); + let mat = MockMatch::new(settings, 4, 1); + + mat.start().await; + + let game = mat.game(2); + game.quit_game().await; + let id = game.state.read().await.id; + + mat.tick().await; + + mat.assert_all_states(|s| { + if s.id != id { + assert!( + s.get_ping(id).is_none(), + "Game {} has not removed 2 from pings", + s.id + ); + assert!( + s.get_caught(id).is_none(), + "Game {} has not removed 2 from caught state", + s.id + ); + } + }) + .await; + } } diff --git a/backend/src/game/events.rs b/manhunt-logic/src/game_events.rs similarity index 66% rename from backend/src/game/events.rs rename to manhunt-logic/src/game_events.rs index f9d7ce9..b74da8f 100644 --- a/backend/src/game/events.rs +++ b/manhunt-logic/src/game_events.rs @@ -1,6 +1,10 @@ use serde::{Deserialize, Serialize}; -use super::{location::Location, state::PlayerPing, Id, UtcDT}; +use crate::{ + game::{Id, UtcDT}, + game_state::PlayerPing, + location::Location, +}; /// An event used between players to update state #[derive(Debug, Clone, Serialize, Deserialize, specta::Type)] @@ -17,11 +21,4 @@ pub enum GameEvent { /// Contains location history of the given player, used after the game to sync location /// histories PostGameSync(Id, Vec<(UtcDT, Location)>), - /// A player has been disconnected and removed from the game (because of error or otherwise). - /// The player should be removed from all state - DroppedPlayer(Id), - /// The underlying transport has disconnected - TransportDisconnect, - /// The underlying transport encountered an error - TransportError(String), } diff --git a/backend/src/game/state.rs b/manhunt-logic/src/game_state.rs similarity index 99% rename from backend/src/game/state.rs rename to manhunt-logic/src/game_state.rs index c05fa18..8f8d22f 100644 --- a/backend/src/game/state.rs +++ b/manhunt-logic/src/game_state.rs @@ -10,13 +10,12 @@ use rand_chacha::ChaCha20Rng; use serde::{Deserialize, Serialize}; use uuid::Uuid; -use crate::game::GameEvent; - -use super::{ +use crate::{ + game::{Id, UtcDT}, + game_events::GameEvent, location::Location, powerups::PowerUpType, settings::{GameSettings, PingStartCondition}, - Id, UtcDT, }; #[derive(Debug, Clone, Serialize, Deserialize, specta::Type)] diff --git a/manhunt-logic/src/lib.rs b/manhunt-logic/src/lib.rs new file mode 100644 index 0000000..9a17cb6 --- /dev/null +++ b/manhunt-logic/src/lib.rs @@ -0,0 +1,27 @@ +mod game; +mod game_events; +mod game_state; +mod lobby; +mod location; +mod powerups; +mod profile; +mod settings; +#[cfg(test)] +mod tests; +mod transport; + +pub use game::{Game, StateUpdateSender}; +pub use game_events::GameEvent; +pub use game_state::{GameHistory, GameUiState}; +pub use lobby::{Lobby, LobbyMessage, LobbyState, StartGameInfo}; +pub use location::{Location, LocationService}; +pub use profile::PlayerProfile; +pub use settings::GameSettings; +pub use transport::{MsgPair, Transport, TransportMessage}; + +pub mod prelude { + use anyhow::Error as AnyhowError; + use std::result::Result as StdResult; + pub type Result = StdResult; + pub use anyhow::Context; +} diff --git a/manhunt-logic/src/lobby.rs b/manhunt-logic/src/lobby.rs new file mode 100644 index 0000000..f80cd5c --- /dev/null +++ b/manhunt-logic/src/lobby.rs @@ -0,0 +1,239 @@ +use std::{collections::HashMap, sync::Arc}; + +use anyhow::anyhow; +use serde::{Deserialize, Serialize}; +use tokio::sync::Mutex; +use uuid::Uuid; + +use crate::{ + game::StateUpdateSender, + prelude::*, + profile::PlayerProfile, + settings::GameSettings, + transport::{Transport, TransportMessage}, +}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct StartGameInfo { + pub settings: GameSettings, + pub initial_caught_state: HashMap, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum LobbyMessage { + /// Message sent on a new peer, to sync profiles + PlayerSync(Uuid, PlayerProfile), + /// Message sent on a new peer from the host, to sync game settings + HostPush(GameSettings), + /// Host signals starting the game + StartGame(StartGameInfo), + /// A player has switched teams + PlayerSwitch(Uuid, bool), +} + +#[derive(Clone, Serialize, Deserialize, specta::Type)] +pub struct LobbyState { + profiles: HashMap, + join_code: String, + /// True represents seeker, false hider + teams: HashMap, + self_id: Uuid, + is_host: bool, + settings: GameSettings, +} + +pub struct Lobby { + is_host: bool, + join_code: String, + state: Mutex, + transport: Arc, + state_updates: U, +} + +impl Lobby { + pub async fn new( + join_code: &str, + host: bool, + profile: PlayerProfile, + settings: GameSettings, + state_updates: U, + ) -> Result> { + let transport = T::initialize(join_code, host) + .await + .context("Failed to connect to lobby")?; + + let self_id = transport.self_id(); + + let lobby = Arc::new(Self { + transport, + state_updates, + is_host: host, + join_code: join_code.to_string(), + state: Mutex::new(LobbyState { + teams: HashMap::from_iter([(self_id, false)]), + join_code: join_code.to_string(), + profiles: HashMap::from_iter([(self_id, profile)]), + self_id, + is_host: host, + settings, + }), + }); + + Ok(lobby) + } + + fn emit_state_update(&self) { + self.state_updates.send_update(); + } + + async fn send_transport_message(&self, id: Option, msg: LobbyMessage) { + if let Some(id) = id { + self.transport.send_message_single(id, msg.into()).await + } else { + self.transport.send_message(msg.into()).await + } + } + + async fn signaling_mark_started(&self) { + self.transport.mark_room_started(&self.join_code).await + } + + async fn handle_lobby(&self, msg: LobbyMessage) -> Option { + let mut state = self.state.lock().await; + match msg { + LobbyMessage::PlayerSync(peer, player_profile) => { + state.profiles.insert(peer, player_profile); + } + LobbyMessage::HostPush(game_settings) => { + state.settings = game_settings; + } + LobbyMessage::StartGame(start_game_info) => { + return Some(start_game_info); + } + LobbyMessage::PlayerSwitch(peer, seeker) => { + state.teams.insert(peer, seeker); + } + } + None + } + + async fn handle_message( + &self, + peer: Option, + msg: TransportMessage, + ) -> Option>> { + match msg { + TransportMessage::Disconnected => Some(Ok(None)), + TransportMessage::Error(why) => Some(Err(anyhow!("Transport error: {why}"))), + TransportMessage::Game(game_event) => { + eprintln!("Peer {peer:?} sent a GameEvent???: {game_event:?}"); + None + } + TransportMessage::Lobby(lobby_message) => self + .handle_lobby(*lobby_message) + .await + .map(|start_game| Ok(Some(start_game))), + TransportMessage::PeerConnect(peer) => { + let mut state = self.state.lock().await; + state.teams.insert(peer, false); + let id = state.self_id; + let msg = LobbyMessage::PlayerSync(id, state.profiles[&id].clone()); + drop(state); + self.send_transport_message(Some(peer), msg).await; + if self.is_host { + let state = self.state.lock().await; + let msg = LobbyMessage::HostPush(state.settings.clone()); + drop(state); + self.send_transport_message(Some(peer), msg).await; + } + None + } + TransportMessage::PeerDisconnect(peer) => { + let mut state = self.state.lock().await; + if peer != state.self_id { + state.profiles.remove(&peer); + state.teams.remove(&peer); + } + None + } + } + } + + pub async fn main_loop(&self) -> Result> { + let res = 'lobby: loop { + self.emit_state_update(); + + let msgs = self.transport.receive_messages().await; + + for (peer, msg) in msgs { + if let Some(res) = self.handle_message(peer, msg).await { + break 'lobby res; + } + } + }; + + if res.is_err() { + self.transport.disconnect().await; + } + + res + } + + pub fn clone_transport(&self) -> Arc { + self.transport.clone() + } + + pub async fn clone_state(&self) -> LobbyState { + self.state.lock().await.clone() + } + + pub async fn clone_profiles(&self) -> HashMap { + let state = self.state.lock().await; + state.profiles.clone() + } + + /// Set self as seeker or hider + pub async fn switch_teams(&self, seeker: bool) { + let mut state = self.state.lock().await; + let id = state.self_id; + if let Some(state_seeker) = state.teams.get_mut(&id) { + *state_seeker = seeker; + } + drop(state); + let msg = LobbyMessage::PlayerSwitch(id, seeker); + self.send_transport_message(None, msg).await; + self.emit_state_update(); + } + + /// (Host) Update game settings + pub async fn update_settings(&self, new_settings: GameSettings) { + if self.is_host { + let mut state = self.state.lock().await; + state.settings = new_settings.clone(); + drop(state); + let msg = LobbyMessage::HostPush(new_settings); + self.send_transport_message(None, msg).await; + self.emit_state_update(); + } + } + + /// (Host) Start the game + pub async fn start_game(&self) { + if self.is_host { + let state = self.state.lock().await; + let start_game_info = StartGameInfo { + settings: state.settings.clone(), + initial_caught_state: state.teams.clone(), + }; + drop(state); + let msg = LobbyMessage::StartGame(start_game_info); + self.signaling_mark_started().await; + self.transport.send_self(msg.clone().into()).await; + self.send_transport_message(None, msg).await; + } + } + + pub async fn quit_lobby(&self) { + self.transport.disconnect().await; + } +} diff --git a/backend/src/game/location.rs b/manhunt-logic/src/location.rs similarity index 100% rename from backend/src/game/location.rs rename to manhunt-logic/src/location.rs diff --git a/backend/src/game/powerups.rs b/manhunt-logic/src/powerups.rs similarity index 100% rename from backend/src/game/powerups.rs rename to manhunt-logic/src/powerups.rs diff --git a/manhunt-logic/src/profile.rs b/manhunt-logic/src/profile.rs new file mode 100644 index 0000000..25bd212 --- /dev/null +++ b/manhunt-logic/src/profile.rs @@ -0,0 +1,7 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Clone, Default, Debug, Serialize, Deserialize, specta::Type)] +pub struct PlayerProfile { + display_name: String, + pfp_base64: Option, +} diff --git a/backend/src/game/settings.rs b/manhunt-logic/src/settings.rs similarity index 100% rename from backend/src/game/settings.rs rename to manhunt-logic/src/settings.rs diff --git a/manhunt-logic/src/tests.rs b/manhunt-logic/src/tests.rs new file mode 100644 index 0000000..4c817fd --- /dev/null +++ b/manhunt-logic/src/tests.rs @@ -0,0 +1,119 @@ +use std::{collections::HashMap, sync::Arc}; + +use tokio::sync::{Mutex, mpsc}; +use uuid::Uuid; + +use crate::{ + MsgPair, StateUpdateSender, Transport, TransportMessage, + location::{Location, LocationService}, + prelude::*, +}; + +type GameEventRx = mpsc::Receiver; +type GameEventTx = mpsc::Sender; + +pub struct MockTransport { + id: Uuid, + rx: Mutex, + txs: HashMap, +} + +impl MockTransport { + pub fn create_mesh(players: u32) -> (Vec, Vec) { + let uuids = (0..players) + .map(|_| uuid::Uuid::new_v4()) + .collect::>(); + let channels = (0..players) + .map(|_| tokio::sync::mpsc::channel(10)) + .collect::>(); + let txs = channels + .iter() + .enumerate() + .map(|(i, (tx, _))| (uuids[i], tx.clone())) + .collect::>(); + + let transports = channels + .into_iter() + .enumerate() + .map(|(i, (_tx, rx))| Self::new(uuids[i], rx, txs.clone())) + .collect::>(); + + (uuids, transports) + } + + fn new(id: Uuid, rx: GameEventRx, txs: HashMap) -> Self { + Self { + id, + rx: Mutex::new(rx), + txs, + } + } +} + +impl Transport for MockTransport { + async fn initialize(_code: &str, _host: bool) -> Result> { + let (_, rx) = mpsc::channel(5); + Ok(Arc::new(Self { + id: Uuid::default(), + rx: Mutex::new(rx), + txs: HashMap::default(), + })) + } + + async fn disconnect(&self) { + self.send_message(TransportMessage::PeerDisconnect(self.id)) + .await; + } + + async fn receive_messages(&self) -> impl Iterator { + let mut rx = self.rx.lock().await; + let mut buf = Vec::with_capacity(20); + rx.recv_many(&mut buf, 20).await; + buf.into_iter() + } + + async fn send_message(&self, msg: TransportMessage) { + for (_id, tx) in self.txs.iter().filter(|(id, _)| **id != self.id) { + tx.send((Some(self.id), msg.clone())) + .await + .expect("Failed to send msg"); + } + } + + async fn send_message_single(&self, peer: Uuid, msg: TransportMessage) { + if let Some(tx) = self.txs.get(&peer) { + tx.send((Some(self.id), msg)) + .await + .expect("Failed to send msg"); + } + } + + async fn send_self(&self, msg: TransportMessage) { + self.txs[&self.id] + .send((Some(self.id), msg)) + .await + .expect("Failed to send msg"); + } + + fn self_id(&self) -> Uuid { + self.id + } +} + +pub struct MockLocation; + +impl LocationService for MockLocation { + fn get_loc(&self) -> Option { + Some(crate::location::Location { + lat: 0.0, + long: 0.0, + heading: None, + }) + } +} + +pub struct DummySender; + +impl StateUpdateSender for DummySender { + fn send_update(&self) {} +} diff --git a/manhunt-logic/src/transport.rs b/manhunt-logic/src/transport.rs new file mode 100644 index 0000000..53c350e --- /dev/null +++ b/manhunt-logic/src/transport.rs @@ -0,0 +1,69 @@ +use std::sync::Arc; + +use super::{game_events::GameEvent, lobby::LobbyMessage}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub enum TransportMessage { + /// Message related to the actual game + /// Boxed for space reasons + Game(Box), + /// Message related to the pre-game lobby + Lobby(Box), + /// Internal message when peer connects + PeerConnect(Uuid), + /// Internal message when peer disconnects + PeerDisconnect(Uuid), + /// Event sent when the transport gets disconnected, used to help consumers know when to stop + /// consuming messages. Note this should represent a success state, the disconnect was + /// triggered by user action. + Disconnected, + /// Event when the transport encounters a critical error and needs to disconnect. + Error(String), +} + +impl From for TransportMessage { + fn from(v: GameEvent) -> Self { + Self::Game(Box::new(v)) + } +} + +impl From for TransportMessage { + fn from(v: LobbyMessage) -> Self { + Self::Lobby(Box::new(v)) + } +} + +pub type MsgPair = (Option, TransportMessage); + +pub trait Transport: Send + Sync { + /// Start the transport loop, This is expected to spawn a new job that will loop until + /// cancelled or an error occurs. + fn initialize( + code: &str, + host: bool, + ) -> impl std::future::Future, anyhow::Error>> + Send; + /// Get the local user's ID + fn self_id(&self) -> Uuid; + /// Check if a room is open to join, non-host players will call this with a code + fn room_joinable(&self, _code: &str) -> impl std::future::Future + Send { + async { true } + } + /// Request a room be marked unjoinable (due to a game starting), the host user will call this. + fn mark_room_started(&self, _code: &str) -> impl Future { + async {} + } + /// Receive an event + fn receive_messages(&self) -> impl Future>; + /// Send a message to a specific peer + fn send_message_single(&self, peer: Uuid, msg: TransportMessage) -> impl Future; + /// Send a message to all other peers + fn send_message(&self, msg: TransportMessage) -> impl Future; + /// Send a message to the local user + fn send_self(&self, msg: TransportMessage) -> impl Future; + /// Disconnect from the transport + fn disconnect(&self) -> impl Future { + async {} + } +} diff --git a/manhunt-transport/Cargo.toml b/manhunt-transport/Cargo.toml new file mode 100644 index 0000000..53e6bba --- /dev/null +++ b/manhunt-transport/Cargo.toml @@ -0,0 +1,20 @@ +[package] +name = "manhunt-transport" +version = "0.1.0" +edition = "2024" + +[dependencies] +anyhow = "1.0.98" +futures = "0.3.31" +log = "0.4.27" +matchbox_protocol = "0.12.0" +matchbox_socket = "0.12.0" +rmp-serde = "1.3.0" +serde = { version = "1.0.219", features = ["derive"] } +tokio = { version = "1.45.1", features = ["macros", "sync", "time", "rt"] } +tokio-util = "0.7.15" +uuid = { version = "1.17.0", features = ["serde"] } +manhunt-logic = { version = "0.1.0", path = "../manhunt-logic" } +rand = { version = "0.9.1", features = ["thread_rng"] } +reqwest = { version = "0.12.20", default-features = false, features = ["charset", "http2", "rustls-tls", "system-proxy"] } +const-str = "0.6.2" diff --git a/manhunt-transport/src/lib.rs b/manhunt-transport/src/lib.rs new file mode 100644 index 0000000..be777ba --- /dev/null +++ b/manhunt-transport/src/lib.rs @@ -0,0 +1,6 @@ +mod matchbox; +mod packets; +mod server; + +pub use matchbox::MatchboxTransport; +pub use server::{generate_join_code, room_exists}; diff --git a/manhunt-transport/src/matchbox.rs b/manhunt-transport/src/matchbox.rs new file mode 100644 index 0000000..b19587f --- /dev/null +++ b/manhunt-transport/src/matchbox.rs @@ -0,0 +1,294 @@ +use std::{pin::Pin, sync::Arc, time::Duration}; + +use anyhow::{Context, anyhow}; +use futures::FutureExt; +use log::error; +use matchbox_socket::{Error as SocketError, PeerId, PeerState, WebRtcSocket}; +use tokio::{ + sync::{Mutex, mpsc}, + task::yield_now, +}; +use tokio_util::sync::CancellationToken; +use uuid::Uuid; + +use manhunt_logic::{Transport, TransportMessage, prelude::*}; + +use crate::{packets::PacketHandler, server}; + +type QueuePair = (mpsc::Sender, Mutex>); +type MsgPair = (Option, TransportMessage); +type Queue = QueuePair; + +pub struct MatchboxTransport { + my_id: Uuid, + incoming: Queue, + outgoing: Queue, + cancel_token: CancellationToken, +} + +type LoopFutRes = Result<(), SocketError>; + +fn map_socket_error(err: SocketError) -> anyhow::Error { + match err { + SocketError::ConnectionFailed(e) => anyhow!("Connection to server failed: {e:?}"), + SocketError::Disconnected(e) => anyhow!("Connection to server lost: {e:?}"), + } +} + +impl MatchboxTransport { + pub async fn new(join_code: &str, is_host: bool) -> Result> { + let (itx, irx) = mpsc::channel(15); + let (otx, orx) = mpsc::channel(15); + + let ws_url = server::room_url(join_code, is_host); + + let (mut socket, mut loop_fut) = WebRtcSocket::new_reliable(&ws_url); + + let res = loop { + tokio::select! { + id = Self::wait_for_id(&mut socket) => { + if let Some(id) = id { + break Ok(id); + } + }, + res = &mut loop_fut => { + break Err(match res { + Ok(_) => anyhow!("Transport disconnected unexpectedly"), + Err(err) => map_socket_error(err) + }); + } + } + } + .context("While trying to join the lobby"); + + match res { + Ok(my_id) => { + let transport = Arc::new(Self { + my_id, + incoming: (itx, Mutex::new(irx)), + outgoing: (otx, Mutex::new(orx)), + cancel_token: CancellationToken::new(), + }); + + tokio::spawn({ + let transport = transport.clone(); + async move { + transport.main_loop(socket, loop_fut).await; + } + }); + + Ok(transport) + } + Err(why) => { + drop(socket); + loop_fut.await.context("While disconnecting")?; + Err(why) + } + } + } + + async fn wait_for_id(socket: &mut WebRtcSocket) -> Option { + if let Some(id) = socket.id() { + Some(id.0) + } else { + yield_now().await; + None + } + } + + async fn push_incoming(&self, id: Option, msg: TransportMessage) { + self.incoming + .0 + .send((id, msg)) + .await + .expect("Failed to push to incoming queue"); + } + + async fn push_many_incoming(&self, msgs: Vec) { + let senders = self + .incoming + .0 + .reserve_many(msgs.len()) + .await + .expect("Failed to reserve in incoming queue"); + + for (sender, msg) in senders.into_iter().zip(msgs.into_iter()) { + sender.send(msg); + } + } + + async fn main_loop( + &self, + mut socket: WebRtcSocket, + loop_fut: Pin + Send + 'static>>, + ) { + let loop_fut = async { + let msg = match loop_fut.await { + Ok(_) => TransportMessage::Disconnected, + Err(e) => { + let msg = map_socket_error(e).to_string(); + TransportMessage::Error(msg) + } + }; + self.push_incoming(None, msg).await; + } + .fuse(); + + tokio::pin!(loop_fut); + + let mut interval = tokio::time::interval(Duration::from_secs(1)); + + let mut outgoing_rx = self.outgoing.1.lock().await; + const MAX_MSG_SEND: usize = 30; + let mut message_buffer = Vec::with_capacity(MAX_MSG_SEND); + + let mut packet_handler = PacketHandler::default(); + + loop { + self.handle_peers(&mut socket).await; + + self.handle_recv(&mut socket, &mut packet_handler).await; + + tokio::select! { + biased; + + _ = self.cancel_token.cancelled() => { + break; + } + + _ = &mut loop_fut => { + break; + } + + _ = outgoing_rx.recv_many(&mut message_buffer, MAX_MSG_SEND) => { + let peers = socket.connected_peers().collect::>(); + self.handle_send(&mut socket, &peers, &mut message_buffer).await; + } + + _ = interval.tick() => { + continue; + } + } + } + } + + async fn handle_peers(&self, socket: &mut WebRtcSocket) { + for (peer, state) in socket.update_peers() { + let msg = match state { + PeerState::Connected => TransportMessage::PeerConnect(peer.0), + PeerState::Disconnected => TransportMessage::PeerDisconnect(peer.0), + }; + self.push_incoming(Some(peer.0), msg).await; + } + } + + async fn handle_send( + &self, + socket: &mut WebRtcSocket, + all_peers: &[PeerId], + messages: &mut Vec, + ) { + let encoded_messages = messages.drain(..).filter_map(|(id, msg)| { + match PacketHandler::message_to_packets(&msg) { + Ok(packets) => Some((id, packets)), + Err(why) => { + error!("Error encoding message to packets: {why:?}"); + None + } + } + }); + + let channel = socket.channel_mut(0); + + for (peer, packets) in encoded_messages { + if let Some(peer) = peer { + for packet in packets { + channel.send(packet.into_boxed_slice(), PeerId(peer)); + } + } else { + for packet in packets { + let boxed = packet.into_boxed_slice(); + for peer in all_peers { + channel.send(boxed.clone(), *peer); + } + } + } + } + } + + async fn handle_recv(&self, socket: &mut WebRtcSocket, handler: &mut PacketHandler) { + let data = socket.channel_mut(0).receive(); + let messages = data + .into_iter() + .filter_map( + |(peer, bytes)| match handler.consume_packet(peer.0, bytes.into_vec()) { + Ok(msg) => msg.map(|msg| (Some(peer.0), msg)), + Err(why) => { + error!("Error receiving message: {why}"); + None + } + }, + ) + .collect(); + self.push_many_incoming(messages).await; + } + + pub async fn send_transport_message(&self, peer: Option, msg: TransportMessage) { + self.outgoing + .0 + .send((peer, msg)) + .await + .expect("Failed to add to outgoing queue"); + } + + pub async fn recv_transport_messages(&self) -> Vec { + let mut incoming_rx = self.incoming.1.lock().await; + let mut buffer = Vec::with_capacity(60); + incoming_rx.recv_many(&mut buffer, 60).await; + buffer + } + + pub fn cancel(&self) { + self.cancel_token.cancel(); + } +} + +impl Transport for MatchboxTransport { + fn self_id(&self) -> Uuid { + self.my_id + } + + async fn receive_messages(&self) -> impl Iterator { + self.recv_transport_messages().await.into_iter() + } + + async fn send_message_single(&self, peer: Uuid, msg: TransportMessage) { + self.send_transport_message(Some(peer), msg).await; + } + + async fn send_message(&self, msg: TransportMessage) { + self.send_transport_message(None, msg).await; + } + + async fn send_self(&self, msg: TransportMessage) { + self.push_incoming(Some(self.my_id), msg).await; + } + + async fn room_joinable(&self, code: &str) -> bool { + server::room_exists(code).await.unwrap_or(false) + } + + async fn mark_room_started(&self, code: &str) { + if let Err(why) = server::mark_room_started(code).await { + error!("Failed to mark room {code} as started: {why:?}"); + } + } + + async fn disconnect(&self) { + self.cancel(); + } + + async fn initialize(code: &str, host: bool) -> Result> { + Self::new(code, host).await + } +} diff --git a/manhunt-transport/src/packets.rs b/manhunt-transport/src/packets.rs new file mode 100644 index 0000000..0fd0d82 --- /dev/null +++ b/manhunt-transport/src/packets.rs @@ -0,0 +1,224 @@ +use std::collections::HashMap; + +use anyhow::{anyhow, bail}; +use manhunt_logic::{prelude::*, TransportMessage}; +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +type PacketEncoded = Vec; +type PacketSet = Vec>; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Packet { + remaining_packets: SeqHeader, + data: Vec, +} + +type SeqHeader = u64; +const SEQ_HEADER_SIZE: usize = size_of::(); + +const MATCHBOX_MAX_SIZE: usize = 65535; +const PACKET_SIZE: usize = MATCHBOX_MAX_SIZE - SEQ_HEADER_SIZE; +const MAX_NUM_PACKETS: u64 = u64::MAX - 1; + +impl Packet { + pub fn from_raw_bytes(mut bytes: PacketEncoded) -> Result { + // First [SEQ_HEADER_SIZE] bytes are our sequence header, in little endian. + if bytes.len() > SEQ_HEADER_SIZE { + let rest = bytes.split_off(SEQ_HEADER_SIZE); + let header = bytes; + let header = header + .try_into() + .map_err(|_| anyhow!("Couldn't parse sequence header"))?; + let remaining_packets = SeqHeader::from_le_bytes(header); + // Remaining bytes are the data + Ok(Self { + remaining_packets, + data: rest, + }) + } else { + bail!("Incoming packet is not long enough"); + } + } + + pub fn into_bytes(self) -> PacketEncoded { + let header_encoded = self.remaining_packets.to_le_bytes(); + header_encoded + .into_iter() + .chain(self.data) + .collect::>() + } + + fn packets_needed(len: u64) -> Result { + if len >= MAX_NUM_PACKETS.saturating_mul(PACKET_SIZE as u64) { + bail!("Message is too long, refusing to send"); + } else if len == 0 { + bail!("Message is empty"); + } else { + Ok(len.div_ceil(PACKET_SIZE as u64)) + } + } +} + +#[derive(Debug, Clone, Default)] +pub struct PacketHandler { + partials: HashMap, +} + +impl PacketHandler { + fn message_to_bytes(msg: &TransportMessage) -> Result> { + rmp_serde::to_vec(&msg).context("Failed to serialize message") + } + + fn message_from_bytes(msg: &[u8]) -> Result { + rmp_serde::from_slice(msg).context("Failed to deserialize message") + } + + pub fn message_to_packets(msg: &TransportMessage) -> Result { + let mut bytes = Self::message_to_bytes(msg)?; + let needed_packets = Packet::packets_needed(bytes.len() as u64)?; + let mut packets = Vec::with_capacity(needed_packets as usize); + for i in 1..=needed_packets { + let remaining_packets = needed_packets - i; + let mut data = bytes.split_off(bytes.len().min(PACKET_SIZE)); + std::mem::swap(&mut data, &mut bytes); + packets.push( + Packet { + remaining_packets, + data, + } + .into_bytes(), + ); + } + + if !bytes.is_empty() { + bail!("Bytes not emptied?"); + } + + Ok(packets) + } + + fn decode_packet_set(set: PacketSet) -> Result { + let combined_bytes = set.into_iter().flatten().collect::>(); + Self::message_from_bytes(&combined_bytes) + } + + pub fn consume_packet( + &mut self, + peer: Uuid, + bytes: PacketEncoded, + ) -> Result> { + match Packet::from_raw_bytes(bytes).context("Failed to decode packet") { + Ok(Packet { + remaining_packets, + data, + }) => { + if remaining_packets == 0 { + let res = if let Some(mut partial) = self.partials.remove(&peer) { + partial.push(data); + Self::decode_packet_set(partial) + } else { + Self::message_from_bytes(&data) + }; + + Some(res).transpose() + } else { + let partial = self + .partials + .entry(peer) + .or_insert_with(|| Vec::with_capacity(remaining_packets as usize + 1)); + partial.push(data); + Ok(None) + } + } + Err(why) => { + // Remove current partial message if we received an invalid packet as the entire + // sequence will now be wrong. + self.partials.remove(&peer); + Err(why) + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_packets_needed() { + assert_eq!(Packet::packets_needed(5).unwrap(), 1, "5 bytes, one packet"); + assert_eq!( + Packet::packets_needed((MATCHBOX_MAX_SIZE + 12) as u64).unwrap(), + 2, + "MAX + 12 bytes, two packets" + ); + assert!( + Packet::packets_needed(0).is_err(), + "Empty packets disallowed" + ); + assert!( + Packet::packets_needed(u64::MAX).is_err(), + "Too many packets disallowed" + ); + } + + #[test] + fn test_basic_packet_handling() { + let mut handler = PacketHandler::default(); + + let msg = TransportMessage::Disconnected; + let data = + PacketHandler::message_to_packets(&msg).expect("Failed to make message into bytes"); + + assert_eq!(data.len(), 1); + + let data = data.into_iter().next().unwrap(); + + println!("dat: {data:?}"); + + let decoded = handler + .consume_packet(Uuid::default(), data) + .expect("Failed to load message from bytes") + .expect("Message not complete despite being less than PACKET_SIZE"); + + assert!( + matches!(decoded, TransportMessage::Disconnected), + "Transport message does not match input" + ); + } + + #[test] + fn test_multipart() { + // Adding random amount to make sure we account for remainders, etc. + let really_big_string = "a".repeat(MATCHBOX_MAX_SIZE * 5 + 35); + let really_big_message = TransportMessage::Error(really_big_string.clone()); + + let packets = + PacketHandler::message_to_packets(&really_big_message).expect("Failed to encode"); + + assert!(packets.len() > 1, "Saving in one packet"); + + let mut handler = PacketHandler::default(); + let mut res = None; + + for pack in packets { + assert!( + pack.len() <= MATCHBOX_MAX_SIZE, + "Packets aren't small enough, {} > {}", + pack.len(), + MATCHBOX_MAX_SIZE + ); + + res = handler + .consume_packet(Uuid::default(), pack) + .expect("Failed to decode"); + } + + if let Some(TransportMessage::Error(s)) = res { + assert_eq!(s, really_big_string, "internal strings aren't equal"); + } else { + panic!("Decoded is the wrong type or wasn't completed"); + } + } +} diff --git a/backend/src/server.rs b/manhunt-transport/src/server.rs similarity index 86% rename from backend/src/server.rs rename to manhunt-transport/src/server.rs index 2686323..2b99070 100644 --- a/backend/src/server.rs +++ b/manhunt-transport/src/server.rs @@ -1,6 +1,6 @@ use reqwest::StatusCode; -use crate::prelude::*; +use manhunt_logic::prelude::*; const fn server_host() -> &'static str { if let Some(host) = option_env!("SIGNAL_SERVER_HOST") { @@ -27,19 +27,11 @@ const fn server_secure() -> bool { } const fn server_ws_proto() -> &'static str { - if server_secure() { - "wss" - } else { - "ws" - } + if server_secure() { "wss" } else { "ws" } } const fn server_http_proto() -> &'static str { - if server_secure() { - "https" - } else { - "http" - } + if server_secure() { "https" } else { "http" } } const SERVER_HOST: &str = server_host(); @@ -77,3 +69,10 @@ pub async fn mark_room_started(code: &str) -> Result { .context("Server returned error")?; Ok(()) } + +pub fn generate_join_code() -> String { + // 5 character sequence of A-Z + (0..5) + .map(|_| (b'A' + rand::random_range(0..26)) as char) + .collect::() +}