1use alloc::boxed::Box;
7use alloc::collections::VecDeque;
8use alloc::rc::{Rc, Weak};
9use alloc::vec::Vec;
10use core::any::Any;
11use core::cell::Cell;
12use core::marker::PhantomData;
13use core::mem;
14
15use crate::event::registry::{EventId, EventRetention};
16use crate::operation::StageOperation;
17use crate::world::{EventReadError, WorldError, WorldOwner};
18
19#[allow(dead_code)]
20pub(crate) struct EventStorage {
21 channels: Vec<EventChannel>,
22}
23
24struct EventChannel {
25 entries: EventEntries,
26 active_len: usize,
27 next_sequence: u64,
28 oldest_retained: u64,
29 retention: EventRetention,
30 cursors: Vec<Weak<Cell<u64>>>,
31 reader_ops_since_prune: u8,
32 closed: bool,
33}
34
35struct EventEntry {
36 sequence: u64,
37 payload: Box<dyn Any>,
38}
39
40enum EventEntries {
41 Linear(Vec<EventEntry>),
42 Ring(VecDeque<EventEntry>),
43}
44
45const READER_PRUNE_INTERVAL: u8 = 128;
46const LINEAR_BOUNDED_CAPACITY_MAX: usize = 16;
47
48#[derive(Copy, Clone, Debug, Eq, PartialEq)]
50pub enum EventReaderStart {
51 OldestRetained,
53 FromNow,
55}
56
57pub struct EventReader<E> {
61 owner: WorldOwner,
62 pub(crate) event_id: EventId,
63 cursor: Rc<Cell<u64>>,
64 last_payload: Option<Box<dyn Any>>,
65 _marker: PhantomData<E>,
66}
67
68impl EventStorage {
69 pub fn new(capacity: usize) -> Self {
70 Self {
71 channels: Vec::with_capacity(capacity),
72 }
73 }
74
75 pub fn ensure_channel(&mut self, index: usize, retention: EventRetention) {
76 while self.channels.len() <= index {
77 self.channels
78 .push(EventChannel::new(EventRetention::Manual));
79 }
80 self.channels[index].set_retention(retention);
81 }
82
83 pub fn send<E: Clone + 'static>(
84 &mut self,
85 event_id: &EventId,
86 event: E,
87 ) -> Result<(), WorldError> {
88 let channel = self.channels.get_mut(event_id.index()).ok_or_else(|| {
89 WorldError::UnregisteredEvent {
90 name: alloc::format!("event {}", event_id.index()),
91 }
92 })?;
93 if channel.closed {
94 return Err(WorldError::EventChannelClosed);
95 }
96 channel.maybe_prune_readers();
97 let sequence = match channel.next_sequence.checked_add(1) {
98 Some(sequence) => sequence,
99 None => {
100 channel.closed = true;
101 return Err(WorldError::EventChannelClosed);
102 }
103 };
104 channel.next_sequence = sequence;
105 channel.push_event(sequence, event);
106 channel.enforce_retention();
107 Ok(())
108 }
109
110 pub fn read_next<'a, E: Clone + 'static>(
111 &mut self,
112 owner: &WorldOwner,
113 reader: &'a mut EventReader<E>,
114 ) -> Result<Option<&'a E>, EventReadError> {
115 reader.validate_owner(owner)?;
116 reader
117 .event_id
118 .validate_owner(owner)
119 .map_err(|_| EventReadError::OwnerMismatch {
120 name: alloc::format!("event {}", reader.event_id.index()),
121 })?;
122 let channel = self
123 .channels
124 .get_mut(reader.event_id.index())
125 .ok_or_else(|| EventReadError::UnregisteredEvent {
126 name: alloc::format!("event {}", reader.event_id.index()),
127 })?;
128 let cursor = reader.cursor.get();
129 if cursor < channel.oldest_retained {
130 let dropped = channel.oldest_retained - cursor;
131 reader.cursor.set(channel.oldest_retained);
132 return Err(EventReadError::Lagged { dropped });
133 }
134 let position = channel.position_after(cursor);
135 let Some(position) = position else {
136 if channel.closed {
137 return Err(EventReadError::ChannelClosed);
138 }
139 return Ok(None);
140 };
141 let entry = channel
142 .entries
143 .get(position)
144 .expect("active event position must be present");
145 let sequence = entry.sequence;
146 let event = entry
147 .payload
148 .downcast_ref::<E>()
149 .ok_or_else(|| EventReadError::UnregisteredEvent {
150 name: alloc::format!("event {}", reader.event_id.index()),
151 })?
152 .clone();
153 reader.last_payload = Some(match reader.last_payload.take() {
154 Some(mut payload) => match payload.downcast_mut::<E>() {
155 Some(slot) => {
156 *slot = event;
157 payload
158 }
159 None => Box::new(event),
160 },
161 None => Box::new(event),
162 });
163 reader.cursor.set(sequence);
164 channel.maybe_prune_readers();
165 Ok(reader
166 .last_payload
167 .as_ref()
168 .and_then(|payload| payload.downcast_ref::<E>()))
169 }
170
171 pub fn create_reader<E: Clone + 'static>(
172 &mut self,
173 owner: WorldOwner,
174 event_id: EventId,
175 start: EventReaderStart,
176 ) -> Result<EventReader<E>, WorldError> {
177 event_id
178 .validate_owner(&owner)
179 .map_err(map_registration_owner_error)?;
180 let channel = self.channels.get_mut(event_id.index()).ok_or_else(|| {
181 WorldError::UnregisteredEvent {
182 name: alloc::format!("event {}", event_id.index()),
183 }
184 })?;
185 let cursor_value = match start {
186 EventReaderStart::OldestRetained => channel
187 .first_active()
188 .map(|entry| entry.sequence.saturating_sub(1))
189 .unwrap_or(channel.oldest_retained),
190 EventReaderStart::FromNow => channel.next_sequence,
191 };
192 let cursor = Rc::new(Cell::new(cursor_value));
193 channel.cursors.push(Rc::downgrade(&cursor));
194 channel.maybe_prune_readers();
195 Ok(EventReader {
196 owner,
197 event_id,
198 cursor,
199 last_payload: None,
200 _marker: PhantomData,
201 })
202 }
203
204 pub fn fork_reader<E: Clone + 'static>(
205 &mut self,
206 owner: &WorldOwner,
207 reader: &EventReader<E>,
208 ) -> Result<EventReader<E>, WorldError> {
209 if let Err(error) = reader.validate_owner(owner) {
210 return Err(map_read_owner_error(error));
211 }
212 reader
213 .event_id
214 .validate_owner(owner)
215 .map_err(map_registration_owner_error)?;
216 let channel = self
217 .channels
218 .get_mut(reader.event_id.index())
219 .ok_or_else(|| WorldError::UnregisteredEvent {
220 name: alloc::format!("event {}", reader.event_id.index()),
221 })?;
222 let cursor = Rc::new(Cell::new(reader.cursor.get()));
223 channel.cursors.push(Rc::downgrade(&cursor));
224 channel.maybe_prune_readers();
225 Ok(EventReader {
226 owner: reader.owner.clone(),
227 event_id: reader.event_id.clone(),
228 cursor,
229 last_payload: None,
230 _marker: PhantomData,
231 })
232 }
233
234 #[cfg(test)]
235 pub(crate) fn clear_channels_for_test(&mut self) {
236 self.channels.clear();
237 }
238
239 #[cfg(test)]
240 pub(crate) fn set_channel_state_for_test(
241 &mut self,
242 index: usize,
243 next_sequence: u64,
244 closed: bool,
245 ) {
246 if let Some(channel) = self.channels.get_mut(index) {
247 if channel.active_len == 0 {
248 channel.oldest_retained = next_sequence;
249 for cursor in channel.cursors.iter().filter_map(Weak::upgrade) {
250 cursor.set(next_sequence);
251 }
252 }
253 channel.next_sequence = next_sequence;
254 channel.closed = closed;
255 }
256 }
257
258 pub fn clear_frame(&mut self, operation: StageOperation) {
259 for channel in &mut self.channels {
260 if matches!(channel.retention, EventRetention::Frame(owner) if owner == operation) {
261 channel.recycle_payloads();
262 channel.oldest_retained = channel.next_sequence;
263 channel.prune_readers();
264 }
265 }
266 }
267}
268
269impl EventChannel {
270 fn new(retention: EventRetention) -> Self {
271 Self {
272 entries: EventEntries::new(retention),
273 active_len: 0,
274 next_sequence: 0,
275 oldest_retained: 0,
276 retention,
277 cursors: Vec::new(),
278 reader_ops_since_prune: 0,
279 closed: false,
280 }
281 }
282
283 fn enforce_retention(&mut self) {
284 match self.retention {
285 EventRetention::Bounded(capacity) => {
286 while self.active_len > capacity {
287 self.entries.recycle_oldest();
288 self.active_len -= 1;
289 }
290 self.refresh_oldest_retained();
291 }
292 EventRetention::Frame(_) | EventRetention::Manual => {}
293 }
294 }
295
296 fn push_event<E: 'static>(&mut self, sequence: u64, event: E) {
297 let overwrite_oldest = matches!(
298 self.retention,
299 EventRetention::Bounded(capacity)
300 if capacity != 0
301 && capacity <= LINEAR_BOUNDED_CAPACITY_MAX
302 && self.active_len == capacity
303 );
304 if overwrite_oldest {
305 self.entries
306 .overwrite_oldest_linear(self.active_len, sequence, event);
307 return;
308 }
309 if let Some(entry) = self.entries.get_mut(self.active_len) {
310 if let Some(slot) = entry.payload.downcast_mut::<E>() {
311 *slot = event;
312 entry.sequence = sequence;
313 self.active_len += 1;
314 return;
315 }
316 entry.payload = Box::new(event);
317 entry.sequence = sequence;
318 } else {
319 self.entries.push(EventEntry {
320 sequence,
321 payload: Box::new(event),
322 });
323 }
324 self.active_len += 1;
325 }
326
327 fn recycle_payloads(&mut self) {
328 self.active_len = 0;
329 }
330
331 fn maybe_prune_readers(&mut self) {
332 self.reader_ops_since_prune = self.reader_ops_since_prune.saturating_add(1);
333 if self.reader_ops_since_prune >= READER_PRUNE_INTERVAL {
334 self.prune_readers();
335 }
336 }
337
338 fn prune_readers(&mut self) {
339 let mut index = 0;
340 while index < self.cursors.len() {
341 if self.cursors[index].strong_count() == 0 {
342 self.cursors.swap_remove(index);
343 } else {
344 index += 1;
345 }
346 }
347 self.reader_ops_since_prune = 0;
348 }
349
350 fn refresh_oldest_retained(&mut self) {
351 self.oldest_retained = self
352 .first_active()
353 .map(|entry| entry.sequence.saturating_sub(1))
354 .unwrap_or(self.next_sequence);
355 }
356
357 fn first_active(&self) -> Option<&EventEntry> {
358 (self.active_len != 0)
359 .then(|| self.entries.first())
360 .flatten()
361 }
362
363 fn position_after(&self, sequence: u64) -> Option<usize> {
364 let offset = sequence.checked_sub(self.oldest_retained)?;
365 let position = usize::try_from(offset).ok()?;
366 (position < self.active_len).then_some(position)
367 }
368
369 fn set_retention(&mut self, retention: EventRetention) {
370 self.entries.reconfigure(retention);
371 self.retention = retention;
372 }
373}
374
375impl EventEntries {
376 fn new(retention: EventRetention) -> Self {
377 if Self::uses_ring(retention) {
378 Self::Ring(VecDeque::with_capacity(16))
379 } else {
380 Self::Linear(Vec::with_capacity(16))
381 }
382 }
383
384 fn uses_ring(retention: EventRetention) -> bool {
385 matches!(
386 retention,
387 EventRetention::Bounded(capacity) if capacity > LINEAR_BOUNDED_CAPACITY_MAX
388 )
389 }
390
391 fn reconfigure(&mut self, retention: EventRetention) {
392 let wants_ring = Self::uses_ring(retention);
393 if matches!(self, Self::Ring(_)) == wants_ring {
394 return;
395 }
396 *self = match mem::replace(self, Self::Linear(Vec::new())) {
397 Self::Linear(entries) => Self::Ring(VecDeque::from(entries)),
398 Self::Ring(entries) => Self::Linear(entries.into_iter().collect()),
399 };
400 }
401
402 #[cfg(test)]
403 fn len(&self) -> usize {
404 match self {
405 Self::Linear(entries) => entries.len(),
406 Self::Ring(entries) => entries.len(),
407 }
408 }
409
410 fn first(&self) -> Option<&EventEntry> {
411 match self {
412 Self::Linear(entries) => entries.first(),
413 Self::Ring(entries) => entries.front(),
414 }
415 }
416
417 fn get(&self, index: usize) -> Option<&EventEntry> {
418 match self {
419 Self::Linear(entries) => entries.get(index),
420 Self::Ring(entries) => entries.get(index),
421 }
422 }
423
424 fn get_mut(&mut self, index: usize) -> Option<&mut EventEntry> {
425 match self {
426 Self::Linear(entries) => entries.get_mut(index),
427 Self::Ring(entries) => entries.get_mut(index),
428 }
429 }
430
431 fn push(&mut self, entry: EventEntry) {
432 match self {
433 Self::Linear(entries) => entries.push(entry),
434 Self::Ring(entries) => entries.push_back(entry),
435 }
436 }
437
438 fn recycle_oldest(&mut self) {
439 match self {
440 Self::Linear(entries) => {
441 let entry = entries.remove(0);
442 entries.push(entry);
443 }
444 Self::Ring(entries) => entries.rotate_left(1),
445 }
446 }
447
448 fn overwrite_oldest_linear<E: 'static>(&mut self, active_len: usize, sequence: u64, event: E) {
449 let Self::Linear(entries) = self else {
450 unreachable!("small bounded event channels use linear storage");
451 };
452 entries[..active_len].rotate_left(1);
453 let entry = &mut entries[active_len - 1];
454 if let Some(slot) = entry.payload.downcast_mut::<E>() {
455 *slot = event;
456 } else {
457 entry.payload = Box::new(event);
458 }
459 entry.sequence = sequence;
460 }
461}
462
463impl<E: Clone + 'static> EventReader<E> {
464 pub fn fork(&mut self, world: &mut crate::world::World) -> Result<Self, WorldError> {
466 world.fork_event_reader(self)
467 }
468
469 pub(crate) fn validate_owner(&self, owner: &WorldOwner) -> Result<(), EventReadError> {
470 if self.owner.same(owner) {
471 Ok(())
472 } else {
473 Err(EventReadError::OwnerMismatch {
474 name: alloc::string::String::from("event reader"),
475 })
476 }
477 }
478}
479
480fn map_read_owner_error(error: EventReadError) -> WorldError {
481 match error {
482 EventReadError::OwnerMismatch { name } => WorldError::UnregisteredEvent { name },
483 EventReadError::UnregisteredEvent { name } => WorldError::UnregisteredEvent { name },
484 EventReadError::Lagged { .. } | EventReadError::ChannelClosed => {
485 WorldError::UnregisteredEvent {
486 name: alloc::string::String::from("invalid reader state"),
487 }
488 }
489 }
490}
491
492fn map_registration_owner_error(
493 error: crate::event::registry::EventRegistrationError,
494) -> WorldError {
495 match error {
496 crate::event::registry::EventRegistrationError::TypeConflict { name, .. } => {
497 WorldError::UnregisteredEvent { name }
498 }
499 crate::event::registry::EventRegistrationError::InvalidCapacity => {
500 WorldError::UnregisteredEvent {
501 name: alloc::string::String::from("invalid event capacity"),
502 }
503 }
504 }
505}
506
507#[cfg(test)]
508mod tests {
509 use alloc::string::String;
510
511 use super::*;
512 use crate::event::EventOptions;
513 use crate::world::WorldBuilder;
514
515 #[derive(Clone, Debug, PartialEq)]
516 struct Damage(u32);
517
518 #[derive(Clone, Debug, PartialEq)]
519 struct Other(u32);
520
521 #[test]
522 fn storage_send_read_fork_and_map_errors() {
523 let mut builder = WorldBuilder::new();
524 let event_id = builder
525 .add_event::<Damage>(EventOptions::manual())
526 .expect("register");
527 let owner = builder.owner_for_test();
528 let mut storage = EventStorage::new(1);
529 storage.ensure_channel(event_id.index(), EventRetention::Manual);
530
531 assert!(matches!(
532 storage.send(&EventId::new(owner.clone(), 99), Damage(1)),
533 Err(WorldError::UnregisteredEvent { .. })
534 ));
535
536 storage.send(&event_id, Damage(1)).expect("one");
537
538 let mut wrong = storage
539 .create_reader::<Other>(
540 owner.clone(),
541 event_id.clone(),
542 EventReaderStart::OldestRetained,
543 )
544 .expect("wrong reader");
545 assert!(matches!(
546 storage.read_next(&owner, &mut wrong),
547 Err(EventReadError::UnregisteredEvent { .. })
548 ));
549
550 let mut reader = storage
551 .create_reader::<Damage>(
552 owner.clone(),
553 event_id.clone(),
554 EventReaderStart::OldestRetained,
555 )
556 .expect("reader");
557 assert!(storage
558 .read_next(&owner, &mut reader)
559 .expect("read")
560 .is_some());
561
562 let other_owner = WorldOwner::new();
563 assert!(matches!(
564 storage.fork_reader(&other_owner, &reader),
565 Err(WorldError::UnregisteredEvent { .. })
566 ));
567 assert!(matches!(
568 map_read_owner_error(EventReadError::OwnerMismatch {
569 name: String::from("reader")
570 }),
571 WorldError::UnregisteredEvent { .. }
572 ));
573 assert!(matches!(
574 map_registration_owner_error(
575 crate::event::registry::EventRegistrationError::TypeConflict {
576 name: String::from("Damage"),
577 existing: String::from("a"),
578 requested: String::from("b"),
579 }
580 ),
581 WorldError::UnregisteredEvent { .. }
582 ));
583 assert!(matches!(
584 map_registration_owner_error(
585 crate::event::registry::EventRegistrationError::InvalidCapacity
586 ),
587 WorldError::UnregisteredEvent { .. }
588 ));
589 }
590
591 #[test]
592 fn dropping_last_reader_does_not_clear_frame_payloads() {
593 let mut builder = WorldBuilder::new();
594 let event_id = builder
595 .add_event::<Damage>(EventOptions::frame(StageOperation::Update))
596 .expect("register");
597 let owner = builder.owner_for_test();
598 let mut storage = EventStorage::new(1);
599 storage.ensure_channel(
600 event_id.index(),
601 EventRetention::Frame(StageOperation::Update),
602 );
603 storage.send(&event_id, Damage(1)).expect("send");
604 {
605 let _reader = storage
606 .create_reader::<Damage>(
607 owner.clone(),
608 event_id.clone(),
609 EventReaderStart::OldestRetained,
610 )
611 .expect("reader");
612 }
613 storage
614 .send(&event_id, Damage(2))
615 .expect("send after reader drop");
616 let mut reader = storage
617 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
618 .expect("late");
619 for expected in [1, 2] {
620 assert_eq!(
621 storage
622 .read_next(&owner, &mut reader)
623 .expect("read")
624 .map(|d| d.0),
625 Some(expected)
626 );
627 }
628 }
629
630 #[test]
631 fn send_rejects_closed_channel_and_recycles_wrong_payload_type() {
632 let mut builder = WorldBuilder::new();
633 let event_id = builder
634 .add_event::<Damage>(EventOptions::manual())
635 .expect("register");
636 let mut storage = EventStorage::new(1);
637 storage.ensure_channel(event_id.index(), EventRetention::Manual);
638 storage.send(&event_id, Other(9)).expect("warm pool");
639 storage.send(&event_id, Damage(1)).expect("typed reuse");
640 storage.channels[event_id.index()].closed = true;
641 assert!(matches!(
642 storage.send(&event_id, Damage(2)),
643 Err(WorldError::EventChannelClosed)
644 ));
645 }
646
647 #[test]
648 fn create_reader_and_read_next_reject_unregistered_channel() {
649 let owner = WorldOwner::new();
650 let mut storage = EventStorage::new(0);
651 let bogus = EventId::new(owner.clone(), 0);
652 assert!(matches!(
653 storage.create_reader::<Damage>(
654 owner.clone(),
655 bogus.clone(),
656 EventReaderStart::OldestRetained
657 ),
658 Err(WorldError::UnregisteredEvent { .. })
659 ));
660 let mut reader = EventReader::<Damage> {
661 owner: owner.clone(),
662 event_id: bogus,
663 cursor: Rc::new(Cell::new(0)),
664 last_payload: None,
665 _marker: PhantomData,
666 };
667 assert!(matches!(
668 storage.read_next(&owner, &mut reader),
669 Err(EventReadError::UnregisteredEvent { .. })
670 ));
671 }
672
673 #[test]
674 fn map_read_owner_error_covers_lagged_and_closed() {
675 assert!(matches!(
676 map_read_owner_error(EventReadError::Lagged { dropped: 2 }),
677 WorldError::UnregisteredEvent { .. }
678 ));
679 assert!(matches!(
680 map_read_owner_error(EventReadError::ChannelClosed),
681 WorldError::UnregisteredEvent { .. }
682 ));
683 }
684
685 #[test]
686 fn read_next_rejects_event_id_owner_mismatch() {
687 let mut builder = WorldBuilder::new();
688 let event_id = builder
689 .add_event::<Damage>(EventOptions::manual())
690 .expect("register");
691 let owner = builder.owner_for_test();
692 let mut storage = EventStorage::new(1);
693 storage.ensure_channel(event_id.index(), EventRetention::Manual);
694 storage.send(&event_id, Damage(1)).expect("send");
695
696 let mut reader = storage
697 .create_reader::<Damage>(
698 owner.clone(),
699 event_id.clone(),
700 EventReaderStart::OldestRetained,
701 )
702 .expect("reader");
703 reader.event_id = EventId::new(WorldOwner::new(), event_id.index() as u32);
704 assert!(matches!(
705 storage.read_next(&owner, &mut reader),
706 Err(EventReadError::OwnerMismatch { name })
707 if name == alloc::format!("event {}", event_id.index())
708 ));
709 }
710
711 #[test]
712 fn storage_reuse_and_entry_variants_cover_all_layout_paths() {
713 let mut builder = WorldBuilder::new();
714 let event_id = builder
715 .add_event::<Damage>(EventOptions::frame(StageOperation::Update))
716 .expect("register");
717 let owner = builder.owner_for_test();
718 let mut storage = EventStorage::new(1);
719 storage.ensure_channel(
720 event_id.index(),
721 EventRetention::Frame(StageOperation::Update),
722 );
723 storage.send(&event_id, Damage(1)).expect("seed");
724 storage.set_channel_state_for_test(event_id.index(), 6, false);
725 storage.clear_frame(StageOperation::Update);
726 storage
727 .send(&event_id, Other(2))
728 .expect("replace recycled type");
729
730 let mut reader = storage
731 .create_reader::<Other>(
732 owner.clone(),
733 event_id.clone(),
734 EventReaderStart::OldestRetained,
735 )
736 .expect("reader");
737 reader.last_payload = Some(Box::new(Damage(9)));
738 assert_eq!(
739 storage
740 .read_next(&owner, &mut reader)
741 .expect("read")
742 .map(|event| event.0),
743 Some(2)
744 );
745
746 storage.clear_frame(StageOperation::Update);
747 storage.set_channel_state_for_test(event_id.index(), 7, false);
748 assert_eq!(storage.channels[0].oldest_retained, 7);
749 assert_eq!(reader.cursor.get(), 7);
750 storage.set_channel_state_for_test(99, 8, false);
751
752 let mut channel = EventChannel::new(EventRetention::Bounded(2));
753 channel.push_event(0, Damage(1));
754 channel.push_event(1, Damage(2));
755 channel.push_event(2, Other(3));
756 assert!(channel.position_after(u64::MAX).is_none());
757
758 let mut ring = EventEntries::new(EventRetention::Bounded(LINEAR_BOUNDED_CAPACITY_MAX + 1));
759 ring.push(EventEntry {
760 sequence: 0,
761 payload: Box::new(Damage(1)),
762 });
763 assert!(ring.get_mut(0).is_some());
764 ring.recycle_oldest();
765
766 let mut linear = EventEntries::new(EventRetention::Manual);
767 linear.push(EventEntry {
768 sequence: 0,
769 payload: Box::new(Damage(4)),
770 });
771 linear.recycle_oldest();
772 }
773
774 #[test]
775 #[should_panic(expected = "small bounded event channels use linear storage")]
776 fn linear_overwrite_invariant_rejects_ring_storage() {
777 let mut ring = EventEntries::new(EventRetention::Bounded(LINEAR_BOUNDED_CAPACITY_MAX + 1));
778 ring.overwrite_oldest_linear(1, 1, Damage(1));
779 }
780
781 #[test]
782 fn fork_reader_rejects_unregistered_channel() {
783 let mut builder = WorldBuilder::new();
784 let event_id = builder
785 .add_event::<Damage>(EventOptions::manual())
786 .expect("register");
787 let owner = builder.owner_for_test();
788 let mut storage = EventStorage::new(1);
789 storage.ensure_channel(event_id.index(), EventRetention::Manual);
790 let reader = storage
791 .create_reader::<Damage>(
792 owner.clone(),
793 event_id.clone(),
794 EventReaderStart::OldestRetained,
795 )
796 .expect("reader");
797 storage.clear_channels_for_test();
798 assert!(matches!(
799 storage.fork_reader(&owner, &reader),
800 Err(WorldError::UnregisteredEvent { name })
801 if name == alloc::format!("event {}", event_id.index())
802 ));
803 }
804
805 #[test]
806 fn reader_progress_does_not_clear_frame_payloads() {
807 let mut builder = WorldBuilder::new();
808 let event_id = builder
809 .add_event::<Damage>(EventOptions::frame(StageOperation::Update))
810 .expect("register");
811 let owner = builder.owner_for_test();
812 let mut storage = EventStorage::new(1);
813 storage.ensure_channel(
814 event_id.index(),
815 EventRetention::Frame(StageOperation::Update),
816 );
817 storage.send(&event_id, Damage(1)).expect("one");
818 storage.send(&event_id, Damage(2)).expect("two");
819 storage.send(&event_id, Damage(3)).expect("three");
820 let reader = storage
821 .create_reader::<Damage>(
822 owner.clone(),
823 event_id.clone(),
824 EventReaderStart::OldestRetained,
825 )
826 .expect("reader");
827 reader.cursor.set(3);
828 storage.send(&event_id, Damage(4)).expect("four");
829 let mut reader = storage
830 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
831 .expect("fresh");
832 for expected in [1, 2, 3, 4] {
833 assert_eq!(
834 storage
835 .read_next(&owner, &mut reader)
836 .expect("read")
837 .map(|d| d.0),
838 Some(expected)
839 );
840 }
841 }
842
843 #[test]
844 fn map_read_owner_error_covers_unregistered_event() {
845 assert!(matches!(
846 map_read_owner_error(EventReadError::UnregisteredEvent {
847 name: String::from("Damage")
848 }),
849 WorldError::UnregisteredEvent { name }
850 if name == "Damage"
851 ));
852 }
853
854 #[test]
855 fn two_readers_both_observe_same_manual_event() {
856 let mut builder = WorldBuilder::new();
857 let event_id = builder
858 .add_event::<Damage>(EventOptions::manual())
859 .expect("register");
860 let owner = builder.owner_for_test();
861 let mut storage = EventStorage::new(1);
862 storage.ensure_channel(event_id.index(), EventRetention::Manual);
863 storage.send(&event_id, Damage(42)).expect("send");
864
865 let mut reader_a = storage
866 .create_reader::<Damage>(
867 owner.clone(),
868 event_id.clone(),
869 EventReaderStart::OldestRetained,
870 )
871 .expect("reader a");
872 let mut reader_b = storage
873 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
874 .expect("reader b");
875
876 assert_eq!(
877 storage
878 .read_next(&owner, &mut reader_a)
879 .expect("read a")
880 .map(|d| d.0),
881 Some(42)
882 );
883 assert_eq!(
884 storage
885 .read_next(&owner, &mut reader_b)
886 .expect("read b")
887 .map(|d| d.0),
888 Some(42)
889 );
890 }
891
892 #[test]
893 fn forked_reader_replays_independently_of_parent() {
894 let mut builder = WorldBuilder::new();
895 let event_id = builder
896 .add_event::<Damage>(EventOptions::manual())
897 .expect("register");
898 let owner = builder.owner_for_test();
899 let mut storage = EventStorage::new(1);
900 storage.ensure_channel(event_id.index(), EventRetention::Manual);
901 storage.send(&event_id, Damage(1)).expect("one");
902 storage.send(&event_id, Damage(2)).expect("two");
903
904 let mut parent = storage
905 .create_reader::<Damage>(
906 owner.clone(),
907 event_id.clone(),
908 EventReaderStart::OldestRetained,
909 )
910 .expect("parent");
911 let mut fork = storage.fork_reader(&owner, &parent).expect("fork");
912 let _ = storage
913 .read_next(&owner, &mut parent)
914 .expect("parent consumes one");
915
916 assert_eq!(
917 storage
918 .read_next(&owner, &mut fork)
919 .expect("fork first")
920 .map(|d| d.0),
921 Some(1)
922 );
923 assert_eq!(
924 storage
925 .read_next(&owner, &mut fork)
926 .expect("fork second")
927 .map(|d| d.0),
928 Some(2)
929 );
930 assert_eq!(
931 storage
932 .read_next(&owner, &mut parent)
933 .expect("parent second")
934 .map(|d| d.0),
935 Some(2)
936 );
937 }
938
939 #[test]
940 fn late_oldest_retained_reader_receives_all_manual_events_in_order() {
941 let mut builder = WorldBuilder::new();
942 let event_id = builder
943 .add_event::<Damage>(EventOptions::manual())
944 .expect("register");
945 let owner = builder.owner_for_test();
946 let mut storage = EventStorage::new(1);
947 storage.ensure_channel(event_id.index(), EventRetention::Manual);
948 storage.send(&event_id, Damage(10)).expect("one");
949 storage.send(&event_id, Damage(20)).expect("two");
950
951 let mut late = storage
952 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
953 .expect("late");
954 assert_eq!(
955 storage
956 .read_next(&owner, &mut late)
957 .expect("first")
958 .map(|d| d.0),
959 Some(10)
960 );
961 assert_eq!(
962 storage
963 .read_next(&owner, &mut late)
964 .expect("second")
965 .map(|d| d.0),
966 Some(20)
967 );
968 }
969
970 #[test]
971 fn bounded_retention_without_reader_keeps_latest_events() {
972 let mut builder = WorldBuilder::new();
973 let event_id = builder
974 .add_event::<Damage>(EventOptions::bounded(2).expect("bounded"))
975 .expect("register");
976 let owner = builder.owner_for_test();
977 let mut storage = EventStorage::new(1);
978 storage.ensure_channel(event_id.index(), EventRetention::Bounded(2));
979 storage.send(&event_id, Damage(1)).expect("one");
980 storage.send(&event_id, Damage(2)).expect("two");
981 storage.send(&event_id, Damage(3)).expect("three");
982 let channel = &storage.channels[event_id.index()];
983 assert_eq!(channel.active_len, 2);
984 assert_eq!(channel.entries.len(), 2);
985
986 let mut late = storage
987 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
988 .expect("late");
989 assert_eq!(
990 storage
991 .read_next(&owner, &mut late)
992 .expect("second retained")
993 .map(|d| d.0),
994 Some(2)
995 );
996 assert_eq!(
997 storage
998 .read_next(&owner, &mut late)
999 .expect("third retained")
1000 .map(|d| d.0),
1001 Some(3)
1002 );
1003 assert!(storage
1004 .read_next(&owner, &mut late)
1005 .expect("drain")
1006 .is_none());
1007 }
1008
1009 #[test]
1010 fn bounded_retention_reports_lag_for_slow_reader() {
1011 let mut builder = WorldBuilder::new();
1012 let event_id = builder
1013 .add_event::<Damage>(EventOptions::bounded(1).expect("bounded"))
1014 .expect("register");
1015 let owner = builder.owner_for_test();
1016 let mut storage = EventStorage::new(1);
1017 storage.ensure_channel(event_id.index(), EventRetention::Bounded(1));
1018 let mut reader = storage
1019 .create_reader::<Damage>(
1020 owner.clone(),
1021 event_id.clone(),
1022 EventReaderStart::OldestRetained,
1023 )
1024 .expect("reader");
1025 storage.send(&event_id, Damage(1)).expect("one");
1026 storage.send(&event_id, Damage(2)).expect("two");
1027 assert!(matches!(
1028 storage.read_next(&owner, &mut reader),
1029 Err(EventReadError::Lagged { dropped: 1 })
1030 ));
1031 }
1032
1033 fn assert_ring_overflow_order_and_lag(capacity: usize, total: usize) {
1034 let mut builder = WorldBuilder::new();
1035 let event_id = builder
1036 .add_event::<Damage>(EventOptions::bounded(capacity).expect("bounded"))
1037 .expect("register");
1038 let owner = builder.owner_for_test();
1039 let mut storage = EventStorage::new(1);
1040 storage.ensure_channel(event_id.index(), EventRetention::Bounded(capacity));
1041 let mut slow = storage
1042 .create_reader::<Damage>(
1043 owner.clone(),
1044 event_id.clone(),
1045 EventReaderStart::OldestRetained,
1046 )
1047 .expect("slow reader");
1048
1049 for value in 0..total {
1050 storage
1051 .send(&event_id, Damage(value as u32))
1052 .expect("overflowing send");
1053 }
1054
1055 let channel = &storage.channels[event_id.index()];
1056 assert!(matches!(&channel.entries, EventEntries::Ring(_)));
1057 assert_eq!(channel.active_len, capacity);
1058 assert_eq!(channel.entries.len(), capacity + 1);
1059 let dropped = (total - capacity) as u64;
1060 assert_eq!(
1061 storage.read_next(&owner, &mut slow),
1062 Err(EventReadError::Lagged { dropped })
1063 );
1064 assert_eq!(
1065 storage
1066 .read_next(&owner, &mut slow)
1067 .expect("slow reader catches up")
1068 .map(|event| event.0),
1069 Some((total - capacity) as u32)
1070 );
1071
1072 let mut oldest = storage
1073 .create_reader::<Damage>(owner.clone(), event_id, EventReaderStart::OldestRetained)
1074 .expect("oldest retained reader");
1075 for expected in (total - capacity)..total {
1076 assert_eq!(
1077 storage
1078 .read_next(&owner, &mut oldest)
1079 .expect("retained read")
1080 .map(|event| event.0),
1081 Some(expected as u32)
1082 );
1083 }
1084 assert!(storage
1085 .read_next(&owner, &mut oldest)
1086 .expect("retained events drained")
1087 .is_none());
1088 }
1089
1090 #[test]
1091 fn ring_capacity_17_preserves_order_and_exact_lag_after_repeated_overflow() {
1092 assert_ring_overflow_order_and_lag(17, 100);
1093 }
1094
1095 #[test]
1096 fn ring_capacity_256_preserves_order_and_exact_lag_after_repeated_overflow() {
1097 assert_ring_overflow_order_and_lag(256, 1_024);
1098 }
1099
1100 #[test]
1101 fn frame_clear_reports_lag_and_resets_oldest_reader_start() {
1102 let mut builder = WorldBuilder::new();
1103 let event_id = builder
1104 .add_event::<Damage>(EventOptions::frame(StageOperation::Update))
1105 .expect("register");
1106 let owner = builder.owner_for_test();
1107 let mut storage = EventStorage::new(1);
1108 storage.ensure_channel(
1109 event_id.index(),
1110 EventRetention::Frame(StageOperation::Update),
1111 );
1112 let mut existing = storage
1113 .create_reader::<Damage>(owner.clone(), event_id.clone(), EventReaderStart::FromNow)
1114 .expect("existing");
1115 storage.send(&event_id, Damage(1)).expect("one");
1116 storage.send(&event_id, Damage(2)).expect("two");
1117 storage.clear_frame(StageOperation::Update);
1118
1119 assert!(matches!(
1120 storage.read_next(&owner, &mut existing),
1121 Err(EventReadError::Lagged { dropped: 2 })
1122 ));
1123 assert!(storage
1124 .read_next(&owner, &mut existing)
1125 .expect("caught up")
1126 .is_none());
1127
1128 let mut late = storage
1129 .create_reader::<Damage>(
1130 owner.clone(),
1131 event_id.clone(),
1132 EventReaderStart::OldestRetained,
1133 )
1134 .expect("late");
1135 assert!(storage
1136 .read_next(&owner, &mut late)
1137 .expect("starts at boundary")
1138 .is_none());
1139 storage.send(&event_id, Damage(3)).expect("three");
1140 assert_eq!(
1141 storage
1142 .read_next(&owner, &mut late)
1143 .expect("new frame")
1144 .map(|event| event.0),
1145 Some(3)
1146 );
1147 }
1148
1149 #[test]
1150 fn frame_clear_keeps_entry_slots_for_reuse() {
1151 let mut builder = WorldBuilder::new();
1152 let event_id = builder
1153 .add_event::<Damage>(EventOptions::frame(StageOperation::Update))
1154 .expect("register");
1155 let mut storage = EventStorage::new(1);
1156 storage.ensure_channel(
1157 event_id.index(),
1158 EventRetention::Frame(StageOperation::Update),
1159 );
1160 for value in 0..4 {
1161 storage.send(&event_id, Damage(value)).expect("send");
1162 }
1163 let channel = &storage.channels[event_id.index()];
1164 assert_eq!(channel.active_len, 4);
1165 assert_eq!(channel.entries.len(), 4);
1166
1167 storage.clear_frame(StageOperation::Update);
1168 let channel = &storage.channels[event_id.index()];
1169 assert_eq!(channel.active_len, 0);
1170 assert_eq!(channel.entries.len(), 4);
1171
1172 storage.send(&event_id, Damage(9)).expect("reuse");
1173 let channel = &storage.channels[event_id.index()];
1174 assert_eq!(channel.active_len, 1);
1175 assert_eq!(channel.entries.len(), 4);
1176 assert_eq!(
1177 channel
1178 .entries
1179 .first()
1180 .and_then(|entry| entry.payload.downcast_ref::<Damage>()),
1181 Some(&Damage(9))
1182 );
1183 }
1184
1185 #[test]
1186 fn entry_storage_adapts_to_retention_without_losing_active_events() {
1187 let mut builder = WorldBuilder::new();
1188 let event_id = builder
1189 .add_event::<Damage>(EventOptions::manual())
1190 .expect("register");
1191 let mut storage = EventStorage::new(1);
1192 storage.ensure_channel(event_id.index(), EventRetention::Manual);
1193 storage.send(&event_id, Damage(1)).expect("first");
1194 storage.send(&event_id, Damage(2)).expect("second");
1195 assert!(matches!(
1196 &storage.channels[event_id.index()].entries,
1197 EventEntries::Linear(_)
1198 ));
1199
1200 storage.ensure_channel(
1201 event_id.index(),
1202 EventRetention::Bounded(LINEAR_BOUNDED_CAPACITY_MAX + 1),
1203 );
1204 let channel = &storage.channels[event_id.index()];
1205 assert!(matches!(&channel.entries, EventEntries::Ring(_)));
1206 assert_eq!(channel.active_len, 2);
1207 assert_eq!(
1208 channel
1209 .entries
1210 .get(0)
1211 .and_then(|entry| entry.payload.downcast_ref::<Damage>()),
1212 Some(&Damage(1))
1213 );
1214
1215 storage.ensure_channel(
1216 event_id.index(),
1217 EventRetention::Bounded(LINEAR_BOUNDED_CAPACITY_MAX),
1218 );
1219 let channel = &storage.channels[event_id.index()];
1220 assert!(matches!(&channel.entries, EventEntries::Linear(_)));
1221 assert_eq!(channel.active_len, 2);
1222 assert_eq!(
1223 channel
1224 .entries
1225 .get(1)
1226 .and_then(|entry| entry.payload.downcast_ref::<Damage>()),
1227 Some(&Damage(2))
1228 );
1229 }
1230
1231 #[test]
1232 fn dropped_reader_slots_are_pruned_within_the_operation_budget() {
1233 let mut builder = WorldBuilder::new();
1234 let event_id = builder
1235 .add_event::<Damage>(EventOptions::manual())
1236 .expect("register");
1237 let owner = builder.owner_for_test();
1238 let mut storage = EventStorage::new(1);
1239 storage.ensure_channel(event_id.index(), EventRetention::Manual);
1240
1241 for _ in 0..usize::from(READER_PRUNE_INTERVAL) {
1242 let reader = storage
1243 .create_reader::<Damage>(owner.clone(), event_id.clone(), EventReaderStart::FromNow)
1244 .expect("reader");
1245 drop(reader);
1246 }
1247 assert!(storage.channels[event_id.index()].cursors.len() <= 1);
1248
1249 for _ in 1..usize::from(READER_PRUNE_INTERVAL) {
1250 let reader = storage
1251 .create_reader::<Damage>(owner.clone(), event_id.clone(), EventReaderStart::FromNow)
1252 .expect("reader");
1253 drop(reader);
1254 }
1255 storage
1256 .send(&event_id, Damage(1))
1257 .expect("prune checkpoint");
1258 assert!(storage.channels[event_id.index()].cursors.is_empty());
1259 }
1260
1261 #[test]
1262 fn from_now_reader_skips_manual_history() {
1263 let mut builder = WorldBuilder::new();
1264 let event_id = builder
1265 .add_event::<Damage>(EventOptions::manual())
1266 .expect("register");
1267 let owner = builder.owner_for_test();
1268 let mut storage = EventStorage::new(1);
1269 storage.ensure_channel(event_id.index(), EventRetention::Manual);
1270 storage.send(&event_id, Damage(1)).expect("one");
1271 storage.send(&event_id, Damage(2)).expect("two");
1272 let mut slow = storage
1273 .create_reader::<Damage>(
1274 owner.clone(),
1275 event_id.clone(),
1276 EventReaderStart::OldestRetained,
1277 )
1278 .expect("slow");
1279 let _ = storage.read_next(&owner, &mut slow).expect("consume one");
1280 storage.send(&event_id, Damage(3)).expect("three");
1281 let mut fast = storage
1282 .create_reader::<Damage>(owner.clone(), event_id.clone(), EventReaderStart::FromNow)
1283 .expect("fast");
1284 storage.send(&event_id, Damage(4)).expect("four");
1285 assert_eq!(
1286 storage
1287 .read_next(&owner, &mut fast)
1288 .expect("read")
1289 .map(|d| d.0),
1290 Some(4)
1291 );
1292 }
1293}