dnet_rpc/consumer/
value.rs1use std::{
4 pin::Pin,
5 sync::{Arc, Mutex},
6 task::{Context, Poll},
7};
8
9use futures::{channel::oneshot, future::FusedFuture, ready, Future, FutureExt};
10use pin_project::{pin_project, pinned_drop};
11
12use crate::parts::consumer::RequestSender;
13
14use super::{Aborter, ResultSender};
15
16#[derive(Debug)]
18#[pin_project(PinnedDrop)]
19pub struct ValueRequest<Request, T> {
20 sender: RequestSender<Request, T>,
21 id: u64,
22 request: Option<Request>,
23 receiver: Option<oneshot::Receiver<super::Result<T>>>,
24 aborter: Option<Aborter<Request, T>>,
25 abort_receiver: Option<oneshot::Receiver<()>>,
26}
27
28impl<Request, T> ValueRequest<Request, T> {
29 pub fn new(sender: RequestSender<Request, T>, id: u64, request: Request) -> Self {
34 ValueRequest {
35 sender,
36 id,
37 request: Some(request),
38 receiver: None,
39 aborter: None,
40 abort_receiver: None,
41 }
42 }
43
44 pub fn id(&self) -> u64 {
46 self.id
47 }
48
49 pub fn aborter(&mut self) -> Aborter<Request, T> {
51 let aborter = self.aborter.get_or_insert_with(|| {
52 let (abort_sender, abort_receiver) = oneshot::channel();
53 self.abort_receiver = Some(abort_receiver);
54 Aborter {
55 id: self.id,
56 sender: self.sender.clone(),
57 abort_sender: Arc::new(Mutex::new(Some(abort_sender))),
58 }
59 });
60 aborter.clone()
61 }
62}
63
64impl<Request, T> Future for ValueRequest<Request, T> {
65 type Output = super::Result<T>;
66
67 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
68 let mut me = self.project();
69 if let Some(abort_receiver) = me.abort_receiver {
70 if let Poll::Ready(result) = abort_receiver.poll_unpin(cx) {
71 if result.is_ok() {
72 me.request.take();
73 me.receiver.take();
74 return Poll::Ready(Err(super::super::Error::Aborted));
75 }
76 }
77 }
78
79 if let Some(request) = me.request.take() {
80 let (sender, mut receiver) = oneshot::channel();
81 let sender = ResultSender::Value(sender);
82 me.sender.send(*me.id, request, sender).map_err(|_| {
83 me.receiver.take();
84 super::super::Error::Shutdown
85 })?;
86 match receiver.poll_unpin(cx) {
87 Poll::Ready(result) => {
88 me.receiver.take();
89 Poll::Ready(result.map_err(|_| super::super::Error::Dropped)?)
90 }
91 Poll::Pending => {
92 *me.receiver = Some(receiver);
93 Poll::Pending
94 }
95 }
96 } else if let Some(receiver) = &mut me.receiver {
97 let result = ready!(receiver.poll_unpin(cx));
98 me.receiver.take();
99 Poll::Ready(result.map_err(|_| super::super::Error::Dropped)?)
100 } else {
101 Poll::Pending
102 }
103 }
104}
105
106impl<Request, T> FusedFuture for ValueRequest<Request, T> {
107 fn is_terminated(&self) -> bool {
108 self.request.is_none() && self.receiver.is_none()
109 }
110}
111
112#[pinned_drop]
113impl<Request, T> PinnedDrop for ValueRequest<Request, T> {
114 fn drop(self: Pin<&mut Self>) {
115 if !self.is_terminated() {
116 self.sender.abort(self.id);
117 }
118 }
119}