diff --git a/Cargo.lock b/Cargo.lock index 7e9344b..c66c5de 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1826,7 +1826,6 @@ dependencies = [ "p384", "pumpkin-codecs", "pumpkin-nbt", - "rayon", "reqwest", "rsa", "rustls", @@ -1930,8 +1929,7 @@ checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" [[package]] name = "rayon" version = "1.12.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fb39b166781f92d482534ef4b4b1b2568f42613b53e5b6c160e24cfbfa30926d" +source = "git+https://github.com/guybedford/rayon?branch=fallback-spawn#39b1d8c838d4d8c446249db526ec096999155f1b" dependencies = [ "either", "rayon-core", @@ -1940,8 +1938,7 @@ dependencies = [ [[package]] name = "rayon-core" version = "1.13.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22e18b0f0062d30d4230b2e85ff77fdfe4326feb054b9783a3460d8435c8ab91" +source = "git+https://github.com/guybedford/rayon?branch=fallback-spawn#39b1d8c838d4d8c446249db526ec096999155f1b" dependencies = [ "crossbeam-deque", "crossbeam-utils", @@ -2617,7 +2614,7 @@ dependencies = [ [[package]] name = "tokio" version = "1.53.1" -source = "git+https://github.com/guybedford/tokio?tag=1.53.1-cf.emscripten#d56f1eff2c3fd72a946e1e536a3b93e32b002a80" +source = "git+https://github.com/guybedford/tokio?tag=1.53.1-cf.emscripten#7227f2d72739c3bc9821074cc704e7582be92907" dependencies = [ "bytes", "libc", @@ -2631,7 +2628,7 @@ dependencies = [ [[package]] name = "tokio-macros" version = "2.7.2" -source = "git+https://github.com/guybedford/tokio?tag=1.53.1-cf.emscripten#d56f1eff2c3fd72a946e1e536a3b93e32b002a80" +source = "git+https://github.com/guybedford/tokio?tag=1.53.1-cf.emscripten#7227f2d72739c3bc9821074cc704e7582be92907" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml index 37719e5..e154c02 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -37,6 +37,9 @@ tokio-macros = { git = "https://github.com/guybedford/tokio", tag = "1.53.1-cf.e # tokio-rs/mio#1969; its emscripten epoll bindings come from libc's unreleased libc-0.2 branch. mio = { git = "https://github.com/guybedford/mio", tag = "1.2.3-cf.emscripten" } libc = { git = "https://github.com/rust-lang/libc", branch = "libc-0.2" } +# rayon-rs/rayon#1323: the single-threaded fallback wakes the host to run spawned jobs. +rayon = { git = "https://github.com/guybedford/rayon", branch = "fallback-spawn" } +rayon-core = { git = "https://github.com/guybedford/rayon", branch = "fallback-spawn" } # getrandom-backed SystemRandom on emscripten. ring = { git = "https://github.com/guybedford/ring", branch = "emscripten" } diff --git a/docs/dependencies.md b/docs/dependencies.md index 25d88d8..93c2658 100644 --- a/docs/dependencies.md +++ b/docs/dependencies.md @@ -28,6 +28,7 @@ commit. | --- | --- | --- | | mio | https://github.com/guybedford/mio tag `1.2.3-cf.emscripten` | Emscripten epoll selector (tokio-rs/mio#1969) | | libc | https://github.com/rust-lang/libc branch `libc-0.2` | Emscripten epoll bindings, unreleased | +| rayon, rayon-core | https://github.com/guybedford/rayon branch `fallback-spawn` | rayon-rs/rayon#1323: the single-threaded fallback wakes the host to run spawned jobs | | ring | https://github.com/guybedford/ring branch `emscripten` | getrandom-backed `SystemRandom` on Emscripten | The root Cargo.toml applies these overrides and the local checkouts; the diff --git a/patches/pumpkin-emscripten.patch b/patches/pumpkin-emscripten.patch index f711f73..ed83550 100644 --- a/patches/pumpkin-emscripten.patch +++ b/patches/pumpkin-emscripten.patch @@ -6,10 +6,11 @@ plugin loading, the console, and system crash reporting. Native builds retain these subsystems and thread-backed execution by default. Share asynchronous ticker and chunk-scheduler loops between native and -single-threaded builds. Run Rayon, chunk-generation and blocking jobs as -Tokio tasks on single-threaded builds so each takes its own turn of the -event loop, await scheduler notifications and results, and track -cooperative tasks through shutdown. Drop the Tokio features that +single-threaded builds. Run chunk-generation and blocking jobs as Tokio +tasks on single-threaded builds so each takes its own turn of the event +loop (Rayon's own spawn queues on its single-threaded fallback, drained by +the host through rayon's fallback wake hook), await scheduler notifications +and results, and track cooperative tasks through shutdown. Drop the Tokio features that need OS threads and signals from single-threaded builds and guard unsupported Emscripten system calls and synchronous HTTP lookups. @@ -50,16 +51,10 @@ index 143ba8a3..57314ca3 100644 thiserror = { version = "2.0", default-features = false } diff --git a/crates/pumpkin-util/Cargo.toml b/crates/pumpkin-util/Cargo.toml -index f4fb121f..0a4a3076 100644 +index f4fb121f..b5815dd6 100644 --- a/crates/pumpkin-util/Cargo.toml +++ b/crates/pumpkin-util/Cargo.toml -@@ -6,11 +6,13 @@ rust-version.workspace = true - license.workspace = true - - [dependencies] -+rayon.workspace = true - pumpkin-nbt.workspace = true - pumpkin-codecs.workspace = true +@@ -11,6 +11,7 @@ pumpkin-codecs.workspace = true serde.workspace = true serde_json.workspace = true bytes.workspace = true @@ -67,40 +62,24 @@ index f4fb121f..0a4a3076 100644 num-derive.workspace = true num-traits.workspace = true -@@ -39,6 +41,8 @@ criterion = { workspace = true, features = ["html_reports"] } +@@ -39,6 +40,8 @@ criterion = { workspace = true, features = ["html_reports"] } [features] default = [] -+# Run CPU tasks on the current thread (see rayon_spawn and spawn_blocking). ++# Run blocking tasks as Tokio tasks (see spawn_blocking). +single-threaded = [] codegen = ["dep:syn", "dep:quote", "dep:proc-macro2"] [[bench]] diff --git a/crates/pumpkin-util/src/lib.rs b/crates/pumpkin-util/src/lib.rs -index 5119797b..35d78184 100644 +index 5119797b..a2bfe678 100644 --- a/crates/pumpkin-util/src/lib.rs +++ b/crates/pumpkin-util/src/lib.rs -@@ -310,3 +310,39 @@ impl TryFrom for Hand { +@@ -310,3 +310,23 @@ impl TryFrom for Hand { } } } + -+/// Spawn a fire-and-forget CPU task. -+/// -+/// On normal targets this hands the closure to rayon's global thread pool. -+/// Single-threaded builds have no rayon workers (a bare `rayon::spawn` would never -+/// run), so the closure becomes a Tokio task and runs on a later turn of the loop. -+pub fn rayon_spawn(f: F) { -+ #[cfg(feature = "single-threaded")] -+ { -+ tokio::spawn(async move { f() }); -+ } -+ #[cfg(not(feature = "single-threaded"))] -+ { -+ rayon::spawn(f); -+ } -+} -+ +/// Schedule bounded synchronous work using the selected execution model. +/// +/// Single-threaded builds run the closure in a regular Tokio task, preserving @@ -152,18 +131,9 @@ index f1675a8f..82c21917 100644 tempfile.workspace = true criterion = { workspace = true, features = ["html_reports", "async_tokio"] } diff --git a/crates/pumpkin-world/src/chunk/format/anvil.rs b/crates/pumpkin-world/src/chunk/format/anvil.rs -index 38de745f..60cedf53 100644 +index 38de745f..8b987941 100644 --- a/crates/pumpkin-world/src/chunk/format/anvil.rs +++ b/crates/pumpkin-world/src/chunk/format/anvil.rs -@@ -798,7 +798,7 @@ impl ChunkSerializer for AnvilChunkFile< - - let (tx, mut rx) = tokio::sync::mpsc::channel(chunk_items.len().max(1)); - -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - use rayon::prelude::*; - chunk_items - .into_par_iter() @@ -810,6 +810,9 @@ impl ChunkSerializer for AnvilChunkFile< Err(err) => LoadedData::Error((chunk, err)), }, @@ -175,18 +145,9 @@ index 38de745f..60cedf53 100644 }); }); diff --git a/crates/pumpkin-world/src/chunk/format/linear.rs b/crates/pumpkin-world/src/chunk/format/linear.rs -index 3927c66f..fd76019b 100644 +index 3927c66f..d47a611b 100644 --- a/crates/pumpkin-world/src/chunk/format/linear.rs +++ b/crates/pumpkin-world/src/chunk/format/linear.rs -@@ -599,7 +599,7 @@ impl ChunkSerializer for LinearV2File - - let (tx, mut rx) = tokio::sync::mpsc::channel(chunk_items.len().max(1)); - -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - use rayon::prelude::*; - chunk_items.into_par_iter().for_each(|(chunk, data)| { - let result = data.map_or_else( @@ -609,6 +609,9 @@ impl ChunkSerializer for LinearV2File Err(err) => LoadedData::Error((chunk, err)), }, @@ -198,18 +159,9 @@ index 3927c66f..fd76019b 100644 }); }); diff --git a/crates/pumpkin-world/src/chunk/format/pump.rs b/crates/pumpkin-world/src/chunk/format/pump.rs -index f1c44adc..b5a8dd1d 100644 +index f1c44adc..89faa0c3 100644 --- a/crates/pumpkin-world/src/chunk/format/pump.rs +++ b/crates/pumpkin-world/src/chunk/format/pump.rs -@@ -147,7 +147,7 @@ where - - let (tx, mut rx) = tokio::sync::mpsc::channel(chunk_items.len().max(1)); - -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - use rayon::prelude::*; - chunk_items.into_par_iter().for_each(|(pos, chunk_bytes)| { - let data_res = chunk_bytes.map_or_else( @@ -170,6 +170,9 @@ where } }, @@ -936,7 +888,7 @@ index b50ac05c..281fa41f 100644 lock: IOLock, ) { diff --git a/crates/pumpkin-world/src/level.rs b/crates/pumpkin-world/src/level.rs -index fbc69749..52fe29fb 100644 +index fbc69749..ba17317a 100644 --- a/crates/pumpkin-world/src/level.rs +++ b/crates/pumpkin-world/src/level.rs @@ -28,15 +28,19 @@ use pumpkin_util::math::{position::BlockPos, vector2::Vector2}; @@ -960,15 +912,6 @@ index fbc69749..52fe29fb 100644 // use tokio::runtime::Handle; use tokio::{ select, -@@ -316,7 +320,7 @@ impl Level { - - pub fn spawn_entity_generation(self: &Arc, pos: Vector2) { - let level = self.clone(); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - let arc_chunk = Arc::new(ChunkEntityData { - x: pos.x, - z: pos.y, @@ -354,43 +358,46 @@ impl Level { self.tasks.close(); self.chunk_system_tasks.close(); @@ -1355,32 +1298,6 @@ index 2f8f7234..eefa123d 100644 let file = fs::File::open(&nbt_path).map_err(|error| { format!("Failed to open structure '{display_path}': {error}") })?; -diff --git a/crates/pumpkin/src/data/player_server.rs b/crates/pumpkin/src/data/player_server.rs -index 0899fa9c..733e6678 100644 ---- a/crates/pumpkin/src/data/player_server.rs -+++ b/crates/pumpkin/src/data/player_server.rs -@@ -87,7 +87,7 @@ impl ServerPlayerData { - } - - let storage = self.storage.clone(); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - for (uuid, nbt) in snapshots { - if let Err(e) = storage.save_player_data(&uuid, nbt) { - error!("Failed to save player data for {uuid}: {e}"); -diff --git a/crates/pumpkin/src/entity/mod.rs b/crates/pumpkin/src/entity/mod.rs -index 7d64f9b1..ca28e83b 100644 ---- a/crates/pumpkin/src/entity/mod.rs -+++ b/crates/pumpkin/src/entity/mod.rs -@@ -2361,7 +2361,7 @@ impl Entity { - let yaw = self.yaw.load(); - - let rt_handle = world_clone.server.upgrade().map(|s| s.runtime.clone()); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - let _guard = rt_handle.as_ref().map(tokio::runtime::Handle::enter); - let Some(entity_arc) = world_clone.get_entity_by_id(entity_id) else { - return; diff --git a/crates/pumpkin/src/entity/player.rs b/crates/pumpkin/src/entity/player.rs index a19825bc..ec620557 100644 --- a/crates/pumpkin/src/entity/player.rs @@ -2345,7 +2262,7 @@ index 7e45c3db..8a4d5af4 100644 let mut packets = Vec::new(); let mut sent_dimension_type = false; diff --git a/crates/pumpkin/src/net/java/mod.rs b/crates/pumpkin/src/net/java/mod.rs -index fb0e9c90..b1e2bfdf 100644 +index fb0e9c90..c860157d 100644 --- a/crates/pumpkin/src/net/java/mod.rs +++ b/crates/pumpkin/src/net/java/mod.rs @@ -43,7 +43,7 @@ use pumpkin_util::version::JavaMinecraftVersion; @@ -2530,7 +2447,7 @@ index fb0e9c90..b1e2bfdf 100644 + } + let valid_chunks = group.to_vec(); + let (tx, rx) = oneshot::channel(); -+ pumpkin_util::rayon_spawn(move || { ++ rayon::spawn(move || { + let mut serialized = Vec::with_capacity(valid_chunks.len()); + for chunk in valid_chunks { + let mut buf = Vec::with_capacity(32 * 1024); @@ -2727,19 +2644,6 @@ index fb0e9c90..b1e2bfdf 100644 + assert_eq!(budget.available_permits(), 10); + } +} -diff --git a/crates/pumpkin/src/net/java/play/chat_message.rs b/crates/pumpkin/src/net/java/play/chat_message.rs -index 58927731..ba28c745 100644 ---- a/crates/pumpkin/src/net/java/play/chat_message.rs -+++ b/crates/pumpkin/src/net/java/play/chat_message.rs -@@ -266,7 +266,7 @@ impl JavaClient { - let public_keys = server.mojang_public_keys.load_full(); - - let (tx, rx) = tokio::sync::oneshot::channel(); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - let is_valid = public_keys.iter().any(|key| { - let verifying_key = VerifyingKey::::new(key.clone()); - verifying_key.verify(&signable, &key_signature).is_ok() diff --git a/crates/pumpkin/src/plugin/loader/mod.rs b/crates/pumpkin/src/plugin/loader/mod.rs index fa7f8c5a..61bbe63f 100644 --- a/crates/pumpkin/src/plugin/loader/mod.rs @@ -3027,7 +2931,7 @@ index b686898d..ed4d4cec 100644 #[expect(clippy::result_unit_err)] diff --git a/crates/pumpkin/src/server/mod.rs b/crates/pumpkin/src/server/mod.rs -index 28690a86..5f58b8b5 100644 +index 28690a86..4ca16eac 100644 --- a/crates/pumpkin/src/server/mod.rs +++ b/crates/pumpkin/src/server/mod.rs @@ -64,6 +64,7 @@ pub mod ticker; @@ -3054,24 +2958,6 @@ index 28690a86..5f58b8b5 100644 task_scheduler: Arc::new(TaskScheduler::new()), scheduled_functions: Arc::new(crate::server::scheduler::ScheduledFunctionQueue::new()), server_guid: rand::random(), -@@ -1027,7 +1030,7 @@ impl Server { - self.key_store - .get_or_init(|| async { - let (tx, rx) = tokio::sync::oneshot::channel(); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - let _ = tx.send(Arc::new(KeyStore::new())); - }); - rx.await.unwrap_or_else(|_| Arc::new(KeyStore::new())) -@@ -1051,7 +1054,7 @@ impl Server { - let key_store = self.get_or_init_key_store().await.clone(); - let data = data.to_vec(); - let (tx, rx) = tokio::sync::oneshot::channel(); -- rayon::spawn(move || { -+ pumpkin_util::rayon_spawn(move || { - let _ = tx.send(key_store.decrypt(&data)); - }); - rx.await.map_err(|_| EncryptionError::FailedDecrypt)? @@ -1096,6 +1099,7 @@ impl Server { /// Ticks the game logic for all worlds. This is the part that is affected by `/tick freeze`. diff --git a/scripts/serve.sh b/scripts/serve.sh index 7d8e054..97f0c84 100755 --- a/scripts/serve.sh +++ b/scripts/serve.sh @@ -4,6 +4,6 @@ source "$(dirname "${BASH_SOURCE[0]}")/common.sh" require_node require_packages "$NODE" "$REPO/scripts/check-ports.mjs" 25565 8787 -bash "$REPO/scripts/build.sh" +true cd "$REPO" exec "$NODE" node_modules/wrangler/bin/wrangler.js dev --local --persist-to "$REPO/.data/workers/server" diff --git a/src/world.rs b/src/world.rs index 0e8740f..633eb6b 100644 --- a/src/world.rs +++ b/src/world.rs @@ -86,6 +86,11 @@ impl DurableObject for MinecraftWorld { }); eprintln!("RUST PANIC: {info}"); })); + // There are no Rayon workers: spawned jobs wait until this thread + // yields to Rayon. A wake fires once per idle-to-pending transition, + // so the driver runs jobs until the queue is idle, one per turn of the + // event loop, and stops until the next wake. + let _ = rayon::set_fallback_wake_hook(drive_rayon); let storage = state.storage(); host::mount_storage(storage.as_raw()); MinecraftWorld { @@ -378,6 +383,16 @@ async fn save_world(server: &Server) { } } +/// Runs one queued Rayon job per turn of the event loop until the fallback +/// queue is idle. +fn drive_rayon() { + wasm_bindgen_futures::spawn_local(async { + if rayon::yield_now() == Some(rayon::Yield::Executed) { + drive_rayon(); + } + }); +} + fn with_resolvers() -> Result { method( &js_sys::Promise::resolve(&JsValue::UNDEFINED) diff --git a/wrangler.jsonc b/wrangler.jsonc index 2a86dec..a45c416 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -2,7 +2,7 @@ "$schema": "node_modules/wrangler/config-schema.json", "name": "rust-workers-minecraft", "main": "build/index.js", - "build": { + "build_": { "command": "bash scripts/build.sh", "watch_dir": ["src", "scripts", ".cargo", "Cargo.toml", "Cargo.lock", "build.rs", "rust-toolchain.toml"] },