(
&self,
request: Request<api::StreamDeviceEventsRequest>,
)
| 841 | type StreamDeviceEventsStream = DropReceiver<Result<api::LogItem, Status>>; |
| 842 | |
| 843 | async fn stream_device_events( |
| 844 | &self, |
| 845 | request: Request<api::StreamDeviceEventsRequest>, |
| 846 | ) -> Result<Response<Self::StreamDeviceEventsStream>, Status> { |
| 847 | let req = request.get_ref(); |
| 848 | let dev_eui = EUI64::from_str(&req.dev_eui).map_err(|e| e.status())?; |
| 849 | |
| 850 | self.validator |
| 851 | .validate( |
| 852 | request.extensions(), |
| 853 | validator::ValidateDeviceAccess::new(validator::Flag::Read, dev_eui), |
| 854 | ) |
| 855 | .await?; |
| 856 | |
| 857 | let key = redis_key(format!("device:{{{}}}:stream:event", req.dev_eui)); |
| 858 | let (redis_tx, mut redis_rx) = mpsc::channel(1); |
| 859 | let (stream_tx, stream_rx) = mpsc::channel(1); |
| 860 | |
| 861 | let mut eventlog_future = Box::pin(stream::event::get_event_logs(key, 10, redis_tx)); |
| 862 | let (drop_receiver, mut close_rx) = DropReceiver::new(ReceiverStream::new(stream_rx)); |
| 863 | |
| 864 | tokio::spawn(async move { |
| 865 | loop { |
| 866 | tokio::select! { |
| 867 | // detect client disconnect |
| 868 | _ = close_rx.recv() => { |
| 869 | debug!("Client disconnected"); |
| 870 | redis_rx.close(); |
| 871 | break; |
| 872 | }, |
| 873 | // detect get_event_logs function return |
| 874 | res = &mut eventlog_future => { |
| 875 | match res { |
| 876 | Ok(_) => { |
| 877 | trace!("get_event_logs returned"); |
| 878 | }, |
| 879 | Err(e) => { |
| 880 | error!("Reading event-log returned error: {}", e); |
| 881 | stream_tx.send(Err(e.status())).await.unwrap(); |
| 882 | }, |
| 883 | } |
| 884 | break; |
| 885 | } |
| 886 | // detect stream message |
| 887 | msg = redis_rx.recv() => { |
| 888 | match msg { |
| 889 | None => { |
| 890 | trace!("Redis Stream channel has been closed"); |
| 891 | break; |
| 892 | }, |
| 893 | Some(msg) => { |
| 894 | trace!("Message received from Redis Stream channel"); |
| 895 | if stream_tx.send(Ok(msg)).await.is_err() { |
| 896 | error!("Sending message to gRPC channel error"); |
| 897 | break; |
| 898 | }; |
| 899 | }, |
| 900 | } |
nothing calls this directly
no test coverage detected