1use std::sync::Arc;
11use std::time::Duration;
12
13use futures::StreamExt;
14use kube::{
15 api::{Api, Patch, PatchParams},
16 runtime::{
17 controller::Action, predicates, reflector, watcher, watcher::Config, Controller,
18 WatchStreamExt,
19 },
20 Client, CustomResource, ResourceExt,
21};
22
23use schemars::JsonSchema;
24use serde::{Deserialize, Serialize};
25use serde_json::json;
26use thiserror::Error;
27
28use crate::{Condition, LavaArchitectureSpec, Phase, Source};
29
30#[derive(CustomResource, Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
34#[kube(
35 group = "lava.pleme.io",
36 version = "v1alpha1",
37 kind = "LavaArchitecture",
38 namespaced,
39 status = "LavaArchitectureStatusCR"
40)]
41#[serde(rename_all = "camelCase")]
42pub struct LavaArchitectureSpecCR {
43 pub source: SourceCR,
44 #[serde(default)]
45 pub bindings: std::collections::BTreeMap<String, String>,
46 #[serde(default)]
47 pub gate: Option<String>,
48 #[serde(default = "default_engine")]
49 pub engine: String,
50}
51
52fn default_engine() -> String {
53 "embedded".to_string()
54}
55
56#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq, Default)]
73#[serde(rename_all = "camelCase")]
74pub struct SourceCR {
75 #[serde(default)]
76 pub inline: Option<String>,
77 #[serde(default)]
78 pub name: Option<String>,
79 #[serde(default)]
80 pub url: Option<String>,
81 #[serde(default)]
82 pub rev: Option<String>,
83 #[serde(default)]
84 pub path: Option<String>,
85}
86
87impl SourceCR {
88 #[must_use]
91 pub fn variant(&self) -> Option<crate::Source> {
92 if let Some(inline) = &self.inline {
93 Some(crate::Source::Inline { inline: inline.clone() })
94 } else if let Some(name) = &self.name {
95 Some(crate::Source::Name { name: name.clone() })
96 } else if let (Some(url), Some(rev), Some(path)) = (&self.url, &self.rev, &self.path) {
97 Some(crate::Source::Git {
98 url: url.clone(),
99 rev: rev.clone(),
100 path: path.clone(),
101 })
102 } else {
103 None
104 }
105 }
106}
107
108#[derive(Debug, Default, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
109#[serde(rename_all = "camelCase")]
110pub struct LavaArchitectureStatusCR {
111 pub phase: Option<String>,
112 #[serde(default)]
113 pub conditions: Vec<ConditionCR>,
114 pub last_synthesized_hash: Option<String>,
115 pub last_applied_at: Option<String>,
116}
117
118#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
119#[serde(rename_all = "camelCase")]
120pub struct ConditionCR {
121 #[serde(rename = "type")]
122 pub kind: String,
123 pub status: String,
124 pub reason: Option<String>,
125 pub message: Option<String>,
126}
127
128impl From<&Condition> for ConditionCR {
129 fn from(c: &Condition) -> Self {
130 Self {
131 kind: c.kind.clone(),
132 status: c.status.clone(),
133 reason: c.reason.clone(),
134 message: c.message.clone(),
135 }
136 }
137}
138
139#[derive(CustomResource, Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
145#[kube(
146 group = "lava.pleme.io",
147 version = "v1alpha1",
148 kind = "RemediationPolicy",
149 namespaced
150)]
151#[serde(rename_all = "camelCase")]
152pub struct RemediationPolicySpec {
153 pub cosmetic: String,
154 pub functional: String,
155 pub critical: String,
156 #[serde(default)]
157 pub escalation: Option<EscalationLadderCR>,
158}
159
160#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
161#[serde(rename_all = "camelCase")]
162pub struct EscalationLadderCR {
163 pub tiers: Vec<EscalationTierCR>,
164}
165
166#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
167#[serde(rename_all = "camelCase")]
168pub struct EscalationTierCR {
169 pub target: NotifyTargetCR,
170 #[serde(default = "default_tier_wait")]
174 pub wait_before_next: String,
175}
176
177fn default_tier_wait() -> String {
178 "PT15M".to_string()
179}
180
181#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
188#[serde(rename_all = "camelCase")]
189pub struct NotifyTargetCR {
190 pub kind: String,
192 #[serde(default)]
193 pub webhook_secret_ref: Option<String>,
194 #[serde(default)]
195 pub topic: Option<String>,
196 #[serde(default)]
197 pub service_key_secret_ref: Option<String>,
198 #[serde(default)]
199 pub address: Option<String>,
200 #[serde(default)]
201 pub url: Option<String>,
202 #[serde(default)]
203 pub secret_ref: Option<String>,
204}
205
206impl RemediationPolicySpec {
207 #[must_use]
211 pub fn to_policy(&self) -> lava_anomaly::RemediationPolicy {
212 use lava_anomaly::RemediationAction;
213 let parse = |s: &str| -> RemediationAction {
214 match s {
215 "NoOp" => RemediationAction::NoOp,
216 "Alert" => RemediationAction::Alert,
217 "AutoCorrect" => RemediationAction::AutoCorrect,
218 "RequireApproval" => RemediationAction::RequireApproval,
219 "Escalate" => RemediationAction::Escalate,
220 _ => RemediationAction::Alert,
221 }
222 };
223 lava_anomaly::RemediationPolicy {
224 cosmetic: parse(&self.cosmetic),
225 functional: parse(&self.functional),
226 critical: parse(&self.critical),
227 escalation: None, }
230 }
231}
232
233#[derive(CustomResource, Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
237#[kube(
238 group = "lava.pleme.io",
239 version = "v1alpha1",
240 kind = "LavaArchitectureDependency",
241 namespaced
242)]
243#[serde(rename_all = "camelCase")]
244pub struct LavaArchitectureDependencySpec {
245 pub from: ResourceRefCR,
246 pub to: ResourceRefCR,
247 pub kind: String, #[serde(default = "default_require_phase")]
249 pub require_phase: String,
250}
251
252fn default_require_phase() -> String {
253 "Applied".to_string()
254}
255
256#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
257#[serde(rename_all = "camelCase")]
258pub struct ResourceRefCR {
259 pub cluster: String,
260 pub namespace: String,
261 pub name: String,
262}
263
264impl LavaArchitectureDependencySpec {
265 #[must_use]
268 pub fn to_lib(&self) -> lava_dependency::LavaArchitectureDependency {
269 use lava_dependency::{DependencyKind, LavaArchitectureDependency};
270 use lava_outcome_chain::ResourceAddress;
271 let kind = match self.kind.as_str() {
272 "Influences" => DependencyKind::Influences,
273 _ => DependencyKind::BlocksOn,
274 };
275 LavaArchitectureDependency {
276 from: ResourceAddress::new(&self.from.cluster, &self.from.namespace, &self.from.name),
277 to: ResourceAddress::new(&self.to.cluster, &self.to.namespace, &self.to.name),
278 kind,
279 require_phase: self.require_phase.clone(),
280 }
281 }
282}
283
284#[derive(Debug, Error)]
285pub enum ControllerError {
286 #[error("kube: {0}")]
287 Kube(#[from] kube::Error),
288 #[error("finalizer: {0}")]
289 Finalizer(String),
290}
291
292pub type SynthesizeFn = Arc<
296 dyn Fn(
297 &Source,
298 &indexmap::IndexMap<String, String>,
299 Option<&str>,
300 ) -> Result<serde_json::Value, String>
301 + Send
302 + Sync,
303>;
304
305#[derive(Clone)]
317pub struct Context {
318 pub client: Client,
319 pub synthesize: SynthesizeFn,
320 pub chain: std::sync::Arc<
321 std::sync::Mutex<
322 lava_outcome_chain::OutcomeChain<
323 lava_outcome_chain::OutcomePayload,
324 lava_outcome_chain::InMemorySink<lava_outcome_chain::OutcomePayload>,
325 lava_outcome_chain::NoSigning,
326 >,
327 >,
328 >,
329}
330
331impl Context {
332 #[must_use]
336 pub fn with_in_memory_chain(client: Client, synthesize: SynthesizeFn) -> Self {
337 Self {
338 client,
339 synthesize,
340 chain: crate::viggy_loop::shared_in_memory_chain(),
341 }
342 }
343}
344
345pub async fn reconcile_one(
354 obj: Arc<LavaArchitecture>,
355 ctx: Arc<Context>,
356) -> Result<Action, ControllerError> {
357 use crate::viggy_loop::{engine_with_default_router, LavaPromessaController};
358 use lava_anomaly::RemediationPolicy;
359 use lava_drift::{DriftDetector, PlannerBackend, PlannerError};
360 use lava_outcome_chain::ResourceAddress;
361
362 let ns = obj.namespace().unwrap_or_else(|| "default".to_string());
363 let name = obj.name_any();
364 let api: Api<LavaArchitecture> = Api::namespaced(ctx.client.clone(), &ns);
365
366 let spec = to_lib_spec(&obj.spec);
367 let address = ResourceAddress::new(
368 std::env::var("LAVA_OPERATOR_CLUSTER").unwrap_or_else(|_| "local".into()),
369 ns.clone(),
370 name.clone(),
371 );
372
373 let source_text = match &spec.source {
376 crate::Source::Inline { inline } => inline.clone(),
377 crate::Source::Name { name } => name.clone(),
378 crate::Source::Git { path, .. } => path.clone(),
379 };
380
381 #[cfg(feature = "magma-bridge")]
386 let detector = DriftDetector::new(crate::magma_bridge::EmbeddedMagmaPlanner);
387
388 #[cfg(not(feature = "magma-bridge"))]
389 let detector = {
390 struct CallbackPlanner {
391 synthesize: SynthesizeFn,
392 source: crate::Source,
393 }
394 impl PlannerBackend for CallbackPlanner {
395 fn plan(
396 &self,
397 _src: &str,
398 bindings: &indexmap::IndexMap<String, String>,
399 ) -> Result<Vec<lava_drift::DriftFinding>, PlannerError> {
400 (self.synthesize)(&self.source, bindings, None)
401 .map(|_| Vec::new())
402 .map_err(PlannerError::Plan)
403 }
404 }
405 DriftDetector::new(CallbackPlanner {
406 synthesize: ctx.synthesize.clone(),
407 source: spec.source.clone(),
408 })
409 };
410
411 let controller_ = LavaPromessaController {
412 source_text,
413 source_address: address.clone(),
414 detector,
415 chain: ctx.chain.clone(),
416 };
417
418 let engine = engine_with_default_router(controller_, RemediationPolicy::default());
419 let bindings = spec.bindings.clone();
422 let report = engine.tick(address, bindings);
423
424 let conditions: Vec<ConditionCR> = report
425 .beats
426 .iter()
427 .map(|b| ConditionCR {
428 kind: format!("Beat.{}", b.beat.as_str()),
429 status: match b.status {
430 lava_viggy::BeatStatus::Ok => "True".to_string(),
431 lava_viggy::BeatStatus::Skipped => "Unknown".to_string(),
432 lava_viggy::BeatStatus::Failed => "False".to_string(),
433 },
434 reason: Some(format!("{:?}", b.status)),
435 message: b.message.clone(),
436 })
437 .collect();
438
439 let status = LavaArchitectureStatusCR {
440 phase: Some(report.final_phase.as_str().to_string()),
441 conditions,
442 last_synthesized_hash: None,
443 last_applied_at: Some(report.ended_at.to_rfc3339()),
444 };
445 let patch = json!({
449 "apiVersion": "lava.pleme.io/v1alpha1",
450 "kind": "LavaArchitecture",
451 "status": status,
452 });
453 api.patch_status(&name, &PatchParams::apply("lava-operator").force(), &Patch::Apply(patch))
454 .await?;
455
456 let requeue = report
457 .decision
458 .as_ref()
459 .map(|d| Duration::from_secs(d.requeue_after.num_seconds().max(1) as u64))
460 .unwrap_or_else(|| Duration::from_secs(30));
461 Ok(Action::requeue(requeue))
462}
463
464pub fn error_policy(
465 _obj: Arc<LavaArchitecture>,
466 _err: &ControllerError,
467 _ctx: Arc<Context>,
468) -> Action {
469 Action::requeue(Duration::from_secs(60))
470}
471
472pub async fn run(synthesize: SynthesizeFn) -> Result<(), kube::Error> {
478 let client = Client::try_default().await?;
479 let api: Api<LavaArchitecture> = Api::all(client.clone());
480 let ctx = Arc::new(Context::with_in_memory_chain(client, synthesize));
481 let (reader, writer) = reflector::store();
491 let stream = watcher(api, Config::default())
492 .default_backoff()
493 .reflect(writer)
494 .applied_objects()
495 .predicate_filter(predicates::generation);
496 Controller::for_stream(stream, reader)
497 .run(reconcile_one, error_policy, ctx)
498 .for_each(|res| async move {
499 match res {
500 Ok((obj, _)) => tracing::info!(?obj, "reconciled"),
501 Err(e) => tracing::warn!(error = %e, "reconcile error"),
502 }
503 })
504 .await;
505 Ok(())
506}
507
508fn to_lib_spec(spec: &LavaArchitectureSpecCR) -> LavaArchitectureSpec {
509 let mut bindings = indexmap::IndexMap::new();
510 for (k, v) in &spec.bindings {
511 bindings.insert(k.clone(), v.clone());
512 }
513 LavaArchitectureSpec {
514 source: spec.source.variant().unwrap_or_else(|| Source::Inline {
515 inline: "(deflava-architecture empty :inputs () :resources ())".to_string(),
516 }),
517 bindings,
518 gate: spec.gate.clone(),
519 engine: spec.engine.clone(),
520 }
521}
522
523#[must_use]
526pub const fn phase_str(p: Phase) -> &'static str {
527 p.as_str()
528}
529
530
531#[cfg(test)]
532mod tests {
533 use super::*;
534
535 #[test]
536 fn to_lib_spec_round_trips_inline_source() {
537 let cr = LavaArchitectureSpecCR {
538 source: SourceCR {
539 inline: Some("(x)".into()),
540 ..Default::default()
541 },
542 bindings: std::collections::BTreeMap::from_iter([(
543 "name".to_string(),
544 "prod".to_string(),
545 )]),
546 gate: Some("vpc".into()),
547 engine: "embedded".into(),
548 };
549 let lib = to_lib_spec(&cr);
550 assert!(matches!(lib.source, Source::Inline { .. }));
551 assert_eq!(lib.bindings["name"], "prod");
552 assert_eq!(lib.gate.as_deref(), Some("vpc"));
553 }
554
555 #[test]
556 fn condition_cr_converts_from_lib_condition() {
557 let c = Condition::ok("Synthesized", "RenderOk");
558 let cr: ConditionCR = (&c).into();
559 assert_eq!(cr.kind, "Synthesized");
560 assert_eq!(cr.status, "True");
561 assert_eq!(cr.reason.as_deref(), Some("RenderOk"));
562 }
563
564 #[test]
565 fn remediation_policy_spec_to_policy_parses_action_strings() {
566 let spec = RemediationPolicySpec {
567 cosmetic: "NoOp".into(),
568 functional: "AutoCorrect".into(),
569 critical: "Escalate".into(),
570 escalation: None,
571 };
572 let p = spec.to_policy();
573 assert_eq!(p.cosmetic, lava_anomaly::RemediationAction::NoOp);
574 assert_eq!(p.functional, lava_anomaly::RemediationAction::AutoCorrect);
575 assert_eq!(p.critical, lava_anomaly::RemediationAction::Escalate);
576 }
577
578 #[test]
579 fn remediation_policy_spec_degrades_unknown_action_to_alert() {
580 let spec = RemediationPolicySpec {
581 cosmetic: "Moonwalk".into(),
582 functional: "Yodel".into(),
583 critical: "Apply".into(),
584 escalation: None,
585 };
586 let p = spec.to_policy();
587 assert_eq!(p.cosmetic, lava_anomaly::RemediationAction::Alert);
588 assert_eq!(p.functional, lava_anomaly::RemediationAction::Alert);
589 assert_eq!(p.critical, lava_anomaly::RemediationAction::Alert);
590 }
591
592 #[test]
593 fn dependency_spec_to_lib_maps_kind_and_addresses() {
594 let spec = LavaArchitectureDependencySpec {
595 from: ResourceRefCR {
596 cluster: "rio".into(),
597 namespace: "lava-system".into(),
598 name: "app".into(),
599 },
600 to: ResourceRefCR {
601 cluster: "rio".into(),
602 namespace: "lava-system".into(),
603 name: "vpc".into(),
604 },
605 kind: "Influences".into(),
606 require_phase: "Applied".into(),
607 };
608 let d = spec.to_lib();
609 assert_eq!(d.kind, lava_dependency::DependencyKind::Influences);
610 assert_eq!(d.from.namespace, "lava-system");
611 assert_eq!(d.to.name, "vpc");
612 }
613
614 #[test]
615 fn dependency_spec_defaults_kind_to_blocks_on_for_unknown_kind() {
616 let spec = LavaArchitectureDependencySpec {
617 from: ResourceRefCR {
618 cluster: "rio".into(),
619 namespace: "lava-system".into(),
620 name: "app".into(),
621 },
622 to: ResourceRefCR {
623 cluster: "rio".into(),
624 namespace: "lava-system".into(),
625 name: "vpc".into(),
626 },
627 kind: "unrecognized".into(),
628 require_phase: "Applied".into(),
629 };
630 let d = spec.to_lib();
631 assert_eq!(d.kind, lava_dependency::DependencyKind::BlocksOn);
632 }
633
634 #[test]
635 fn phase_str_covers_every_variant() {
636 for p in [
637 Phase::Pending,
638 Phase::Synthesized,
639 Phase::Planned,
640 Phase::Applied,
641 Phase::Drifted,
642 Phase::Reconverging,
643 Phase::Finalizing,
644 Phase::Failed,
645 ] {
646 assert!(!phase_str(p).is_empty(), "{p:?}");
647 }
648 assert_eq!(phase_str(Phase::Pending), "Pending");
649 assert_eq!(phase_str(Phase::Synthesized), "Synthesized");
650 assert_eq!(phase_str(Phase::Planned), "Planned");
651 assert_eq!(phase_str(Phase::Applied), "Applied");
652 assert_eq!(phase_str(Phase::Failed), "Failed");
653 }
654}