diff --git a/Cargo.lock b/Cargo.lock index b99d3be76..6c95f0177 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -366,6 +366,7 @@ dependencies = [ "asap_sketchlib", "prost", "serde", + "serde_json", "thiserror 1.0.69", ] diff --git a/control_plane/src/main.rs b/control_plane/src/main.rs index 6d9e0c1db..cede8fbe0 100644 --- a/control_plane/src/main.rs +++ b/control_plane/src/main.rs @@ -1,4 +1,3 @@ -use control_plane::accuracy; use control_plane::backend_client; use control_plane::emit; use control_plane::epsilon_alloc; @@ -25,6 +24,7 @@ use axum::{ routing::{get, post}, Json, Router, }; +use serde::{Deserialize, Serialize}; use serde_json::json; use std::collections::HashMap; use std::sync::Arc; @@ -552,6 +552,10 @@ async fn main() { let app = Router::new() .route("/api/v1/plan", post(handle_plan)) + .route( + "/api/v1/physical-plan/compile-and-publish", + post(handle_compile_and_publish_physical_plan), + ) .route("/api/v1/plan/auto", post(handle_plan_auto)) .route("/api/v1/plan/pareto", post(handle_pareto)) .route("/api/v1/plan/:metric", get(handle_get_plan)) @@ -573,6 +577,180 @@ async fn main() { axum::serve(listener, app).await.unwrap(); } +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct PhysicalPlanQueryRequest { + query_id: String, + query_string: String, + metric: String, + window_secs: u64, + #[serde(default)] + group_by: Vec, + accuracy: types_v2::AccuracyTarget, +} + +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct CompileAndPublishPhysicalPlanRequest { + queries: Vec, + collector_ids: Vec, + capability_snapshot_id: String, + #[serde(default)] + evidence: HashMap, + planner_revision: String, + max_evidence_age_ms: u64, + #[serde(default = "default_physical_plan_timeout_ms")] + apply_timeout_ms: u64, +} + +fn default_physical_plan_timeout_ms() -> u64 { + 10_000 +} + +#[derive(Debug, Serialize)] +struct CompileAndPublishPhysicalPlanResponse { + plan_id: u64, + generated_at_unix_ms: u64, + collector_ids: Vec, +} + +/// Compile one Planner IR decision into the matching backend/Collector views +/// and install them in dependency order. Unlike the legacy planning endpoint, +/// this MVP boundary is fail-closed: the Collector plan is never published +/// unless the backend accepted the exact matching BackendPlan first. +async fn handle_compile_and_publish_physical_plan( + State(st): State, + Json(request): Json, +) -> impl IntoResponse { + let (bundle, collector_ids, apply_timeout) = match compile_physical_plan_request(request) { + Ok(compiled) => compiled, + Err(response) => return response.into_response(), + }; + + let Some(backend) = st.backend_client.as_ref() else { + return ( + StatusCode::SERVICE_UNAVAILABLE, + "CONTROLLER_BACKEND_ENDPOINT is required for physical-plan publication".to_string(), + ) + .into_response(); + }; + if let Err(error) = st + .opamp + .ensure_collector_plan_targets(&bundle.collector_plans, apply_timeout) + .await + { + return ( + StatusCode::BAD_GATEWAY, + format!("collector physical-plan preflight failed: {error}"), + ) + .into_response(); + } + if let Err(error) = backend + .post_backend_plan_typed(bundle.backend_plan.encode_to_vec()) + .await + { + return ( + StatusCode::BAD_GATEWAY, + format!("backend rejected physical plan: {error}"), + ) + .into_response(); + } + if let Err(error) = st + .opamp + .publish_collector_plans(&bundle.collector_plans, apply_timeout) + .await + { + return ( + StatusCode::BAD_GATEWAY, + format!("collector physical-plan publication failed: {error}"), + ) + .into_response(); + } + + Json(CompileAndPublishPhysicalPlanResponse { + plan_id: bundle.envelope.plan_id, + generated_at_unix_ms: bundle.envelope.generated_at_unix_ms, + collector_ids, + }) + .into_response() +} + +// Keep Planner's Rc-backed rewrite DAG outside the async handler's future. +// Only the Send-safe compiled bundle crosses an await point. +fn compile_physical_plan_request( + request: CompileAndPublishPhysicalPlanRequest, +) -> Result< + ( + physical::compiler::CompiledPlanBundle, + Vec, + Duration, + ), + (StatusCode, String), +> { + if request.queries.is_empty() || request.collector_ids.is_empty() { + return Err(( + StatusCode::UNPROCESSABLE_ENTITY, + "queries and collector_ids must both be non-empty".to_string(), + )); + } + if request.max_evidence_age_ms == 0 || request.apply_timeout_ms == 0 { + return Err(( + StatusCode::UNPROCESSABLE_ENTITY, + "max_evidence_age_ms and apply_timeout_ms must be non-zero".to_string(), + )); + } + + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64; + let mut queries = Vec::with_capacity(request.queries.len()); + for query in request.queries { + if query.query_id.trim().is_empty() + || query.metric.trim().is_empty() + || query.window_secs == 0 + { + return Err(( + StatusCode::UNPROCESSABLE_ENTITY, + "query_id, metric, and window_secs must be non-empty/non-zero".to_string(), + )); + } + let expr = match parse_query_expr_canonical(&query.query_string, query.accuracy.clone()) { + Ok(expr) => expr, + Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string())), + }; + queries.push(physical::compiler::PlanningQuery { + query_id: query.query_id, + expr, + source: intent_algebra::Source::TimeSeries { + metric: query.metric, + }, + window_secs: query.window_secs, + group_by: query.group_by, + accuracy: query.accuracy, + }); + } + + let bundle = match physical::compiler::PhysicalCompiler.compile( + physical::compiler::PlanningRequest { + queries, + evidence: request.evidence, + planner_revision: request.planner_revision, + }, + physical::compiler::DeploymentEnvironment { + collector_ids: request.collector_ids.clone(), + capability_snapshot_id: request.capability_snapshot_id, + observed_at_unix_ms: now, + max_evidence_age_ms: request.max_evidence_age_ms, + }, + ) { + Ok(bundle) => bundle, + Err(error) => return Err((StatusCode::UNPROCESSABLE_ENTITY, error.to_string())), + }; + let apply_timeout = Duration::from_millis(request.apply_timeout_ms); + Ok((bundle, request.collector_ids, apply_timeout)) +} + // ── Handlers ────────────────────────────────────────────────────────────────── async fn handle_plan(State(st): State, Json(spec): Json) -> impl IntoResponse { @@ -2864,17 +3042,18 @@ mod api_tests { } /// Acceptance test: PR #339 (planner) ↔ PR #340 (emitter) stitch - /// produces the 5-sketch routing-connector wire shape when the - /// workload registry covers the six MVP §46 contract metrics. + /// produces the evidence-safe routing-connector wire shape. TopK is + /// intentionally absent here: the bootstrap path has no membership + /// evidence and the latest Planner contract fails it closed. /// /// Asserts: - /// - All 5 sketch processors loaded under `processors:`. + /// - All four evidence-safe sketch processors are loaded. /// - `routing` lives in `connectors:` (NOT `processors:`). /// - All 6 named pipelines emitted (raw_passthrough + 5 sketches). /// - Each metric routed to its expected pipeline via /// `name == "..."`. #[tokio::test] - async fn bootstrap_emits_5sketch_routing_for_six_contract_metrics() { + async fn bootstrap_omits_topk_without_membership_evidence() { let _env = EnvVarGuard::set(physical::stage_split::ENV_USE_TYPED_STAGE_SPLIT, "1"); let (_, app, tmp) = test_app_with_six_contract_metrics(); @@ -2888,13 +3067,14 @@ mod api_tests { let yaml = String::from_utf8(body.to_vec()).unwrap(); std::fs::remove_file(&tmp).ok(); - // ── Contract 1: all 5 sketch processors loaded ──────────────────── - for proc in ["ddsketch:", "KLL:", "HLL:", "countsketch:", "countmin:"] { + // ── Contract 1: all evidence-safe sketch processors loaded ──────── + for proc in ["ddsketch:", "KLL:", "HLL:", "countmin:"] { assert!( yaml.contains(proc), "missing top-level sketch processor `{proc}`\n{yaml}" ); } + assert!(!yaml.contains("countsketch:")); // ── Contract 2: routing in connectors, not processors ───────────── let connectors_idx = yaml @@ -2926,14 +3106,13 @@ mod api_tests { "metrics/ddsketch_path:", // http_latency_ms "metrics/kll_path:", // request_size_bytes "metrics/hll_path:", // unique_users_per_min - "metrics/countsketch_path:", // top_endpoint_qps "metrics/countminsketch_path:", // endpoint_request_freq ] { assert!(yaml.contains(pl), "missing pipeline `{pl}`\n{yaml}"); } // ── Contract 4: each sketched metric carries an OTTL condition ── - // The 5 sketched metrics must each have a `name == "..."` + // Each evidence-safe sketched metric must have a routing rule. // rule in the routing connector. // `http_requests_total` (raw) does NOT need a rule — it falls // through to the default `metrics/raw_passthrough` pipeline. @@ -2941,7 +3120,6 @@ mod api_tests { "http_latency_ms", "request_size_bytes", "unique_users_per_min", - "top_endpoint_qps", "endpoint_request_freq", ] { let needle = format!("name == \\\"{sketched}\\\""); @@ -2954,7 +3132,7 @@ mod api_tests { } } - // ── Stitching-gap regression: live mvp-workload.yaml binds all 5 sketches ── + // ── Bootstrap respects the latest Planner evidence gate ────────────── // // Reproduces the live demo gap (3 of 6 contract metrics silently dropped // because `WorkloadEntry` didn't carry `sketch_family_override` and @@ -3049,26 +3227,25 @@ mod api_tests { (state, router, tmp_path) } - /// Pinning regression: the routing table emitted by the bootstrap - /// endpoint must cover all 5 sketched contract metrics — DDSketch - /// (`http_requests_total_latency_ms`), KLL (`request_size_bytes`), - /// HLL (`unique_users_per_min`), CountSketch (`top_endpoint_qps`), - /// CountMinSketch (`endpoint_request_freq`). + /// The legacy bootstrap input has no TopK separation evidence, so the + /// latest Planner contract must omit CountSketch while retaining the four + /// independently executable materializations. /// /// Without the fix, this test fails with only 2 sketched routes /// (DDSketch + KLL); HLL / CountSketch / CountMinSketch silently drop. #[tokio::test] - async fn bootstrap_routing_table_covers_all_five_sketches_for_live_mvp_yaml() { + async fn bootstrap_routing_table_excludes_unevidenced_topk() { let _env = EnvVarGuard::set(physical::stage_split::ENV_USE_TYPED_STAGE_SPLIT, "1"); let (state, app, tmp) = test_app_with_live_mvp_workload_metrics(); - // ── Direct check: collect_metric_to_family produces 5 entries ──── + // TopK has no membership-separation evidence in this legacy + // bootstrap input, so only four materializations are executable. let map = emit::collect_metric_to_family(&state.workload_registry, &state.workload_store); assert_eq!( map.len(), - 5, - "metric_to_family should have 5 sketched entries (raw declines), got {map:?}", + 4, + "metric_to_family should have 4 evidence-safe entries, got {map:?}", ); // ASAPCollector#400: values are now SETs of families. For the // demo workload each metric is queried by exactly one capability, @@ -3077,7 +3254,6 @@ mod api_tests { ("http_requests_total_latency_ms", "{DDSketch}"), ("request_size_bytes", "{Kll}"), ("unique_users_per_min", "{Hll}"), - ("top_endpoint_qps", "{CountSketch}"), ("endpoint_request_freq", "{Cms}"), ] { let got = map @@ -3105,7 +3281,6 @@ mod api_tests { "http_requests_total_latency_ms", "request_size_bytes", "unique_users_per_min", - "top_endpoint_qps", "endpoint_request_freq", ] { let needle = format!("name == \\\"{sketched}\\\""); @@ -3116,6 +3291,7 @@ mod api_tests { "missing routing rule for `{sketched}`\n{yaml}" ); } + assert!(!yaml.contains("name == \"top_endpoint_qps\"")); } // ── Regression: archive tier covers all 5 sketched metrics ──────────────── @@ -3138,7 +3314,7 @@ mod api_tests { // against a mock backend, captures every body, and asserts the // final swap covers all 5 sketched metrics simultaneously. #[tokio::test] - async fn storage_routing_cumulative_push_covers_all_5_sketched_metrics() { + async fn storage_routing_cumulative_push_covers_all_planned_metrics() { // Activate the typed-stage-split path (the only path that // emits storage-routing JSON; the legacy path no-ops). let _env = EnvVarGuard::set(physical::stage_split::ENV_USE_TYPED_STAGE_SPLIT, "1"); @@ -3182,14 +3358,13 @@ mod api_tests { .route("/api/v1/plan", axum::routing::post(handle_plan)) .with_state(state.clone()); - // The 5 sketched contract metrics from MVP §46. Each gets a + // The four evidence-safe sketched contract metrics. Each gets a // separate POST /api/v1/plan, mirroring the demo's // per-workload plan-emit cycle. let sketched = [ "http_requests_total_latency_ms", // DDSketch "request_size_bytes", // KLL "unique_users_per_min", // HLL - "top_endpoint_qps", // CountSketch "endpoint_request_freq", // CountMinSketch ]; @@ -3210,7 +3385,7 @@ mod api_tests { } // Drain the mock sink: every plan-emit must have produced - // exactly one body (5 plans → 5 bodies). + // exactly one body. let bodies = sink.lock().unwrap().clone(); assert_eq!( bodies.len(), @@ -3221,7 +3396,7 @@ mod api_tests { // The LAST captured body is the one the backend will leave // installed (the swap is destructive — last write wins). It - // MUST list ALL 5 sketched metrics, otherwise the swap would + // MUST list all planned metrics, otherwise the swap would // erase the routing entries for the metrics planned earlier // in the sequence and the backend would default them to // `sketch_store` → archive_miss for those metrics' archive @@ -3342,7 +3517,7 @@ mod api_tests { .route("/api/v1/plan", axum::routing::post(handle_plan)) .with_state(state.clone()); - // The 5 sketched contract metrics from MVP §46 — the SAME + // The evidence-safe sketched contract metrics — the same // set the sibling `storage_routing_cumulative_push_...` test // exercises. Each gets a separate `POST /api/v1/plan` with // the metric-name → classified sketch family from @@ -3355,7 +3530,6 @@ mod api_tests { "http_requests_total_latency_ms", "request_size_bytes", "unique_users_per_min", - "top_endpoint_qps", "endpoint_request_freq", ]; @@ -3376,7 +3550,7 @@ mod api_tests { } // Drain the mock sink: every plan-emit must have produced - // exactly one streaming-config body (5 plans → 5 bodies). + // exactly one streaming-config body. let bodies = sink.lock().unwrap().clone(); assert_eq!( bodies.len(), @@ -3387,7 +3561,7 @@ mod api_tests { // The LAST body is the one the data plane's swap installs // (the swap is destructive — last write wins). It MUST list - // aggregations for ALL 5 sketched metrics, otherwise the + // aggregations for all planned metrics, otherwise the // swap erases the earlier metrics' rows and queries against // them fail with `…capability not satisfied` — the // streaming-config analogue of the storage-routing diff --git a/control_plane/src/opamp/mod.rs b/control_plane/src/opamp/mod.rs index 75b0f8ec4..0459bf7c5 100644 --- a/control_plane/src/opamp/mod.rs +++ b/control_plane/src/opamp/mod.rs @@ -9,8 +9,9 @@ /// /// The server pushes `RemoteConfig` as an OpAMP `ServerToAgent.remote_config` /// message (protobuf binary frame) and receives `AgentToServer` status reports. -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::sync::Arc; +use std::time::Duration; use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade}; use axum::extract::State; @@ -20,7 +21,7 @@ use futures_util::sink::SinkExt; use futures_util::stream::StreamExt; use prost::Message as ProstMessage; use serde::{Deserialize, Serialize}; -use tokio::sync::{mpsc, RwLock}; +use tokio::sync::{mpsc, Notify, RwLock}; use tracing::{info, warn}; /// Generated OpAMP protobuf types (from proto/opamp.proto). @@ -46,6 +47,50 @@ pub struct AgentStatus { pub error: Option, } +pub const COLLECTOR_PLAN_CAPABILITY: &str = "io.projectasap.collector-plan.v1"; +pub const COLLECTOR_PLAN_MESSAGE: &str = "collector_plan"; +pub const PLAN_STATUS_MESSAGE: &str = "plan_status"; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +pub enum CollectorPlanStatusKind { + Applied, + Failed, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +pub struct CollectorPlanStatus { + pub plan_id: u64, + pub status: CollectorPlanStatusKind, + #[serde(default)] + pub error: Option, +} + +#[derive(Debug, thiserror::Error)] +pub enum CollectorPlanPublishError { + #[error("collector {collector_id} did not advertise {COLLECTOR_PLAN_CAPABILITY}")] + CapabilityTimeout { collector_id: String }, + #[error("collector {collector_id} disconnected before plan publication")] + Disconnected { collector_id: String }, + #[error("compiled bundle contains duplicate collector target {collector_id}")] + DuplicateTarget { collector_id: String }, + #[error("collector {collector_id} plan serialization failed: {source}")] + Serialize { + collector_id: String, + #[source] + source: serde_json::Error, + }, + #[error("collector {collector_id} did not report plan {plan_id} before timeout")] + StatusTimeout { collector_id: String, plan_id: u64 }, + #[error("collector {collector_id} rejected plan {plan_id}: {error}")] + Rejected { + collector_id: String, + plan_id: u64, + error: String, + }, +} + // ── Role ────────────────────────────────────────────────────────────────────── /// Role of a connected OTel collector. @@ -79,13 +124,28 @@ impl AgentRole { // ── Server ──────────────────────────────────────────────────────────────────── -type AgentMap = HashMap, AgentRole)>; +#[derive(Debug, Clone)] +enum OutboundMessage { + RemoteConfig(RemoteConfig), + CollectorPlan(Vec), +} + +struct AgentConnection { + tx: mpsc::Sender, + role: AgentRole, + custom_capabilities: HashSet, +} + +type AgentMap = HashMap; +type PlanStatusMap = HashMap<(String, u64), CollectorPlanStatus>; pub type OnConnectFn = Arc; pub type OnDisconnectFn = Arc; pub struct OpampServer { agents: Arc>, + plan_statuses: Arc>, + state_changed: Arc, on_connect: Option, on_disconnect: Option, } @@ -94,6 +154,8 @@ impl Default for OpampServer { fn default() -> Self { Self { agents: Arc::new(RwLock::new(HashMap::new())), + plan_statuses: Arc::new(RwLock::new(HashMap::new())), + state_changed: Arc::new(Notify::new()), on_connect: None, on_disconnect: None, } @@ -104,6 +166,8 @@ impl Clone for OpampServer { fn clone(&self) -> Self { Self { agents: Arc::clone(&self.agents), + plan_statuses: Arc::clone(&self.plan_statuses), + state_changed: Arc::clone(&self.state_changed), on_connect: self.on_connect.clone(), on_disconnect: self.on_disconnect.clone(), } @@ -160,10 +224,15 @@ impl OpampServer { /// Pushes a config to a specific agent. Returns false if not connected. pub async fn push(&self, agent_id: &str, cfg: RemoteConfig) -> bool { - if let Some((tx, _)) = self.agents.read().await.get(agent_id) { - tx.send(cfg).await.is_ok() - } else { - false + let tx = self + .agents + .read() + .await + .get(agent_id) + .map(|connection| connection.tx.clone()); + match tx { + Some(tx) => tx.send(OutboundMessage::RemoteConfig(cfg)).await.is_ok(), + None => false, } } @@ -182,7 +251,7 @@ impl OpampServer { .read() .await .iter() - .filter(|(_, (_, r))| *r == role) + .filter(|(_, connection)| connection.role == role) .map(|(id, _)| id.clone()) .collect(); for id in ids { @@ -201,9 +270,150 @@ impl OpampServer { .read() .await .iter() - .map(|(id, (_, role))| (id.clone(), role.clone())) + .map(|(id, connection)| (id.clone(), connection.role.clone())) .collect() } + + /// Publish all per-target physical plans and require an exact semantic + /// APPLIED report from every Collector. Config hashes are deliberately not + /// accepted as plan activation evidence. + pub async fn publish_collector_plans( + &self, + plans: &[crate::physical::compiler::CollectorPlan], + timeout: Duration, + ) -> Result, CollectorPlanPublishError> { + self.ensure_collector_plan_targets(plans, timeout).await?; + + for plan in plans { + let body = serde_json::to_vec(plan).map_err(|source| { + CollectorPlanPublishError::Serialize { + collector_id: plan.collector_id.clone(), + source, + } + })?; + self.plan_statuses + .write() + .await + .remove(&(plan.collector_id.clone(), plan.envelope.plan_id)); + let sent = self.send_collector_plan(&plan.collector_id, body).await; + if !sent { + return Err(CollectorPlanPublishError::Disconnected { + collector_id: plan.collector_id.clone(), + }); + } + } + + let mut reports = Vec::with_capacity(plans.len()); + for plan in plans { + let report = self + .wait_for_plan_status(&plan.collector_id, plan.envelope.plan_id, timeout) + .await + .ok_or_else(|| CollectorPlanPublishError::StatusTimeout { + collector_id: plan.collector_id.clone(), + plan_id: plan.envelope.plan_id, + })?; + if report.status != CollectorPlanStatusKind::Applied { + return Err(CollectorPlanPublishError::Rejected { + collector_id: plan.collector_id.clone(), + plan_id: report.plan_id, + error: report.error.clone().unwrap_or_else(|| "unspecified".into()), + }); + } + reports.push(report); + } + Ok(reports) + } + + /// Validate target uniqueness and wait for every target to advertise the + /// physical-plan capability. The orchestrator calls this before changing + /// backend state, then publication repeats it to close disconnect races. + pub async fn ensure_collector_plan_targets( + &self, + plans: &[crate::physical::compiler::CollectorPlan], + timeout: Duration, + ) -> Result<(), CollectorPlanPublishError> { + let mut targets = HashSet::with_capacity(plans.len()); + for plan in plans { + if !targets.insert(plan.collector_id.as_str()) { + return Err(CollectorPlanPublishError::DuplicateTarget { + collector_id: plan.collector_id.clone(), + }); + } + } + for plan in plans { + if !self + .wait_for_capability(&plan.collector_id, COLLECTOR_PLAN_CAPABILITY, timeout) + .await + { + return Err(CollectorPlanPublishError::CapabilityTimeout { + collector_id: plan.collector_id.clone(), + }); + } + } + + Ok(()) + } + + async fn send_collector_plan(&self, agent_id: &str, body: Vec) -> bool { + let agents = self.agents.read().await; + let Some(connection) = agents.get(agent_id) else { + return false; + }; + if !connection + .custom_capabilities + .contains(COLLECTOR_PLAN_CAPABILITY) + { + return false; + } + let tx = connection.tx.clone(); + drop(agents); + tx.send(OutboundMessage::CollectorPlan(body)).await.is_ok() + } + + async fn wait_for_capability( + &self, + agent_id: &str, + capability: &str, + timeout: Duration, + ) -> bool { + let deadline = tokio::time::Instant::now() + timeout; + loop { + let changed = self.state_changed.notified(); + if self + .agents + .read() + .await + .get(agent_id) + .is_some_and(|agent| agent.custom_capabilities.contains(capability)) + { + return true; + } + let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() || tokio::time::timeout(remaining, changed).await.is_err() { + return false; + } + } + } + + async fn wait_for_plan_status( + &self, + agent_id: &str, + plan_id: u64, + timeout: Duration, + ) -> Option { + let key = (agent_id.to_string(), plan_id); + let deadline = tokio::time::Instant::now() + timeout; + loop { + let changed = self.state_changed.notified(); + if let Some(status) = self.plan_statuses.read().await.get(&key).cloned() { + return Some(status); + } + let remaining = deadline.saturating_duration_since(tokio::time::Instant::now()); + if remaining.is_zero() || tokio::time::timeout(remaining, changed).await.is_err() { + return None; + } + } + } } async fn handle_socket( @@ -212,11 +422,16 @@ async fn handle_socket( role: AgentRole, srv: Arc, ) { - let (tx, mut rx) = mpsc::channel::(16); - srv.agents - .write() - .await - .insert(agent_id.clone(), (tx, role.clone())); + let (tx, mut rx) = mpsc::channel::(16); + srv.agents.write().await.insert( + agent_id.clone(), + AgentConnection { + tx, + role: role.clone(), + custom_capabilities: HashSet::new(), + }, + ); + srv.state_changed.notify_waiters(); info!(agent = %agent_id, ?role, "agent connected"); if let Some(cb) = &srv.on_connect { @@ -235,9 +450,17 @@ async fn handle_socket( // conformance so the agent never has to take the legacy path. let writer_id = agent_id.clone(); let write_task = tokio::spawn(async move { - while let Some(cfg) = rx.recv().await { - // Build standard OpAMP ServerToAgent with RemoteConfig. - let server_to_agent = encode_remote_config(&cfg); + while let Some(message) = rx.recv().await { + let (server_to_agent, description) = match message { + OutboundMessage::RemoteConfig(cfg) => ( + encode_remote_config(&cfg), + format!("remote config hash={}", cfg.config_hash), + ), + OutboundMessage::CollectorPlan(body) => ( + encode_collector_plan(body), + "typed CollectorPlan".to_string(), + ), + }; let mut payload = Vec::new(); if server_to_agent.encode(&mut payload).is_err() { warn!(agent = %writer_id, "failed to encode OpAMP protobuf"); @@ -251,7 +474,7 @@ async fn handle_socket( if ws_tx.send(Message::Binary(buf.into())).await.is_err() { break; } - info!(agent = %writer_id, hash = %cfg.config_hash, "config pushed (OpAMP protobuf)"); + info!(agent = %writer_id, message = %description, "message pushed (OpAMP protobuf)"); } }); @@ -290,6 +513,33 @@ async fn handle_socket( match opamp_proto::AgentToServer::decode(payload) { Ok(ats) => { info!(agent = %agent_id, "received AgentToServer (OpAMP protobuf)"); + if let Some(capabilities) = &ats.custom_capabilities { + if let Some(connection) = srv.agents.write().await.get_mut(&agent_id) { + connection.custom_capabilities = + capabilities.capabilities.iter().cloned().collect(); + } + srv.state_changed.notify_waiters(); + } + if let Some(message) = &ats.custom_message { + if message.capability == COLLECTOR_PLAN_CAPABILITY + && message.r#type == PLAN_STATUS_MESSAGE + { + match serde_json::from_slice::(&message.data) { + Ok(status) => { + srv.plan_statuses + .write() + .await + .insert((agent_id.clone(), status.plan_id), status); + srv.state_changed.notify_waiters(); + } + Err(error) => warn!( + agent = %agent_id, + %error, + "invalid typed CollectorPlan status" + ), + } + } + } // Log effective config if reported. Triggered by // the `ReportFullState` flag on our outgoing // ServerToAgent (see `encode_remote_config`). @@ -377,6 +627,7 @@ async fn handle_socket( write_task.abort(); srv.agents.write().await.remove(&agent_id); + srv.state_changed.notify_waiters(); info!(agent = %agent_id, "agent disconnected"); if let Some(cb) = &srv.on_disconnect { @@ -453,13 +704,27 @@ fn encode_remote_config(cfg: &RemoteConfig) -> opamp_proto::ServerToAgent { } } +fn encode_collector_plan(body: Vec) -> opamp_proto::ServerToAgent { + opamp_proto::ServerToAgent { + capabilities: 0x01, + custom_capabilities: Some(opamp_proto::CustomCapabilities { + capabilities: vec![COLLECTOR_PLAN_CAPABILITY.into()], + }), + custom_message: Some(opamp_proto::CustomMessage { + capability: COLLECTOR_PLAN_CAPABILITY.into(), + r#type: COLLECTOR_PLAN_MESSAGE.into(), + data: body, + }), + ..Default::default() + } +} + // ── Tests ───────────────────────────────────────────────────────────────────── #[cfg(test)] mod tests { use super::*; use axum::{routing::get, Router}; - use tokio::net::TcpListener; #[tokio::test] async fn no_agents_push_returns_false() { @@ -715,4 +980,140 @@ mod tests { "backend-role client must not receive agent-role push" ); } + + fn test_collector_plan( + collector_id: &str, + plan_id: u64, + ) -> crate::physical::compiler::CollectorPlan { + crate::physical::compiler::CollectorPlan { + collector_id: collector_id.into(), + envelope: crate::physical::compiler::PlanEnvelope { + plan_id, + generated_at_unix_ms: 1, + planner_revision: crate::physical::compiler::PLANNER_REVISION.into(), + capability_snapshot_id: "caps-1".into(), + }, + materializations: vec![crate::physical::compiler::CollectorMaterialization { + query_id: "q".into(), + metric: "requests".into(), + algorithm: "hll".into(), + parameters: serde_json::json!({"precision": 14}), + group_by: vec!["service".into()], + window_secs: 60, + evidence_source: None, + }], + } + } + + async fn send_agent_message( + ws: &mut tokio_tungstenite::WebSocketStream< + tokio_tungstenite::MaybeTlsStream, + >, + message: opamp_proto::AgentToServer, + ) { + use futures_util::SinkExt; + let mut protobuf = Vec::new(); + message.encode(&mut protobuf).unwrap(); + let mut frame = vec![0]; + frame.extend(protobuf); + ws.send(tokio_tungstenite::tungstenite::Message::Binary(frame)) + .await + .unwrap(); + } + + #[tokio::test] + async fn typed_plan_publication_waits_for_capability_and_exact_applied_status() { + let (srv, addr) = start_server().await; + let mut agent_ws = connect_ws_client(addr, "edge-a", "agent").await; + + send_agent_message( + &mut agent_ws, + opamp_proto::AgentToServer { + custom_capabilities: Some(opamp_proto::CustomCapabilities { + capabilities: vec![COLLECTOR_PLAN_CAPABILITY.into()], + }), + ..Default::default() + }, + ) + .await; + + let publisher = Arc::clone(&srv); + let plan = test_collector_plan("edge-a", 42); + let publish = tokio::spawn(async move { + publisher + .publish_collector_plans(&[plan], Duration::from_secs(2)) + .await + }); + + let frame = tokio::time::timeout(Duration::from_secs(2), agent_ws.next()) + .await + .expect("typed plan timed out") + .unwrap() + .unwrap(); + let message = decode_server_to_agent_frame(&frame.into_data()); + let custom = message.custom_message.expect("custom message"); + assert_eq!(custom.capability, COLLECTOR_PLAN_CAPABILITY); + assert_eq!(custom.r#type, COLLECTOR_PLAN_MESSAGE); + let wire_plan: crate::physical::compiler::CollectorPlan = + serde_json::from_slice(&custom.data).unwrap(); + assert_eq!(wire_plan.collector_id, "edge-a"); + assert_eq!(wire_plan.envelope.plan_id, 42); + + let status = serde_json::to_vec(&CollectorPlanStatus { + plan_id: 42, + status: CollectorPlanStatusKind::Applied, + error: None, + }) + .unwrap(); + send_agent_message( + &mut agent_ws, + opamp_proto::AgentToServer { + custom_message: Some(opamp_proto::CustomMessage { + capability: COLLECTOR_PLAN_CAPABILITY.into(), + r#type: PLAN_STATUS_MESSAGE.into(), + data: status, + }), + ..Default::default() + }, + ) + .await; + + let reports = publish.await.unwrap().unwrap(); + assert_eq!(reports.len(), 1); + assert_eq!(reports[0].plan_id, 42); + assert_eq!(reports[0].status, CollectorPlanStatusKind::Applied); + } + + #[tokio::test] + async fn typed_plan_publication_fails_without_advertised_capability() { + let (srv, addr) = start_server().await; + let _agent_ws = connect_ws_client(addr, "edge-a", "agent").await; + let error = srv + .publish_collector_plans( + &[test_collector_plan("edge-a", 42)], + Duration::from_millis(30), + ) + .await + .unwrap_err(); + assert!(matches!( + error, + CollectorPlanPublishError::CapabilityTimeout { collector_id } + if collector_id == "edge-a" + )); + } + + #[tokio::test] + async fn typed_plan_publication_rejects_duplicate_targets_before_sending() { + let srv = OpampServer::new(); + let plan = test_collector_plan("edge-a", 42); + let error = srv + .publish_collector_plans(&[plan.clone(), plan], Duration::from_millis(1)) + .await + .unwrap_err(); + assert!(matches!( + error, + CollectorPlanPublishError::DuplicateTarget { collector_id } + if collector_id == "edge-a" + )); + } } diff --git a/control_plane/src/physical/compiler.rs b/control_plane/src/physical/compiler.rs new file mode 100644 index 000000000..292957a71 --- /dev/null +++ b/control_plane/src/physical/compiler.rs @@ -0,0 +1,460 @@ +//! Backend-owned physical compilation over ASAPPlanner's selected post-ASAP IR. +//! +//! Planner owns semantic alternatives and guarantees. This module owns the +//! deployment decision: evidence freshness, target capabilities, windows, the +//! Collector execution projection, and the matching BackendPlan. + +use std::collections::HashMap; + +use asap_aware_mapping::{ + AccuracyEvidenceProvider, DefaultAccuracyModel, EqualSplitAllocator, PropagationStats, +}; +use planner_types::post_asap::{ + CompositionOperator, SketchQuery, SummaryExpr, SummaryFamilyType, SummaryNode, +}; +use planner_types::pre_asap::QueryExpr; +use serde::{Deserialize, Serialize}; +use serde_json::{json, Value}; +use thiserror::Error; + +use crate::backend_plan::{self, BackendPlan}; +use crate::emit::monitor::MonitorIntent; +use crate::intent_algebra::Source; +use crate::physical::colored_dag::emitter::{ + AggregationInput, BackendAggregation, BackendReadout, BackendStageConfig, +}; +use crate::sketch_algebra::cost_model::ControlPlaneCostModel; +use crate::types_v2::AccuracyTarget; + +pub const PLANNER_REVISION: &str = "3afcba68f4e8397fb81e2be988f47120f63f7a39"; + +#[derive(Debug, Clone)] +pub struct PlanningQuery { + pub query_id: String, + pub expr: QueryExpr, + pub source: Source, + pub window_secs: u64, + /// Label names are deployment metadata because Planner's canonical IR + /// currently carries positional column IDs at this boundary. + pub group_by: Vec, + pub accuracy: AccuracyTarget, +} + +#[derive(Debug, Clone, Default)] +pub struct PlanningRequest { + pub queries: Vec, + pub evidence: HashMap, + pub planner_revision: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct TopKMembershipEvidence { + pub selected_lower_bound: f64, + pub excluded_upper_bound: f64, + pub interval_failure_probability: f64, + pub observed_at_unix_ms: u64, + pub source: String, +} + +#[derive(Debug, Clone)] +pub struct DeploymentEnvironment { + pub collector_ids: Vec, + pub capability_snapshot_id: String, + pub observed_at_unix_ms: u64, + pub max_evidence_age_ms: u64, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct PlanEnvelope { + pub plan_id: u64, + pub generated_at_unix_ms: u64, + pub planner_revision: String, + pub capability_snapshot_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct CollectorMaterialization { + pub query_id: String, + pub metric: String, + pub algorithm: String, + pub parameters: Value, + pub group_by: Vec, + pub window_secs: u64, + pub evidence_source: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +pub struct CollectorPlan { + pub collector_id: String, + pub envelope: PlanEnvelope, + pub materializations: Vec, +} + +#[derive(Debug, Clone)] +pub struct CompiledPlanBundle { + pub envelope: PlanEnvelope, + pub collector_plans: Vec, + pub backend_plan: BackendPlan, +} + +#[derive(Debug, Error)] +pub enum CompileError { + #[error("planner revision mismatch: request={request}, compiler={compiler}")] + PlannerRevision { + request: String, + compiler: &'static str, + }, + #[error("query {query_id}: {reason}")] + Query { query_id: String, reason: String }, + #[error("query {query_id}: TopK evidence is stale or invalid: {reason}")] + InvalidEvidence { query_id: String, reason: String }, + #[error("failed to construct BackendPlan: {0}")] + BackendPlan(#[from] anyhow::Error), +} + +struct QueryEvidence<'a>(Option<&'a TopKMembershipEvidence>); + +impl AccuracyEvidenceProvider for QueryEvidence<'_> { + fn propagation_stats( + &self, + op: &CompositionOperator, + _family: &SummaryFamilyType, + _query: Option<&SketchQuery>, + ) -> PropagationStats { + match (op, self.0) { + (CompositionOperator::TopKSelection, Some(e)) => PropagationStats { + topk_selected_lower_bound: Some(e.selected_lower_bound), + topk_excluded_upper_bound: Some(e.excluded_upper_bound), + topk_interval_failure_probability: Some(e.interval_failure_probability), + ..Default::default() + }, + _ => PropagationStats::default(), + } + } +} + +#[derive(Debug, Default)] +pub struct PhysicalCompiler; + +impl PhysicalCompiler { + pub fn compile( + &self, + request: PlanningRequest, + environment: DeploymentEnvironment, + ) -> Result { + if request.planner_revision != PLANNER_REVISION { + return Err(CompileError::PlannerRevision { + request: request.planner_revision, + compiler: PLANNER_REVISION, + }); + } + + let mut aggregations = Vec::with_capacity(request.queries.len()); + let mut readouts = Vec::with_capacity(request.queries.len()); + let mut collector_materializations = Vec::with_capacity(request.queries.len()); + + for query in &request.queries { + let evidence = request.evidence.get(&query.query_id); + if let Some(e) = evidence { + validate_evidence(&query.query_id, e, &environment)?; + } + let model = ControlPlaneCostModel::new(query.accuracy.clone()); + let node = crate::planner_selection::select_summary_with_evidence( + &query.expr, + &model, + &DefaultAccuracyModel, + &EqualSplitAllocator, + &QueryEvidence(evidence), + ) + .map_err(|error| CompileError::Query { + query_id: query.query_id.clone(), + reason: error.to_string(), + })?; + let selected = extract_selected(&node).ok_or_else(|| CompileError::Query { + query_id: query.query_id.clone(), + reason: "selected plan has no executable sketch materialization/readout".into(), + })?; + let metric = match &query.source { + Source::TimeSeries { metric } => metric.clone(), + Source::Table { .. } => { + return Err(CompileError::Query { + query_id: query.query_id.clone(), + reason: "MVP physical compiler supports time-series sources only".into(), + }) + } + }; + let aggregation_id = format!("{}:{}", query.query_id, metric); + let kind = asap_types::SummaryKind::from(selected.kind.clone()); + let params = asap_types::SummaryParams::from(selected.params.clone()); + aggregations.push(BackendAggregation { + aggregation_id: aggregation_id.clone(), + metric_name: metric.clone(), + sketch_kind: kind, + sketch_params: params, + window_secs: query.window_secs, + spatial_filter: String::new(), + grouping: query.group_by.clone(), + item_label: None, + aggregation_input: AggregationInput::SketchEnvelope, + agg_type_override: None, + }); + readouts.push(BackendReadout { + aggregation_id, + op: selected.readout.clone(), + }); + collector_materializations.push(CollectorMaterialization { + query_id: query.query_id.clone(), + metric, + algorithm: format!("{:?}", selected.kind.algorithm()).to_ascii_lowercase(), + parameters: sketch_params_json(&selected.params), + group_by: query.group_by.clone(), + window_secs: query.window_secs, + evidence_source: evidence.map(|e| e.source.clone()), + }); + } + + let plan_id = stable_plan_id(&collector_materializations); + let envelope = PlanEnvelope { + plan_id, + generated_at_unix_ms: environment.observed_at_unix_ms, + planner_revision: PLANNER_REVISION.into(), + capability_snapshot_id: environment.capability_snapshot_id, + }; + let backend_plan = backend_plan::from_stage_config( + &BackendStageConfig { + aggregations, + readouts, + }, + &Vec::::new(), + plan_id, + environment.observed_at_unix_ms, + )?; + let collector_plans = environment + .collector_ids + .into_iter() + .map(|collector_id| CollectorPlan { + collector_id, + envelope: envelope.clone(), + materializations: collector_materializations.clone(), + }) + .collect(); + Ok(CompiledPlanBundle { + envelope, + collector_plans, + backend_plan, + }) + } +} + +fn validate_evidence( + query_id: &str, + evidence: &TopKMembershipEvidence, + env: &DeploymentEnvironment, +) -> Result<(), CompileError> { + let age = env + .observed_at_unix_ms + .saturating_sub(evidence.observed_at_unix_ms); + let valid = evidence.selected_lower_bound.is_finite() + && evidence.excluded_upper_bound.is_finite() + && evidence.selected_lower_bound > evidence.excluded_upper_bound + && (0.0..=1.0).contains(&evidence.interval_failure_probability) + && !evidence.source.trim().is_empty() + && age <= env.max_evidence_age_ms; + if valid { + Ok(()) + } else { + Err(CompileError::InvalidEvidence { + query_id: query_id.into(), + reason: format!("margin/failure/source invalid or age {age}ms exceeds policy"), + }) + } +} + +struct SelectedSketch { + kind: planner_types::post_asap::SketchKind, + params: planner_types::post_asap::SketchParams, + readout: SketchQuery, +} + +fn extract_selected(node: &SummaryNode) -> Option { + let SummaryExpr::SummaryEstimate { + summary_input, + query, + } = &node.expr + else { + return None; + }; + let SummaryExpr::SummaryAgg { + family: SummaryFamilyType::Sketch(kind, _), + .. + } = &summary_input.expr + else { + return None; + }; + Some(SelectedSketch { + kind: kind.clone(), + params: kind.params().clone(), + readout: query.clone(), + }) +} + +fn sketch_params_json(params: &planner_types::post_asap::SketchParams) -> Value { + use planner_types::post_asap::SketchParams as P; + match params { + P::Kll { k } => json!({"k": k}), + P::Cms { width, depth } => json!({"width": width, "depth": depth}), + P::Hll { precision } => json!({"precision": precision}), + P::DDSketch { alpha } => json!({"alpha": alpha}), + P::CmsWithHeap { + width, + depth, + heap_size, + } => json!({"width": width, "depth": depth, "heap_size": heap_size}), + P::Kmv { k } | P::Theta { k } => json!({"k": k}), + P::CountSketch { width, depth } => json!({"width": width, "depth": depth}), + P::CountSketchWithHeap { + width, + depth, + heap_size, + } => json!({"width": width, "depth": depth, "heap_size": heap_size}), + } +} + +fn stable_plan_id(materializations: &[CollectorMaterialization]) -> u64 { + use std::hash::{Hash, Hasher}; + let bytes = serde_json::to_vec(materializations).unwrap_or_default(); + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + bytes.hash(&mut hasher); + hasher.finish() +} + +#[cfg(test)] +mod tests { + use super::*; + + fn environment(now: u64) -> DeploymentEnvironment { + DeploymentEnvironment { + collector_ids: vec!["edge-a".into(), "edge-b".into()], + capability_snapshot_id: "caps-7".into(), + observed_at_unix_ms: now, + max_evidence_age_ms: 60_000, + } + } + + fn request(query_id: &str, promql: &str) -> PlanningRequest { + let accuracy = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let parsed = crate::query_parser::parse_query_expr_canonical(promql, accuracy.clone()) + .expect("canonical query"); + let expr = if promql.starts_with("topk(") { + use crate::optimizer::engine::{DefaultCostModel, RewriteRule, TopKFusion}; + let cost = DefaultCostModel { + raw_bytes_per_sec: 1.0, + deployment: None, + }; + TopKFusion + .try_rewrite(parsed.clone(), &cost) + .unwrap_or(parsed) + } else { + parsed + }; + PlanningRequest { + queries: vec![PlanningQuery { + query_id: query_id.into(), + expr, + source: Source::TimeSeries { metric: "m".into() }, + window_secs: 60, + group_by: vec![], + accuracy, + }], + evidence: HashMap::new(), + planner_revision: PLANNER_REVISION.into(), + } + } + + #[test] + fn compiles_one_decision_into_matching_collector_and_backend_views() { + let bundle = PhysicalCompiler + .compile( + request("q-quantile", "quantile_over_time(0.99, m[1m])"), + environment(10_000), + ) + .expect("compile"); + assert_eq!(bundle.collector_plans.len(), 2); + assert_eq!(bundle.backend_plan.materializations.len(), 1); + assert_eq!(bundle.backend_plan.routing.len(), 1); + for plan in &bundle.collector_plans { + assert_eq!(plan.envelope, bundle.envelope); + assert_eq!(plan.materializations[0].metric, "m"); + assert_eq!(plan.materializations[0].window_secs, 60); + assert!(matches!( + plan.materializations[0].algorithm.as_str(), + "ddsketch" | "kll" + )); + } + } + + #[test] + fn topk_fails_closed_without_membership_evidence() { + let error = PhysicalCompiler + .compile(request("q-topk", "topk(5, m)"), environment(10_000)) + .expect_err("missing certificate must fail"); + assert!(matches!(error, CompileError::Query { .. })); + } + + #[test] + fn stale_topk_evidence_is_rejected_before_planner_selection() { + let mut request = request("q-topk", "topk(5, m)"); + request.evidence.insert( + "q-topk".into(), + TopKMembershipEvidence { + selected_lower_bound: 101.0, + excluded_upper_bound: 100.0, + interval_failure_probability: 0.005, + observed_at_unix_ms: 1, + source: "runtime-margin-monitor".into(), + }, + ); + let error = PhysicalCompiler + .compile(request, environment(100_000)) + .expect_err("stale certificate must fail"); + assert!(matches!(error, CompileError::InvalidEvidence { .. })); + } + + #[test] + fn fresh_topk_evidence_enables_physical_compilation() { + let mut request = request("q-topk", "topk(5, m)"); + request.evidence.insert( + "q-topk".into(), + TopKMembershipEvidence { + selected_lower_bound: 101.0, + excluded_upper_bound: 100.0, + interval_failure_probability: 0.005, + observed_at_unix_ms: 9_500, + source: "runtime-margin-monitor".into(), + }, + ); + let bundle = PhysicalCompiler + .compile(request, environment(10_000)) + .expect("certified TopK compiles"); + assert_eq!(bundle.backend_plan.materializations.len(), 1); + assert_eq!( + bundle.collector_plans[0].materializations[0] + .evidence_source + .as_deref(), + Some("runtime-margin-monitor") + ); + } + + #[test] + fn planner_revision_is_part_of_the_compile_contract() { + let mut request = request("q", "quantile_over_time(0.9, m[1m])"); + request.planner_revision = "different".into(); + assert!(matches!( + PhysicalCompiler.compile(request, environment(10_000)), + Err(CompileError::PlannerRevision { .. }) + )); + } +} diff --git a/control_plane/src/physical/mod.rs b/control_plane/src/physical/mod.rs index 1fa4b3f1d..993f13e51 100644 --- a/control_plane/src/physical/mod.rs +++ b/control_plane/src/physical/mod.rs @@ -22,6 +22,7 @@ pub mod allocator; pub mod colored_dag; +pub mod compiler; pub mod plan; pub mod planner; pub mod sketch_catalog; diff --git a/control_plane/src/planner_selection.rs b/control_plane/src/planner_selection.rs index 4bd389c85..bc706a9e8 100644 --- a/control_plane/src/planner_selection.rs +++ b/control_plane/src/planner_selection.rs @@ -7,7 +7,8 @@ use std::rc::Rc; use asap_aware_mapping::{ - CostModel, Replacement, ReplacementStrategy, SketchAlgorithmStrategy, TargetSubDAG, + AccuracyBudgetAllocator, AccuracyEvidenceProvider, AccuracyModel, CostModel, Replacement, + ReplacementStrategy, SketchAlgorithmStrategy, TargetSubDAG, }; use planner_types::post_asap::{ SummaryExpr, SummaryFamilyType, SummaryField, SummaryNode, SummarySchema, @@ -69,6 +70,33 @@ pub fn select_summary_default(expr: &QueryExpr) -> Result, Selec select_summary(expr, &asap_aware_mapping::DefaultCostModel) } +/// Select from Planner's legal candidates with deployment-supplied accuracy +/// models and typed evidence (for example a TopK membership certificate). +pub fn select_summary_with_evidence( + expr: &QueryExpr, + cost_model: &dyn CostModel, + accuracy_model: &dyn AccuracyModel, + allocator: &dyn AccuracyBudgetAllocator, + evidence: &dyn AccuracyEvidenceProvider, +) -> Result, SelectionError> { + let root = Rc::new(expr.clone()); + let strategy = SketchAlgorithmStrategy::with_models_and_evidence( + cost_model, + accuracy_model, + allocator, + evidence, + ); + let candidate = strategy + .replacements(&TargetSubDAG::new(&root)) + .into_iter() + .next() + .ok_or(SelectionError::NoLegalCandidate)?; + match candidate.replacement { + Replacement::Summary(node) => Ok(node), + Replacement::Rewrite(_) => Err(SelectionError::UnexpectedRewrite), + } +} + /// Select a legal summary when Planner offers one, otherwise preserve the /// subtree explicitly. This mirrors the removed single-tree binder's /// conservative behavior and is appropriate for serving fallbacks; physical diff --git a/docs/developer_docs/control-plane/physical-compiler.md b/docs/developer_docs/control-plane/physical-compiler.md index 7a8a9ef4f..0eaf96562 100644 --- a/docs/developer_docs/control-plane/physical-compiler.md +++ b/docs/developer_docs/control-plane/physical-compiler.md @@ -1,16 +1,16 @@ # Developing the Planner adapter and physical compiler -> Interface status: target public API. Existing migration modules must converge -> on this boundary. +> Interface status: implemented MVP API in +> `control_plane::physical::compiler`. ## Current implementation boundary -The production path has not yet converged on the public interfaces below. It -currently consumes ASAPPlanner types pinned in `control_plane/Cargo.toml`, then -uses `physical::colored_dag::StageAllocator`, `ThreeStageEmitter`, and -`backend_plan::from_stage_config` to produce agent YAML and `BackendPlan`. -Callers must not treat the target `PlanningRequest`, `PhysicalCompiler`, or -`CompiledPlanBundle` examples below as implemented Rust APIs. +The compiler consumes ASAPPlanner types pinned to revision `3afcba6`, selects +from Planner's legal candidate space with backend-owned cost and evidence +inputs, and emits one `CompiledPlanBundle`. The bundle contains a CollectorPlan +for every target collector and the matching BackendPlan. Legacy +`StageAllocator`/`ThreeStageEmitter` paths remain for older publication flows; +they are not a second semantic planner. ASAPPlanner owns logical semantics and summary selection. In particular, a deployment override may choose only a family compatible with the selected @@ -27,17 +27,18 @@ The control plane has three public layers: PlanningRequest | v -PlannerAdapter ----------> SelectedLogicalPlan - | - v +ASAPPlanner candidate selection + | + v PhysicalCompiler -------> CompiledPlanBundle | | v v CollectorPlan BackendPlan ``` -- **Planner adapter** owns the typed call to ASAPPlanner. It supplies the whole - workload and receives one selected logical plan without copying Planner IR. +- **Planner selection boundary** is + `planner_selection::select_summary_with_evidence`. It enumerates Planner's + candidates and commits only a legal candidate. - **Physical compiler** adds backend-owned placement, windows, transport, and runtime capabilities without changing logical semantics. - **Plan bundle** is the only output passed to publication. CollectorPlan and @@ -49,28 +50,15 @@ remain public ASAPPlanner interfaces. Runtime publication is documented in ## 2. Public interfaces and definitions -### Planner adapter - -```rust -pub trait PlannerAdapter { - type Error; - - fn select( - &self, - request: PlanningRequest, - ) -> Result; -} -``` +### Planning request `PlanningRequest` is backend-owned request context around Planner's canonical -workload value: +per-query IR: ```rust pub struct PlanningRequest { - pub workload: asap_planner::Workload, - pub schema: asap_planner::SchemaCatalog, - pub constraints: asap_planner::PlanningConstraints, - pub cost_inputs: asap_planner::CostInputs, + pub queries: Vec, + pub evidence: HashMap, pub planner_revision: String, } ``` @@ -79,20 +67,21 @@ Input definitions: | Field | Definition | | --- | --- | -| `workload` | Complete workload; shared queries must not be split into independent calls. | -| `schema` | Source/label/type information required to bind queries. | -| `constraints` | Accuracy and logical requirements supplied by the caller. | -| `cost_inputs` | Measured/declared logical cost inputs; unknown values stay unknown. | +| `queries` | Canonical `QueryExpr`, source, window, grouping labels, accuracy, and stable query ID. | +| `evidence` | Optional typed TopK membership certificates keyed by query ID. | | `planner_revision` | Immutable Planner build/revision used for reproducibility. | -`SelectedLogicalPlan` wraps Planner's public selected post-ASAP workload plan -and correlation metadata; it does not define another DAG: +TopK evidence is accepted only when its selected lower bound is strictly above +the excluded upper bound, its failure probability is valid, its source is +non-empty, and its observation is fresh under `DeploymentEnvironment`. ```rust -pub struct SelectedLogicalPlan { - pub workload_plan: asap_planner::SelectedWorkloadPlan, - pub planner_revision: String, - pub query_ids: Vec, +pub struct TopKMembershipEvidence { + pub selected_lower_bound: f64, + pub excluded_upper_bound: f64, + pub interval_failure_probability: f64, + pub observed_at_unix_ms: u64, + pub source: String, } ``` @@ -102,37 +91,26 @@ planning from each implementing their own query-to-summary mapping. ### Physical compiler ```rust -pub trait PhysicalCompiler { - type Error; - - fn compile( +impl PhysicalCompiler { + pub fn compile( &self, - selected: SelectedLogicalPlan, + request: PlanningRequest, environment: DeploymentEnvironment, - policy: RuntimePolicy, - ) -> Result; + ) -> Result; } ``` ```rust pub struct DeploymentEnvironment { - pub topology: DeploymentTopology, - pub collectors: Vec, - pub backend: BackendTarget, + pub collector_ids: Vec, pub capability_snapshot_id: String, -} - -pub struct RuntimePolicy { - pub activation: Timestamp, - pub expiry: Option, - pub freshness: FreshnessPolicy, - pub retention: RetentionPolicy, - pub transmission: TransmissionPolicy, + pub observed_at_unix_ms: u64, + pub max_evidence_age_ms: u64, } pub struct CompiledPlanBundle { pub envelope: PlanEnvelope, - pub collector_plans: Vec, + pub collector_plans: Vec, // complete per-target projections pub backend_plan: BackendPlan, } ``` @@ -141,16 +119,16 @@ Supporting public types: | Type | Definition | | --- | --- | -| `DeploymentTopology` | Runtime stages, network relationships, and isolation boundaries available for placement. | -| `CollectorTarget` | Collector identity, edge assignment, endpoint reference, and advertised capability snapshot. | -| `BackendTarget` | Data-plane identity, endpoint reference, storage routes, and advertised capabilities. | -| `FreshnessPolicy` | Maximum readiness lag, watermark, and allowed-lateness requirements. | -| `RetentionPolicy` | Duration and lifecycle rules for active/draining materializations. | -| `TransmissionPolicy` | Allowed raw/full/delta modes, cadence, encoding, and checkpoint limits. | -| `PlanEnvelope` | Shared `plan_id`, `plan_version`, activation/expiry, backend compatibility, and Planner revision. | -| `CollectorPlan` | Versioned public YAML execution contract owned by ASAPCollector. | +| `DeploymentEnvironment` | Target collector IDs, capability snapshot identity, planning time, and evidence freshness policy. | +| `PlanEnvelope` | Shared deterministic `plan_id`, generation time, capability snapshot, and Planner revision. | +| `CollectorPlan` | Serializable execution projection consumed by ASAPCollector. | | `BackendPlan` | Versioned public data-plane materialization/routing contract defined in this repository. | +Current MVP limits are explicit: time-series sources and sketch +materializations are supported; table sources and non-sketch selected families +return `CompileError`. Runtime activation/expiry and richer topology placement +remain publication-layer work and are not claimed by this compiler API. + Output definitions: | Output | Definition | diff --git a/docs/developer_docs/control-plane/plan-publication.md b/docs/developer_docs/control-plane/plan-publication.md index 6d62ec483..47b26a718 100644 --- a/docs/developer_docs/control-plane/plan-publication.md +++ b/docs/developer_docs/control-plane/plan-publication.md @@ -1,268 +1,80 @@ -# Developing runtime plan publication +# MVP physical-plan publication -> Interface status: target public API. OpAMP transport exists today; semantic -> CollectorPlan application/reporting is still incomplete. +This page describes the implemented MVP contract. The control plane compiles +ASAPPlanner's selected post-ASAP IR once and projects that decision into one +typed `BackendPlan` plus one target-specific `CollectorPlan` per Collector. -## 1. Code architecture +## API -Publication begins only after physical compilation returns a complete bundle: +`POST /api/v1/physical-plan/compile-and-publish` accepts: -```text -CompiledPlanBundle - | - v -PlanPublisher - | | - v v -CollectorClient BackendPlanClient - | | - v v -CollectorReport BackendPlanReport - \ / - v v - ActivationResult -``` - -`CollectorClient` is the ASAPQuery-side counterpart of ASAPCollector's -authoritative -[`opamp-config-push.md`](https://github.com/ProjectASAP/ASAPCollector/blob/main/docs/developer_docs/asapcollector/opamp-config-push.md). -This repository does not redefine CollectorPlan fields. - -Publication also does not reinterpret either plan. The physical compiler owns -their contents; the publisher owns delivery, staging, activation barriers, -report correlation, and rollback. - -## 2. Public interfaces and definitions - -### Runtime clients - -```rust -pub trait CollectorPlanClient { - type Error; - - async fn stage( - &self, - target: CollectorTarget, - plan: CollectorPlan, - ) -> Result; - - async fn activate( - &self, - target: CollectorTarget, - plan_id: &str, - plan_version: u64, - ) -> Result; - - async fn abort( - &self, - target: CollectorTarget, - plan_id: &str, - plan_version: u64, - ) -> Result; -} - -pub trait BackendPlanClient { - type Error; - - async fn stage( - &self, - target: BackendTarget, - plan: BackendPlan, - ) -> Result; - - async fn activate( - &self, - target: BackendTarget, - plan_id: &str, - plan_version: u64, - ) -> Result; - - async fn abort( - &self, - target: BackendTarget, - plan_id: &str, - plan_version: u64, - ) -> Result; -} -``` - -Collector transport requirements come directly from the corresponding -ASAPCollector interface: - -- OpAMP `AgentRemoteConfig`/`AgentConfigMap`; -- exact entry name `asap-collector-plan.yaml`; -- YAML CollectorPlan with `content_type: application/yaml`; -- OpAMP `config_hash` identifies bytes, not cross-runtime plan semantics; and -- `RemoteConfigStatus.APPLIED` is delivery/application evidence, not semantic - activation evidence. - -The payload delivered to each target is different even though it shares one -bundle envelope: - -| Target | Published artifact | Components that must stage it | -| --- | --- | --- | -| Each ASAPCollector | That target's `CollectorPlan` YAML in `asap-collector-plan.yaml` | OpAMP receiver, plan validator, collection/precompute runtime, window manager, transmission/export pipeline | -| ASAPQuery data plane | One typed `BackendPlan` | Plan manager, ingest/precompute engine, SID/metadata registry, summary store, query router/readout and fallback adapter | +- `queries`: query ID, PromQL, metric, window seconds, grouping labels, and a + typed `AccuracyTarget`; +- `collector_ids`: the required OpAMP agent IDs; +- `capability_snapshot_id` and the exact `planner_revision`; +- optional per-query TopK evidence, with `max_evidence_age_ms`; and +- `apply_timeout_ms` (default 10000). -The Collector report proves the expected materialization definitions are -installed and executable. It does not enumerate future per-label SIDs. The -Backend report proves matching descriptors, ingest/storage routes, and query -routes are staged. SIDs are allocated or resolved later as concrete -materialized series arrive. +Unknown JSON fields, empty target/query sets, zero windows/timeouts, stale +evidence, and a Planner revision mismatch are rejected. The response is only +successful after every target has applied the same generated `plan_id`. -### Application reports - -```rust -pub enum ApplicationStatus { - Rejected, - Staged, - Active, - Expired, - Failed, -} - -pub struct ApplicationError { - pub code: String, - pub path: String, - pub message: String, -} - -pub struct CollectorApplicationReport { - pub plan_id: String, - pub plan_version: u64, - pub backend_compat: String, - pub remote_config_hash: Vec, - pub status: ApplicationStatus, - pub active_materialization_ids: Vec, - pub effective_capability_hash: String, - pub observed_at: Timestamp, - pub activated_at: Option, - pub errors: Vec, -} - -pub struct BackendApplicationReport { - pub plan_id: String, - pub plan_version: u64, - pub backend_compat: String, - pub status: ApplicationStatus, - pub active_materialization_ids: Vec, - pub observed_at: Timestamp, - pub activated_at: Option, - pub errors: Vec, -} -``` - -The collector report corresponds to capability -`io.asap.collector.plan.v1`, message type `application_report`. Unknown report -versions or missing required fields are errors. - -Why reports are separate from transport acknowledgement: the MVP must prove the -runtime applied the intended semantic plan, not merely that bytes arrived. - -### Publisher - -```rust -pub trait PlanPublisher { - type Error; - - async fn publish( - &self, - bundle: CompiledPlanBundle, - ) -> Result; - - async fn rollback( - &self, - plan_id: &str, - plan_version: u64, - ) -> Result; -} - -pub struct ActivationResult { - pub plan_id: String, - pub plan_version: u64, - pub status: ApplicationStatus, - pub collector_reports: Vec, - pub backend_report: BackendApplicationReport, -} -``` - -`publish` returns `Active` only when every required runtime reports the same -plan/version/compatibility and expected materializations. Partial staging is an -error result and keeps the prior valid plan authoritative. - -### Publication and activation sequence +## Installation order and failure semantics ```text -1. validate complete CompiledPlanBundle -2. stage BackendPlan (ingest/store/query routes not yet authoritative) -3. stage every CollectorPlan (collection/export not yet authoritative) -4. correlate reports with envelope, targets, capabilities, and expected materializations -5. wait for the common activation time and all readiness conditions -6. issue activation for the common time; install backend routing before allowing Collector export -7. observe emitted plan/materialization IDs and retain the previous version for drain/rollback +validate request and compile one bundle + | + v +preflight every Collector capability + | + v +POST typed BackendPlan protobuf; require 2xx + | + v +publish target-specific CollectorPlans over OpAMP + | + v +require exact (agent_id, plan_id, APPLIED) from every target ``` -Backend staging precedes Collector staging so the consumer can validate the -contract before new producers are allowed to emit. Staging order alone is not -activation: both sides remain gated by the shared activation time and the -publisher's readiness decision. If any required target rejects or times out, -the publisher aborts the new version on every staged target; the previous -unexpired bundle remains authoritative. +Preflight happens before backend mutation. Publication repeats capability +validation to close disconnect races. A missing backend endpoint, backend +non-2xx, Collector disconnect, timeout, malformed report, wrong plan ID, or +`FAILED` status fails the request. This path intentionally does not inherit the +legacy replanner's best-effort behavior. -`stage`, `activate`, and `abort` are distinct semantic operations even when a -transport implements them as versioned remote-config updates. An `APPLIED` -transport acknowledgement for staged bytes must not be mapped directly to -`Active`. Distributed activation cannot be literally instantaneous, so the -backend installs the accepting ingest view first, both sides use the common -activation timestamp, and query routing becomes authoritative only after the -required active reports correlate. During that bounded transition the backend -may accept new-version payloads without serving queries from an incomplete -new-version view. +The MVP endpoint installs the backend before enabling new producers. It does +not claim distributed atomic activation or rollback; those remain post-MVP +work. If a Collector fails after backend installation, the backend has a +superset accepting view but the request fails and no success is reported. -For replacement, Collector state created under the old materialization drains -according to its lifecycle policy, while the backend retains the corresponding -SID metadata and query route for the declared drain horizon. A changed summary -semantic produces a new materialization identity and new SIDs; publication -must not relabel old state into the new definition. +## OpAMP custom capability -## 3. Adding and verifying functionality +The exact capability and message types shared with ASAPCollector are: -### Add another collector transport +| Field | Value | +| --- | --- | +| capability | `io.projectasap.collector-plan.v1` | +| server-to-agent message type | `collector_plan` | +| agent-to-server message type | `plan_status` | +| payload encoding | UTF-8 JSON | -1. Implement `CollectorPlanClient`; keep CollectorPlan unchanged. -2. Preserve plan identity separately from transport byte identity. -3. Map transport errors to structured client errors. -4. Verify identical re-delivery is idempotent and conflicting bytes for the same - plan/version are rejected. +`collector_plan` is the serialized `physical::compiler::CollectorPlan`. A +status payload is strict JSON: -Interpretation: a successful `stage` report is not global activation; only -`PlanPublisher::publish` can return an active bundle. - -### Add an application status or report field - -1. Version the public report schema/capability. -2. Define required/optional behavior and compatibility. -3. Update collector and backend clients together. -4. Verify older readers reject unknown required semantics rather than defaulting. - -### Add rollback policy +```json +{"plan_id": 42, "status": "APPLIED", "error": null} +``` -1. Select only a retained complete bundle through `rollback`. -2. Stage both runtime sides like a normal publication. -3. Verify the result reports the restored version and all materializations. -4. Verify failed rollback leaves the current active bundle unchanged. +`status` is exactly `APPLIED` or `FAILED`. Transport/config acknowledgements +are not accepted as semantic plan evidence. Reports are scoped to the agent ID +of the WebSocket connection, preventing one Collector from acknowledging +another Collector's plan. -### Required output checks +## Backend wire contract -- reports match bundle identity and expected materialization sets; -- stale/expired/conflicting versions fail; -- collector-only or backend-only success never returns `Active`; -- report artifacts are machine-readable by the MVP harness; and -- post-activation emitted state carries the activated identities; -- the BackendPlan declares ingest/storage for every CollectorPlan output; -- every backend query route references a staged materialization descriptor; -- Collector reports refer to materialization definitions, not runtime SIDs; -- no Collector begins new-version export before compatible backend readiness; - and -- failed activation sends an abort/rollback action to every target that staged - the candidate version. +The matching `BackendPlan` uses the protobuf contract documented in +`control_plane/docs/design-backend-plan-wire-format.md` and is sent to +`POST /api/v1/backend-plan` with `application/x-protobuf`. Both projections +carry the same numeric `plan_id`; the compiler, not either transport, owns the +materialization choice.