Skip to main content

moirai/event/
queue.rs

1//! Event channel storage, retention enforcement, and independent typed readers.
2//!
3//! [`EventStorage`] backs send/read/fork operations. [`EventReader`] tracks a private cursor and
4//! clones payloads so multiple readers can observe the same channel independently.
5
6use 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/// Explicit reader cursor policy when creating an [`EventReader`].
49#[derive(Copy, Clone, Debug, Eq, PartialEq)]
50pub enum EventReaderStart {
51    /// Begin at the oldest payload still retained by the channel.
52    OldestRetained,
53    /// Begin after the current channel tail, skipping prior history.
54    FromNow,
55}
56
57/// Independent typed event reader whose reads own cloned payloads.
58///
59/// Forked readers replay from their own cursor without advancing the parent.
60pub 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    /// Creates a sibling reader sharing this channel with an independent cursor.
465    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}