diff --git a/crates/openshell-driver-kubernetes/Cargo.toml b/crates/openshell-driver-kubernetes/Cargo.toml index 714b7d05c9..f72fae635d 100644 --- a/crates/openshell-driver-kubernetes/Cargo.toml +++ b/crates/openshell-driver-kubernetes/Cargo.toml @@ -38,6 +38,7 @@ notify = "8" [dev-dependencies] temp-env = "0.3" +tokio = { workspace = true, features = ["test-util"] } toml = { workspace = true } [lints] diff --git a/crates/openshell-driver-kubernetes/README.md b/crates/openshell-driver-kubernetes/README.md index ec37c4e67a..f5f0081bb8 100644 --- a/crates/openshell-driver-kubernetes/README.md +++ b/crates/openshell-driver-kubernetes/README.md @@ -47,7 +47,9 @@ state and platform events into the shared compute-driver protobuf surface used by the gateway. Kubernetes API calls use explicit timeouts so gRPC handlers do not block -indefinitely when the API server is slow or unavailable. +indefinitely when the API server is slow or unavailable. Resource and Event +watches recover in place with API-friendly backoff after transient watcher +errors, avoiding a gateway-side watch restart and its associated watch gap. ## Workspace Persistence diff --git a/crates/openshell-driver-kubernetes/src/driver.rs b/crates/openshell-driver-kubernetes/src/driver.rs index ddc7fe2a4f..6f4244c1ac 100644 --- a/crates/openshell-driver-kubernetes/src/driver.rs +++ b/crates/openshell-driver-kubernetes/src/driver.rs @@ -26,6 +26,7 @@ use kube::api::{ }; use kube::core::gvk::GroupVersionKind; use kube::core::{DynamicObject, ObjectMeta}; +use kube::runtime::WatchStreamExt; use kube::runtime::watcher::{self, Event}; use kube::{Client, Error as KubeError}; use openshell_core::driver_mounts; @@ -1782,8 +1783,16 @@ impl KubernetesComputeDriver { .await?; let event_api: Api = Api::namespaced(self.watch_client.clone(), &namespace); let watcher_config = watcher::Config::default().labels(&openshell_sandbox_label_selector()); - let mut sandbox_stream = watcher::watcher(agent_sandbox_api.api, watcher_config).boxed(); - let mut event_stream = watcher::watcher(event_api, watcher::Config::default()).boxed(); + let mut sandbox_stream = recovering_watcher_stream( + watcher::watcher(agent_sandbox_api.api, watcher_config), + "sandbox-resource", + ) + .boxed(); + let mut event_stream = recovering_watcher_stream( + watcher::watcher(event_api, watcher::Config::default()), + "kubernetes-event", + ) + .boxed(); let (tx, rx) = mpsc::channel(256); tokio::spawn(async move { @@ -1792,8 +1801,8 @@ impl KubernetesComputeDriver { loop { tokio::select! { - result = sandbox_stream.try_next() => match result { - Ok(Some(Event::Applied(obj))) => { + event = sandbox_stream.next() => match event { + Some(Event::Applied(obj)) => { if let Ok((kube_name, sandbox)) = sandbox_from_object(&namespace, obj) { update_indexes(&mut sandbox_name_to_id, &mut agent_pod_to_id, &kube_name, &sandbox); let event = WatchSandboxesEvent { @@ -1806,7 +1815,7 @@ impl KubernetesComputeDriver { } } } - Ok(Some(Event::Deleted(obj))) => { + Some(Event::Deleted(obj)) => { if is_openshell_managed(&obj) && let Ok(sandbox_id) = sandbox_id_from_object(&obj) { @@ -1821,7 +1830,7 @@ impl KubernetesComputeDriver { } } } - Ok(Some(Event::Restarted(objs))) => { + Some(Event::Restarted(objs)) => { for obj in objs { if let Ok((kube_name, sandbox)) = sandbox_from_object(&namespace, obj) { update_indexes(&mut sandbox_name_to_id, &mut agent_pod_to_id, &kube_name, &sandbox); @@ -1836,19 +1845,15 @@ impl KubernetesComputeDriver { } } } - Ok(None) => { + None => { let _ = tx.send(Err(KubernetesDriverError::Message( "sandbox watcher stream ended unexpectedly".to_string() ))).await; break; } - Err(err) => { - let _ = tx.send(Err(KubernetesDriverError::Message(err.to_string()))).await; - break; - } }, - result = event_stream.try_next() => match result { - Ok(Some(Event::Applied(obj))) => { + event = event_stream.next() => match event { + Some(Event::Applied(obj)) => { if let Some((sandbox_id, event)) = map_kube_event_to_platform( &sandbox_name_to_id, &agent_pod_to_id, @@ -1864,20 +1869,16 @@ impl KubernetesComputeDriver { } } } - Ok(Some(Event::Deleted(_))) => {} - Ok(Some(Event::Restarted(_))) => { + Some(Event::Deleted(_)) => {} + Some(Event::Restarted(_)) => { debug!(namespace = %namespace, "Kubernetes event watcher restarted"); } - Ok(None) => { + None => { let _ = tx.send(Err(KubernetesDriverError::Message( "kubernetes event watcher stream ended".to_string() ))).await; break; } - Err(err) => { - let _ = tx.send(Err(KubernetesDriverError::Message(err.to_string()))).await; - break; - } }, () = tx.closed() => break, } @@ -1895,15 +1896,59 @@ impl KubernetesComputeDriver { Self::cluster_wide_sandbox_api(self.watch_client.clone(), sandbox_api_version); let selector = self.openshell_sandbox_selector(); let watcher_config = watcher::Config::default().labels(&selector); - let mut sandbox_stream = watcher::watcher(cluster_api.api, watcher_config).boxed(); - let (tx, rx) = mpsc::channel(256); - let default_namespace = self.config.namespace.clone(); + let sandbox_stream = recovering_watcher_stream( + watcher::watcher(cluster_api.api, watcher_config), + "sandbox-resource", + ) + .boxed(); - tokio::spawn(async move { - loop { - tokio::select! { - result = sandbox_stream.try_next() => match result { - Ok(Some(Event::Applied(obj))) => { + Ok(cluster_wide_watch_stream( + sandbox_stream, + self.config.namespace.clone(), + )) + } +} + +fn cluster_wide_watch_stream(mut sandbox_stream: S, default_namespace: String) -> WatchStream +where + S: Stream> + Send + Unpin + 'static, +{ + let (tx, rx) = mpsc::channel(256); + + tokio::spawn(async move { + loop { + tokio::select! { + event = sandbox_stream.next() => match event { + Some(Event::Applied(obj)) => { + let ns = obj.metadata.namespace.clone() + .unwrap_or_else(|| default_namespace.clone()); + if let Ok((_kube_name, sandbox)) = sandbox_from_object(&ns, obj) { + let event = WatchSandboxesEvent { + payload: Some(watch_sandboxes_event::Payload::Sandbox( + WatchSandboxesSandboxEvent { sandbox: Some(sandbox) } + )), + }; + if tx.send(Ok(event)).await.is_err() { + break; + } + } + } + Some(Event::Deleted(obj)) => { + if is_openshell_managed(&obj) + && let Ok(sandbox_id) = sandbox_id_from_object(&obj) + { + let event = WatchSandboxesEvent { + payload: Some(watch_sandboxes_event::Payload::Deleted( + WatchSandboxesDeletedEvent { sandbox_id } + )), + }; + if tx.send(Ok(event)).await.is_err() { + break; + } + } + } + Some(Event::Restarted(objs)) => { + for obj in objs { let ns = obj.metadata.namespace.clone() .unwrap_or_else(|| default_namespace.clone()); if let Ok((_kube_name, sandbox)) = sandbox_from_object(&ns, obj) { @@ -1913,58 +1958,61 @@ impl KubernetesComputeDriver { )), }; if tx.send(Ok(event)).await.is_err() { - break; - } - } - } - Ok(Some(Event::Deleted(obj))) => { - if is_openshell_managed(&obj) - && let Ok(sandbox_id) = sandbox_id_from_object(&obj) - { - let event = WatchSandboxesEvent { - payload: Some(watch_sandboxes_event::Payload::Deleted( - WatchSandboxesDeletedEvent { sandbox_id } - )), - }; - if tx.send(Ok(event)).await.is_err() { - break; + return; } } } - Ok(Some(Event::Restarted(objs))) => { - for obj in objs { - let ns = obj.metadata.namespace.clone() - .unwrap_or_else(|| default_namespace.clone()); - if let Ok((_kube_name, sandbox)) = sandbox_from_object(&ns, obj) { - let event = WatchSandboxesEvent { - payload: Some(watch_sandboxes_event::Payload::Sandbox( - WatchSandboxesSandboxEvent { sandbox: Some(sandbox) } - )), - }; - if tx.send(Ok(event)).await.is_err() { - return; - } - } - } - } - Ok(None) => { - let _ = tx.send(Err(KubernetesDriverError::Message( - "sandbox watcher stream ended unexpectedly".to_string() - ))).await; - break; - } - Err(err) => { - let _ = tx.send(Err(KubernetesDriverError::Message(err.to_string()))).await; - break; - } - }, - () = tx.closed() => break, - } + } + None => { + let _ = tx.send(Err(KubernetesDriverError::Message( + "sandbox watcher stream ended unexpectedly".to_string() + ))).await; + break; + } + }, + () = tx.closed() => break, } - }); + } + }); - Ok(Box::pin(ReceiverStream::new(rx))) - } + Box::pin(ReceiverStream::new(rx)) +} + +fn recovering_watcher_stream( + stream: S, + watcher: &'static str, +) -> impl Stream> +where + S: Stream, E>>, + E: std::fmt::Display, +{ + continue_on_watcher_errors(stream.default_backoff(), watcher) +} + +/// Drop kube-runtime watcher errors after logging them so continued polling can +/// drive its built-in relist and recovery state machine. The production adapter +/// above applies backoff first to avoid hot-looping on persistent API failures. +fn continue_on_watcher_errors( + stream: S, + watcher: &'static str, +) -> impl Stream> +where + S: Stream, E>>, + E: std::fmt::Display, +{ + stream.filter_map(move |result| { + futures::future::ready(match result { + Ok(event) => Some(event), + Err(err) => { + warn!( + watcher, + error = %err, + "Kubernetes watcher stream error; waiting for kube-runtime recovery" + ); + None + } + }) + }) } fn should_try_next_sandbox_api_version(err: &KubeError) -> bool { @@ -4568,6 +4616,119 @@ mod tests { }) } + fn expired_watch_error() -> watcher::Error { + watcher::Error::WatchError(kube::core::ErrorResponse { + status: "Failure".to_string(), + message: "too old resource version".to_string(), + reason: "Expired".to_string(), + code: 410, + }) + } + + #[tokio::test] + async fn sandbox_watcher_error_does_not_hide_restarted_recovery_event() { + let recovered = DynamicObject { + types: None, + metadata: ObjectMeta { + name: Some("recovered-sandbox".to_string()), + ..Default::default() + }, + data: serde_json::json!({}), + }; + let source = futures::stream::iter([ + Err(expired_watch_error()), + Ok(Event::Restarted(vec![recovered])), + ]); + let mut stream = continue_on_watcher_errors(source, "sandbox-resource"); + + let event = stream + .next() + .await + .expect("410 Expired must not terminate the watcher stream"); + let Event::Restarted(objects) = event else { + panic!("expected kube-runtime recovery to emit Restarted"); + }; + assert_eq!(objects.len(), 1); + assert_eq!( + objects[0].metadata.name.as_deref(), + Some("recovered-sandbox") + ); + assert!( + stream.next().await.is_none(), + "source closure must be preserved" + ); + } + + #[tokio::test(start_paused = true)] + async fn outward_watch_stream_survives_expired_error_and_backoff_recovery() { + let recovered = DynamicObject { + types: None, + metadata: ObjectMeta { + name: Some("recovered-sandbox".to_string()), + namespace: Some("recovered-namespace".to_string()), + labels: Some(BTreeMap::from([ + (LABEL_SANDBOX_ID.to_string(), "sandbox-id".to_string()), + (LABEL_SANDBOX_NAME.to_string(), "sandbox-name".to_string()), + (LABEL_SANDBOX_WORKSPACE.to_string(), "workspace".to_string()), + ( + LABEL_MANAGED_BY.to_string(), + LABEL_MANAGED_BY_VALUE.to_string(), + ), + ])), + ..Default::default() + }, + data: serde_json::json!({}), + }; + let source = futures::stream::iter([ + Err(expired_watch_error()), + Ok(Event::Restarted(vec![recovered])), + ]) + .chain(futures::stream::pending()); + let sandbox_stream = recovering_watcher_stream(source, "sandbox-resource").boxed(); + let mut outward = cluster_wide_watch_stream(sandbox_stream, "default".to_string()); + + let event = outward + .next() + .await + .expect("outward stream must stay open through recovery") + .expect("recoverable watcher error must not reach the outward stream"); + let Some(watch_sandboxes_event::Payload::Sandbox(event)) = event.payload else { + panic!("expected recovered sandbox event"); + }; + let sandbox = event.sandbox.expect("sandbox payload must be populated"); + assert_eq!(sandbox.id, "sandbox-id"); + assert_eq!(sandbox.namespace, "recovered-namespace"); + + let next = outward.next(); + futures::pin_mut!(next); + assert!( + futures::poll!(next).is_pending(), + "outward stream must remain open after the recovered event" + ); + } + + #[tokio::test] + async fn kubernetes_event_watcher_error_does_not_hide_restarted_recovery_event() { + let source = futures::stream::iter([ + Err(expired_watch_error()), + Ok(Event::Restarted(vec![KubeEventObj::default()])), + ]); + let mut stream = continue_on_watcher_errors(source, "kubernetes-event"); + + let event = stream + .next() + .await + .expect("410 Expired must not terminate the watcher stream"); + let Event::Restarted(events) = event else { + panic!("expected kube-runtime recovery to emit Restarted"); + }; + assert_eq!(events.len(), 1); + assert!( + stream.next().await.is_none(), + "source closure must be preserved" + ); + } + #[test] fn sandbox_api_version_probe_retries_on_structured_and_raw_404() { let structured = kube_api_error(404, "could not find the requested resource");