From 2175e9b31cd4b875631fa9d53d889a55772f3e9a Mon Sep 17 00:00:00 2001 From: Jeroen Schweitzer Date: Sat, 11 Apr 2026 00:08:41 +0200 Subject: [PATCH] feat(simulation): multi-threaded executor + EntityRng (#843 Part B) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Enable Bevy multi-threaded executor via bevy_tasks multi_threaded feature. Systems within the same TickPhase that don't share mutable resources now run in parallel automatically. Add EntityRng component — per-entity ChaCha20Rng seeded from world_seed + StableId via splitmix64 mixing. More deterministic than shared SimRng (order-independent). Migrate all monologue systems (4 of 13 SimRng consumers) to EntityRng, removing contention that serialized them against conversation/dialogue systems. Add rayon dependency (infrastructure only, no par_iter calls yet). SimRng retained for world-level randomness: conversation pairing, knowledge transfer, dialogue, ticker, storyteller. Co-Authored-By: Claude Opus 4.6 (1M context) --- server/Cargo.lock | 49 ++++++++++++++++++++++++++++++ server/Cargo.toml | 2 ++ server/src/main.rs | 18 +++++++++-- server/src/simulation/monologue.rs | 34 ++++++++++++--------- server/src/simulation/rng.rs | 44 +++++++++++++++++++++++++-- 5 files changed, 128 insertions(+), 19 deletions(-) diff --git a/server/Cargo.lock b/server/Cargo.lock index d8913abde..d15f9473c 100644 --- a/server/Cargo.lock +++ b/server/Cargo.lock @@ -306,6 +306,7 @@ dependencies = [ "async-task", "atomic-waker", "bevy_platform", + "concurrent-queue", "crossbeam-queue", "derive_more", "futures-lite", @@ -485,6 +486,25 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "crossbeam-deque" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9dd111b7b7f7d55b72c0a6ae361660ee5853c9af73f70c3c2ef6858b950e2e51" +dependencies = [ + "crossbeam-epoch", + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-epoch" +version = "0.9.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b82ac4a3c2ca9c3460964f020e1402edd5753411d7737aa39c3714ad1b5420e" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-queue" version = "0.3.12" @@ -581,6 +601,12 @@ dependencies = [ "serde", ] +[[package]] +name = "either" +version = "1.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "48c757948c5ede0e46177b7add2e67155f70e33c07fea8284df6576da70b3719" + [[package]] name = "equivalent" version = "1.0.2" @@ -605,6 +631,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab" dependencies = [ "concurrent-queue", + "parking", "pin-project-lite", ] @@ -1094,6 +1121,26 @@ dependencies = [ "getrandom", ] +[[package]] +name = "rayon" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "368f01d005bf8fd9b1206fb6fa653e6c4a81ceb1466406b81792d87c5677a58f" +dependencies = [ + "either", + "rayon-core", +] + +[[package]] +name = "rayon-core" +version = "1.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22e18b0f0062d30d4230b2e85ff77fdfe4326feb054b9783a3460d8435c8ab91" +dependencies = [ + "crossbeam-deque", + "crossbeam-utils", +] + [[package]] name = "regex-automata" version = "0.4.14" @@ -1260,6 +1307,7 @@ version = "0.1.34" dependencies = [ "bevy_app", "bevy_ecs", + "bevy_tasks", "bincode", "clap", "crossbeam-channel", @@ -1267,6 +1315,7 @@ dependencies = [ "pathfinding", "rand", "rand_chacha", + "rayon", "rmp-serde", "ron", "rusqlite", diff --git a/server/Cargo.toml b/server/Cargo.toml index 1bc0f1046..617e3ca6e 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -6,6 +6,8 @@ edition = "2021" [dependencies] bevy_ecs = "0.18" bevy_app = "0.18" +bevy_tasks = { version = "0.18", features = ["multi_threaded"] } +rayon = "1" serde = { version = "1", features = ["derive"] } serde_yaml = "0.9" ron = "0.8" diff --git a/server/src/main.rs b/server/src/main.rs index b13edfb53..b34113bd3 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -199,7 +199,7 @@ fn main() { std::process::exit(1); } } else { - setup_proof_room(&mut app, archetype); + setup_proof_room(&mut app, archetype, seed); } tracing::info!( @@ -395,9 +395,11 @@ fn dump_schedule_graph() { fn setup_proof_room( app: &mut App, archetype: settled_reach_server::bridge::types::CharacterArchetype, + world_seed: u64, ) { use settled_reach_server::knowledge::registry::EntityRegistry; use settled_reach_server::knowledge::KnowledgeGraph; + use settled_reach_server::simulation::rng::EntityRng; use settled_reach_server::npc::relationships::{RelationshipEdge, RelationshipGraph}; use settled_reach_server::npc::{ Contentment, DailyRoutine, Npc, RelationshipKind, RoutineEntry, ToleranceThreshold, Want, @@ -459,7 +461,10 @@ fn setup_proof_room( PlayerMoveCooldown::default(), )) .id(); - registry.register(player); + let player_sid = registry.register(player); + app.world_mut() + .entity_mut(player) + .insert(EntityRng::from_seed_and_id(world_seed, player_sid.0)); // NPC 1: Dock worker at (16,13) — behind wall, full routine let npc1 = app @@ -507,6 +512,9 @@ fn setup_proof_room( )) .id(); let npc1_sid = registry.register(npc1); + app.world_mut() + .entity_mut(npc1) + .insert(EntityRng::from_seed_and_id(world_seed, npc1_sid.0)); // NPC 2: Field tech at (14,18) — visible to player, has routine let npc2 = app @@ -544,6 +552,9 @@ fn setup_proof_room( )) .id(); let npc2_sid = registry.register(npc2); + app.world_mut() + .entity_mut(npc2) + .insert(EntityRng::from_seed_and_id(world_seed, npc2_sid.0)); // NPC 3: Guard at (18,14) — stationary, no routine let npc3 = app @@ -565,6 +576,9 @@ fn setup_proof_room( )) .id(); let npc3_sid = registry.register(npc3); + app.world_mut() + .entity_mut(npc3) + .insert(EntityRng::from_seed_and_id(world_seed, npc3_sid.0)); // Populate RelationshipGraph with a few edges { diff --git a/server/src/simulation/monologue.rs b/server/src/simulation/monologue.rs index 09b628bc2..9c423fca2 100644 --- a/server/src/simulation/monologue.rs +++ b/server/src/simulation/monologue.rs @@ -18,7 +18,7 @@ use crate::knowledge::{ContradictionDetectedQueue, EntityRegistry}; use crate::perception::interpretation::ObservationTrigger; use crate::simulation::conversation::NpcName; use crate::simulation::movement::{PlayerCharacter, TilePosition}; -use crate::simulation::rng::SimRng; +use crate::simulation::rng::EntityRng; use crate::simulation::time::SimulationTime; use crate::storyteller::EngagementRecord; @@ -242,17 +242,17 @@ impl SprintAnomalyQueue { /// System ordering: after trigger_monologue, before compute_observer_snapshot. pub fn process_sprint_anomaly_monologue( time: Res, - mut rng: ResMut, mut query: Query< ( &mut SprintAnomalyQueue, &mut MonologueBuffer, &mut MonologueState, + &mut EntityRng, ), With, >, ) { - let Ok((mut queue, mut buffer, mut state)) = query.single_mut() else { + let Ok((mut queue, mut buffer, mut state, mut entity_rng)) = query.single_mut() else { return; }; @@ -262,7 +262,7 @@ pub fn process_sprint_anomaly_monologue( } if let Some(_entity_id) = queue.take_ready(time.tick) { - let index = rng.rng.random_range(0..ANOMALY_LINES.len()); + let index = entity_rng.rng.random_range(0..ANOMALY_LINES.len()); let (id, text) = ANOMALY_LINES[index]; buffer.event = Some(MonologueEvent { @@ -295,18 +295,19 @@ pub fn process_sprint_anomaly_monologue( /// System ordering: after trigger_monologue, before process_sprint_anomaly_monologue. pub fn trigger_recognition_monologue( time: Res, - mut rng: ResMut, mut query: Query< ( &mut crate::perception::cognitive_delay::CognitiveDelay, &mut MonologueBuffer, &mut MonologueState, + &mut EntityRng, ), With, >, anomaly_markers: Query<(), With>, ) { - let Ok((mut cognitive_delay, mut buffer, mut state)) = query.single_mut() else { + let Ok((mut cognitive_delay, mut buffer, mut state, mut entity_rng)) = query.single_mut() + else { return; }; @@ -339,7 +340,7 @@ pub fn trigger_recognition_monologue( return; }; - let i = rng.rng.random_range(0..RECOGNITION_LINES.len()); + let i = entity_rng.rng.random_range(0..RECOGNITION_LINES.len()); let (id, text) = ( RECOGNITION_LINES[i].0.to_string(), RECOGNITION_LINES[i].1.to_string(), @@ -420,7 +421,6 @@ fn sound_range_tiles(range: &crate::knowledge::types::SoundRange) -> u32 { #[allow(clippy::too_many_arguments)] pub fn trigger_event_monologue( time: Res, - mut rng: ResMut, observation_queue: Option>, sound_queue: Option>, mut post_conv_queue: ResMut, @@ -430,6 +430,7 @@ pub fn trigger_event_monologue( &mut MonologueState, &mut MonologueBuffer, Option<&crate::simulation::conversation::ConversationEventBuffer>, + &mut EntityRng, ), With, >, @@ -440,7 +441,9 @@ pub fn trigger_event_monologue( // Saved for NPC attribution (engagement tracking #570) and trigger detection. let post_conv_npcs: Vec = post_conv_queue.drain(); - let Ok((player_pos, mut state, mut buffer, conv_buffer_opt)) = query.single_mut() else { + let Ok((player_pos, mut state, mut buffer, conv_buffer_opt, mut entity_rng)) = + query.single_mut() + else { return; }; @@ -489,7 +492,7 @@ pub fn trigger_event_monologue( let Some(trigger) = trigger else { return }; - let (id, text) = select_hardcoded_fallback(trigger, &mut rng.rng); + let (id, text) = select_hardcoded_fallback(trigger, &mut entity_rng.rng); buffer.event = Some(MonologueEvent { id: id.clone(), @@ -584,7 +587,6 @@ fn has_hear_sound_event( /// Remove this stub when those systems' ordering constraints are refactored. pub fn trigger_monologue( _time: Res, - _rng: ResMut, _query: Query< (&TilePosition, &mut MonologueState, &mut MonologueBuffer), With, @@ -635,10 +637,12 @@ pub(crate) fn resolve_name( pub fn process_contradiction_monologue( time: Res, _registry: Res, - mut rng: ResMut, mut contradiction_queue: ResMut, _npc_names: Query<&NpcName>, - mut player_query: Query<(&mut MonologueBuffer, &mut MonologueState), With>, + mut player_query: Query< + (&mut MonologueBuffer, &mut MonologueState, &mut EntityRng), + With, + >, ) { // _registry and _npc_names are available for future triggers needing live name resolution // via resolve_name(). ContradictionDetected uses pre-resolved names from the event payload. @@ -647,7 +651,7 @@ pub fn process_contradiction_monologue( return; } - let Ok((mut buffer, mut state)) = player_query.single_mut() else { + let Ok((mut buffer, mut state, mut entity_rng)) = player_query.single_mut() else { contradiction_queue.drain(); return; }; @@ -664,7 +668,7 @@ pub fn process_contradiction_monologue( return; }; - let template_idx = rng.rng.random_range(0..CONTRADICTION_TEMPLATE_LINES.len()); + let template_idx = entity_rng.rng.random_range(0..CONTRADICTION_TEMPLATE_LINES.len()); let text = CONTRADICTION_TEMPLATE_LINES[template_idx] .replace("{source}", &event.source_display_name) .replace("{subject}", &event.subject_display_name); diff --git a/server/src/simulation/rng.rs b/server/src/simulation/rng.rs index c02650a1e..b2198b824 100644 --- a/server/src/simulation/rng.rs +++ b/server/src/simulation/rng.rs @@ -6,8 +6,11 @@ use bevy_ecs::prelude::*; use rand::SeedableRng; use rand_chacha::ChaCha20Rng; -/// Simulation RNG resource -/// ChaCha20 RNG with stored seed for deterministic replay +/// Simulation RNG resource — shared world-level randomness. +/// +/// Used for non-entity randomness: storyteller dice, world events, ticker rotation. +/// For entity-level randomness, use [`EntityRng`] instead — it enables parallelism +/// because each entity has its own independent, deterministic RNG stream. #[derive(Resource)] pub struct SimRng { pub rng: ChaCha20Rng, @@ -29,6 +32,43 @@ impl SimRng { } } +/// Per-entity RNG component — deterministic, parallel-safe (#843). +/// +/// Each entity gets its own ChaCha20 stream seeded from `world_seed + StableId`. +/// Systems that need randomness for a specific entity query `&mut EntityRng` +/// instead of `ResMut`. This unlocks parallelism: Bevy can run systems +/// on disjoint entity sets concurrently. +/// +/// Determinism: same `world_seed` + same `StableId` = same RNG sequence, +/// regardless of thread execution order. +#[derive(Component)] +pub struct EntityRng { + pub rng: ChaCha20Rng, +} + +impl EntityRng { + /// Create a new EntityRng from a world seed and entity stable ID. + /// + /// Uses splitmix64 mixing to combine the two seeds with good avalanche + /// properties — avoids correlated sequences for entities with similar IDs. + pub fn from_seed_and_id(world_seed: u64, stable_id: u64) -> Self { + let mixed = splitmix64(world_seed.wrapping_add(stable_id)); + Self { + rng: ChaCha20Rng::seed_from_u64(mixed), + } + } +} + +/// Splitmix64 mixing function — good avalanche properties for seed derivation. +/// Ensures that similar inputs (e.g., consecutive StableIds) produce +/// uncorrelated output seeds. +fn splitmix64(mut x: u64) -> u64 { + x = x.wrapping_add(0x9e3779b97f4a7c15); + x = (x ^ (x >> 30)).wrapping_mul(0xbf58476d1ce4e5b9); + x = (x ^ (x >> 27)).wrapping_mul(0x94d049bb133111eb); + x ^ (x >> 31) +} + #[cfg(test)] mod tests { use super::*;