Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 4 additions & 7 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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" }

Expand Down
1 change: 1 addition & 0 deletions docs/dependencies.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
150 changes: 18 additions & 132 deletions patches/pumpkin-emscripten.patch
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -50,57 +51,35 @@ 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
+tokio = { workspace = true, features = ["rt"] }

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<i32> for Hand {
@@ -310,3 +310,23 @@ impl TryFrom<i32> 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: FnOnce() + Send + 'static>(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
Expand Down Expand Up @@ -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<S: SingleChunkDataSerializer + 'static> 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<S: SingleChunkDataSerializer + 'static> ChunkSerializer for AnvilChunkFile<
Err(err) => LoadedData::Error((chunk, err)),
},
Expand All @@ -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<S: SingleChunkDataSerializer + 'static> ChunkSerializer for LinearV2File<S>

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<S: SingleChunkDataSerializer + 'static> ChunkSerializer for LinearV2File<S>
Err(err) => LoadedData::Error((chunk, err)),
},
Expand All @@ -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
}
},
Expand Down Expand Up @@ -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};
Expand All @@ -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<Self>, pos: Vector2<i32>) {
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();
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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::<Sha1>::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
Expand Down Expand Up @@ -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;
Expand All @@ -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`.
Expand Down
2 changes: 1 addition & 1 deletion scripts/serve.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
15 changes: 15 additions & 0 deletions src/world.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<JsValue, JsValue> {
method(
&js_sys::Promise::resolve(&JsValue::UNDEFINED)
Expand Down
2 changes: 1 addition & 1 deletion wrangler.jsonc
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
},
Expand Down
Loading