Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
69 commits
Select commit Hold shift + click to select a range
292b15c
Restack PR #756 with implementation before standalone acceptance
zzylol Sep 26, 2026
80d17d4
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
a3ce0a7
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
79981d6
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
3f7c5b9
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
8906e35
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
8305a78
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
ec3600e
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
a37df74
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
2e4f65c
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
b92af97
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
77835c5
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
2c38ecb
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
13660da
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
e0ae1c2
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
5d39a91
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
2de7864
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
b24fa79
Merge branch 'standalone/stack-765' into standalone/stack-756
zzylol Sep 26, 2026
37d7438
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
68917b8
merge: preserve SDS binding validation in diagnostics branch
zzylol Sep 26, 2026
f3a353c
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
2d83f9f
fix(diagnostics): report deployed output identity for SDS coordinates
zzylol Sep 26, 2026
85b4972
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
626601c
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
80596af
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
5382f21
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
9627a32
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
8fa4973
Merge branch 'impl/sds-stack-765' into impl/sds-stack-756
zzylol Sep 26, 2026
67082a1
feat: trace deployment selection and bound SDS lifecycle
zzylol Sep 27, 2026
fccb0e6
Merge current workload costing into runtime diagnostics
zzylol Sep 27, 2026
f77585a
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
6df1563
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
b1f0062
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
79d9fc1
feat: trace installed physical compilation recovery and input binding
zzylol Sep 27, 2026
d25e935
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
486f9d4
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
01584c0
feat: trace native SDS publication and bound recovery
zzylol Sep 27, 2026
a5cb52c
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
1518443
feat: trace installed current-series native execution stages
zzylol Sep 27, 2026
701224f
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 27, 2026
3632c14
feat: trace Planner snapshot heap candidate compilation
zzylol Sep 27, 2026
a5de159
Merge native Rate ranking and instrument physical input execution
zzylol Sep 27, 2026
0f17e15
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
292406b
Trace native maintenance, bound heap reads and shared physical execution
zzylol Sep 28, 2026
2cb8c3d
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
1d90397
Document native Rate aggregation diagnostics across maintenance and q…
zzylol Sep 28, 2026
acb8c7b
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
349befc
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
ba3d443
merge: align runtime diagnostics with dataset-bound SDS contracts
zzylol Sep 28, 2026
55bdcbe
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
19344c6
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
a51a25b
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
b558222
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
50e6b89
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
7865cb9
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
36f934a
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
43c0423
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
6daa50a
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
4326810
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
631d648
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
330062f
Keep recovery diagnostics while propagating metadata errors
zzylol Sep 28, 2026
a57dd8d
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
9d7bf7d
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
f4f9e52
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
86dd68a
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
300b7db
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
3dbaadb
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
6abf84c
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
5e7cfb9
Merge branch 'impl/sds-stack-761' into impl/sds-stack-756
zzylol Sep 28, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions control_plane/src/backend_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -268,12 +268,16 @@ impl BackendClient {
}

