1#[cfg(feature = "multi_threaded")]
2use crate::message::MessageParIter;
3use crate::{
4 message::{Message, MessageCursor, MessageIterator, MessageIteratorWithId, Messages},
5 system::{Local, Res, SystemParam, SystemParamValidationError},
6};
7
8#[derive(const _: () =
{
type __StructFieldsAlias<'w, 's, M> =
(Local<'s, MessageCursor<M>>, Res<'w, Messages<M>>);
#[doc(hidden)]
pub struct FetchState<M: Message> {
state: <__StructFieldsAlias<'static, 'static, M> as
bevy_ecs::system::SystemParam>::State,
}
unsafe impl<M: Message> bevy_ecs::system::SystemParam for
MessageReader<'_, '_, M> {
type State = FetchState<M>;
type Item<'w, 's> = MessageReader<'w, 's, M>;
fn init_state(world: &mut bevy_ecs::world::World) -> Self::State {
FetchState {
state: <__StructFieldsAlias<'_, '_, M> as
bevy_ecs::system::SystemParam>::init_state(world),
}
}
fn init_access(state: &Self::State,
system_meta: &mut bevy_ecs::system::SystemMeta,
component_access_set: &mut bevy_ecs::query::FilteredAccessSet,
world: &mut bevy_ecs::world::World) {
<__StructFieldsAlias<'_, '_, M> as
bevy_ecs::system::SystemParam>::init_access(&state.state,
system_meta, component_access_set, world);
}
fn apply(state: &mut Self::State,
system_meta: &bevy_ecs::system::SystemMeta,
world: &mut bevy_ecs::world::World) {
<__StructFieldsAlias<'_, '_, M> as
bevy_ecs::system::SystemParam>::apply(&mut state.state,
system_meta, world);
}
fn queue(state: &mut Self::State,
system_meta: &bevy_ecs::system::SystemMeta,
world: bevy_ecs::world::DeferredWorld) {
<__StructFieldsAlias<'_, '_, M> as
bevy_ecs::system::SystemParam>::queue(&mut state.state,
system_meta, world);
}
#[inline]
unsafe fn get_param<'w,
's>(state: &'s mut Self::State,
system_meta: &bevy_ecs::system::SystemMeta,
world:
bevy_ecs::world::unsafe_world_cell::UnsafeWorldCell<'w>,
change_tick: bevy_ecs::change_detection::Tick)
->
::core::result::Result<Self::Item<'w, 's>,
bevy_ecs::system::SystemParamValidationError> {
let (fieldreader, fieldmessages) = &mut state.state;
let fieldreader =
unsafe {
<Local<'s, MessageCursor<M>> as
bevy_ecs::system::SystemParam>::get_param(fieldreader,
system_meta, world, change_tick)
}.map_err(|err|
bevy_ecs::system::SystemParamValidationError::new::<Self>(err.skipped,
err.message, "::reader"))?;
let fieldmessages =
unsafe {
<Res<'w, Messages<M>> as
bevy_ecs::system::SystemParam>::get_param(fieldmessages,
system_meta, world, change_tick)
}.map_err(|err|
bevy_ecs::system::SystemParamValidationError::new::<Self>(err.skipped,
"Message not initialized", "::messages"))?;
::core::result::Result::Ok(MessageReader {
reader: fieldreader,
messages: fieldmessages,
})
}
}
unsafe impl<'w, 's, M: Message> bevy_ecs::system::ReadOnlySystemParam
for MessageReader<'w, 's, M> where
Local<'s,
MessageCursor<M>>: bevy_ecs::system::ReadOnlySystemParam,
Res<'w, Messages<M>>: bevy_ecs::system::ReadOnlySystemParam {}
};SystemParam, #[automatically_derived]
impl<'w, 's, M: ::core::fmt::Debug + Message> ::core::fmt::Debug for
MessageReader<'w, 's, M> {
#[inline]
fn fmt(&self, f: &mut ::core::fmt::Formatter) -> ::core::fmt::Result {
::core::fmt::Formatter::debug_struct_field2_finish(f, "MessageReader",
"reader", &self.reader, "messages", &&self.messages)
}
}Debug)]
34pub struct MessageReader<'w, 's, M: Message> {
35 pub(super) reader: Local<'s, MessageCursor<M>>,
36 #[system_param(validation_message = "Message not initialized")]
37 messages: Res<'w, Messages<M>>,
38}
39
40impl<'w, 's, M: Message> MessageReader<'w, 's, M> {
41 pub fn read(&mut self) -> MessageIterator<'_, M> {
45 self.reader.read(&self.messages)
46 }
47
48 pub fn read_with_id(&mut self) -> MessageIteratorWithId<'_, M> {
50 self.reader.read_with_id(&self.messages)
51 }
52
53 #[cfg(feature = "multi_threaded")]
89 pub fn par_read(&mut self) -> MessageParIter<'_, M> {
90 self.reader.par_read(&self.messages)
91 }
92
93 pub fn len(&self) -> usize {
95 self.reader.len(&self.messages)
96 }
97
98 pub fn is_empty(&self) -> bool {
120 self.reader.is_empty(&self.messages)
121 }
122
123 pub fn clear(&mut self) {
130 self.reader.clear(&self.messages);
131 }
132}
133
134#[derive(#[automatically_derived]
impl<'w, 's, M: ::core::fmt::Debug + Message> ::core::fmt::Debug for
PopulatedMessageReader<'w, 's, M> {
#[inline]
fn fmt(&self, f: &mut ::core::fmt::Formatter) -> ::core::fmt::Result {
::core::fmt::Formatter::debug_tuple_field1_finish(f,
"PopulatedMessageReader", &&self.0)
}
}Debug)]
141pub struct PopulatedMessageReader<'w, 's, M: Message>(MessageReader<'w, 's, M>);
142
143impl<'w, 's, M: Message> core::ops::Deref for PopulatedMessageReader<'w, 's, M> {
144 type Target = MessageReader<'w, 's, M>;
145
146 fn deref(&self) -> &Self::Target {
147 &self.0
148 }
149}
150
151impl<'w, 's, M: Message> core::ops::DerefMut for PopulatedMessageReader<'w, 's, M> {
152 fn deref_mut(&mut self) -> &mut Self::Target {
153 &mut self.0
154 }
155}
156
157unsafe impl<'w, 's, M: Message> SystemParam for PopulatedMessageReader<'w, 's, M> {
159 type State = <MessageReader<'w, 's, M> as SystemParam>::State;
160 type Item<'world, 'state> = PopulatedMessageReader<'world, 'state, M>;
161
162 fn init_state(world: &mut crate::prelude::World) -> Self::State {
163 MessageReader::<M>::init_state(world)
164 }
165
166 fn init_access(
167 state: &Self::State,
168 system_meta: &mut crate::system::SystemMeta,
169 component_access_set: &mut crate::query::FilteredAccessSet,
170 world: &mut crate::prelude::World,
171 ) {
172 MessageReader::<M>::init_access(state, system_meta, component_access_set, world);
173 }
174
175 unsafe fn get_param<'world, 'state>(
176 state: &'state mut Self::State,
177 system_meta: &crate::system::SystemMeta,
178 world: crate::world::unsafe_world_cell::UnsafeWorldCell<'world>,
179 change_tick: crate::change_detection::Tick,
180 ) -> Result<Self::Item<'world, 'state>, SystemParamValidationError> {
181 let reader = unsafe { MessageReader::get_param(state, system_meta, world, change_tick)? };
183 if reader.is_empty() {
184 Err(SystemParamValidationError::skipped::<Self>(
185 "message queue is empty",
186 ))
187 } else {
188 Ok(PopulatedMessageReader(reader))
189 }
190 }
191}
192
193#[cfg(test)]
194mod tests {
195 use core::sync::atomic::{AtomicBool, Ordering};
196
197 use super::*;
198 use crate::message::MessageRegistry;
199 use crate::prelude::*;
200 use bevy_platform::sync::Arc;
201
202 #[test]
203 fn test_populated_message_reader() {
204 let system_ran = Arc::new(AtomicBool::new(false));
205
206 let mut world = World::new();
207 MessageRegistry::register_message::<TheMessage>(&mut world);
208
209 let mut schedule = Schedule::default();
210 schedule.add_systems({
211 let system_ran = system_ran.clone();
212 move |mut _reader: PopulatedMessageReader<TheMessage>| {
213 system_ran.store(true, Ordering::SeqCst);
214 }
215 });
216
217 schedule.run(&mut world);
218 assert!(
219 !system_ran.load(Ordering::SeqCst),
220 "system with PopulatedMessageReader should have been skipped"
221 );
222
223 world.write_message(TheMessage);
224 schedule.run(&mut world);
225 assert!(
226 system_ran.load(Ordering::SeqCst),
227 "system with PopulatedMessageReader should NOT have been skipped"
228 );
229
230 #[derive(Message)]
231 struct TheMessage;
232 }
233}