diff --git a/CLAUDE.md b/CLAUDE.md index bacfd24..c3772c4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -113,7 +113,7 @@ cargo fmt && cargo clippy -- -D warnings && cargo test | `axum` 0.8 | HTTP server | | `cdevents-sdk` | CDEvents emission | | `reqwest` | HTTP client | -| `ulid` | ULID generation for FALSE Protocol | +| `false-protocol` | FALSE Protocol occurrence types | | `chrono` | Timestamp handling | | `thiserror` | Error types | | `async-trait` | Async trait support | diff --git a/Cargo.lock b/Cargo.lock index 8884c92..0a90671 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -631,6 +631,17 @@ dependencies = [ "pin-project-lite", ] +[[package]] +name = "false-protocol" +version = "0.1.0" +dependencies = [ + "chrono", + "regex", + "serde", + "serde_json", + "ulid", +] + [[package]] name = "fastrand" version = "2.3.0" @@ -1386,6 +1397,7 @@ dependencies = [ "cdevents-sdk", "chrono", "cloudevents-sdk", + "false-protocol", "futures", "gateway-api", "k8s-openapi", @@ -1408,7 +1420,6 @@ dependencies = [ "toml", "tracing", "tracing-subscriber", - "ulid", "uuid", "x509-parser", ] diff --git a/Cargo.toml b/Cargo.toml index 7a39f9d..bd3f9ef 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -59,8 +59,8 @@ async-trait = "0.1" # Prometheus metrics prometheus = "0.13" -# FALSE Protocol (ULID generation for occurrence IDs) -ulid = "1" +# FALSE Protocol occurrence types +false-protocol = { path = "../false-protocol/rust" } [dev-dependencies] serde_yaml = "0.9" diff --git a/README.md b/README.md index 416b83f..f90a560 100644 --- a/README.md +++ b/README.md @@ -33,8 +33,9 @@ Part of the [False Systems](https://github.com/false-systems) toolchain. ## Quick Start ```bash -# Clone and build +# Clone KULTA and its sibling dependency git clone https://github.com/false-systems/kulta +git clone https://github.com/false-systems/false-protocol cd kulta cargo build --release diff --git a/deploy/crd.yaml b/deploy/crd.yaml index a570671..eb1af1d 100644 --- a/deploy/crd.yaml +++ b/deploy/crd.yaml @@ -53,6 +53,33 @@ spec: Compatible with Argo Rollouts API for easy migration' properties: + advisor: + default: + level: 'Off' + timeoutSeconds: 0 + description: AI advisor configuration for progressive AI adoption + properties: + endpoint: + description: URL of AI advisory service (e.g., MCP endpoint, HTTP + API) + nullable: true + type: string + level: + default: 'Off' + description: Advisor integration level + enum: + - 'Off' + - Context + - Advised + - Planned + - Driven + type: string + timeoutSeconds: + description: Timeout for advisory calls in seconds + format: uint64 + minimum: 0.0 + type: integer + type: object maxSurge: description: 'Maximum number of pods that can be scheduled above the desired number during update. @@ -9865,6 +9892,11 @@ spec: - timestamp type: object type: array + lastDecisionSource: + description: Source of last analysis decision (Threshold, Advisor, + Human) + nullable: true + type: string message: description: Human-readable message nullable: true @@ -19780,6 +19812,11 @@ spec: - timestamp type: object type: array + lastDecisionSource: + description: Source of last analysis decision (Threshold, Advisor, + Human) + nullable: true + type: string message: description: Human-readable message nullable: true diff --git a/src/controller/advisor.rs b/src/controller/advisor.rs new file mode 100644 index 0000000..d184780 --- /dev/null +++ b/src/controller/advisor.rs @@ -0,0 +1,279 @@ +//! AI advisor integration for progressive delivery decisions +//! +//! Follows the same trait-based pattern as `MetricsQuerier` (prometheus.rs): +//! - `AnalysisAdvisor` trait for abstraction +//! - `NoOpAdvisor` for Level 0/1 (no advisory calls) +//! - `HttpAdvisor` for Level 2+ (calls external AI service) +//! - `MockAdvisor` for testing +//! +//! The advisor never overrides threshold decisions at Level 2 — it only +//! provides recommendations that are logged alongside the threshold result. + +use crate::crd::rollout::{Recommendation, RecommendedAction}; +use async_trait::async_trait; +use serde::{Deserialize, Serialize}; +use std::time::Duration; +use thiserror::Error; + +#[derive(Debug, Error)] +pub enum AdvisorError { + #[error("Advisory service unreachable: {0}")] + Unreachable(String), + + #[error("Advisory service returned invalid response: {0}")] + InvalidResponse(String), + + #[error("Advisory call timed out after {0:?}")] + Timeout(Duration), +} + +/// Everything the advisor needs to make a recommendation +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct AnalysisContext { + pub rollout_name: String, + pub namespace: String, + pub strategy: String, + pub current_step: Option, + pub current_weight: Option, + pub metrics_healthy: bool, + pub phase: String, + pub history: Vec, +} + +/// Trait for AI advisory integration +/// +/// Production code uses `HttpAdvisor` which calls an external AI service. +/// Tests use `MockAdvisor` which returns preconfigured responses. +/// Default is `NoOpAdvisor` which returns Continue with zero confidence. +#[async_trait] +pub trait AnalysisAdvisor: Send + Sync { + /// Request a recommendation from the advisor + async fn advise(&self, context: &AnalysisContext) -> Result; + + /// Downcast support for testing + fn as_any(&self) -> &dyn std::any::Any; +} + +/// No-op advisor for Level 0/1 (default) +/// +/// Returns Continue with zero confidence — the threshold decision is used as-is. +pub struct NoOpAdvisor; + +#[async_trait] +impl AnalysisAdvisor for NoOpAdvisor { + async fn advise(&self, _ctx: &AnalysisContext) -> Result { + Ok(Recommendation { + action: RecommendedAction::Continue, + confidence: 0.0, + reasoning: "no advisor configured".into(), + }) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +/// HTTP-based advisor for Level 2+ (production) +/// +/// Calls an external AI advisory service with rollout context +/// and returns the recommendation. Times out gracefully. +pub struct HttpAdvisor { + client: reqwest::Client, + endpoint: String, + timeout: Duration, +} + +impl HttpAdvisor { + pub fn new(endpoint: String, timeout: Duration) -> Self { + let client = reqwest::Client::builder() + .timeout(timeout) + .build() + .unwrap_or_default(); + Self { + client, + endpoint, + timeout, + } + } +} + +#[async_trait] +impl AnalysisAdvisor for HttpAdvisor { + async fn advise(&self, context: &AnalysisContext) -> Result { + let response = self + .client + .post(&self.endpoint) + .json(context) + .send() + .await + .map_err(|e| { + if e.is_timeout() { + AdvisorError::Timeout(self.timeout) + } else { + AdvisorError::Unreachable(e.to_string()) + } + })?; + + let recommendation: Recommendation = response + .json() + .await + .map_err(|e| AdvisorError::InvalidResponse(e.to_string()))?; + + Ok(recommendation) + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +/// Mock advisor for testing +/// +/// Returns a preconfigured recommendation. Thread-safe via Arc>. +#[cfg(test)] +pub struct MockAdvisor { + pub response: std::sync::Arc>>, + pub call_count: std::sync::Arc, +} + +#[cfg(test)] +impl MockAdvisor { + pub fn new(recommendation: Recommendation) -> Self { + Self { + response: std::sync::Arc::new(std::sync::Mutex::new(Ok(recommendation))), + call_count: std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)), + } + } + + pub fn new_failing(error_msg: &str) -> Self { + Self { + response: std::sync::Arc::new(std::sync::Mutex::new(Err(error_msg.to_string()))), + call_count: std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)), + } + } + + pub fn calls(&self) -> u32 { + self.call_count.load(std::sync::atomic::Ordering::Relaxed) + } +} + +#[cfg(test)] +#[async_trait] +impl AnalysisAdvisor for MockAdvisor { + async fn advise(&self, _ctx: &AnalysisContext) -> Result { + self.call_count + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + + let guard = self + .response + .lock() + .map_err(|_| AdvisorError::Unreachable("lock poisoned".into()))?; + + match &*guard { + Ok(rec) => Ok(rec.clone()), + Err(msg) => Err(AdvisorError::Unreachable(msg.clone())), + } + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_noop_advisor_returns_continue() { + let advisor = NoOpAdvisor; + let ctx = AnalysisContext { + rollout_name: "my-app".into(), + namespace: "default".into(), + strategy: "canary".into(), + current_step: Some(1), + current_weight: Some(20), + metrics_healthy: true, + phase: "Progressing".into(), + history: vec![], + }; + + let rec = advisor.advise(&ctx).await.unwrap(); + assert_eq!(rec.action, RecommendedAction::Continue); + assert_eq!(rec.confidence, 0.0); + } + + #[tokio::test] + async fn test_mock_advisor_returns_configured_response() { + let advisor = MockAdvisor::new(Recommendation { + action: RecommendedAction::Rollback, + confidence: 0.95, + reasoning: "high error rate detected".into(), + }); + + let ctx = AnalysisContext { + rollout_name: "my-app".into(), + namespace: "default".into(), + strategy: "canary".into(), + current_step: Some(2), + current_weight: Some(40), + metrics_healthy: false, + phase: "Progressing".into(), + history: vec![], + }; + + let rec = advisor.advise(&ctx).await.unwrap(); + assert_eq!(rec.action, RecommendedAction::Rollback); + assert_eq!(rec.confidence, 0.95); + assert_eq!(advisor.calls(), 1); + } + + #[tokio::test] + async fn test_mock_advisor_tracks_call_count() { + let advisor = MockAdvisor::new(Recommendation { + action: RecommendedAction::Continue, + confidence: 0.8, + reasoning: "looks good".into(), + }); + + let ctx = AnalysisContext { + rollout_name: "test".into(), + namespace: "default".into(), + strategy: "canary".into(), + current_step: None, + current_weight: None, + metrics_healthy: true, + phase: "Progressing".into(), + history: vec![], + }; + + let _ = advisor.advise(&ctx).await; + let _ = advisor.advise(&ctx).await; + let _ = advisor.advise(&ctx).await; + assert_eq!(advisor.calls(), 3); + } + + #[tokio::test] + async fn test_mock_advisor_error() { + let advisor = MockAdvisor::new_failing("connection refused"); + + let ctx = AnalysisContext { + rollout_name: "test".into(), + namespace: "default".into(), + strategy: "canary".into(), + current_step: None, + current_weight: None, + metrics_healthy: true, + phase: "Progressing".into(), + history: vec![], + }; + + let result = advisor.advise(&ctx).await; + assert!(result.is_err()); + assert!(result + .unwrap_err() + .to_string() + .contains("connection refused")); + } +} diff --git a/src/controller/cdevents_test.rs b/src/controller/cdevents_test.rs index 6f057b9..3beeda6 100644 --- a/src/controller/cdevents_test.rs +++ b/src/controller/cdevents_test.rs @@ -38,6 +38,7 @@ async fn test_emit_service_deployed_on_initialization() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, // No status yet - this is a new rollout }; @@ -139,6 +140,7 @@ async fn test_emit_service_upgraded_on_step_progression() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -247,6 +249,7 @@ async fn test_emit_service_rolledback_on_failure() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -361,6 +364,7 @@ async fn test_emit_service_published_on_completion() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -469,6 +473,7 @@ async fn test_cdevent_contains_kulta_custom_data() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -548,6 +553,7 @@ async fn test_simple_strategy_emits_deployed_and_published() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -623,6 +629,7 @@ async fn test_blue_green_emits_deployed_on_preview() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -701,6 +708,7 @@ async fn test_blue_green_emits_published_on_promotion() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -791,6 +799,7 @@ async fn test_emit_experiment_concluded_event() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { phase: Some(Phase::Experimenting), @@ -816,6 +825,7 @@ async fn test_emit_experiment_concluded_event() { winner: Some(ABVariant::B), conclusion_reason: Some(ABConclusionReason::ConsensusReached), }), + last_decision_source: None, ..Default::default() }; @@ -845,12 +855,12 @@ async fn test_emit_experiment_concluded_event() { let kulta = &json["customData"]["kulta"]; assert_eq!(kulta["strategy"], "ab-testing"); assert_eq!(kulta["experiment"]["winner"], "B"); - assert_eq!( - kulta["experiment"]["conclusion_reason"], - "ConsensusReached" - ); + assert_eq!(kulta["experiment"]["conclusion_reason"], "ConsensusReached"); assert_eq!(kulta["experiment"]["sample_size_a"], 5000); - assert!(!kulta["experiment"]["metrics"].as_array().unwrap().is_empty()); + assert!(!kulta["experiment"]["metrics"] + .as_array() + .unwrap() + .is_empty()); } // Test A/B initialization event (None → Experimenting = service.deployed) @@ -893,6 +903,7 @@ async fn test_emit_service_deployed_on_ab_initialization() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, // No previous status → initialization }; diff --git a/src/controller/mod.rs b/src/controller/mod.rs index d62098b..d0097d7 100644 --- a/src/controller/mod.rs +++ b/src/controller/mod.rs @@ -1,3 +1,4 @@ +pub mod advisor; pub mod cdevents; pub mod clock; pub mod occurrence; diff --git a/src/controller/occurrence.rs b/src/controller/occurrence.rs index 3f7d3a9..d2516c7 100644 --- a/src/controller/occurrence.rs +++ b/src/controller/occurrence.rs @@ -7,136 +7,18 @@ //! cross-tool correlation by AHTI. //! //! KULTA emits both CDEvents (standard) and FALSE Protocol occurrences (AHTI integration). +//! +//! Types are provided by the `false-protocol` crate — KULTA only contains +//! the mapping logic from rollout state to occurrences. use crate::controller::clock::Clock; -use crate::crd::rollout::{Phase, Rollout}; +use crate::crd::rollout::{Phase, Recommendation, Rollout}; use chrono::{DateTime, Utc}; -use serde::Serialize; +use false_protocol::{Entity, Error as OccurrenceError, Occurrence, Outcome, Severity}; use std::collections::HashMap; use std::sync::Arc; use tracing::warn; -/// FALSE Protocol occurrence -#[derive(Debug, Serialize)] -pub struct Occurrence { - pub id: String, - pub timestamp: String, - pub source: String, - #[serde(rename = "type")] - pub occurrence_type: String, - pub severity: Severity, - pub outcome: Outcome, - pub context: OccurrenceContext, - #[serde(skip_serializing_if = "Option::is_none")] - pub error: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub reasoning: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub history: Option, - #[serde(skip_serializing_if = "HashMap::is_empty", default)] - pub data: HashMap, - #[serde(skip_serializing_if = "Vec::is_empty", default)] - pub entities: Vec, -} - -/// Severity levels -#[derive(Debug, Serialize, Clone)] -#[serde(rename_all = "snake_case")] -pub enum Severity { - Debug, - Info, - Warning, - Error, - Critical, -} - -/// Outcome of the occurrence -#[derive(Debug, Serialize, Clone)] -#[serde(rename_all = "snake_case")] -pub enum Outcome { - Success, - Failure, - Timeout, - InProgress, - Unknown, -} - -/// Occurrence context -#[derive(Debug, Serialize)] -pub struct OccurrenceContext { - #[serde(skip_serializing_if = "Option::is_none")] - pub cluster: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub namespace: Option, - #[serde(skip_serializing_if = "Vec::is_empty", default)] - pub correlation_keys: Vec, -} - -/// Correlation key for cross-tool linking -#[derive(Debug, Serialize)] -pub struct CorrelationKey { - #[serde(rename = "type")] - pub key_type: String, - pub value: String, -} - -/// AI-native error block -#[derive(Debug, Default, Serialize)] -pub struct OccurrenceError { - pub code: String, - pub what_failed: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub why_it_matters: Option, - #[serde(skip_serializing_if = "Vec::is_empty", default)] - pub possible_causes: Vec, - #[serde(skip_serializing_if = "Option::is_none")] - pub suggested_fix: Option, -} - -/// AI-native reasoning block -#[derive(Debug, Serialize)] -pub struct Reasoning { - pub summary: String, - pub explanation: String, - pub confidence: f64, - #[serde(skip_serializing_if = "Vec::is_empty", default)] - pub recommendations: Vec, -} - -/// History of steps taken -#[derive(Debug, Serialize)] -pub struct History { - #[serde(skip_serializing_if = "Option::is_none")] - pub duration_ms: Option, - pub steps: Vec, -} - -/// Individual step in history -#[derive(Debug, Serialize)] -pub struct HistoryStep { - pub description: String, - pub status: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub duration_ms: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub error: Option, -} - -/// Entity reference with version tracking -#[derive(Debug, Serialize)] -pub struct Entity { - #[serde(rename = "type")] - pub entity_type: String, - pub id: String, - pub name: String, - pub version: String, - pub observed_at: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub namespace: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub source_of_truth: Option, -} - /// Map phase transition to occurrence type suffix /// /// Returns just the action suffix (e.g., "failed", "completed"). @@ -215,7 +97,10 @@ pub fn emit_occurrence( } }; let now = clock.now(); - let occurrence = build_occurrence(rollout, old_phase, new_phase, strategy, now); + let occurrence = match build_occurrence(rollout, old_phase, new_phase, strategy, now) { + Some(occ) => occ, + None => return, + }; let json = match serde_json::to_string(&occurrence) { Ok(j) => j, @@ -232,15 +117,18 @@ pub fn emit_occurrence( } } -/// Build an occurrence from rollout state +/// Build an occurrence from rollout state. +/// +/// Returns `None` if the crate's validation rejects the occurrence type +/// (should not happen with well-formed strategy names, but we never +/// fail reconciliation on occurrence emission). fn build_occurrence( rollout: &Rollout, old_phase: Option<&Phase>, new_phase: &Phase, strategy: &str, now: DateTime, -) -> Occurrence { - // name and namespace are validated by emit_occurrence before calling this +) -> Option { let name = rollout.metadata.name.as_deref().unwrap_or("unknown"); let namespace = rollout.metadata.namespace.as_deref().unwrap_or("unknown"); let uid = rollout.metadata.uid.as_deref().unwrap_or(""); @@ -269,58 +157,104 @@ fn build_occurrence( .as_ref() .and_then(|s| s.message.clone()) .unwrap_or_else(|| "Rollout failed".to_string()); + + let current_weight = rollout.status.as_ref().and_then(|s| s.current_weight); + let current_step = rollout.status.as_ref().and_then(|s| s.current_step_index); + + // Build rich context: what_failed includes traffic context + let what_failed = match (current_weight, current_step) { + (Some(weight), Some(step)) => format!( + "{} for {} {} at step {}/{} ({}% traffic)", + message, + name, + strategy, + step + 1, + rollout + .spec + .strategy + .canary + .as_ref() + .map(|c| c.steps.len()) + .unwrap_or(0), + weight, + ), + _ => format!("Rollout {} failed during {} deployment", name, strategy), + }; + + // Richer possible_causes based on failure message + let mut possible_causes = vec![message.clone()]; + if message.contains("metrics exceeded") || message.contains("error rate") { + possible_causes.push(format!("New code path in {} handlers", name)); + possible_causes.push("Downstream service degradation".to_string()); + } + if message.contains("deadline exceeded") { + possible_causes.push("Pods failing readiness probes".to_string()); + possible_causes.push("Image pull failures or resource constraints".to_string()); + } + + let suggested_fix = if message.contains("metrics exceeded") { + format!( + "Rollback {} to stable, check dependent service health, review recent changes", + name + ) + } else { + format!( + "Check metrics and pod logs for {}, consider manual rollback", + name + ) + }; + Some(OccurrenceError { code: "ROLLOUT_FAILED".to_string(), - what_failed: format!("Rollout {} failed during {} deployment", name, strategy), + what_failed, why_it_matters: Some(format!( - "Service {} in namespace {} may be serving degraded traffic", - name, namespace - )), - possible_causes: vec![message], - suggested_fix: Some(format!( - "Check metrics for {} and consider manual rollback", - name + "Service {} in namespace {} may be serving degraded traffic to {}% of requests", + name, + namespace, + current_weight.unwrap_or(0), )), + possible_causes, + suggested_fix: Some(suggested_fix), + ..Default::default() }) } else { None }; - Occurrence { - id: ulid::Ulid::new().to_string(), - timestamp: now.to_rfc3339(), - source: "kulta".to_string(), - occurrence_type: occurrence_type.to_string(), - severity, - outcome, - context: OccurrenceContext { - cluster: std::env::var("KULTA_CLUSTER_NAME").ok(), - namespace: Some(namespace.to_string()), - correlation_keys: vec![ - CorrelationKey { - key_type: "deployment".to_string(), - value: name.to_string(), - }, - CorrelationKey { - key_type: "namespace".to_string(), - value: namespace.to_string(), - }, - ], - }, - error, - reasoning: None, // AHTI adds reasoning, not KULTA - history: None, - data, - entities: vec![Entity { - entity_type: "rollout".to_string(), - id: uid.to_string(), - name: name.to_string(), - version: resource_version.to_string(), - observed_at: now.to_rfc3339(), - namespace: Some(namespace.to_string()), - source_of_truth: Some("k8s-api".to_string()), - }], + let mut entity = Entity::from_k8s("rollout", uid, name, namespace, resource_version); + entity.observed_at = now; + + let mut occ = match Occurrence::new("kulta", &occurrence_type) { + Ok(o) => o, + Err(errs) => { + warn!( + errors = ?errs, + occurrence_type = %occurrence_type, + "Failed to construct FALSE Protocol occurrence (non-fatal)" + ); + return None; + } + }; + + occ.timestamp = now; + occ = occ + .severity(severity) + .outcome(outcome) + .in_namespace(namespace) + .correlate("deployment", name) + .correlate("namespace", namespace) + .with_entity(entity) + .with_data(data); + + if let Ok(cluster) = std::env::var("KULTA_CLUSTER_NAME") { + occ = occ.in_cluster(&cluster); + } + + if let Some(err) = error { + occ = occ.with_error(err); } + + Some(occ) } /// Get the occurrence output directory. @@ -363,6 +297,87 @@ fn write_occurrence(json: &str) -> std::io::Result<()> { Ok(()) } +/// Emit a FALSE Protocol occurrence for an advisor consultation (Level 2+) +/// +/// Emits `{strategy}.advisor.recommendation` events that record what the +/// advisor recommended alongside the threshold decision. +pub fn emit_advisor_occurrence( + rollout: &Rollout, + strategy: &str, + recommendation: &Recommendation, + threshold_healthy: bool, + clock: &Arc, +) { + let name = match rollout.metadata.name.as_deref() { + Some(n) => n, + None => return, + }; + let namespace = match rollout.metadata.namespace.as_deref() { + Some(ns) => ns, + None => return, + }; + let uid = rollout.metadata.uid.as_deref().unwrap_or(""); + let resource_version = rollout.metadata.resource_version.as_deref().unwrap_or("0"); + let now = clock.now(); + + let prefix = match strategy { + "blue_green" => "bluegreen", + "ab_testing" => "abtesting", + "simple" => "rolling", + other => other, + }; + let occurrence_type = format!("{}.advisor.recommendation", prefix); + + let mut occ = match Occurrence::new("kulta", &occurrence_type) { + Ok(o) => o, + Err(errs) => { + warn!(errors = ?errs, "Failed to construct advisor occurrence (non-fatal)"); + return; + } + }; + + let mut data = HashMap::new(); + data.insert( + "advisor".to_string(), + serde_json::json!({ + "action": recommendation.action, + "confidence": recommendation.confidence, + "reasoning": recommendation.reasoning, + "threshold_healthy": threshold_healthy, + "threshold_prevails": true, + }), + ); + + let mut entity = Entity::from_k8s("rollout", uid, name, namespace, resource_version); + entity.observed_at = now; + + occ.timestamp = now; + occ = occ + .severity(Severity::Info) + .outcome(Outcome::InProgress) + .in_namespace(namespace) + .correlate("deployment", name) + .correlate("namespace", namespace) + .with_entity(entity) + .with_data(data); + + if let Ok(cluster) = std::env::var("KULTA_CLUSTER_NAME") { + occ = occ.in_cluster(&cluster); + } + + let json = match serde_json::to_string(&occ) { + Ok(j) => j, + Err(e) => { + warn!(error = %e, "Failed to serialize advisor occurrence (non-fatal)"); + return; + } + }; + + if let Err(e) = write_occurrence(&json) { + warn!(error = %e, "Failed to write advisor occurrence (non-fatal)"); + } +} + #[cfg(test)] #[allow(clippy::unwrap_used, clippy::expect_used)] mod tests { @@ -405,6 +420,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -428,12 +444,10 @@ mod tests { phase_to_occurrence_suffix(Some(&Phase::Progressing), &Phase::Paused), "paused" ); - // Concluded maps to completed assert_eq!( phase_to_occurrence_suffix(Some(&Phase::Experimenting), &Phase::Concluded), "completed" ); - // Paused from non-Progressing phase still maps to paused assert_eq!( phase_to_occurrence_suffix(Some(&Phase::Preview), &Phase::Paused), "paused" @@ -442,22 +456,18 @@ mod tests { #[test] fn test_build_occurrence_type_strategy_prefixes() { - // Canary strategy assert_eq!( build_occurrence_type("canary", None, &Phase::Progressing), "canary.rollout.progressing" ); - // Blue-green strategy assert_eq!( build_occurrence_type("blue_green", None, &Phase::Completed), "bluegreen.rollout.completed" ); - // A/B testing strategy assert_eq!( build_occurrence_type("ab_testing", Some(&Phase::Experimenting), &Phase::Failed), "abtesting.rollout.failed" ); - // Simple strategy assert_eq!( build_occurrence_type("simple", None, &Phase::Completed), "rolling.rollout.completed" @@ -469,16 +479,16 @@ mod tests { let rollout = test_rollout(); let now = Utc::now(); - let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now); + let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now).unwrap(); assert_eq!(occ.source, "kulta"); assert_eq!(occ.occurrence_type, "canary.rollout.progressing"); - assert!(matches!(occ.severity, Severity::Info)); - assert!(matches!(occ.outcome, Outcome::InProgress)); + assert_eq!(occ.severity, Severity::Info); + assert_eq!(occ.outcome, Outcome::InProgress); assert!(occ.error.is_none()); - assert_eq!(occ.entities.len(), 1); - assert_eq!(occ.entities[0].name, "my-app"); - assert_eq!(occ.entities[0].version, "rv-456"); + assert_eq!(occ.context.entities.len(), 1); + assert_eq!(occ.context.entities[0].name, "my-app"); + assert_eq!(occ.context.entities[0].version, "rv-456"); assert_eq!(occ.context.namespace.as_deref(), Some("production")); } @@ -493,11 +503,12 @@ mod tests { &Phase::Failed, "canary", now, - ); + ) + .unwrap(); assert_eq!(occ.occurrence_type, "canary.rollout.failed"); - assert!(matches!(occ.severity, Severity::Error)); - assert!(matches!(occ.outcome, Outcome::Failure)); + assert_eq!(occ.severity, Severity::Error); + assert_eq!(occ.outcome, Outcome::Failure); assert!(occ.error.is_some()); let err = occ.error.as_ref().unwrap(); assert_eq!(err.code, "ROLLOUT_FAILED"); @@ -510,14 +521,14 @@ mod tests { let rollout = test_rollout(); let now = Utc::now(); - let occ = build_occurrence(&rollout, None, &Phase::Completed, "simple", now); + let occ = build_occurrence(&rollout, None, &Phase::Completed, "simple", now).unwrap(); let json = serde_json::to_string(&occ).expect("Should serialize"); - // Verify key fields present in JSON assert!(json.contains("\"source\":\"kulta\"")); assert!(json.contains("\"type\":\"rolling.rollout.completed\"")); assert!(json.contains("\"severity\":\"info\"")); assert!(json.contains("\"outcome\":\"success\"")); + assert!(json.contains("\"protocol_version\":\"1.0\"")); // No error block for success assert!(!json.contains("\"error\"")); // No reasoning (AHTI's job) @@ -529,7 +540,7 @@ mod tests { let rollout = test_rollout(); let now = Utc::now(); - let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now); + let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now).unwrap(); // ULID is 26 characters, uppercase alphanumeric assert_eq!(occ.id.len(), 26); @@ -548,15 +559,15 @@ mod tests { #[test] fn test_build_occurrence_with_missing_metadata() { let mut rollout = test_rollout(); - rollout.metadata = ObjectMeta::default(); // no name, namespace, uid, resource_version + rollout.metadata = ObjectMeta::default(); let now = Utc::now(); - let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now); + let occ = build_occurrence(&rollout, None, &Phase::Progressing, "canary", now).unwrap(); - assert_eq!(occ.entities[0].name, "unknown"); + assert_eq!(occ.context.entities[0].name, "unknown"); assert_eq!(occ.context.namespace.as_deref(), Some("unknown")); - assert_eq!(occ.entities[0].id, ""); // uid defaults to "" - assert_eq!(occ.entities[0].version, "0"); // resource_version defaults to "0" + assert_eq!(occ.context.entities[0].id, ""); + assert_eq!(occ.context.entities[0].version, "0"); } #[test] @@ -581,7 +592,6 @@ mod tests { #[test] fn test_phase_to_occurrence_suffix_initializing() { - // Initializing maps to "progressing" (catch-all) assert_eq!( phase_to_occurrence_suffix(None, &Phase::Initializing), "progressing" @@ -611,9 +621,123 @@ mod tests { &Phase::Failed, "canary", now, - ); + ) + .unwrap(); let err = occ.error.as_ref().unwrap(); assert!(err.possible_causes[0].contains("High error rate")); } + + #[test] + fn test_build_occurrence_failed_with_metrics_exceeded_has_rich_context() { + use crate::crd::rollout::{ + CanaryStep, CanaryStrategy, RolloutStatus, RolloutStrategy as RolloutStrategySpec, + }; + + let mut rollout = test_rollout(); + rollout.spec.strategy = RolloutStrategySpec { + canary: Some(CanaryStrategy { + canary_service: "my-app-canary".into(), + stable_service: "my-app-stable".into(), + port: None, + steps: vec![ + CanaryStep { + set_weight: Some(20), + pause: None, + }, + CanaryStep { + set_weight: Some(50), + pause: None, + }, + CanaryStep { + set_weight: Some(100), + pause: None, + }, + ], + traffic_routing: None, + analysis: None, + }), + blue_green: None, + simple: None, + ab_testing: None, + }; + rollout.status = Some(RolloutStatus { + phase: Some(Phase::Progressing), + current_weight: Some(20), + current_step_index: Some(0), + message: Some("Rollback triggered: metrics exceeded thresholds".into()), + ..Default::default() + }); + + let now = Utc::now(); + let occ = build_occurrence( + &rollout, + Some(&Phase::Progressing), + &Phase::Failed, + "canary", + now, + ) + .unwrap(); + + let err = occ.error.as_ref().unwrap(); + // Rich what_failed includes step and weight context + assert!( + err.what_failed.contains("step 1/3"), + "what_failed: {}", + err.what_failed + ); + assert!( + err.what_failed.contains("20% traffic"), + "what_failed: {}", + err.what_failed + ); + // Multiple possible causes for metrics failures + assert!(err.possible_causes.len() > 1); + // why_it_matters includes traffic percentage + assert!(err.why_it_matters.as_ref().unwrap().contains("20%")); + // suggested_fix mentions rollback + assert!(err.suggested_fix.as_ref().unwrap().contains("Rollback")); + } + + #[test] + fn test_build_occurrence_failed_deadline_exceeded_has_rich_context() { + let mut rollout = test_rollout(); + rollout.status = Some(crate::crd::rollout::RolloutStatus { + message: Some("Progress deadline exceeded: no progress made in 600 seconds".into()), + ..Default::default() + }); + + let now = Utc::now(); + let occ = build_occurrence( + &rollout, + Some(&Phase::Progressing), + &Phase::Failed, + "canary", + now, + ) + .unwrap(); + + let err = occ.error.as_ref().unwrap(); + assert!(err + .possible_causes + .iter() + .any(|c| c.contains("readiness probes"))); + } + + #[test] + fn test_emit_advisor_occurrence_does_not_panic() { + use crate::crd::rollout::{Recommendation, RecommendedAction}; + + let rollout = test_rollout(); + let clock: Arc = Arc::new(MockClock::new(Utc::now())); + + let recommendation = Recommendation { + action: RecommendedAction::Continue, + confidence: 0.85, + reasoning: "metrics look healthy, no anomalies detected".into(), + }; + + // Should not panic even if file write fails in test env + emit_advisor_occurrence(&rollout, "canary", &recommendation, true, &clock); + } } diff --git a/src/controller/prometheus_ab.rs b/src/controller/prometheus_ab.rs index 842293c..47bd4fc 100644 --- a/src/controller/prometheus_ab.rs +++ b/src/controller/prometheus_ab.rs @@ -484,20 +484,16 @@ mod tests { #[test] fn test_calculate_ab_significance_both_zero_rates() { - let result = calculate_ab_significance( - 0.0, 0.0, 10000, 10000, 0.95, - &ABMetricDirection::Lower, - ); + let result = + calculate_ab_significance(0.0, 0.0, 10000, 10000, 0.95, &ABMetricDirection::Lower); assert!(!result.is_significant); assert!((result.effect_size - 0.0).abs() < 0.001); } #[test] fn test_calculate_ab_significance_rate_a_zero_rate_b_positive() { - let result = calculate_ab_significance( - 0.0, 0.05, 10000, 10000, 0.95, - &ABMetricDirection::Lower, - ); + let result = + calculate_ab_significance(0.0, 0.05, 10000, 10000, 0.95, &ABMetricDirection::Lower); // effect_size should be 1.0 when rate_a is 0 and rate_b > 0 assert!((result.effect_size - 1.0).abs() < 0.001); } @@ -506,10 +502,7 @@ mod tests { fn test_calculate_ab_significance_se_zero_guard() { // Both rates identical and non-zero with same sample sizes → se could be very small // but with truly identical rates, z_score = 0, so not significant - let result = calculate_ab_significance( - 0.5, 0.5, 100, 100, 0.95, - &ABMetricDirection::Lower, - ); + let result = calculate_ab_significance(0.5, 0.5, 100, 100, 0.95, &ABMetricDirection::Lower); assert!(!result.is_significant); } diff --git a/src/controller/rollout/reconcile.rs b/src/controller/rollout/reconcile.rs index 79c8bc1..59b1c72 100644 --- a/src/controller/rollout/reconcile.rs +++ b/src/controller/rollout/reconcile.rs @@ -1,7 +1,8 @@ +use crate::controller::advisor::{AnalysisAdvisor, AnalysisContext, NoOpAdvisor}; use crate::controller::cdevents::emit_status_change_event; use crate::controller::occurrence::emit_occurrence; use crate::controller::prometheus::MetricsQuerier; -use crate::crd::rollout::{Phase, Rollout, RolloutStatus}; +use crate::crd::rollout::{AdvisorLevel, Phase, Rollout, RolloutStatus}; use crate::server::LeaderState; use chrono::{DateTime, Utc}; use kube::api::{Api, Patch, PatchParams}; @@ -48,6 +49,7 @@ pub struct Context { pub client: kube::Client, pub cdevents_sink: Arc, pub prometheus_client: Arc, + pub advisor: Arc, pub clock: Arc, /// Optional leader state for multi-replica deployments /// When Some, reconciliation is skipped if not the leader @@ -70,6 +72,7 @@ impl Context { client, cdevents_sink: Arc::new(cdevents_sink), prometheus_client: Arc::new(prometheus_client), + advisor: Arc::new(NoOpAdvisor), clock, leader_state: None, metrics, @@ -92,6 +95,7 @@ impl Context { client, cdevents_sink: Arc::new(cdevents_sink), prometheus_client: Arc::new(prometheus_client), + advisor: Arc::new(NoOpAdvisor), clock, leader_state: Some(leader_state), metrics, @@ -130,6 +134,7 @@ impl Context { client, cdevents_sink: Arc::new(crate::controller::cdevents::MockEventSink::new()), prometheus_client: Arc::new(crate::controller::prometheus::MockPrometheusClient::new()), + advisor: Arc::new(NoOpAdvisor), clock: Arc::new(crate::controller::clock::SystemClock), leader_state: None, metrics: None, @@ -148,6 +153,7 @@ impl Context { client: mock.client, cdevents_sink: mock.cdevents_sink, prometheus_client: mock.prometheus_client, + advisor: mock.advisor, clock: mock.clock, leader_state: Some(leader_state), metrics: None, @@ -227,6 +233,55 @@ pub async fn reconcile(rollout: Arc, ctx: Arc) -> Result { + info!( + rollout = ?name, + advisor_action = ?recommendation.action, + confidence = recommendation.confidence, + reasoning = %recommendation.reasoning, + threshold_healthy = is_healthy, + "Advisor recommendation received (threshold decision prevails)" + ); + // Emit advisor recommendation occurrence + crate::controller::occurrence::emit_advisor_occurrence( + &rollout, + strategy.name(), + &recommendation, + is_healthy, + &ctx.clock, + ); + } + Err(e) => { + warn!( + rollout = ?name, + error = %e, + "Advisor consultation failed, falling back to threshold decision" + ); + } + } + } + if !is_healthy { warn!(rollout = ?name, "Metrics unhealthy, triggering rollback"); @@ -309,6 +364,7 @@ pub async fn reconcile(rollout: Arc, ctx: Arc) -> Result Rollout { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -193,6 +196,7 @@ fn create_test_rollout_with_blue_green() -> Rollout { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -282,6 +286,7 @@ fn test_ab_testing_creates_variant_replicasets() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -431,6 +436,7 @@ fn create_test_rollout_with_canary() -> Rollout { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -499,6 +505,7 @@ async fn test_reconcile_creates_stable_replicaset() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -616,6 +623,7 @@ async fn test_build_replicaset_spec() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -705,6 +713,7 @@ async fn test_reconcile_creates_canary_replicaset() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -800,6 +809,7 @@ async fn test_replicaset_has_kulta_managed_label() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -921,6 +931,7 @@ async fn test_build_both_stable_and_canary_replicasets() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -1056,6 +1067,7 @@ async fn test_calculate_traffic_weights_step0() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), // First step: 20% canary @@ -1109,6 +1121,7 @@ async fn test_calculate_traffic_weights_step1() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(1), // Second step: 50% canary @@ -1156,6 +1169,7 @@ async fn test_calculate_traffic_weights_no_step() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, // No status yet, default to 100% stable }; @@ -1206,6 +1220,7 @@ async fn test_calculate_traffic_weights_complete() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(1), // Last step: 100% canary @@ -1253,6 +1268,7 @@ async fn test_calculate_traffic_weights_beyond_steps() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(5), // Beyond available steps (only 1 step) @@ -1300,6 +1316,7 @@ async fn test_build_httproute_backend_weights() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), // 20% canary @@ -1365,6 +1382,7 @@ async fn test_convert_to_gateway_api_backend_refs() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), // 20% canary @@ -1422,6 +1440,7 @@ async fn test_gateway_api_backend_refs_no_canary_strategy() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -1473,6 +1492,7 @@ async fn test_initialize_rollout_status() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, // No status yet - should be initialized }; @@ -1527,6 +1547,7 @@ async fn test_initialize_sets_progress_started_at() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -1587,6 +1608,7 @@ async fn test_should_progress_to_next_step() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1645,6 +1667,7 @@ async fn test_should_not_progress_when_paused() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1697,6 +1720,7 @@ async fn test_advance_to_next_step() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1762,6 +1786,7 @@ async fn test_advance_preserves_progress_started_at() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1821,6 +1846,7 @@ async fn test_advance_to_final_step() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1886,6 +1912,7 @@ async fn test_compute_desired_status_for_new_rollout() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, // No status - should be initialized }; @@ -1939,6 +1966,7 @@ async fn test_compute_desired_status_progresses_step() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -1997,6 +2025,7 @@ async fn test_compute_desired_status_respects_pause() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -3148,6 +3177,7 @@ async fn test_evaluate_rollout_metrics_healthy() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -3236,6 +3266,7 @@ async fn test_evaluate_rollout_metrics_unhealthy() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -3309,6 +3340,7 @@ async fn test_evaluate_rollout_metrics_no_analysis_config() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { current_step_index: Some(0), @@ -3391,6 +3423,7 @@ async fn test_evaluate_rollout_metrics_skips_during_warmup() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { replicas: 3, @@ -3472,6 +3505,7 @@ async fn test_evaluate_rollout_metrics_runs_after_warmup() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { replicas: 3, @@ -3552,6 +3586,7 @@ async fn test_evaluate_rollout_metrics_no_warmup_configured() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { replicas: 3, @@ -3623,6 +3658,7 @@ async fn test_blue_green_builds_httproute_backend_refs() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { phase: Some(Phase::Preview), @@ -3699,6 +3735,7 @@ async fn test_blue_green_httproute_after_promotion() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: Some(RolloutStatus { phase: Some(Phase::Completed), @@ -4008,10 +4045,7 @@ async fn test_evaluate_ab_max_duration_exceeded() { assert!(result.should_conclude); assert!(result.winner.is_none()); - assert_eq!( - result.reason, - Some(ABConclusionReason::MaxDurationExceeded) - ); + assert_eq!(result.reason, Some(ABConclusionReason::MaxDurationExceeded)); } /// Max duration NOT exceeded → continues to analysis @@ -4098,14 +4132,8 @@ async fn test_evaluate_ab_prometheus_sample_query_failure() { let prom = MockPrometheusClient::new(); prom.enqueue_error(crate::controller::prometheus::PrometheusError::NoData); - let rollout = create_ab_rollout_with_analysis( - &started, - Phase::Experimenting, - None, - None, - None, - None, - ); + let rollout = + create_ab_rollout_with_analysis(&started, Phase::Experimenting, None, None, None, None); let ctx = create_test_context_with_prometheus(prom, now); let result = evaluate_ab_experiment(&rollout, &ctx).await.unwrap(); @@ -4124,14 +4152,8 @@ async fn test_evaluate_ab_prometheus_error_rate_failure() { prom.enqueue_response(1000.0); // sample B prom.enqueue_error(crate::controller::prometheus::PrometheusError::NoData); // rate A fails - let rollout = create_ab_rollout_with_analysis( - &started, - Phase::Experimenting, - None, - None, - None, - None, - ); + let rollout = + create_ab_rollout_with_analysis(&started, Phase::Experimenting, None, None, None, None); let ctx = create_test_context_with_prometheus(prom, now); let result = evaluate_ab_experiment(&rollout, &ctx).await.unwrap(); @@ -4205,14 +4227,8 @@ async fn test_evaluate_ab_no_significance() { async fn test_evaluate_ab_no_analysis_config() { let now = Utc::now(); let started = (now - chrono::Duration::hours(1)).to_rfc3339(); - let mut rollout = create_ab_rollout_with_analysis( - &started, - Phase::Experimenting, - None, - None, - None, - None, - ); + let mut rollout = + create_ab_rollout_with_analysis(&started, Phase::Experimenting, None, None, None, None); // Remove the analysis config if let Some(ab) = &mut rollout.spec.strategy.ab_testing { ab.analysis = None; diff --git a/src/controller/strategies/ab_testing.rs b/src/controller/strategies/ab_testing.rs index 3989ff8..ff39231 100644 --- a/src/controller/strategies/ab_testing.rs +++ b/src/controller/strategies/ab_testing.rs @@ -201,6 +201,7 @@ impl RolloutStrategy for ABTestingStrategyHandler { winner: None, conclusion_reason: None, }), + last_decision_source: None, ..Default::default() } } @@ -453,6 +454,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: phase.map(|p| RolloutStatus { phase: Some(p), @@ -746,6 +748,7 @@ mod tests { winner: Some(ABVariant::B), conclusion_reason: Some(ABConclusionReason::ConsensusReached), }), + last_decision_source: None, ..Default::default() }); @@ -771,6 +774,7 @@ mod tests { winner: None, conclusion_reason: None, // No conclusion yet }), + last_decision_source: None, ..Default::default() }); diff --git a/src/controller/strategies/blue_green.rs b/src/controller/strategies/blue_green.rs index 17eb58f..cb7a747 100644 --- a/src/controller/strategies/blue_green.rs +++ b/src/controller/strategies/blue_green.rs @@ -193,6 +193,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } diff --git a/src/controller/strategies/canary.rs b/src/controller/strategies/canary.rs index 01bd38a..2d48f88 100644 --- a/src/controller/strategies/canary.rs +++ b/src/controller/strategies/canary.rs @@ -169,6 +169,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: current_weight.map(|weight| crate::crd::rollout::RolloutStatus { phase: Some(Phase::Progressing), @@ -183,6 +184,7 @@ mod tests { progress_started_at: None, decisions: vec![], ab_experiment: None, + last_decision_source: None, }), } } diff --git a/src/controller/strategies/mod.rs b/src/controller/strategies/mod.rs index 78161d0..9e3c7f8 100644 --- a/src/controller/strategies/mod.rs +++ b/src/controller/strategies/mod.rs @@ -364,6 +364,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } diff --git a/src/controller/strategies/simple.rs b/src/controller/strategies/simple.rs index 10ef363..51b5b55 100644 --- a/src/controller/strategies/simple.rs +++ b/src/controller/strategies/simple.rs @@ -94,6 +94,7 @@ impl RolloutStrategy for SimpleStrategyHandler { progress_started_at: None, decisions: vec![], ab_experiment: None, + last_decision_source: None, } } @@ -161,6 +162,7 @@ mod tests { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } diff --git a/src/crd/conversion.rs b/src/crd/conversion.rs index 31425f2..1397b1b 100644 --- a/src/crd/conversion.rs +++ b/src/crd/conversion.rs @@ -60,6 +60,7 @@ pub fn convert_to_v1alpha1(spec: &v1beta1::RolloutSpec) -> v1alpha1::RolloutSpec max_surge: spec.max_surge.clone(), max_unavailable: spec.max_unavailable.clone(), progress_deadline_seconds: spec.progress_deadline_seconds, + advisor: Default::default(), } } diff --git a/src/crd/conversion_test.rs b/src/crd/conversion_test.rs index 392cfcc..076d12a 100644 --- a/src/crd/conversion_test.rs +++ b/src/crd/conversion_test.rs @@ -19,6 +19,7 @@ fn test_v1alpha1_to_v1beta1_adds_default_max_surge() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }; let v1beta1_spec = convert_to_v1beta1(&v1alpha1_spec); @@ -38,6 +39,7 @@ fn test_v1alpha1_to_v1beta1_adds_default_max_unavailable() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }; let v1beta1_spec = convert_to_v1beta1(&v1alpha1_spec); @@ -57,6 +59,7 @@ fn test_v1alpha1_to_v1beta1_adds_default_progress_deadline() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }; let v1beta1_spec = convert_to_v1beta1(&v1alpha1_spec); @@ -91,6 +94,7 @@ fn test_v1alpha1_to_v1beta1_preserves_existing_fields() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }; let v1beta1_spec = convert_to_v1beta1(&v1alpha1_spec); @@ -185,6 +189,7 @@ fn test_roundtrip_v1alpha1_to_v1beta1_to_v1alpha1() { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }; let converted = convert_to_v1beta1(&original); diff --git a/src/crd/rollout.rs b/src/crd/rollout.rs index 6315af5..1387ec9 100644 --- a/src/crd/rollout.rs +++ b/src/crd/rollout.rs @@ -55,6 +55,14 @@ pub struct RolloutSpec { skip_serializing_if = "Option::is_none" )] pub progress_deadline_seconds: Option, + + /// AI advisor configuration for progressive AI adoption + #[serde(default, skip_serializing_if = "is_default_advisor_config")] + pub advisor: AdvisorConfig, +} + +fn is_default_advisor_config(c: &AdvisorConfig) -> bool { + c.level == AdvisorLevel::Off && c.endpoint.is_none() && c.timeout_seconds == 10 } fn default_replicas() -> i32 { @@ -532,6 +540,10 @@ pub struct RolloutStatus { /// A/B experiment status (only for abTesting strategy) #[serde(rename = "abExperiment", skip_serializing_if = "Option::is_none")] pub ab_experiment: Option, + + /// Source of last analysis decision (Threshold, Advisor, Human) + #[serde(rename = "lastDecisionSource", skip_serializing_if = "Option::is_none")] + pub last_decision_source: Option, } /// A/B experiment status tracking @@ -614,6 +626,81 @@ pub enum ABConclusionReason { ConsensusReached, } +/// AI advisor integration level +/// +/// Progressive adoption ladder for AI-assisted rollout decisions. +/// Each level adds capability while preserving backward compatibility. +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, Default, PartialEq, Eq)] +pub enum AdvisorLevel { + /// No AI integration (default, current behavior) + #[default] + Off, + /// Rich FALSE Protocol events with AI-native context fields + Context, + /// AI analyzes metrics and returns recommendations (threshold still decides) + Advised, + /// AI proposes rollout strategy (human approves) — future + Planned, + /// AI drives the loop with human override — future + Driven, +} + +/// Configuration for the AI advisor +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, Default)] +pub struct AdvisorConfig { + /// Advisor integration level + #[serde(default)] + pub level: AdvisorLevel, + + /// URL of AI advisory service (e.g., MCP endpoint, HTTP API) + #[serde(skip_serializing_if = "Option::is_none")] + pub endpoint: Option, + + /// Timeout for advisory calls in seconds + #[serde( + rename = "timeoutSeconds", + default = "default_advisor_timeout", + skip_serializing_if = "is_default_advisor_timeout" + )] + pub timeout_seconds: u64, +} + +fn default_advisor_timeout() -> u64 { + 10 +} + +fn is_default_advisor_timeout(v: &u64) -> bool { + *v == 10 +} + +/// What the advisor recommends after analysis +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct Recommendation { + pub action: RecommendedAction, + pub confidence: f64, + pub reasoning: String, +} + +/// Recommended action from the advisor +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub enum RecommendedAction { + Continue, + Pause, + Rollback, + Advance { + #[serde(rename = "toWeight")] + to_weight: u32, + }, +} + +/// Tracks where a decision came from +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub enum DecisionSource { + Threshold, + Advisor { confidence: String }, + Human, +} + #[cfg(test)] #[path = "rollout_test.rs"] mod tests; diff --git a/src/crd/rollout_test.rs b/src/crd/rollout_test.rs index b48061e..f334b97 100644 --- a/src/crd/rollout_test.rs +++ b/src/crd/rollout_test.rs @@ -407,6 +407,7 @@ fn test_ab_experiment_status_serialization() { winner: None, conclusion_reason: None, }), + last_decision_source: None, ..Default::default() }; diff --git a/tests/seppo_integration_test.rs b/tests/seppo_integration_test.rs index 4b87f1c..5c494f8 100644 --- a/tests/seppo_integration_test.rs +++ b/tests/seppo_integration_test.rs @@ -258,6 +258,7 @@ async fn test_canary_full_lifecycle(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -394,6 +395,7 @@ async fn test_canary_pause_and_promote(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -495,6 +497,7 @@ async fn test_status_decisions_tracking(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -582,6 +585,7 @@ async fn test_blue_green_promotion(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -717,6 +721,7 @@ async fn test_blue_green_auto_promotion(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -811,6 +816,7 @@ async fn test_httproute_weight_updates(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -918,6 +924,7 @@ async fn test_simple_strategy_lifecycle(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -1008,6 +1015,7 @@ async fn test_image_update_triggers_rollout(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; @@ -1060,6 +1068,7 @@ async fn test_image_update_triggers_rollout(ctx: TestContext) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, }; diff --git a/tests/stress_test.rs b/tests/stress_test.rs index 5afb7b7..266c48d 100644 --- a/tests/stress_test.rs +++ b/tests/stress_test.rs @@ -118,6 +118,7 @@ fn create_rollout(name: &str, namespace: &str, replicas: i32, image: &str) -> Ro max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -180,6 +181,7 @@ fn create_rollout_with_pauses( max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, } @@ -821,6 +823,7 @@ async fn test_edge_minimal_steps(ctx: Context) { max_surge: None, max_unavailable: None, progress_deadline_seconds: None, + advisor: Default::default(), }, status: None, };