1use 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
30pub 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
50pub 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
70lazy_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
636fn 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(¤t.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
854lazy_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#[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
1202pub 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#[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#[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 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 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}