1use behavior::{
4 ActionItem, ActionItemResult, Delivery, EndpointAddress, EstablishedDelivery,
5 EstablishedRecipient, InterpretItem, InterpretSends, Interpretation, InterpretationProgress,
6 ItemSettlement, Own, Protocol, Recipient, RecipientAddress, SendEffects, SendInput,
7 SendSettlements, SendsFor, SettledItem, SourceCustody, SourceProgress, SourceSettlementCustody,
8};
9use core::future::Future;
10
11mod sealed {
12 pub trait DeliveryRoute {}
13}
14
15pub trait DeliveryRoute: sealed::DeliveryRoute + Sized {
49 type Protocol: Protocol;
51 type Sends: SendEffects;
53
54 fn deliver(self, message: <Self::Protocol as Protocol>::Msg) -> Self::Sends;
56}
57
58impl<P: Protocol> sealed::DeliveryRoute for Recipient<P> {}
59impl<P: Protocol> DeliveryRoute for Recipient<P> {
60 type Protocol = P;
61 type Sends = Vec<Delivery<P>>;
62
63 fn deliver(self, message: P::Msg) -> Self::Sends {
64 vec![Delivery::new(self, message)]
65 }
66}
67
68impl<P> sealed::DeliveryRoute for EstablishedRecipient<P>
69where
70 P: Protocol,
71 P::Addr: EndpointAddress,
72{
73}
74impl<P> DeliveryRoute for EstablishedRecipient<P>
75where
76 P: Protocol,
77 P::Addr: EndpointAddress,
78{
79 type Protocol = P;
80 type Sends = Vec<EstablishedDelivery<P>>;
81
82 fn deliver(self, message: P::Msg) -> Self::Sends {
83 vec![EstablishedDelivery::new(self, message)]
84 }
85}
86
87pub enum ReplyRoute<P>
95where
96 P: Protocol,
97 P::Addr: EndpointAddress,
98{
99 Logical(Recipient<P>),
101 Established(EstablishedRecipient<P>),
103}
104
105impl<P> Clone for ReplyRoute<P>
106where
107 P: Protocol,
108 P::Addr: EndpointAddress,
109{
110 fn clone(&self) -> Self {
111 match self {
112 Self::Logical(recipient) => Self::Logical(*recipient),
113 Self::Established(recipient) => Self::Established(recipient.clone()),
114 }
115 }
116}
117
118impl<P> ReplyRoute<P>
119where
120 P: Protocol,
121 P::Addr: EndpointAddress,
122{
123 #[must_use]
125 pub const fn logical(recipient: Recipient<P>) -> Self {
126 Self::Logical(recipient)
127 }
128
129 #[must_use]
131 pub const fn established(recipient: EstablishedRecipient<P>) -> Self {
132 Self::Established(recipient)
133 }
134}
135
136impl<P> From<Recipient<P>> for ReplyRoute<P>
137where
138 P: Protocol,
139 P::Addr: EndpointAddress,
140{
141 fn from(recipient: Recipient<P>) -> Self {
142 Self::logical(recipient)
143 }
144}
145
146impl<P> From<EstablishedRecipient<P>> for ReplyRoute<P>
147where
148 P: Protocol,
149 P::Addr: EndpointAddress,
150{
151 fn from(recipient: EstablishedRecipient<P>) -> Self {
152 Self::established(recipient)
153 }
154}
155
156pub enum ReplyDelivery<Logical, Established> {
158 Logical(Logical),
160 Established(Established),
162}
163
164impl<Logical, Established> behavior::ClassifySettlement for ReplyDelivery<Logical, Established>
165where
166 Logical: behavior::ClassifySettlement,
167 Established: behavior::ClassifySettlement,
168{
169 fn settlement_status(&self) -> behavior::SettlementStatus {
170 match self {
171 Self::Logical(delivery) => delivery.settlement_status(),
172 Self::Established(delivery) => delivery.settlement_status(),
173 }
174 }
175}
176
177impl<Logical: Clone, Established: Clone> Clone for ReplyDelivery<Logical, Established> {
178 fn clone(&self) -> Self {
179 match self {
180 Self::Logical(delivery) => Self::Logical(delivery.clone()),
181 Self::Established(delivery) => Self::Established(delivery.clone()),
182 }
183 }
184}
185
186impl<Logical, Established> core::fmt::Debug for ReplyDelivery<Logical, Established>
187where
188 Logical: core::fmt::Debug,
189 Established: core::fmt::Debug,
190{
191 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
192 match self {
193 Self::Logical(delivery) => formatter.debug_tuple("Logical").field(delivery).finish(),
194 Self::Established(delivery) => formatter
195 .debug_tuple("Established")
196 .field(delivery)
197 .finish(),
198 }
199 }
200}
201
202impl<Logical, Established> PartialEq for ReplyDelivery<Logical, Established>
203where
204 Logical: PartialEq,
205 Established: PartialEq,
206{
207 fn eq(&self, other: &Self) -> bool {
208 match (self, other) {
209 (Self::Logical(left), Self::Logical(right)) => left == right,
210 (Self::Established(left), Self::Established(right)) => left == right,
211 (Self::Logical(_), Self::Established(_)) | (Self::Established(_), Self::Logical(_)) => {
212 false
213 }
214 }
215 }
216}
217
218impl<Logical: Eq, Established: Eq> Eq for ReplyDelivery<Logical, Established> {}
219
220pub struct ReplyDeliveries<Logical, Established> {
222 deliveries: Vec<ReplyDelivery<Logical, Established>>,
223}
224
225impl<Logical, Established> ReplyDeliveries<Logical, Established> {
226 #[must_use]
228 pub fn new(deliveries: Vec<ReplyDelivery<Logical, Established>>) -> Self {
229 Self { deliveries }
230 }
231
232 #[must_use]
234 pub fn as_slice(&self) -> &[ReplyDelivery<Logical, Established>] {
235 &self.deliveries
236 }
237
238 #[must_use]
240 pub fn into_deliveries(self) -> Vec<ReplyDelivery<Logical, Established>> {
241 self.deliveries
242 }
243}
244
245impl<Logical: Clone, Established: Clone> Clone for ReplyDeliveries<Logical, Established> {
246 fn clone(&self) -> Self {
247 Self::new(self.deliveries.clone())
248 }
249}
250
251impl<Logical, Established> core::fmt::Debug for ReplyDeliveries<Logical, Established>
252where
253 ReplyDelivery<Logical, Established>: core::fmt::Debug,
254{
255 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
256 formatter.debug_list().entries(&self.deliveries).finish()
257 }
258}
259
260impl<Logical, Established> PartialEq for ReplyDeliveries<Logical, Established>
261where
262 ReplyDelivery<Logical, Established>: PartialEq,
263{
264 fn eq(&self, other: &Self) -> bool {
265 self.deliveries == other.deliveries
266 }
267}
268
269impl<Logical, Established> Eq for ReplyDeliveries<Logical, Established> where
270 ReplyDelivery<Logical, Established>: Eq
271{
272}
273
274impl<Logical, Established> SendEffects for ReplyDeliveries<Logical, Established> {
275 fn empty() -> Self {
276 Self::new(Vec::new())
277 }
278
279 fn append(&mut self, mut other: Self) {
280 self.deliveries.append(&mut other.deliveries);
281 }
282}
283
284impl<Logical, Established> behavior::ClassifySettlement for ReplyDeliveries<Logical, Established>
285where
286 Logical: behavior::ClassifySettlement,
287 Established: behavior::ClassifySettlement,
288{
289 fn settlement_status(&self) -> behavior::SettlementStatus {
290 self.deliveries.settlement_status()
291 }
292}
293
294impl<Event, Logical, Established> SendsFor<Event> for ReplyDeliveries<Logical, Established> {}
295
296impl<Logical, Established> SendInput<ReplyDelivery<Logical, Established>, Own>
297 for ReplyDeliveries<Logical, Established>
298{
299 fn emit(&mut self, input: ReplyDelivery<Logical, Established>) {
300 self.deliveries.push(input);
301 }
302}
303
304impl<P> SendSettlements for ReplyDeliveries<Delivery<P>, EstablishedDelivery<P>>
305where
306 P: Protocol,
307 P::Addr: EndpointAddress,
308 Delivery<P>: ActionItem,
309 EstablishedDelivery<P>: ActionItem,
310{
311 type Settlements =
312 ReplyDeliveries<ActionItemResult<Delivery<P>>, ActionItemResult<EstablishedDelivery<P>>>;
313 type SourceCustody = Self::Settlements;
314 type InterpretationCustody = Vec<
315 ReplyDelivery<
316 (
317 Option<Delivery<P>>,
318 Option<
319 ItemSettlement<
320 Delivery<P>,
321 <Delivery<P> as ActionItem>::Accepted,
322 <Delivery<P> as ActionItem>::Rejection,
323 <Delivery<P> as ActionItem>::Prerequisite,
324 >,
325 >,
326 ),
327 (
328 Option<EstablishedDelivery<P>>,
329 Option<
330 ItemSettlement<
331 EstablishedDelivery<P>,
332 <EstablishedDelivery<P> as ActionItem>::Accepted,
333 <EstablishedDelivery<P> as ActionItem>::Rejection,
334 <EstablishedDelivery<P> as ActionItem>::Prerequisite,
335 >,
336 >,
337 ),
338 >,
339 >;
340
341 fn prepare_interpretation(
342 progress: &mut Option<
343 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
344 >,
345 ) {
346 *progress = match progress.take() {
347 Some(InterpretationProgress::Original(original)) => {
348 Some(InterpretationProgress::Interpreting(
349 original
350 .deliveries
351 .into_iter()
352 .map(|delivery| match delivery {
353 ReplyDelivery::Logical(input) => {
354 ReplyDelivery::Logical((Some(input), None))
355 }
356 ReplyDelivery::Established(input) => {
357 ReplyDelivery::Established((Some(input), None))
358 }
359 })
360 .collect(),
361 ))
362 }
363 retained => retained,
364 };
365 }
366
367 fn finish_interpretation(
368 progress: &mut Option<
369 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
370 >,
371 ) {
372 let Some(InterpretationProgress::Interpreting(rows)) = progress.as_ref() else {
373 return;
374 };
375 let corrupt = rows.iter().position(|row| match row {
376 ReplyDelivery::Logical((_, received)) => {
377 matches!(received, Some(ItemSettlement::Corrupt { .. }))
378 }
379 ReplyDelivery::Established((_, received)) => {
380 matches!(received, Some(ItemSettlement::Corrupt { .. }))
381 }
382 });
383 let complete = rows.iter().enumerate().all(|(index, row)| {
384 let after_corrupt = corrupt.is_some_and(|corrupt| index > corrupt);
385 match row {
386 ReplyDelivery::Logical((input, received)) => {
387 if after_corrupt {
388 input.is_some() && received.is_none()
389 } else {
390 input.is_none() && received.is_some()
391 }
392 }
393 ReplyDelivery::Established((input, received)) => {
394 if after_corrupt {
395 input.is_some() && received.is_none()
396 } else {
397 input.is_none() && received.is_some()
398 }
399 }
400 }
401 });
402 if !complete {
403 return;
404 }
405 let Some(InterpretationProgress::Interpreting(rows)) = progress.take() else {
406 return;
407 };
408 let mut remaining = rows.into_iter();
409 let mut settled = Vec::with_capacity(remaining.len());
410 while let Some(row) = remaining.next() {
411 match row {
412 ReplyDelivery::Logical((None, Some(received))) => {
413 settled.push(ReplyDelivery::Logical(SettledItem::Attempted(received)))
414 }
415 ReplyDelivery::Logical((Some(input), None)) => {
416 settled.push(ReplyDelivery::Logical(SettledItem::Unattempted(input)))
417 }
418 ReplyDelivery::Established((None, Some(received))) => {
419 settled.push(ReplyDelivery::Established(SettledItem::Attempted(received)))
420 }
421 ReplyDelivery::Established((Some(input), None)) => {
422 settled.push(ReplyDelivery::Established(SettledItem::Unattempted(input)))
423 }
424 row => {
425 let restored = settled
426 .into_iter()
427 .map(|settled| match settled {
428 ReplyDelivery::Logical(SettledItem::Attempted(received)) => {
429 ReplyDelivery::Logical((None, Some(received)))
430 }
431 ReplyDelivery::Logical(SettledItem::Unattempted(input)) => {
432 ReplyDelivery::Logical((Some(input), None))
433 }
434 ReplyDelivery::Established(SettledItem::Attempted(received)) => {
435 ReplyDelivery::Established((None, Some(received)))
436 }
437 ReplyDelivery::Established(SettledItem::Unattempted(input)) => {
438 ReplyDelivery::Established((Some(input), None))
439 }
440 })
441 .chain(core::iter::once(row))
442 .chain(remaining)
443 .collect();
444 *progress = Some(InterpretationProgress::Interpreting(restored));
445 return;
446 }
447 }
448 }
449 let settled = ReplyDeliveries::new(settled);
450 *progress = Some(InterpretationProgress::Completed(match corrupt {
451 Some(_) => Interpretation::Corrupt(settled),
452 None => Interpretation::Complete(settled),
453 }));
454 }
455
456 fn unattempted(
457 progress: &mut Option<
458 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
459 >,
460 ) {
461 *progress = match progress.take() {
462 Some(InterpretationProgress::Original(original)) => Some(
463 InterpretationProgress::Completed(Interpretation::Complete(ReplyDeliveries::new(
464 original
465 .deliveries
466 .into_iter()
467 .map(|delivery| match delivery {
468 ReplyDelivery::Logical(input) => {
469 ReplyDelivery::Logical(SettledItem::Unattempted(input))
470 }
471 ReplyDelivery::Established(input) => {
472 ReplyDelivery::Established(SettledItem::Unattempted(input))
473 }
474 })
475 .collect(),
476 ))),
477 ),
478 retained => retained,
479 };
480 }
481}
482
483impl<Host, RootEvent, P> SourceSettlementCustody<Host, RootEvent>
484 for ReplyDeliveries<ActionItemResult<Delivery<P>>, ActionItemResult<EstablishedDelivery<P>>>
485where
486 P: Protocol,
487 P::Addr: EndpointAddress,
488 Delivery<P>: ActionItem,
489 EstablishedDelivery<P>: ActionItem,
490{
491 type Custody = Self;
492
493 fn prepare_source(progress: &mut Option<SourceProgress<Self, Self::Custody>>) {
494 *progress = match progress.take() {
495 Some(SourceProgress::Original(original)) => Some(SourceProgress::Completed(
496 SourceCustody::Exhausted(original),
497 )),
498 retained => retained,
499 };
500 }
501
502 fn offer_next_to_source(
503 _: &mut Self::Custody,
504 _: &mut Host,
505 ) -> impl core::future::Future<Output = ()> + Send {
506 core::future::ready(())
507 }
508
509 fn finish_source(progress: &mut Option<SourceProgress<Self, Self::Custody>>) {
510 *progress = match progress.take() {
511 Some(SourceProgress::Original(original) | SourceProgress::Offering(original)) => Some(
512 SourceProgress::Completed(SourceCustody::Exhausted(original)),
513 ),
514 retained => retained,
515 };
516 }
517}
518
519impl<Interpreter, RootEvent, Path, P> InterpretSends<Interpreter, RootEvent, Path>
520 for ReplyDeliveries<Delivery<P>, EstablishedDelivery<P>>
521where
522 Interpreter: InterpretItem<Delivery<P>, RootEvent, Path>
523 + InterpretItem<EstablishedDelivery<P>, RootEvent, Path>
524 + Send,
525 P: Protocol,
526 P::Addr: EndpointAddress,
527 P::Addr: Send,
528 P::Msg: Send,
529 <P::Addr as RecipientAddress>::Established<P>: Send,
530{
531 fn interpret(
532 progress: &mut Option<
533 InterpretationProgress<Self, Self::InterpretationCustody, Self::Settlements>,
534 >,
535 interpreter: &mut Interpreter,
536 ) -> impl Future<Output = ()> + Send {
537 async move {
538 Self::prepare_interpretation(progress);
539 let Some(InterpretationProgress::Interpreting(rows)) = progress else {
540 return;
541 };
542 for row in rows {
543 match row {
544 ReplyDelivery::Logical((input, received)) => {
545 match (input.as_ref(), received.as_ref()) {
546 (None, Some(ItemSettlement::Corrupt { .. })) => break,
547 (None, Some(_)) => continue,
548 (Some(_), None) => {}
549 _ => return,
550 }
551 <Interpreter as InterpretItem<Delivery<P>, RootEvent, Path>>::interpret_item(interpreter,input,received).await;
552 match (input.as_ref(), received.as_ref()) {
553 (None, Some(ItemSettlement::Corrupt { .. })) => break,
554 (None, Some(_)) => {}
555 _ => return,
556 }
557 }
558 ReplyDelivery::Established((input, received)) => {
559 match (input.as_ref(), received.as_ref()) {
560 (None, Some(ItemSettlement::Corrupt { .. })) => break,
561 (None, Some(_)) => continue,
562 (Some(_), None) => {}
563 _ => return,
564 }
565 <Interpreter as InterpretItem<EstablishedDelivery<P>, RootEvent, Path>>::interpret_item(interpreter,input,received).await;
566 match (input.as_ref(), received.as_ref()) {
567 (None, Some(ItemSettlement::Corrupt { .. })) => break,
568 (None, Some(_)) => {}
569 _ => return,
570 }
571 }
572 }
573 }
574 Self::finish_interpretation(progress);
575 }
576 }
577}
578
579impl<P> sealed::DeliveryRoute for ReplyRoute<P>
580where
581 P: Protocol,
582 P::Addr: EndpointAddress,
583{
584}
585impl<P> DeliveryRoute for ReplyRoute<P>
586where
587 P: Protocol,
588 P::Addr: EndpointAddress,
589{
590 type Protocol = P;
591 type Sends = ReplyDeliveries<Delivery<P>, EstablishedDelivery<P>>;
592
593 fn deliver(self, message: P::Msg) -> Self::Sends {
594 let delivery = match self {
595 Self::Logical(recipient) => ReplyDelivery::Logical(Delivery::new(recipient, message)),
596 Self::Established(recipient) => {
597 ReplyDelivery::Established(EstablishedDelivery::new(recipient, message))
598 }
599 };
600 ReplyDeliveries::new(vec![delivery])
601 }
602}