1#[cfg(test)]
4use behavior::Recipient;
5use behavior::{Actions, Address, Behavior, BehaviorActed, BehaviorBase, Never, NoBirths, User};
6use thiserror::Error;
7
8use crate::DeliveryRoute;
9
10#[derive(Debug, Clone, PartialEq, Eq)]
12pub enum AcknowledgementState<P> {
13 Pending {
15 remaining: Vec<P>,
17 acknowledged: Vec<P>,
19 },
20 Completed,
22 Cancelled,
24}
25
26#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct AcknowledgementRecord<K, P> {
29 pub key: K,
31 pub state: AcknowledgementState<P>,
33}
34
35#[derive(Debug, Clone, PartialEq, Eq)]
37pub enum AcknowledgementInput<K, P> {
38 Acknowledge { key: K, participant: P },
40 Cancel { key: K },
42}
43
44#[derive(Debug, Error, Clone, PartialEq, Eq)]
46pub enum AcknowledgementError<K, P> {
47 #[error("acknowledgement key already exists")]
49 Existing {
50 key: K,
52 participants: Vec<P>,
54 },
55 #[error("acknowledgement key is unknown")]
57 Unknown(AcknowledgementInput<K, P>),
58 #[error("participant is not required by this acknowledgement")]
60 UnexpectedParticipant {
61 key: K,
63 participant: P,
65 },
66 #[error("participant acknowledgement is a duplicate")]
68 DuplicateParticipant {
69 key: K,
71 participant: P,
73 },
74 #[error("acknowledgement lifecycle is already complete")]
76 Completed(AcknowledgementInput<K, P>),
77 #[error("acknowledgement lifecycle is cancelled")]
79 Cancelled(AcknowledgementInput<K, P>),
80}
81
82#[derive(Debug, Clone, PartialEq, Eq)]
84pub enum AcknowledgementOutcome<K, P> {
85 Started {
87 key: K,
89 remaining: usize,
91 },
92 Acknowledged {
94 key: K,
96 participant: P,
98 remaining: usize,
100 },
101 Completed {
103 key: K,
105 },
106 Cancelled {
108 key: K,
110 },
111 Rejected(AcknowledgementError<K, P>),
113}
114
115pub enum AcknowledgementMessage<K, P, Route> {
117 Begin {
119 key: K,
121 participants: Vec<P>,
123 reply_to: Route,
125 },
126 Acknowledge {
128 key: K,
130 participant: P,
132 reply_to: Route,
134 },
135 Cancel {
137 key: K,
139 reply_to: Route,
141 },
142}
143
144pub struct Acknowledgements<
158 A: Address,
159 K,
160 P,
161 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
162> {
163 records: Vec<AcknowledgementRecord<K, P>>,
164 marker: core::marker::PhantomData<fn() -> (A, Route)>,
165}
166
167impl<A, K, P, Route> Acknowledgements<A, K, P, Route>
168where
169 A: Address,
170 K: Clone + Eq,
171 P: Clone + Eq,
172 Route:
173 DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
174{
175 #[must_use]
177 pub const fn new() -> Self {
178 Self {
179 records: Vec::new(),
180 marker: core::marker::PhantomData,
181 }
182 }
183
184 #[must_use]
186 pub fn records(&self) -> &[AcknowledgementRecord<K, P>] {
187 &self.records
188 }
189
190 fn result(
191 reply_to: Route,
192 outcome: AcknowledgementOutcome<K, P>,
193 ) -> Actions<A, Never, Route::Sends, NoBirths> {
194 Actions::send(reply_to.deliver(outcome))
195 }
196
197 fn begin(
198 &mut self,
199 key: K,
200 participants: Vec<P>,
201 reply_to: Route,
202 ) -> Actions<A, Never, Route::Sends, NoBirths> {
203 if self.records.iter().any(|record| record.key == key) {
204 return Self::result(
205 reply_to,
206 AcknowledgementOutcome::Rejected(AcknowledgementError::Existing {
207 key,
208 participants,
209 }),
210 );
211 }
212 let mut distinct = Vec::new();
213 for participant in participants {
214 if !distinct.contains(&participant) {
215 distinct.push(participant);
216 }
217 }
218 let remaining = distinct.len();
219 let state = if distinct.is_empty() {
220 AcknowledgementState::Completed
221 } else {
222 AcknowledgementState::Pending {
223 remaining: distinct,
224 acknowledged: Vec::new(),
225 }
226 };
227 self.records.push(AcknowledgementRecord {
228 key: key.clone(),
229 state,
230 });
231 let outcome = if remaining == 0 {
232 AcknowledgementOutcome::Completed { key }
233 } else {
234 AcknowledgementOutcome::Started { key, remaining }
235 };
236 Self::result(reply_to, outcome)
237 }
238
239 fn acknowledge(
240 &mut self,
241 key: K,
242 participant: P,
243 reply_to: Route,
244 ) -> Actions<A, Never, Route::Sends, NoBirths> {
245 let Some(record) = self.records.iter_mut().find(|record| record.key == key) else {
246 return Self::result(
247 reply_to,
248 AcknowledgementOutcome::Rejected(AcknowledgementError::Unknown(
249 AcknowledgementInput::Acknowledge { key, participant },
250 )),
251 );
252 };
253 let (remaining, acknowledged) = match &mut record.state {
254 AcknowledgementState::Pending {
255 remaining,
256 acknowledged,
257 } => (remaining, acknowledged),
258 AcknowledgementState::Completed => {
259 return Self::result(
260 reply_to,
261 AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
262 AcknowledgementInput::Acknowledge { key, participant },
263 )),
264 );
265 }
266 AcknowledgementState::Cancelled => {
267 return Self::result(
268 reply_to,
269 AcknowledgementOutcome::Rejected(AcknowledgementError::Cancelled(
270 AcknowledgementInput::Acknowledge { key, participant },
271 )),
272 );
273 }
274 };
275 if acknowledged.contains(&participant) {
276 return Self::result(
277 reply_to,
278 AcknowledgementOutcome::Rejected(AcknowledgementError::DuplicateParticipant {
279 key,
280 participant,
281 }),
282 );
283 }
284 let Some(index) = remaining
285 .iter()
286 .position(|required| required == &participant)
287 else {
288 return Self::result(
289 reply_to,
290 AcknowledgementOutcome::Rejected(AcknowledgementError::UnexpectedParticipant {
291 key,
292 participant,
293 }),
294 );
295 };
296 remaining.remove(index);
297 acknowledged.push(participant.clone());
298 if remaining.is_empty() {
299 record.state = AcknowledgementState::Completed;
300 Self::result(reply_to, AcknowledgementOutcome::Completed { key })
301 } else {
302 Self::result(
303 reply_to,
304 AcknowledgementOutcome::Acknowledged {
305 key,
306 participant,
307 remaining: remaining.len(),
308 },
309 )
310 }
311 }
312
313 fn cancel(&mut self, key: K, reply_to: Route) -> Actions<A, Never, Route::Sends, NoBirths> {
314 let Some(record) = self.records.iter_mut().find(|record| record.key == key) else {
315 return Self::result(
316 reply_to,
317 AcknowledgementOutcome::Rejected(AcknowledgementError::Unknown(
318 AcknowledgementInput::Cancel { key },
319 )),
320 );
321 };
322 match record.state {
323 AcknowledgementState::Pending { .. } => {
324 record.state = AcknowledgementState::Cancelled;
325 Self::result(reply_to, AcknowledgementOutcome::Cancelled { key })
326 }
327 AcknowledgementState::Completed => Self::result(
328 reply_to,
329 AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
330 AcknowledgementInput::Cancel { key },
331 )),
332 ),
333 AcknowledgementState::Cancelled => Self::result(
334 reply_to,
335 AcknowledgementOutcome::Rejected(AcknowledgementError::Cancelled(
336 AcknowledgementInput::Cancel { key },
337 )),
338 ),
339 }
340 }
341}
342
343impl<A, K, P, Route> Default for Acknowledgements<A, K, P, Route>
344where
345 A: Address,
346 K: Clone + Eq,
347 P: Clone + Eq,
348 Route:
349 DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
350{
351 fn default() -> Self {
352 Self::new()
353 }
354}
355
356impl<A, K, P, Route> BehaviorBase for Acknowledgements<A, K, P, Route>
357where
358 A: Address,
359 Route:
360 DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
361{
362 type Base = Self;
363 fn base(&self) -> &Self {
364 self
365 }
366}
367
368impl<A, K, P, Route> behavior::Protocol for Acknowledgements<A, K, P, Route>
369where
370 A: Address,
371 Route:
372 DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
373{
374 type Addr = A;
375 type Msg = AcknowledgementMessage<K, P, Route>;
376}
377
378impl<A, K, P, Route> Behavior for Acknowledgements<A, K, P, Route>
379where
380 A: Address,
381 K: Clone + Eq,
382 P: Clone + Eq,
383 Route:
384 DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = AcknowledgementOutcome<K, P>>>,
385 Route::Sends: behavior::SendsFor<User<A, AcknowledgementMessage<K, P, Route>>>,
386{
387 type Protocol = Self;
388 type Event = User<A, behavior::BehaviorMessage<Self>>;
389 type Sends = Route::Sends;
390 type Ph = Never;
391 type Error = Never;
392 type Birth = NoBirths;
393
394 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
395 Ok(match event.message {
396 AcknowledgementMessage::Begin {
397 key,
398 participants,
399 reply_to,
400 } => self.begin(key, participants, reply_to),
401 AcknowledgementMessage::Acknowledge {
402 key,
403 participant,
404 reply_to,
405 } => self.acknowledge(key, participant, reply_to),
406 AcknowledgementMessage::Cancel { key, reply_to } => self.cancel(key, reply_to),
407 })
408 }
409}
410
411#[cfg(test)]
412mod tests {
413 use super::*;
414 use crate::Activate as _;
415 use behavior::MailAddr;
416
417 struct Reply;
418 impl behavior::Protocol for Reply {
419 type Addr = MailAddr;
420 type Msg = AcknowledgementOutcome<u8, u8>;
421 }
422
423 impl Behavior for Reply {
424 type Protocol = Self;
425 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
426 type Sends = Vec<Never>;
427 type Ph = Never;
428 type Error = Never;
429 type Birth = NoBirths;
430 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
431 Ok(Actions::cont())
432 }
433 }
434
435 type Subject = Acknowledgements<MailAddr, u8, u8, Recipient<Reply>>;
436 fn reply() -> Recipient<Reply> {
437 Recipient::global(MailAddr(1))
438 }
439
440 #[test]
441 fn membership_is_normalized_and_completion_is_terminal() {
442 let mut subject = (Subject::new()).initialize().unwrap().behavior;
443 let started = subject
444 .receive(
445 MailAddr(9),
446 AcknowledgementMessage::Begin {
447 key: 7,
448 participants: vec![1, 1, 2],
449 reply_to: reply(),
450 },
451 )
452 .unwrap();
453 assert!(matches!(
454 started.sends[0].message,
455 AcknowledgementOutcome::Started {
456 key: 7,
457 remaining: 2
458 }
459 ));
460 let first = subject
461 .receive(
462 MailAddr(9),
463 AcknowledgementMessage::Acknowledge {
464 key: 7,
465 participant: 1,
466 reply_to: reply(),
467 },
468 )
469 .unwrap();
470 assert!(matches!(
471 first.sends[0].message,
472 AcknowledgementOutcome::Acknowledged { remaining: 1, .. }
473 ));
474 let completed = subject
475 .receive(
476 MailAddr(9),
477 AcknowledgementMessage::Acknowledge {
478 key: 7,
479 participant: 2,
480 reply_to: reply(),
481 },
482 )
483 .unwrap();
484 assert!(matches!(
485 completed.sends[0].message,
486 AcknowledgementOutcome::Completed { key: 7 }
487 ));
488 let stale = subject
489 .receive(
490 MailAddr(9),
491 AcknowledgementMessage::Acknowledge {
492 key: 7,
493 participant: 2,
494 reply_to: reply(),
495 },
496 )
497 .unwrap();
498 assert!(matches!(
499 stale.sends[0].message,
500 AcknowledgementOutcome::Rejected(AcknowledgementError::Completed(
501 AcknowledgementInput::Acknowledge {
502 key: 7,
503 participant: 2,
504 }
505 ))
506 ));
507 }
508
509 #[test]
510 fn rejection_and_cancellation_do_not_conflate_phases() {
511 let mut subject = (Subject::new()).initialize().unwrap().behavior;
512 let started = subject
513 .receive(
514 MailAddr(9),
515 AcknowledgementMessage::Begin {
516 key: 1,
517 participants: vec![3],
518 reply_to: reply(),
519 },
520 )
521 .unwrap();
522 assert!(matches!(
523 started.sends[0].message,
524 AcknowledgementOutcome::Started {
525 key: 1,
526 remaining: 1
527 }
528 ));
529 let unexpected = subject
530 .receive(
531 MailAddr(9),
532 AcknowledgementMessage::Acknowledge {
533 key: 1,
534 participant: 4,
535 reply_to: reply(),
536 },
537 )
538 .unwrap();
539 assert!(matches!(
540 unexpected.sends[0].message,
541 AcknowledgementOutcome::Rejected(AcknowledgementError::UnexpectedParticipant { .. })
542 ));
543 let cancelled = subject
544 .receive(
545 MailAddr(9),
546 AcknowledgementMessage::Cancel {
547 key: 1,
548 reply_to: reply(),
549 },
550 )
551 .unwrap();
552 assert!(matches!(
553 cancelled.sends[0].message,
554 AcknowledgementOutcome::Cancelled { key: 1 }
555 ));
556 assert!(matches!(
557 subject.records()[0].state,
558 AcknowledgementState::Cancelled
559 ));
560 }
561}