Skip to content
Draft
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

8 changes: 6 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,14 @@ The backend is being rebuilt on the Planner's physical DAG:
|---|---|
| Control plane: plan a PromQL workload, write `plan.json` | done |
| Data plane: load and validate `plan.json`, accept Remote Write | done |
| Stored summaries (SDS) | next |
| Precompute engine: run ingestion-time nodes per pane | planned |
| Stored summaries (SDS): one definition per computation, panes by end time | done |
| Precompute engine: build ingestion-time outputs once per tumbling pane | done |
| Query engine: run query-time nodes over stored panes and raw data | planned |

Panes close `--allowed-lateness-ms` (default 15 s) after the newest sample
passes their end; later samples for them are dropped. The first pane after
startup is skipped because it may be partial.

Until the query engine lands, the data plane forwards every query to
Prometheus and marks the response `x-asap-execution: forwarded`. SQL and
ClickHouse are not supported yet.
Expand Down
1 change: 1 addition & 0 deletions data_plane/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ axum.workspace = true
reqwest.workspace = true
prost = "0.13"
snap = "1.1"
futures = "0.3"

[dev-dependencies]
tower = { version = "0.5", features = ["util"] }
Expand Down
100 changes: 100 additions & 0 deletions data_plane/src/engine.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
//! The engine thread: owns the installed plan and everything that executes
//! it. Planner IR holds `Rc`s, so the plan cannot be shared across the HTTP
//! server's threads; handlers send the engine commands instead.

use std::sync::mpsc::{sync_channel, Receiver, SyncSender, TrySendError};

use asap_executor::physical_planner::frontier_from_timing;
use control_plane::InstalledPlan;
use tokio::sync::oneshot;

use crate::precompute::{Precompute, PrecomputeStats, RawSample};
use crate::sds::Sds;

/// Commands queued for the engine before ingestion pushes back.
const QUEUE: usize = 1024;

pub enum Command {
Ingest(Vec<RawSample>),
Stats(oneshot::Sender<PrecomputeStats>),
}

#[derive(Debug, thiserror::Error)]
#[error("the engine is busy")]
pub struct Busy;

/// The HTTP side of the engine thread.
#[derive(Clone)]
pub struct EngineHandle {
commands: SyncSender<Command>,
}

impl EngineHandle {
/// Starts the engine for the plan serialized in `plan_json`. Fails if the
/// plan is invalid or the data plane cannot execute it.
pub fn spawn(plan_json: String, allowed_lateness_ms: i64) -> anyhow::Result<Self> {
let (commands, receiver) = sync_channel(QUEUE);
let (ready, started) = std::sync::mpsc::channel();
std::thread::Builder::new()
.name("engine".into())
.spawn(move || match Engine::new(&plan_json, allowed_lateness_ms) {
Ok(engine) => {
let _ = ready.send(Ok(()));
engine.serve(receiver);
}
Err(error) => {
let _ = ready.send(Err(error));
}
})?;
started.recv()??;
Ok(Self { commands })
}

/// Queues `samples` for precompute, or fails if the engine is behind.
pub fn try_ingest(&self, samples: Vec<RawSample>) -> Result<(), Busy> {
match self.commands.try_send(Command::Ingest(samples)) {
Ok(()) => Ok(()),
Err(TrySendError::Full(_) | TrySendError::Disconnected(_)) => Err(Busy),
}
}

/// Precompute counters, once every command queued before has run.
pub async fn stats(&self) -> PrecomputeStats {
let (reply, stats) = oneshot::channel();
let _ = self.commands.send(Command::Stats(reply));
stats.await.unwrap_or_default()
}
}

struct Engine {
precompute: Precompute,
}

impl Engine {
fn new(plan_json: &str, allowed_lateness_ms: i64) -> anyhow::Result<Self> {
let plan: InstalledPlan = serde_json::from_str(plan_json)?;
plan.validate()?;
let dag = plan.dag.dag;
let frontier = frontier_from_timing(&dag)?;
let sds = Sds::from_dag(&dag, &frontier)?;
tracing::info!(
definitions = sds.definitions.len(),
outputs = sds.outputs.len(),
"stored summaries"
);
Ok(Self {
precompute: Precompute::new(dag, &sds, allowed_lateness_ms)?,
})
}

fn serve(mut self, commands: Receiver<Command>) {
for command in commands {
match command {
Command::Ingest(samples) => self.precompute.ingest(&samples),
Command::Stats(reply) => {
let _ = reply.send(self.precompute.stats);
}
}
}
}
}
69 changes: 57 additions & 12 deletions data_plane/src/lib.rs
Original file line number Diff line number Diff line change
@@ -1,13 +1,16 @@
//! Backend data plane: receives Prometheus Remote Write samples and serves
//! the Prometheus query API for an installed plan.
//!
//! The installed plan is ASAPPlanner's selected `PhysicalASAPDAG`. Executing
//! it is not implemented yet: samples are validated and counted, and every
//! query is forwarded to Prometheus.
//! The installed plan is ASAPPlanner's selected `PhysicalASAPDAG`. The
//! engine thread builds its ingestion-time outputs pane by pane; queries are
//! still forwarded to Prometheus.

