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
40 changes: 31 additions & 9 deletions src/config/unit_test/mod.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
#![warn(clippy::pedantic)]

// should match vector-unit-test-tests feature
#[cfg(all(
test,
Expand Down Expand Up @@ -64,6 +66,8 @@ pub struct UnitTestResult {
}

impl UnitTest {
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_panics_doc, reason = "Panic documentation deferred")]
pub async fn run(self) -> UnitTestResult {
let diff = config::ConfigDiff::initial(&self.config);
let (topology, _) = RunningTopology::start_validated(self.config, diff, self.pieces)
Expand Down Expand Up @@ -102,6 +106,8 @@ fn init_log_schema_from_paths(
Ok(())
}

// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
pub async fn build_unit_tests_main(
paths: &[ConfigPath],
signal_handler: &mut signal::SignalHandler,
Expand All @@ -121,12 +127,14 @@ pub async fn build_unit_tests_main(
build_unit_tests(config_builder).await
}

// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
pub async fn build_unit_tests(
mut config_builder: ConfigBuilder,
) -> Result<Vec<UnitTest>, Vec<String>> {
// Sanitize config by removing existing sources and sinks
config_builder.sources = Default::default();
config_builder.sinks = Default::default();
config_builder.sources = IndexMap::default();
config_builder.sinks = IndexMap::default();

let test_definitions = std::mem::take(&mut config_builder.tests);
let mut tests = Vec::new();
Expand Down Expand Up @@ -172,6 +180,9 @@ pub struct UnitTestBuildMetadata {
}

impl UnitTestBuildMetadata {
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
#[allow(clippy::missing_panics_doc, reason = "Panic documentation deferred")]
pub fn initialize(config_builder: &mut ConfigBuilder) -> Result<Self, Vec<String>> {
// A unique id used to name test sources and sinks to avoid name clashes
let random_id = Uuid::new_v4().to_string();
Expand All @@ -189,7 +200,7 @@ impl UnitTestBuildMetadata {

// Map a test source to every transform
let mut template_sources = IndexMap::new();
for (key, transform) in config_builder.transforms.iter_mut() {
for (key, transform) in &mut config_builder.transforms {
let test_source_id = source_ids
.get(key)
.expect("Missing test source for a transform")
Expand Down Expand Up @@ -235,6 +246,9 @@ impl UnitTestBuildMetadata {
}

/// Convert test inputs into sources for use in a unit testing topology
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
#[allow(clippy::missing_panics_doc, reason = "Panic documentation deferred")]
pub fn hydrate_into_sources(
&self,
inputs: &[TestInput],
Expand Down Expand Up @@ -268,6 +282,9 @@ impl UnitTestBuildMetadata {
}

/// Convert test outputs into sinks for use in a unit testing topology
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
#[allow(clippy::missing_panics_doc, reason = "Panic documentation deferred")]
pub fn hydrate_into_sinks(
&self,
test_name: &str,
Expand Down Expand Up @@ -296,7 +313,7 @@ impl UnitTestBuildMetadata {
let sink_ids = ids.clone();
let sink_config = UnitTestSinkConfig {
test_name: test_name.to_string(),
transform_ids: ids.iter().map(|id| id.to_string()).collect(),
transform_ids: ids.iter().map(std::string::ToString::to_string).collect(),
result_tx: Arc::new(Mutex::new(Some(tx))),
check: UnitTestSinkCheck::Checks {
conditions: built.conditions,
Expand Down Expand Up @@ -327,7 +344,7 @@ impl UnitTestBuildMetadata {
.map(|(transform_ids, sink_config)| {
let transform_ids_str = transform_ids
.iter()
.map(|s| s.to_string())
.map(std::string::ToString::to_string)
.collect::<Vec<_>>();
let sink_ids = transform_ids
.iter()
Expand Down Expand Up @@ -446,7 +463,7 @@ async fn build_unit_test(
config_builder.global.wildcard_matching.unwrap_or_default(),
)?;
let valid_outputs = graph.output_map()?;
for (_, transform) in config_builder.transforms.iter_mut() {
for (_, transform) in &mut config_builder.transforms {
let inputs = std::mem::take(&mut transform.inputs);
transform.inputs = inputs
.into_iter()
Expand Down Expand Up @@ -476,7 +493,7 @@ async fn build_unit_test(
/// consumed but its other outputs are left unconsumed.
///
/// To avoid warning logs that occur when building such topologies, we construct
/// a NoOp sink here whose sole purpose is to consume any "loose end" outputs.
/// a `NoOp` sink here whose sole purpose is to consume any "loose end" outputs.
fn get_loose_end_outputs_sink(config: &ConfigBuilder) -> Option<SinkOuter<String>> {
let config = config.clone();
let transform_ids = config.transforms.iter().flat_map(|(key, transform)| {
Expand Down Expand Up @@ -508,7 +525,7 @@ fn get_loose_end_outputs_sink(config: &ConfigBuilder) -> Option<SinkOuter<String
None
} else {
let noop_sink = UnitTestSinkConfig {
test_name: "".to_string(),
test_name: String::new(),
transform_ids: vec![],
result_tx: Arc::new(Mutex::new(None)),
check: UnitTestSinkCheck::NoOp,
Expand Down Expand Up @@ -545,7 +562,7 @@ fn build_and_validate_inputs(
errors.push(format!(
"inputs[{index}]: unable to locate target transform '{}'",
input.insert_at
))
));
}
}

Expand All @@ -562,6 +579,11 @@ pub(super) struct BuiltOutput {
pub(super) conditions: Vec<Vec<Condition>>,
}

// https://github.com/vectordotdev/vector/issues/23659
#[allow(
clippy::default_trait_access,
reason = "Preserve inferred default types"
)]
fn build_outputs(
test_outputs: &[TestOutput],
) -> Result<IndexMap<Vec<OutputId>, BuiltOutput>, Vec<String>> {
Expand Down
4 changes: 2 additions & 2 deletions src/config/unit_test/unit_test_components.rs
Original file line number Diff line number Diff line change
Expand Up @@ -240,9 +240,9 @@ impl StreamSink<Event> for UnitTestSink {
let mut check_errors = Vec::new();
for (j, condition) in check.iter().enumerate() {
let mut condition_errors = Vec::new();
for event in output_events.iter() {
for event in &output_events {
match condition.check_with_context(event.clone()).0 {
Ok(_) => {
Ok(()) => {
condition_errors.clear();
break;
}
Expand Down
4 changes: 4 additions & 0 deletions src/config/unix.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
#![warn(clippy::pedantic)]

use std::cell::RefCell;

use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -41,6 +43,8 @@ impl<T> UnixOnly<T> {
/// Pass the closure with `#[cfg(unix)]` so Unix-only APIs in its body are not
/// compiled on other targets. Context is accepted on every target and is
/// dropped, together with the configuration, if Unix is unavailable.
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
pub fn on_unix<C, R>(
self,
context: C,
Expand Down
14 changes: 11 additions & 3 deletions src/config/validation.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
#![warn(clippy::pedantic)]

use std::{collections::HashMap, path::PathBuf};

use futures_util::{FutureExt, StreamExt, TryFutureExt, TryStreamExt, stream};
Expand Down Expand Up @@ -88,6 +90,8 @@ pub fn check_names<'a, I: Iterator<Item = &'a ComponentKey>>(names: I) -> Result
}
}

// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::items_after_statements, reason = "Helper relocation deferred")]
pub fn check_shape(config: &ConfigBuilder) -> Result<(), Vec<String>> {
let mut errors = vec![];

Expand Down Expand Up @@ -214,7 +218,7 @@ pub fn check_values(config: &ConfigBuilder) -> Result<(), Vec<String>> {
/// does not have a named output with the name [`DEFAULT_OUTPUT`]
pub fn check_outputs(config: &ConfigBuilder) -> Result<(), Vec<String>> {
let mut errors = Vec::new();
for (key, source) in config.sources.iter() {
for (key, source) in &config.sources {
let outputs = source.inner.outputs(config.schema.log_namespace());
if outputs
.iter()
Expand All @@ -227,7 +231,7 @@ pub fn check_outputs(config: &ConfigBuilder) -> Result<(), Vec<String>> {
}
}

for (key, transform) in config.transforms.iter() {
for (key, transform) in &config.transforms {
// Structural validation: reserved names, duplicate routes, invalid sample rates.
// These checks run during config compilation. Transforms that need the schema/enrichment
// context must implement validate_with_context(), called later in validate.rs.
Expand Down Expand Up @@ -345,7 +349,10 @@ pub async fn check_buffer_preconditions(config: &Config) -> Result<(), Vec<Strin
let mut errors = Vec::new();

for (mountpoint, buffers) in mountpoint_buffer_mapping {
let buffer_max_size_total: u64 = buffers.iter().map(|usage| usage.max_size()).sum();
let buffer_max_size_total: u64 = buffers
.iter()
.map(vector_lib::buffers::config::DiskUsage::max_size)
.sum();
let mountpoint_total_capacity = mountpoints
.get(&mountpoint)
.copied()
Expand Down Expand Up @@ -383,6 +390,7 @@ async fn process_partitions(partitions: Vec<Partition>) -> heim::Result<IndexMap
.await
}

#[must_use]
pub fn warnings(config: &Config) -> Vec<String> {
let mut warnings = vec![];

Expand Down
45 changes: 25 additions & 20 deletions src/config/watcher.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
#![warn(clippy::pedantic)]

use notify::{EventKind, RecursiveMode, recommended_watcher};
use std::{
collections::{HashMap, HashSet},
Expand Down Expand Up @@ -64,10 +66,13 @@ impl Watcher {
}
}

/// Sends a ReloadFromDisk or ReloadEnrichmentTables on config_path changes.
/// Sends a `ReloadFromDisk` or `ReloadEnrichmentTables` on `config_path` changes.
/// Accumulates file changes until no change for given duration has occurred.
/// Has best effort guarantee of detecting all file changes from the end of
/// this function until the main thread stops.
// https://github.com/vectordotdev/vector/issues/23659
#[allow(clippy::missing_errors_doc, reason = "Error documentation deferred")]
#[allow(clippy::too_many_lines, reason = "Preserve existing control flow")]
pub fn spawn_thread<'a>(
watcher_conf: WatcherConfig,
signal_tx: crate::signal::SignalTx,
Expand Down Expand Up @@ -124,7 +129,7 @@ pub fn spawn_thread<'a>(
let changed_components: HashMap<_, _> = component_configs
.clone()
.into_iter()
.flat_map(|p| p.contains(&changed_paths))
.filter_map(|p| p.contains(&changed_paths))
.collect();

// We need to read paths to resolve any inode changes that may have happened.
Expand All @@ -137,7 +142,17 @@ pub fn spawn_thread<'a>(
debug!(message = "Reloaded paths.");

info!("Configuration file changed.");
if !changed_components.is_empty() {
if changed_components.is_empty() {
_ = signal_tx
.send(crate::signal::SignalTo::ReloadFromDisk)
.map_err(|error| {
error!(
message = "Unable to reload configuration file. Restart Vector to reload it.",
cause = %error,
internal_log_rate_limit = false,
);
});
} else {
info!(
"Component {:?} configuration changed.",
changed_components.keys()
Expand All @@ -154,7 +169,7 @@ pub fn spawn_thread<'a>(
message = "Unable to reload enrichment tables.",
cause = %error,
internal_log_rate_limit = false,
)
);
});
} else {
_ = signal_tx
Expand All @@ -166,22 +181,12 @@ pub fn spawn_thread<'a>(
message = "Unable to reload component configuration. Restart Vector to reload it.",
cause = %error,
internal_log_rate_limit = false,
)
);
});
}
} else {
_ = signal_tx
.send(crate::signal::SignalTo::ReloadFromDisk)
.map_err(|error| {
error!(
message = "Unable to reload configuration file. Restart Vector to reload it.",
cause = %error,
internal_log_rate_limit = false,
)
});
}
} else {
debug!(message = "Ignoring event.", event = ?event)
debug!(message = "Ignoring event.", event = ?event);
}
}
}
Expand All @@ -198,7 +203,7 @@ pub fn spawn_thread<'a>(
// determine if anything changed.
info!("Speculating that configuration files have changed.");
_ = signal_tx.send(crate::signal::SignalTo::ReloadFromDisk).map_err(|error| {
error!(message = "Unable to reload configuration file. Restart Vector to reload it.", cause = %error)
error!(message = "Unable to reload configuration file. Restart Vector to reload it.", cause = %error);
});
}
}
Expand Down Expand Up @@ -263,7 +268,7 @@ mod tests {
trace_init();

let delay = Duration::from_secs(3);
let dir = temp_dir().to_path_buf();
let dir = temp_dir().clone();
let watcher_conf = WatcherConfig::RecommendedWatcher;
let component_file_path = vec![dir.join("tls.cert"), dir.join("tls.key")];
let http_component = ComponentKey::from("http");
Expand Down Expand Up @@ -324,7 +329,7 @@ mod tests {
trace_init();

let delay = Duration::from_secs(3);
let dir = temp_dir().to_path_buf();
let dir = temp_dir().clone();
let file_path = dir.join("vector.toml");
let watcher_conf = WatcherConfig::RecommendedWatcher;

Expand Down Expand Up @@ -403,7 +408,7 @@ mod tests {
trace_init();

let delay = Duration::from_secs(3);
let dir = temp_dir().to_path_buf();
let dir = temp_dir().clone();
let sub_dir = dir.join("sources");
let file_path = sub_dir.join("input.toml");
let watcher_conf = WatcherConfig::RecommendedWatcher;
Expand Down
Loading