(
&self,
request: Request<api::StreamDeviceFramesRequest>,
)
| 773 | type StreamDeviceFramesStream = DropReceiver<Result<api::LogItem, Status>>; |
| 774 | |
| 775 | async fn stream_device_frames( |
| 776 | &self, |
| 777 | request: Request<api::StreamDeviceFramesRequest>, |
| 778 | ) -> Result<Response<Self::StreamDeviceFramesStream>, Status> { |
| 779 | let req = request.get_ref(); |
| 780 | let dev_eui = EUI64::from_str(&req.dev_eui).map_err(|e| e.status())?; |
| 781 | |
| 782 | self.validator |
| 783 | .validate( |
| 784 | request.extensions(), |
| 785 | validator::ValidateDeviceAccess::new(validator::Flag::Read, dev_eui), |
| 786 | ) |
| 787 | .await?; |
| 788 | |
| 789 | let key = redis_key(format!("device:{{{}}}:stream:frame", req.dev_eui)); |
| 790 | let (redis_tx, mut redis_rx) = mpsc::channel(1); |
| 791 | let (stream_tx, stream_rx) = mpsc::channel(1); |
| 792 | |
| 793 | let mut framelog_future = Box::pin(stream::frame::get_frame_logs(key, 10, redis_tx)); |
| 794 | let (drop_receiver, mut close_rx) = DropReceiver::new(ReceiverStream::new(stream_rx)); |
| 795 | |
| 796 | tokio::spawn(async move { |
| 797 | loop { |
| 798 | tokio::select! { |
| 799 | // detect client disconnect |
| 800 | _ = close_rx.recv() => { |
| 801 | debug!("Client disconnected"); |
| 802 | redis_rx.close(); |
| 803 | break; |
| 804 | } |
| 805 | // detect get_frame_logs function return |
| 806 | res = &mut framelog_future => { |
| 807 | match res { |
| 808 | Ok(_) => { |
| 809 | trace!("get_frame_logs returned"); |
| 810 | }, |
| 811 | Err(e) => { |
| 812 | error!("Reading frame-log returned error: {}", e); |
| 813 | stream_tx.send(Err(e.status())).await.unwrap(); |
| 814 | }, |
| 815 | } |
| 816 | break; |
| 817 | } |
| 818 | // detect stream message |
| 819 | msg = redis_rx.recv() => { |
| 820 | match msg { |
| 821 | None => { |
| 822 | trace!("Redis Stream channel has been closed"); |
| 823 | break; |
| 824 | }, |
| 825 | Some(msg) => { |
| 826 | trace!("Message received from Redis Stream channel"); |
| 827 | if stream_tx.send(Ok(msg)).await.is_err() { |
| 828 | error!("Sending message to gRPC channel error"); |
| 829 | break; |
| 830 | }; |
| 831 | }, |
| 832 | } |
nothing calls this directly
no test coverage detected