Skip to main content

dnet_rpc/consumer/
mod.rs

1//! Consumer related functionality.
2
3mod value;
4use std::{
5    sync::{Arc, Mutex},
6    time::Duration,
7};
8
9use futures::{channel::oneshot, future::pending};
10use serde::{Deserialize, Serialize};
11pub use value::*;
12
13mod stream;
14pub use stream::*;
15
16use crate::{
17    parts::consumer::{RequestSender, ResultSender},
18    producer,
19};
20
21#[allow(unused_imports)]
22use super::ShutdownType;
23
24/// Error handler for errors that may occur while sending or receiving
25/// messages through transport.
26pub type ErrorHandler<Request, Error> = dnet_utils::pipe::ErrorHandler<Message<Request>, Error>;
27
28/// Helper trait for consumer transports.
29pub trait Transport<Request, Response, Error>:
30    crate::Transport<producer::Message<Response>, self::Message<Request>, Error> + Unpin
31{
32}
33impl<T, Request, Response, Error> Transport<Request, Response, Error> for T where
34    T: crate::Transport<producer::Message<Response>, self::Message<Request>, Error> + Unpin
35{
36}
37
38/// Request id.
39pub type RequestId = u64;
40
41/// Consumer message.
42#[derive(Debug, Clone, Serialize, Deserialize)]
43pub struct Message<Request> {
44    /// Request id.
45    pub id: RequestId,
46
47    /// Message payload.
48    pub payload: Payload<Request>,
49}
50
51/// Consumer message payload.
52#[derive(Debug, Clone, Serialize, Deserialize)]
53pub enum Payload<Request> {
54    /// Request arguments.
55    Request(Request),
56
57    /// Abort request.
58    Abort,
59}
60
61#[cfg(target_arch = "wasm32")]
62mod transferable {
63    use std::marker::PhantomData;
64
65    use dnet_base::Codec;
66    use dnet_js::{wrapper::WrapperLikeTransferable, IntoTransferable};
67
68    use crate::consumer::RequestId;
69
70    /// Consumer message.
71    #[derive(Debug, Clone, IntoTransferable)]
72    pub struct Message<C, Request>
73    where
74        Request: WrapperLikeTransferable<C>,
75        C: Codec,
76    {
77        /// Request id.
78        pub id: RequestId,
79
80        /// Message payload.
81        #[into_transferable]
82        pub payload: Payload<C, Request>,
83    }
84
85    /// Consumer message payload.
86    #[derive(Debug, Clone, IntoTransferable)]
87    pub enum Payload<C, Request>
88    where
89        Request: WrapperLikeTransferable<C>,
90        C: Codec,
91    {
92        /// Request arguments.
93        Request {
94            #[into_transferable]
95            request: Request,
96            _codec: PhantomData<C>,
97        },
98
99        /// Abort request.
100        Abort,
101    }
102
103    impl<C, Request> From<super::Message<Request>> for Message<C, Request>
104    where
105        Request: WrapperLikeTransferable<C>,
106        C: Codec,
107    {
108        fn from(value: super::Message<Request>) -> Self {
109            Self {
110                id: value.id,
111                payload: value.payload.into(),
112            }
113        }
114    }
115
116    impl<C, Request> From<super::Payload<Request>> for Payload<C, Request>
117    where
118        Request: WrapperLikeTransferable<C>,
119        C: Codec,
120    {
121        fn from(value: super::Payload<Request>) -> Self {
122            match value {
123                super::Payload::Request(request) => Self::Request {
124                    request,
125                    _codec: PhantomData,
126                },
127                super::Payload::Abort => Self::Abort,
128            }
129        }
130    }
131}
132
133#[cfg(target_arch = "wasm32")]
134impl<C, Request> From<transferable::Message<C, Request>> for Message<Request>
135where
136    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
137    C: dnet_base::Codec,
138{
139    fn from(value: transferable::Message<C, Request>) -> Self {
140        Self {
141            id: value.id,
142            payload: value.payload.into(),
143        }
144    }
145}
146
147#[cfg(target_arch = "wasm32")]
148impl<C, Request> From<transferable::Payload<C, Request>> for Payload<Request>
149where
150    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
151    C: dnet_base::Codec,
152{
153    fn from(value: transferable::Payload<C, Request>) -> Self {
154        match value {
155            transferable::Payload::Request { request, .. } => Self::Request(request),
156            transferable::Payload::Abort => Self::Abort,
157        }
158    }
159}
160
161#[cfg(target_arch = "wasm32")]
162impl<C, Request>
163    dnet_js::IntoTransferable<
164        dnet_js::wrapper::Context<C>,
165        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
166    > for Message<Request>
167where
168    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
169    C: dnet_base::Codec,
170{
171    type Output = transferable::MessageWrapper<C, Request>;
172
173    fn into_transferable(self) -> Self::Output {
174        transferable::Message::from(self).into_transferable()
175    }
176}
177
178#[cfg(target_arch = "wasm32")]
179impl<C, Request>
180    dnet_js::FromTransferable<
181        dnet_js::wrapper::Context<C>,
182        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
183    > for Message<Request>
184where
185    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
186    C: dnet_base::Codec,
187{
188    type Input = transferable::MessageWrapper<C, Request>;
189
190    fn from_transferable(input: Self::Input) -> Self {
191        use dnet_utils::unwrap::Unwrap;
192
193        input.unwrap().into()
194    }
195}
196
197/// Result type for consumer methods.
198pub type Result<T> = std::result::Result<T, super::Error>;
199
200/// Configuration for [consume] method.
201///
202/// [consume]: self::Consume::consume
203pub struct Configuration {
204    /// When this future resolves all pending consumer requests error out with
205    /// returned [ShutdownType]-related error.
206    pub shutdown: Box<dyn crate::Shutdown>,
207
208    /// Consumer timeout - closes used transport if consumer is idle
209    /// (there are no pending requests) and consumer is not used for
210    /// specified duration.
211    pub timeout: Option<Duration>,
212}
213
214impl Default for Configuration {
215    fn default() -> Self {
216        Self {
217            shutdown: Box::new(pending()),
218            timeout: Default::default(),
219        }
220    }
221}
222
223/// Responsible for aborting requests/streams.
224#[derive(Debug)]
225pub struct Aborter<Request, T> {
226    id: RequestId,
227    sender: RequestSender<Request, T>,
228    abort_sender: Arc<Mutex<Option<oneshot::Sender<()>>>>,
229}
230
231impl<Request, Output> Aborter<Request, Output> {
232    /// Abort value request/stream request/stream.
233    ///
234    /// **NOTE** regarding streams: Values that were already received from the producer
235    /// before abort has been completed will still be returned by the stream.
236    pub fn abort(self) {
237        let mut abort_sender = self.abort_sender.lock().unwrap();
238        if let Some(abort_sender) = abort_sender.take() {
239            self.sender.abort(self.id);
240            let _ = abort_sender.send(());
241        }
242    }
243
244    /// Id of the request this [Aborter] can abort.
245    pub fn request_id(&self) -> u64 {
246        self.id
247    }
248}
249
250impl<Request, Output> Clone for Aborter<Request, Output> {
251    fn clone(&self) -> Self {
252        Self {
253            id: self.id,
254            sender: self.sender.clone(),
255            abort_sender: self.abort_sender.clone(),
256        }
257    }
258}
259
260/// Trait implemented by consumers.
261///
262/// It is used to create consumers.
263///
264/// You should never have to implement it manually - use macros.
265pub trait Consume<Consumer> {
266    /// Request message type.
267    type Request;
268
269    /// Response message type.
270    type Response;
271
272    /// Create consumer using given transport and configuration.
273    fn consume<Transport, Error>(
274        transport: Transport,
275        configuration: Configuration,
276        error_handler: ErrorHandler<Self::Request, Error>,
277    ) -> Consumer
278    where
279        Transport: self::Transport<Self::Request, Self::Response, Error>,
280        Error: crate::TransportError;
281}