1use crate::{
4 marshal::coding::types::CodedBlock,
5 types::{coding::Commitment, Round},
6 CertifiableBlock,
7};
8use commonware_actor::mailbox::{Overflow, Policy, Sender};
9use commonware_coding::Scheme as CodingScheme;
10use commonware_cryptography::{Hasher, PublicKey};
11use commonware_utils::channel::oneshot;
12use std::{collections::VecDeque, sync::Arc};
13
14pub(crate) enum Message<B, C, H, P>
18where
19 B: CertifiableBlock,
20 C: CodingScheme,
21 H: Hasher,
22 P: PublicKey,
23{
24 Proposed {
26 block: Arc<CodedBlock<B, C, H>>,
28 round: Round,
30 },
31 Discovered {
33 commitment: Commitment,
35 leader: P,
37 round: Round,
39 },
40 Notarized {
46 commitment: Commitment,
48 round: Round,
50 },
51 GetByCommitment {
53 commitment: Commitment,
55 response: oneshot::Sender<Option<Arc<CodedBlock<B, C, H>>>>,
57 },
58 GetByDigest {
60 digest: B::Digest,
62 response: oneshot::Sender<Option<Arc<CodedBlock<B, C, H>>>>,
64 },
65 SubscribeAssignedShardVerified {
76 commitment: Commitment,
78 response: oneshot::Sender<()>,
80 },
81 SubscribeByCommitment {
84 commitment: Commitment,
86 response: oneshot::Sender<Arc<CodedBlock<B, C, H>>>,
88 },
89 SubscribeByDigest {
92 digest: B::Digest,
94 response: oneshot::Sender<Arc<CodedBlock<B, C, H>>>,
96 },
97 Prune {
99 through: Commitment,
101 },
102}
103
104impl<B, C, H, P> Message<B, C, H, P>
105where
106 B: CertifiableBlock,
107 C: CodingScheme,
108 H: Hasher,
109 P: PublicKey,
110{
111 pub(crate) fn response_closed(&self) -> bool {
112 match self {
113 Self::GetByCommitment { response, .. } | Self::GetByDigest { response, .. } => {
114 response.is_closed()
115 }
116 Self::SubscribeAssignedShardVerified { response, .. } => response.is_closed(),
117 Self::SubscribeByCommitment { response, .. }
118 | Self::SubscribeByDigest { response, .. } => response.is_closed(),
119 Self::Proposed { .. }
120 | Self::Discovered { .. }
121 | Self::Notarized { .. }
122 | Self::Prune { .. } => false,
123 }
124 }
125}
126
127pub(crate) struct Pending<B, C, H, P>(VecDeque<Message<B, C, H, P>>)
128where
129 B: CertifiableBlock,
130 C: CodingScheme,
131 H: Hasher,
132 P: PublicKey;
133
134impl<B, C, H, P> Default for Pending<B, C, H, P>
135where
136 B: CertifiableBlock,
137 C: CodingScheme,
138 H: Hasher,
139 P: PublicKey,
140{
141 fn default() -> Self {
142 Self(VecDeque::new())
143 }
144}
145
146impl<B, C, H, P> Overflow<Message<B, C, H, P>> for Pending<B, C, H, P>
147where
148 B: CertifiableBlock,
149 C: CodingScheme,
150 H: Hasher,
151 P: PublicKey,
152{
153 fn is_empty(&self) -> bool {
154 self.0.is_empty()
155 }
156
157 fn drain<F>(&mut self, mut push: F)
158 where
159 F: FnMut(Message<B, C, H, P>) -> Option<Message<B, C, H, P>>,
160 {
161 while let Some(message) = self.0.pop_front() {
162 if message.response_closed() {
163 continue;
164 }
165
166 if let Some(message) = push(message) {
167 self.0.push_front(message);
168 break;
169 }
170 }
171 }
172}
173
174impl<B, C, H, P> Policy for Message<B, C, H, P>
175where
176 B: CertifiableBlock,
177 C: CodingScheme,
178 H: Hasher,
179 P: PublicKey,
180{
181 type Overflow = Pending<B, C, H, P>;
182
183 fn handle(overflow: &mut Self::Overflow, message: Self) {
184 if message.response_closed() {
185 return;
186 }
187
188 overflow.0.push_back(message);
189 }
190}
191
192#[derive(Clone)]
196pub struct Mailbox<B, C, H, P>
197where
198 B: CertifiableBlock,
199 C: CodingScheme,
200 H: Hasher,
201 P: PublicKey,
202{
203 pub(super) sender: Sender<Message<B, C, H, P>>,
204}
205
206impl<B, C, H, P> Mailbox<B, C, H, P>
207where
208 B: CertifiableBlock,
209 C: CodingScheme,
210 H: Hasher,
211 P: PublicKey,
212{
213 pub(crate) const fn new(sender: Sender<Message<B, C, H, P>>) -> Self {
215 Self { sender }
216 }
217
218 pub fn proposed(&self, round: Round, block: CodedBlock<B, C, H>) {
220 self.proposed_shared(round, Arc::new(block));
221 }
222
223 pub(crate) fn proposed_shared(&self, round: Round, block: Arc<CodedBlock<B, C, H>>) {
224 let _ = self.sender.enqueue(Message::Proposed { block, round });
225 }
226
227 pub fn discovered(&self, commitment: Commitment, leader: P, round: Round) {
229 let _ = self.sender.enqueue(Message::Discovered {
230 commitment,
231 leader,
232 round,
233 });
234 }
235
236 pub fn notarized(&self, commitment: Commitment, round: Round) {
243 let _ = self
244 .sender
245 .enqueue(Message::Notarized { commitment, round });
246 }
247
248 pub async fn get(&self, commitment: Commitment) -> Option<Arc<CodedBlock<B, C, H>>> {
250 let (response, receiver) = oneshot::channel();
251 let _ = self.sender.enqueue(Message::GetByCommitment {
252 commitment,
253 response,
254 });
255 receiver.await.ok().flatten()
256 }
257
258 pub async fn get_by_digest(&self, digest: B::Digest) -> Option<Arc<CodedBlock<B, C, H>>> {
260 let (response, receiver) = oneshot::channel();
261 let _ = self
262 .sender
263 .enqueue(Message::GetByDigest { digest, response });
264 receiver.await.ok().flatten()
265 }
266
267 pub fn subscribe_assigned_shard_verified(
278 &self,
279 commitment: Commitment,
280 ) -> oneshot::Receiver<()> {
281 let (responder, receiver) = oneshot::channel();
282 let _ = self
283 .sender
284 .enqueue(Message::SubscribeAssignedShardVerified {
285 commitment,
286 response: responder,
287 });
288 receiver
289 }
290
291 pub fn subscribe(&self, commitment: Commitment) -> oneshot::Receiver<Arc<CodedBlock<B, C, H>>> {
293 let (responder, receiver) = oneshot::channel();
294 let _ = self.sender.enqueue(Message::SubscribeByCommitment {
295 commitment,
296 response: responder,
297 });
298 receiver
299 }
300
301 pub fn subscribe_by_digest(
303 &self,
304 digest: B::Digest,
305 ) -> oneshot::Receiver<Arc<CodedBlock<B, C, H>>> {
306 let (responder, receiver) = oneshot::channel();
307 let _ = self.sender.enqueue(Message::SubscribeByDigest {
308 digest,
309 response: responder,
310 });
311 receiver
312 }
313
314 pub fn prune(&self, through: Commitment) {
316 let _ = self.sender.enqueue(Message::Prune { through });
317 }
318}