(&mut self)
| 359 | } |
| 360 | |
| 361 | async fn try_recv(&mut self) { |
| 362 | loop { |
| 363 | if self.ended { |
| 364 | trace!("response stream has ended"); |
| 365 | return; |
| 366 | } |
| 367 | |
| 368 | let cancel_token = self.cancel_token.clone(); |
| 369 | tokio::select! { |
| 370 | res = self.recv() => { |
| 371 | let _ = self.event_tx.send(res).await.map_err(|err| error!(?err, "failed to send event to channel")); |
| 372 | }, |
| 373 | _ = cancel_token.cancelled() => { |
| 374 | debug!("response parser was cancelled"); |
| 375 | let err = self.error(RecvErrorKind::Cancelled); |
| 376 | *self.request_metadata.lock().await = Some(err.request_metadata.clone()); |
| 377 | let _ = self.event_tx.send(Err(err)).await.map_err(|err| error!(?err, "failed to send error to channel")); |
| 378 | return; |
| 379 | }, |
| 380 | } |
| 381 | } |
| 382 | } |
| 383 | |
| 384 | /// Consumes the associated [ConverseStreamResponse] until a valid [ResponseEvent] is parsed. |
| 385 | async fn recv(&mut self) -> Result<ResponseEvent, RecvError> { |
no test coverage detected