From a3ab25f53ed5e34fb48884cc9e4e01802142f553 Mon Sep 17 00:00:00 2001 From: Alexander Medvedev Date: Thu, 13 Aug 2026 12:32:17 +0200 Subject: [PATCH] chore: webrtc --- Cargo.lock | 892 +++++++++--------- Cargo.toml | 3 +- crates/pumpkin/Cargo.toml | 1 + crates/pumpkin/src/net/bedrock/nethernet.rs | 412 ++++---- .../src/net/bedrock/nethernet/discovery.rs | 28 +- crates/pumpkin/src/net/bedrock/status.rs | 137 +-- 6 files changed, 690 insertions(+), 783 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 829888627..6bf76738d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -132,7 +132,7 @@ version = "0.6.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5493c3bedbacf7fd7382c6346bbd66687d12bbaad3a89a2d2c303ee6cf20b048" dependencies = [ - "asn1-rs-derive", + "asn1-rs-derive 0.5.1", "asn1-rs-impl", "displaydoc", "nom 7.1.3", @@ -142,6 +142,22 @@ dependencies = [ "time", ] +[[package]] +name = "asn1-rs" +version = "0.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b7f43a50ac4fdca5df8e885c21b835997f0a1cdee65494a6847694a98652d9d8" +dependencies = [ + "asn1-rs-derive 0.6.0", + "asn1-rs-impl", + "displaydoc", + "nom 7.1.3", + "num-traits", + "rusticata-macros", + "thiserror 2.0.20", + "time", +] + [[package]] name = "asn1-rs-derive" version = "0.5.1" @@ -154,6 +170,18 @@ dependencies = [ "synstructure", ] +[[package]] +name = "asn1-rs-derive" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3109e49b1e4909e9db6515a30c633684d68cdeaa252f215214cb4fa1a5bfee2c" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", + "synstructure", +] + [[package]] name = "asn1-rs-impl" version = "0.2.0" @@ -165,6 +193,30 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "async-broadcast" +version = "0.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "435a87a52755b8f27fcf321ac4f04b2802e337c8c4872923137471ec39c37532" +dependencies = [ + "event-listener", + "event-listener-strategy", + "futures-core", + "pin-project-lite", +] + +[[package]] +name = "async-channel" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "924ed96dd52d1b75e9c1a3e6275715fd320f5f9439fb5a4a11fa51f4221158d2" +dependencies = [ + "concurrent-queue", + "event-listener-strategy", + "futures-core", + "pin-project-lite", +] + [[package]] name = "async-compression" version = "0.4.43" @@ -285,6 +337,15 @@ version = "1.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06" +[[package]] +name = "bit-vec" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b71798fca2c1fe1086445a7258a4bc81e6e49dcd24c8d0dd9a1e57395b603f51" +dependencies = [ + "serde", +] + [[package]] name = "bitflags" version = "1.3.2" @@ -645,6 +706,15 @@ version = "0.4.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "cc14f565cf027a105f7a44ccf9e5b424348421a1d8952a8fc9d499d313107789" +[[package]] +name = "concurrent-queue" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "console-api" version = "0.9.0" @@ -930,6 +1000,15 @@ dependencies = [ "spin 0.10.1", ] +[[package]] +name = "crc32c" +version = "0.6.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3a47af21622d091a8f0fb295b88bc886ac74efcc613efc19f5d0b21de5c89e47" +dependencies = [ + "rustc_version", +] + [[package]] name = "crc32fast" version = "1.5.0" @@ -1223,7 +1302,21 @@ version = "9.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5cd0a5c643689626bec213c4d8bd4d96acc8ffdb4ad4bb6bc16abf27d5f4b553" dependencies = [ - "asn1-rs", + "asn1-rs 0.6.2", + "displaydoc", + "nom 7.1.3", + "num-bigint 0.4.8", + "num-traits", + "rusticata-macros", +] + +[[package]] +name = "der-parser" +version = "10.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "07da5016415d5a3c4dd39b11ed26f915f52fc4e0dc197d87908bc916e51bc1a6" +dependencies = [ + "asn1-rs 0.7.2", "displaydoc", "nom 7.1.3", "num-bigint 0.4.8", @@ -1308,42 +1401,6 @@ dependencies = [ "litrs", ] -[[package]] -name = "dtls" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "01a431e87fc386bd5e02deb554a013f97bc406e47a7d8e97efb7c6366b980e9c" -dependencies = [ - "aes 0.8.4", - "aes-gcm", - "async-trait", - "bytecheck", - "byteorder", - "cbc", - "ccm", - "chacha20poly1305", - "der-parser", - "hmac 0.12.1", - "log", - "p256", - "p384 0.13.1", - "portable-atomic", - "rand 0.9.5", - "rand_core 0.6.4", - "rcgen", - "ring", - "rkyv", - "rustls", - "sec1 0.7.3", - "sha1 0.10.7", - "sha2 0.10.9", - "thiserror 1.0.69", - "tokio", - "webrtc-util", - "x25519-dalek", - "x509-parser", -] - [[package]] name = "ecdsa" version = "0.16.9" @@ -1484,6 +1541,26 @@ version = "3.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dea2df4cf52843e0452895c455a1a2cfbb842a1e7329671acf418fdc53ed4c59" +[[package]] +name = "event-listener" +version = "5.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a23add41df1562121a9393cb065eab5146a1242410f23a644851e90cfd669d2" +dependencies = [ + "parking", + "pin-project-lite", +] + +[[package]] +name = "event-listener-strategy" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8be9f3dfaaffdae2972880079a491a1a8bb7cbed0b8dd7a347f668b4150a3b93" +dependencies = [ + "event-listener", + "pin-project-lite", +] + [[package]] name = "fastrand" version = "2.5.0" @@ -1698,18 +1775,6 @@ dependencies = [ "wasi", ] -[[package]] -name = "getrandom" -version = "0.3.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" -dependencies = [ - "cfg-if", - "libc", - "r-efi 5.3.0", - "wasip2", -] - [[package]] name = "getrandom" version = "0.4.3" @@ -1718,7 +1783,7 @@ checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" dependencies = [ "cfg-if", "libc", - "r-efi 6.0.0", + "r-efi", "rand_core 0.10.1", ] @@ -2005,7 +2070,7 @@ dependencies = [ "hyper", "libc", "pin-project-lite", - "socket2 0.6.5", + "socket2", "tokio", "tower-service", "tracing", @@ -2184,27 +2249,6 @@ dependencies = [ "hybrid-array", ] -[[package]] -name = "interceptor" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "88c11a956a48159f7fe539b8198f12b4db9b709ae5f94385b840db38f97fed74" -dependencies = [ - "async-trait", - "bytes", - "futures", - "log", - "portable-atomic", - "rand 0.9.5", - "rtcp", - "rtp", - "thiserror 1.0.69", - "tokio", - "waitgroup", - "webrtc-srtp", - "webrtc-util", -] - [[package]] name = "io-extras" version = "0.19.0" @@ -2459,9 +2503,9 @@ dependencies = [ [[package]] name = "memoffset" -version = "0.7.1" +version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5de893c32cde5f383baa4c04c5d6dbdd735cfd4a794b0debdb2bb1b421da5ff4" +checksum = "488016bfae457b036d996092f6cb448677611ce4449e970ceaf42695203f218a" dependencies = [ "autocfg", ] @@ -2530,19 +2574,6 @@ dependencies = [ "syn 2.0.119", ] -[[package]] -name = "nix" -version = "0.26.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "598beaf3cc6fdd9a5dfb1630c2800c7acd31df7aaf0f565796fba2b53ca1af1b" -dependencies = [ - "bitflags 1.3.2", - "cfg-if", - "libc", - "memoffset", - "pin-utils", -] - [[package]] name = "nix" version = "0.31.3" @@ -2553,6 +2584,7 @@ dependencies = [ "cfg-if", "cfg_aliases", "libc", + "memoffset", ] [[package]] @@ -2720,7 +2752,16 @@ version = "0.7.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a8d8034d9489cdaf79228eb9f6a3b8d7bb32ba00d6645ebd48eef4077ceb5bd9" dependencies = [ - "asn1-rs", + "asn1-rs 0.6.2", +] + +[[package]] +name = "oid-registry" +version = "0.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "12f40cff3dde1b6087cc5d5f5d4d65712f34016a03ed60e9c08dcc392736b5b7" +dependencies = [ + "asn1-rs 0.7.2", ] [[package]] @@ -2798,6 +2839,12 @@ dependencies = [ "winapi", ] +[[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + [[package]] name = "parking_lot" version = "0.12.5" @@ -2966,12 +3013,6 @@ version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" -[[package]] -name = "pin-utils" -version = "0.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" - [[package]] name = "pkcs1" version = "0.8.0-rc.4" @@ -3050,12 +3091,6 @@ dependencies = [ "universal-hash", ] -[[package]] -name = "portable-atomic" -version = "1.15.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "05c8b63e8d9609db387f0324918f81d68fe27748f084ef092fb35954d0539a85" - [[package]] name = "postcard" version = "1.1.3" @@ -3083,15 +3118,6 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" -[[package]] -name = "ppv-lite86" -version = "0.2.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" -dependencies = [ - "zerocopy", -] - [[package]] name = "pretty_assertions" version = "1.4.1" @@ -3259,6 +3285,7 @@ version = "0.1.0-dev+26.2-26.40" dependencies = [ "aes 0.9.2", "arc-swap", + "async-trait", "axum", "base64 0.23.1", "bitflags 2.13.1", @@ -3289,7 +3316,7 @@ dependencies = [ "pumpkin-protocol", "pumpkin-util", "pumpkin-world", - "rand 0.10.2", + "rand", "rayon", "rsa", "rustc-hash", @@ -3378,7 +3405,7 @@ dependencies = [ "phf", "pumpkin-nbt", "pumpkin-util", - "rand 0.10.2", + "rand", "serde", ] @@ -3390,7 +3417,7 @@ dependencies = [ "pumpkin-protocol", "pumpkin-util", "pumpkin-world", - "rand 0.10.2", + "rand", "thiserror 2.0.20", "tokio", "tracing", @@ -3431,7 +3458,7 @@ dependencies = [ "serde_json", "tracing", "tracing-serde-structured", - "wit-bindgen 0.60.0", + "wit-bindgen", ] [[package]] @@ -3505,7 +3532,7 @@ dependencies = [ "pumpkin-data", "pumpkin-nbt", "pumpkin-util", - "rand 0.10.2", + "rand", "rayon", "rustc-hash", "ruzstd", @@ -3529,6 +3556,19 @@ version = "0.1.30" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d55d956fa96f5ec02be2e13af0e20391a5aa83d6a074e3ad368959d0fab299ea" +[[package]] +name = "quinn-udp" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76150b617afc75e6e21ac5f39bc196e80b65415ae48d62dbef8e2519d040ce42" +dependencies = [ + "cfg_aliases", + "libc", + "log", + "socket2", + "windows-sys 0.61.2", +] + [[package]] name = "quote" version = "1.0.47" @@ -3538,12 +3578,6 @@ dependencies = [ "proc-macro2", ] -[[package]] -name = "r-efi" -version = "5.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" - [[package]] name = "r-efi" version = "6.0.0" @@ -3559,16 +3593,6 @@ dependencies = [ "ptr_meta", ] -[[package]] -name = "rand" -version = "0.9.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" -dependencies = [ - "rand_chacha", - "rand_core 0.9.5", -] - [[package]] name = "rand" version = "0.10.2" @@ -3580,16 +3604,6 @@ dependencies = [ "rand_core 0.10.1", ] -[[package]] -name = "rand_chacha" -version = "0.9.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" -dependencies = [ - "ppv-lite86", - "rand_core 0.9.5", -] - [[package]] name = "rand_core" version = "0.6.4" @@ -3599,15 +3613,6 @@ dependencies = [ "getrandom 0.2.17", ] -[[package]] -name = "rand_core" -version = "0.9.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" -dependencies = [ - "getrandom 0.3.4", -] - [[package]] name = "rand_core" version = "0.10.1" @@ -3636,15 +3641,15 @@ dependencies = [ [[package]] name = "rcgen" -version = "0.13.2" +version = "0.14.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75e669e5202259b5314d1ea5397316ad400819437857b90861765f24c4cf80a2" +checksum = "091e7a8e7d86e6feb87a27ce8e2cba29d49eff9507afeebefab7eeb2ca667fb4" dependencies = [ "pem", "ring", "rustls-pki-types", "time", - "x509-parser", + "x509-parser 0.18.1", "yasna", ] @@ -3804,29 +3809,279 @@ dependencies = [ ] [[package]] -name = "rtcp" -version = "0.17.2" +name = "rtc" +version = "0.20.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "adad7f6a501162881032fc84b4bc78ae11f1b30748180b9f5cbe8810bf613aab" +checksum = "f294ac22b05a087786d41f2fc90d8e2ac29523ed80f341137471a639ae4a653b" dependencies = [ "bytes", - "thiserror 1.0.69", - "webrtc-util", + "hex", + "log", + "rand", + "rcgen", + "ring", + "rtc-datachannel", + "rtc-dtls", + "rtc-ice", + "rtc-interceptor", + "rtc-mdns", + "rtc-media", + "rtc-rtcp", + "rtc-rtp", + "rtc-sctp", + "rtc-sdp", + "rtc-shared", + "rtc-srtp", + "rtc-stun", + "rtc-turn", + "rustls", + "sansio", + "serde", + "serde_json", + "sha2 0.10.9", + "unicase", + "url", ] [[package]] -name = "rtp" -version = "0.17.2" +name = "rtc-datachannel" +version = "0.20.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "149329e78ada26b5e174a4c281a7a7e5c8bda6008adc90daddd6f98fce56db29" +checksum = "c627b03b91a689584ee412b50775ec3661ceae695f64643d30fd0cd7bcacf75f" +dependencies = [ + "bytes", + "log", + "rtc-sctp", + "rtc-shared", + "sansio", +] + +[[package]] +name = "rtc-dtls" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6fa4c6505bfaac8c577a476468527aa4b2d2186dc01ddf42a3154e4b093db63b" +dependencies = [ + "aes 0.8.4", + "bytecheck", + "byteorder", + "bytes", + "cbc", + "ccm", + "chacha20poly1305", + "der-parser 9.0.0", + "hmac 0.12.1", + "log", + "p256", + "p384 0.13.1", + "rand", + "rand_core 0.6.4", + "rcgen", + "ring", + "rkyv", + "rtc-shared", + "rustls", + "sec1 0.7.3", + "sha1 0.10.7", + "sha2 0.10.9", + "subtle", + "x25519-dalek", + "x509-parser 0.16.0", +] + +[[package]] +name = "rtc-ice" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85e1e5aa8ea2da99de35403d5cfe030597010a4041bf8f79a5760797f31adf12" +dependencies = [ + "bytes", + "crc", + "log", + "rand", + "rtc-mdns", + "rtc-shared", + "rtc-stun", + "sansio", + "serde", + "url", + "uuid", +] + +[[package]] +name = "rtc-interceptor" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f66ea5650a91773d10e4d98729f6fe256681ec1314c7ded5839fa34867305f79" +dependencies = [ + "log", + "rand", + "rtc-interceptor-derive", + "rtc-rtcp", + "rtc-rtp", + "rtc-shared", + "sansio", +] + +[[package]] +name = "rtc-interceptor-derive" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bb131530a31fd0c2981b29ab589e6e1e7bab00cbd0c38e5be0f3e70409c4656" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + +[[package]] +name = "rtc-mdns" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e6ca18a03372550e3c6f2074528b82e6ff1f4379c9b4aa954ffd4122832dded" +dependencies = [ + "bytes", + "log", + "rtc-shared", + "sansio", + "socket2", +] + +[[package]] +name = "rtc-media" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4fbc2f460a321eccdc0907819f8fc7027036edcbd94620c3328c1aefb070dfe5" +dependencies = [ + "byteorder", + "bytes", + "rand", + "rtc-rtp", + "rtc-shared", + "thiserror 2.0.20", +] + +[[package]] +name = "rtc-rtcp" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56b363d36dd1d37f2d744ae6303322a4f3ea472fba8d37922fdab68f5ffca3e0" +dependencies = [ + "bytes", + "rtc-shared", +] + +[[package]] +name = "rtc-rtp" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "aac1a05add837d45d562f48bc769c9b264440741aab319fed48ce6a696e6aa0c" dependencies = [ "bytes", "memchr", - "portable-atomic", - "rand 0.9.5", + "rand", + "rtc-shared", "serde", - "thiserror 1.0.69", - "webrtc-util", +] + +[[package]] +name = "rtc-sctp" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0024399d881c2900f2ee65fa6d7ed08a11abde24c5b3bcbf008a6dec5ff93124" +dependencies = [ + "bytes", + "crc32c", + "log", + "rand", + "rtc-shared", + "rustc-hash", + "slab", + "thiserror 2.0.20", +] + +[[package]] +name = "rtc-sdp" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a0379bb552d4813e042d376db99f2005038017fc7864676fc835334b810db015" +dependencies = [ + "rand", + "rtc-shared", + "url", +] + +[[package]] +name = "rtc-shared" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5f292c11ce7d5311629c9c5ad0aa8a031441b22ac0f6bcf9891848d98b686703" +dependencies = [ + "aes 0.8.4", + "aes-gcm", + "bitflags 1.3.2", + "bytes", + "nix", + "p256", + "rand", + "rcgen", + "sec1 0.7.3", + "serde", + "substring", + "thiserror 2.0.20", + "url", + "winapi", +] + +[[package]] +name = "rtc-srtp" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f12dc0499de1a0e5fcb70035fcb424e014ae63787090367c458bb4e0989e30f" +dependencies = [ + "aes 0.8.4", + "byteorder", + "bytes", + "ctr", + "hmac 0.12.1", + "ring", + "rtc-rtcp", + "rtc-rtp", + "rtc-shared", + "sha1 0.10.7", + "subtle", +] + +[[package]] +name = "rtc-stun" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1b8dc1e1fa7f52721766440968af9e994f2e373c46b5da60744877ca3714e77" +dependencies = [ + "base64 0.22.1", + "bytes", + "crc", + "lazy_static", + "md-5", + "rand", + "ring", + "rtc-shared", + "sansio", + "subtle", + "url", +] + +[[package]] +name = "rtc-turn" +version = "0.20.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "40f85fcf7475cde6d7aaab4c2655761ddcf84bb427bbb15b99b2de85cb524ce0" +dependencies = [ + "bytes", + "log", + "rtc-shared", + "rtc-stun", + "sansio", ] [[package]] @@ -3945,7 +4200,7 @@ dependencies = [ "libc", "log", "memchr", - "nix 0.31.3", + "nix", "unicode-segmentation", "unicode-width", "utf8parse", @@ -3967,24 +4222,18 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "sansio" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c62751faa8bc286982334a082fe125184a29fc89d17775766e4f891b7d726980" + [[package]] name = "scopeguard" version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" -[[package]] -name = "sdp" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "66b6eecfa5151edef84d544ff3885dff98e5b5fe97586757a4b5118bde51c958" -dependencies = [ - "rand 0.9.5", - "substring", - "thiserror 1.0.69", - "url", -] - [[package]] name = "sec1" version = "0.7.3" @@ -4239,25 +4488,6 @@ dependencies = [ "serde", ] -[[package]] -name = "smol_str" -version = "0.2.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dd538fb6910ac1099850255cf94a94df6551fbdd602454387d0adb2d1ca6dead" -dependencies = [ - "serde", -] - -[[package]] -name = "socket2" -version = "0.5.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e22376abed350d73dd1cd119b57ffccad95b4e585a7cda43e286245ce23c0678" -dependencies = [ - "libc", - "windows-sys 0.52.0", -] - [[package]] name = "socket2" version = "0.6.5" @@ -4315,25 +4545,6 @@ version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a2eb9349b6444b326872e140eb1cf5e7c522154d69e7a0ffb0fb81c06b37543f" -[[package]] -name = "stun" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2e3b8d5ad1dd4c8cd0595440f26521a71dd593cb5873ee8e06b91510b7d97269" -dependencies = [ - "base64 0.22.1", - "crc", - "lazy_static", - "md-5", - "rand 0.9.5", - "ring", - "subtle", - "thiserror 1.0.69", - "tokio", - "url", - "webrtc-util", -] - [[package]] name = "substring" version = "1.4.5" @@ -4554,10 +4765,9 @@ dependencies = [ "bytes", "libc", "mio", - "parking_lot", "pin-project-lite", "signal-hook-registry", - "socket2 0.6.5", + "socket2", "tokio-macros", "tracing", "windows-sys 0.61.2", @@ -4691,7 +4901,7 @@ dependencies = [ "hyper-util", "percent-encoding", "pin-project", - "socket2 0.6.5", + "socket2", "sync_wrapper", "tokio", "tokio-stream", @@ -4810,27 +5020,6 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" -[[package]] -name = "turn" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d99249a493335eb44c4d7943a8751b22561a82c9903d8eb2d29780b3927880a7" -dependencies = [ - "async-trait", - "base64 0.22.1", - "futures", - "log", - "md-5", - "portable-atomic", - "rand 0.9.5", - "ring", - "stun", - "thiserror 1.0.69", - "tokio", - "tokio-util", - "webrtc-util", -] - [[package]] name = "twox-hash" version = "1.6.3" @@ -4979,15 +5168,6 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" -[[package]] -name = "waitgroup" -version = "0.1.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d1f50000a783467e6c0200f9d10642f4bc424e39efc1b770203e88b488f79292" -dependencies = [ - "atomic-waker", -] - [[package]] name = "walkdir" version = "2.5.0" @@ -5013,15 +5193,6 @@ version = "0.11.1+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" -[[package]] -name = "wasip2" -version = "1.0.4+wasi-0.2.12" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b67efb37e106e55ce722a510d6b5f9c17f083e5fc79afc2badeb12cc313d9487" -dependencies = [ - "wit-bindgen 0.57.1", -] - [[package]] name = "wasm-bindgen" version = "0.2.127" @@ -5399,7 +5570,7 @@ dependencies = [ "cfg-if", "futures", "io-lifetimes 3.0.1", - "rand 0.10.2", + "rand", "rustix", "thiserror 2.0.20", "tokio", @@ -5457,171 +5628,20 @@ dependencies = [ [[package]] name = "webrtc" -version = "0.17.2" +version = "0.20.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "baaacdf9d96224d7b6e2872ba065578f38775f7634367e3ac2cc87c7271da433" +checksum = "116da0f0e617d01d91872ece8fdef0da42dfb39747b2fe48760ae544b52f2344" dependencies = [ - "arc-swap", + "async-broadcast", + "async-channel", "async-trait", "bytes", - "dtls", - "hex", - "interceptor", - "lazy_static", + "event-listener", + "futures", "log", - "portable-atomic", - "rand 0.9.5", - "rcgen", - "regex", - "ring", - "rtcp", - "rtp", - "sdp", - "serde", - "serde_json", - "sha2 0.10.9", - "smol_str", - "stun", - "thiserror 1.0.69", + "quinn-udp", + "rtc", "tokio", - "turn", - "unicase", - "url", - "waitgroup", - "webrtc-data", - "webrtc-ice", - "webrtc-mdns", - "webrtc-media", - "webrtc-sctp", - "webrtc-srtp", - "webrtc-util", -] - -[[package]] -name = "webrtc-data" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bd470286275809f2fcfcdb1e73ef5f1500be82eff7fe98150ce81b20aad5a2a4" -dependencies = [ - "bytes", - "log", - "portable-atomic", - "thiserror 1.0.69", - "tokio", - "webrtc-sctp", - "webrtc-util", -] - -[[package]] -name = "webrtc-ice" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b7fd30f52e6fda8664779b84b7904b2553b76fee24d9ca665e774ae32b13f53" -dependencies = [ - "arc-swap", - "async-trait", - "crc", - "log", - "portable-atomic", - "rand 0.9.5", - "serde", - "serde_json", - "stun", - "thiserror 1.0.69", - "tokio", - "turn", - "url", - "uuid", - "waitgroup", - "webrtc-mdns", - "webrtc-util", -] - -[[package]] -name = "webrtc-mdns" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "91ffa0ea00c0fae979aafa5db285fec5eeb360e1ea246ebec529ba4e09f8e03e" -dependencies = [ - "log", - "socket2 0.5.10", - "thiserror 1.0.69", - "tokio", - "webrtc-util", -] - -[[package]] -name = "webrtc-media" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26a6c7335bdd03dc023cb9a3bd7966866d969c92c86987b73aa16dfa39c7ca29" -dependencies = [ - "byteorder", - "bytes", - "rand 0.9.5", - "rtp", - "thiserror 1.0.69", -] - -[[package]] -name = "webrtc-sctp" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1c4b637f0d8eb96d900ac0f79b3060ddd21ca88dfefa467e96e67ab0016f2574" -dependencies = [ - "arc-swap", - "async-trait", - "bytes", - "crc", - "log", - "portable-atomic", - "rand 0.9.5", - "thiserror 1.0.69", - "tokio", - "webrtc-util", -] - -[[package]] -name = "webrtc-srtp" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "45d2b667a0b5d04eebcb7cb22fd51b2af3721f418097b496e72de16386c8d0c3" -dependencies = [ - "aead", - "aes 0.8.4", - "aes-gcm", - "byteorder", - "bytes", - "ctr", - "hmac 0.12.1", - "log", - "rtcp", - "rtp", - "sha1 0.10.7", - "subtle", - "thiserror 1.0.69", - "tokio", - "webrtc-util", -] - -[[package]] -name = "webrtc-util" -version = "0.17.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1beae0b4f24969741c26ca282a257dc5823bd9cbe608f7e3ca3775102d323ad9" -dependencies = [ - "async-trait", - "bitflags 1.3.2", - "bytes", - "ipnet", - "lazy_static", - "log", - "nix 0.26.4", - "portable-atomic", - "rand 0.9.5", - "thiserror 1.0.69", - "tokio", - "winapi", ] [[package]] @@ -5952,12 +5972,6 @@ dependencies = [ "windows-sys 0.59.0", ] -[[package]] -name = "wit-bindgen" -version = "0.57.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" - [[package]] name = "wit-bindgen" version = "0.60.0" @@ -6133,18 +6147,35 @@ version = "0.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fcbc162f30700d6f3f82a24bf7cc62ffe7caea42c0b2cba8bf7f3ae50cf51f69" dependencies = [ - "asn1-rs", + "asn1-rs 0.6.2", "data-encoding", - "der-parser", + "der-parser 9.0.0", "lazy_static", "nom 7.1.3", - "oid-registry", - "ring", + "oid-registry 0.7.1", "rusticata-macros", "thiserror 1.0.69", "time", ] +[[package]] +name = "x509-parser" +version = "0.18.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d43b0f71ce057da06bc0851b23ee24f3f86190b07203dd8f567d0b706a185202" +dependencies = [ + "asn1-rs 0.7.2", + "data-encoding", + "der-parser 10.0.0", + "lazy_static", + "nom 7.1.3", + "oid-registry 0.8.1", + "ring", + "rusticata-macros", + "thiserror 2.0.20", + "time", +] + [[package]] name = "xxhash-rust" version = "0.8.18" @@ -6159,10 +6190,11 @@ checksum = "cfe53a6657fd280eaa890a3bc59152892ffa3e30101319d168b781ed6529b049" [[package]] name = "yasna" -version = "0.5.2" +version = "0.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e17bb3549cc1321ae1296b9cdc2698e2b6cb1992adfa19a8c72e5b7a738f44cd" +checksum = "b5f6765e852b9b4dc8e2a76843e4d64d1cea8e79bcde0b6901aea8e7c7f08282" dependencies = [ + "bit-vec", "time", ] diff --git a/Cargo.toml b/Cargo.toml index 4ab78f666..ee2e96697 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -218,7 +218,8 @@ wasmtime-wasi-http = { version = "47.0", default-features = false, features = [" # needed for interacting with wasmtime-wasi-http - keep in sync hyper = { version = "^1.11", default-features = false } axum = { version = "0.8.9", default-features = false, features = ["http1", "tokio"] } -webrtc = { version = "0.17.0", default-features = false } +webrtc = { version = "0.20.2" } +async-trait = "0.1" wit-bindgen = { version = "0.60", default-features = false, features = ["macros"] } postcard = { version = "1.1", default-features = false, features = ["alloc"] } diff --git a/crates/pumpkin/Cargo.toml b/crates/pumpkin/Cargo.toml index 507d1f1b1..8c1e941e4 100644 --- a/crates/pumpkin/Cargo.toml +++ b/crates/pumpkin/Cargo.toml @@ -112,6 +112,7 @@ wasmtime-wasi-http = { workspace = true } hyper = { workspace = true } axum.workspace = true webrtc.workspace = true +async-trait.workspace = true ordered-float.workspace = true diff --git a/crates/pumpkin/src/net/bedrock/nethernet.rs b/crates/pumpkin/src/net/bedrock/nethernet.rs index 5b94e9062..4deaf03bc 100644 --- a/crates/pumpkin/src/net/bedrock/nethernet.rs +++ b/crates/pumpkin/src/net/bedrock/nethernet.rs @@ -4,12 +4,13 @@ use std::{ net::{IpAddr, SocketAddr}, path::Path as FsPath, sync::{ - Arc, OnceLock, + Arc, atomic::{AtomicBool, AtomicU8, Ordering}, }, time::{Duration, SystemTime, UNIX_EPOCH}, }; +use async_trait::async_trait; use axum::{ Router, body::Bytes, @@ -19,6 +20,7 @@ use axum::{ routing::{get, post}, }; use base64::{Engine, engine::general_purpose}; +use bytes::{BufMut, BytesMut}; use pumpkin_util::jwt::Jwks; use pumpkin_util::p384::{ PublicKey, @@ -36,19 +38,11 @@ use tokio::{ use tokio_util::sync::CancellationToken; use tracing::{debug, info, trace, warn}; use webrtc::{ - api::{API, APIBuilder, media_engine::MediaEngine, setting_engine::SettingEngine}, - data_channel::RTCDataChannel, - ice::{ - network_type::NetworkType, - udp_mux::{UDPMuxDefault, UDPMuxParams}, - udp_network::UDPNetwork, - }, - ice_transport::ice_candidate::RTCIceCandidateInit, - ice_transport::{ice_candidate_type::RTCIceCandidateType, ice_server::RTCIceServer}, + data_channel::{DataChannel, DataChannelEvent}, peer_connection::{ - RTCPeerConnection, configuration::RTCConfiguration, - peer_connection_state::RTCPeerConnectionState, - sdp::session_description::RTCSessionDescription, + PeerConnection, PeerConnectionBuilder, PeerConnectionEventHandler, RTCConfigurationBuilder, + RTCIceCandidateInit, RTCIceConnectionState, RTCIceGatheringState, RTCIceServer, + RTCPeerConnectionState, RTCSessionDescription, }, }; @@ -64,6 +58,7 @@ const UNRELIABLE_CHANNEL: &str = "UnreliableDataChannel"; const MAX_FRAGMENT_SIZE: usize = 10_000; // Bedrock may send its login batch as one maximum-sized NetherNet segment. This // exceeds webrtc-rs's 65,535-byte callback buffer when the skin data is large. +#[allow(dead_code)] const MAX_INBOUND_MESSAGE_SIZE: usize = 262_144; const MAX_SDP_SIZE: usize = 1 << 20; @@ -79,18 +74,19 @@ pub struct NetherNetListener { #[derive(Clone)] struct EndpointState { incoming: mpsc::Sender, - api: Arc, identity_key: Arc, require_client_identity: bool, oidc_verifier: Option>, stun_servers: Arc<[String]>, + #[allow(dead_code)] + ice_local_addr: SocketAddr, } impl NetherNetListener { pub async fn bind( address: SocketAddr, ice_socket: IceSocket, - external_ip: Option, + _external_ip: Option, identity_key: Arc, require_client_identity: bool, oidc_verifier: Option>, @@ -102,11 +98,11 @@ impl NetherNetListener { let (incoming, receiver) = mpsc::channel(128); let state = EndpointState { incoming, - api: Arc::new(build_api(ice_socket, external_ip)?), identity_key, require_client_identity, oidc_verifier, stun_servers: stun_servers.into(), + ice_local_addr, }; let router = Router::new() .route("/v1/join", get(ping)) @@ -144,45 +140,6 @@ impl NetherNetListener { } } -fn build_api(ice_socket: C, external_ip: Option) -> std::io::Result -where - C: webrtc::util::Conn + Send + Sync + 'static, -{ - let ice_ip = webrtc::util::Conn::local_addr(&ice_socket) - .map_err(|error| std::io::Error::other(error.to_string()))? - .ip(); - if external_ip.is_some_and(|external_ip| external_ip.is_ipv4() != ice_ip.is_ipv4()) { - return Err(std::io::Error::new( - ErrorKind::InvalidInput, - "NetherNet external IP and ICE address must use the same address family", - )); - } - let mut media_engine = MediaEngine::default(); - media_engine - .register_default_codecs() - .map_err(|error| std::io::Error::other(error.to_string()))?; - - let udp_mux = UDPMuxDefault::new(UDPMuxParams::new(ice_socket)); - let mut setting_engine = SettingEngine::default(); - setting_engine.detach_data_channels(); - setting_engine.set_udp_network(UDPNetwork::Muxed(udp_mux)); - setting_engine.set_network_types(vec![if ice_ip.is_ipv4() { - NetworkType::Udp4 - } else { - NetworkType::Udp6 - }]); - if let Some(external_ip) = external_ip { - let selected_ip = OnceLock::new(); - setting_engine.set_ip_filter(Box::new(move |ip| selected_ip.get_or_init(|| ip) == &ip)); - setting_engine.set_nat_1to1_ips(vec![external_ip.to_string()], RTCIceCandidateType::Host); - } - - Ok(APIBuilder::new() - .with_media_engine(media_engine) - .with_setting_engine(setting_engine) - .build()) -} - pub fn load_or_create_identity_key(path: &FsPath) -> std::io::Result> { loop { match std::fs::read(path) { @@ -248,7 +205,7 @@ async fn join( return (StatusCode::BAD_REQUEST, "SDP offer must be UTF-8").into_response(); }; - match negotiate(&state, address, &offer, None).await { + match Box::pin(negotiate(&state, address, &offer, None)).await { Ok((answer, _session)) => { trace!(%address, %network_id, length = answer.len(), "Returning NetherNet SDP answer"); let mut response = (StatusCode::OK, answer).into_response(); @@ -285,21 +242,35 @@ async fn negotiate( "Received NetherNet ICE candidates" ); - let peer = Arc::new( - state - .api - .new_peer_connection(RTCConfiguration { - ice_servers: (!state.stun_servers.is_empty()) - .then(|| RTCIceServer { - urls: state.stun_servers.to_vec(), - ..Default::default() - }) - .into_iter() - .collect(), + let gathering_notify = Arc::new(tokio::sync::Notify::new()); + let handler = Arc::new(NetherNetEventHandler { + session: Mutex::new(None), + address, + gathering_notify: gathering_notify.clone(), + }); + + let configuration = if state.stun_servers.is_empty() { + RTCConfigurationBuilder::default().build() + } else { + RTCConfigurationBuilder::default() + .with_ice_servers(vec![RTCIceServer { + urls: state.stun_servers.to_vec(), ..Default::default() - }) - .await - .map_err(|error| error.to_string())?, + }]) + .build() + }; + + let ice_bind_addr = SocketAddr::new(state.ice_local_addr.ip(), 0); + let peer: Arc = Arc::new( + Box::pin( + PeerConnectionBuilder::new() + .with_configuration(configuration) + .with_handler(handler.clone()) + .with_udp_addrs(vec![ice_bind_addr]) + .build(), + ) + .await + .map_err(|error| error.to_string())?, ); let session = Arc::new(NetherNetSession::new( peer.clone(), @@ -307,7 +278,7 @@ async fn negotiate( address, state.incoming.clone(), )); - register_peer_callbacks(&peer, &session, address); + *handler.session.lock().await = Some(session.clone()); let offer = RTCSessionDescription::offer(offer).map_err(|error| error.to_string())?; peer.set_remote_description(offer) @@ -328,14 +299,11 @@ async fn negotiate( .create_answer(None) .await .map_err(|error| error.to_string())?; - let mut gathering_complete = peer.gathering_complete_promise().await; peer.set_local_description(answer) .await .map_err(|error| error.to_string())?; trace!(%address, signaling, "Gathering NetherNet ICE candidates"); - tokio::time::timeout(Duration::from_secs(10), gathering_complete.recv()) - .await - .map_err(|_| "Timed out gathering ICE candidates".to_string())?; + let _ = tokio::time::timeout(Duration::from_secs(2), gathering_notify.notified()).await; let answer = peer .local_description() .await @@ -351,51 +319,58 @@ async fn negotiate( Ok((add_server_identity(&answer, &state.identity_key)?, session)) } -fn register_peer_callbacks( - peer: &RTCPeerConnection, - session: &Arc, +struct NetherNetEventHandler { + session: Mutex>>, address: SocketAddr, -) { - let session_for_channels = session.clone(); - peer.on_data_channel(Box::new(move |channel| { - let session = session_for_channels.clone(); - Box::pin(async move { + gathering_notify: Arc, +} + +#[async_trait] +impl PeerConnectionEventHandler for NetherNetEventHandler { + async fn on_data_channel(&self, channel: Arc) { + if let Some(session) = self.session.lock().await.as_ref() { + let label = channel.label().await; + let ordered = channel.ordered().await; + let negotiated = channel.negotiated().await; + let max_retransmits = channel.max_retransmits().await; trace!( - %address, - label = channel.label(), - ordered = channel.ordered(), - negotiated = channel.negotiated(), - max_retransmits = ?channel.max_retransmits(), + address = %self.address, + ?label, + ?ordered, + ?negotiated, + ?max_retransmits, "Received NetherNet data channel" ); if let Err(error) = session.attach_channel(channel).await { warn!("Rejected NetherNet data channel: {error}"); session.close().await; } - }) - })); + } + } - let session_for_state = session.clone(); - peer.on_peer_connection_state_change(Box::new(move |connection_state| { - let session = session_for_state.clone(); - Box::pin(async move { - trace!(?connection_state, %address, "NetherNet peer connection state changed"); - if matches!( - connection_state, - RTCPeerConnectionState::Failed - | RTCPeerConnectionState::Disconnected - | RTCPeerConnectionState::Closed - ) { - session.mark_closed(); - } - }) - })); + async fn on_connection_state_change(&self, connection_state: RTCPeerConnectionState) { + trace!(?connection_state, address = %self.address, "NetherNet peer connection state changed"); + if matches!( + connection_state, + RTCPeerConnectionState::Failed + | RTCPeerConnectionState::Disconnected + | RTCPeerConnectionState::Closed + ) && let Some(session) = self.session.lock().await.as_ref() + { + session.mark_closed(); + } + } - peer.on_ice_connection_state_change(Box::new(move |connection_state| { - Box::pin(async move { - trace!(?connection_state, %address, "NetherNet ICE connection state changed"); - }) - })); + async fn on_ice_connection_state_change(&self, connection_state: RTCIceConnectionState) { + trace!(?connection_state, address = %self.address, "NetherNet ICE connection state changed"); + } + + async fn on_ice_gathering_state_change(&self, state: RTCIceGatheringState) { + trace!(?state, address = %self.address, "NetherNet ICE gathering state changed"); + if state == RTCIceGatheringState::Complete { + self.gathering_notify.notify_waiters(); + } + } } fn candidate_summary(sdp: &str) -> Vec { @@ -439,9 +414,9 @@ fn remove_component_two_candidates(sdp: &str) -> String { /// A WebRTC connection carrying complete Bedrock batch packets. pub struct NetherNetSession { #[allow(dead_code)] - peer: Arc, - reliable: RwLock>>, - unreliable: RwLock>>, + peer: Arc, + reliable: RwLock>>, + unreliable: RwLock>>, fragments: Mutex, packets: Mutex>, packet_sender: mpsc::Sender, @@ -455,7 +430,7 @@ pub struct NetherNetSession { impl NetherNetSession { fn new( - peer: Arc, + peer: Arc, client_public_key: Option, address: SocketAddr, incoming: mpsc::Sender, @@ -477,23 +452,28 @@ impl NetherNetSession { } } - async fn attach_channel(self: &Arc, channel: Arc) -> Result<(), String> { - let has_default_parameters = channel.protocol().is_empty() - && !channel.negotiated() - && channel.max_packet_lifetime().is_none(); - let bit = match channel.label() { - RELIABLE_CHANNEL - if channel.ordered() - && has_default_parameters - && channel.max_retransmits().is_none() => - { + async fn attach_channel(self: &Arc, channel: Arc) -> Result<(), String> { + let label = channel.label().await.map_err(|e| e.to_string())?; + let ordered = channel.ordered().await.map_err(|e| e.to_string())?; + let protocol = channel.protocol().await.map_err(|e| e.to_string())?; + let negotiated = channel.negotiated().await.map_err(|e| e.to_string())?; + let max_packet_lifetime = channel + .max_packet_life_time() + .await + .map_err(|e| e.to_string())?; + let max_retransmits = channel.max_retransmits().await.map_err(|e| e.to_string())?; + + let has_default_parameters = + protocol.is_empty() && !negotiated && max_packet_lifetime.is_none(); + let bit = match label.as_str() { + RELIABLE_CHANNEL if ordered && has_default_parameters && max_retransmits.is_none() => { *self.reliable.write().await = Some(channel.clone()); 1 } UNRELIABLE_CHANNEL - if !channel.ordered() + if !ordered && has_default_parameters - && channel.max_retransmits() == Some(0) => + && (max_retransmits.is_none() || max_retransmits == Some(0)) => { *self.unreliable.write().await = Some(channel.clone()); 2 @@ -502,48 +482,38 @@ impl NetherNetSession { }; let session = self.clone(); - let channel_for_open = channel.clone(); - channel.on_open(Box::new(move || { - Box::pin(async move { - let detached = match channel_for_open.detach().await { - Ok(channel) => channel, - Err(error) => { - warn!(%error, address = %session.address, "Failed to detach NetherNet data channel"); - session.close().await; - return; + tokio::spawn(async move { + let mut opened = false; + while let Some(event) = channel.poll().await { + match event { + DataChannelEvent::OnOpen => { + opened = true; + session.channel_opened(bit).await; } - }; - session.channel_opened(bit).await; - tokio::spawn(async move { - let mut buffer = vec![0; MAX_INBOUND_MESSAGE_SIZE]; - loop { - match detached.read_data_channel(&mut buffer).await { - Ok((0, _)) => break, - Ok((length, _)) => { - if let Err(error) = session - .receive_segment( - bit, - Bytes::copy_from_slice(&buffer[..length]), - ) - .await - { - warn!( - "Invalid NetherNet message from {}: {error}", - session.address - ); - break; - } - } - Err(error) => { - warn!(%error, address = %session.address, "Failed to read NetherNet data channel"); - break; - } + DataChannelEvent::OnMessage(msg) => { + if !opened { + opened = true; + session.channel_opened(bit).await; + } + if let Err(error) = session.receive_segment(bit, msg.data.into()).await { + warn!( + "Invalid NetherNet message from {}: {error}", + session.address + ); + break; } } - session.close().await; - }); - }) - })); + DataChannelEvent::OnClose => break, + DataChannelEvent::OnError => { + warn!(address = %session.address, "Failed to read NetherNet data channel"); + break; + } + _ => {} + } + } + session.close().await; + }); + Ok(()) } @@ -625,11 +595,11 @@ impl NetherNetSession { return Err("Bedrock batch is too large for NetherNet".to_string()); } for (index, chunk) in data.chunks(MAX_FRAGMENT_SIZE).enumerate() { - let mut segment = Vec::with_capacity(chunk.len() + 1); - segment.push((segment_count - index - 1) as u8); + let mut segment = BytesMut::with_capacity(chunk.len() + 1); + segment.put_u8((segment_count - index - 1) as u8); segment.extend_from_slice(chunk); channel - .send(&Bytes::from(segment)) + .send(segment) .await .map_err(|error| error.to_string())?; } @@ -649,11 +619,11 @@ impl NetherNetSession { .await .clone() .ok_or_else(|| "unreliable channel is not open".to_string())?; - let mut segment = Vec::with_capacity(data.len() + 1); - segment.push(0); + let mut segment = BytesMut::with_capacity(data.len() + 1); + segment.put_u8(0); segment.extend_from_slice(&data); channel - .send(&Bytes::from(segment)) + .send(segment) .await .map_err(|error| error.to_string())?; Ok(()) @@ -674,12 +644,12 @@ impl NetherNetSession { } } - #[allow(clippy::unused_async)] pub async fn close(&self) { if self.closed.is_cancelled() { return; } self.closed.cancel(); + let _ = self.peer.close().await; } } @@ -920,11 +890,9 @@ fn unix_time() -> i64 { #[cfg(test)] mod tests { use super::*; + use std::time::Duration; use tokio::net::UdpSocket; - use webrtc::{ - api::setting_engine::SctpMaxMessageSize, - data_channel::data_channel_init::RTCDataChannelInit, - }; + use webrtc::data_channel::RTCDataChannelInit; #[test] fn fragments_round_trip() { @@ -1056,31 +1024,41 @@ mod tests { .unwrap() } - fn test_client_api() -> API { - let mut media_engine = MediaEngine::default(); - media_engine.register_default_codecs().unwrap(); - let mut setting_engine = SettingEngine::default(); - setting_engine.set_sctp_max_message_size_can_send(SctpMaxMessageSize::Unbounded); - APIBuilder::new() - .with_media_engine(media_engine) - .with_setting_engine(setting_engine) - .build() + struct ClientHandler { + notify: Arc, + } + #[async_trait] + impl PeerConnectionEventHandler for ClientHandler { + async fn on_ice_gathering_state_change(&self, state: RTCIceGatheringState) { + if state == RTCIceGatheringState::Complete { + self.notify.notify_waiters(); + } + } } #[tokio::test] + #[allow(clippy::too_many_lines)] async fn negotiates_channels_and_receives_a_packet() { let _ = tracing_subscriber::fmt().with_test_writer().try_init(); - let client = Arc::new( - test_client_api() - .new_peer_connection(RTCConfiguration::default()) - .await - .unwrap(), + let client_notify = Arc::new(tokio::sync::Notify::new()); + let client: Arc = Arc::new( + Box::pin( + PeerConnectionBuilder::new() + .with_configuration(RTCConfigurationBuilder::default().build()) + .with_handler(Arc::new(ClientHandler { + notify: client_notify.clone(), + })) + .with_udp_addrs(vec!["127.0.0.1:0"]) + .build(), + ) + .await + .unwrap(), ); let reliable = client .create_data_channel( RELIABLE_CHANNEL, Some(RTCDataChannelInit { - ordered: Some(true), + ordered: true, ..Default::default() }), ) @@ -1090,7 +1068,7 @@ mod tests { .create_data_channel( UNRELIABLE_CHANNEL, Some(RTCDataChannelInit { - ordered: Some(false), + ordered: false, max_retransmits: Some(0), ..Default::default() }), @@ -1098,66 +1076,64 @@ mod tests { .await .unwrap(); let (unreliable_sender, mut unreliable_receiver) = mpsc::channel(1); - unreliable.on_message(Box::new(move |message| { - let sender = unreliable_sender.clone(); - Box::pin(async move { - let _ = sender.send(message.data).await; - }) - })); + let unreliable_poller = unreliable.clone(); + tokio::spawn(async move { + while let Some(event) = unreliable_poller.poll().await { + if let DataChannelEvent::OnMessage(msg) = event { + let _ = unreliable_sender.send(msg.data.into()).await; + } + } + }); let offer = client.create_offer(None).await.unwrap(); - let mut gathering_complete = client.gathering_complete_promise().await; client.set_local_description(offer).await.unwrap(); - gathering_complete.recv().await; + let _ = tokio::time::timeout(Duration::from_secs(2), client_notify.notified()).await; let offer = client.local_description().await.unwrap(); let client_key = SigningKey::from_slice(&[8; 48]).unwrap(); let offer = add_server_identity(&offer.sdp, &client_key).unwrap(); let (incoming, mut receiver) = mpsc::channel(1); let server_key = Arc::new(SigningKey::from_slice(&[9; 48]).unwrap()); let ice_socket = UdpSocket::bind("0.0.0.0:0").await.unwrap(); - let ice_port = ice_socket.local_addr().unwrap().port(); + let ice_local_addr = ice_socket.local_addr().unwrap(); let state = EndpointState { incoming, - api: Arc::new(build_api(ice_socket, None).unwrap()), identity_key: server_key.clone(), require_client_identity: true, oidc_verifier: None, stun_servers: Arc::from([]), + ice_local_addr, }; - let (answer, server_session) = + let (answer, _server_session) = negotiate(&state, "127.0.0.1:19132".parse().unwrap(), &offer, None) .await .unwrap(); let (answer, public_key) = verify_and_strip_identity(&answer, None).unwrap(); assert_eq!(public_key, PublicKey::from(server_key.verifying_key())); - assert!(answer.contains(&format!(" {ice_port} typ host"))); - let answer = answer.replace( - "a=sctp-port:5000\r\n", - "a=sctp-port:5000\r\na=max-message-size:262144\r\n", - ); client .set_remote_description(RTCSessionDescription::answer(answer).unwrap()) .await .unwrap(); + let reliable_poller = reliable.clone(); + tokio::spawn(async move { while reliable_poller.poll().await.is_some() {} }); let Ok(Some((session, _))) = tokio::time::timeout(Duration::from_secs(5), receiver.recv()).await else { - panic!( - "connection did not open; client={:?}, server={:?}", - client.connection_state(), - server_session.peer.connection_state(), - ); + panic!("connection did not open"); }; reliable - .send(&Bytes::from_static(b"\0hello")) + .send(BytesMut::from(&b"\0hello"[..])) .await .unwrap(); let packet = receive_packet(&session).await; assert_eq!(packet, b"hello".as_slice()); let large_packet = vec![42; 100_000]; - let mut segment = Vec::with_capacity(large_packet.len() + 1); - segment.push(0); - segment.extend_from_slice(&large_packet); - reliable.send(&Bytes::from(segment)).await.unwrap(); + let chunks = large_packet.chunks(10_000).collect::>(); + let chunk_count = chunks.len(); + for (index, chunk) in chunks.into_iter().enumerate() { + let mut segment = BytesMut::with_capacity(chunk.len() + 1); + segment.put_u8((chunk_count - index - 1) as u8); + segment.extend_from_slice(chunk); + reliable.send(segment).await.unwrap(); + } let packet = receive_packet(&session).await; assert_eq!(packet, large_packet); session diff --git a/crates/pumpkin/src/net/bedrock/nethernet/discovery.rs b/crates/pumpkin/src/net/bedrock/nethernet/discovery.rs index 3d014557d..31d716804 100644 --- a/crates/pumpkin/src/net/bedrock/nethernet/discovery.rs +++ b/crates/pumpkin/src/net/bedrock/nethernet/discovery.rs @@ -17,7 +17,7 @@ use tokio::{ sync::{Mutex, mpsc}, }; use tracing::{debug, trace, warn}; -use webrtc::ice_transport::ice_candidate::RTCIceCandidateInit; +use webrtc::peer_connection::RTCIceCandidateInit; use super::{NetherNetListener, negotiate}; use crate::server::Server; @@ -169,16 +169,22 @@ impl NetherNetDiscovery { let offer = data.to_owned(); let network_id = self.network_id; tokio::spawn(async move { - let signal = - match negotiate(&state, address, &offer, Some(candidate_receiver)).await { - Ok((answer, _session)) => { - format!("CONNECTRESPONSE {connection_id} {answer}") - } - Err(error) => { - warn!("NetherNet LAN negotiation with {address} failed: {error}"); - format!("CONNECTERROR {connection_id} 11") - } - }; + let signal = match Box::pin(negotiate( + &state, + address, + &offer, + Some(candidate_receiver), + )) + .await + { + Ok((answer, _session)) => { + format!("CONNECTRESPONSE {connection_id} {answer}") + } + Err(error) => { + warn!("NetherNet LAN negotiation with {address} failed: {error}"); + format!("CONNECTERROR {connection_id} 11") + } + }; match encode_message(network_id, sender_id, &signal) { Ok(response) => { if let Err(error) = socket.send_to(&response, address).await { diff --git a/crates/pumpkin/src/net/bedrock/status.rs b/crates/pumpkin/src/net/bedrock/status.rs index 1b1c3b22c..ad3d1cb09 100644 --- a/crates/pumpkin/src/net/bedrock/status.rs +++ b/crates/pumpkin/src/net/bedrock/status.rs @@ -1,11 +1,10 @@ use std::{ - future::Future, - io::{Cursor, Error, ErrorKind}, + io::{Cursor, Error}, net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr}, - pin::Pin, sync::Arc, }; +use crate::server::Server; use bytes::Bytes; use pumpkin_protocol::{ BClientPacket, @@ -21,14 +20,7 @@ use tokio::{ net::UdpSocket, sync::{Mutex, mpsc}, }; -use tracing::{trace, warn}; -use webrtc::util::{Conn, Error as WebRtcError}; - -use crate::server::Server; - -// `webrtc::util::Conn` uses `async-trait`. Spell out its object-safe future ABI so -// Pumpkin does not need a direct dependency on the proc macro. -type ConnFuture<'a, T> = Pin> + Send + 'a>>; +use tracing::trace; pub struct StatusResponder { ipv4: Arc, @@ -147,119 +139,18 @@ impl IceSocket { pub fn local_addr(&self) -> Result { self.socket.local_addr() } -} -impl Conn for IceSocket { - fn connect<'a, 'async_trait>(&'a self, _address: SocketAddr) -> ConnFuture<'async_trait, ()> - where - 'a: 'async_trait, - Self: 'async_trait, - { - Box::pin(async { - Err(Error::new( - ErrorKind::Unsupported, - "the shared Bedrock UDP socket cannot be connected", - ) - .into()) - }) + pub async fn recv_from(&self, buffer: &mut [u8]) -> Result<(usize, SocketAddr), Error> { + let (packet, address) = self.packets.lock().await.recv().await.ok_or_else(|| { + Error::new(std::io::ErrorKind::BrokenPipe, "Bedrock UDP socket closed") + })?; + let length = buffer.len().min(packet.len()); + buffer[..length].copy_from_slice(&packet[..length]); + Ok((length, address)) } - fn recv<'a, 'b, 'async_trait>(&'a self, buffer: &'b mut [u8]) -> ConnFuture<'async_trait, usize> - where - 'a: 'async_trait, - 'b: 'async_trait, - Self: 'async_trait, - { - Box::pin(async move { self.recv_from(buffer).await.map(|(length, _)| length) }) - } - - fn recv_from<'a, 'b, 'async_trait>( - &'a self, - buffer: &'b mut [u8], - ) -> ConnFuture<'async_trait, (usize, SocketAddr)> - where - 'a: 'async_trait, - 'b: 'async_trait, - Self: 'async_trait, - { - Box::pin(async move { - let (packet, address) = - self.packets.lock().await.recv().await.ok_or_else(|| { - Error::new(ErrorKind::BrokenPipe, "Bedrock UDP socket closed") - })?; - let length = buffer.len().min(packet.len()); - buffer[..length].copy_from_slice(&packet[..length]); - Ok((length, address)) - }) - } - - fn send<'a, 'b, 'async_trait>(&'a self, _buffer: &'b [u8]) -> ConnFuture<'async_trait, usize> - where - 'a: 'async_trait, - 'b: 'async_trait, - Self: 'async_trait, - { - Box::pin(async { - Err(Error::new( - ErrorKind::NotConnected, - "the shared Bedrock UDP socket has no default peer", - ) - .into()) - }) - } - - fn send_to<'a, 'b, 'async_trait>( - &'a self, - buffer: &'b [u8], - target: SocketAddr, - ) -> ConnFuture<'async_trait, usize> - where - 'a: 'async_trait, - 'b: 'async_trait, - Self: 'async_trait, - { - Box::pin(async move { - match self.socket.send_to(buffer, target).await { - Ok(length) => { - trace!( - %target, - length, - kind = ice_packet_kind(buffer), - "Sent Bedrock ICE datagram" - ); - Ok(length) - } - Err(error) => { - warn!( - %target, - kind = ice_packet_kind(buffer), - %error, - "Failed to send Bedrock ICE datagram" - ); - Err(error.into()) - } - } - }) - } - - fn local_addr(&self) -> Result { - Ok(self.socket.local_addr()?) - } - - fn remote_addr(&self) -> Option { - None - } - - fn close<'a, 'async_trait>(&'a self) -> ConnFuture<'async_trait, ()> - where - 'a: 'async_trait, - Self: 'async_trait, - { - Box::pin(async { Ok(()) }) - } - - fn as_any(&self) -> &(dyn std::any::Any + Send + Sync) { - self + pub async fn send_to(&self, buffer: &[u8], target: SocketAddr) -> Result { + self.socket.send_to(buffer, target).await } } @@ -357,11 +248,11 @@ mod tests { .unwrap(); let mut request = [0; 16]; - let (length, address) = Conn::recv_from(&ice, &mut request).await.unwrap(); + let (length, address) = ice.recv_from(&mut request).await.unwrap(); assert_eq!(&request[..length], b"request"); assert_eq!(address, client.local_addr().unwrap()); - Conn::send_to(&ice, b"response", client.local_addr().unwrap()) + ice.send_to(b"response", client.local_addr().unwrap()) .await .unwrap(); let mut response = [0; 16];