diff --git a/controller/src/backend_client.rs b/controller/src/backend_client.rs new file mode 100644 index 00000000..c79b4a7f --- /dev/null +++ b/controller/src/backend_client.rs @@ -0,0 +1,190 @@ +//! HTTP client that pushes a freshly-generated `StreamingConfig` YAML +//! to the ASAPQuery-backend's `POST /api/v1/streaming-config` endpoint. +//! +//! This is the controller-side **producer** of the PR E phase 1 / phase 2 +//! hot-reload contract that landed in ASAPQuery-backend PRs #10 and #12. +//! The replanner calls into this module immediately after generating a +//! new plan so the backend's active `StreamingConfig` is updated without +//! a restart and subsequent queries observe the new aggregation layout. +//! +//! The client is **fire-and-forget at the call site** — the replanner +//! awaits the POST but doesn't block its own return on the outcome. +//! Errors are logged at WARN; the controller is expected to be tolerant +//! of transient backend unavailability because the next replan cycle +//! will try again with the latest plan. + +use std::time::Duration; + +use anyhow::{Context, Result}; +use reqwest::Client; +use tracing::{debug, warn}; + +/// Minimal HTTP client for ASAPQuery-backend's streaming-config endpoint. +/// Built once at controller startup from the `CONTROLLER_BACKEND_ENDPOINT` +/// environment variable and shared via `Arc` with the replanner. +#[derive(Debug, Clone)] +pub struct BackendClient { + endpoint: String, + http: Client, +} + +impl BackendClient { + /// Construct a client pointing at the backend's plan-push endpoint. + /// `endpoint` should be the full URL, e.g. + /// `http://backend.svc:8088/api/v1/streaming-config`. + /// + /// A 5-second timeout bounds the duration a slow or unreachable + /// backend can stall the replanner — consistent with the symmetric + /// 5-second timeout on ASAPQuery-backend's `HttpControllerClient` + /// (the reverse direction in the same loop). + pub fn new(endpoint: impl Into) -> Self { + let http = Client::builder() + .timeout(Duration::from_secs(5)) + .build() + .unwrap_or_else(|_| Client::new()); + Self { + endpoint: endpoint.into(), + http, + } + } + + /// Construct with an explicit `reqwest::Client`. Used by tests that + /// need to inject a mock-server URL without reconfiguring the + /// timeout setup. + pub fn with_http(endpoint: impl Into, http: Client) -> Self { + Self { + endpoint: endpoint.into(), + http, + } + } + + pub fn endpoint(&self) -> &str { + &self.endpoint + } + + /// POST the given `StreamingConfig` YAML to the backend. Returns + /// `Ok(())` on any 2xx status, otherwise an error carrying the + /// status code and response body. The caller (typically + /// [`Replanner::replan_metric`]) logs the error and moves on — the + /// next replan cycle will retry with the latest plan. + pub async fn push_streaming_config(&self, yaml: String) -> Result<()> { + debug!( + endpoint = %self.endpoint, + yaml_bytes = yaml.len(), + "pushing streaming-config YAML to ASAPQuery-backend" + ); + let resp = self + .http + .post(&self.endpoint) + .header("content-type", "application/x-yaml") + .body(yaml) + .send() + .await + .context("failed to POST streaming-config to backend")?; + + let status = resp.status(); + if status.is_success() { + Ok(()) + } else { + let body = resp.text().await.unwrap_or_default(); + Err(anyhow::anyhow!( + "backend returned {} for streaming-config POST: {}", + status, + body + )) + } + } +} + +/// Fire-and-forget convenience helper used by the replanner. Logs +/// errors at WARN and never propagates them — the replanner should +/// never fail an entire replan because the backend was temporarily +/// unreachable. +pub async fn push_or_log(client: &BackendClient, metric: &str, yaml: String) { + match client.push_streaming_config(yaml).await { + Ok(()) => { + debug!(metric, endpoint = %client.endpoint, "streaming-config push succeeded"); + } + Err(e) => { + warn!( + metric, + endpoint = %client.endpoint, + error = %e, + "streaming-config push to ASAPQuery-backend failed; \ + next replan cycle will retry" + ); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use axum::extract::State; + use axum::routing::post; + use axum::Router; + use std::sync::{Arc as StdArc, Mutex}; + + #[derive(Clone)] + struct SharedSink(StdArc>>); + + async fn start_mock_backend(sink: SharedSink, status: axum::http::StatusCode) -> String { + let app = Router::new() + .route( + "/api/v1/streaming-config", + post( + move |State(sink): State, body: axum::body::Bytes| async move { + let yaml = String::from_utf8_lossy(&body).to_string(); + sink.0.lock().unwrap().push(yaml); + status + }, + ), + ) + .with_state(sink); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + tokio::time::sleep(Duration::from_millis(50)).await; + format!("http://{addr}/api/v1/streaming-config") + } + + #[tokio::test] + async fn success_path_round_trips_yaml() { + let sink = SharedSink(StdArc::new(Mutex::new(Vec::new()))); + let url = start_mock_backend(sink.clone(), axum::http::StatusCode::OK).await; + + let client = BackendClient::new(url); + let yaml = "aggregations:\n - aggregationId: 42\n metric: cpu\n".to_string(); + client + .push_streaming_config(yaml.clone()) + .await + .expect("push ok"); + + let received = sink.0.lock().unwrap(); + assert_eq!(received.len(), 1); + assert_eq!(received[0], yaml); + } + + #[tokio::test] + async fn non_2xx_status_is_reported_as_error() { + let sink = SharedSink(StdArc::new(Mutex::new(Vec::new()))); + let url = + start_mock_backend(sink.clone(), axum::http::StatusCode::INTERNAL_SERVER_ERROR).await; + + let client = BackendClient::new(url); + let result = client.push_streaming_config("whatever".to_string()).await; + assert!(result.is_err(), "expected error on 500, got {result:?}"); + let msg = result.unwrap_err().to_string(); + assert!(msg.contains("500"), "error msg should mention 500: {msg}"); + } + + #[tokio::test] + async fn push_or_log_swallows_errors() { + // Point at an unreachable port so the request fails fast. + let client = BackendClient::new("http://127.0.0.1:1/api/v1/streaming-config"); + // Must not panic or propagate — fire-and-forget semantics. + push_or_log(&client, "cpu_usage", "content".to_string()).await; + } +} diff --git a/controller/src/config/asapquery_backend.rs b/controller/src/config/asapquery_backend.rs new file mode 100644 index 00000000..665472d2 --- /dev/null +++ b/controller/src/config/asapquery_backend.rs @@ -0,0 +1,299 @@ +//! Convert a [`CollectionPlan`] into the YAML shape ASAPQuery-backend's +//! `POST /api/v1/streaming-config` endpoint accepts (the same format its +//! `StreamingConfig::from_yaml_data` parser consumes at startup). +//! +//! This is **separate** from `config::backend` (which produces OTel YAML +//! for a backend OTel collector running merge processors). The two +//! consumers are different: +//! +//! * `config::backend` — OTel collector, expects +//! `processors: { ddsketch_merge: {...} }` + `service.pipelines`. +//! * `config::asapquery_backend` (this module) — ASAPQuery-backend +//! query engine, expects +//! `aggregations: [{ aggregationId, aggregationType, metric, labels, +//! parameters, windowSize, windowType, spatialFilter }]`. +//! +//! Both are generated from the same `CollectionPlan` fields but target +//! different services. The replanner pushes the OTel YAML via OpAMP to +//! backend-role collectors and pushes this one via HTTP to the +//! ASAPQuery-backend's `/api/v1/streaming-config` endpoint. + +use std::time::Duration; + +use anyhow::{Context, Result}; + +use crate::types::{AgentCollectorConfig, CollectionPlan, SketchType}; + +/// Stable aggregation ID used when the planner has no explicit id to +/// assign. The ASAPQuery-backend uses `u64` agg IDs; we derive one +/// deterministically from the metric name so the same metric always +/// maps to the same id across successive pushes (otherwise the backend +/// would grow unbounded as each replan introduces a new agg_id). +pub fn deterministic_agg_id(metric: &str) -> u64 { + use std::collections::hash_map::DefaultHasher; + use std::hash::{Hash, Hasher}; + let mut h = DefaultHasher::new(); + metric.hash(&mut h); + // Bias away from 0 so the id space is [1, u64::MAX]; 0 is reserved + // in some of the backend's existing test fixtures as a sentinel. + h.finish().saturating_add(1) +} + +/// Generate the `StreamingConfig` YAML for the ASAPQuery-backend from a +/// single-metric `CollectionPlan`. Produces a one-element `aggregations` +/// list — the backend's endpoint will merge this into its active config +/// (add on conflict, replace on same id). +/// +/// # Errors +/// +/// Returns an error if the plan is missing a window (the backend's +/// config parser rejects zero-window aggregations) or if YAML +/// serialization fails. +pub fn generate_streaming_config_yaml(metric: &str, plan: &CollectionPlan) -> Result { + let agg = &plan.agent_config; + let window_secs = agg + .window_duration + .map(|d: Duration| d.as_secs()) + .unwrap_or(0); + if window_secs == 0 { + anyhow::bail!( + "generate_streaming_config_yaml: plan for metric {metric} has \ + no window_duration; ASAPQuery-backend rejects zero-window aggregations" + ); + } + + // Parameters map: copy sketch-type-specific params (K for KLL, + // epsilon/delta for CountMin, etc.) into the string-keyed YAML map + // the backend expects. We serialize via serde_yaml to pick up + // SketchParams' own Serialize impl and then re-parse into a + // generic Mapping so we can embed it. + let params_yaml = serde_yaml::to_value(&agg.sketch_params) + .context("serialize SketchParams for ASAPQuery streaming config")?; + + let agg_type_str = map_sketch_type_to_agg_type(&agg.sketch_type); + + let aggregation = serde_yaml::Mapping::from_iter([ + ( + serde_yaml::Value::from("aggregationId"), + serde_yaml::Value::from(deterministic_agg_id(metric)), + ), + ( + serde_yaml::Value::from("aggregationType"), + serde_yaml::Value::from(agg_type_str), + ), + ( + serde_yaml::Value::from("aggregationSubType"), + serde_yaml::Value::from(""), + ), + ( + serde_yaml::Value::from("metric"), + serde_yaml::Value::from(metric), + ), + ( + serde_yaml::Value::from("labels"), + labels_mapping(&agg.aggregate_by), + ), + (serde_yaml::Value::from("parameters"), params_yaml), + ( + serde_yaml::Value::from("windowSize"), + serde_yaml::Value::from(window_secs), + ), + ( + serde_yaml::Value::from("windowType"), + serde_yaml::Value::from("tumbling"), + ), + ( + serde_yaml::Value::from("spatialFilter"), + serde_yaml::Value::from(normalize_spatial_filter(&agg.label_matchers)), + ), + ]); + + let top = serde_yaml::Mapping::from_iter([( + serde_yaml::Value::from("aggregations"), + serde_yaml::Value::Sequence(vec![serde_yaml::Value::Mapping(aggregation)]), + )]); + + serde_yaml::to_string(&serde_yaml::Value::Mapping(top)) + .context("serialize ASAPQuery streaming-config YAML") +} + +/// Map the controller's `SketchType` to the backend's +/// `AggregationType::Display` string. These strings must match what the +/// backend's `FromStr for AggregationType` in +/// `promql_utilities::query_logics::enums` accepts — hence the variant +/// names rather than the collector factory names (e.g. `"DatasketchesKLL"` +/// not `"KLL"`). +fn map_sketch_type_to_agg_type(t: &SketchType) -> &'static str { + match t { + SketchType::DDSketch => "DDSketch", + SketchType::KLL => "DatasketchesKLL", + SketchType::HLL => "HLL", + SketchType::CountSketch => "CountSketch", + SketchType::CountMinSketch => "CountMinSketch", + } +} + +/// Build the `labels` sub-mapping the backend expects. All three lists +/// exist because the backend's parser reads them separately for +/// key-value / spatial-rollup distinction; today the controller only +/// tracks `aggregate_by` (grouping), so rollup and aggregated stay +/// empty and are populated in a follow-up when the cost model starts +/// producing richer label metadata. +fn labels_mapping(aggregate_by: &[String]) -> serde_yaml::Value { + serde_yaml::Value::Mapping(serde_yaml::Mapping::from_iter([ + ( + serde_yaml::Value::from("grouping"), + serde_yaml::Value::Sequence( + aggregate_by + .iter() + .cloned() + .map(serde_yaml::Value::from) + .collect(), + ), + ), + ( + serde_yaml::Value::from("rollup"), + serde_yaml::Value::Sequence(vec![]), + ), + ( + serde_yaml::Value::from("aggregated"), + serde_yaml::Value::Sequence(vec![]), + ), + ])) +} + +/// Join the controller's `label_matchers` list (each shaped like +/// `"key=value"`) into a single comma-separated string the backend's +/// spatial-filter parser accepts. When the list is empty, returns an +/// empty string (the backend treats that as "no spatial filter"). +fn normalize_spatial_filter(label_matchers: &[String]) -> String { + label_matchers.join(",") +} + +// ─── Unused-warning suppression for types that are referenced only +// inside the unit tests below. This keeps the module self-contained +// even when the rest of the controller crate's cfg(test) surface grows. +#[allow(dead_code)] +fn _type_check(_: &AgentCollectorConfig) {} + +#[cfg(test)] +mod tests { + use super::*; + use crate::types::{ + BackendCollectorConfig, CollectionPlan, DeltaDecision, GatewayCollectorConfig, OutputMode, + ProcessorMode, SketchParams, TransmissionCostSummary, + }; + use std::time::Duration; + + fn dummy_plan(sketch_type: SketchType) -> CollectionPlan { + CollectionPlan { + agent_config: AgentCollectorConfig { + output_mode: OutputMode::Sketch, + sketch_type: sketch_type.clone(), + sketch_params: SketchParams::default(), + aggregate_by: vec!["host".to_string(), "service".to_string()], + label_matchers: vec!["env=prod".to_string()], + window_duration: Some(Duration::from_secs(30)), + mode: ProcessorMode::Window, + enable_self_monitoring: false, + transmit_sketch: true, + drop_original: true, + enable_series_id: false, + series_id_ttl_secs: 0, + delta_transmission: false, + delta_threshold: 0.0, + }, + gateway_config: GatewayCollectorConfig { passthrough: true }, + backend_config: BackendCollectorConfig { + merge_sketch_type: sketch_type, + group_by: vec![], + }, + precompute: vec![], + valid_until: chrono::Utc::now() + chrono::Duration::seconds(300), + delta_decision: DeltaDecision::default(), + transmission_cost_summary: TransmissionCostSummary::default(), + staged_plan: None, + } + } + + #[test] + fn deterministic_id_is_stable_across_calls() { + assert_eq!( + deterministic_agg_id("cpu_usage"), + deterministic_agg_id("cpu_usage") + ); + assert_ne!( + deterministic_agg_id("cpu_usage"), + deterministic_agg_id("mem_usage") + ); + // Id is biased away from 0 so test fixtures that use 0 as a + // sentinel don't accidentally collide. + assert_ne!(deterministic_agg_id("any"), 0); + } + + #[test] + fn yaml_round_trips_through_serde_yaml() { + let plan = dummy_plan(SketchType::DDSketch); + let yaml = generate_streaming_config_yaml("cpu_usage", &plan).expect("yaml ok"); + let parsed: serde_yaml::Value = serde_yaml::from_str(&yaml).expect("re-parse ok"); + + let aggs = parsed["aggregations"].as_sequence().expect("sequence"); + assert_eq!(aggs.len(), 1); + let a = &aggs[0]; + assert_eq!(a["aggregationType"], serde_yaml::Value::from("DDSketch")); + assert_eq!(a["metric"], serde_yaml::Value::from("cpu_usage")); + assert_eq!(a["windowSize"], serde_yaml::Value::from(30u64)); + assert_eq!(a["windowType"], serde_yaml::Value::from("tumbling")); + assert_eq!(a["spatialFilter"], serde_yaml::Value::from("env=prod")); + + let grouping = a["labels"]["grouping"].as_sequence().expect("grouping seq"); + let grouping: Vec<&str> = grouping.iter().filter_map(|v| v.as_str()).collect(); + assert_eq!(grouping, vec!["host", "service"]); + } + + #[test] + fn maps_all_sketch_types() { + assert_eq!( + map_sketch_type_to_agg_type(&SketchType::DDSketch), + "DDSketch" + ); + assert_eq!( + map_sketch_type_to_agg_type(&SketchType::KLL), + "DatasketchesKLL", + "KLL must map to the backend's enum variant name, not the factory name" + ); + assert_eq!(map_sketch_type_to_agg_type(&SketchType::HLL), "HLL"); + assert_eq!( + map_sketch_type_to_agg_type(&SketchType::CountSketch), + "CountSketch" + ); + assert_eq!( + map_sketch_type_to_agg_type(&SketchType::CountMinSketch), + "CountMinSketch" + ); + } + + #[test] + fn rejects_plan_without_window_duration() { + let mut plan = dummy_plan(SketchType::HLL); + plan.agent_config.window_duration = None; + let err = generate_streaming_config_yaml("m", &plan).expect_err("should error"); + assert!( + err.to_string().contains("window_duration"), + "error should mention window_duration: {err}" + ); + } + + #[test] + fn spatial_filter_joins_label_matchers() { + let mut plan = dummy_plan(SketchType::DDSketch); + plan.agent_config.label_matchers = + vec!["env=prod".to_string(), "region=us-east".to_string()]; + let yaml = generate_streaming_config_yaml("m", &plan).expect("ok"); + let parsed: serde_yaml::Value = serde_yaml::from_str(&yaml).unwrap(); + assert_eq!( + parsed["aggregations"][0]["spatialFilter"], + serde_yaml::Value::from("env=prod,region=us-east") + ); + } +} diff --git a/controller/src/config/mod.rs b/controller/src/config/mod.rs index 9660cb94..183096dd 100644 --- a/controller/src/config/mod.rs +++ b/controller/src/config/mod.rs @@ -1,9 +1,11 @@ pub mod agent; +pub mod asapquery_backend; pub mod backend; pub mod precompute; pub mod workloads; pub use agent::generate_agent_config; +pub use asapquery_backend::generate_streaming_config_yaml; pub use backend::{generate_backend_config, generate_backend_config_staged}; pub use precompute::{should_precompute, build_precompute_jobs, PrecomputeClient}; pub use workloads::WorkloadRegistry; diff --git a/controller/src/main.rs b/controller/src/main.rs index eaf948a5..53bbaa3d 100644 --- a/controller/src/main.rs +++ b/controller/src/main.rs @@ -1,5 +1,6 @@ mod algebra; mod analyzer; +mod backend_client; mod config; mod monitor; mod opamp; @@ -73,6 +74,7 @@ async fn main() { .and_then(|v| v.parse().ok()) .unwrap_or(60u64), ); + let backend_endpoint = std::env::var("CONTROLLER_BACKEND_ENDPOINT").ok(); // ── SP-5: Online EMA cost store ─────────────────────────────────────────── let online_store = init_online_store(); @@ -225,14 +227,31 @@ async fn main() { } // ── Replanner — closes the SP-8 feedback loop ───────────────────────────── - let replanner = Arc::new(Replanner::new( - Arc::clone(&planner), - Arc::clone(&plan_store), - Arc::clone(&workload_store), - Arc::clone(&opamp_srv), - Arc::clone(&scraper), - opamp_ep.clone(), - )); + let replanner = { + let mut r = Replanner::new( + Arc::clone(&planner), + Arc::clone(&plan_store), + Arc::clone(&workload_store), + Arc::clone(&opamp_srv), + Arc::clone(&scraper), + opamp_ep.clone(), + ); + if let Some(endpoint) = backend_endpoint.as_ref() { + info!( + endpoint = %endpoint, + "ASAPQuery-backend StreamingConfig push enabled" + ); + r = r.with_backend_client(Arc::new(backend_client::BackendClient::new( + endpoint.clone(), + ))); + } else { + info!( + "ASAPQuery-backend StreamingConfig push disabled \ + (set CONTROLLER_BACKEND_ENDPOINT= to enable)" + ); + } + Arc::new(r) + }; // Bind the late-binding cells so callbacks can reach the replanner and registry. *replanner_cell.write().await = Some(Arc::clone(&replanner)); *registry_cell.write().await = Some(Arc::clone(&workload_registry)); diff --git a/controller/src/replan.rs b/controller/src/replan.rs index 8700b0a0..4ba0dd7a 100644 --- a/controller/src/replan.rs +++ b/controller/src/replan.rs @@ -20,7 +20,11 @@ use std::time::Duration; use tokio::sync::RwLock; use tracing::{info, warn}; -use crate::config::{generate_agent_config, generate_backend_config, build_precompute_jobs}; +use crate::backend_client::{push_or_log, BackendClient}; +use crate::config::{ + build_precompute_jobs, generate_agent_config, generate_backend_config, + generate_streaming_config_yaml, +}; use crate::monitor::Scraper; use crate::opamp::{AgentRole, OpampServer, RemoteConfig}; use crate::planner::BaselinePlanner; @@ -43,6 +47,15 @@ pub struct Replanner { opamp: Arc, scraper: Arc, opamp_endpoint: String, + /// Optional client for pushing newly-generated `StreamingConfig` + /// YAML to the ASAPQuery-backend's `/api/v1/streaming-config` + /// endpoint. When present, every successful replan POSTs the new + /// plan to the backend in addition to the existing OpAMP pushes + /// to agent-role and backend-role collectors. Configured via the + /// `CONTROLLER_BACKEND_ENDPOINT` env var; defaults to `None` so + /// existing deployments that don't yet run ASAPQuery-backend + /// behave exactly as before. + backend_client: Option>, /// Maps agent_id → metric_name so violation callbacks can look up which /// metric a particular agent is serving. agent_to_metric: Arc>>, @@ -64,10 +77,21 @@ impl Replanner { opamp, scraper, opamp_endpoint: opamp_endpoint.into(), + backend_client: None, agent_to_metric: Arc::new(RwLock::new(HashMap::new())), } } + /// Attach a [`BackendClient`] so every replan also pushes the new + /// `StreamingConfig` YAML to the ASAPQuery-backend via HTTP. + /// Builder-style — call during controller startup in `main.rs`. + /// Without this call, replans continue to push only via OpAMP and + /// the ASAPQuery-backend (if running) keeps its startup config. + pub fn with_backend_client(mut self, client: Arc) -> Self { + self.backend_client = Some(client); + self + } + // ── Agent registry ──────────────────────────────────────────────────────── /// Record that `agent_id` is serving `metric`. Called from `handle_plan` @@ -157,6 +181,27 @@ impl Replanner { ).await; } + // Push the ASAPQuery-backend StreamingConfig YAML via HTTP if a + // backend client is configured. This is the producer side of the + // ASAPQuery PR E hot-reload contract: the backend receives the + // new plan on its /api/v1/streaming-config endpoint and makes it + // visible to the next query without restarting. + if let Some(backend_client) = self.backend_client.as_ref() { + match generate_streaming_config_yaml(metric, &plan) { + Ok(yaml) => { + push_or_log(backend_client, metric, yaml).await; + } + Err(e) => { + warn!( + metric, + error = %e, + "failed to build ASAPQuery streaming-config YAML — \ + skipping backend HTTP push for this replan cycle" + ); + } + } + } + // Update scraper endpoint sketch types for correct EMA attribution. let sketch_type = plan.agent_config.sketch_type; for agent_id in self.opamp.connected_agents().await {