Skip to main content

vynil_core/
k8s.rs

1//! Kubernetes handlers.
2//!
3//! Feature `k8s` only (implies `rhai`). Provides `K8sGeneric` (CRUD via SSA + discovery cache),
4//! `K8sObject` (per-object helpers), `K8sRaw` (raw API discovery) and typed workload helpers
5//! (`K8sDeploy`, `K8sJob`, …). The API is Rhai-facing throughout (methods return [`crate::RhaiRes`]).
6//!
7//! The discovery cache is wired via `OnceLock` function pointers injected by the consumer (see
8//! `context_is_wired` / `set_get_client` …) — the crate itself does not assume a kubeconfig.
9
10use std::sync::OnceLock;
11
12use crate::{Error, Result, RhaiRes, rhai_err, rhai_err_str};
13use k8s_openapi::api::{
14    apps::v1::{DaemonSet, Deployment, StatefulSet},
15    batch::v1::Job,
16};
17use kube::{
18    Client, ResourceExt,
19    api::{
20        Api, DeleteParams, DynamicObject, ListParams, ObjectList, PartialObjectMeta, Patch, PatchParams,
21        PostParams,
22    },
23    discovery::{ApiCapabilities, ApiResource, Discovery, Scope},
24    runtime::wait::{Condition, await_condition, conditions},
25};
26use rhai::{Dynamic, Engine, serde::to_dynamic};
27use serde_json::json;
28use tokio::sync::RwLock;
29
30// ── Context function pointers (initialized by common at startup) ─────────────
31
32pub static GET_CLIENT: OnceLock<Box<dyn Fn() -> Client + Send + Sync>> = OnceLock::new();
33pub static GET_LABELS: OnceLock<Box<dyn Fn() -> Option<serde_json::Value> + Send + Sync>> = OnceLock::new();
34pub static GET_OWNER: OnceLock<Box<dyn Fn() -> Option<serde_json::Value> + Send + Sync>> = OnceLock::new();
35pub static GET_OWNER_NS: OnceLock<Box<dyn Fn() -> Option<String> + Send + Sync>> = OnceLock::new();
36
37pub fn set_get_client(f: Box<dyn Fn() -> Client + Send + Sync>) {
38    GET_CLIENT.set(f).ok();
39}
40pub fn set_get_labels(f: Box<dyn Fn() -> Option<serde_json::Value> + Send + Sync>) {
41    GET_LABELS.set(f).ok();
42}
43pub fn set_get_owner(f: Box<dyn Fn() -> Option<serde_json::Value> + Send + Sync>) {
44    GET_OWNER.set(f).ok();
45}
46pub fn set_get_owner_ns(f: Box<dyn Fn() -> Option<String> + Send + Sync>) {
47    GET_OWNER_NS.set(f).ok();
48}
49
50/// True if the client name (crate::set_client_name) and the 4 context accessors have been
51/// injected (see common::context::wire_core_k8s).
52pub fn context_is_wired() -> bool {
53    GET_CLIENT.get().is_some()
54        && crate::client_name_is_set()
55        && GET_LABELS.get().is_some()
56        && GET_OWNER.get().is_some()
57        && GET_OWNER_NS.get().is_some()
58}
59
60fn call_get_labels() -> Option<serde_json::Value> {
61    GET_LABELS.get().map(|f| f()).unwrap_or(None)
62}
63fn call_get_owner() -> Option<serde_json::Value> {
64    GET_OWNER.get().map(|f| f()).unwrap_or(None)
65}
66fn call_get_owner_ns() -> Option<String> {
67    GET_OWNER_NS.get().map(|f| f()).unwrap_or(None)
68}
69
70// ── k8sgeneric ───────────────────────────────────────────────────────────────
71
72lazy_static::lazy_static! {
73    pub static ref CLIENT: Client = {
74        let f = GET_CLIENT.get().expect("k8s context not initialized");
75        tokio::task::block_in_place(|| {
76            tokio::runtime::Handle::current().block_on(async move { f() })
77        })
78    };
79}
80
81fn aggregated_apiservice_group(spec: &serde_json::Value) -> Option<String> {
82    spec.get("service").filter(|s| !s.is_null())?;
83    spec.get("group")
84        .and_then(|g| g.as_str())
85        .filter(|g| !g.is_empty())
86        .map(|g| g.to_string())
87}
88
89async fn excluded_apiservice_groups() -> Vec<String> {
90    let ar = ApiResource {
91        group: "apiregistration.k8s.io".to_string(),
92        version: "v1".to_string(),
93        api_version: "apiregistration.k8s.io/v1".to_string(),
94        kind: "APIService".to_string(),
95        plural: "apiservices".to_string(),
96    };
97    let api: Api<DynamicObject> = Api::all_with(CLIENT.clone(), &ar);
98    match api.list(&ListParams::default()).await {
99        Ok(list) => list
100            .items
101            .iter()
102            .filter_map(|obj| aggregated_apiservice_group(obj.data.get("spec")?))
103            .collect(),
104        Err(e) => {
105            tracing::warn!("E_DISCOVERY_WARN: cannot list APIServices ({e}), proceeding without exclusions");
106            vec![]
107        }
108    }
109}
110
111async fn async_populate_cache() -> Discovery {
112    let excluded = excluded_apiservice_groups().await;
113    let excluded_refs: Vec<&str> = excluded.iter().map(|s| s.as_str()).collect();
114    Discovery::new(CLIENT.clone())
115        .exclude(&excluded_refs)
116        .run()
117        .await
118        .expect("create discovery (excluding api-services)")
119}
120
121fn populate_cache() -> Discovery {
122    tokio::task::block_in_place(|| {
123        tokio::runtime::Handle::current().block_on(async move {
124            let excluded = excluded_apiservice_groups().await;
125            let excluded_refs: Vec<&str> = excluded.iter().map(|s| s.as_str()).collect();
126            Discovery::new(CLIENT.clone())
127                .exclude(&excluded_refs)
128                .run()
129                .await
130                .expect("create discovery (excluding api-services)")
131        })
132    })
133}
134
135lazy_static::lazy_static! {
136    pub static ref CACHE: RwLock<Discovery> = RwLock::new(populate_cache());
137}
138
139pub fn update_cache() {
140    tokio::task::block_in_place(|| {
141        tokio::runtime::Handle::current().block_on(async move {
142            match tokio::time::timeout(std::time::Duration::from_secs(60), async_populate_cache()).await {
143                Ok(discovery) => {
144                    *CACHE.write().await = discovery;
145                }
146                Err(_) => {
147                    tracing::warn!(
148                        "E_DISCOVERY_TIMEOUT: update_k8s_crd_cache exceeded 30s, keeping old cache"
149                    );
150                }
151            }
152        })
153    })
154}
155
156type DynObjCondition = Box<dyn Fn(&DynamicObject) -> Result<bool, Box<rhai::EvalAltResult>>>;
157
158#[derive(Clone, Debug)]
159pub struct K8sObject {
160    pub api: Api<DynamicObject>,
161    pub obj: PartialObjectMeta,
162    pub kind: String,
163}
164impl K8sObject {
165    pub fn rhai_delete(&mut self) -> RhaiRes<()> {
166        tokio::task::block_in_place(|| {
167            tokio::runtime::Handle::current().block_on(async move {
168                self.api
169                    .delete(&self.obj.name_any(), &DeleteParams::foreground())
170                    .await
171                    .map_err(Error::KubeError)
172                    .map(|_| ())
173            })
174        })
175        .map_err(rhai_err)
176    }
177
178    pub fn rhai_wait_deleted(&mut self, timeout: i64) -> RhaiRes<()> {
179        let name = self.obj.name_any();
180        let uid = self
181            .obj
182            .uid()
183            .ok_or_else(|| rhai_err_str(format!("cannot wait for deletion of {name}: uid is missing")))?;
184        tokio::task::block_in_place(|| {
185            tokio::runtime::Handle::current().block_on(async move {
186                let cond = await_condition(self.api.clone(), &name, conditions::is_deleted(&uid));
187                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
188                    .await
189                    .map_err(Error::Elapsed)
190            })
191        })
192        .map_err(rhai_err)?
193        .map_err(|e| rhai_err(Error::KubeWaitError(e)))
194        .map(|_| ())
195    }
196
197    pub fn get_metadata(&mut self) -> RhaiRes<Dynamic> {
198        let v = serde_json::to_value(self.obj.metadata.clone())
199            .map_err(|e| rhai_err(Error::SerializationError(e)))?;
200        to_dynamic(v)
201    }
202
203    pub fn get_kind(&mut self) -> String {
204        if let Some(t) = self.obj.types.clone() {
205            t.kind
206        } else {
207            "".to_string()
208        }
209    }
210
211    pub fn original_kind(&mut self) -> String {
212        self.kind.clone()
213    }
214
215    pub fn is_condition(cond: String) -> impl Condition<DynamicObject> {
216        move |obj: Option<&DynamicObject>| {
217            let Some(conditions) = obj
218                .and_then(|o| o.data.get("status"))
219                .and_then(|s| s.get("conditions"))
220                .and_then(|c| c.as_array())
221            else {
222                return false;
223            };
224            conditions.iter().any(|c| {
225                c.get("type").and_then(|t| t.as_str()) == Some(cond.as_str())
226                    && c.get("status").and_then(|s| s.as_str()) == Some("True")
227            })
228        }
229    }
230
231    pub fn wait_condition(&mut self, condition: String, timeout: i64) -> RhaiRes<()> {
232        let name = self.obj.name_any();
233        let cond = await_condition(self.api.clone(), &name, Self::is_condition(condition));
234        tokio::task::block_in_place(|| {
235            tokio::runtime::Handle::current().block_on(async move {
236                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
237                    .await
238                    .map_err(Error::Elapsed)
239            })
240        })
241        .map_err(rhai_err)?
242        .map_err(Error::KubeWaitError)
243        .map_err(rhai_err)?;
244        Ok(())
245    }
246
247    pub fn is_status(prop: String) -> impl Condition<DynamicObject> {
248        move |obj: Option<&DynamicObject>| {
249            obj.and_then(|o| o.data.get("status"))
250                .and_then(|s| s.get(prop.as_str()))
251                .and_then(|v| v.as_bool())
252                .unwrap_or(false)
253        }
254    }
255
256    pub fn have_status(prop: String) -> impl Condition<DynamicObject> {
257        move |obj: Option<&DynamicObject>| {
258            obj.and_then(|o| o.data.get("status"))
259                .and_then(|s| s.get(prop.as_str()))
260                .map(|v| !v.is_null())
261                .unwrap_or(false)
262        }
263    }
264
265    pub fn have_status_value(prop: String, value: String) -> impl Condition<DynamicObject> {
266        move |obj: Option<&DynamicObject>| {
267            obj.and_then(|o| o.data.get("status"))
268                .and_then(|s| s.get(prop.as_str()))
269                .and_then(|v| v.as_str())
270                .map(|v| v == value.as_str())
271                .unwrap_or(false)
272        }
273    }
274
275    pub fn wait_status(&mut self, prop: String, timeout: i64) -> RhaiRes<()> {
276        let name = self.obj.name_any();
277        tracing::debug!("wait_status({}) for {} {}", &prop, self.kind, name);
278        let cond = await_condition(self.api.clone(), &name, Self::is_status(prop));
279        tokio::task::block_in_place(|| {
280            tokio::runtime::Handle::current().block_on(async move {
281                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
282                    .await
283                    .map_err(Error::Elapsed)
284            })
285        })
286        .map_err(rhai_err)?
287        .map_err(Error::KubeWaitError)
288        .map_err(rhai_err)?;
289        Ok(())
290    }
291
292    pub fn wait_status_prop(&mut self, prop: String, timeout: i64) -> RhaiRes<()> {
293        let name = self.obj.name_any();
294        tracing::debug!("wait_status({}) for {} {}", &prop, self.kind, name);
295        let cond = await_condition(self.api.clone(), &name, Self::have_status(prop));
296        tokio::task::block_in_place(|| {
297            tokio::runtime::Handle::current().block_on(async move {
298                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
299                    .await
300                    .map_err(Error::Elapsed)
301            })
302        })
303        .map_err(rhai_err)?
304        .map_err(Error::KubeWaitError)
305        .map_err(rhai_err)?;
306        Ok(())
307    }
308
309    pub fn wait_status_string(&mut self, prop: String, value: String, timeout: i64) -> RhaiRes<()> {
310        let name = self.obj.name_any();
311        tracing::debug!("wait_status({}) for {} {}", &prop, self.kind, name);
312        let cond = await_condition(self.api.clone(), &name, Self::have_status_value(prop, value));
313        tokio::task::block_in_place(|| {
314            tokio::runtime::Handle::current().block_on(async move {
315                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
316                    .await
317                    .map_err(Error::Elapsed)
318            })
319        })
320        .map_err(rhai_err)?
321        .map_err(Error::KubeWaitError)
322        .map_err(rhai_err)?;
323        Ok(())
324    }
325
326    pub fn is_for(cond: DynObjCondition) -> impl Condition<DynamicObject> {
327        move |obj: Option<&DynamicObject>| {
328            if let Some(dynobj) = &obj
329                && dynobj.data.is_object()
330            {
331                return cond(dynobj).unwrap_or_else(|e| {
332                    tracing::warn!("wait_for closure error: {:?}", e);
333                    false
334                });
335            }
336            false
337        }
338    }
339
340    pub fn wait_for(&mut self, condition: DynObjCondition, timeout: i64) -> RhaiRes<()> {
341        let name = self.obj.name_any();
342        let cond = await_condition(self.api.clone(), &name, Self::is_for(condition));
343        tokio::task::block_in_place(|| {
344            tokio::runtime::Handle::current().block_on(async move {
345                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
346                    .await
347                    .map_err(Error::Elapsed)
348            })
349        })
350        .map_err(rhai_err)?
351        .map_err(Error::KubeWaitError)
352        .map_err(rhai_err)?;
353        Ok(())
354    }
355}
356
357#[derive(Clone, Debug)]
358pub struct K8sGeneric {
359    pub api: Option<Api<DynamicObject>>,
360    pub ns: Option<String>,
361    pub scope: Scope,
362    pub kind: String,
363}
364
365impl K8sGeneric {
366    #[must_use]
367    pub fn new(name: &str, ns: Option<String>) -> K8sGeneric {
368        tokio::task::block_in_place(|| {
369            tokio::runtime::Handle::current().block_on(async move {
370                if let Some((res, cap)) = CACHE
371                    .read()
372                    .await
373                    .groups()
374                    .flat_map(|group| {
375                        group
376                            .resources_by_stability()
377                            .into_iter()
378                            .map(move |res: (ApiResource, ApiCapabilities)| (group, res))
379                    })
380                    .filter(|(_, (res, _))| {
381                        name.eq_ignore_ascii_case(&res.kind) || name.eq_ignore_ascii_case(&res.plural)
382                    })
383                    .min_by_key(|(group, _res)| group.name())
384                    .map(|(_, res)| res)
385                {
386                    tracing::debug!("K8sGeneric::new Using {}/{}/{}", res.group, res.version, res.kind);
387                    let api = if cap.scope == Scope::Cluster || ns.is_none() {
388                        Api::all_with(CLIENT.clone(), &res)
389                    } else if let Some(namespace) = ns.clone() {
390                        Api::namespaced_with(CLIENT.clone(), &namespace, &res)
391                    } else {
392                        Api::default_namespaced_with(CLIENT.clone(), &res)
393                    };
394                    K8sGeneric {
395                        api: Some(api),
396                        ns,
397                        scope: cap.scope,
398                        kind: res.kind,
399                    }
400                } else {
401                    K8sGeneric {
402                        api: None,
403                        ns: None,
404                        scope: Scope::Cluster,
405                        kind: String::new(),
406                    }
407                }
408            })
409        })
410    }
411
412    #[must_use]
413    pub fn new_api_version(api_group: &str, version: &str, name: &str, ns: Option<String>) -> K8sGeneric {
414        tokio::task::block_in_place(|| {
415            tokio::runtime::Handle::current().block_on(async move {
416                if let Some((res, cap)) = CACHE
417                    .read()
418                    .await
419                    .groups()
420                    .flat_map(|group| {
421                        group
422                            .resources_by_stability()
423                            .into_iter()
424                            .map(move |res: (ApiResource, ApiCapabilities)| (group, res))
425                    })
426                    .filter(|(group, (res, _))| {
427                        group.name() == api_group
428                            && res.version == version
429                            && (name.eq_ignore_ascii_case(&res.kind)
430                                || name.eq_ignore_ascii_case(&res.plural))
431                    })
432                    .min_by_key(|(group, _res)| group.name())
433                    .map(|(_, res)| res)
434                {
435                    tracing::debug!(
436                        "K8sGeneric::new_api_version Using {}/{}/{}",
437                        res.group,
438                        res.version,
439                        res.kind
440                    );
441                    let api = if cap.scope == Scope::Cluster || ns.is_none() {
442                        Api::all_with(CLIENT.clone(), &res)
443                    } else if let Some(namespace) = ns.clone() {
444                        Api::namespaced_with(CLIENT.clone(), &namespace, &res)
445                    } else {
446                        Api::default_namespaced_with(CLIENT.clone(), &res)
447                    };
448                    K8sGeneric {
449                        api: Some(api),
450                        ns,
451                        scope: cap.scope,
452                        kind: res.kind,
453                    }
454                } else {
455                    K8sGeneric {
456                        api: None,
457                        ns: None,
458                        scope: Scope::Cluster,
459                        kind: String::new(),
460                    }
461                }
462            })
463        })
464    }
465
466    pub fn new_ns(name: String, ns: String) -> K8sGeneric {
467        K8sGeneric::new(name.as_str(), Some(ns))
468    }
469
470    pub fn new_global(name: String) -> K8sGeneric {
471        K8sGeneric::new(name.as_str(), None)
472    }
473
474    pub fn new_group_ns(api_version: String, name: String, ns: String) -> K8sGeneric {
475        let arr = api_version.split("/").collect::<Vec<&str>>();
476        if arr.len() > 1 {
477            K8sGeneric::new_api_version(arr[0], arr[1], name.as_str(), Some(ns))
478        } else {
479            K8sGeneric::new(name.as_str(), Some(ns))
480        }
481    }
482
483    pub fn rhai_get_scope(&mut self) -> String {
484        if self.scope == Scope::Cluster {
485            "cluster".to_string()
486        } else {
487            "namespace".to_string()
488        }
489    }
490
491    pub fn exist(&self) -> bool {
492        self.api.is_some()
493    }
494
495    pub fn rhai_exist(&mut self) -> RhaiRes<Dynamic> {
496        to_dynamic(self.api.is_some())
497    }
498
499    pub fn list(&self) -> Result<ObjectList<DynamicObject>> {
500        if let Some(api) = self.api.clone() {
501            tokio::task::block_in_place(|| {
502                tokio::runtime::Handle::current()
503                    .block_on(async move { api.list(&ListParams::default()).await.map_err(Error::KubeError) })
504            })
505        } else {
506            Err(Error::UnsupportedMethod)
507        }
508    }
509
510    pub fn rhai_list(&mut self) -> RhaiRes<Dynamic> {
511        let res = self.list().map_err(rhai_err)?;
512        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
513        to_dynamic(v)
514    }
515
516    pub fn list_labels(&self, labels: String) -> Result<ObjectList<DynamicObject>> {
517        if let Some(api) = self.api.clone() {
518            tokio::task::block_in_place(|| {
519                tokio::runtime::Handle::current().block_on(async move {
520                    let mut lp = ListParams::default();
521                    lp = lp.labels(&labels);
522                    api.list(&lp).await.map_err(Error::KubeError)
523                })
524            })
525        } else {
526            Err(Error::UnsupportedMethod)
527        }
528    }
529
530    pub fn rhai_list_labels(&mut self, labels: String) -> RhaiRes<Dynamic> {
531        let res = self.list_labels(labels).map_err(rhai_err)?;
532        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
533        to_dynamic(v)
534    }
535
536    pub fn list_meta(&self) -> Result<ObjectList<PartialObjectMeta>> {
537        if let Some(api) = self.api.clone() {
538            tokio::task::block_in_place(|| {
539                tokio::runtime::Handle::current().block_on(async move {
540                    api.list_metadata(&ListParams::default())
541                        .await
542                        .map_err(Error::KubeError)
543                })
544            })
545        } else {
546            Err(Error::UnsupportedMethod)
547        }
548    }
549
550    pub fn rhai_list_meta(&mut self) -> RhaiRes<Dynamic> {
551        let res = self.list_meta().map_err(rhai_err)?;
552        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
553        to_dynamic(v)
554    }
555
556    pub fn get(&self, name: &str) -> Result<DynamicObject> {
557        if let Some(api) = self.api.clone() {
558            tokio::task::block_in_place(|| {
559                tokio::runtime::Handle::current()
560                    .block_on(async move { api.get(name).await.map_err(Error::KubeError) })
561            })
562        } else {
563            Err(Error::UnsupportedMethod)
564        }
565    }
566
567    pub fn rhai_get(&mut self, name: String) -> RhaiRes<Dynamic> {
568        let res = self.get(&name).map_err(rhai_err)?;
569        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
570        to_dynamic(v)
571    }
572
573    pub fn get_meta(&self, name: &str) -> Result<PartialObjectMeta> {
574        if let Some(api) = self.api.clone() {
575            tokio::task::block_in_place(|| {
576                tokio::runtime::Handle::current()
577                    .block_on(async move { api.get_metadata(name).await.map_err(Error::KubeError) })
578            })
579        } else {
580            Err(Error::UnsupportedMethod)
581        }
582    }
583
584    pub fn rhai_get_meta(&mut self, name: String) -> RhaiRes<Dynamic> {
585        let res = self.get_meta(&name).map_err(rhai_err)?;
586        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
587        to_dynamic(v)
588    }
589
590    pub fn rhai_get_obj(&mut self, name: String) -> RhaiRes<K8sObject> {
591        let Some(api) = self.api.clone() else {
592            return Err(rhai_err(Error::UnsupportedMethod));
593        };
594        let res = self.get_meta(&name).map_err(rhai_err)?;
595        Ok(K8sObject {
596            api,
597            obj: res,
598            kind: self.kind.clone(),
599        })
600    }
601
602    pub fn delete(&self, name: &str) -> Result<()> {
603        if let Some(api) = self.api.clone() {
604            tokio::task::block_in_place(|| {
605                tokio::runtime::Handle::current().block_on(async move {
606                    api.delete(name, &DeleteParams::foreground())
607                        .await
608                        .map_err(Error::KubeError)
609                        .map(|_| ())
610                })
611            })
612        } else {
613            Err(Error::UnsupportedMethod)
614        }
615    }
616
617    pub fn rhai_delete(&mut self, name: String) -> RhaiRes<()> {
618        self.delete(&name).map_err(rhai_err)
619    }
620
621    fn inject_labels_and_owner(
622        &self,
623        handle: serde_json::Map<String, serde_json::Value>,
624    ) -> serde_json::Map<String, serde_json::Value> {
625        prepare_handle(
626            handle,
627            call_get_labels(),
628            call_get_owner(),
629            call_get_owner_ns(),
630            self.ns.clone(),
631            self.scope == Scope::Namespaced,
632        )
633    }
634}
635
636// ── Pure helper — testable without a live K8s client ─────────────────────────
637
638fn prepare_handle(
639    mut handle: serde_json::Map<String, serde_json::Value>,
640    labels: Option<serde_json::Value>,
641    owner: Option<serde_json::Value>,
642    owner_ns: Option<String>,
643    my_ns: Option<String>,
644    is_namespaced: bool,
645) -> serde_json::Map<String, serde_json::Value> {
646    if !handle.contains_key("metadata") || !handle["metadata"].is_object() {
647        handle.insert("metadata".to_string(), json!({}));
648    }
649    if let Some(labels) = labels {
650        if !handle["metadata"].as_object().unwrap().contains_key("labels") {
651            handle["metadata"]
652                .as_object_mut()
653                .unwrap()
654                .insert("labels".to_string(), json!({}));
655        } else if !handle["metadata"].as_object_mut().unwrap()["labels"].is_object() {
656            handle["metadata"].as_object_mut().unwrap().remove_entry("labels");
657            handle["metadata"]
658                .as_object_mut()
659                .unwrap()
660                .insert("labels".to_string(), json!({}));
661        }
662        if let Some(label_map) = labels.as_object() {
663            for (k, v) in label_map {
664                if !handle["metadata"].as_object_mut().unwrap()["labels"]
665                    .as_object_mut()
666                    .unwrap()
667                    .keys()
668                    .any(|name| name == k)
669                {
670                    handle["metadata"].as_object_mut().unwrap()["labels"]
671                        .as_object_mut()
672                        .unwrap()
673                        .insert(k.to_string(), v.clone());
674                }
675            }
676        }
677    }
678    if is_namespaced
679        && let Some(owner) = owner
680        && let Some(ns) = owner_ns
681        && let Some(mine) = my_ns
682        && ns == mine
683    {
684        if handle["metadata"]
685            .as_object()
686            .unwrap()
687            .contains_key("ownerReferences")
688        {
689            handle["metadata"].as_object_mut().unwrap()["ownerReferences"]
690                .as_array_mut()
691                .unwrap()
692                .push(owner);
693        } else {
694            handle["metadata"]
695                .as_object_mut()
696                .unwrap()
697                .insert("ownerReferences".to_string(), vec![owner].into());
698        }
699    }
700    handle
701}
702
703impl K8sGeneric {
704    pub fn create(&self, data: serde_json::Map<String, serde_json::Value>) -> Result<DynamicObject> {
705        if let Some(api) = self.api.clone() {
706            let handle = self.inject_labels_and_owner(data);
707            tokio::task::block_in_place(|| {
708                tokio::runtime::Handle::current().block_on(async move {
709                    match serde_json::from_value(handle.into()) {
710                        Ok(obj) => api
711                            .create(&PostParams::default(), &obj)
712                            .await
713                            .map_err(Error::KubeError),
714                        Err(e) => Err(Error::SerializationError(e)),
715                    }
716                })
717            })
718        } else {
719            Err(Error::UnsupportedMethod)
720        }
721    }
722
723    pub fn rhai_create(&mut self, data: rhai::Dynamic) -> RhaiRes<Dynamic> {
724        let data = rhai::serde::from_dynamic(&data)?;
725        let res = self.create(data).map_err(|e: Error| rhai_err(e))?;
726        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
727        to_dynamic(v)
728    }
729
730    pub fn replace(
731        &self,
732        name: &str,
733        data: serde_json::Map<String, serde_json::Value>,
734    ) -> Result<DynamicObject> {
735        if let Some(api) = self.api.clone() {
736            let handle = self.inject_labels_and_owner(data);
737            tokio::task::block_in_place(|| {
738                tokio::runtime::Handle::current().block_on(async move {
739                    match serde_json::from_value(handle.into()) {
740                        Ok(obj) => api
741                            .replace(name, &PostParams::default(), &obj)
742                            .await
743                            .map_err(Error::KubeError),
744                        Err(e) => Err(Error::SerializationError(e)),
745                    }
746                })
747            })
748        } else {
749            Err(Error::UnsupportedMethod)
750        }
751    }
752
753    pub fn rhai_replace(&mut self, name: String, data: rhai::Dynamic) -> RhaiRes<Dynamic> {
754        let data = rhai::serde::from_dynamic(&data)?;
755        let res = self.replace(&name, data).map_err(|e: Error| rhai_err(e))?;
756        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
757        to_dynamic(v)
758    }
759
760    pub fn patch(
761        &self,
762        name: &str,
763        patch_data: serde_json::Map<String, serde_json::Value>,
764    ) -> Result<DynamicObject> {
765        if let Some(api) = self.api.clone() {
766            let handle = self.inject_labels_and_owner(patch_data);
767            tokio::task::block_in_place(|| {
768                tokio::runtime::Handle::current().block_on(async move {
769                    api.patch(
770                        name,
771                        &PatchParams::apply(&crate::get_client_name()).force(),
772                        &Patch::Apply(handle),
773                    )
774                    .await
775                    .map_err(Error::KubeError)
776                })
777            })
778        } else {
779            Err(Error::UnsupportedMethod)
780        }
781    }
782
783    pub fn rhai_patch(&mut self, name: String, data: rhai::Dynamic) -> RhaiRes<Dynamic> {
784        let data = rhai::serde::from_dynamic(&data)?;
785        let res = self.patch(&name, data).map_err(|e: Error| rhai_err(e))?;
786        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
787        to_dynamic(v)
788    }
789
790    pub fn apply(
791        &self,
792        name: &str,
793        patch_data: serde_json::Map<String, serde_json::Value>,
794    ) -> Result<DynamicObject> {
795        if let Some(api) = self.api.clone() {
796            let kind = patch_data
797                .get("kind")
798                .and_then(|k| k.as_str())
799                .unwrap_or("")
800                .to_string();
801            let handle = self.inject_labels_and_owner(patch_data);
802            let api_for_get = api.clone();
803            tokio::task::block_in_place(|| {
804                tokio::runtime::Handle::current().block_on(async move {
805                    match api
806                        .patch(
807                            name,
808                            &PatchParams::apply(&crate::get_client_name()).force(),
809                            &Patch::Apply(handle),
810                        )
811                        .await
812                    {
813                        Ok(obj) => Ok(obj),
814                        Err(e) => {
815                            if kind == "Job"
816                                && e.to_string().contains("immutable")
817                                && let Ok(current) = api_for_get.get(name).await
818                                && job_is_completed(&current.data)
819                            {
820                                tracing::debug!(
821                                    "Job {name} spec.template immutable but already completed — skipping"
822                                );
823                                return Ok(current);
824                            }
825                            Err(Error::KubeError(e))
826                        }
827                    }
828                })
829            })
830        } else {
831            Err(Error::UnsupportedMethod)
832        }
833    }
834
835    pub fn rhai_apply(&mut self, name: String, data: rhai::Dynamic) -> RhaiRes<Dynamic> {
836        let data = rhai::serde::from_dynamic(&data)?;
837        let res = self.apply(&name, data).map_err(|e: Error| rhai_err(e))?;
838        let v = serde_json::to_value(res).map_err(|e| rhai_err(Error::SerializationError(e)))?;
839        to_dynamic(v)
840    }
841}
842
843fn job_is_completed(data: &serde_json::Value) -> bool {
844    let status = match data.get("status") {
845        Some(s) => s,
846        None => return false,
847    };
848    if status.get("succeeded").and_then(|v| v.as_i64()).unwrap_or(0) > 0 {
849        return true;
850    }
851    status.get("completionTime").is_some()
852}
853
854// ── k8sraw ───────────────────────────────────────────────────────────────────
855
856lazy_static::lazy_static! {
857    pub static ref RAW_CLIENT: Client = {
858        let f = GET_CLIENT.get().expect("k8s context not initialized");
859        tokio::task::block_in_place(|| tokio::runtime::Handle::current().block_on(async move { f() }))
860    };
861}
862
863#[derive(Clone)]
864pub struct K8sRaw {
865    pub client: Client,
866}
867
868impl Default for K8sRaw {
869    fn default() -> Self {
870        Self::new()
871    }
872}
873
874impl K8sRaw {
875    pub fn new() -> Self {
876        Self {
877            client: RAW_CLIENT.clone(),
878        }
879    }
880
881    pub async fn get_url(&self, url: String) -> Result<serde_json::Value> {
882        let req = http::Request::get(url)
883            .body(Default::default())
884            .map_err(Error::RawHTTP)?;
885        let resp = self
886            .client
887            .request::<serde_json::Value>(req)
888            .await
889            .map_err(Error::KubeError)?;
890        Ok(resp)
891    }
892
893    pub async fn get_url_as_disco(&self, url: String) -> Result<serde_json::Value> {
894        let req = http::Request::get(url)
895            .header("Accept", "application/json;g=apidiscovery.k8s.io;v=v2;as=APIGroupDiscoveryList,application/json;g=apidiscovery.k8s.io;v=v2beta1;as=APIGroupDiscoveryList,application/json")
896            .body(Default::default()).map_err(Error::RawHTTP)?;
897        let resp = self
898            .client
899            .request::<serde_json::Value>(req)
900            .await
901            .map_err(Error::KubeError)?;
902        Ok(resp)
903    }
904
905    pub async fn get_api_version(&self) -> Result<serde_json::Value> {
906        self.get_url("/version".to_string()).await
907    }
908
909    pub async fn get_api_resources(&self) -> Result<serde_json::Value> {
910        self.get_url_as_disco("/apis".to_string()).await
911    }
912
913    pub fn rhai_get_url(&mut self, url: String) -> RhaiRes<Dynamic> {
914        tokio::task::block_in_place(|| {
915            tokio::runtime::Handle::current().block_on(async move {
916                let res = self.get_url(url).await.map_err(rhai_err)?;
917                let v = serde_json::to_string(&res)
918                    .map_err(Error::SerializationError)
919                    .map_err(rhai_err)?;
920                serde_json::from_str(&v)
921                    .map_err(Error::SerializationError)
922                    .map_err(rhai_err)
923            })
924        })
925    }
926
927    pub fn rhai_get_api_version(&mut self) -> RhaiRes<Dynamic> {
928        tokio::task::block_in_place(|| {
929            tokio::runtime::Handle::current().block_on(async move {
930                let ver = self.get_api_version().await.map_err(rhai_err)?;
931                let v = serde_json::to_string(&ver)
932                    .map_err(Error::SerializationError)
933                    .map_err(rhai_err)?;
934                serde_json::from_str(&v)
935                    .map_err(Error::SerializationError)
936                    .map_err(rhai_err)
937            })
938        })
939    }
940
941    pub fn rhai_get_api_resources(&mut self) -> RhaiRes<Dynamic> {
942        tokio::task::block_in_place(|| {
943            tokio::runtime::Handle::current().block_on(async move {
944                let ver = self.get_api_resources().await.map_err(rhai_err)?;
945                let v = serde_json::to_string(&ver)
946                    .map_err(Error::SerializationError)
947                    .map_err(rhai_err)?;
948                serde_json::from_str(&v)
949                    .map_err(Error::SerializationError)
950                    .map_err(rhai_err)
951            })
952        })
953    }
954}
955
956// ── k8sworkload ──────────────────────────────────────────────────────────────
957
958#[derive(Clone, Debug)]
959pub struct K8sDaemonSet {
960    pub api: Api<DaemonSet>,
961    pub obj: DaemonSet,
962}
963impl K8sDaemonSet {
964    pub fn is_deamonset_available() -> impl Condition<DaemonSet> {
965        |obj: Option<&DaemonSet>| {
966            if let Some(ds) = &obj
967                && let Some(s) = &ds.status
968            {
969                return s.desired_number_scheduled == s.number_available.unwrap_or(0);
970            }
971            false
972        }
973    }
974
975    pub fn get_deamonset(namespace: String, name: String) -> RhaiRes<K8sDaemonSet> {
976        let api: Api<DaemonSet> = Api::namespaced(CLIENT.clone(), &namespace);
977        let d = tokio::task::block_in_place(|| {
978            tokio::runtime::Handle::current()
979                .block_on(async move { api.get(&name).await.map_err(Error::KubeError) })
980        })
981        .map_err(rhai_err)?;
982        Ok(K8sDaemonSet {
983            api: Api::namespaced(CLIENT.clone(), &namespace),
984            obj: d,
985        })
986    }
987
988    pub fn get_metadata(&mut self) -> RhaiRes<Dynamic> {
989        let v =
990            serde_json::to_string(&self.obj.metadata).map_err(|e| rhai_err(Error::SerializationError(e)))?;
991        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
992    }
993
994    pub fn get_spec(&mut self) -> RhaiRes<Dynamic> {
995        let v = serde_json::to_string(&self.obj.spec).map_err(|e| rhai_err(Error::SerializationError(e)))?;
996        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
997    }
998
999    pub fn get_status(&mut self) -> RhaiRes<Dynamic> {
1000        let v =
1001            serde_json::to_string(&self.obj.status).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1002        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1003    }
1004
1005    pub fn wait_available(&mut self, timeout: i64) -> RhaiRes<()> {
1006        let name = self.obj.name_any();
1007        let cond = await_condition(self.api.clone(), &name, Self::is_deamonset_available());
1008        tokio::task::block_in_place(|| {
1009            tokio::runtime::Handle::current().block_on(async move {
1010                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
1011                    .await
1012                    .map_err(Error::Elapsed)
1013            })
1014        })
1015        .map_err(rhai_err)?
1016        .map_err(|e| rhai_err(Error::KubeWaitError(e)))?;
1017        Ok(())
1018    }
1019}
1020
1021#[derive(Clone, Debug)]
1022pub struct K8sStatefulSet {
1023    pub api: Api<StatefulSet>,
1024    pub obj: StatefulSet,
1025}
1026impl K8sStatefulSet {
1027    pub fn is_sts_available() -> impl Condition<StatefulSet> {
1028        |obj: Option<&StatefulSet>| {
1029            if let Some(sts) = &obj
1030                && let Some(spec) = &sts.spec
1031                && let Some(s) = &sts.status
1032            {
1033                return spec.replicas.unwrap_or(1) == s.available_replicas.unwrap_or(0);
1034            }
1035            false
1036        }
1037    }
1038
1039    pub fn get_sts(namespace: String, name: String) -> RhaiRes<K8sStatefulSet> {
1040        let api: Api<StatefulSet> = Api::namespaced(CLIENT.clone(), &namespace);
1041        let d = tokio::task::block_in_place(|| {
1042            tokio::runtime::Handle::current()
1043                .block_on(async move { api.get(&name).await.map_err(Error::KubeError) })
1044        })
1045        .map_err(rhai_err)?;
1046        Ok(K8sStatefulSet {
1047            api: Api::namespaced(CLIENT.clone(), &namespace),
1048            obj: d,
1049        })
1050    }
1051
1052    pub fn get_metadata(&mut self) -> RhaiRes<Dynamic> {
1053        let v =
1054            serde_json::to_string(&self.obj.metadata).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1055        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1056    }
1057
1058    pub fn get_spec(&mut self) -> RhaiRes<Dynamic> {
1059        let v = serde_json::to_string(&self.obj.spec).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1060        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1061    }
1062
1063    pub fn get_status(&mut self) -> RhaiRes<Dynamic> {
1064        let v =
1065            serde_json::to_string(&self.obj.status).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1066        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1067    }
1068
1069    pub fn wait_available(&mut self, timeout: i64) -> RhaiRes<()> {
1070        let name = self.obj.name_any();
1071        let cond = await_condition(self.api.clone(), &name, Self::is_sts_available());
1072        tokio::task::block_in_place(|| {
1073            tokio::runtime::Handle::current().block_on(async move {
1074                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
1075                    .await
1076                    .map_err(Error::Elapsed)
1077            })
1078        })
1079        .map_err(rhai_err)?
1080        .map_err(|e| rhai_err(Error::KubeWaitError(e)))?;
1081        Ok(())
1082    }
1083}
1084
1085#[derive(Clone, Debug)]
1086pub struct K8sDeploy {
1087    pub api: Api<Deployment>,
1088    pub obj: Deployment,
1089}
1090impl K8sDeploy {
1091    pub fn is_deploy_available() -> impl Condition<Deployment> {
1092        |obj: Option<&Deployment>| {
1093            if let Some(job) = &obj
1094                && let Some(s) = &job.status
1095                && let Some(conds) = &s.conditions
1096                && let Some(pcond) = conds.iter().find(|c| c.type_ == "Available")
1097            {
1098                return pcond.status == "True";
1099            }
1100            false
1101        }
1102    }
1103
1104    pub fn get_deployment(namespace: String, name: String) -> RhaiRes<K8sDeploy> {
1105        let api: Api<Deployment> = Api::namespaced(CLIENT.clone(), &namespace);
1106        let d = tokio::task::block_in_place(|| {
1107            tokio::runtime::Handle::current()
1108                .block_on(async move { api.get(&name).await.map_err(Error::KubeError) })
1109        })
1110        .map_err(rhai_err)?;
1111        Ok(K8sDeploy {
1112            api: Api::namespaced(CLIENT.clone(), &namespace),
1113            obj: d,
1114        })
1115    }
1116
1117    pub fn get_metadata(&mut self) -> RhaiRes<Dynamic> {
1118        let v =
1119            serde_json::to_string(&self.obj.metadata).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1120        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1121    }
1122
1123    pub fn get_spec(&mut self) -> RhaiRes<Dynamic> {
1124        let v = serde_json::to_string(&self.obj.spec).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1125        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1126    }
1127
1128    pub fn get_status(&mut self) -> RhaiRes<Dynamic> {
1129        let v =
1130            serde_json::to_string(&self.obj.status).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1131        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1132    }
1133
1134    pub fn wait_available(&mut self, timeout: i64) -> RhaiRes<()> {
1135        let name = self.obj.name_any();
1136        let cond = await_condition(self.api.clone(), &name, Self::is_deploy_available());
1137        tokio::task::block_in_place(|| {
1138            tokio::runtime::Handle::current().block_on(async move {
1139                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
1140                    .await
1141                    .map_err(Error::Elapsed)
1142            })
1143        })
1144        .map_err(rhai_err)?
1145        .map_err(|e| rhai_err(Error::KubeWaitError(e)))?;
1146        Ok(())
1147    }
1148}
1149
1150#[derive(Clone, Debug)]
1151pub struct K8sJob {
1152    pub api: Api<Job>,
1153    pub obj: Job,
1154}
1155impl K8sJob {
1156    pub fn get_job(namespace: String, name: String) -> RhaiRes<K8sJob> {
1157        let api: Api<Job> = Api::namespaced(CLIENT.clone(), &namespace);
1158        let j = tokio::task::block_in_place(|| {
1159            tokio::runtime::Handle::current()
1160                .block_on(async move { api.get(&name).await.map_err(Error::KubeError) })
1161        })
1162        .map_err(rhai_err)?;
1163        Ok(K8sJob {
1164            api: Api::namespaced(CLIENT.clone(), &namespace),
1165            obj: j,
1166        })
1167    }
1168
1169    pub fn get_metadata(&mut self) -> RhaiRes<Dynamic> {
1170        let v =
1171            serde_json::to_string(&self.obj.metadata).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1172        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1173    }
1174
1175    pub fn get_spec(&mut self) -> RhaiRes<Dynamic> {
1176        let v = serde_json::to_string(&self.obj.spec).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1177        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1178    }
1179
1180    pub fn get_status(&mut self) -> RhaiRes<Dynamic> {
1181        let v =
1182            serde_json::to_string(&self.obj.status).map_err(|e| rhai_err(Error::SerializationError(e)))?;
1183        serde_json::from_str(&v).map_err(|e| rhai_err(Error::SerializationError(e)))
1184    }
1185
1186    pub fn wait_done(&mut self, timeout: i64) -> RhaiRes<()> {
1187        let name = self.obj.name_any();
1188        let cond = await_condition(self.api.clone(), &name, conditions::is_job_completed());
1189        tokio::task::block_in_place(|| {
1190            tokio::runtime::Handle::current().block_on(async move {
1191                tokio::time::timeout(std::time::Duration::from_secs(timeout as u64), cond)
1192                    .await
1193                    .map_err(Error::Elapsed)
1194            })
1195        })
1196        .map_err(rhai_err)?
1197        .map_err(|e| rhai_err(Error::KubeWaitError(e)))?;
1198        Ok(())
1199    }
1200}
1201
1202// ── Rhai registration ────────────────────────────────────────────────────────
1203
1204pub fn k8sgeneric_rhai_register(engine: &mut Engine) {
1205    use crate::{register_k8s_generic, register_k8s_object};
1206    engine
1207        .register_type_with_name::<DynamicObject>("DynamicObject")
1208        .register_get("data", |obj: &mut DynamicObject| -> Dynamic {
1209            Dynamic::from(obj.data.clone())
1210        });
1211    register_k8s_object!(engine, K8sObject);
1212    register_k8s_generic!(
1213        engine,
1214        K8sGeneric,
1215        K8sObject,
1216        K8sGeneric::new_global,
1217        K8sGeneric::new_ns,
1218        K8sGeneric::new_group_ns
1219    );
1220}
1221
1222pub fn k8sraw_rhai_register(engine: &mut Engine) {
1223    use crate::register_k8s_raw;
1224    register_k8s_raw!(engine, K8sRaw, K8sRaw::new);
1225}
1226
1227pub fn k8sworkload_rhai_register(engine: &mut Engine) {
1228    engine
1229        .register_type_with_name::<K8sDeploy>("K8sDeploy")
1230        .register_fn("get_deployment", K8sDeploy::get_deployment)
1231        .register_get("metadata", K8sDeploy::get_metadata)
1232        .register_get("spec", K8sDeploy::get_spec)
1233        .register_get("status", K8sDeploy::get_status)
1234        .register_fn("wait_available", K8sDeploy::wait_available);
1235    engine
1236        .register_type_with_name::<K8sDaemonSet>("K8sDaemonSet")
1237        .register_fn("get_deamonset", K8sDaemonSet::get_deamonset)
1238        .register_get("metadata", K8sDaemonSet::get_metadata)
1239        .register_get("spec", K8sDaemonSet::get_spec)
1240        .register_get("status", K8sDaemonSet::get_status)
1241        .register_fn("wait_available", K8sDaemonSet::wait_available);
1242    engine
1243        .register_type_with_name::<K8sStatefulSet>("K8sStatefulSet")
1244        .register_fn("get_statefulset", K8sStatefulSet::get_sts)
1245        .register_get("metadata", K8sStatefulSet::get_metadata)
1246        .register_get("spec", K8sStatefulSet::get_spec)
1247        .register_get("status", K8sStatefulSet::get_status)
1248        .register_fn("wait_available", K8sStatefulSet::wait_available);
1249    engine
1250        .register_type_with_name::<K8sJob>("K8sJob")
1251        .register_fn("get_job", K8sJob::get_job)
1252        .register_get("metadata", K8sJob::get_metadata)
1253        .register_get("spec", K8sJob::get_spec)
1254        .register_get("status", K8sJob::get_status)
1255        .register_fn("wait_done", K8sJob::wait_done);
1256}
1257
1258// ── Macros ───────────────────────────────────────────────────────────────────
1259
1260#[macro_export]
1261macro_rules! register_k8s_object {
1262    ($engine:expr, $type:ty) => {{
1263        let _delete: fn(&mut $type) -> $crate::RhaiRes<()> = <$type>::rhai_delete;
1264        let _wait_deleted: fn(&mut $type, i64) -> $crate::RhaiRes<()> = <$type>::rhai_wait_deleted;
1265        let _get_kind: fn(&mut $type) -> String = <$type>::get_kind;
1266        let _original_kind: fn(&mut $type) -> String = <$type>::original_kind;
1267        let _get_metadata: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::get_metadata;
1268        let _wait_condition: fn(&mut $type, String, i64) -> $crate::RhaiRes<()> = <$type>::wait_condition;
1269        let _wait_status: fn(&mut $type, String, i64) -> $crate::RhaiRes<()> = <$type>::wait_status;
1270        let _wait_status_prop: fn(&mut $type, String, i64) -> $crate::RhaiRes<()> = <$type>::wait_status_prop;
1271        let _wait_status_string: fn(&mut $type, String, String, i64) -> $crate::RhaiRes<()> =
1272            <$type>::wait_status_string;
1273        $engine
1274            .register_type_with_name::<$type>("K8sObject")
1275            .register_get("kind", _get_kind)
1276            .register_get("original_kind", _original_kind)
1277            .register_get("metadata", _get_metadata)
1278            .register_fn("delete", _delete)
1279            .register_fn("wait_condition", _wait_condition)
1280            .register_fn("wait_status", _wait_status)
1281            .register_fn("wait_status_prop", _wait_status_prop)
1282            .register_fn("wait_status_string", _wait_status_string)
1283            .register_fn("wait_deleted", _wait_deleted)
1284    }};
1285}
1286
1287#[macro_export]
1288macro_rules! register_k8s_generic {
1289    ($engine:expr, $type:ty, $obj_type:ty,
1290     $new_global:expr, $new_ns:expr, $new_group_ns:expr) => {{
1291        let _scope: fn(&mut $type) -> String = <$type>::rhai_get_scope;
1292        let _exist: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_exist;
1293        let _list: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_list;
1294        let _list_labels: fn(&mut $type, String) -> $crate::RhaiRes<rhai::Dynamic> =
1295            <$type>::rhai_list_labels;
1296        let _list_meta: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_list_meta;
1297        let _get: fn(&mut $type, String) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_get;
1298        let _get_meta: fn(&mut $type, String) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_get_meta;
1299        let _get_obj: fn(&mut $type, String) -> $crate::RhaiRes<$obj_type> = <$type>::rhai_get_obj;
1300        let _delete: fn(&mut $type, String) -> $crate::RhaiRes<()> = <$type>::rhai_delete;
1301        let _create: fn(&mut $type, rhai::Dynamic) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_create;
1302        let _replace: fn(&mut $type, String, rhai::Dynamic) -> $crate::RhaiRes<rhai::Dynamic> =
1303            <$type>::rhai_replace;
1304        let _patch: fn(&mut $type, String, rhai::Dynamic) -> $crate::RhaiRes<rhai::Dynamic> =
1305            <$type>::rhai_patch;
1306        let _apply: fn(&mut $type, String, rhai::Dynamic) -> $crate::RhaiRes<rhai::Dynamic> =
1307            <$type>::rhai_apply;
1308        $engine
1309            .register_type_with_name::<$type>("K8sGeneric")
1310            .register_fn("k8s_resource", $new_global)
1311            .register_fn("k8s_resource", $new_ns)
1312            .register_fn("k8s_resource", $new_group_ns)
1313            .register_fn("list", _list)
1314            .register_fn("list", _list_labels)
1315            .register_fn("update_k8s_crd_cache", $crate::k8s::update_cache)
1316            .register_fn("list_meta", _list_meta)
1317            .register_fn("get", _get)
1318            .register_fn("get_meta", _get_meta)
1319            .register_fn("get_obj", _get_obj)
1320            .register_fn("delete", _delete)
1321            .register_fn("create", _create)
1322            .register_fn("replace", _replace)
1323            .register_fn("patch", _patch)
1324            .register_fn("apply", _apply)
1325            .register_fn("exist", _exist)
1326            .register_get("scope", _scope)
1327    }};
1328}
1329
1330#[macro_export]
1331macro_rules! register_k8s_raw {
1332    ($engine:expr, $type:ty, $new:expr) => {{
1333        let _get_url: fn(&mut $type, String) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_get_url;
1334        let _get_version: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> = <$type>::rhai_get_api_version;
1335        let _get_api_resources: fn(&mut $type) -> $crate::RhaiRes<rhai::Dynamic> =
1336            <$type>::rhai_get_api_resources;
1337        $engine
1338            .register_type_with_name::<$type>("K8sRaw")
1339            .register_fn("new_k8s_raw", $new)
1340            .register_fn("get_url", _get_url)
1341            .register_fn("get_cluster_version", _get_version)
1342            .register_fn("get_api_resources", _get_api_resources)
1343    }};
1344}
1345
1346// ── Tests ────────────────────────────────────────────────────────────────────
1347
1348#[cfg(test)]
1349mod tests {
1350    use super::*;
1351
1352    #[test]
1353    fn discovery_excludes_aggregated_keeps_local_apiservices() {
1354        let aggregated = serde_json::json!({
1355            "group": "metrics.k8s.io",
1356            "service": { "name": "metrics-server", "namespace": "kube-system" }
1357        });
1358        assert_eq!(
1359            aggregated_apiservice_group(&aggregated),
1360            Some("metrics.k8s.io".to_string())
1361        );
1362
1363        let local_null = serde_json::json!({ "group": "storage.k8s.io", "service": null });
1364        assert_eq!(aggregated_apiservice_group(&local_null), None);
1365        let local_absent = serde_json::json!({ "group": "apps" });
1366        assert_eq!(aggregated_apiservice_group(&local_absent), None);
1367
1368        let empty_group = serde_json::json!({ "group": "", "service": { "name": "x" } });
1369        assert_eq!(aggregated_apiservice_group(&empty_group), None);
1370    }
1371
1372    #[test]
1373    fn test_job_is_completed_succeeded() {
1374        let data = serde_json::json!({"status": {"succeeded": 1}});
1375        assert!(job_is_completed(&data));
1376    }
1377
1378    #[test]
1379    fn test_job_is_completed_zero_succeeded() {
1380        let data = serde_json::json!({"status": {"succeeded": 0}});
1381        assert!(!job_is_completed(&data));
1382    }
1383
1384    #[test]
1385    fn test_job_is_completed_completion_time() {
1386        let data = serde_json::json!({"status": {"completionTime": "2024-01-01T00:00:00Z"}});
1387        assert!(job_is_completed(&data));
1388    }
1389
1390    #[test]
1391    fn test_job_is_completed_no_status() {
1392        let data = serde_json::json!({});
1393        assert!(!job_is_completed(&data));
1394    }
1395
1396    #[test]
1397    fn test_job_is_completed_still_running() {
1398        let data = serde_json::json!({"status": {"active": 1, "succeeded": 0}});
1399        assert!(!job_is_completed(&data));
1400    }
1401
1402    // ── prepare_handle ────────────────────────────────────────────────────────
1403
1404    fn map(v: serde_json::Value) -> serde_json::Map<String, serde_json::Value> {
1405        match v {
1406            serde_json::Value::Object(m) => m,
1407            _ => panic!("expected object"),
1408        }
1409    }
1410
1411    #[test]
1412    fn prepare_handle_inserts_metadata_when_absent() {
1413        let input = map(serde_json::json!({"kind": "ConfigMap", "spec": {}}));
1414        let out = prepare_handle(input, None, None, None, None, false);
1415        assert!(out.contains_key("metadata"), "metadata must be added");
1416        assert!(out["metadata"].is_object());
1417    }
1418
1419    #[test]
1420    fn prepare_handle_inserts_metadata_when_not_object() {
1421        let input = map(serde_json::json!({"metadata": "bad-string", "kind": "ConfigMap"}));
1422        let out = prepare_handle(input, None, None, None, None, false);
1423        assert!(
1424            out["metadata"].is_object(),
1425            "non-object metadata must be replaced with {{}}"
1426        );
1427    }
1428
1429    #[test]
1430    fn prepare_handle_preserves_existing_metadata() {
1431        let input = map(serde_json::json!({"metadata": {"name": "foo"}, "kind": "ConfigMap"}));
1432        let out = prepare_handle(input, None, None, None, None, false);
1433        assert_eq!(out["metadata"]["name"], "foo");
1434    }
1435
1436    #[test]
1437    fn prepare_handle_injects_labels_when_missing() {
1438        let input = map(serde_json::json!({"kind": "ConfigMap"}));
1439        let labels = serde_json::json!({"app": "myapp", "tier": "backend"});
1440        let out = prepare_handle(input, Some(labels), None, None, None, false);
1441        assert_eq!(out["metadata"]["labels"]["app"], "myapp");
1442        assert_eq!(out["metadata"]["labels"]["tier"], "backend");
1443    }
1444
1445    #[test]
1446    fn prepare_handle_labels_do_not_override_existing() {
1447        let input = map(serde_json::json!({"metadata": {"labels": {"app": "existing"}}}));
1448        let labels = serde_json::json!({"app": "override-attempt", "extra": "v"});
1449        let out = prepare_handle(input, Some(labels), None, None, None, false);
1450        assert_eq!(
1451            out["metadata"]["labels"]["app"], "existing",
1452            "existing label must not be overridden"
1453        );
1454        assert_eq!(
1455            out["metadata"]["labels"]["extra"], "v",
1456            "new label must be injected"
1457        );
1458    }
1459
1460    #[test]
1461    fn prepare_handle_owner_ref_injected_when_same_ns() {
1462        let input = map(serde_json::json!({"kind": "ConfigMap"}));
1463        let owner = serde_json::json!({"apiVersion": "v1", "kind": "Pod", "name": "owner", "uid": "abc"});
1464        let out = prepare_handle(
1465            input,
1466            None,
1467            Some(owner.clone()),
1468            Some("mynamespace".to_string()),
1469            Some("mynamespace".to_string()),
1470            true,
1471        );
1472        let refs = out["metadata"]["ownerReferences"].as_array().unwrap();
1473        assert_eq!(refs.len(), 1);
1474        assert_eq!(refs[0]["uid"], "abc");
1475    }
1476
1477    #[test]
1478    fn prepare_handle_owner_ref_not_injected_when_different_ns() {
1479        let input = map(serde_json::json!({"kind": "ConfigMap"}));
1480        let owner = serde_json::json!({"apiVersion": "v1", "kind": "Pod", "name": "owner", "uid": "abc"});
1481        let out = prepare_handle(
1482            input,
1483            None,
1484            Some(owner),
1485            Some("other-ns".to_string()),
1486            Some("mynamespace".to_string()),
1487            true,
1488        );
1489        assert!(
1490            out["metadata"]
1491                .as_object()
1492                .unwrap()
1493                .get("ownerReferences")
1494                .is_none()
1495        );
1496    }
1497
1498    // ── Condition closures ────────────────────────────────────────────────────
1499
1500    fn dynobj(extra: serde_json::Value) -> DynamicObject {
1501        let mut base = serde_json::json!({
1502            "apiVersion": "v1",
1503            "kind": "Foo",
1504            "metadata": {"name": "test"}
1505        });
1506        if let (Some(obj), serde_json::Value::Object(fields)) = (base.as_object_mut(), extra) {
1507            obj.extend(fields);
1508        }
1509        serde_json::from_value(base).unwrap()
1510    }
1511
1512    #[test]
1513    fn is_condition_matches_ready_true() {
1514        let obj = dynobj(serde_json::json!({
1515            "status": {"conditions": [{"type": "Ready", "status": "True"}]}
1516        }));
1517        assert!(K8sObject::is_condition("Ready".to_string()).matches_object(Some(&obj)));
1518    }
1519
1520    #[test]
1521    fn is_condition_no_match_wrong_type() {
1522        let obj = dynobj(serde_json::json!({
1523            "status": {"conditions": [{"type": "Available", "status": "True"}]}
1524        }));
1525        assert!(!K8sObject::is_condition("Ready".to_string()).matches_object(Some(&obj)));
1526    }
1527
1528    #[test]
1529    fn is_condition_no_match_status_false() {
1530        let obj = dynobj(serde_json::json!({
1531            "status": {"conditions": [{"type": "Ready", "status": "False"}]}
1532        }));
1533        assert!(!K8sObject::is_condition("Ready".to_string()).matches_object(Some(&obj)));
1534    }
1535
1536    #[test]
1537    fn is_condition_missing_conditions_returns_false() {
1538        let obj = dynobj(serde_json::json!({"status": {}}));
1539        assert!(!K8sObject::is_condition("Ready".to_string()).matches_object(Some(&obj)));
1540    }
1541
1542    #[test]
1543    fn is_condition_missing_status_returns_false() {
1544        let obj = dynobj(serde_json::json!({}));
1545        assert!(!K8sObject::is_condition("Ready".to_string()).matches_object(Some(&obj)));
1546    }
1547
1548    #[test]
1549    fn is_condition_on_none_returns_false() {
1550        assert!(!K8sObject::is_condition("Ready".to_string()).matches_object(None));
1551    }
1552
1553    #[test]
1554    fn is_status_true_boolean() {
1555        let obj = dynobj(serde_json::json!({"status": {"healthy": true}}));
1556        assert!(K8sObject::is_status("healthy".to_string()).matches_object(Some(&obj)));
1557    }
1558
1559    #[test]
1560    fn is_status_false_boolean() {
1561        let obj = dynobj(serde_json::json!({"status": {"healthy": false}}));
1562        assert!(!K8sObject::is_status("healthy".to_string()).matches_object(Some(&obj)));
1563    }
1564
1565    #[test]
1566    fn is_status_missing_prop() {
1567        let obj = dynobj(serde_json::json!({"status": {}}));
1568        assert!(!K8sObject::is_status("healthy".to_string()).matches_object(Some(&obj)));
1569    }
1570
1571    #[test]
1572    fn have_status_non_null() {
1573        let obj = dynobj(serde_json::json!({"status": {"phase": "Running"}}));
1574        assert!(K8sObject::have_status("phase".to_string()).matches_object(Some(&obj)));
1575    }
1576
1577    #[test]
1578    fn have_status_null_value() {
1579        let obj = dynobj(serde_json::json!({"status": {"phase": null}}));
1580        assert!(!K8sObject::have_status("phase".to_string()).matches_object(Some(&obj)));
1581    }
1582
1583    #[test]
1584    fn have_status_missing_prop() {
1585        let obj = dynobj(serde_json::json!({"status": {}}));
1586        assert!(!K8sObject::have_status("phase".to_string()).matches_object(Some(&obj)));
1587    }
1588
1589    #[test]
1590    fn have_status_value_matches() {
1591        let obj = dynobj(serde_json::json!({"status": {"phase": "Running"}}));
1592        assert!(
1593            K8sObject::have_status_value("phase".to_string(), "Running".to_string())
1594                .matches_object(Some(&obj))
1595        );
1596    }
1597
1598    #[test]
1599    fn have_status_value_no_match() {
1600        let obj = dynobj(serde_json::json!({"status": {"phase": "Pending"}}));
1601        assert!(
1602            !K8sObject::have_status_value("phase".to_string(), "Running".to_string())
1603                .matches_object(Some(&obj))
1604        );
1605    }
1606
1607    #[test]
1608    fn have_status_value_missing_prop() {
1609        let obj = dynobj(serde_json::json!({"status": {}}));
1610        assert!(
1611            !K8sObject::have_status_value("phase".to_string(), "Running".to_string())
1612                .matches_object(Some(&obj))
1613        );
1614    }
1615}