use crate::{
change_detection::MaybeLocation,
system::{Command, SystemBuffer, SystemMeta},
world::{DeferredWorld, World},
};
use alloc::vec::Vec;
use bevy_ptr::{OwningPtr, Unaligned};
use core::{
fmt::Debug,
mem::{size_of, MaybeUninit},
ptr::NonNull,
};
use log::warn;
#[cfg(feature = "std")]
use crate::error::{BevyError, ErrorContext};
#[cfg(feature = "std")]
use alloc::boxed::Box;
#[cfg(feature = "std")]
use bevy_utils::DebugName;
#[cfg(feature = "std")]
use std::panic::{catch_unwind, resume_unwind, AssertUnwindSafe};
struct CommandMeta {
consume_command_and_get_size:
unsafe fn(value: OwningPtr<Unaligned>, world: Option<&mut World>, cursor: &mut usize),
}
pub struct CommandQueue {
pub(crate) bytes: Vec<MaybeUninit<u8>>,
pub(crate) caller: MaybeLocation,
warn_on_unapplied: bool,
}
impl Default for CommandQueue {
#[track_caller]
fn default() -> Self {
Self {
bytes: Default::default(),
caller: MaybeLocation::caller(),
warn_on_unapplied: true,
}
}
}
impl Debug for CommandQueue {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("CommandQueue")
.field("len_bytes", &self.bytes.len())
.field("caller", &self.caller)
.finish_non_exhaustive()
}
}
unsafe impl Send for CommandQueue {}
unsafe impl Sync for CommandQueue {}
impl CommandQueue {
#[track_caller]
pub fn silent() -> Self {
CommandQueue {
bytes: Default::default(),
caller: MaybeLocation::caller(),
warn_on_unapplied: false,
}
}
#[inline]
pub fn push<C: Command<Out = ()>>(&mut self, command: C) {
#[repr(C, packed)]
struct Packed<C: Command<Out = ()>> {
meta: CommandMeta,
command: C,
}
let meta = CommandMeta {
consume_command_and_get_size: |command, mut world, cursor| {
*cursor += size_of::<C>();
let command: C = unsafe { command.read_unaligned() };
let f = || {
match world.as_deref_mut() {
Some(world) => {
command.apply(world);
world.flush();
}
None => drop(command),
}
};
#[cfg(feature = "std")]
{
let result = catch_unwind(AssertUnwindSafe(f));
if let Err(payload) = result {
let name = DebugName::type_name::<C>();
handle_panic_payload(world, payload, name);
}
}
#[cfg(not(feature = "std"))]
(f)();
},
};
let old_len = self.bytes.len();
self.bytes.reserve(size_of::<Packed<C>>());
let ptr = unsafe { self.bytes.as_mut_ptr().add(old_len) };
unsafe {
ptr.cast::<Packed<C>>()
.write_unaligned(Packed { meta, command });
}
unsafe {
self.bytes.set_len(old_len + size_of::<Packed<C>>());
}
}
#[inline]
pub fn apply(&mut self, world: &mut World) {
world.flush_commands();
let mut runner = unsafe { CommandQueueRunner::new((self, world), |(queue, _)| queue, 0) };
runner.run(|(_, world)| Some(world));
}
pub fn append(&mut self, other: &mut CommandQueue) {
self.bytes.append(&mut other.bytes);
}
#[inline]
pub fn is_empty(&self) -> bool {
self.bytes.is_empty()
}
pub(crate) fn len(&self) -> usize {
self.bytes.len()
}
pub fn silence_drop_warning(&mut self) {
self.warn_on_unapplied = false;
}
}
impl Drop for CommandQueue {
fn drop(&mut self) {
if !self.bytes.is_empty() && self.warn_on_unapplied {
if let Some(caller) = self.caller.into_option() {
warn!("CommandQueue has un-applied commands being dropped. Did you forget to call SystemState::apply? caller:{caller:?}");
} else {
warn!("CommandQueue has un-applied commands being dropped. Did you forget to call SystemState::apply?");
}
}
unsafe { drop(CommandQueueRunner::new(self, |queue| queue, 0)) };
}
}
impl SystemBuffer for CommandQueue {
#[inline]
fn apply(&mut self, _system_meta: &SystemMeta, world: &mut World) {
#[cfg(feature = "trace")]
let _span_guard = _system_meta.commands_span.enter();
self.apply(world);
}
#[inline]
fn queue(&mut self, _system_meta: &SystemMeta, mut world: DeferredWorld) {
world.commands().append(self);
}
}
pub(crate) struct CommandQueueRunner<D, F>
where
F: Fn(&mut D) -> &mut CommandQueue,
{
data: D,
command_queue: F,
local_cursor: usize,
start: usize,
stop: usize,
}
impl<D, F> CommandQueueRunner<D, F>
where
F: Fn(&mut D) -> &mut CommandQueue,
{
pub unsafe fn new(mut data: D, command_queue: F, start: usize) -> Self {
let stop = command_queue(&mut data).len();
Self {
data,
command_queue,
local_cursor: start,
start,
stop,
}
}
pub fn run(&mut self, world: impl Fn(&mut D) -> Option<&mut World>) {
while self.local_cursor < self.stop {
let command_queue = (self.command_queue)(&mut self.data);
let meta = unsafe {
command_queue
.bytes
.as_mut_ptr()
.add(self.local_cursor)
.cast::<CommandMeta>()
.read_unaligned()
};
self.local_cursor += size_of::<CommandMeta>();
let cmd = unsafe {
OwningPtr::<Unaligned>::new(NonNull::new_unchecked(
command_queue
.bytes
.as_mut_ptr()
.add(self.local_cursor)
.cast(),
))
};
unsafe {
(meta.consume_command_and_get_size)(
cmd,
world(&mut self.data),
&mut self.local_cursor,
);
}
}
}
}
#[cfg(feature = "std")]
#[cold]
fn handle_panic_payload(
world: Option<&mut World>,
payload: Box<dyn core::any::Any + Send>,
name: DebugName,
) {
let Some(world) = world else {
resume_unwind(payload)
};
let error = BevyError::panic("Command panicked", payload);
world.fallback_error_handler()(error, ErrorContext::Command { name });
}
impl<D, F> Drop for CommandQueueRunner<D, F>
where
F: Fn(&mut D) -> &mut CommandQueue,
{
fn drop(&mut self) {
self.run(|_| None);
let command_queue = (self.command_queue)(&mut self.data);
unsafe { command_queue.bytes.set_len(self.start) };
}
}
#[cfg(test)]
mod test {
use super::*;
use crate::{
component::Component,
error::{BevyError, ErrorContext, FallbackErrorHandler},
resource::Resource,
};
use alloc::{
borrow::ToOwned,
string::{String, ToString},
sync::Arc,
};
use core::{
panic::AssertUnwindSafe,
sync::atomic::{AtomicU32, Ordering},
};
use std::sync::Mutex;
#[cfg(miri)]
use alloc::format;
struct DropCheck(Arc<AtomicU32>);
impl DropCheck {
fn new() -> (Self, Arc<AtomicU32>) {
let drops = Arc::new(AtomicU32::new(0));
(Self(drops.clone()), drops)
}
}
impl Drop for DropCheck {
fn drop(&mut self) {
self.0.fetch_add(1, Ordering::Relaxed);
}
}
impl Command for DropCheck {
type Out = ();
fn apply(self, _: &mut World) {}
}
#[test]
fn test_command_queue_inner_drop() {
let mut queue = CommandQueue::default();
let (dropcheck_a, drops_a) = DropCheck::new();
let (dropcheck_b, drops_b) = DropCheck::new();
queue.push(dropcheck_a);
queue.push(dropcheck_b);
assert_eq!(drops_a.load(Ordering::Relaxed), 0);
assert_eq!(drops_b.load(Ordering::Relaxed), 0);
let mut world = World::new();
queue.apply(&mut world);
assert_eq!(drops_a.load(Ordering::Relaxed), 1);
assert_eq!(drops_b.load(Ordering::Relaxed), 1);
}
#[test]
fn test_command_queue_inner_drop_early() {
let mut queue = CommandQueue::default();
let (dropcheck_a, drops_a) = DropCheck::new();
let (dropcheck_b, drops_b) = DropCheck::new();
queue.push(dropcheck_a);
queue.push(dropcheck_b);
assert_eq!(drops_a.load(Ordering::Relaxed), 0);
assert_eq!(drops_b.load(Ordering::Relaxed), 0);
drop(queue);
assert_eq!(drops_a.load(Ordering::Relaxed), 1);
assert_eq!(drops_b.load(Ordering::Relaxed), 1);
}
#[derive(Component)]
struct A;
struct SpawnCommand;
impl Command for SpawnCommand {
type Out = ();
fn apply(self, world: &mut World) {
world.spawn(A);
}
}
#[test]
fn test_command_queue_inner() {
let mut queue = CommandQueue::default();
queue.push(SpawnCommand);
queue.push(SpawnCommand);
let mut world = World::new();
queue.apply(&mut world);
assert_eq!(world.query::<&A>().query(&world).count(), 2);
queue.apply(&mut world);
assert_eq!(world.query::<&A>().query(&world).count(), 2);
}
#[expect(
dead_code,
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."
)]
struct PanicCommand(String);
impl Command for PanicCommand {
type Out = ();
fn apply(self, _: &mut World) {
panic!("command is panicking");
}
}
#[test]
fn test_command_queue_inner_panic_safe_panic() {
let mut queue = CommandQueue::default();
queue.push(PanicCommand("I panic!".to_owned()));
queue.push(SpawnCommand);
let mut world = World::new();
let _ = catch_unwind(AssertUnwindSafe(|| {
queue.apply(&mut world);
}));
queue.push(SpawnCommand);
queue.push(SpawnCommand);
queue.apply(&mut world);
assert_eq!(world.query::<&A>().query(&world).count(), 2);
}
#[test]
fn test_command_queue_inner_panic_safe_handled() {
let mut queue = CommandQueue::default();
queue.push(PanicCommand("I panic!".to_owned()));
queue.push(SpawnCommand);
fn record_last_error(error: BevyError, context: ErrorContext) {
*LAST_ERROR.lock().unwrap() = Some((error, context));
}
static LAST_ERROR: Mutex<Option<(BevyError, ErrorContext)>> = Mutex::new(None);
*LAST_ERROR.lock().unwrap() = None;
let mut world = World::new();
world.insert_resource(FallbackErrorHandler(record_last_error));
queue.apply(&mut world);
queue.push(SpawnCommand);
queue.push(SpawnCommand);
queue.apply(&mut world);
assert_eq!(world.query::<&A>().query(&world).count(), 3);
let (error, context) = LAST_ERROR.lock().unwrap().take().unwrap();
assert!(error.to_string().contains("Command panicked"));
let name = DebugName::type_name::<PanicCommand>();
assert_eq!(context, ErrorContext::Command { name });
}
#[test]
fn test_command_queue_inner_nested_panic_safe_panic() {
#[derive(Resource, Default)]
struct Order(Vec<usize>);
let mut world = World::new();
world.init_resource::<Order>();
fn add_index(index: usize) -> impl Command {
move |world: &mut World| world.resource_mut::<Order>().0.push(index)
}
world.commands().queue(add_index(1));
world.commands().queue(|world: &mut World| {
world.commands().queue(add_index(2));
world.commands().queue(PanicCommand("I panic!".to_owned()));
world.commands().queue(add_index(3));
world.flush_commands();
});
world.commands().queue(add_index(4));
let _ = catch_unwind(AssertUnwindSafe(|| {
world.flush_commands();
}));
world.commands().queue(add_index(5));
world.flush_commands();
assert_eq!(&world.resource::<Order>().0, &[1, 2, 5]);
}
#[test]
fn test_command_queue_inner_nested_panic_safe_handled() {
#[derive(Resource, Default)]
struct Order(Vec<usize>);
fn record_last_error(error: BevyError, context: ErrorContext) {
*LAST_ERROR.lock().unwrap() = Some((error, context));
}
static LAST_ERROR: Mutex<Option<(BevyError, ErrorContext)>> = Mutex::new(None);
*LAST_ERROR.lock().unwrap() = None;
let mut world = World::new();
world.init_resource::<Order>();
world.insert_resource(FallbackErrorHandler(record_last_error));
fn add_index(index: usize) -> impl Command {
move |world: &mut World| world.resource_mut::<Order>().0.push(index)
}
world.commands().queue(add_index(1));
world.commands().queue(|world: &mut World| {
world.commands().queue(add_index(2));
world.commands().queue(PanicCommand("I panic!".to_owned()));
world.commands().queue(add_index(3));
world.flush_commands();
});
world.commands().queue(add_index(4));
world.flush_commands();
world.commands().queue(add_index(5));
world.flush_commands();
assert_eq!(&world.resource::<Order>().0, &[1, 2, 3, 4, 5]);
let (error, context) = LAST_ERROR.lock().unwrap().take().unwrap();
assert!(error.to_string().contains("Command panicked"));
let name = DebugName::type_name_of_val(&PanicCommand(String::new()).handle_error());
assert_eq!(context, ErrorContext::Command { name });
}
fn assert_is_send_impl(_: impl Send) {}
fn assert_is_send(command: impl Command) {
assert_is_send_impl(command);
}
#[test]
fn test_command_is_send() {
assert_is_send(SpawnCommand);
}
#[expect(
dead_code,
reason = "This struct is used to test how the CommandQueue reacts to padding added by rust's compiler."
)]
struct CommandWithPadding(u8, u16);
impl Command for CommandWithPadding {
type Out = ();
fn apply(self, _: &mut World) {}
}
#[cfg(miri)]
#[test]
fn test_uninit_bytes() {
let mut queue = CommandQueue::default();
queue.push(CommandWithPadding(0, 0));
let _ = format!("{:?}", queue.bytes);
}
}