1use crate::alloc_prelude::*;
3
4use alloc::collections::VecDeque;
5use core::marker::PhantomData;
6
7#[cfg_attr(feature = "serde-serialize", derive(Serialize, Deserialize))]
9#[derive(Clone)]
10pub struct Subscription<T> {
11 id: u32,
13 _phantom: PhantomData<T>,
14}
15
16#[cfg_attr(feature = "serde-serialize", derive(Serialize, Deserialize))]
17#[derive(Clone)]
18struct PubSubCursor {
19 id: u32,
21 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#[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 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 pub fn publish(&mut self, message: T) {
79 if self.offsets.is_empty() {
80 return;
82 }
83
84 self.messages.push_back(message);
85 }
86
87 #[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 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 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 pub fn ack(&mut self, sub: &Subscription<T>) {
125 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 while self.offsets.front() == Some(&u32::MAX) {
137 self.offsets.pop_front();
138 self.deleted_offsets += 1;
139 }
140
141 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 self.reset_shifts();
156 }
157 }
158}