Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
4 changes: 4 additions & 0 deletions control_plane/proto/backend_plan.proto
Original file line number Diff line number Diff line change
Expand Up @@ -191,4 +191,8 @@ message BackendPlan {
map<uint64, Materialization> materializations = 3;
repeated RoutingEntry routing = 4;
repeated MonitorSpec monitors = 5;
uint64 plan_version = 6;
uint64 activation_unix_ms = 7;
optional uint64 expiry_unix_ms = 8;
string backend_compat = 9;
}
29 changes: 29 additions & 0 deletions control_plane/src/backend_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -391,6 +391,35 @@ impl BackendClient {
Err(classify_http_status(status, body, "PhysicalPlan POST"))
}
}

pub async fn activate_physical_plan(
&self,
plan_id: u64,
plan_version: u64,
) -> std::result::Result<(), BackendPostError> {
let url = format!("{}/activate", derive_physical_plan_url(&self.endpoint));
let response = self
.http
.post(&url)
.json(&serde_json::json!({
"plan_id": plan_id,
"plan_version": plan_version,
}))
.send()
.await
.map_err(classify_reqwest_error)?;
let status = response.status();
if status.is_success() {
Ok(())
} else {
let body = response.text().await.unwrap_or_default();
Err(classify_http_status(
status,
body,
"PhysicalPlan activation POST",
))
}
}
}

