| 112 | #[async_trait] |
| 113 | impl<S: MessageSender> Transport for MultiplexedTransport<S> { |
| 114 | async fn request<P, R>( |
| 115 | &self, |
| 116 | peer_id: &PublicKey, |
| 117 | request: &RequestObject<P>, |
| 118 | ) -> Result<JsonRpcResponse<R>> |
| 119 | where |
| 120 | P: Serialize + Send + Sync, |
| 121 | R: DeserializeOwned + Send, |
| 122 | { |
| 123 | let id = request.id.as_ref().ok_or(Error::MissingId)?; |
| 124 | let payload = serde_json::to_vec(request)?; |
| 125 | |
| 126 | // Register pending before sending |
| 127 | let rx = self.pending().insert(id.clone()).await; |
| 128 | |
| 129 | // Send via backend |
| 130 | if let Err(e) = self.sender.send(peer_id, &payload).await { |
| 131 | self.pending.remove(id).await; |
| 132 | return Err(e); |
| 133 | }; |
| 134 | |
| 135 | let response_bytes = tokio::time::timeout(self.timeout, rx) |
| 136 | .await |
| 137 | .map_err(|_| { |
| 138 | let pending = self.pending.clone(); |
| 139 | let id = id.clone(); |
| 140 | tokio::spawn(async move { pending.remove(&id).await }); |
| 141 | Error::Timeout |
| 142 | })? |
| 143 | .map_err(|_| Error::Internal("channel closed unexpectedly".into()))?; |
| 144 | |
| 145 | Ok(serde_json::from_slice(&response_bytes)?) |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | #[cfg(test)] |