From 462f80e798447dd5bc2ce0b38a80d59ce9740d1d Mon Sep 17 00:00:00 2001 From: Zeying Zhu Date: Wed, 15 Apr 2026 15:37:27 -0400 Subject: [PATCH] =?UTF-8?q?feat:=20controller=20BackendClient=20=E2=80=94?= =?UTF-8?q?=20push=20StreamingConfig=20to=20ASAPQuery?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes the producer side of the ASAPQuery PR E hot-reload contract. After a replan, the controller now (optionally) POSTs the new plan as a StreamingConfig YAML to ASAPQuery-backend's /api/v1/streaming-config endpoint, paralleling its existing OpAMP push to agent-role and backend-role collectors. ## Why ASAPQuery PRs #10 and #12 landed the endpoint and made it observable at query time, but the backend's active StreamingConfig only changes when something POSTs to it. Manual curl was the only producer until now. With this PR the controller closes the loop: query hits SimpleEngine miss → ASAPQuery PR #11 fires fire-and-forget POST to controller → this controller runs the planner, generates a new plan → THIS PR pushes the plan to the backend's streaming-config endpoint → next query re-snapshots and finds a match (via PR #12 phase 2) The data path (DataCollector agent → OTLP sketch → backend ingest) was already covered by the existing OpAMP push. This PR covers the control path back to the query-side state the backend holds. ## What's new ### `controller::config::asapquery_backend` Converts a `CollectionPlan` + metric name into the YAML shape `StreamingConfig::from_yaml_data` consumes (aggregations: [{aggregationId, aggregationType, metric, labels, parameters, windowSize, windowType, spatialFilter}]). Key details: * Maps `SketchType` → backend `AggregationType::Display` string (notably KLL → "DatasketchesKLL", NOT the factory string "KLL") * Derives a deterministic `u64` `aggregationId` from the metric name so repeat pushes for the same metric update (not duplicate) the backend's agg map * Rejects zero-window plans (the backend parser does too) * Joins the controller's `label_matchers` list into the backend's comma-separated `spatialFilter` string Separate from existing `config::backend` which emits OTel YAML (for a backend OTel collector running merge processors). The two consumers are different services consuming different formats. ### `controller::backend_client` Thin HTTP client wrapping reqwest::Client. Methods: * `BackendClient::new(endpoint)` — 5-second timeout (symmetric with ASAPQuery's HttpControllerClient in PR #11) * `BackendClient::push_streaming_config(yaml) -> Result<()>` — POSTs with content-type application/x-yaml, maps non-2xx to Err * `push_or_log(client, metric, yaml)` — fire-and-forget helper used by the replanner, never propagates errors (next replan cycle retries) ### `Replanner::with_backend_client(client)` builder New optional field. When set, `replan_metric()` calls `generate_streaming_config_yaml(metric, &plan)` after the OpAMP pushes and fire-and-forgets the result via the client. Without the builder call, replans behave exactly as before — existing deployments that don't yet run ASAPQuery-backend are unaffected. ### `main.rs` wiring New env var `CONTROLLER_BACKEND_ENDPOINT`. When set, constructs a `BackendClient` and attaches it to the Replanner via the new builder. Logs whether the feature is enabled at startup. ## Tests ### `config::asapquery_backend` (5 new) * `deterministic_id_is_stable_across_calls` — same metric → same id, different metrics → different ids, never returns 0 * `yaml_round_trips_through_serde_yaml` — generate, re-parse, assert every field (metric, window_size, window_type, spatial filter, grouping labels) * `maps_all_sketch_types` — pins the SketchType → AggregationType name mapping for all 5 sketch types, including the KLL → DatasketchesKLL gotcha * `rejects_plan_without_window_duration` — zero-window plan is rejected with a clear error * `spatial_filter_joins_label_matchers` — multiple matchers are joined with commas ### `backend_client` (3 new) * `success_path_round_trips_yaml` — spawns a local axum mock that records POST bodies, asserts the client posts the exact YAML and the server receives it * `non_2xx_status_is_reported_as_error` — mock returns 500, client maps to formatted Err containing the status * `push_or_log_swallows_errors` — points at an unreachable port, verifies the fire-and-forget helper does not propagate the failure (replan must never abort on backend unavailability) ## Validation * cargo check --all-targets: clean * cargo test --bin controller: 353 passed (up 8 from main) * cargo fmt --check on new files only: clean (pre-existing fmt drift in unrelated files not touched by this PR) ## Follow-ups * **Rate-limit / dedupe** — the controller may push the same plan repeatedly across rapid replans. The backend tolerates this idempotently via deterministic agg_ids, but a "push only on hash change" filter would cut unnecessary HTTP traffic. * **Retry with backoff** — fire-and-forget is fine for the common case. A transient backend outage currently drops one replan cycle; a small retry queue would recover it. * **Merge semantics on the backend side** — today the backend's POST handler REPLACES the entire StreamingConfig. When the controller pushes for metric A, any aggregations the backend had for metric B are wiped. A follow-up should either (a) restrict the controller to push full snapshots, or (b) extend the backend endpoint to support partial updates keyed by metric / agg_id. * **Labels.rollup / aggregated** — today we leave those empty because the controller doesn't track rollup dimensions separately. When the cost model starts producing richer label metadata, wire it through. Co-Authored-By: Claude Opus 4.6 (1M context) --- controller/src/backend_client.rs | 190 +++++++++++++ controller/src/config/asapquery_backend.rs | 299 +++++++++++++++++++++ controller/src/config/mod.rs | 2 + controller/src/main.rs | 35 ++- controller/src/replan.rs | 47 +++- 5 files changed, 564 insertions(+), 9 deletions(-) create mode 100644 controller/src/backend_client.rs create mode 100644 controller/src/config/asapquery_backend.rs 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 {