fn derive_physical_plan_url(endpoint: &str) -> String {
Expand Down
4 changes: 4 additions & 0 deletions control_plane/src/backend_plan/from_stage_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,10 @@ pub fn from_stage_config(
Ok(BackendPlan {
plan_id,
generated_at_unix_ms,
plan_version: 1,
activation_unix_ms: generated_at_unix_ms,
expiry_unix_ms: None,
backend_compat: super::BACKEND_COMPAT.into(),
materializations,
routing,
monitors: plan_monitors,
Expand Down
64 changes: 64 additions & 0 deletions control_plane/src/backend_plan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ use asap_types::{AggregationType, MonitorSpec, PolicyFingerprint};
use prost::Message as _;
use thiserror::Error;

pub const BACKEND_COMPAT: &str = "asap-query-backend.v1";

use crate::physical::runtime_capability::Capability;
use asap_types::enums::WindowKind;
use planner_types::post_asap::{
Expand Down Expand Up @@ -82,6 +84,27 @@ pub enum ValidationError {
IncompatibleRoute { fingerprint: u64 },
#[error("stale plan generation: incoming={incoming}, active={active}")]
StaleGeneration { incoming: u64, active: u64 },
#[error("stale plan version for plan {plan_id}: incoming={incoming}, active={active}")]
StalePlanVersion {
plan_id: u64,
incoming: u64,
active: u64,
},
#[error("plan {plan_id} version {plan_version} was reused with different content")]
ReusedPlanVersion { plan_id: u64, plan_version: u64 },
#[error("non-bootstrap plan must have a non-zero plan version")]
ZeroPlanVersion,
#[error("non-bootstrap plan must have an activation time")]
MissingActivation,
#[error("plan expiry {expiry} is not after activation {activation}")]
InvalidExpiry { activation: u64, expiry: u64 },
#[error("non-bootstrap plan must declare backend compatibility")]
MissingBackendCompat,
#[error("unsupported backend compatibility `{actual}`; expected `{expected}`")]
UnsupportedBackendCompat {
actual: String,
expected: &'static str,
},
}

// ── WindowSpec ───────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -772,6 +795,10 @@ pub struct BackendPlan {
/// existing convention for the pre-`BackendPlan` wire format).
pub plan_id: u64,
pub generated_at_unix_ms: u64,
pub plan_version: u64,
pub activation_unix_ms: u64,
pub expiry_unix_ms: Option<u64>,
pub backend_compat: String,
pub materializations: HashMap<PolicyFingerprint, Materialization>,
pub routing: Vec<RoutingEntry>,
pub monitors: Vec<MonitorSpec>,
Expand All @@ -782,6 +809,10 @@ impl From<&BackendPlan> for proto::BackendPlan {
proto::BackendPlan {
plan_id: p.plan_id,
generated_at_unix_ms: p.generated_at_unix_ms,
plan_version: p.plan_version,
activation_unix_ms: p.activation_unix_ms,
expiry_unix_ms: p.expiry_unix_ms,
backend_compat: p.backend_compat.clone(),
materializations: p
.materializations
.iter()
Expand Down Expand Up @@ -809,6 +840,10 @@ impl TryFrom<proto::BackendPlan> for BackendPlan {
Ok(BackendPlan {
plan_id: p.plan_id,
generated_at_unix_ms: p.generated_at_unix_ms,
plan_version: p.plan_version,
activation_unix_ms: p.activation_unix_ms,
expiry_unix_ms: p.expiry_unix_ms,
backend_compat: p.backend_compat,
materializations,
routing,
monitors: p.monitors.into_iter().map(Into::into).collect(),
Expand All @@ -832,6 +867,31 @@ impl BackendPlan {
/// Validate cross-references and invariants required before a decoded
/// plan may become visible to ingest or query readers.
pub fn validate(&self) -> Result<(), ValidationError> {
if self.plan_id != 0 {
if self.plan_version == 0 {
return Err(ValidationError::ZeroPlanVersion);
}
if self.activation_unix_ms == 0 {
return Err(ValidationError::MissingActivation);
}
if self.backend_compat.trim().is_empty() {
return Err(ValidationError::MissingBackendCompat);
}
if self.backend_compat != BACKEND_COMPAT {
return Err(ValidationError::UnsupportedBackendCompat {
actual: self.backend_compat.clone(),
expected: BACKEND_COMPAT,
});
}
if let Some(expiry) = self.expiry_unix_ms {
if expiry <= self.activation_unix_ms {
return Err(ValidationError::InvalidExpiry {
activation: self.activation_unix_ms,
expiry,
});
}
}
}
for (key, materialization) in &self.materializations {
if *key != materialization.fingerprint {
return Err(ValidationError::FingerprintMismatch {
Expand Down Expand Up @@ -976,6 +1036,10 @@ mod tests {
BackendPlan {
plan_id: 42,
generated_at_unix_ms: 1_735_000_000_000,
plan_version: 1,
activation_unix_ms: 1_735_000_000_000,
expiry_unix_ms: None,
backend_compat: "asap-query-backend.v1".into(),
materializations,
routing: vec![
RoutingEntry {
Expand Down
7 changes: 7 additions & 0 deletions control_plane/src/emit/backend_push.rs
Original file line number Diff line number Diff line change
Expand Up @@ -247,6 +247,9 @@ async fn push_documents_coupled(
plan_id: crate::backend_plan::BackendPlan::decode(&plan_bytes)
.map(|plan| plan.plan_id)
.unwrap_or_default(),
plan_version: crate::backend_plan::BackendPlan::decode(&plan_bytes)
.map(|plan| plan.plan_version)
.unwrap_or_default(),
entries: Default::default(),
};

Expand Down Expand Up @@ -454,7 +457,11 @@ async fn push_cumulative_entries(
let precompute_plan = PrecomputePlan {
envelope: PlanEnvelope {
plan_id,
plan_version: 1,
generated_at_unix_ms,
activation_unix_ms: generated_at_unix_ms,
expiry_unix_ms: None,
backend_compat: "asap-query-backend.v1".into(),
planner_revision: crate::physical::compiler::PLANNER_REVISION.into(),
capability_snapshot_id: "replanner".into(),
},
Expand Down
49 changes: 49 additions & 0 deletions control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -595,6 +595,10 @@ struct CompileAndPublishPhysicalPlanRequest {
evidence: HashMap<String, physical::compiler::TopKMembershipEvidence>,
planner_revision: String,
max_evidence_age_ms: u64,
plan_version: u64,
activation_unix_ms: u64,
expiry_unix_ms: Option<u64>,
backend_compat: String,
#[serde(default = "default_physical_plan_timeout_ms")]
apply_timeout_ms: u64,
}
Expand All @@ -606,6 +610,8 @@ fn default_physical_plan_timeout_ms() -> u64 {
#[derive(Debug, Serialize)]
struct CompileAndPublishPhysicalPlanResponse {
plan_id: u64,
plan_version: u64,
status: &'static str,
generated_at_unix_ms: u64,
collector_ids: Vec<String>,
}
Expand Down Expand Up @@ -667,9 +673,36 @@ async fn handle_compile_and_publish_physical_plan(
)
.into_response();
}
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 {
return (
StatusCode::GATEWAY_TIMEOUT,
"activation time exceeds apply_timeout_ms; backend remains staged".to_string(),
)
.into_response();
}
if activation_wait > 0 {
tokio::time::sleep(Duration::from_millis(activation_wait)).await;
}
if let Err(error) = backend
.activate_physical_plan(bundle.envelope.plan_id, bundle.envelope.plan_version)
.await
{
return (
StatusCode::BAD_GATEWAY,
format!("backend physical-plan activation failed: {error}"),
)
.into_response();
}

Json(CompileAndPublishPhysicalPlanResponse {
plan_id: bundle.envelope.plan_id,
plan_version: bundle.envelope.plan_version,
status: "active",
generated_at_unix_ms: bundle.envelope.generated_at_unix_ms,
collector_ids,
})
Expand All @@ -693,6 +726,18 @@ fn compile_physical_plan_request(
"max_evidence_age_ms and apply_timeout_ms must be non-zero".to_string(),
));
}
if request.plan_version == 0
|| request.activation_unix_ms == 0
|| request.backend_compat != control_plane::backend_plan::BACKEND_COMPAT
|| request
.expiry_unix_ms
.is_some_and(|expiry| expiry <= request.activation_unix_ms)
{
return Err((
StatusCode::UNPROCESSABLE_ENTITY,
"plan_version, activation_unix_ms, backend_compat and expiry are invalid".to_string(),
));
}

let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
Expand Down Expand Up @@ -748,6 +793,10 @@ fn compile_physical_plan_request(
capability_snapshot_id: request.capability_snapshot_id,
observed_at_unix_ms: now,
max_evidence_age_ms: request.max_evidence_age_ms,
plan_version: request.plan_version,
activation_unix_ms: request.activation_unix_ms,
expiry_unix_ms: request.expiry_unix_ms,
backend_compat: request.backend_compat,
},
) {
Ok(bundle) => bundle,
Expand Down
Loading
Loading