Skip to main content

dnet_rpc/producer/
mod.rs

1//! Producer related functionality.
2
3pub mod abortable;
4
5use std::time::Duration;
6
7use dportable::JoinHandle;
8use futures::future::pending;
9use serde::{Deserialize, Serialize};
10
11use crate::consumer;
12
13use super::ShutdownType;
14
15/// Error handler for errors that may occur while sending or receiving
16/// messages through transport.
17pub type ErrorHandler<Response, Error> = dnet_utils::pipe::ErrorHandler<Message<Response>, Error>;
18
19/// Helper trait for producer transports.
20pub trait Transport<Request, Response, Error>:
21    crate::Transport<consumer::Message<Request>, self::Message<Response>, Error> + Unpin
22{
23}
24impl<T, Request, Response, Error> Transport<Request, Response, Error> for T where
25    T: crate::Transport<consumer::Message<Request>, self::Message<Response>, Error> + Unpin
26{
27}
28
29/// Configuration for [produce] method.
30///
31/// [produce]: self::Produce::produce
32pub struct Configuration {
33    /// When this future resolves producers stops producing and returns.
34    ///
35    /// **NOTE**: By default it is set with [futures::future::Pending], but it doesn't
36    /// mean producer will never stop - it will still stop when transport is closed,
37    /// or it can be stopped by provided `send_error_callback` or `receive_error_callback`.
38    ///
39    /// This future is intended for triggering shutdown manually.
40    pub shutdown: Box<dyn crate::Shutdown>,
41
42    /// Optional duration after which producer will shutdown if it will not receive
43    /// any messages during that period (and there are no requests pending).
44    ///
45    /// **NOTE**: It resets on every message received - so a producer that has received messages
46    /// in the past can still time out when a period of `duration` length with no further messages
47    /// occurs.
48    pub timeout: Option<Duration>,
49}
50
51impl Default for Configuration {
52    fn default() -> Self {
53        Self {
54            shutdown: Box::new(pending()),
55            timeout: Default::default(),
56        }
57    }
58}
59
60/// Producer message.
61#[derive(Debug, Clone, Serialize, Deserialize)]
62pub enum Message<Response> {
63    /// Response to request.
64    Response {
65        /// Request id.
66        id: u64,
67
68        /// Response value.
69        response: Response,
70    },
71
72    /// Producer was shutdown with [ShutdownType::Aborted].
73    Aborted,
74
75    /// Producer was shutdown with [ShutdownType::Shutdown].
76    Shutdown,
77}
78
79/// Stream response wrapper.
80#[derive(Debug, Clone, Serialize, Deserialize)]
81pub enum StreamResponse<T> {
82    /// Stream open.
83    Open,
84
85    /// Next stream item.
86    Item(T),
87
88    /// Stream closed (has no more items).
89    Closed,
90}
91
92#[cfg(target_arch = "wasm32")]
93mod transferable {
94    use std::marker::PhantomData;
95
96    use dnet_base::Codec;
97    use dnet_js::{wrapper::WrapperLikeTransferable, IntoTransferable};
98
99    /// Producer message.
100    #[derive(Debug, Clone, IntoTransferable)]
101    pub enum Message<C, Response>
102    where
103        Response: WrapperLikeTransferable<C>,
104        C: Codec,
105    {
106        /// Response to request.
107        Response {
108            id: u64,
109            #[into_transferable]
110            response: Response,
111            _codec: PhantomData<C>,
112        },
113
114        /// Producer was shutdown with [ShutdownType::Aborted].
115        Aborted,
116
117        /// Producer was shutdown with [ShutdownType::Shutdown].
118        Shutdown,
119    }
120
121    impl<C, Response> From<super::Message<Response>> for Message<C, Response>
122    where
123        Response: WrapperLikeTransferable<C>,
124        C: Codec,
125    {
126        fn from(value: super::Message<Response>) -> Self {
127            match value {
128                super::Message::Response { id, response, .. } => Self::Response {
129                    id,
130                    response,
131                    _codec: PhantomData,
132                },
133                super::Message::Aborted => Self::Aborted,
134                super::Message::Shutdown => Self::Shutdown,
135            }
136        }
137    }
138
139    /// Stream response wrapper.
140    #[derive(Debug, Clone, IntoTransferable)]
141    pub enum StreamResponse<C, T>
142    where
143        T: WrapperLikeTransferable<C>,
144        C: Codec,
145    {
146        /// Stream open.
147        Open,
148
149        /// Next stream item.
150        Item {
151            #[into_transferable]
152            item: T,
153            _codec: PhantomData<C>,
154        },
155
156        /// Stream closed (has no more items).
157        Closed,
158    }
159
160    impl<C, T> From<super::StreamResponse<T>> for StreamResponse<C, T>
161    where
162        T: WrapperLikeTransferable<C>,
163        C: Codec,
164    {
165        fn from(value: super::StreamResponse<T>) -> Self {
166            match value {
167                super::StreamResponse::Open => Self::Open,
168                super::StreamResponse::Item(item) => Self::Item {
169                    item,
170                    _codec: PhantomData,
171                },
172                super::StreamResponse::Closed => Self::Closed,
173            }
174        }
175    }
176}
177
178#[cfg(target_arch = "wasm32")]
179impl<C, Response> From<transferable::Message<C, Response>> for Message<Response>
180where
181    Response: dnet_js::wrapper::WrapperLikeTransferable<C>,
182    C: dnet_base::Codec,
183{
184    fn from(value: transferable::Message<C, Response>) -> Self {
185        match value {
186            transferable::Message::Response { id, response, .. } => Self::Response { id, response },
187            transferable::Message::Aborted => Self::Aborted,
188            transferable::Message::Shutdown => Self::Shutdown,
189        }
190    }
191}
192
193#[cfg(target_arch = "wasm32")]
194impl<C, Request>
195    dnet_js::IntoTransferable<
196        dnet_js::wrapper::Context<C>,
197        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
198    > for Message<Request>
199where
200    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
201    C: dnet_base::Codec,
202{
203    type Output = transferable::MessageWrapper<C, Request>;
204
205    fn into_transferable(self) -> Self::Output {
206        transferable::Message::from(self).into_transferable()
207    }
208}
209
210#[cfg(target_arch = "wasm32")]
211impl<C, Request>
212    dnet_js::FromTransferable<
213        dnet_js::wrapper::Context<C>,
214        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
215    > for Message<Request>
216where
217    Request: dnet_js::wrapper::WrapperLikeTransferable<C>,
218    C: dnet_base::Codec,
219{
220    type Input = transferable::MessageWrapper<C, Request>;
221
222    fn from_transferable(input: Self::Input) -> Self {
223        use dnet_utils::unwrap::Unwrap;
224
225        input.unwrap().into()
226    }
227}
228
229#[cfg(target_arch = "wasm32")]
230impl<C, T> From<transferable::StreamResponse<C, T>> for StreamResponse<T>
231where
232    T: dnet_js::wrapper::WrapperLikeTransferable<C>,
233    C: dnet_base::Codec,
234{
235    fn from(value: transferable::StreamResponse<C, T>) -> Self {
236        match value {
237            transferable::StreamResponse::Open => Self::Open,
238            transferable::StreamResponse::Item { item, .. } => Self::Item(item),
239            transferable::StreamResponse::Closed => Self::Closed,
240        }
241    }
242}
243
244#[cfg(target_arch = "wasm32")]
245impl<C, T>
246    dnet_js::IntoTransferable<
247        dnet_js::wrapper::Context<C>,
248        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
249    > for StreamResponse<T>
250where
251    T: dnet_js::wrapper::WrapperLikeTransferable<C>,
252    C: dnet_base::Codec,
253{
254    type Output = transferable::StreamResponseWrapper<C, T>;
255
256    fn into_transferable(self) -> Self::Output {
257        transferable::StreamResponse::from(self).into_transferable()
258    }
259}
260
261#[cfg(target_arch = "wasm32")]
262impl<C, T>
263    dnet_js::FromTransferable<
264        dnet_js::wrapper::Context<C>,
265        dnet_js::wrapper::Error<<C as dnet_base::Encode>::Error, <C as dnet_base::Decode>::Error>,
266    > for StreamResponse<T>
267where
268    T: dnet_js::wrapper::WrapperLikeTransferable<C>,
269    C: dnet_base::Codec,
270{
271    type Input = transferable::StreamResponseWrapper<C, T>;
272
273    fn from_transferable(input: Self::Input) -> Self {
274        use dnet_utils::unwrap::Unwrap;
275
276        input.unwrap().into()
277    }
278}
279
280/// Trait implemented by producers.
281///
282/// You should never have to implement it manually - use derive macro.
283pub trait Produce: Sized {
284    /// Request message type.
285    type Request;
286    /// Response message type.
287    type Response;
288
289    /// Produce using given transport and configuration.
290    fn produce<Transport, Error>(
291        self,
292        transport: Transport,
293        configuration: Configuration,
294        error_handler: ErrorHandler<Self::Response, Error>,
295    ) -> JoinHandle<ShutdownType>
296    where
297        Transport: self::Transport<Self::Request, Self::Response, Error>,
298        Error: crate::TransportError;
299}