(self: Arc<Self>, mut cancel: watch::Receiver<bool>)
| 972 | } |
| 973 | |
| 974 | async fn watch_loop(self: Arc<Self>, mut cancel: watch::Receiver<bool>) { |
| 975 | loop { |
| 976 | let mut stream = match self |
| 977 | .driver |
| 978 | .watch_sandboxes(Request::new(WatchSandboxesRequest {})) |
| 979 | .await |
| 980 | { |
| 981 | Ok(response) => response.into_inner(), |
| 982 | Err(err) => { |
| 983 | warn!(error = %err, "Compute driver watch stream failed to start"); |
| 984 | tokio::select! { |
| 985 | () = tokio::time::sleep(Duration::from_secs(2)) => {} |
| 986 | _ = cancel.changed() => return, |
| 987 | } |
| 988 | continue; |
| 989 | } |
| 990 | }; |
| 991 | |
| 992 | let mut restart = false; |
| 993 | loop { |
| 994 | tokio::select! { |
| 995 | item = stream.next() => { |
| 996 | match item { |
| 997 | Some(Ok(event)) => { |
| 998 | if let Err(err) = self.apply_watch_event(event).await { |
| 999 | warn!(error = %err, "Failed to apply compute driver event"); |
| 1000 | } |
| 1001 | } |
| 1002 | Some(Err(err)) => { |
| 1003 | warn!(error = %err, "Compute driver watch stream errored"); |
| 1004 | restart = true; |
| 1005 | break; |
| 1006 | } |
| 1007 | None => break, |
| 1008 | } |
| 1009 | } |
| 1010 | _ = cancel.changed() => return, |
| 1011 | } |
| 1012 | } |
| 1013 | |
| 1014 | if !restart { |
| 1015 | warn!("Compute driver watch stream ended unexpectedly"); |
| 1016 | } |
| 1017 | tokio::select! { |
| 1018 | () = tokio::time::sleep(Duration::from_secs(2)) => {} |
| 1019 | _ = cancel.changed() => return, |
| 1020 | } |
| 1021 | } |
| 1022 | } |
| 1023 | |
| 1024 | async fn reconcile_loop(self: Arc<Self>, mut cancel: watch::Receiver<bool>) { |
| 1025 | loop { |
no test coverage detected