1mod 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
24pub type ErrorHandler<Request, Error> = dnet_utils::pipe::ErrorHandler<Message<Request>, Error>;
27
28pub 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
38pub type RequestId = u64;
40
41#[derive(Debug, Clone, Serialize, Deserialize)]
43pub struct Message<Request> {
44 pub id: RequestId,
46
47 pub payload: Payload<Request>,
49}
50
51#[derive(Debug, Clone, Serialize, Deserialize)]
53pub enum Payload<Request> {
54 Request(Request),
56
57 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 #[derive(Debug, Clone, IntoTransferable)]
72 pub struct Message<C, Request>
73 where
74 Request: WrapperLikeTransferable<C>,
75 C: Codec,
76 {
77 pub id: RequestId,
79
80 #[into_transferable]
82 pub payload: Payload<C, Request>,
83 }
84
85 #[derive(Debug, Clone, IntoTransferable)]
87 pub enum Payload<C, Request>
88 where
89 Request: WrapperLikeTransferable<C>,
90 C: Codec,
91 {
92 Request {
94 #[into_transferable]
95 request: Request,
96 _codec: PhantomData<C>,
97 },
98
99 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
197pub type Result<T> = std::result::Result<T, super::Error>;
199
200pub struct Configuration {
204 pub shutdown: Box<dyn crate::Shutdown>,
207
208 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#[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 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 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
260pub trait Consume<Consumer> {
266 type Request;
268
269 type Response;
271
272 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}