Skip to main content

bevy_ecs/schedule/executor/
multi_threaded.rs

1use alloc::{boxed::Box, vec::Vec};
2use bevy_platform::cell::SyncUnsafeCell;
3use bevy_platform::sync::Arc;
4use bevy_tasks::{ComputeTaskPool, Scope, TaskPool, ThreadExecutor};
5use concurrent_queue::ConcurrentQueue;
6use core::{any::Any, panic::AssertUnwindSafe};
7use fixedbitset::FixedBitSet;
8use std::sync::{Mutex, MutexGuard};
9
10#[cfg(feature = "trace")]
11use tracing::{info_span, Span};
12
13use crate::{
14    error::{BevyError, ErrorContext, ErrorHandler, Result},
15    prelude::Resource,
16    schedule::{
17        is_apply_deferred, ConditionWithAccess, SystemExecutor, SystemSchedule, SystemWithAccess,
18    },
19    system::{BoxedSystem, RunSystemError, ScheduleSystem},
20    world::{unsafe_world_cell::UnsafeWorldCell, World},
21};
22#[cfg(feature = "hotpatching")]
23use crate::{prelude::DetectChanges, HotPatchChanges};
24
25use super::__rust_begin_short_backtrace;
26
27/// Borrowed data used by the [`MultiThreadedExecutor`].
28struct Environment<'env, 'sys> {
29    executor: &'env MultiThreadedExecutor,
30    systems: &'sys [SyncUnsafeCell<SystemWithAccess>],
31    conditions: SyncUnsafeCell<Conditions<'sys>>,
32    world_cell: UnsafeWorldCell<'env>,
33}
34
35struct Conditions<'a> {
36    system_conditions: &'a mut [Vec<ConditionWithAccess>],
37    set_conditions: &'a mut [Vec<ConditionWithAccess>],
38    sets_with_conditions_of_systems: &'a [FixedBitSet],
39    systems_in_sets_with_conditions: &'a [FixedBitSet],
40}
41
42impl<'env, 'sys> Environment<'env, 'sys> {
43    fn new(
44        executor: &'env MultiThreadedExecutor,
45        schedule: &'sys mut SystemSchedule,
46        world: &'env mut World,
47    ) -> Self {
48        Environment {
49            executor,
50            systems: SyncUnsafeCell::from_mut(schedule.systems.as_mut_slice()).as_slice_of_cells(),
51            conditions: SyncUnsafeCell::new(Conditions {
52                system_conditions: &mut schedule.system_conditions,
53                set_conditions: &mut schedule.set_conditions,
54                sets_with_conditions_of_systems: &schedule.sets_with_conditions_of_systems,
55                systems_in_sets_with_conditions: &schedule.systems_in_sets_with_conditions,
56            }),
57            world_cell: world.as_unsafe_world_cell(),
58        }
59    }
60}
61
62/// Per-system data used by the [`MultiThreadedExecutor`].
63// Copied here because it can't be read from the system when it's running.
64struct SystemTaskMetadata {
65    /// The set of systems whose `component_access_set()` conflicts with this one.
66    conflicting_systems: FixedBitSet,
67    /// The set of systems whose `component_access_set()` conflicts with this system's conditions.
68    /// Note that this is separate from `conflicting_systems` to handle the case where
69    /// a system is skipped by an earlier system set condition or system stepping,
70    /// and needs access to run its conditions but not for itself.
71    condition_conflicting_systems: FixedBitSet,
72    /// Indices of the systems that directly depend on the system.
73    dependents: Vec<usize>,
74    /// Is `true` if the system does not access `!Send` data.
75    is_send: bool,
76    /// Is `true` if the system is exclusive.
77    is_exclusive: bool,
78}
79
80/// The result of running a system that is sent across a channel.
81struct SystemResult {
82    system_index: usize,
83}
84
85/// Runs the schedule using a thread pool. Non-conflicting systems can run in parallel.
86pub struct MultiThreadedExecutor {
87    /// The running state, protected by a mutex so that a reference to the executor can be shared across tasks.
88    state: Mutex<ExecutorState>,
89    /// Queue of system completion events.
90    system_completion: ConcurrentQueue<SystemResult>,
91    /// Setting when true applies deferred system buffers after all systems have run
92    apply_final_deferred: bool,
93    /// When set, tells the executor that a thread has panicked.
94    panic_payload: Mutex<Option<Box<dyn Any + Send>>>,
95    starting_systems: FixedBitSet,
96    /// Cached tracing span
97    #[cfg(feature = "trace")]
98    executor_span: Span,
99}
100
101/// The state of the executor while running.
102pub struct ExecutorState {
103    /// Metadata for scheduling and running system tasks.
104    system_task_metadata: Vec<SystemTaskMetadata>,
105    /// The set of systems whose `component_access_set()` conflicts with this system set's conditions.
106    set_condition_conflicting_systems: Vec<FixedBitSet>,
107    /// Returns `true` if a system with non-`Send` access is running.
108    local_thread_running: bool,
109    /// Returns `true` if an exclusive system is running.
110    exclusive_running: bool,
111    /// The number of systems that are running.
112    num_running_systems: usize,
113    /// The number of dependencies each system has that have not been satisfied. A dependency
114    /// is satisfied when the predecessor completes.
115    num_dependencies_remaining: Vec<usize>,
116    /// System sets whose conditions have been evaluated.
117    evaluated_sets: FixedBitSet,
118    /// Systems that have no remaining dependencies and are waiting to run.
119    ready_systems: FixedBitSet,
120    /// copy of `ready_systems`
121    ready_systems_copy: FixedBitSet,
122    /// Systems that are running.
123    running_systems: FixedBitSet,
124    /// Systems that got skipped.
125    skipped_systems: FixedBitSet,
126    /// Systems whose conditions have been evaluated and were run or skipped.
127    completed_systems: FixedBitSet,
128    /// Systems that have run but have not had their buffers applied.
129    unapplied_systems: FixedBitSet,
130}
131
132/// References to data required by the executor.
133/// This is copied to each system task so that can invoke the executor when they complete.
134// These all need to outlive 'scope in order to be sent to new tasks,
135// and keeping them all in a struct means we can use lifetime elision.
136#[derive(Copy, Clone)]
137struct Context<'scope, 'env, 'sys> {
138    environment: &'env Environment<'env, 'sys>,
139    scope: &'scope Scope<'scope, 'env, ()>,
140    error_handler: ErrorHandler,
141}
142
143impl Default for MultiThreadedExecutor {
144    fn default() -> Self {
145        Self::new()
146    }
147}
148
149impl SystemExecutor for MultiThreadedExecutor {
150    fn init(&mut self, schedule: &SystemSchedule) {
151        let state = self.state.get_mut().unwrap();
152        // pre-allocate space
153        let sys_count = schedule.system_ids.len();
154        let set_count = schedule.set_ids.len();
155
156        self.system_completion = ConcurrentQueue::bounded(sys_count.max(1));
157        self.starting_systems = FixedBitSet::with_capacity(sys_count);
158        state.evaluated_sets = FixedBitSet::with_capacity(set_count);
159        state.ready_systems = FixedBitSet::with_capacity(sys_count);
160        state.ready_systems_copy = FixedBitSet::with_capacity(sys_count);
161        state.running_systems = FixedBitSet::with_capacity(sys_count);
162        state.completed_systems = FixedBitSet::with_capacity(sys_count);
163        state.skipped_systems = FixedBitSet::with_capacity(sys_count);
164        state.unapplied_systems = FixedBitSet::with_capacity(sys_count);
165
166        state.system_task_metadata = Vec::with_capacity(sys_count);
167        for index in 0..sys_count {
168            state.system_task_metadata.push(SystemTaskMetadata {
169                conflicting_systems: FixedBitSet::with_capacity(sys_count),
170                condition_conflicting_systems: FixedBitSet::with_capacity(sys_count),
171                dependents: schedule.system_dependents[index].clone(),
172                is_send: schedule.systems[index].system.is_send(),
173                is_exclusive: schedule.systems[index].access.is_exclusive(),
174            });
175            // A system with no dependencies is a starting system.
176            if schedule.system_dependencies[index] == 0 {
177                self.starting_systems.insert(index);
178            }
179        }
180
181        {
182            #[cfg(feature = "trace")]
183            let _span = info_span!("calculate conflicting systems").entered();
184            for index1 in 0..sys_count {
185                let system1 = &schedule.systems[index1];
186                for index2 in 0..index1 {
187                    let system2 = &schedule.systems[index2];
188                    if !system2.access.is_compatible(&system1.access) {
189                        state.system_task_metadata[index1]
190                            .conflicting_systems
191                            .insert(index2);
192                        state.system_task_metadata[index2]
193                            .conflicting_systems
194                            .insert(index1);
195                    }
196                }
197
198                for index2 in 0..sys_count {
199                    let system2 = &schedule.systems[index2];
200                    if schedule.system_conditions[index1]
201                        .iter()
202                        .any(|condition| !system2.access.is_compatible(&condition.access))
203                    {
204                        state.system_task_metadata[index1]
205                            .condition_conflicting_systems
206                            .insert(index2);
207                    }
208                }
209            }
210
211            state.set_condition_conflicting_systems.clear();
212            state.set_condition_conflicting_systems.reserve(set_count);
213            for set_idx in 0..set_count {
214                let mut conflicting_systems = FixedBitSet::with_capacity(sys_count);
215                for sys_index in 0..sys_count {
216                    let system = &schedule.systems[sys_index];
217                    if schedule.set_conditions[set_idx]
218                        .iter()
219                        .any(|condition| !system.access.is_compatible(&condition.access))
220                    {
221                        conflicting_systems.insert(sys_index);
222                    }
223                }
224                state
225                    .set_condition_conflicting_systems
226                    .push(conflicting_systems);
227            }
228        }
229
230        state.num_dependencies_remaining = Vec::with_capacity(sys_count);
231    }
232
233    fn run(
234        &mut self,
235        schedule: &mut SystemSchedule,
236        world: &mut World,
237        _skip_systems: Option<&FixedBitSet>,
238        error_handler: ErrorHandler,
239    ) {
240        let state = self.state.get_mut().unwrap();
241        // reset counts
242        if schedule.systems.is_empty() {
243            return;
244        }
245        state.num_running_systems = 0;
246        state
247            .num_dependencies_remaining
248            .clone_from(&schedule.system_dependencies);
249        state.ready_systems.clone_from(&self.starting_systems);
250
251        // If stepping is enabled, make sure we skip those systems that should
252        // not be run.
253        #[cfg(feature = "bevy_debug_stepping")]
254        if let Some(skipped_systems) = _skip_systems {
255            debug_assert_eq!(skipped_systems.len(), state.completed_systems.len());
256            // mark skipped systems as completed
257            state.completed_systems |= skipped_systems;
258
259            // signal the dependencies for each of the skipped systems, as
260            // though they had run
261            for system_index in skipped_systems.ones() {
262                state.signal_dependents(system_index);
263                state.ready_systems.remove(system_index);
264            }
265        }
266
267        let thread_executor = world
268            .get_resource::<MainThreadExecutor>()
269            .map(|e| e.0.clone());
270        let thread_executor = thread_executor.as_deref();
271
272        let environment = &Environment::new(self, schedule, world);
273
274        ComputeTaskPool::get_or_init(TaskPool::default).scope_with_executor(
275            false,
276            thread_executor,
277            |scope| {
278                let context = Context {
279                    environment,
280                    scope,
281                    error_handler,
282                };
283
284                // The first tick won't need to process finished systems, but we still need to run the loop in
285                // tick_executor() in case a system completes while the first tick still holds the mutex.
286                context.tick_executor();
287            },
288        );
289
290        // End the borrows of self and world in environment by copying out the reference to systems.
291        let systems = environment.systems;
292
293        let state = self.state.get_mut().unwrap();
294        if self.apply_final_deferred {
295            // Do one final apply buffers after all systems have completed
296            // Commands should be applied while on the scope's thread, not the executor's thread
297            let res = apply_deferred(&state.unapplied_systems, systems, world, error_handler);
298            if let Err(payload) = res {
299                let panic_payload = self.panic_payload.get_mut().unwrap();
300                *panic_payload = Some(payload);
301            }
302            state.unapplied_systems.clear();
303        }
304
305        // check to see if there was a panic
306        let payload = self.panic_payload.get_mut().unwrap();
307        if let Some(payload) = payload.take() {
308            std::panic::resume_unwind(payload);
309        }
310
311        debug_assert!(state.ready_systems.is_clear());
312        debug_assert!(state.running_systems.is_clear());
313        state.evaluated_sets.clear();
314        state.skipped_systems.clear();
315        state.completed_systems.clear();
316    }
317
318    fn set_apply_final_deferred(&mut self, value: bool) {
319        self.apply_final_deferred = value;
320    }
321}
322
323impl<'scope, 'env: 'scope, 'sys> Context<'scope, 'env, 'sys> {
324    fn system_completed(&self, system_index: usize, res: Result<(), Box<dyn Any + Send>>) {
325        // tell the executor that the system finished
326        self.environment
327            .executor
328            .system_completion
329            .push(SystemResult { system_index })
330            .unwrap_or_else(|error| unreachable!("{}", error));
331        if let Err(payload) = res {
332            // set the payload to propagate the error
333            let mut panic_payload = self.environment.executor.panic_payload.lock().unwrap();
334            *panic_payload = Some(payload);
335        }
336        self.tick_executor();
337    }
338
339    #[expect(
340        clippy::mut_from_ref,
341        reason = "Field is only accessed here and is guarded by lock with a documented safety comment"
342    )]
343    fn try_lock<'a>(&'a self) -> Option<(&'a mut Conditions<'sys>, MutexGuard<'a, ExecutorState>)> {
344        let guard = self.environment.executor.state.try_lock().ok()?;
345        // SAFETY: This is an exclusive access as no other location fetches conditions mutably, and
346        // is synchronized by the lock on the executor state.
347        let conditions = unsafe { &mut *self.environment.conditions.get() };
348        Some((conditions, guard))
349    }
350
351    fn tick_executor(&self) {
352        // Ensure that the executor handles any events pushed to the system_completion queue by this thread.
353        // If this thread acquires the lock, the executor runs after the push() and they are processed.
354        // If this thread does not acquire the lock, then the is_empty() check on the other thread runs
355        // after the lock is released, which is after try_lock() failed, which is after the push()
356        // on this thread, so the is_empty() check will see the new events and loop.
357        loop {
358            let Some((conditions, mut guard)) = self.try_lock() else {
359                return;
360            };
361            guard.tick(self, conditions);
362            // Make sure we drop the guard before checking system_completion.is_empty(), or we could lose events.
363            drop(guard);
364            if self.environment.executor.system_completion.is_empty() {
365                return;
366            }
367        }
368    }
369}
370
371impl MultiThreadedExecutor {
372    /// Creates a new `multi_threaded` executor for use with a [`Schedule`].
373    ///
374    /// [`Schedule`]: crate::schedule::Schedule
375    pub fn new() -> Self {
376        Self {
377            state: Mutex::new(ExecutorState::new()),
378            system_completion: ConcurrentQueue::unbounded(),
379            starting_systems: FixedBitSet::new(),
380            apply_final_deferred: true,
381            panic_payload: Mutex::new(None),
382            #[cfg(feature = "trace")]
383            executor_span: info_span!("multithreaded executor"),
384        }
385    }
386}
387
388impl ExecutorState {
389    fn new() -> Self {
390        Self {
391            system_task_metadata: Vec::new(),
392            set_condition_conflicting_systems: Vec::new(),
393            num_running_systems: 0,
394            num_dependencies_remaining: Vec::new(),
395            local_thread_running: false,
396            exclusive_running: false,
397            evaluated_sets: FixedBitSet::new(),
398            ready_systems: FixedBitSet::new(),
399            ready_systems_copy: FixedBitSet::new(),
400            running_systems: FixedBitSet::new(),
401            skipped_systems: FixedBitSet::new(),
402            completed_systems: FixedBitSet::new(),
403            unapplied_systems: FixedBitSet::new(),
404        }
405    }
406
407    fn tick(&mut self, context: &Context, conditions: &mut Conditions) {
408        #[cfg(feature = "trace")]
409        let _span = context.environment.executor.executor_span.enter();
410
411        for result in context.environment.executor.system_completion.try_iter() {
412            self.finish_system_and_handle_dependents(result);
413        }
414
415        // SAFETY:
416        // - `finish_system_and_handle_dependents` has updated the currently running systems.
417        // - `rebuild_active_access` locks access for all currently running systems.
418        unsafe {
419            self.spawn_system_tasks(context, conditions);
420        }
421    }
422
423    /// # Safety
424    /// - Caller must ensure that `self.ready_systems` does not contain any systems that
425    ///   have been mutably borrowed (such as the systems currently running).
426    /// - `world_cell` must have permission to access all world data (not counting
427    ///   any world data that is claimed by systems currently running on this executor).
428    unsafe fn spawn_system_tasks(&mut self, context: &Context, conditions: &mut Conditions) {
429        if self.exclusive_running {
430            return;
431        }
432
433        #[cfg(feature = "hotpatching")]
434        #[expect(
435            clippy::undocumented_unsafe_blocks,
436            reason = "This actually could result in UB if a system tries to mutate
437            `HotPatchChanges`. We allow this as the resource only exists with the `hotpatching` feature.
438            and `hotpatching` should never be enabled in release."
439        )]
440        #[cfg(feature = "hotpatching")]
441        let hotpatch_tick = unsafe {
442            context
443                .environment
444                .world_cell
445                .get_resource_ref::<HotPatchChanges>()
446        }
447        .map(|r| r.last_changed())
448        .unwrap_or_default();
449
450        // can't borrow since loop mutably borrows `self`
451        let mut ready_systems = core::mem::take(&mut self.ready_systems_copy);
452
453        // Skipping systems may cause their dependents to become ready immediately.
454        // If that happens, we need to run again immediately or we may fail to spawn those dependents.
455        let mut check_for_new_ready_systems = true;
456        while check_for_new_ready_systems {
457            check_for_new_ready_systems = false;
458
459            ready_systems.clone_from(&self.ready_systems);
460
461            for system_index in ready_systems.ones() {
462                debug_assert!(!self.running_systems.contains(system_index));
463                // SAFETY: Caller assured that these systems are not running.
464                // Therefore, no other reference to this system exists and there is no aliasing.
465                let system =
466                    &mut unsafe { &mut *context.environment.systems[system_index].get() }.system;
467
468                #[cfg(feature = "hotpatching")]
469                if hotpatch_tick.is_newer_than(
470                    system.get_last_run(),
471                    context.environment.world_cell.change_tick(),
472                ) {
473                    system.refresh_hotpatch();
474                }
475
476                if !self.can_run(system_index, conditions) {
477                    // NOTE: exclusive systems with ambiguities are susceptible to
478                    // being significantly displaced here (compared to single-threaded order)
479                    // if systems after them in topological order can run
480                    // if that becomes an issue, `break;` if exclusive system
481                    continue;
482                }
483
484                self.ready_systems.remove(system_index);
485
486                // SAFETY: `can_run` returned true, which means that:
487                // - There can be no systems running whose accesses would conflict with any conditions.
488                if unsafe {
489                    !self.should_run(
490                        system_index,
491                        system,
492                        conditions,
493                        context.environment.world_cell,
494                        context.error_handler,
495                    )
496                } {
497                    self.skip_system_and_signal_dependents(system_index);
498                    // signal_dependents may have set more systems to ready.
499                    check_for_new_ready_systems = true;
500                    continue;
501                }
502
503                self.running_systems.insert(system_index);
504                self.num_running_systems += 1;
505
506                if self.system_task_metadata[system_index].is_exclusive {
507                    // SAFETY: `can_run` returned true for this system,
508                    // which means no systems are currently borrowed.
509                    unsafe {
510                        self.spawn_exclusive_system_task(context, system_index);
511                    }
512                    check_for_new_ready_systems = false;
513                    break;
514                }
515
516                // SAFETY:
517                // - Caller ensured no other reference to this system exists.
518                // - `system_task_metadata[system_index].is_exclusive` is `false`,
519                //   so `System::is_exclusive` returned `false` when we called it.
520                // - `can_run` returned true, so no systems with conflicting world access are running.
521                unsafe {
522                    self.spawn_system_task(context, system_index);
523                }
524            }
525        }
526
527        // give back
528        self.ready_systems_copy = ready_systems;
529    }
530
531    fn can_run(&mut self, system_index: usize, conditions: &mut Conditions) -> bool {
532        let system_meta = &self.system_task_metadata[system_index];
533        if system_meta.is_exclusive && self.num_running_systems > 0 {
534            return false;
535        }
536
537        if !system_meta.is_send && self.local_thread_running {
538            return false;
539        }
540
541        // TODO: an earlier out if world's archetypes did not change
542        for set_idx in conditions.sets_with_conditions_of_systems[system_index]
543            .difference(&self.evaluated_sets)
544        {
545            if !self.set_condition_conflicting_systems[set_idx].is_disjoint(&self.running_systems) {
546                return false;
547            }
548        }
549
550        if !system_meta
551            .condition_conflicting_systems
552            .is_disjoint(&self.running_systems)
553        {
554            return false;
555        }
556
557        if !self.skipped_systems.contains(system_index)
558            && !system_meta
559                .conflicting_systems
560                .is_disjoint(&self.running_systems)
561        {
562            return false;
563        }
564
565        true
566    }
567
568    /// # Safety
569    /// * `world` must have permission to read any world data required by
570    ///   the system's conditions: this includes conditions for the system
571    ///   itself, and conditions for any of the system's sets.
572    unsafe fn should_run(
573        &mut self,
574        system_index: usize,
575        system: &mut ScheduleSystem,
576        conditions: &mut Conditions,
577        world: UnsafeWorldCell,
578        error_handler: ErrorHandler,
579    ) -> bool {
580        let mut should_run = !self.skipped_systems.contains(system_index);
581
582        for set_idx in conditions.sets_with_conditions_of_systems[system_index].ones() {
583            if self.evaluated_sets.contains(set_idx) {
584                continue;
585            }
586
587            // Evaluate the system set's conditions.
588            // SAFETY:
589            // - The caller ensures that `world` has permission to read any data
590            //   required by the conditions.
591            let set_conditions_met = unsafe {
592                evaluate_and_fold_conditions(
593                    &mut conditions.set_conditions[set_idx],
594                    world,
595                    error_handler,
596                    system,
597                    true,
598                )
599            };
600
601            if !set_conditions_met {
602                self.skipped_systems
603                    .union_with(&conditions.systems_in_sets_with_conditions[set_idx]);
604            }
605
606            should_run &= set_conditions_met;
607            self.evaluated_sets.insert(set_idx);
608        }
609
610        // Evaluate the system's conditions.
611        // SAFETY:
612        // - The caller ensures that `world` has permission to read any data
613        //   required by the conditions.
614        let system_conditions_met = unsafe {
615            evaluate_and_fold_conditions(
616                &mut conditions.system_conditions[system_index],
617                world,
618                error_handler,
619                system,
620                false,
621            )
622        };
623
624        if !system_conditions_met {
625            self.skipped_systems.insert(system_index);
626        }
627
628        should_run &= system_conditions_met;
629
630        should_run
631    }
632
633    /// # Safety
634    /// - Caller must not alias systems that are running.
635    /// - `is_exclusive` must have returned `false` for the specified system.
636    /// - `world` must have permission to access the world data
637    ///   used by the specified system.
638    unsafe fn spawn_system_task(&mut self, context: &Context, system_index: usize) {
639        let system = &context.environment.systems[system_index];
640        // Move the full context object into the new future.
641        let context = *context;
642
643        let system_meta = &self.system_task_metadata[system_index];
644
645        let task = async move {
646            let res = handle_errors(
647                |system| {
648                    // SAFETY:
649                    // - The caller ensures that we have permission to
650                    // access the world data used by the system.
651                    // - `is_exclusive` returned false
652                    unsafe {
653                        __rust_begin_short_backtrace::run_unsafe(
654                            system,
655                            context.environment.world_cell,
656                        )
657                    }
658                },
659                // SAFETY: this system is not running, no other reference exists
660                unsafe { &mut (*system.get()).system },
661                context.error_handler,
662                "System panicked",
663            );
664            context.system_completed(system_index, res);
665        };
666
667        if system_meta.is_send {
668            context.scope.spawn(task);
669        } else {
670            self.local_thread_running = true;
671            context.scope.spawn_on_external(task);
672        }
673    }
674
675    /// # Safety
676    /// Caller must ensure no systems are currently borrowed.
677    unsafe fn spawn_exclusive_system_task(&mut self, context: &Context, system_index: usize) {
678        let system = &context.environment.systems[system_index];
679        // Move the full context object into the new future.
680        let context = *context;
681
682        // SAFETY: this system is not running, no other reference exists
683        if is_apply_deferred(unsafe { &*(*system.get()).system }) {
684            // TODO: avoid allocation
685            let unapplied_systems = self.unapplied_systems.clone();
686            self.unapplied_systems.clear();
687            let task = async move {
688                // SAFETY: `can_run` returned true for this system, which means
689                // that no other systems currently have access to the world.
690                let world = unsafe { context.environment.world_cell.world_mut() };
691                let res = apply_deferred(
692                    &unapplied_systems,
693                    context.environment.systems,
694                    world,
695                    context.error_handler,
696                );
697                context.system_completed(system_index, res);
698            };
699
700            context.scope.spawn_on_scope(task);
701        } else {
702            let task = async move {
703                // SAFETY: `can_run` returned true for this system, which means
704                // that no other systems currently have access to the world.
705                let world = unsafe { context.environment.world_cell.world_mut() };
706                let res = handle_errors(
707                    |system| __rust_begin_short_backtrace::run(system, world),
708                    // SAFETY: this system is not running, no other reference exists
709                    unsafe { &mut (*system.get()).system },
710                    context.error_handler,
711                    "Exclusive system panicked",
712                );
713                context.system_completed(system_index, res);
714            };
715
716            context.scope.spawn_on_scope(task);
717        }
718
719        self.exclusive_running = true;
720        self.local_thread_running = true;
721    }
722
723    fn finish_system_and_handle_dependents(&mut self, result: SystemResult) {
724        let SystemResult { system_index, .. } = result;
725
726        if self.system_task_metadata[system_index].is_exclusive {
727            self.exclusive_running = false;
728        }
729
730        if self.system_task_metadata[system_index].is_exclusive
731            || !self.system_task_metadata[system_index].is_send
732        {
733            self.local_thread_running = false;
734        }
735
736        debug_assert!(self.num_running_systems >= 1);
737        self.num_running_systems -= 1;
738        self.running_systems.remove(system_index);
739        self.completed_systems.insert(system_index);
740        self.unapplied_systems.insert(system_index);
741
742        self.signal_dependents(system_index);
743    }
744
745    fn skip_system_and_signal_dependents(&mut self, system_index: usize) {
746        self.completed_systems.insert(system_index);
747        self.signal_dependents(system_index);
748    }
749
750    /// Called when `system_index` completes, satisfying one dependency for each of its
751    /// dependents and marking any that become ready to run.
752    fn signal_dependents(&mut self, system_index: usize) {
753        for &dep_idx in &self.system_task_metadata[system_index].dependents {
754            let remaining = &mut self.num_dependencies_remaining[dep_idx];
755            debug_assert!(*remaining >= 1);
756            *remaining -= 1;
757            if *remaining == 0 && !self.completed_systems.contains(dep_idx) {
758                self.ready_systems.insert(dep_idx);
759            }
760        }
761    }
762}
763
764fn apply_deferred(
765    unapplied_systems: &FixedBitSet,
766    systems: &[SyncUnsafeCell<SystemWithAccess>],
767    world: &mut World,
768    error_handler: ErrorHandler,
769) -> Result<(), Box<dyn Any + Send>> {
770    for system_index in unapplied_systems.ones() {
771        // SAFETY: none of these systems are running, no other references exist
772        let system = &mut unsafe { &mut *systems[system_index].get() }.system;
773        handle_errors(
774            |system| {
775                system.apply_deferred(world);
776                Ok(())
777            },
778            system,
779            error_handler,
780            "Encountered a panic while applying system buffers",
781        )?;
782    }
783    Ok(())
784}
785
786/// # Safety
787/// - `world` must have permission to read any world data
788///   required by `conditions`.
789unsafe fn evaluate_and_fold_conditions(
790    conditions: &mut [ConditionWithAccess],
791    world: UnsafeWorldCell,
792    error_handler: ErrorHandler,
793    for_system: &ScheduleSystem,
794    on_set: bool,
795) -> bool {
796    #[expect(
797        clippy::unnecessary_fold,
798        reason = "Short-circuiting here would prevent conditions from mutating their own state as needed."
799    )]
800    conditions
801        .iter_mut()
802        .map(|ConditionWithAccess { condition, .. }| {
803            let potential_unwind = std::panic::catch_unwind(AssertUnwindSafe(||
804                // SAFETY:
805                // - The caller ensures that `world` has permission to read any data
806                //   required by the condition.
807                unsafe {__rust_begin_short_backtrace::readonly_run_unsafe(&mut **condition, world)}
808            ));
809            match potential_unwind {
810                // Let the error handler handle the panic
811                Err(payload) => {
812                    __rust_begin_short_backtrace::error_handler(
813                        error_handler,
814                        BevyError::panic(
815                            "Encountered panic",
816                            payload,
817                        ),
818                        ErrorContext::RunCondition {
819                            name: condition.name(),
820                            last_run: condition.get_last_run(),
821                            system: for_system.name(),
822                            on_set,
823                        },
824                    ); false},
825                // Condition returned an error, let the error handler handle it
826                Ok(Err(RunSystemError::Failed(err))) => {
827                    __rust_begin_short_backtrace::error_handler(
828                        error_handler,
829                        err,
830                        ErrorContext::RunCondition {
831                            name: condition.name(),
832                            last_run: condition.get_last_run(),
833                            system: for_system.name(),
834                            on_set,
835                        },
836                    ); false
837                },
838                Ok(Err(RunSystemError::Skipped(_))) => false,
839                Ok(Ok(result)) => result,
840            }
841        })
842        .fold(true, |acc, res| acc && res)
843}
844
845/// Handle a potential panic or failed system by invoking the error handler
846/// and/or returning a panic payload with which to resume unwinding.
847fn handle_errors(
848    f: impl FnOnce(&mut BoxedSystem) -> Result<(), RunSystemError>,
849    system: &mut BoxedSystem,
850    error_handler: ErrorHandler,
851    error_message: &str,
852) -> Result<(), Box<dyn Any + Send>> {
853    let potential_unwind = std::panic::catch_unwind(AssertUnwindSafe(|| f(system)));
854    match potential_unwind {
855        // Let the error handler handle the panic, passing on any panic it throws
856        Err(payload) => std::panic::catch_unwind(AssertUnwindSafe(|| {
857            __rust_begin_short_backtrace::error_handler(
858                error_handler,
859                BevyError::panic(error_message, payload),
860                ErrorContext::System {
861                    name: system.name(),
862                    last_run: system.get_last_run(),
863                },
864            );
865        })),
866        // System returned an error, let the error handler handle it, passing on any panic it throws
867        Ok(Err(RunSystemError::Failed(err))) => std::panic::catch_unwind(AssertUnwindSafe(|| {
868            __rust_begin_short_backtrace::error_handler(
869                error_handler,
870                err,
871                ErrorContext::System {
872                    name: system.name(),
873                    last_run: system.get_last_run(),
874                },
875            );
876        })),
877        // Success (or skipped system)
878        _ => Ok(()),
879    }
880}
881
882/// New-typed [`ThreadExecutor`] [`Resource`] that is used to run systems on the main thread
883#[derive(Resource, Clone)]
884pub struct MainThreadExecutor(pub Arc<ThreadExecutor<'static>>);
885
886impl Default for MainThreadExecutor {
887    fn default() -> Self {
888        Self::new()
889    }
890}
891
892impl MainThreadExecutor {
893    /// Creates a new executor that can be used to run systems on the main thread.
894    pub fn new() -> Self {
895        MainThreadExecutor(TaskPool::get_thread_executor())
896    }
897}
898
899#[cfg(test)]
900mod tests {
901    use alloc::string::String;
902    use core::{
903        panic::AssertUnwindSafe,
904        sync::atomic::{AtomicBool, Ordering::Relaxed},
905    };
906    use std::panic::catch_unwind;
907
908    use crate::{
909        change_detection::Tick,
910        error::{BevyError, ErrorContext, FallbackErrorHandler},
911        prelude::Resource,
912        schedule::{IntoScheduleConfigs, MultiThreadedExecutor, Schedule},
913        system::{
914            Commands, NonSendMut, SystemAccess, SystemMeta, SystemParam, SystemParamValidationError,
915        },
916        world::{unsafe_world_cell::UnsafeWorldCell, World},
917    };
918
919    #[derive(Resource)]
920    struct R;
921
922    struct ExclusiveMarker;
923
924    // SAFETY: No world data is accessed.
925    unsafe impl SystemParam for ExclusiveMarker {
926        type State = ();
927        type Item<'world, 'state> = ExclusiveMarker;
928
929        fn init_state(_world: &mut World) -> Self::State {}
930
931        fn init_access(
932            _state: &Self::State,
933            system_meta: &mut SystemMeta,
934            system_access: &mut SystemAccess,
935            _world: &mut World,
936        ) {
937            system_access.require_exclusive_access::<Self>(system_meta);
938        }
939
940        unsafe fn get_param<'world, 'state>(
941            _state: &'state mut Self::State,
942            _system_meta: &SystemMeta,
943            _world: UnsafeWorldCell<'world>,
944            _change_tick: Tick,
945        ) -> Result<Self::Item<'world, 'state>, SystemParamValidationError> {
946            Ok(ExclusiveMarker)
947        }
948    }
949
950    #[test]
951    fn skipped_systems_notify_dependents() {
952        let mut world = World::new();
953        let mut schedule = Schedule::default();
954        schedule.set_executor(MultiThreadedExecutor::new());
955        schedule.add_systems(
956            (
957                (|| {}).run_if(|| false),
958                // This system depends on a system that is always skipped.
959                |mut commands: Commands| {
960                    commands.insert_resource(R);
961                },
962            )
963                .chain(),
964        );
965        schedule.run(&mut world);
966        assert!(world.get_resource::<R>().is_some());
967    }
968
969    /// Regression test for case where exclusive system left local thread state as not
970    /// cleared and prevented subsequent non-send system runs
971    #[test]
972    fn exclusive_system_reporting_send_releases_the_local_thread() {
973        #[derive(Default)]
974        struct NonSendMarker(bool);
975
976        let mut world = World::new();
977        world.insert_non_send(NonSendMarker::default());
978
979        let mut schedule = Schedule::default();
980        schedule.set_executor(MultiThreadedExecutor::new());
981        schedule.add_systems(
982            (
983                |_: ExclusiveMarker| {},
984                |mut marker: NonSendMut<NonSendMarker>| {
985                    marker.0 = true;
986                },
987            )
988                .chain(),
989        );
990
991        schedule.run(&mut world);
992        assert!(world.non_send::<NonSendMarker>().0);
993    }
994
995    /// Regression test for a weird bug flagged by MIRI in
996    /// `spawn_exclusive_system_task`, related to a `&mut World` being captured
997    /// inside an `async` block and somehow remaining alive even after its last use.
998    #[test]
999    fn check_spawn_exclusive_system_task_miri() {
1000        let mut world = World::new();
1001        let mut schedule = Schedule::default();
1002        schedule.set_executor(MultiThreadedExecutor::new());
1003        schedule.add_systems(((|_: Commands| {}), |_: Commands| {}).chain());
1004        schedule.run(&mut world);
1005    }
1006
1007    #[test]
1008    fn panic_to_error() {
1009        let mut world = World::new();
1010
1011        let mut schedule_error = Schedule::default();
1012        schedule_error.set_executor(MultiThreadedExecutor::new());
1013        schedule_error.add_systems(|| Err(BevyError::ignore("")));
1014
1015        let mut schedule_panic = Schedule::default();
1016        schedule_panic.set_executor(MultiThreadedExecutor::new());
1017        schedule_panic.add_systems(|| {
1018            panic!("System's panic payload");
1019        });
1020
1021        static HANDLER_CALLED: AtomicBool = AtomicBool::new(false);
1022        fn handle(_: BevyError, ctx: ErrorContext) {
1023            assert!(matches!(ctx, ErrorContext::System { .. }));
1024            HANDLER_CALLED.store(true, Relaxed);
1025        }
1026        world.insert_resource(FallbackErrorHandler(handle));
1027
1028        // System error
1029        schedule_error.run(&mut world);
1030        assert!(HANDLER_CALLED.load(Relaxed));
1031
1032        // System panic
1033        HANDLER_CALLED.store(false, Relaxed);
1034        schedule_panic.run(&mut world);
1035        assert!(HANDLER_CALLED.load(Relaxed));
1036
1037        const PANIC_PAYLOAD: &str = "UwU";
1038        fn panic(_: BevyError, ctx: ErrorContext) {
1039            assert!(matches!(ctx, ErrorContext::System { .. }));
1040            panic!("{}", PANIC_PAYLOAD);
1041        }
1042        world.insert_resource(FallbackErrorHandler(panic));
1043
1044        // System error, handler panic
1045        let result = catch_unwind(AssertUnwindSafe(|| schedule_error.run(&mut world)));
1046        let payload = result.unwrap_err();
1047        assert_eq!(
1048            payload
1049                .downcast_ref::<String>()
1050                .map(String::as_str)
1051                .unwrap_or_else(|| payload.downcast_ref::<&str>().unwrap()),
1052            PANIC_PAYLOAD
1053        );
1054
1055        // System panic, handler panic
1056        let result = catch_unwind(AssertUnwindSafe(|| schedule_panic.run(&mut world)));
1057        let payload = result.unwrap_err();
1058        assert_eq!(
1059            payload
1060                .downcast_ref::<String>()
1061                .map(String::as_str)
1062                .unwrap_or_else(|| payload.downcast_ref::<&str>().unwrap()),
1063            PANIC_PAYLOAD
1064        );
1065
1066        static SYSTEM_RAN: AtomicBool = AtomicBool::new(false);
1067        let system = || {
1068            SYSTEM_RAN.store(true, Relaxed);
1069        };
1070        let mut schedule_condition_error = Schedule::default();
1071        schedule_condition_error.set_executor(MultiThreadedExecutor::new());
1072        schedule_condition_error.add_systems(system.run_if(|| Err(BevyError::ignore(""))));
1073
1074        let mut schedule_condition_panic = Schedule::default();
1075        schedule_condition_panic.set_executor(MultiThreadedExecutor::new());
1076        schedule_condition_panic.add_systems(system.run_if(|| {
1077            panic!("Condition's panic payload");
1078        }));
1079
1080        world.insert_resource(FallbackErrorHandler(handle));
1081
1082        // Condition error
1083        schedule_error.run(&mut world);
1084        assert!(HANDLER_CALLED.load(Relaxed));
1085        assert!(!SYSTEM_RAN.load(Relaxed));
1086
1087        // Condition panic
1088        HANDLER_CALLED.store(false, Relaxed);
1089        schedule_panic.run(&mut world);
1090        assert!(HANDLER_CALLED.load(Relaxed));
1091        assert!(!SYSTEM_RAN.load(Relaxed));
1092
1093        world.insert_resource(FallbackErrorHandler(panic));
1094
1095        // Condition error, handler panic
1096        let result = catch_unwind(AssertUnwindSafe(|| schedule_error.run(&mut world)));
1097        let payload = result.unwrap_err();
1098        assert_eq!(
1099            payload
1100                .downcast_ref::<String>()
1101                .map(String::as_str)
1102                .unwrap_or_else(|| payload.downcast_ref::<&str>().unwrap()),
1103            PANIC_PAYLOAD
1104        );
1105        assert!(!SYSTEM_RAN.load(Relaxed));
1106
1107        // Condition panic, handler panic
1108        let result = catch_unwind(AssertUnwindSafe(|| schedule_panic.run(&mut world)));
1109        let payload = result.unwrap_err();
1110        assert_eq!(
1111            payload
1112                .downcast_ref::<String>()
1113                .map(String::as_str)
1114                .unwrap_or_else(|| payload.downcast_ref::<&str>().unwrap()),
1115            PANIC_PAYLOAD
1116        );
1117        assert!(!SYSTEM_RAN.load(Relaxed));
1118    }
1119}