diff --git a/controller/src/config/mod.rs b/controller/src/config/mod.rs index 7fc9cb95..3690254a 100644 --- a/controller/src/config/mod.rs +++ b/controller/src/config/mod.rs @@ -3,6 +3,8 @@ pub mod asapquery_backend; pub mod backend; pub mod precompute; pub mod stage_config; +pub mod stage_config_otap; +pub mod stage_config_telegraf; pub mod workloads; pub use agent::generate_agent_config; @@ -13,4 +15,168 @@ pub use stage_config::{ emit_backend_config_json, emit_backend_storage_routing, emit_backend_storage_routing_with_prometheus, emit_edge_yaml, emit_gateway_yaml, }; +pub use stage_config_otap::emit_otap_dag_yaml; +pub use stage_config_telegraf::emit_telegraf_toml; pub use workloads::WorkloadRegistry; + +use crate::stage_split::emitter::EdgeStageConfig; +use anyhow::Result; + +/// Phase ε.1.5 — which edge runtime an agent identifies as. +/// +/// Today every agent the controller has built for runs the OTel-collector +/// (`Sketchcollector`); Phase ε.1.5 adds the two new runtime variants the +/// per-runtime emitters target. The runtime is reported by the agent on +/// OpAMP `on_connect` (header `X-Agent-Runtime`); when absent (legacy +/// agents) the controller defaults to `Sketchcollector` so the existing +/// behaviour is preserved. +/// +/// Phase ε.1.5 commits the enum + emit-dispatch function. Threading the +/// runtime through OpAMP `on_connect` and into the typed L5 emit path +/// is a follow-up — the emitters can be exercised in isolation today +/// (the Phase ε.1.5 test suite does exactly that). +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum AgentRuntime { + /// Default — OTel-collector contrib build (existing behaviour). + Sketchcollector, + /// otap-dataflow Rust runtime. + Sketchotap, + /// Telegraf runtime. + Sketchtelegraf, +} + +impl Default for AgentRuntime { + fn default() -> Self { + AgentRuntime::Sketchcollector + } +} + +impl AgentRuntime { + /// Parse an `X-Agent-Runtime` header value. Recognises + /// `sketchcollector` / `sketchotap` / `sketchtelegraf` + /// (case-insensitive); any other value (including the empty string) + /// defaults to `Sketchcollector` so legacy agents keep working. + pub fn from_header(value: &str) -> Self { + match value.trim().to_lowercase().as_str() { + "sketchotap" | "otap" => AgentRuntime::Sketchotap, + "sketchtelegraf" | "telegraf" => AgentRuntime::Sketchtelegraf, + _ => AgentRuntime::Sketchcollector, + } + } +} + +/// Phase ε.1.5 — dispatch the edge emit by agent runtime. Mirrors +/// `emit_edge_yaml`'s `(cfg, opamp_endpoint) -> String` shape; the OTAP +/// and Telegraf emitters take an additional optional Prometheus URL +/// override which we pass through `prometheus_url`. +/// +/// `prometheus_url` is the Mode-3 destination override: +/// * `Sketchcollector` → ignored (the OTel-collector emitter already +/// reads `${ASAP_PROMETHEUS_OTLP_URL}` at runtime); +/// * `Sketchotap` → OTLP HTTP URL passed to `emit_otap_dag_yaml`; +/// * `Sketchtelegraf` → remote-write URL passed to `emit_telegraf_toml`. +pub fn emit_for_runtime( + runtime: AgentRuntime, + cfg: &EdgeStageConfig, + opamp_endpoint: &str, + prometheus_url: Option<&str>, +) -> Result { + match runtime { + AgentRuntime::Sketchcollector => emit_edge_yaml(cfg, opamp_endpoint), + AgentRuntime::Sketchotap => emit_otap_dag_yaml(cfg, opamp_endpoint, prometheus_url), + AgentRuntime::Sketchtelegraf => emit_telegraf_toml(cfg, prometheus_url), + } +} + +#[cfg(test)] +mod runtime_tests { + use super::*; + + #[test] + fn agent_runtime_from_header_recognises_three_values() { + assert_eq!(AgentRuntime::from_header("sketchcollector"), AgentRuntime::Sketchcollector); + assert_eq!(AgentRuntime::from_header("sketchotap"), AgentRuntime::Sketchotap); + assert_eq!(AgentRuntime::from_header("sketchtelegraf"), AgentRuntime::Sketchtelegraf); + } + + #[test] + fn agent_runtime_from_header_short_aliases() { + assert_eq!(AgentRuntime::from_header("otap"), AgentRuntime::Sketchotap); + assert_eq!(AgentRuntime::from_header("telegraf"), AgentRuntime::Sketchtelegraf); + } + + #[test] + fn agent_runtime_from_header_default_is_sketchcollector() { + assert_eq!(AgentRuntime::from_header(""), AgentRuntime::Sketchcollector); + assert_eq!(AgentRuntime::from_header("garbage"), AgentRuntime::Sketchcollector); + } + + #[test] + fn emit_for_runtime_default_matches_emit_edge_yaml() { + use crate::sketch_algebra::params::{DDSketchParams, SketchParams}; + use crate::stage_split::emitter::{EdgeSketchProcessor, ExportTarget}; + use crate::stage_split::stage_id::StageId; + use crate::sketch_algebra::params::SketchKind; + + let cfg = EdgeStageConfig { + source_metric: Some("m".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: vec![EdgeSketchProcessor { + processor_name: "ddsketchprocessor".to_string(), + sketch_kind: SketchKind::DDSketch, + sketch_params: SketchParams::DDSketch(DDSketchParams { alpha: 0.01 }), + aggregation_id: "agg0".to_string(), + }], + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + }; + + let collector = emit_for_runtime( + AgentRuntime::Sketchcollector, &cfg, "ws://ctrl/v1/opamp", None, + ).expect("collector emit ok"); + let direct = emit_edge_yaml(&cfg, "ws://ctrl/v1/opamp").expect("direct emit ok"); + assert_eq!(collector, direct, "Sketchcollector dispatch must equal emit_edge_yaml"); + } + + #[test] + fn emit_for_runtime_otap_yields_dag_yaml() { + use crate::stage_split::emitter::ExportTarget; + use crate::stage_split::stage_id::StageId; + + let cfg = EdgeStageConfig { + source_metric: Some("m".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + }; + let yaml = emit_for_runtime( + AgentRuntime::Sketchotap, &cfg, "ws://ctrl/v1/opamp", None, + ).expect("otap emit ok"); + // OTAP-specific token. + assert!(yaml.contains("otel_dataflow/v1"), "expected OTAP DAG version\n{yaml}"); + } + + #[test] + fn emit_for_runtime_telegraf_yields_toml() { + use crate::stage_split::emitter::ExportTarget; + use crate::stage_split::stage_id::StageId; + + let cfg = EdgeStageConfig { + source_metric: Some("m".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + }; + let toml = emit_for_runtime( + AgentRuntime::Sketchtelegraf, &cfg, "ws://ctrl/v1/opamp", None, + ).expect("telegraf emit ok"); + // Telegraf-specific token. + assert!(toml.contains("[[inputs.opentelemetry]]"), "expected Telegraf TOML header\n{toml}"); + } +} diff --git a/controller/src/config/stage_config_otap.rs b/controller/src/config/stage_config_otap.rs new file mode 100644 index 00000000..fbaf811b --- /dev/null +++ b/controller/src/config/stage_config_otap.rs @@ -0,0 +1,516 @@ +//! Phase ε.1.5 — OTAP Dataflow DAG YAML emitter (per-runtime mirror of +//! [`super::stage_config::emit_edge_yaml`]). +//! +//! `sketchotap` uses the otap-dataflow Rust runtime; its config surface is +//! a DAG YAML where `nodes..type` is a registered plugin URN +//! (e.g. `receiver:otlp`, `exporter:otlp_http`, +//! `urn:otel:exporter:otlp_http`). The OTLP HTTP exporter ships in +//! `otel-arrow/rust/otap-dataflow/crates/core-nodes/src/exporters/otlp_http_exporter/` +//! and registers under the URN `urn:otel:exporter:otlp_http` via +//! `linkme`'s `distributed_slice(OTAP_EXPORTER_FACTORIES)`. +//! +//! Three modes (mirror the [`crate::planner::wire_cost::BindMode`] enum +//! Phase ε.1 introduced): +//! +//! 1. `SketchAtEdge` — DAG includes a sketch processor node between the +//! OTLP receiver and the OTLP gRPC exporter to the gateway. (The +//! `asap_sketches` plugin lives in the otap-patch tree; we wire its +//! type URN here without depending on its source.) +//! 2. `RawAtEdgeSketchAtBackend` — passthrough DAG: receiver → exporter. +//! No sketch processor. Egress is OTLP gRPC to the gateway, which +//! builds sketches at backend ingest. +//! 3. `RawAtEdgePrometheusArchive` — passthrough DAG: receiver → OTLP +//! HTTP exporter pointed at Prometheus's native OTLP receiver +//! (`/api/v1/otlp/v1/metrics`). +//! +//! The function consumes the same [`EdgeStageConfig`] the OTel-collector +//! emitter does — Phase ε.1.5 keeps the typed L5 plan as the single +//! source of truth across all three runtimes. The emitter dispatches per +//! `EdgeStageConfig` via the same `prometheus_archive_metrics` / +//! `sketch_processors` signals the OTel-collector emitter uses (Mode 3 +//! ↔ `prometheus_archive_metrics` non-empty; Mode 1 ↔ `sketch_processors` +//! non-empty; Mode 2 ↔ both empty + the bind decision lives in the +//! upstream stage_split — the edge YAML for Mode 2 is identical to a +//! plain raw passthrough at this layer). + +use anyhow::{Context, Result}; +use serde::Serialize; +use serde_yaml::{Mapping, Value}; +use std::collections::BTreeMap; + +use crate::sketch_algebra::params::{SketchKind, SketchParams}; +use crate::stage_split::emitter::{EdgeStageConfig, EdgeSketchProcessor, ExportTarget}; +use crate::stage_split::stage_id::StageId; + +/// Default URL for Prometheus's native OTLP HTTP receiver. +/// Matches `super::stage_config::emit_edge_yaml`'s placeholder so the +/// three runtime emitters agree on the wire endpoint. +pub const DEFAULT_PROMETHEUS_OTLP_URL: &str = + "http://prometheus:9090/api/v1/otlp/v1/metrics"; + +/// URN of the OTLP HTTP exporter registered by +/// `otel-arrow/rust/otap-dataflow/crates/core-nodes/src/exporters/otlp_http_exporter/`. +const URN_OTLP_HTTP_EXPORTER: &str = "exporter:otlp_http"; + +/// URN of the OTLP gRPC exporter registered by +/// `otel-arrow/rust/otap-dataflow/crates/core-nodes/src/exporters/otlp_grpc_exporter/`. +const URN_OTLP_GRPC_EXPORTER: &str = "exporter:otlp_grpc"; + +/// URN of the OTLP receiver (gRPC + HTTP). +const URN_OTLP_RECEIVER: &str = "receiver:otlp"; + +/// URN of the asap_sketches processor registered in the otap-patch tree. +/// Phase ε.1.5 wires the URN abstractly; the binary side lands the +/// plugin in `otap-patch/plugins/asap_sketches/` per +/// [`docs/design-asap-otap-rust-integration.md`]. +const URN_ASAP_SKETCHES_PROCESSOR: &str = "processor:asap_sketches"; + +// ── DAG YAML structural types ──────────────────────────────────────────────── +// +// These mirror the otap-dataflow `engine`/`groups`/`pipelines`/`nodes` +// schema in `otel-arrow/rust/otap-dataflow/configs/*.yaml`. We keep a +// minimal set to round-trip; the full schema (channel_capacity policies, +// engine settings, etc.) is left at struct defaults — Phase ε.1.5 only +// commits the wiring shape, not policy. + +#[derive(Serialize)] +struct OtapDag { + version: String, + engine: BTreeMap, + groups: BTreeMap, +} + +#[derive(Serialize)] +struct Group { + pipelines: BTreeMap, +} + +#[derive(Serialize)] +struct PipelineDef { + nodes: BTreeMap, + connections: Vec, +} + +#[derive(Serialize)] +struct NodeDef { + #[serde(rename = "type")] + kind: String, + config: Value, +} + +#[derive(Serialize)] +struct Connection { + from: String, + to: String, +} + +// ── Public API ─────────────────────────────────────────────────────────────── + +/// Build the OTAP-Dataflow DAG YAML for the `sketchotap` runtime from a +/// typed L5 [`EdgeStageConfig`]. +/// +/// `opamp_endpoint` is the controller's WebSocket URL; reserved for a +/// future `extension:opamp` node when the otap-dataflow runtime grows +/// OpAMP support (today the otap-dataflow `engine` block has no +/// extension model, so we accept the param for shape-parity with +/// [`super::stage_config::emit_edge_yaml`] and ignore it). +/// +/// `prometheus_otlp_url` overrides the default Prometheus OTLP HTTP +/// endpoint when present (for Mode 3 metrics). `None` falls back to +/// [`DEFAULT_PROMETHEUS_OTLP_URL`]. +pub fn emit_otap_dag_yaml( + cfg: &EdgeStageConfig, + _opamp_endpoint: &str, + prometheus_otlp_url: Option<&str>, +) -> Result { + let mut nodes: BTreeMap = BTreeMap::new(); + let mut connections: Vec = Vec::new(); + + // ── OTLP receiver ──────────────────────────────────────────────────────── + // Both gRPC + HTTP listeners — matches `emit_edge_yaml`'s shape so + // the three runtimes accept the same upstream traffic. + let receiver_cfg: Value = serde_yaml::from_str( + "protocols:\n grpc:\n listening_addr: \"0.0.0.0:4317\"\n http:\n listening_addr: \"0.0.0.0:4318\"\n", + ) + .context("parse OTAP otlp receiver block")?; + nodes.insert( + "receiver".to_string(), + NodeDef { + kind: URN_OTLP_RECEIVER.to_string(), + config: receiver_cfg, + }, + ); + + // ── Mode dispatch ──────────────────────────────────────────────────────── + let has_prometheus_archive = !cfg.prometheus_archive_metrics.is_empty(); + let has_sketch = !cfg.sketch_processors.is_empty(); + + if has_prometheus_archive { + // Mode 3 — Prometheus archive: passthrough → otlp_http exporter + // pointed at Prometheus's native OTLP receiver. We do NOT also + // emit a sketch node; Mode 3 metrics are the whole edge stream + // for that pipeline. (When a single agent is hosting Mode 3 + + // Mode 1/2 metrics simultaneously, the upstream typed splitter + // produces two `EdgeStageConfig`s — one per mode bucket — and + // we emit two pipelines side-by-side via the otap-dataflow + // multi-pipeline `pipelines:` map. Phase ε.1.5 ships the + // single-pipeline case; the multi-pipeline case is an upstream + // splitter concern.) + let prom_url = + prometheus_otlp_url.unwrap_or(DEFAULT_PROMETHEUS_OTLP_URL); + let exp_cfg = build_otlp_http_exporter_config(prom_url); + nodes.insert( + "exporter".to_string(), + NodeDef { + kind: URN_OTLP_HTTP_EXPORTER.to_string(), + config: exp_cfg, + }, + ); + connections.push(Connection { + from: "receiver".to_string(), + to: "exporter".to_string(), + }); + } else if has_sketch { + // Mode 1 — sketch at edge. Insert one processor per + // `EdgeSketchProcessor`; chain them serially between receiver + // and the gateway-bound OTLP gRPC exporter. + let mut prev = "receiver".to_string(); + for (i, sp) in cfg.sketch_processors.iter().enumerate() { + let name = format!("sketch_{i}"); + nodes.insert( + name.clone(), + NodeDef { + kind: URN_ASAP_SKETCHES_PROCESSOR.to_string(), + config: build_asap_sketches_config(sp, cfg.window_secs), + }, + ); + connections.push(Connection { + from: prev.clone(), + to: name.clone(), + }); + prev = name; + } + let endpoint = resolve_export_endpoint("gateway", &cfg.exporter_target); + nodes.insert( + "exporter".to_string(), + NodeDef { + kind: URN_OTLP_GRPC_EXPORTER.to_string(), + config: build_otlp_grpc_exporter_config(&endpoint), + }, + ); + connections.push(Connection { + from: prev, + to: "exporter".to_string(), + }); + } else { + // Mode 2 — raw at edge → sketch at backend. Passthrough DAG. + let endpoint = resolve_export_endpoint("gateway", &cfg.exporter_target); + nodes.insert( + "exporter".to_string(), + NodeDef { + kind: URN_OTLP_GRPC_EXPORTER.to_string(), + config: build_otlp_grpc_exporter_config(&endpoint), + }, + ); + connections.push(Connection { + from: "receiver".to_string(), + to: "exporter".to_string(), + }); + } + + let mut pipelines = BTreeMap::new(); + pipelines.insert( + "main".to_string(), + PipelineDef { nodes, connections }, + ); + + let mut groups = BTreeMap::new(); + groups.insert("default".to_string(), Group { pipelines }); + + let dag = OtapDag { + version: "otel_dataflow/v1".to_string(), + engine: BTreeMap::new(), + groups, + }; + + serde_yaml::to_string(&dag).context("serialize OTAP DAG YAML") +} + +// ── Internals ──────────────────────────────────────────────────────────────── + +fn resolve_export_endpoint(default_host: &str, target: &ExportTarget) -> String { + match target { + ExportTarget::Endpoint(s) => s.clone(), + ExportTarget::Stage(StageId::Edge) => "edge:4317".to_string(), + ExportTarget::Stage(StageId::Gateway) => format!("{default_host}:4317"), + ExportTarget::Stage(StageId::Backend) => format!("{default_host}:4317"), + } +} + +/// Build the OTLP HTTP exporter `config:` block. Matches the +/// `crates/core-nodes/src/exporters/otlp_http_exporter/config.rs` schema +/// — `endpoint` (base URL) plus an explicit `metrics_endpoint` so the +/// Prometheus path `/api/v1/otlp/v1/metrics` round-trips verbatim. +fn build_otlp_http_exporter_config(metrics_url: &str) -> Value { + // Derive the bare endpoint from the metrics URL: drop the path. For + // typical inputs this is `http://prometheus:9090`. + let base = match metrics_url.find("/api/") { + Some(i) => &metrics_url[..i], + None => metrics_url, + }; + let yaml = format!( + "endpoint: \"{base}\"\nmetrics_endpoint: \"{metrics_url}\"\nhttp:\n request_timeout: \"30s\"\nclient_pool_size: 1\n", + ); + serde_yaml::from_str(&yaml).expect("inline OTLP HTTP exporter config is valid YAML") +} + +/// Build the OTLP gRPC exporter `config:` block. The otap-dataflow +/// `otlp_grpc` exporter uses `grpc_endpoint` as the field name (see +/// `configs/otlp-otlp.yaml`). +fn build_otlp_grpc_exporter_config(endpoint: &str) -> Value { + let url = if endpoint.starts_with("http://") || endpoint.starts_with("https://") { + endpoint.to_string() + } else { + format!("http://{endpoint}") + }; + let yaml = format!("grpc_endpoint: \"{url}\"\ntimeout: \"15s\"\n"); + serde_yaml::from_str(&yaml).expect("inline OTLP gRPC exporter config is valid YAML") +} + +/// Build the per-edge-processor `asap_sketches` config block. Mirrors +/// the same fields the OTel-collector emitter writes +/// (`super::stage_config::build_edge_processor_block`) so the binary +/// side can share a single schema across the OTel + OTAP runtimes. +fn build_asap_sketches_config(sp: &EdgeSketchProcessor, window_secs: Option) -> Value { + let mut m = Mapping::new(); + if let Some(w) = window_secs { + m.insert("mode".into(), Value::String("window".to_string())); + m.insert("window_duration".into(), Value::String(format!("{w}s"))); + } else { + m.insert("mode".into(), Value::String("batch".to_string())); + } + m.insert( + "aggregation_id".into(), + Value::String(sp.aggregation_id.clone()), + ); + m.insert("sketch_kind".into(), Value::String(sketch_kind_tag(&sp.sketch_kind).into())); + match &sp.sketch_params { + SketchParams::Kll(p) => { + m.insert("k".into(), Value::Number((p.k as u64).into())); + } + SketchParams::DDSketch(p) => { + m.insert("relative_accuracy".into(), Value::Number(p.alpha.into())); + m.insert("delta_transmission".into(), Value::Bool(true)); + } + SketchParams::Hll(_p) => { + m.insert("delta_transmission".into(), Value::Bool(true)); + } + SketchParams::Cms(p) => { + m.insert("rows".into(), Value::Number((p.d as u64).into())); + m.insert("columns".into(), Value::Number((p.w as u64).into())); + m.insert("delta_transmission".into(), Value::Bool(true)); + } + SketchParams::CountSketch(p) => { + let epsilon = std::f64::consts::E / (p.w as f64); + let delta = 2f64.powi(-(p.d as i32)); + m.insert("epsilon".into(), Value::Number(epsilon.into())); + m.insert("delta".into(), Value::Number(delta.into())); + m.insert("delta_transmission".into(), Value::Bool(true)); + } + } + Value::Mapping(m) +} + +fn sketch_kind_tag(kind: &SketchKind) -> &'static str { + match kind { + SketchKind::Kll => "kll", + SketchKind::DDSketch => "ddsketch", + SketchKind::Hll => "hll", + SketchKind::Cms => "cms", + SketchKind::CountSketch => "count_sketch", + } +} + +// ── Tests ──────────────────────────────────────────────────────────────────── + +#[cfg(test)] +mod tests { + use super::*; + use crate::sketch_algebra::params::DDSketchParams; + use crate::stage_split::emitter::{EdgeSketchProcessor, PrometheusArchiveMetric}; + + /// Minimal struct-stub used to validate the emitted DAG parses as the + /// otap-dataflow schema. We don't pull in the otap-df-config crate + /// here (it would add an enormous dependency footprint to the + /// controller); instead we verify the top-level shape (`version`, + /// `groups`, `pipelines`, `nodes`, `connections`) round-trips. + #[derive(Debug, serde::Deserialize)] + struct OtapDagStub { + version: String, + #[allow(dead_code)] + engine: serde_yaml::Value, + groups: BTreeMap, + } + + #[derive(Debug, serde::Deserialize)] + struct GroupStub { + pipelines: BTreeMap, + } + + #[derive(Debug, serde::Deserialize)] + struct PipelineStub { + nodes: BTreeMap, + connections: Vec, + } + + #[derive(Debug, serde::Deserialize)] + struct NodeStub { + #[serde(rename = "type")] + kind: String, + #[allow(dead_code)] + config: serde_yaml::Value, + } + + #[derive(Debug, serde::Deserialize)] + struct ConnectionStub { + from: String, + to: String, + } + + fn ddsketch_edge_cfg_mode1() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: vec![EdgeSketchProcessor { + processor_name: "ddsketchprocessor".to_string(), + sketch_kind: SketchKind::DDSketch, + sketch_params: SketchParams::DDSketch(DDSketchParams { alpha: 0.01 }), + aggregation_id: "agg0".to_string(), + }], + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + } + } + + fn raw_edge_cfg_mode2() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + } + } + + fn prom_edge_cfg_mode3() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: vec![PrometheusArchiveMetric { + metric: "http_request_duration_seconds".to_string(), + window_secs: Some(60), + label_proj: vec!["service.name".to_string()], + }], + } + } + + /// Mode 1 snapshot — sketch at edge: receiver → asap_sketches → + /// otlp_grpc exporter to gateway. + #[test] + fn otap_dag_mode1_sketch_at_edge_shape() { + let yaml = + emit_otap_dag_yaml(&ddsketch_edge_cfg_mode1(), "ws://ctrl/v1/opamp", None) + .expect("emit_otap_dag_yaml ok"); + let dag: OtapDagStub = serde_yaml::from_str(&yaml).expect("DAG parses"); + assert_eq!(dag.version, "otel_dataflow/v1"); + let pipe = dag.groups.get("default").unwrap().pipelines.get("main").unwrap(); + // Receiver + sketch + exporter == 3 nodes. + assert_eq!(pipe.nodes.len(), 3, "expected 3 nodes\n{yaml}"); + assert_eq!(pipe.nodes.get("receiver").unwrap().kind, URN_OTLP_RECEIVER); + assert_eq!( + pipe.nodes.get("sketch_0").unwrap().kind, + URN_ASAP_SKETCHES_PROCESSOR + ); + assert_eq!(pipe.nodes.get("exporter").unwrap().kind, URN_OTLP_GRPC_EXPORTER); + // Connections: receiver → sketch_0 → exporter. + assert_eq!(pipe.connections.len(), 2); + assert_eq!(pipe.connections[0].from, "receiver"); + assert_eq!(pipe.connections[0].to, "sketch_0"); + assert_eq!(pipe.connections[1].from, "sketch_0"); + assert_eq!(pipe.connections[1].to, "exporter"); + // Endpoint contains gateway:4317. + assert!(yaml.contains("gateway:4317"), "missing gateway endpoint\n{yaml}"); + } + + /// Mode 2 snapshot — raw at edge: receiver → otlp_grpc exporter. + /// No sketch node; the gateway / backend will build sketches. + #[test] + fn otap_dag_mode2_raw_at_edge_shape() { + let yaml = emit_otap_dag_yaml(&raw_edge_cfg_mode2(), "ws://ctrl/v1/opamp", None) + .expect("emit_otap_dag_yaml ok"); + let dag: OtapDagStub = serde_yaml::from_str(&yaml).expect("DAG parses"); + let pipe = dag.groups.get("default").unwrap().pipelines.get("main").unwrap(); + assert_eq!(pipe.nodes.len(), 2, "expected receiver + exporter only\n{yaml}"); + assert_eq!(pipe.nodes.get("receiver").unwrap().kind, URN_OTLP_RECEIVER); + assert_eq!(pipe.nodes.get("exporter").unwrap().kind, URN_OTLP_GRPC_EXPORTER); + // Direct connection. + assert_eq!(pipe.connections.len(), 1); + assert_eq!(pipe.connections[0].from, "receiver"); + assert_eq!(pipe.connections[0].to, "exporter"); + // No sketch processor in YAML. + assert!( + !yaml.contains(URN_ASAP_SKETCHES_PROCESSOR), + "Mode 2 must not include a sketch processor\n{yaml}" + ); + } + + /// Mode 3 snapshot — Prometheus archive: receiver → otlp_http + /// exporter pointed at `/api/v1/otlp/v1/metrics`. + #[test] + fn otap_dag_mode3_prometheus_archive_shape() { + let yaml = emit_otap_dag_yaml(&prom_edge_cfg_mode3(), "ws://ctrl/v1/opamp", None) + .expect("emit_otap_dag_yaml ok"); + let dag: OtapDagStub = serde_yaml::from_str(&yaml).expect("DAG parses"); + let pipe = dag.groups.get("default").unwrap().pipelines.get("main").unwrap(); + assert_eq!(pipe.nodes.len(), 2); + assert_eq!(pipe.nodes.get("exporter").unwrap().kind, URN_OTLP_HTTP_EXPORTER); + // Path round-trips verbatim. + assert!( + yaml.contains("/api/v1/otlp/v1/metrics"), + "missing Prometheus OTLP path\n{yaml}" + ); + // No sketch processor. + assert!( + !yaml.contains(URN_ASAP_SKETCHES_PROCESSOR), + "Mode 3 must not include a sketch processor\n{yaml}" + ); + } + + /// Mode 3 with override URL — caller can redirect to a non-default + /// Prometheus instance (`https://prom-prod:9090/...`). + #[test] + fn otap_dag_mode3_url_override() { + let yaml = emit_otap_dag_yaml( + &prom_edge_cfg_mode3(), + "ws://ctrl/v1/opamp", + Some("https://prom-prod:9090/api/v1/otlp/v1/metrics"), + ) + .expect("emit_otap_dag_yaml ok"); + assert!( + yaml.contains("https://prom-prod:9090/api/v1/otlp/v1/metrics"), + "override URL not propagated\n{yaml}" + ); + // The base endpoint should drop the path. (serde_yaml elides + // quotes around scalar strings that don't need them, so we + // match the unquoted form.) + assert!( + yaml.contains("endpoint: https://prom-prod:9090\n"), + "base endpoint not derived\n{yaml}" + ); + } +} diff --git a/controller/src/config/stage_config_telegraf.rs b/controller/src/config/stage_config_telegraf.rs new file mode 100644 index 00000000..0f1b3505 --- /dev/null +++ b/controller/src/config/stage_config_telegraf.rs @@ -0,0 +1,481 @@ +//! Phase ε.1.5 — Telegraf TOML emitter (per-runtime mirror of +//! [`super::stage_config::emit_edge_yaml`]). +//! +//! `sketchtelegraf` is the Telegraf-runtime variant of the ASAP edge +//! agent. Telegraf consumes TOML; the relevant plugins are: +//! +//! * `inputs.opentelemetry` — OTLP gRPC / HTTP receiver (port 4317 / +//! 4318) — provides the upstream OTLP stream. +//! * `processors.allsketches` — the Telegraf-side streaming sketch +//! processor patched into `telegraf-patch/processors/`. Mirror of the +//! OTel-collector `*sketchprocessor` family. +//! * `outputs.opentelemetry` — OTLP **gRPC** exporter to the gateway +//! (Mode 1 / Mode 2). Telegraf's shipped plugin is gRPC-only — see +//! `telegraf/plugins/outputs/opentelemetry/opentelemetry.go`. There is +//! no `protocol = "http/protobuf"` field in the upstream plugin. +//! * `outputs.http` — generic HTTP POST output. For Mode 3 we use this +//! with the `prometheusremotewrite` serializer so the agent can ship +//! raw samples to Prometheus's remote-write endpoint +//! (`/api/v1/write`). This is a documented Phase ε.1.5 deviation +//! from the OTel-collector path's `otlphttp/prometheus` exporter: +//! Telegraf has no OTLP-HTTP serializer, but Prometheus's +//! remote-write endpoint accepts the same physical archive that the +//! OTLP receiver writes to, so the **archive contents end up +//! identical**. The wire framing differs; the storage outcome does +//! not. +//! +//! Three modes (mirror the [`crate::planner::wire_cost::BindMode`] enum +//! Phase ε.1 introduced): +//! +//! 1. `SketchAtEdge` — `[[processors.allsketches]]` between +//! `[[inputs.opentelemetry]]` and `[[outputs.opentelemetry]]`. +//! 2. `RawAtEdgeSketchAtBackend` — passthrough: input → output. No +//! sketch processor at the edge. +//! 3. `RawAtEdgePrometheusArchive` — passthrough: input → +//! `[[outputs.http]]` with `data_format = "prometheusremotewrite"` +//! pointed at Prometheus's `/api/v1/write`. + +use anyhow::{Context, Result}; + +use crate::sketch_algebra::params::{SketchKind, SketchParams}; +use crate::stage_split::emitter::{EdgeStageConfig, EdgeSketchProcessor, ExportTarget}; +use crate::stage_split::stage_id::StageId; + +/// Default Prometheus remote-write URL for Mode 3 — Telegraf doesn't +/// support OTLP-HTTP egress, so we land in the same Prometheus archive +/// via remote-write instead. The URL maps to the same Prometheus instance +/// the OTel-collector emitter targets via OTLP HTTP — Prometheus accepts +/// both ingest paths and stores into the same TSDB. +pub const DEFAULT_PROMETHEUS_REMOTE_WRITE_URL: &str = + "http://prometheus:9090/api/v1/write"; + +// ── Public API ─────────────────────────────────────────────────────────────── + +/// Build the Telegraf TOML for the `sketchtelegraf` runtime from a +/// typed L5 [`EdgeStageConfig`]. +/// +/// `prometheus_remote_write_url` overrides the default Prometheus +/// remote-write endpoint when present (for Mode 3 metrics). `None` +/// falls back to [`DEFAULT_PROMETHEUS_REMOTE_WRITE_URL`]. +pub fn emit_telegraf_toml( + cfg: &EdgeStageConfig, + prometheus_remote_write_url: Option<&str>, +) -> Result { + let mut out = String::new(); + out.push_str("# Generated by ASAP controller (Phase ε.1.5).\n"); + out.push_str("# Mode-2/3 detection follows the EdgeStageConfig signals.\n\n"); + + // ── inputs.opentelemetry ───────────────────────────────────────────────── + // Telegraf's OTLP input listens on the default OTLP ports. + out.push_str("[[inputs.opentelemetry]]\n"); + out.push_str(" service_address = \"0.0.0.0:4317\"\n"); + out.push_str(" http_service_address = \"0.0.0.0:4318\"\n"); + out.push_str(" timeout = \"5s\"\n"); + out.push('\n'); + + // ── Mode dispatch ──────────────────────────────────────────────────────── + let has_prometheus_archive = !cfg.prometheus_archive_metrics.is_empty(); + let has_sketch = !cfg.sketch_processors.is_empty(); + + if has_prometheus_archive { + // Mode 3 — Prometheus archive: passthrough, then remote-write + // to Prometheus. We do NOT include `[[processors.allsketches]]`. + let url = + prometheus_remote_write_url.unwrap_or(DEFAULT_PROMETHEUS_REMOTE_WRITE_URL); + emit_outputs_http_remote_write(&mut out, url); + } else if has_sketch { + // Mode 1 — sketch at edge. One `[[processors.allsketches]]` per + // sketch processor; outputs.opentelemetry to gateway. + for sp in &cfg.sketch_processors { + emit_processors_allsketches(&mut out, sp, cfg.window_secs); + } + let endpoint = resolve_export_endpoint("gateway", &cfg.exporter_target); + emit_outputs_opentelemetry(&mut out, &endpoint); + } else { + // Mode 2 — raw at edge. Passthrough; outputs.opentelemetry + // ships raw OTLP to the gateway. + let endpoint = resolve_export_endpoint("gateway", &cfg.exporter_target); + emit_outputs_opentelemetry(&mut out, &endpoint); + } + + // Validate that what we emitted parses as TOML — catches malformed + // table headers / quoted strings before the agent boots. + let _: toml_minimal::Document = toml_minimal::Document::parse(&out) + .with_context(|| format!("generated Telegraf TOML failed minimal parse: {out}"))?; + + Ok(out) +} + +// ── Internals ──────────────────────────────────────────────────────────────── + +fn resolve_export_endpoint(default_host: &str, target: &ExportTarget) -> String { + match target { + ExportTarget::Endpoint(s) => s.clone(), + ExportTarget::Stage(StageId::Edge) => "edge:4317".to_string(), + ExportTarget::Stage(StageId::Gateway) => format!("{default_host}:4317"), + ExportTarget::Stage(StageId::Backend) => format!("{default_host}:4317"), + } +} + +fn emit_outputs_opentelemetry(out: &mut String, endpoint: &str) { + // `outputs.opentelemetry` is gRPC-only; `service_address` takes + // `host:port` (no scheme). Compression defaults to gzip. + out.push_str("[[outputs.opentelemetry]]\n"); + out.push_str(&format!(" service_address = \"{endpoint}\"\n")); + out.push_str(" timeout = \"5s\"\n"); + out.push_str(" compression = \"gzip\"\n"); + out.push('\n'); +} + +fn emit_outputs_http_remote_write(out: &mut String, url: &str) { + // Telegraf's `outputs.http` POSTs the serialized batch to the URL. + // `data_format = "prometheusremotewrite"` selects the remote-write + // serializer shipped in `telegraf/plugins/serializers/prometheusremotewrite/`. + out.push_str("[[outputs.http]]\n"); + out.push_str(&format!(" url = \"{url}\"\n")); + out.push_str(" method = \"POST\"\n"); + out.push_str(" data_format = \"prometheusremotewrite\"\n"); + out.push_str(" content_encoding = \"snappy\"\n"); + out.push_str(" [outputs.http.headers]\n"); + out.push_str(" Content-Type = \"application/x-protobuf\"\n"); + out.push_str(" X-Prometheus-Remote-Write-Version = \"0.1.0\"\n"); + out.push('\n'); +} + +fn emit_processors_allsketches( + out: &mut String, + sp: &EdgeSketchProcessor, + window_secs: Option, +) { + out.push_str("[[processors.allsketches]]\n"); + let mode = if window_secs.is_some() { "window" } else { "batch" }; + out.push_str(&format!(" mode = \"{mode}\"\n")); + if let Some(w) = window_secs { + out.push_str(&format!(" window_duration = \"{w}s\"\n")); + } + out.push_str(&format!( + " aggregation_id = \"{}\"\n", + sp.aggregation_id + )); + out.push_str(&format!( + " sketch_kind = \"{}\"\n", + sketch_kind_tag(&sp.sketch_kind) + )); + match &sp.sketch_params { + SketchParams::Kll(p) => { + out.push_str(&format!(" k = {}\n", p.k)); + } + SketchParams::DDSketch(p) => { + out.push_str(&format!(" relative_accuracy = {}\n", p.alpha)); + out.push_str(" delta_transmission = true\n"); + } + SketchParams::Hll(_p) => { + out.push_str(" delta_transmission = true\n"); + } + SketchParams::Cms(p) => { + out.push_str(&format!(" rows = {}\n", p.d)); + out.push_str(&format!(" columns = {}\n", p.w)); + out.push_str(" delta_transmission = true\n"); + } + SketchParams::CountSketch(p) => { + let epsilon = std::f64::consts::E / (p.w as f64); + let delta = 2f64.powi(-(p.d as i32)); + out.push_str(&format!(" epsilon = {epsilon}\n")); + out.push_str(&format!(" delta = {delta}\n")); + out.push_str(" delta_transmission = true\n"); + } + } + out.push('\n'); +} + +fn sketch_kind_tag(kind: &SketchKind) -> &'static str { + match kind { + SketchKind::Kll => "kll", + SketchKind::DDSketch => "ddsketch", + SketchKind::Hll => "hll", + SketchKind::Cms => "cms", + SketchKind::CountSketch => "count_sketch", + } +} + +// ── Minimal TOML parser stub ───────────────────────────────────────────────── +// +// The controller crate doesn't depend on a full `toml` crate (the +// existing surface is YAML / JSON only). To validate that our emitted +// TOML is syntactically valid, we ship a tiny purpose-built validator +// that recognizes the subset Telegraf consumes: `[[table.array]]` / +// `[table]` headers, `key = "string"`, `key = number`, `key = bool`, +// indented sub-tables, and `# comments`. This is conservative — it +// rejects malformed table headers / unbalanced quotes / etc., which +// is the failure mode we care about catching before agent boot. + +mod toml_minimal { + use anyhow::{anyhow, Result}; + + #[derive(Debug)] + pub struct Document; + + impl Document { + pub fn parse(src: &str) -> Result { + for (lineno, raw) in src.lines().enumerate() { + let line = strip_comment(raw).trim(); + if line.is_empty() { + continue; + } + if line.starts_with('[') { + if !is_balanced_brackets(line) { + return Err(anyhow!( + "line {}: unbalanced [ in table header: {raw}", + lineno + 1 + )); + } + continue; + } + // key = value + let Some(eq) = line.find('=') else { + return Err(anyhow!( + "line {}: expected `key = value`, got: {raw}", + lineno + 1 + )); + }; + let key = line[..eq].trim(); + let val = line[eq + 1..].trim(); + if key.is_empty() { + return Err(anyhow!("line {}: empty key in: {raw}", lineno + 1)); + } + if val.is_empty() { + return Err(anyhow!("line {}: empty value in: {raw}", lineno + 1)); + } + // Validate value: string (balanced quotes), bool, or + // unquoted scalar (number / fraction). + if val.starts_with('"') { + if !val.ends_with('"') || val.len() < 2 { + return Err(anyhow!( + "line {}: unbalanced \" in: {raw}", + lineno + 1 + )); + } + } + } + Ok(Document) + } + } + + fn strip_comment(line: &str) -> &str { + // Strip a trailing `# ...` comment, but not when inside quotes. + let mut in_str = false; + for (i, c) in line.char_indices() { + match c { + '"' => in_str = !in_str, + '#' if !in_str => return &line[..i], + _ => {} + } + } + line + } + + fn is_balanced_brackets(s: &str) -> bool { + let mut depth = 0i32; + for c in s.chars() { + match c { + '[' => depth += 1, + ']' => depth -= 1, + _ => {} + } + if depth < 0 { + return false; + } + } + depth == 0 + } +} + +// ── Tests ──────────────────────────────────────────────────────────────────── + +#[cfg(test)] +mod tests { + use super::*; + use crate::sketch_algebra::params::{DDSketchParams, KllParams}; + use crate::stage_split::emitter::{EdgeSketchProcessor, PrometheusArchiveMetric}; + + fn ddsketch_edge_cfg_mode1() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: vec![EdgeSketchProcessor { + processor_name: "ddsketchprocessor".to_string(), + sketch_kind: SketchKind::DDSketch, + sketch_params: SketchParams::DDSketch(DDSketchParams { alpha: 0.01 }), + aggregation_id: "agg0".to_string(), + }], + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + } + } + + fn raw_edge_cfg_mode2() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + } + } + + fn prom_edge_cfg_mode3() -> EdgeStageConfig { + EdgeStageConfig { + source_metric: Some("http_request_duration_seconds".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: Vec::new(), + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: vec![PrometheusArchiveMetric { + metric: "http_request_duration_seconds".to_string(), + window_secs: Some(60), + label_proj: vec!["service.name".to_string()], + }], + } + } + + /// Mode 1 snapshot — sketch at edge. + /// `[[inputs.opentelemetry]]` + `[[processors.allsketches]]` + + /// `[[outputs.opentelemetry]]`. + #[test] + fn telegraf_toml_mode1_sketch_at_edge_shape() { + let toml = emit_telegraf_toml(&ddsketch_edge_cfg_mode1(), None) + .expect("emit_telegraf_toml ok"); + assert!(toml.contains("[[inputs.opentelemetry]]"), "missing input\n{toml}"); + assert!( + toml.contains("[[processors.allsketches]]"), + "missing sketch processor\n{toml}" + ); + assert!( + toml.contains("[[outputs.opentelemetry]]"), + "missing output\n{toml}" + ); + assert!( + toml.contains("service_address = \"gateway:4317\""), + "missing gateway endpoint\n{toml}" + ); + // Sketch params preserved. + assert!(toml.contains("relative_accuracy = 0.01"), "missing alpha\n{toml}"); + assert!(toml.contains("sketch_kind = \"ddsketch\""), "wrong kind\n{toml}"); + } + + /// Mode 2 snapshot — raw at edge. No `[[processors.allsketches]]`. + /// `[[outputs.opentelemetry]]` ships raw OTLP to the gateway. + #[test] + fn telegraf_toml_mode2_raw_at_edge_shape() { + let toml = emit_telegraf_toml(&raw_edge_cfg_mode2(), None) + .expect("emit_telegraf_toml ok"); + assert!(toml.contains("[[inputs.opentelemetry]]"), "missing input\n{toml}"); + assert!( + !toml.contains("[[processors.allsketches]]"), + "Mode 2 must not include sketch processor\n{toml}" + ); + assert!( + toml.contains("[[outputs.opentelemetry]]"), + "missing output\n{toml}" + ); + assert!( + toml.contains("service_address = \"gateway:4317\""), + "missing gateway endpoint\n{toml}" + ); + } + + /// Mode 3 snapshot — Prometheus archive. `[[outputs.http]]` POSTs + /// to Prometheus's remote-write endpoint `/api/v1/write` (Telegraf + /// has no OTLP-HTTP serializer; remote-write lands in the same + /// Prometheus TSDB the OTel-collector path lands in via OTLP HTTP). + #[test] + fn telegraf_toml_mode3_prometheus_archive_shape() { + let toml = emit_telegraf_toml(&prom_edge_cfg_mode3(), None) + .expect("emit_telegraf_toml ok"); + assert!(toml.contains("[[inputs.opentelemetry]]"), "missing input\n{toml}"); + assert!( + !toml.contains("[[processors.allsketches]]"), + "Mode 3 must not include sketch processor\n{toml}" + ); + assert!( + toml.contains("[[outputs.http]]"), + "missing http output\n{toml}" + ); + assert!( + toml.contains("url = \"http://prometheus:9090/api/v1/write\""), + "missing Prometheus remote-write URL\n{toml}" + ); + assert!( + toml.contains("data_format = \"prometheusremotewrite\""), + "missing remote-write serializer\n{toml}" + ); + // No outputs.opentelemetry — Mode 3 is a passthrough to + // Prometheus, not the gateway. + assert!( + !toml.contains("[[outputs.opentelemetry]]"), + "Mode 3 must not also export to gateway\n{toml}" + ); + } + + /// Mode 3 with override URL — caller can redirect to a different + /// Prometheus instance. + #[test] + fn telegraf_toml_mode3_url_override() { + let toml = emit_telegraf_toml( + &prom_edge_cfg_mode3(), + Some("https://prom-prod:9090/api/v1/write"), + ) + .expect("emit_telegraf_toml ok"); + assert!( + toml.contains("https://prom-prod:9090/api/v1/write"), + "override URL not propagated\n{toml}" + ); + } + + /// Catch malformed output early — emit_telegraf_toml should produce + /// TOML that round-trips through our minimal validator. + #[test] + fn telegraf_toml_all_modes_parse() { + for (name, cfg) in [ + ("mode1", ddsketch_edge_cfg_mode1()), + ("mode2", raw_edge_cfg_mode2()), + ("mode3", prom_edge_cfg_mode3()), + ] { + let toml = emit_telegraf_toml(&cfg, None) + .unwrap_or_else(|e| panic!("emit failed for {name}: {e}")); + // Parsing happens inside emit_telegraf_toml; if we got Ok, + // parsing succeeded. Spot-check a handful of expected + // tokens defensively. + assert!(toml.contains("inputs.opentelemetry"), "{name}: missing input header"); + } + } + + /// KLL params — k is preserved, no delta_transmission flag (KLL has + /// no delta variant per Implementation.tex). + #[test] + fn telegraf_toml_mode1_kll_no_delta_flag() { + let cfg = EdgeStageConfig { + source_metric: Some("metric".to_string()), + label_filters: Vec::new(), + window_secs: Some(60), + sketch_processors: vec![EdgeSketchProcessor { + processor_name: "kllprocessor".to_string(), + sketch_kind: SketchKind::Kll, + sketch_params: SketchParams::Kll(KllParams { k: 200 }), + aggregation_id: "agg0".to_string(), + }], + exporter_target: ExportTarget::Stage(StageId::Gateway), + prometheus_archive_metrics: Vec::new(), + }; + let toml = emit_telegraf_toml(&cfg, None).expect("emit ok"); + assert!(toml.contains("k = 200"), "k not propagated\n{toml}"); + assert!(toml.contains("sketch_kind = \"kll\""), "kind\n{toml}"); + // KLL has no delta variant — flag must be absent. + assert!( + !toml.contains("delta_transmission"), + "KLL must not have delta_transmission\n{toml}" + ); + } +} diff --git a/docs/control-plane-design.md b/docs/control-plane-design.md index 1d56fa1d..9509dd11 100644 --- a/docs/control-plane-design.md +++ b/docs/control-plane-design.md @@ -347,6 +347,43 @@ The controller is implemented in **Rust**, enabling direct in-process integratio --- +## Per-Runtime Emit Paths (Phase ε.1.5) + +Three edge runtimes ship with ASAP — `sketchcollector` (OTel-collector +contrib build), `sketchotap` (otap-dataflow Rust runtime), and +`sketchtelegraf` (Telegraf runtime). All three accept the same typed L5 +[`EdgeStageConfig`] from the controller's stage_split emitter; each runtime +has its own emit function in `controller/src/config/`: + +| Runtime | Emitter | Output format | +|---|---|---| +| `sketchcollector` | `stage_config::emit_edge_yaml` | OTel-collector YAML | +| `sketchotap` | `stage_config_otap::emit_otap_dag_yaml` | otap-dataflow DAG YAML (`version: otel_dataflow/v1`) | +| `sketchtelegraf` | `stage_config_telegraf::emit_telegraf_toml` | Telegraf TOML | + +`config::emit_for_runtime(runtime, cfg, opamp_endpoint, prometheus_url)` +dispatches by `AgentRuntime`. The runtime is reported by the agent on +OpAMP `on_connect` via the `X-Agent-Runtime` header (`sketchcollector` / +`sketchotap` / `sketchtelegraf`); when absent, the controller defaults +to `Sketchcollector` so legacy agents keep working. + +All three modes from Phase ε.1's [`BindMode`] enum are supported by +all three runtime emitters: + +* **Mode 1 — `SketchAtEdge`**: sketch processor at the edge, OTLP + egress to gateway. +* **Mode 2 — `RawAtEdgeSketchAtBackend`**: passthrough at edge, OTLP + egress to gateway; backend builds sketches at ingest. +* **Mode 3 — `RawAtEdgePrometheusArchive`**: passthrough at edge, + egress to Prometheus's archive. + * `sketchcollector`/`sketchotap` use OTLP HTTP to Prometheus's + native receiver at `/api/v1/otlp/v1/metrics`. + * `sketchtelegraf` uses `outputs.http` with the + `prometheusremotewrite` serializer to `/api/v1/write`. Telegraf + has no OTLP-HTTP serializer in the version we ship; remote-write + lands in the same Prometheus TSDB so the storage outcome is + identical (only the wire framing differs). + ## Out of Scope (for now) - ML-based workload prediction