pub mod engine;
pub mod precompute;
pub mod remote_write;
pub mod sds;

use std::collections::BTreeMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;

Expand All @@ -17,14 +20,20 @@ use axum::http::{header, HeaderMap, HeaderValue, Method, StatusCode};
use axum::response::{IntoResponse, Response};
use axum::routing::{get, post};
use axum::Router;
use control_plane::InstalledPlan;
use engine::EngineHandle;
use precompute::RawSample;

/// Prometheus' staleness marker: a NaN that ends a series, not a sample.
const STALE_NAN_BITS: u64 = 0x7ff0_0000_0000_0002;

/// Response header naming what answered a query.
pub const EXECUTION_HEADER: &str = "x-asap-execution";

pub struct DataPlane {
/// The installed plan's queries, in root order. The plan itself is not
/// `Send` (its operator payloads hold `Rc`s), so it stays with its loader.
/// The installed plan's queries, in root order.
pub queries: Vec<String>,
pub engine: EngineHandle,
/// Base URL of the Prometheus that answers forwarded queries.
pub prometheus_url: String,
pub limits: remote_write::Limits,
Expand All @@ -34,14 +43,23 @@ pub struct DataPlane {
}

impl DataPlane {
pub fn new(queries: Vec<String>, prometheus_url: String) -> Self {
Self {
queries,
/// Installs the plan serialized in `plan_json` (an [`InstalledPlan`]).
/// Panes close `allowed_lateness_ms` after the newest sample passes them.
pub fn new(
plan_json: String,
prometheus_url: String,
allowed_lateness_ms: i64,
) -> anyhow::Result<Self> {
let plan: InstalledPlan = serde_json::from_str(&plan_json)?;
plan.validate()?;
Ok(Self {
queries: plan.queries,
engine: EngineHandle::spawn(plan_json, allowed_lateness_ms)?,
prometheus_url: prometheus_url.trim_end_matches('/').to_string(),
limits: remote_write::Limits::default(),
samples: AtomicU64::new(0),
client: reqwest::Client::new(),
}
})
}
}

Expand All @@ -57,9 +75,36 @@ pub fn router(state: Arc<DataPlane>) -> Router {
async fn write(State(state): State<Arc<DataPlane>>, body: Bytes) -> Response {
match remote_write::decode(&body, &state.limits) {
Ok(request) => {
let samples: usize = request.timeseries.iter().map(|s| s.samples.len()).sum();
state.samples.fetch_add(samples as u64, Ordering::Relaxed);
StatusCode::NO_CONTENT.into_response()
let mut samples = Vec::new();
for series in request.timeseries {
let labels: Arc<BTreeMap<String, String>> = Arc::new(
series
.labels
.into_iter()
.map(|label| (label.name, label.value))
.collect(),
);
samples.extend(
series
.samples
.into_iter()
.filter(|s| s.value.to_bits() != STALE_NAN_BITS)
.map(|s| RawSample {
labels: labels.clone(),
timestamp_ms: s.timestamp,
value: s.value,
}),
);
}
let count = samples.len() as u64;
match state.engine.try_ingest(samples) {
Ok(()) => {
state.samples.fetch_add(count, Ordering::Relaxed);
StatusCode::NO_CONTENT.into_response()
}
// Remote Write senders retry on 5xx.
Err(busy) => (StatusCode::SERVICE_UNAVAILABLE, busy.to_string()).into_response(),
}
}
Err(error) => (StatusCode::BAD_REQUEST, error.to_string()).into_response(),
}
Expand Down
18 changes: 10 additions & 8 deletions data_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ use std::path::PathBuf;
use std::sync::Arc;

use clap::Parser;
use control_plane::InstalledPlan;
use data_plane::{router, DataPlane};

/// Serve Remote Write ingestion and the Prometheus query API for a plan.
Expand All @@ -16,6 +15,9 @@ struct Args {
prometheus_url: String,
#[arg(long, default_value = "0.0.0.0:8088")]
listen: String,
/// How long after the newest sample passes a pane's end it is built.
#[arg(long, default_value_t = 15_000)]
allowed_lateness_ms: i64,
}

#[tokio::main]
Expand All @@ -24,13 +26,13 @@ async fn main() -> anyhow::Result<()> {
.with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
.init();
let args = Args::parse();
let plan = InstalledPlan::load(&args.plan)?;
tracing::info!(
queries = plan.queries.len(),
nodes = plan.dag.dag.nodes.len(),
"installed plan"
);
let state = Arc::new(DataPlane::new(plan.queries, args.prometheus_url));
let plan = std::fs::read_to_string(&args.plan)?;
let state = Arc::new(DataPlane::new(
plan,
args.prometheus_url,
args.allowed_lateness_ms,
)?);
tracing::info!(queries = state.queries.len(), "installed plan");
let listener = tokio::net::TcpListener::bind(&args.listen).await?;
tracing::info!(address = %listener.local_addr()?, "listening");
axum::serve(listener, router(state)).await?;
Expand Down
Loading