(
&self,
request: Request<api::StreamGatewayFramesRequest>,
)
| 706 | type StreamGatewayFramesStream = DropReceiver<Result<api::LogItem, Status>>; |
| 707 | |
| 708 | async fn stream_gateway_frames( |
| 709 | &self, |
| 710 | request: Request<api::StreamGatewayFramesRequest>, |
| 711 | ) -> Result<Response<Self::StreamGatewayFramesStream>, Status> { |
| 712 | let req = request.get_ref(); |
| 713 | let gw_id = EUI64::from_str(&req.gateway_id).map_err(|e| e.status())?; |
| 714 | |
| 715 | self.validator |
| 716 | .validate( |
| 717 | request.extensions(), |
| 718 | validator::ValidateGatewayAccess::new(validator::Flag::Read, gw_id), |
| 719 | ) |
| 720 | .await?; |
| 721 | |
| 722 | let key = redis_key(format!("gw:{{{}}}:stream:frame", req.gateway_id)); |
| 723 | let (redis_tx, mut redis_rx) = mpsc::channel(1); |
| 724 | let (stream_tx, stream_rx) = mpsc::channel(1); |
| 725 | |
| 726 | let mut framelog_future = Box::pin(stream::frame::get_frame_logs(key, 10, redis_tx)); |
| 727 | let (drop_receiver, mut close_rx) = DropReceiver::new(ReceiverStream::new(stream_rx)); |
| 728 | |
| 729 | tokio::spawn(async move { |
| 730 | loop { |
| 731 | tokio::select! { |
| 732 | // detect client disconnect |
| 733 | _ = close_rx.recv() => { |
| 734 | debug!("Client disconnected"); |
| 735 | break; |
| 736 | } |
| 737 | // detect get_frame_logs function return |
| 738 | res = &mut framelog_future => { |
| 739 | match res { |
| 740 | Ok(_) => { |
| 741 | trace!("get_frame_logs returned"); |
| 742 | }, |
| 743 | Err(e) => { |
| 744 | error!("Reading frame-log returned error: {}", e); |
| 745 | stream_tx.send(Err(e.status())).await.unwrap(); |
| 746 | }, |
| 747 | } |
| 748 | break; |
| 749 | } |
| 750 | // detect stream message |
| 751 | msg = redis_rx.recv() => { |
| 752 | match msg { |
| 753 | None => { |
| 754 | trace!("Redis Stream channel has been closed"); |
| 755 | break; |
| 756 | }, |
| 757 | Some(msg) => { |
| 758 | trace!("Message received from Redis Stream channel"); |
| 759 | if stream_tx.send(Ok(msg)).await.is_err() { |
| 760 | error!("Sending message to gRPC channel error"); |
| 761 | break; |
| 762 | }; |
| 763 | }, |
| 764 | } |
| 765 | } |
nothing calls this directly
no test coverage detected