Skip to main content

rapier3d/data/
pubsub.rs

1//! Publish-subscribe mechanism for internal events.
2use crate::alloc_prelude::*;
3
4use alloc::collections::VecDeque;
5use core::marker::PhantomData;
6
7/// A permanent subscription to a pub-sub queue.
8#[cfg_attr(feature = "serde-serialize", derive(Serialize, Deserialize))]
9#[derive(Clone)]
10pub struct Subscription<T> {
11    // Position on the cursor array.
12    id: u32,
13    _phantom: PhantomData<T>,
14}
15
16#[cfg_attr(feature = "serde-serialize", derive(Serialize, Deserialize))]
17#[derive(Clone)]
18struct PubSubCursor {
19    // Position on the offset array.
20    id: u32,
21    // Index of the next message to read.
22    // NOTE: Having this here is not actually necessary because
23    // this value is supposed to be equal to `offsets[self.id]`.
24    // However, we keep it because it lets us avoid one lookup
25    // on the `offsets` array inside of message-polling loops
26    // based on `read_ith`.
27    next: u32,
28}
29
30impl PubSubCursor {
31    fn id(&self, num_deleted: u32) -> usize {
32        (self.id - num_deleted) as usize
33    }
34
35    fn next(&self, num_deleted: u32) -> usize {
36        (self.next - num_deleted) as usize
37    }
38}
39
40/// A pub-sub queue.
41#[cfg_attr(feature = "serde-serialize", derive(Serialize, Deserialize))]
42#[derive(Clone, Default)]
43pub struct PubSub<T> {
44    deleted_messages: u32,
45    deleted_offsets: u32,
46    messages: VecDeque<T>,
47    offsets: VecDeque<u32>,
48    cursors: Vec<PubSubCursor>,
49}
50
51impl<T> PubSub<T> {
52    /// Create a new empty pub-sub queue.
53    pub fn new() -> Self {
54        Self {
55            deleted_offsets: 0,
56            deleted_messages: 0,
57            messages: VecDeque::new(),
58            offsets: VecDeque::new(),
59            cursors: Vec::new(),
60        }
61    }
62
63    fn reset_shifts(&mut self) {
64        for offset in &mut self.offsets {
65            *offset -= self.deleted_messages;
66        }
67
68        for cursor in &mut self.cursors {
69            cursor.id -= self.deleted_offsets;
70            cursor.next -= self.deleted_messages;
71        }
72
73        self.deleted_offsets = 0;
74        self.deleted_messages = 0;
75    }
76
77    /// Publish a new message.
78    pub fn publish(&mut self, message: T) {
79        if self.offsets.is_empty() {
80            // No subscribers, drop the message.
81            return;
82        }
83
84        self.messages.push_back(message);
85    }
86
87    /// Subscribe to the queue.
88    ///
89    /// A subscription cannot be cancelled.
90    #[must_use]
91    pub fn subscribe(&mut self) -> Subscription<T> {
92        let cursor = PubSubCursor {
93            next: self.messages.len() as u32 + self.deleted_messages,
94            id: self.offsets.len() as u32 + self.deleted_offsets,
95        };
96
97        let subscription = Subscription {
98            id: self.cursors.len() as u32,
99            _phantom: PhantomData,
100        };
101
102        self.offsets.push_back(cursor.next);
103        self.cursors.push(cursor);
104        subscription
105    }
106
107    /// Read the i-th message not yet read by the given subscriber.
108    pub fn read_ith(&self, sub: &Subscription<T>, i: usize) -> Option<&T> {
109        let cursor = &self.cursors[sub.id as usize];
110        self.messages.get(cursor.next(self.deleted_messages) + i)
111    }
112
113    /// Get the messages not yet read by the given subscriber.
114    pub fn read(&self, sub: &Subscription<T>) -> impl Iterator<Item = &T> {
115        let cursor = &self.cursors[sub.id as usize];
116        let next = cursor.next(self.deleted_messages);
117
118        self.messages.range(next..)
119    }
120
121    /// Makes the given subscribe acknowledge all the messages in the queue.
122    ///
123    /// A subscriber cannot read acknowledged messages any more.
124    pub fn ack(&mut self, sub: &Subscription<T>) {
125        // Update the cursor.
126        let cursor = &mut self.cursors[sub.id as usize];
127
128        self.offsets[cursor.id(self.deleted_offsets)] = u32::MAX;
129        cursor.id = self.offsets.len() as u32 + self.deleted_offsets;
130
131        cursor.next = self.messages.len() as u32 + self.deleted_messages;
132        self.offsets.push_back(cursor.next);
133
134        // Now clear the messages we don't need to
135        // maintain in memory anymore.
136        while self.offsets.front() == Some(&u32::MAX) {
137            self.offsets.pop_front();
138            self.deleted_offsets += 1;
139        }
140
141        // There must be at least one offset otherwise
142        // that would mean we have no subscribers.
143        let next = self.offsets.front().unwrap();
144        let num_to_delete = *next - self.deleted_messages;
145
146        for _ in 0..num_to_delete {
147            self.messages.pop_front();
148        }
149
150        self.deleted_messages += num_to_delete;
151
152        if self.deleted_messages > u32::MAX / 2 || self.deleted_offsets > u32::MAX / 2 {
153            // Don't let the deleted_* shifts grow indefinitely otherwise
154            // they will end up overflowing, breaking everything.
155            self.reset_shifts();
156        }
157    }
158}