/// Publish one authoritative catalog generation and all plans that reference it.
#[tracing::instrument(level = "debug", target = "asap_runtime_debug", skip_all,
fields(plan_id = publication.transmission_plan.envelope.plan_id,
plan_version = publication.transmission_plan.envelope.plan_version))]
pub async fn post_catalog_plan_typed(
&self,
publication: &crate::physical::publication::PhysicalPlanPublication,
storage_routing: Option<serde_json::Value>,
adaptation_evidence: &[crate::physical::compiler::RuntimeAdaptationEvidence],
) -> std::result::Result<(), BackendPostError> {
tracing::debug!(target: "asap_runtime_debug", "backend plan stage request started");
let body = publication
.install_request(storage_routing, adaptation_evidence.to_vec())
.map_err(|error| BackendPostError::Permanent(anyhow::anyhow!(error)))?;
Expand All @@ -285,6 +289,8 @@ impl BackendClient {
.await
.map_err(classify_reqwest_error)?;
let status = response.status();
tracing::debug!(target: "asap_runtime_debug", http_status = %status,
"backend plan stage response received");
if status.is_success() {
Ok(())
} else {
Expand All @@ -297,11 +303,18 @@ impl BackendClient {
}
}

#[tracing::instrument(
level = "debug",
target = "asap_runtime_debug",
skip_all,
fields(plan_id, plan_version)
)]
pub async fn discard_staged_physical_plan(
&self,
plan_id: u64,
plan_version: u64,
) -> std::result::Result<(), BackendPostError> {
tracing::debug!(target: "asap_runtime_debug", "staged backend plan cleanup requested");
let response = self
.http
.post(format!(
Expand All @@ -324,11 +337,18 @@ impl BackendClient {
}
}

#[tracing::instrument(
level = "debug",
target = "asap_runtime_debug",
skip_all,
fields(plan_id, plan_version)
)]
pub async fn activate_physical_plan(
&self,
plan_id: u64,
plan_version: u64,
) -> std::result::Result<(), BackendPostError> {
tracing::debug!(target: "asap_runtime_debug", "backend plan activation request started");
let url = format!("{}/activate", derive_physical_plan_url(&self.endpoint));
let response = self
.http
Expand All @@ -341,6 +361,8 @@ impl BackendClient {
.await
.map_err(classify_reqwest_error)?;
let status = response.status();
tracing::debug!(target: "asap_runtime_debug", http_status = %status,
"backend plan activation response received");
if status.is_success() {
Ok(())
} else {
Expand Down
84 changes: 80 additions & 4 deletions control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,13 @@ use axum::{
use serde::{Deserialize, Serialize};
use serde_json::json;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tracing::info;

static NEXT_PLAN_CALL_ID: AtomicU64 = AtomicU64::new(1);

use opamp::OpampServer;
use physical::deployment_cost::online as online_cost_model;
use physical::deployment_cost::online::{init_store as init_online_store, OnlineMetricsStore};
Expand Down Expand Up @@ -50,7 +53,16 @@ struct AppState {

#[tokio::main]
async fn main() {
tracing_subscriber::fmt::init();
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
)
.with_file(true)
.with_line_number(true)
.with_target(true)
.with_writer(std::io::stdout)
.init();

let api_addr = std::env::var("CONTROLLER_ADDR").unwrap_or_else(|_| "0.0.0.0:8080".into());
let opamp_addr =
Expand Down Expand Up @@ -271,11 +283,16 @@ async fn compile_and_publish_physical_plan(
mut request: CompileAndPublishPhysicalPlanRequest,
frontend: QueryFrontend,
) -> Response {
let call_id = NEXT_PLAN_CALL_ID.fetch_add(1, Ordering::Relaxed);
let started = std::time::Instant::now();
tracing::debug!(target: "asap_runtime_debug", call_id, ?frontend,
"physical plan compilation requested");
// Serialize typed activations so an older response cannot overwrite the
// catalog recorded after a newer backend activation.
let mut active_catalog = st.active_summary_catalog.lock().await;
if let Some(erp) = &mut request.erp {
if let Err(error) = erp.hydrate_observed_shape(&st.runtime_samples) {
tracing::warn!(call_id, %error, "ERP observation hydration failed");
return (StatusCode::UNPROCESSABLE_ENTITY, error).into_response();
}
let catalog = active_catalog.clone();
Expand All @@ -287,16 +304,37 @@ async fn compile_and_publish_physical_plan(
(bundle, ids, timeout, adaptation, manifests)
}
Ok((None, ..)) => {
tracing::error!(
call_id,
"physical plan compilation produced no deployable plan"
);
return (
StatusCode::INTERNAL_SERVER_ERROR,
"publication requires a selected plan",
)
.into_response()
.into_response();
}
Err(response) => {
tracing::warn!(call_id, status = %response.0, error = %response.1,
"physical plan compilation failed");
return physical_compile_failure(response);
}
Err(response) => return physical_compile_failure(response),
};
info!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
collector_count = bundle.collector_plans.len(),
"physical plan compiled"
);

let Some(backend) = st.backend_client.as_ref() else {
tracing::warn!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
"backend endpoint is not configured"
);
return (
StatusCode::SERVICE_UNAVAILABLE,
"CONTROLLER_BACKEND_ENDPOINT is required for physical-plan publication".to_string(),
Expand All @@ -308,6 +346,8 @@ async fn compile_and_publish_physical_plan(
.ensure_collector_plan_targets(&bundle.collector_plans, apply_timeout)
.await
{
tracing::warn!(call_id, plan_id = bundle.envelope.plan_id, plan_version = bundle.envelope.plan_version,
%error, "collector plan preflight failed");
return (
StatusCode::BAD_GATEWAY,
format!("collector physical-plan preflight failed: {error}"),
Expand All @@ -317,11 +357,14 @@ async fn compile_and_publish_physical_plan(
let publication = match bundle.to_publication_artifact() {
Ok(publication) => publication,
Err(error) => {
tracing::error!(call_id, plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version, %error,
"catalog publication artifact failed");
return (
StatusCode::INTERNAL_SERVER_ERROR,
format!("invalid catalog publication: {error}"),
)
.into_response()
.into_response();
}
};
if let Err(error) = backend
Expand All @@ -332,17 +375,27 @@ async fn compile_and_publish_physical_plan(
)
.await
{
tracing::warn!(call_id, plan_id = bundle.envelope.plan_id, plan_version = bundle.envelope.plan_version,
%error, "backend plan staging failed");
return (
StatusCode::BAD_GATEWAY,
format!("backend rejected physical plan: {error}"),
)
.into_response();
}
info!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
"backend physical plan staged"
);
if let Err(error) = st
.opamp
.publish_collector_plans(&bundle.collector_plans, apply_timeout)
.await
{
tracing::warn!(call_id, plan_id = bundle.envelope.plan_id, plan_version = bundle.envelope.plan_version,
%error, "collector plan publication failed");
let cleanup = backend
.discard_staged_physical_plan(bundle.envelope.plan_id, bundle.envelope.plan_version)
.await;
Expand All @@ -352,12 +405,26 @@ async fn compile_and_publish_physical_plan(
)
.into_response();
}
info!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
"collector plans published"
);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64;
let activation_wait = bundle.envelope.activation_unix_ms.saturating_sub(now);
if activation_wait > apply_timeout.as_millis() as u64 {
tracing::warn!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
activation_wait_ms = activation_wait,
timeout_ms = apply_timeout.as_millis() as u64,
"scheduled activation exceeds apply timeout; backend plan remains staged"
);
return (
StatusCode::GATEWAY_TIMEOUT,
"activation time exceeds apply_timeout_ms; backend remains staged".to_string(),
Expand All @@ -371,6 +438,8 @@ async fn compile_and_publish_physical_plan(
.activate_physical_plan(bundle.envelope.plan_id, bundle.envelope.plan_version)
.await
{
tracing::warn!(call_id, plan_id = bundle.envelope.plan_id, plan_version = bundle.envelope.plan_version,
%error, "backend plan activation failed");
return (
StatusCode::BAD_GATEWAY,
format!("backend physical-plan activation failed: {error}"),
Expand All @@ -379,6 +448,13 @@ async fn compile_and_publish_physical_plan(
}

*active_catalog = Some(Arc::new(bundle.summary_catalog));
info!(
call_id,
plan_id = bundle.envelope.plan_id,
plan_version = bundle.envelope.plan_version,
elapsed_ms = started.elapsed().as_millis() as u64,
"physical plan active"
);

Json(CompileAndPublishPhysicalPlanResponse {
cost_comparison: bundle.cost_comparison,
Expand Down
14 changes: 14 additions & 0 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -532,11 +532,15 @@ pub fn gos_policy_from_accuracy_budget(
})
}

#[tracing::instrument(level = "debug", target = "asap_runtime_debug", skip_all,
fields(plan_id = envelope.plan_id, plan_version = envelope.plan_version,
producer_count = precompute.producers.len()))]
pub fn build_transmission_plan(
envelope: PlanEnvelope,
precompute: &PrecomputePlan,
runtime_policies: &BTreeMap<asap_types::PolicyFingerprint, RuntimeRulePolicy>,
) -> Result<TransmissionPlan, TransmissionPlanError> {
tracing::debug!(target: "asap_runtime_debug", "transmission plan construction started");
if envelope != precompute.envelope {
return Err(TransmissionPlanError::EnvelopeMismatch);
}
Expand Down Expand Up @@ -1000,6 +1004,9 @@ impl BackendLocalPlanningInput {
let asap_aware_mapping::Replacement::Summary(root) = candidate.replacement else {
continue;
};
let _physical = tracing::debug_span!(target: "asap_runtime_debug", "physical_candidate_compile",
stage = "planner.physical_candidate", query_id = %query.query_id,
input_kind = "bound_promql_vector").entered();
let compiled = asap_physical_operators::physical_planner::promql_rows::compile_current_series_readout(&root)
.or_else(|_| asap_physical_operators::physical_planner::promql_rows::compile_rate_ranking(&root).map(|(_, program)| program));
if let Ok(physical) = asap_physical_operators::physical_planner::promql_rows::compile_fixed_window_rate_aggregation(&root) {
Expand Down Expand Up @@ -1308,12 +1315,16 @@ impl DeploymentPlanCompiler {
self.compile_for_frontend(request, environment, QueryFrontend::MetricsQl)
}

#[tracing::instrument(level = "debug", target = "asap_runtime_debug", skip_all,
fields(frontend = ?frontend, plan_version = environment.plan_version,
query_count = request.queries.len()))]
pub fn compile_for_frontend(
&self,
mut request: PhysicalCompilationRequest,
environment: PhysicalDeploymentContext,
frontend: QueryFrontend,
) -> Result<CompiledPhysicalPlan, CompileError> {
tracing::debug!(target: "asap_runtime_debug", stage = "deployment.bind", "backend deployment compiler entered");
environment
.dataset_identity
.validate()
Expand Down Expand Up @@ -2718,6 +2729,8 @@ pub fn select_logical_roots_with_error_resource_profiles(
select_logical_roots_with_trace(queries, roots, evidence, exact_costs, erp).map(|_| ())
}

#[tracing::instrument(level = "debug", target = "asap_runtime_debug", skip_all,
fields(query_count = queries.len(), root_count = roots.len()))]
pub fn select_logical_roots_with_trace(
queries: &mut [QueryCompilationInput],
roots: Vec<Rc<QueryExpr>>,
Expand Down Expand Up @@ -2767,6 +2780,7 @@ fn logical_roots_and_candidates(
now_ms: u64,
mut candidates_out: Option<&mut Vec<Vec<(usize, Rc<SummaryNode>)>>>,
) -> Result<Vec<serde_json::Value>, CompileError> {
tracing::debug!(target: "asap_runtime_debug", stage = "planner.select", "Planner selection entered");
let mut traces = Vec::new();
if roots.len() != queries.len() {
return Err(CompileError::Snapshot(
Expand Down
2 changes: 2 additions & 0 deletions control_plane/src/physical/maintained_population.rs
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,8 @@ pub(super) fn operator(

/// The maintained population is a deployment source; ranking is compiled by
/// Planner before this candidate is priced or installed.
#[tracing::instrument(level = "debug", target = "asap_runtime_debug", skip_all,
fields(stage = "physical.compile_install", query_id = %entry.query_id, input_kind = "current_series_snapshot"), err)]
pub(super) fn install_native_topk(
entry: &mut asap_types::query_plan::QueryPlanEntry,
selected: &std::rc::Rc<SummaryNode>,
Expand Down
16 changes: 16 additions & 0 deletions control_plane/src/physical/workload_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -795,6 +795,17 @@ pub fn select_candidates(
}
}
}
tracing::debug!(target: "asap_runtime_debug", stage = "deployment.candidate_inventory",
candidate_count = candidate_evaluations.len(),
search_coverage = ?materialization_search_coverage,
"bounded candidate evaluation completed; absence is not a cost rejection");
for candidate in &candidate_evaluations {
tracing::debug!(target: "asap_runtime_debug", stage = "deployment.candidate_evaluation",
candidate_id = ?candidate.candidate_id,
physical_candidate_id = ?candidate.physical_candidate_id,
status = ?candidate.status, total_cost = candidate.total_cost,
reason = ?candidate.unavailable_reason, "candidate evaluated");
}
let selected = asap_physical_operators::physical_planner::select_candidate(
priced_candidates,
|(cost, _, manifest, _, _)| {
Expand All @@ -817,6 +828,11 @@ pub fn select_candidates(
})?;
let (_, mut plan, selected_manifest, component_costs, best_index) = selected.candidate;
candidate_evaluations[best_index].status = CandidateEvaluationStatus::Selected;
tracing::debug!(target: "asap_runtime_debug", stage = "deployment.candidate_selected",
plan_id = plan.envelope.plan_id, plan_version = plan.envelope.plan_version,
candidate_id = ?candidate_evaluations[best_index].candidate_id,
total_cost = candidate_evaluations[best_index].total_cost,
"lowest quoted cost selected within the admitted inventory");
plan.cost_comparison = Some(CandidatePlanSelectionReport {
planner_selection_trace,
materialization_search_coverage,
Expand Down
Loading