From aa33e131012fd747f4a82da42c6f46a7b34842e8 Mon Sep 17 00:00:00 2001 From: Thomas Date: Fri, 7 Aug 2026 14:16:02 -0400 Subject: [PATCH 1/2] chore(tests): drop influxdb from prometheus integration tests --- Cargo.toml | 2 +- src/sinks/mod.rs | 2 +- .../remote_write/integration_tests.rs | 171 ++++++++---------- .../prometheus/config/compose.yaml | 15 -- tests/integration/prometheus/config/test.yaml | 1 - 5 files changed, 80 insertions(+), 111 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 36f5611e1f8f6..8ea15a6e99603 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1099,7 +1099,7 @@ nginx-integration-tests = ["sources-nginx_metrics"] opentelemetry-integration-tests = ["sources-opentelemetry", "dep:prost"] postgresql_metrics-integration-tests = ["sources-postgresql_metrics"] postgres_sink-integration-tests = ["sinks-postgres"] -prometheus-integration-tests = ["sinks-prometheus", "sources-prometheus", "sinks-influxdb"] +prometheus-integration-tests = ["sinks-prometheus", "sources-prometheus"] pulsar-integration-tests = ["sinks-pulsar", "sources-pulsar"] redis-integration-tests = ["sinks-redis", "sources-redis"] splunk-integration-tests = ["sinks-splunk_hec"] diff --git a/src/sinks/mod.rs b/src/sinks/mod.rs index 95967c645aa16..ad000ad41de10 100644 --- a/src/sinks/mod.rs +++ b/src/sinks/mod.rs @@ -72,7 +72,7 @@ pub mod honeycomb; pub mod http; #[cfg(feature = "sinks-humio")] pub mod humio; -#[cfg(any(feature = "sinks-influxdb", feature = "prometheus-integration-tests"))] +#[cfg(feature = "sinks-influxdb")] pub mod influxdb; #[cfg(feature = "sinks-kafka")] pub mod kafka; diff --git a/src/sinks/prometheus/remote_write/integration_tests.rs b/src/sinks/prometheus/remote_write/integration_tests.rs index 25de1d2e3bac4..cbb9bd17b9962 100644 --- a/src/sinks/prometheus/remote_write/integration_tests.rs +++ b/src/sinks/prometheus/remote_write/integration_tests.rs @@ -1,111 +1,96 @@ -use std::{collections::HashMap, ops::Range}; - -use serde_json::Value; +use vector_lib::event::EventStatus; use super::tests::*; use crate::{ - config::{SinkConfig, SinkContext}, - event::{Event, metric::MetricValue}, - sinks::{ - influxdb::test_util::{cleanup_v1, format_timestamp, onboarding_v1, query_v1}, - prometheus::remote_write::config::RemoteWriteConfig, - }, - test_util::components::{HTTP_SINK_TAGS, assert_sink_compliance}, - tls::{self, TlsConfig}, + SourceSender, + config::{SinkConfig, SinkContext, SourceConfig, SourceContext}, + event::Event, + sinks::prometheus::remote_write::config::RemoteWriteConfig, + sources::prometheus::PrometheusRemoteWriteConfig, + test_util::{self, addr::next_addr, wait_for_tcp}, + tls::{self, MaybeTlsSettings, TlsEnableableConfig}, }; -const HTTP_URL: &str = "http://influxdb-v1:8086"; -const HTTPS_URL: &str = "https://influxdb-v1-tls:8087"; - #[tokio::test] async fn insert_metrics_over_http() { - insert_metrics(HTTP_URL).await; + insert_metrics(None).await; } #[tokio::test] async fn insert_metrics_over_https() { - insert_metrics(HTTPS_URL).await; -} - -async fn insert_metrics(url: &str) { - assert_sink_compliance(&HTTP_SINK_TAGS, async { - let database = onboarding_v1(url).await; - - let cx = SinkContext::default(); - - let config = RemoteWriteConfig { - endpoint: format!("{url}/api/v1/prom/write?db={database}"), - tls: Some(TlsConfig { - ca_file: Some(tls::TEST_PEM_CA_PATH.into()), - ..Default::default() - }), - ..Default::default() - }; - let events = create_events(0..5, |n| n * 11.0); - - let (sink, _) = config.build(cx).await.expect("error building config"); - sink.run_events(events.clone()).await.unwrap(); - - let result = query(url, &format!("show series on {database}")).await; - - let values = &result["results"][0]["series"][0]["values"]; - assert_eq!(values.as_array().unwrap().len(), 5); - - for event in events { - let metric = event.into_metric(); - let result = query( - url, - &format!(r#"SELECT * FROM "{}".."{}""#, database, metric.name()), - ) - .await; - - let metrics = decode_metrics(&result["results"][0]["series"][0]); - assert_eq!(metrics.len(), 1); - let output = &metrics[0]; - - match metric.value() { - MetricValue::Gauge { value } => { - assert_eq!(output["value"], Value::Number((*value as u32).into())) - } - _ => panic!("Unhandled metric value, fix the test"), - } - for (tag, value) in metric.tags().unwrap().iter_single() { - assert_eq!(output[tag], Value::String(value.to_string())); - } - let timestamp = - format_timestamp(metric.timestamp().unwrap(), chrono::SecondsFormat::Millis); - assert_eq!(output["time"], Value::String(timestamp)); - } - - cleanup_v1(url, &database).await; - }) - .await -} - -async fn query(url: &str, query: &str) -> Value { - let result = query_v1(url, query).await; - let text = result.text().await.unwrap(); - serde_json::from_str(&text).expect("error when parsing InfluxDB response JSON") + insert_metrics(Some(TlsEnableableConfig::test_config())).await; } -fn decode_metrics(data: &Value) -> Vec> { - let data = data.as_object().expect("Data is not an object"); - let columns = data["columns"].as_array().expect("Columns is not an array"); - data["values"] - .as_array() - .expect("Values is not an array") - .iter() - .map(|values| { - columns - .iter() - .zip(values.as_array().unwrap().iter()) - .map(|(column, value)| (column.as_str().unwrap().to_owned(), value.clone())) - .collect() - }) - .collect() +async fn insert_metrics(tls: Option) { + let (_guard, address) = next_addr(); + let (tx, rx) = SourceSender::new_test_finalize(EventStatus::Delivered); + + let proto = MaybeTlsSettings::from_config(tls.as_ref(), true) + .unwrap() + .http_protocol_name(); + + // Start a `prometheus_remote_write` source to act as the remote write receiver. + let tls_yaml = tls.as_ref().map(|_| { + format!( + "tls:\n enabled: true\n ca_file: \"{}\"\n crt_file: \"{}\"\n key_file: \"{}\"", + tls::TEST_PEM_CA_PATH, + tls::TEST_PEM_CRT_PATH, + tls::TEST_PEM_KEY_PATH + ) + }); + let source: PrometheusRemoteWriteConfig = serde_yaml::from_str(&format!( + "address: \"{address}\"\n{}", + tls_yaml.unwrap_or_default() + )) + .unwrap(); + let source = source + .build(SourceContext::new_test(tx, None)) + .await + .expect("source should not fail to build"); + tokio::spawn(source); + wait_for_tcp(address).await; + + let config = RemoteWriteConfig { + endpoint: format!("{proto}://localhost:{}/", address.port()), + tls: tls.map(|tls| tls.options), + ..Default::default() + }; + let events = create_events(0..5, |n| n * 11.0); + let events_copy = events.clone(); + + let cx = SinkContext::default(); + let (sink, _) = config.build(cx).await.expect("error building config"); + + let mut output = test_util::spawn_collect_ready( + async move { + sink.run_events(events_copy).await.unwrap(); + }, + rx, + 1, + ) + .await; + + // The MetricBuffer used by the sink may reorder the metrics, so + // put them back into order before comparing. + output.sort_unstable_by_key(|event| event.as_metric().name().to_owned()); + + assert_eq!(output.len(), events.len()); + + for (sent, received) in events.iter().zip(output.iter()) { + let sent = sent.as_metric(); + let received = received.as_metric(); + assert_eq!(sent.name(), received.name()); + assert_eq!(sent.value(), received.value()); + assert_eq!(sent.tags(), received.tags()); + // Remote write stores timestamps with millisecond precision. + assert_eq!( + sent.timestamp().unwrap().timestamp_millis(), + received.timestamp().unwrap().timestamp_millis() + ); + } } -fn create_events(name_range: Range, value: impl Fn(f64) -> f64) -> Vec { +fn create_events(name_range: std::ops::Range, value: impl Fn(f64) -> f64) -> Vec { name_range .map(move |num| create_event(format!("metric_{num}"), value(num as f64))) .collect() diff --git a/tests/integration/prometheus/config/compose.yaml b/tests/integration/prometheus/config/compose.yaml index 3ff987c6d5189..c5ec8f9740673 100644 --- a/tests/integration/prometheus/config/compose.yaml +++ b/tests/integration/prometheus/config/compose.yaml @@ -1,21 +1,6 @@ version: "3" services: - influxdb-v1: - image: docker.io/influxdb:${CONFIG_INFLUXDB} - environment: - - INFLUXDB_REPORTING_DISABLED=true - influxdb-v1-tls: - image: docker.io/influxdb:${CONFIG_INFLUXDB} - environment: - - INFLUXDB_REPORTING_DISABLED=true - - INFLUXDB_HTTP_HTTPS_ENABLED=true - - INFLUXDB_HTTP_BIND_ADDRESS=:8087 - - INFLUXDB_BIND_ADDRESS=:8089 - - INFLUXDB_HTTP_HTTPS_CERTIFICATE=/etc/ssl/intermediate_server/certs/influxdb-v1-tls-chain.cert.pem - - INFLUXDB_HTTP_HTTPS_PRIVATE_KEY=/etc/ssl/intermediate_server/private/influxdb-v1-tls.key.pem - volumes: - - ../../../data/ca:/etc/ssl:ro prometheus: image: docker.io/prom/prometheus:${CONFIG_PROMETHEUS} command: --config.file=/etc/vector/prometheus.yaml diff --git a/tests/integration/prometheus/config/test.yaml b/tests/integration/prometheus/config/test.yaml index ac79ea17c3d8e..11cba767636a6 100644 --- a/tests/integration/prometheus/config/test.yaml +++ b/tests/integration/prometheus/config/test.yaml @@ -8,7 +8,6 @@ env: matrix: prometheus: ["v2.33.4"] - influxdb: ["1.8"] # changes to these files/paths will invoke the integration test in CI # expressions are evaluated using https://github.com/micromatch/picomatch From d18da95ff03717b49116d01d58ce966b4ab43dda Mon Sep 17 00:00:00 2001 From: Thomas Date: Fri, 7 Aug 2026 14:24:00 -0400 Subject: [PATCH 2/2] chore(tests): restore assert_sink_compliance in prometheus remote_write integration tests --- .../remote_write/integration_tests.rs | 91 ++++++++++++++++++- 1 file changed, 86 insertions(+), 5 deletions(-) diff --git a/src/sinks/prometheus/remote_write/integration_tests.rs b/src/sinks/prometheus/remote_write/integration_tests.rs index cbb9bd17b9962..e27cbabe65b90 100644 --- a/src/sinks/prometheus/remote_write/integration_tests.rs +++ b/src/sinks/prometheus/remote_write/integration_tests.rs @@ -1,3 +1,14 @@ +use std::net::SocketAddr; + +use bytes::Bytes; +use futures::{FutureExt, SinkExt, TryFutureExt, channel::mpsc}; +use futures_util::StreamExt; +use http::request::Parts; +use hyper::{ + Body, Request, Response, Server, + service::{make_service_fn, service_fn}, +}; +use stream_cancel::{Trigger, Tripwire}; use vector_lib::event::EventStatus; use super::tests::*; @@ -7,7 +18,12 @@ use crate::{ event::Event, sinks::prometheus::remote_write::config::RemoteWriteConfig, sources::prometheus::PrometheusRemoteWriteConfig, - test_util::{self, addr::next_addr, wait_for_tcp}, + test_util::{ + self, + addr::next_addr, + components::{HTTP_SINK_TAGS, assert_sink_compliance}, + wait_for_tcp, + }, tls::{self, MaybeTlsSettings, TlsEnableableConfig}, }; @@ -22,6 +38,37 @@ async fn insert_metrics_over_https() { } async fn insert_metrics(tls: Option) { + // Verify the sink emits the expected compliance events while sending to a + // remote write endpoint. A mock server is used as the receiver here (rather + // than Vector's own source) so that only the sink's internal events are + // recorded. + assert_sink_compliance(&HTTP_SINK_TAGS, async { + let (_guard, addr) = next_addr(); + let (rx, trigger, server) = build_mock_server(addr, tls.clone()).await; + tokio::spawn(server); + + let proto = MaybeTlsSettings::from_config(tls.as_ref(), true) + .unwrap() + .http_protocol_name(); + let config = RemoteWriteConfig { + endpoint: format!("{proto}://localhost:{}/write", addr.port()), + tls: tls.clone().map(|tls| tls.options), + ..Default::default() + }; + let events = create_events(0..5, |n| n * 11.0); + let cx = SinkContext::default(); + let (sink, _) = config.build(cx).await.expect("error building config"); + sink.run_events(events.clone()).await.unwrap(); + + drop(trigger); + let requests = rx.collect::>().await; + assert_eq!(requests.len(), 1); + }) + .await; + + // Verify the sink's output is accepted by Vector's own + // `prometheus_remote_write` source and that the metrics round-trip + // correctly. let (_guard, address) = next_addr(); let (tx, rx) = SourceSender::new_test_finalize(EventStatus::Delivered); @@ -29,7 +76,6 @@ async fn insert_metrics(tls: Option) { .unwrap() .http_protocol_name(); - // Start a `prometheus_remote_write` source to act as the remote write receiver. let tls_yaml = tls.as_ref().map(|_| { format!( "tls:\n enabled: true\n ca_file: \"{}\"\n crt_file: \"{}\"\n key_file: \"{}\"", @@ -70,8 +116,6 @@ async fn insert_metrics(tls: Option) { ) .await; - // The MetricBuffer used by the sink may reorder the metrics, so - // put them back into order before comparing. output.sort_unstable_by_key(|event| event.as_metric().name().to_owned()); assert_eq!(output.len(), events.len()); @@ -82,7 +126,6 @@ async fn insert_metrics(tls: Option) { assert_eq!(sent.name(), received.name()); assert_eq!(sent.value(), received.value()); assert_eq!(sent.tags(), received.tags()); - // Remote write stores timestamps with millisecond precision. assert_eq!( sent.timestamp().unwrap().timestamp_millis(), received.timestamp().unwrap().timestamp_millis() @@ -90,6 +133,44 @@ async fn insert_metrics(tls: Option) { } } +/// Builds a mock remote write HTTP server, optionally with TLS enabled. +async fn build_mock_server( + addr: SocketAddr, + tls: Option, +) -> ( + mpsc::Receiver<(Parts, Bytes)>, + Trigger, + impl std::future::Future>, +) { + let (tx, rx) = mpsc::channel(100); + let service = make_service_fn(move |_| { + let tx = tx.clone(); + async move { + Ok::<_, hyper::Error>(service_fn(move |req: Request| { + let mut tx = tx.clone(); + async move { + let (parts, body) = req.into_parts(); + tokio::spawn(async move { + let bytes = http_body::Body::collect(body).await.unwrap().to_bytes(); + tx.send((parts, bytes)).await.unwrap(); + }); + Ok::<_, hyper::Error>(Response::new(Body::empty())) + } + })) + } + }); + + let settings = MaybeTlsSettings::from_config(tls.as_ref(), true).unwrap(); + let listener = settings.bind(&addr).await.unwrap(); + let (trigger, tripwire) = Tripwire::new(); + let server = Server::builder(hyper::server::accept::from_stream(listener.accept_stream())) + .serve(service) + .with_graceful_shutdown(tripwire.then(crate::shutdown::tripwire_handler)) + .map_err(|error| panic!("Server error: {error}")); + + (rx, trigger, server) +} + fn create_events(name_range: std::ops::Range, value: impl Fn(f64) -> f64) -> Vec { name_range .map(move |num| create_event(format!("metric_{num}"), value(num as f64)))