Skip to main content

bevy_ecs/world/
command_queue.rs

1use crate::{
2    change_detection::MaybeLocation,
3    system::{Command, SystemBuffer, SystemMeta},
4    world::{DeferredWorld, World},
5};
6
7use alloc::vec::Vec;
8use bevy_ptr::{OwningPtr, Unaligned};
9use core::{
10    fmt::Debug,
11    mem::{size_of, MaybeUninit},
12    ptr::NonNull,
13};
14use log::warn;
15
16#[cfg(feature = "std")]
17use crate::error::{BevyError, ErrorContext};
18#[cfg(feature = "std")]
19use alloc::boxed::Box;
20#[cfg(feature = "std")]
21use bevy_utils::DebugName;
22#[cfg(feature = "std")]
23use std::panic::{catch_unwind, resume_unwind, AssertUnwindSafe};
24
25struct CommandMeta {
26    /// SAFETY: The `value` must point to a value of type `T: Command`,
27    /// where `T` is some specific type that was used to produce this metadata.
28    ///
29    /// `world` is optional to allow this one function pointer to perform double-duty as a drop.
30    ///
31    /// Advances `cursor` by the size of `T` in bytes.
32    consume_command_and_get_size:
33        unsafe fn(value: OwningPtr<Unaligned>, world: Option<&mut World>, cursor: &mut usize),
34}
35
36/// Densely and efficiently stores a queue of heterogenous types implementing [`Command`].
37// NOTE: [`CommandQueue`] is implemented via a `Vec<MaybeUninit<u8>>` instead of a `Vec<Box<dyn Command>>`
38// as an optimization. Since commands are used frequently in systems as a way to spawn
39// entities/components/resources, and it's not currently possible to parallelize these
40// due to mutable [`World`] access, maximizing performance for [`CommandQueue`] is
41// preferred to simplicity of implementation.
42pub struct CommandQueue {
43    /// This buffer densely stores all queued commands.
44    ///
45    /// For each command, one `CommandMeta` is stored, followed by zero or more bytes
46    /// to store the command itself. To interpret these bytes, a pointer must
47    /// be passed to the corresponding `CommandMeta.apply_command_and_get_size` fn pointer.
48    pub(crate) bytes: Vec<MaybeUninit<u8>>,
49    pub(crate) caller: MaybeLocation,
50    /// Always emit a warning if a command is dropped before it is applied.
51    /// Defaults to `true`.
52    ///
53    /// This setting can be turned off for commands that might be dropped (due to application exit) before those
54    /// commands are applied in ordinary situations, for example delayed commands.
55    warn_on_unapplied: bool,
56}
57
58impl Default for CommandQueue {
59    #[track_caller]
60    fn default() -> Self {
61        Self {
62            bytes: Default::default(),
63            caller: MaybeLocation::caller(),
64            warn_on_unapplied: true,
65        }
66    }
67}
68
69// CommandQueue needs to implement Debug manually, rather than deriving it, because the derived impl just prints
70// [core::mem::maybe_uninit::MaybeUninit<u8>, core::mem::maybe_uninit::MaybeUninit<u8>, ..] for every byte in the vec,
71// which gets extremely verbose very quickly, while also providing no useful information.
72// It is not possible to soundly print the values of the contained bytes, as some of them may be padding or uninitialized (#4863)
73// So instead, the manual impl just prints the length of vec.
74impl Debug for CommandQueue {
75    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
76        f.debug_struct("CommandQueue")
77            .field("len_bytes", &self.bytes.len())
78            .field("caller", &self.caller)
79            .finish_non_exhaustive()
80    }
81}
82
83// SAFETY: All commands [`Command`] implement [`Send`]
84unsafe impl Send for CommandQueue {}
85
86// SAFETY: `&CommandQueue` never gives access to the inner commands.
87unsafe impl Sync for CommandQueue {}
88
89impl CommandQueue {
90    /// Create a queue that does not warn when dropped.
91    #[track_caller]
92    pub fn silent() -> Self {
93        CommandQueue {
94            bytes: Default::default(),
95            caller: MaybeLocation::caller(),
96            warn_on_unapplied: false,
97        }
98    }
99
100    /// Push a [`Command`] onto the queue.
101    #[inline]
102    pub fn push<C: Command<Out = ()>>(&mut self, command: C) {
103        // Stores a command alongside its metadata.
104        // `repr(C)` prevents the compiler from reordering the fields,
105        // while `repr(packed)` prevents the compiler from inserting padding bytes.
106        #[repr(C, packed)]
107        struct Packed<C: Command<Out = ()>> {
108            meta: CommandMeta,
109            command: C,
110        }
111
112        let meta = CommandMeta {
113            consume_command_and_get_size: |command, mut world, cursor| {
114                *cursor += size_of::<C>();
115
116                // Putting the command onto the stack is necessary not just for alignment and to be able to consume it,
117                // but also because applying the command may cause the command queue to reallocate.
118                // SAFETY: According to the invariants of `CommandMeta.consume_command_and_get_size`,
119                // `command` must point to a value of type `C`.
120                let command: C = unsafe { command.read_unaligned() };
121
122                let f = || {
123                    match world.as_deref_mut() {
124                        // Apply command to the provided world...
125                        Some(world) => {
126                            command.apply(world);
127                            // The command may have queued up world commands, which we flush here to ensure they are also picked up.
128                            // If the current command queue already the World Command queue, this will still behave appropriately because the global cursor
129                            // is still at the current `stop`, ensuring only the newly queued Commands will be applied.
130                            world.flush();
131                        }
132                        // ...or discard it.
133                        None => drop(command),
134                    }
135                };
136
137                #[cfg(feature = "std")]
138                {
139                    let result = catch_unwind(AssertUnwindSafe(f));
140                    if let Err(payload) = result {
141                        let name = DebugName::type_name::<C>();
142                        handle_panic_payload(world, payload, name);
143                    }
144                }
145
146                #[cfg(not(feature = "std"))]
147                (f)();
148            },
149        };
150
151        let old_len = self.bytes.len();
152
153        // Reserve enough bytes for both the metadata and the command itself.
154        self.bytes.reserve(size_of::<Packed<C>>());
155
156        // Pointer to the bytes at the end of the buffer.
157        // SAFETY: We know it is within bounds of the allocation, due to the call to `.reserve()`.
158        let ptr = unsafe { self.bytes.as_mut_ptr().add(old_len) };
159
160        // Write the metadata into the buffer, followed by the command.
161        // We are using a packed struct to write them both as one operation.
162        // SAFETY: `ptr` must be non-null, since it is within a non-null buffer.
163        // The call to `reserve()` ensures that the buffer has enough space to fit a value of type `C`,
164        // and it is valid to write any bit pattern since the underlying buffer is of type `MaybeUninit<u8>`.
165        unsafe {
166            ptr.cast::<Packed<C>>()
167                .write_unaligned(Packed { meta, command });
168        }
169
170        // Extend the length of the buffer to include the data we just wrote.
171        // SAFETY: The new length is guaranteed to fit in the vector's capacity,
172        // due to the call to `.reserve()` above.
173        unsafe {
174            self.bytes.set_len(old_len + size_of::<Packed<C>>());
175        }
176    }
177
178    /// Execute the queued [`Command`]s in the world after applying any commands in the world's internal queue.
179    /// This clears the queue.
180    #[inline]
181    pub fn apply(&mut self, world: &mut World) {
182        // flush the world's internal queue
183        world.flush_commands();
184        // SAFETY:
185        // * `self` is always returned
186        // * The first command always start at 0
187        // * `&mut self` prevents all other access to this queue
188        let mut runner = unsafe { CommandQueueRunner::new((self, world), |(queue, _)| queue, 0) };
189        runner.run(|(_, world)| Some(world));
190    }
191
192    /// Take all commands from `other` and append them to `self`, leaving `other` empty
193    pub fn append(&mut self, other: &mut CommandQueue) {
194        self.bytes.append(&mut other.bytes);
195    }
196
197    /// Returns false if there are any commands in the queue
198    #[inline]
199    pub fn is_empty(&self) -> bool {
200        self.bytes.is_empty()
201    }
202
203    /// The number of bytes of commands in the queue.
204    pub(crate) fn len(&self) -> usize {
205        self.bytes.len()
206    }
207
208    /// Silences drop warning if commands are unapplied.
209    pub fn silence_drop_warning(&mut self) {
210        self.warn_on_unapplied = false;
211    }
212}
213
214impl Drop for CommandQueue {
215    fn drop(&mut self) {
216        if !self.bytes.is_empty() && self.warn_on_unapplied {
217            if let Some(caller) = self.caller.into_option() {
218                warn!("CommandQueue has un-applied commands being dropped. Did you forget to call SystemState::apply? caller:{caller:?}");
219            } else {
220                warn!("CommandQueue has un-applied commands being dropped. Did you forget to call SystemState::apply?");
221            }
222        }
223        // Dropping a `CommandQueueRunner` will drop all unapplied commands.
224        // SAFETY:
225        // * `self` is always returned
226        // * The first command always start at 0
227        // * `&mut self` prevents all other access to this queue
228        unsafe { drop(CommandQueueRunner::new(self, |queue| queue, 0)) };
229    }
230}
231
232impl SystemBuffer for CommandQueue {
233    #[inline]
234    fn apply(&mut self, _system_meta: &SystemMeta, world: &mut World) {
235        #[cfg(feature = "trace")]
236        let _span_guard = _system_meta.commands_span.enter();
237        self.apply(world);
238    }
239
240    #[inline]
241    fn queue(&mut self, _system_meta: &SystemMeta, mut world: DeferredWorld) {
242        world.commands().append(self);
243    }
244}
245
246/// A RAII guard used while running commands to ensure
247/// that unapplied commands are dropped during unwind.
248pub(crate) struct CommandQueueRunner<D, F>
249where
250    F: Fn(&mut D) -> &mut CommandQueue,
251{
252    data: D,
253    command_queue: F,
254    local_cursor: usize,
255    start: usize,
256    stop: usize,
257}
258
259impl<D, F> CommandQueueRunner<D, F>
260where
261    F: Fn(&mut D) -> &mut CommandQueue,
262{
263    /// Constructs a new [`CommandQueueRunner`] for the given queue.
264    ///
265    /// This applies or drops commands from `start` to the current end of the queue,
266    /// and will truncate the queue to `start` when dropped.
267    ///
268    /// Stores references to just the [`World`] (when running the world command queue), to just a [`CommandQueue`] (when dropping), or both (when running any other queues).
269    /// In the case of the world's command queue, this allows us to fetch a new reference to the queue after every command, as the command can invalidate it due to receiving the `&mut World`.
270    /// References to other command queues can be stored directly.
271    ///
272    /// # Safety
273    ///
274    /// * `command_queue(&mut data)` must always return the same queue.
275    /// * `start` is the index of the first byte of a command in the queue,
276    ///   or the length of the queue
277    /// * Until the `CommandQueueRunner` is dropped, nothing else may
278    ///   access commands between `start` and `command_queue.len()`
279    pub unsafe fn new(mut data: D, command_queue: F, start: usize) -> Self {
280        let stop = command_queue(&mut data).len();
281        Self {
282            data,
283            command_queue,
284            local_cursor: start,
285            start,
286            stop,
287        }
288    }
289
290    /// Applies or drops the commands in the queue.
291    ///
292    /// If `world` returns [`Some`], this will apply the queued [commands](`Command`) to that `World`.
293    /// If `world` returns [`None`], this will drop the queued [commands](`Command`) (without applying them).
294    pub fn run(&mut self, world: impl Fn(&mut D) -> Option<&mut World>) {
295        while self.local_cursor < self.stop {
296            let command_queue = (self.command_queue)(&mut self.data);
297
298            // We must re-read the pointer to the allocation before each command
299            // as the previous might have cause a reallocation.
300            // SAFETY: The cursor is either at the start of the buffer, or just after the previous command.
301            // Since we know that the cursor is in bounds, it must point to the start of a new command.
302            let meta = unsafe {
303                command_queue
304                    .bytes
305                    .as_mut_ptr()
306                    .add(self.local_cursor)
307                    .cast::<CommandMeta>()
308                    .read_unaligned()
309            };
310
311            // Advance to the bytes just after `meta`, which represent a type-erased command.
312            self.local_cursor += size_of::<CommandMeta>();
313            // Construct an owned pointer to the command.
314            // SAFETY: It is safe to transfer ownership out of `self.bytes`, since the increment of `cursor` above
315            // guarantees that nothing stored in the buffer will get observed after this function ends.
316            // `cmd` points to a valid address of a stored command, so it must be non-null.
317            let cmd = unsafe {
318                OwningPtr::<Unaligned>::new(NonNull::new_unchecked(
319                    command_queue
320                        .bytes
321                        .as_mut_ptr()
322                        .add(self.local_cursor)
323                        .cast(),
324                ))
325            };
326            // SAFETY: The data underneath the cursor must correspond to the type erased in metadata,
327            // since they were stored next to each other by `.push()`.
328            // For ZSTs, the type doesn't matter as long as the pointer is non-null.
329            // This also advances the cursor past the command. For ZSTs, the cursor will not move.
330            // At this point, it will either point to the next `CommandMeta`,
331            // or the cursor will be out of bounds and the loop will end.
332            unsafe {
333                (meta.consume_command_and_get_size)(
334                    cmd,
335                    world(&mut self.data),
336                    &mut self.local_cursor,
337                );
338            }
339        }
340    }
341}
342
343/// Handle a panic thrown within a command.
344///
345/// This is a separate non-generic function so that the panic handling code
346/// is not monomorphized separately for each command type.
347#[cfg(feature = "std")]
348#[cold]
349fn handle_panic_payload(
350    world: Option<&mut World>,
351    payload: Box<dyn core::any::Any + Send>,
352    name: DebugName,
353) {
354    let Some(world) = world else {
355        resume_unwind(payload)
356    };
357    let error = BevyError::panic("Command panicked", payload);
358    world.fallback_error_handler()(error, ErrorContext::Command { name });
359}
360
361impl<D, F> Drop for CommandQueueRunner<D, F>
362where
363    F: Fn(&mut D) -> &mut CommandQueue,
364{
365    fn drop(&mut self) {
366        // Drop any unapplied commands before resetting the length.
367        // If `run` completed successfully then this will do nothing.
368        self.run(|_| None);
369
370        let command_queue = (self.command_queue)(&mut self.data);
371
372        // Reset the buffer: all commands past the original `start` cursor have been applied.
373        // SAFETY: we are setting the length of bytes to the original length, minus the length of the original
374        // list of commands being considered. All bytes remaining in the Vec are still valid, unapplied commands.
375        unsafe { command_queue.bytes.set_len(self.start) };
376    }
377}
378
379#[cfg(test)]
380mod test {
381    use super::*;
382    use crate::{
383        component::Component,
384        error::{BevyError, ErrorContext, FallbackErrorHandler},
385        resource::Resource,
386    };
387    use alloc::{
388        borrow::ToOwned,
389        string::{String, ToString},
390        sync::Arc,
391    };
392    use core::{
393        panic::AssertUnwindSafe,
394        sync::atomic::{AtomicU32, Ordering},
395    };
396    use std::sync::Mutex;
397
398    #[cfg(miri)]
399    use alloc::format;
400
401    struct DropCheck(Arc<AtomicU32>);
402
403    impl DropCheck {
404        fn new() -> (Self, Arc<AtomicU32>) {
405            let drops = Arc::new(AtomicU32::new(0));
406            (Self(drops.clone()), drops)
407        }
408    }
409
410    impl Drop for DropCheck {
411        fn drop(&mut self) {
412            self.0.fetch_add(1, Ordering::Relaxed);
413        }
414    }
415
416    impl Command for DropCheck {
417        type Out = ();
418
419        fn apply(self, _: &mut World) {}
420    }
421
422    #[test]
423    fn test_command_queue_inner_drop() {
424        let mut queue = CommandQueue::default();
425
426        let (dropcheck_a, drops_a) = DropCheck::new();
427        let (dropcheck_b, drops_b) = DropCheck::new();
428
429        queue.push(dropcheck_a);
430        queue.push(dropcheck_b);
431
432        assert_eq!(drops_a.load(Ordering::Relaxed), 0);
433        assert_eq!(drops_b.load(Ordering::Relaxed), 0);
434
435        let mut world = World::new();
436        queue.apply(&mut world);
437
438        assert_eq!(drops_a.load(Ordering::Relaxed), 1);
439        assert_eq!(drops_b.load(Ordering::Relaxed), 1);
440    }
441
442    /// Asserts that inner [commands](`Command`) are dropped on early drop of [`CommandQueue`].
443    /// Originally identified as an issue in [#10676](https://github.com/bevyengine/bevy/issues/10676)
444    #[test]
445    fn test_command_queue_inner_drop_early() {
446        let mut queue = CommandQueue::default();
447
448        let (dropcheck_a, drops_a) = DropCheck::new();
449        let (dropcheck_b, drops_b) = DropCheck::new();
450
451        queue.push(dropcheck_a);
452        queue.push(dropcheck_b);
453
454        assert_eq!(drops_a.load(Ordering::Relaxed), 0);
455        assert_eq!(drops_b.load(Ordering::Relaxed), 0);
456
457        drop(queue);
458
459        assert_eq!(drops_a.load(Ordering::Relaxed), 1);
460        assert_eq!(drops_b.load(Ordering::Relaxed), 1);
461    }
462
463    #[derive(Component)]
464    struct A;
465
466    struct SpawnCommand;
467
468    impl Command for SpawnCommand {
469        type Out = ();
470
471        fn apply(self, world: &mut World) {
472            world.spawn(A);
473        }
474    }
475
476    #[test]
477    fn test_command_queue_inner() {
478        let mut queue = CommandQueue::default();
479
480        queue.push(SpawnCommand);
481        queue.push(SpawnCommand);
482
483        let mut world = World::new();
484        queue.apply(&mut world);
485
486        assert_eq!(world.query::<&A>().query(&world).count(), 2);
487
488        // The previous call to `apply` cleared the queue.
489        // This call should do nothing.
490        queue.apply(&mut world);
491        assert_eq!(world.query::<&A>().query(&world).count(), 2);
492    }
493
494    #[expect(
495        dead_code,
496        reason = "The inner string is used to ensure that, when the PanicCommand gets pushed to the queue, some data is written to the `bytes` vector."
497    )]
498    struct PanicCommand(String);
499    impl Command for PanicCommand {
500        type Out = ();
501
502        fn apply(self, _: &mut World) {
503            panic!("command is panicking");
504        }
505    }
506
507    #[test]
508    fn test_command_queue_inner_panic_safe_panic() {
509        let mut queue = CommandQueue::default();
510
511        queue.push(PanicCommand("I panic!".to_owned()));
512        // This will get skipped due to the panic
513        queue.push(SpawnCommand);
514
515        let mut world = World::new();
516
517        let _ = catch_unwind(AssertUnwindSafe(|| {
518            queue.apply(&mut world);
519        }));
520
521        // Even though the first command panicked, it's still ok to push
522        // more commands.
523        queue.push(SpawnCommand);
524        queue.push(SpawnCommand);
525        queue.apply(&mut world);
526        assert_eq!(world.query::<&A>().query(&world).count(), 2);
527    }
528
529    #[test]
530    fn test_command_queue_inner_panic_safe_handled() {
531        let mut queue = CommandQueue::default();
532
533        queue.push(PanicCommand("I panic!".to_owned()));
534        // This will get run because the fallback error handler
535        // handles the panicking command.
536        queue.push(SpawnCommand);
537
538        fn record_last_error(error: BevyError, context: ErrorContext) {
539            *LAST_ERROR.lock().unwrap() = Some((error, context));
540        }
541        static LAST_ERROR: Mutex<Option<(BevyError, ErrorContext)>> = Mutex::new(None);
542        *LAST_ERROR.lock().unwrap() = None;
543
544        let mut world = World::new();
545        world.insert_resource(FallbackErrorHandler(record_last_error));
546
547        queue.apply(&mut world);
548
549        // Even though the first command panicked, it's still ok to push
550        // more commands.
551        queue.push(SpawnCommand);
552        queue.push(SpawnCommand);
553        queue.apply(&mut world);
554        assert_eq!(world.query::<&A>().query(&world).count(), 3);
555
556        let (error, context) = LAST_ERROR.lock().unwrap().take().unwrap();
557        assert!(error.to_string().contains("Command panicked"));
558        let name = DebugName::type_name::<PanicCommand>();
559        assert_eq!(context, ErrorContext::Command { name });
560    }
561
562    #[test]
563    fn test_command_queue_inner_nested_panic_safe_panic() {
564        #[derive(Resource, Default)]
565        struct Order(Vec<usize>);
566
567        let mut world = World::new();
568        world.init_resource::<Order>();
569
570        fn add_index(index: usize) -> impl Command {
571            move |world: &mut World| world.resource_mut::<Order>().0.push(index)
572        }
573        world.commands().queue(add_index(1));
574        world.commands().queue(|world: &mut World| {
575            world.commands().queue(add_index(2));
576            world.commands().queue(PanicCommand("I panic!".to_owned()));
577            // Everything after here will get skipped due to the panic
578            world.commands().queue(add_index(3));
579            world.flush_commands();
580        });
581        world.commands().queue(add_index(4));
582
583        let _ = catch_unwind(AssertUnwindSafe(|| {
584            world.flush_commands();
585        }));
586
587        world.commands().queue(add_index(5));
588        world.flush_commands();
589        assert_eq!(&world.resource::<Order>().0, &[1, 2, 5]);
590    }
591
592    #[test]
593    fn test_command_queue_inner_nested_panic_safe_handled() {
594        #[derive(Resource, Default)]
595        struct Order(Vec<usize>);
596
597        fn record_last_error(error: BevyError, context: ErrorContext) {
598            *LAST_ERROR.lock().unwrap() = Some((error, context));
599        }
600        static LAST_ERROR: Mutex<Option<(BevyError, ErrorContext)>> = Mutex::new(None);
601        *LAST_ERROR.lock().unwrap() = None;
602
603        let mut world = World::new();
604        world.init_resource::<Order>();
605        world.insert_resource(FallbackErrorHandler(record_last_error));
606
607        fn add_index(index: usize) -> impl Command {
608            move |world: &mut World| world.resource_mut::<Order>().0.push(index)
609        }
610        world.commands().queue(add_index(1));
611        world.commands().queue(|world: &mut World| {
612            world.commands().queue(add_index(2));
613            world.commands().queue(PanicCommand("I panic!".to_owned()));
614            // Everything after here will get run because the
615            // fallback error handler handles the panicking command.
616            world.commands().queue(add_index(3));
617            world.flush_commands();
618        });
619        world.commands().queue(add_index(4));
620
621        world.flush_commands();
622
623        world.commands().queue(add_index(5));
624        world.flush_commands();
625        assert_eq!(&world.resource::<Order>().0, &[1, 2, 3, 4, 5]);
626
627        let (error, context) = LAST_ERROR.lock().unwrap().take().unwrap();
628        assert!(error.to_string().contains("Command panicked"));
629        let name = DebugName::type_name_of_val(&PanicCommand(String::new()).handle_error());
630        assert_eq!(context, ErrorContext::Command { name });
631    }
632
633    // NOTE: `CommandQueue` is `Send` because `Command` is send.
634    // If the `Command` trait gets reworked to be non-send, `CommandQueue`
635    // should be reworked.
636    // This test asserts that Command types are send.
637    fn assert_is_send_impl(_: impl Send) {}
638    fn assert_is_send(command: impl Command) {
639        assert_is_send_impl(command);
640    }
641
642    #[test]
643    fn test_command_is_send() {
644        assert_is_send(SpawnCommand);
645    }
646
647    #[expect(
648        dead_code,
649        reason = "This struct is used to test how the CommandQueue reacts to padding added by rust's compiler."
650    )]
651    struct CommandWithPadding(u8, u16);
652    impl Command for CommandWithPadding {
653        type Out = ();
654
655        fn apply(self, _: &mut World) {}
656    }
657
658    #[cfg(miri)]
659    #[test]
660    fn test_uninit_bytes() {
661        let mut queue = CommandQueue::default();
662        queue.push(CommandWithPadding(0, 0));
663        let _ = format!("{:?}", queue.bytes);
664    }
665}