This is an automated email from the ASF dual-hosted git repository. szaszm pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/nifi-minifi-cpp.git
commit dfcd29022594730244fbf770ee87988edf1274fc Author: Gabor Gyimesi <[email protected]> AuthorDate: Fri Jun 30 07:31:40 2023 +0200 MINIFICPP-2076 Implement logging metrics publisher Closes #1532 Signed-off-by: Marton Szasz <[email protected]> --- METRICS.md | 55 +++++- docker/test/integration/cluster/ContainerStore.py | 3 + .../test/integration/cluster/DockerTestCluster.py | 3 + .../cluster/containers/MinifiContainer.py | 20 +- .../features/MiNiFi_integration_test_driver.py | 3 + .../features/core_functionality.feature | 14 ++ docker/test/integration/features/steps/steps.py | 6 + .../prometheus/PrometheusMetricsPublisher.cpp | 6 +- .../tests/PrometheusMetricsPublisherTest.cpp | 17 +- .../standard-processors/processors/FetchFile.cpp | 40 +--- .../standard-processors/processors/FetchFile.h | 17 +- libminifi/include/core/state/LogMetricsPublisher.h | 57 ++++++ .../include/core/state/MetricsPublisherFactory.h | 4 +- .../include/core/state/MetricsPublisherStore.h | 4 +- libminifi/include/core/state/Value.h | 25 ++- libminifi/include/properties/Configuration.h | 4 + libminifi/include/utils/LogUtils.h | 60 ++++++ libminifi/src/Configuration.cpp | 4 + libminifi/src/core/state/LogMetricsPublisher.cpp | 126 ++++++++++++ .../src/core/state/MetricsPublisherFactory.cpp | 17 +- libminifi/src/core/state/MetricsPublisherStore.cpp | 6 +- libminifi/src/core/state/Value.cpp | 18 +- libminifi/test/unit/LogMetricsPublisherTests.cpp | 217 +++++++++++++++++++++ libminifi/test/unit/MetricsPublisherStoreTests.cpp | 72 +++++++ 24 files changed, 709 insertions(+), 89 deletions(-) diff --git a/METRICS.md b/METRICS.md index f0cc52d54..f384fe570 100644 --- a/METRICS.md +++ b/METRICS.md @@ -42,25 +42,66 @@ Aside from the publisher exposed metrics, metrics are also sent through C2 proto ## Configuration -To configure the a metrics publisher first we have to set which publisher class should be used: +Currently LogMetricsPublisher and PrometheusMetricsPublisher are available that can be configured as metrics publishers. C2 metrics are published through C2 specific properties, see [C2 documentation](C2.md) for more information on that. - # in minifi.properties +The LogMetricsPublisher serializes all the configured metrics into a json output and writes the json to the MiNiFi logs periodically. LogMetricsPublisher follows the conventions of the C2 metrics, and all information that is present in those metrics, including string data, is present in the log metrics as well. An example log entry may look like the following: + + [2023-03-09 15:04:32.268] [org::apache::nifi::minifi::state::LogMetricsPublisher] [info] {"LogMetrics":{"RepositoryMetrics":{"flowfile":{"running":"true","full":"false","size":"0"},"provenance":{"running":"true","full":"false","size":"0"}}}} - nifi.metrics.publisher.class=PrometheusMetricsPublisher +PrometheusMetricsPublisher publishes only numerical metrics to a Prometheus server in Prometheus specific format. This is different from the json format of the C2 and LogMetricsPublisher. -Currently PrometheusMetricsPublisher is the only available publisher in MiNiFi C++ which publishes metrics to a Prometheus server. -To use the publisher a port should also be configured where the metrics will be available to be scraped through: +### Common configuration properties + +To configure the a publisher first we have to specify the class in the properties. One or multiple publisher can be defined in comma separated format: # in minifi.properties - nifi.metrics.publisher.PrometheusMetricsPublisher.port=9936 + nifi.metrics.publisher.class=LogMetricsPublisher -The following option defines which metric classes should be exposed through the metrics publisher in configured with a comma separated value: + # alternatively + + nifi.metrics.publisher.class=LogMetricsPublisher,PrometheusMetricsPublisher + +To define which metrics should be published either the generic or the publisher specific metrics property should be used. The generic metrics are applied to all publishers if no publisher specific metric is specified. # in minifi.properties + # define generic metrics for all selected publisher classes + nifi.metrics.publisher.metrics=QueueMetrics,RepositoryMetrics,GetFileMetrics,DeviceInfoNode,FlowInformation,processorMetrics/Tail.* + # alternatively LogMetricsPublisher will only use the following metrics + + nifi.metrics.publisher.LogMetricsPublisher.metrics=QueueMetrics,RepositoryMetrics + +Additional configuration properties may be required by specific publishers, these are listed below. + +### LogMetricsPublisher + +LogMetricsPublisher requires a logging interval to be configured which states how often the selected metrics should be logged + + # in minifi.properties + + # log the metrics in MiNiFi app logs every 30 seconds + + nifi.metrics.publisher.LogMetricsPublisher.logging.interval=30s + +Optionally LogMetricsPublisher can be configured which log level should the publisher use. The default log level is INFO + + # in minifi.properties + + # change log level to debug + + nifi.metrics.publisher.LogMetricsPublisher.log.level=DEBUG + +### PrometheusMetricsPublisher + +PrometheusMetricsPublisher requires a port to be configured where the metrics will be available to be scraped from: + + # in minifi.properties + + nifi.metrics.publisher.PrometheusMetricsPublisher.port=9936 + An agent identifier should also be defined to identify which agent the metric is exposed from. If not set, the hostname is used as the identifier. # in minifi.properties diff --git a/docker/test/integration/cluster/ContainerStore.py b/docker/test/integration/cluster/ContainerStore.py index 2f9d29aef..b0bb0449a 100644 --- a/docker/test/integration/cluster/ContainerStore.py +++ b/docker/test/integration/cluster/ContainerStore.py @@ -322,6 +322,9 @@ class ContainerStore: def set_controller_socket_properties_in_minifi(self): self.minifi_options.enable_controller_socket = True + def enable_log_metrics_publisher_in_minifi(self): + self.minifi_options.enable_log_metrics_publisher = True + def get_startup_finished_log_entry(self, container_name): container_name = self.get_container_name_with_postfix(container_name) return self.containers[container_name].get_startup_finished_log_entry() diff --git a/docker/test/integration/cluster/DockerTestCluster.py b/docker/test/integration/cluster/DockerTestCluster.py index 6e694712c..6d4ad4dc5 100644 --- a/docker/test/integration/cluster/DockerTestCluster.py +++ b/docker/test/integration/cluster/DockerTestCluster.py @@ -96,6 +96,9 @@ class DockerTestCluster: def set_controller_socket_properties_in_minifi(self): self.container_store.set_controller_socket_properties_in_minifi() + def enable_log_metrics_publisher_in_minifi(self): + self.container_store.enable_log_metrics_publisher_in_minifi() + def get_app_log(self, container_name): container_name = self.container_store.get_container_name_with_postfix(container_name) log_source = self.container_store.log_source(container_name) diff --git a/docker/test/integration/cluster/containers/MinifiContainer.py b/docker/test/integration/cluster/containers/MinifiContainer.py index de1fa6eed..3ce398f7e 100644 --- a/docker/test/integration/cluster/containers/MinifiContainer.py +++ b/docker/test/integration/cluster/containers/MinifiContainer.py @@ -36,6 +36,7 @@ class MinifiOptions: self.use_flow_config_from_url = False self.set_ssl_context_properties = False self.enable_controller_socket = False + self.enable_log_metrics_publisher = False class MinifiContainer(FlowContainer): @@ -117,11 +118,20 @@ class MinifiContainer(FlowContainer): if not self.options.enable_provenance: f.write("nifi.provenance.repository.class.name=NoOpRepository\n") - if self.options.enable_prometheus: - f.write("nifi.metrics.publisher.agent.identifier=Agent1\n") - f.write("nifi.metrics.publisher.class=PrometheusMetricsPublisher\n") - f.write("nifi.metrics.publisher.PrometheusMetricsPublisher.port=9936\n") - f.write("nifi.metrics.publisher.metrics=RepositoryMetrics,QueueMetrics,PutFileMetrics,processorMetrics/Get.*,FlowInformation,DeviceInfoNode,AgentStatus\n") + if self.options.enable_prometheus or self.options.enable_log_metrics_publisher: + classes = [] + if self.options.enable_prometheus: + f.write("nifi.metrics.publisher.agent.identifier=Agent1\n") + f.write("nifi.metrics.publisher.PrometheusMetricsPublisher.port=9936\n") + f.write("nifi.metrics.publisher.PrometheusMetricsPublisher.metrics=RepositoryMetrics,QueueMetrics,PutFileMetrics,processorMetrics/Get.*,FlowInformation,DeviceInfoNode,AgentStatus\n") + classes.append("PrometheusMetricsPublisher") + + if self.options.enable_log_metrics_publisher: + f.write("nifi.metrics.publisher.LogMetricsPublisher.metrics=RepositoryMetrics\n") + f.write("nifi.metrics.publisher.LogMetricsPublisher.logging.interval=1s\n") + classes.append("LogMetricsPublisher") + + f.write("nifi.metrics.publisher.class=" + ",".join(classes) + "\n") if self.options.use_flow_config_from_url: f.write(f"nifi.c2.flow.url=http://minifi-c2-server-{self.feature_context.id}:10090/c2/config?class=minifi-test-class\n") diff --git a/docker/test/integration/features/MiNiFi_integration_test_driver.py b/docker/test/integration/features/MiNiFi_integration_test_driver.py index dd8c3cc18..cf26a1c21 100644 --- a/docker/test/integration/features/MiNiFi_integration_test_driver.py +++ b/docker/test/integration/features/MiNiFi_integration_test_driver.py @@ -388,3 +388,6 @@ class MiNiFi_integration_test: def manifest_can_be_retrieved_through_minifi_controller(self, container_name: str): assert self.cluster.manifest_can_be_retrieved_through_minifi_controller(container_name) or self.cluster.log_app_output() + + def enable_log_metrics_publisher_in_minifi(self): + self.cluster.enable_log_metrics_publisher_in_minifi() diff --git a/docker/test/integration/features/core_functionality.feature b/docker/test/integration/features/core_functionality.feature index e43a6cb6f..8d3352bd4 100644 --- a/docker/test/integration/features/core_functionality.feature +++ b/docker/test/integration/features/core_functionality.feature @@ -70,3 +70,17 @@ Feature: Core flow functionalities When all instances start up Then the peak memory usage of the agent is more than 130 MB in less than 20 seconds And the memory usage of the agent decreases to 70% peak usage in less than 20 seconds + + Scenario: Metrics can be logged + Given a GenerateFlowFile processor + And log metrics publisher is enabled in MiNiFi + When all instances start up + Then the Minifi logs contain the following message: '[info] {' in less than 30 seconds + And the Minifi logs contain the following message: ' "LogMetrics": {' in less than 2 seconds + And the Minifi logs contain the following message: ' "RepositoryMetrics": {' in less than 2 seconds + And the Minifi logs contain the following message: ' "flowfile": {' in less than 2 seconds + And the Minifi logs contain the following message: ' "running": "true",' in less than 2 seconds + And the Minifi logs contain the following message: ' "full": "false",' in less than 2 seconds + And the Minifi logs contain the following message: ' "size": "0"' in less than 2 seconds + And the Minifi logs contain the following message: ' },' in less than 2 seconds + And the Minifi logs contain the following message: ' "provenance": {' in less than 2 seconds diff --git a/docker/test/integration/features/steps/steps.py b/docker/test/integration/features/steps/steps.py index 00f6093ed..20a7b5fa9 100644 --- a/docker/test/integration/features/steps/steps.py +++ b/docker/test/integration/features/steps/steps.py @@ -348,6 +348,11 @@ def step_impl(context): context.test.enable_prometheus_in_minifi() +@given("log metrics publisher is enabled in MiNiFi") +def step_impl(context): + context.test.enable_log_metrics_publisher_in_minifi() + + # HTTP proxy setup @given("the http proxy server is set up") @given("a http proxy server is set up accordingly") @@ -902,6 +907,7 @@ def step_impl(context, query, number_of_rows, timeout_seconds): @then("the Minifi logs contain the following message: \"{log_message}\" in less than {duration}") +@then("the Minifi logs contain the following message: '{log_message}' in less than {duration}") def step_impl(context, log_message, duration): context.test.check_minifi_log_contents(log_message, humanfriendly.parse_timespan(duration)) diff --git a/extensions/prometheus/PrometheusMetricsPublisher.cpp b/extensions/prometheus/PrometheusMetricsPublisher.cpp index 80b65033a..f745249a3 100644 --- a/extensions/prometheus/PrometheusMetricsPublisher.cpp +++ b/extensions/prometheus/PrometheusMetricsPublisher.cpp @@ -70,7 +70,11 @@ void PrometheusMetricsPublisher::loadMetricNodes() { std::vector<state::response::SharedResponseNode> PrometheusMetricsPublisher::getMetricNodes() { gsl_Expects(response_node_loader_ && configuration_); std::vector<state::response::SharedResponseNode> nodes; - if (auto metric_classes_str = configuration_->get(minifi::Configuration::nifi_metrics_publisher_metrics)) { + auto metric_classes_str = configuration_->get(minifi::Configuration::nifi_metrics_publisher_prometheus_metrics_publisher_metrics); + if (!metric_classes_str || metric_classes_str->empty()) { + metric_classes_str = configuration_->get(minifi::Configuration::nifi_metrics_publisher_metrics); + } + if (metric_classes_str && !metric_classes_str->empty()) { auto metric_classes = utils::StringUtils::split(*metric_classes_str, ","); for (const std::string& clazz : metric_classes) { auto response_nodes = response_node_loader_->loadResponseNodes(clazz); diff --git a/extensions/prometheus/tests/PrometheusMetricsPublisherTest.cpp b/extensions/prometheus/tests/PrometheusMetricsPublisherTest.cpp index dcb11082f..365d892ef 100644 --- a/extensions/prometheus/tests/PrometheusMetricsPublisherTest.cpp +++ b/extensions/prometheus/tests/PrometheusMetricsPublisherTest.cpp @@ -93,7 +93,22 @@ TEST_CASE_METHOD(PrometheusPublisherTestFixtureWithRealExposer, "Test prometheus } TEST_CASE_METHOD(PrometheusPublisherTestFixtureWithDummyExposer, "Test adding metrics to exposer", "[prometheusPublisherTest]") { - configuration_->set(Configure::nifi_metrics_publisher_metrics, "QueueMetrics,RepositoryMetrics,DeviceInfoNode,FlowInformation,AgentInformation,InvalidMetrics,GetFileMetrics,GetTCPMetrics"); + SECTION("Define metrics in general metrics property") { + configuration_->set(Configure::nifi_metrics_publisher_metrics, "QueueMetrics,RepositoryMetrics,DeviceInfoNode,FlowInformation,AgentInformation,InvalidMetrics,GetFileMetrics,GetTCPMetrics"); + } + SECTION("Define metrics in publisher specific metrics property") { + configuration_->set(Configure::nifi_metrics_publisher_prometheus_metrics_publisher_metrics, + "QueueMetrics,RepositoryMetrics,DeviceInfoNode,FlowInformation,AgentInformation,InvalidMetrics,GetFileMetrics,GetTCPMetrics"); + } + SECTION("Publisher specific metrics property should be prioritized") { + configuration_->set(Configure::nifi_metrics_publisher_prometheus_metrics_publisher_metrics, + "QueueMetrics,RepositoryMetrics,DeviceInfoNode,FlowInformation,AgentInformation,InvalidMetrics,GetFileMetrics,GetTCPMetrics"); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "QueueMetrics"); + } + SECTION("Empty metrics property should be ignored") { + configuration_->set(Configure::nifi_metrics_publisher_metrics, "QueueMetrics,RepositoryMetrics,DeviceInfoNode,FlowInformation,AgentInformation,InvalidMetrics,GetFileMetrics,GetTCPMetrics"); + configuration_->set(Configure::nifi_metrics_publisher_prometheus_metrics_publisher_metrics, ""); + } configuration_->set(Configure::nifi_metrics_publisher_agent_identifier, "AgentId-1"); publisher_->initialize(configuration_, response_node_loader_); publisher_->loadMetricNodes(); diff --git a/extensions/standard-processors/processors/FetchFile.cpp b/extensions/standard-processors/processors/FetchFile.cpp index 4fcce6ce7..e3810e918 100644 --- a/extensions/standard-processors/processors/FetchFile.cpp +++ b/extensions/standard-processors/processors/FetchFile.cpp @@ -61,16 +61,16 @@ const core::Property FetchFile::MoveConflictStrategy( const core::Property FetchFile::LogLevelWhenFileNotFound( core::PropertyBuilder::createProperty("Log level when file not found") ->withDescription("Log level to use in case the file does not exist when the processor is triggered") - ->withDefaultValue<std::string>(toString(LogLevelOption::LOGGING_ERROR)) - ->withAllowableValues<std::string>(LogLevelOption::values()) + ->withDefaultValue<std::string>(toString(utils::LogUtils::LogLevelOption::LOGGING_ERROR)) + ->withAllowableValues<std::string>(utils::LogUtils::LogLevelOption::values()) ->isRequired(true) ->build()); const core::Property FetchFile::LogLevelWhenPermissionDenied( core::PropertyBuilder::createProperty("Log level when permission denied") ->withDescription("Log level to use in case agent does not have sufficient permissions to read the file") - ->withDefaultValue<std::string>(toString(LogLevelOption::LOGGING_ERROR)) - ->withAllowableValues<std::string>(LogLevelOption::values()) + ->withDefaultValue<std::string>(toString(utils::LogUtils::LogLevelOption::LOGGING_ERROR)) + ->withAllowableValues<std::string>(utils::LogUtils::LogLevelOption::values()) ->isRequired(true) ->build()); @@ -99,8 +99,8 @@ void FetchFile::onSchedule(const std::shared_ptr<core::ProcessContext> &context, throw Exception(PROCESS_SCHEDULE_EXCEPTION, "Move Destination Directory is required when Completion Strategy is set to Move File"); } move_confict_strategy_ = utils::parseEnumProperty<MoveConflictStrategyOption>(*context, MoveConflictStrategy); - log_level_when_file_not_found_ = utils::parseEnumProperty<LogLevelOption>(*context, LogLevelWhenFileNotFound); - log_level_when_permission_denied_ = utils::parseEnumProperty<LogLevelOption>(*context, LogLevelWhenPermissionDenied); + log_level_when_file_not_found_ = utils::parseEnumProperty<utils::LogUtils::LogLevelOption>(*context, LogLevelWhenFileNotFound); + log_level_when_permission_denied_ = utils::parseEnumProperty<utils::LogUtils::LogLevelOption>(*context, LogLevelWhenPermissionDenied); } std::filesystem::path FetchFile::getFileToFetch(core::ProcessContext& context, const std::shared_ptr<core::FlowFile>& flow_file) { @@ -116,30 +116,6 @@ std::filesystem::path FetchFile::getFileToFetch(core::ProcessContext& context, c return std::filesystem::path(file_to_fetch_path) / filename; } -template<typename... Args> -void FetchFile::logWithLevel(LogLevelOption log_level, Args&&... args) const { - switch (log_level.value()) { - case LogLevelOption::LOGGING_TRACE: - logger_->log_trace(std::forward<Args>(args)...); - break; - case LogLevelOption::LOGGING_DEBUG: - logger_->log_debug(std::forward<Args>(args)...); - break; - case LogLevelOption::LOGGING_INFO: - logger_->log_info(std::forward<Args>(args)...); - break; - case LogLevelOption::LOGGING_WARN: - logger_->log_warn(std::forward<Args>(args)...); - break; - case LogLevelOption::LOGGING_ERROR: - logger_->log_error(std::forward<Args>(args)...); - break; - case LogLevelOption::LOGGING_OFF: - default: - break; - } -} - std::filesystem::path FetchFile::getMoveAbsolutePath(const std::filesystem::path& file_name) const { return move_destination_directory_ / file_name; } @@ -209,7 +185,7 @@ void FetchFile::onTrigger(const std::shared_ptr<core::ProcessContext> &context, const auto file_to_fetch_path = getFileToFetch(*context, flow_file); if (!std::filesystem::is_regular_file(file_to_fetch_path)) { - logWithLevel(log_level_when_file_not_found_, "File to fetch was not found: '%s'!", file_to_fetch_path.string()); + utils::LogUtils::logWithLevel(logger_, log_level_when_file_not_found_, "File to fetch was not found: '%s'!", file_to_fetch_path.string()); session->transfer(flow_file, NotFound); return; } @@ -232,7 +208,7 @@ void FetchFile::onTrigger(const std::shared_ptr<core::ProcessContext> &context, session->transfer(flow_file, Success); } catch (const utils::FileReaderCallbackIOError& io_error) { if (io_error.error_code == EACCES) { - logWithLevel(log_level_when_permission_denied_, "Read permission denied for file '%s' to be fetched!", file_to_fetch_path.string()); + utils::LogUtils::logWithLevel(logger_, log_level_when_permission_denied_, "Read permission denied for file '%s' to be fetched!", file_to_fetch_path.string()); session->transfer(flow_file, PermissionDenied); } else { logger_->log_error("Fetching file '%s' failed! %s", file_to_fetch_path.string(), io_error.what()); diff --git a/extensions/standard-processors/processors/FetchFile.h b/extensions/standard-processors/processors/FetchFile.h index 90e0c411d..aa2fc9026 100644 --- a/extensions/standard-processors/processors/FetchFile.h +++ b/extensions/standard-processors/processors/FetchFile.h @@ -25,6 +25,7 @@ #include "core/Property.h" #include "utils/Enum.h" #include "core/logging/LoggerConfiguration.h" +#include "utils/LogUtils.h" namespace org::apache::nifi::minifi::processors { @@ -43,15 +44,6 @@ class FetchFile : public core::Processor { (FAIL, "Fail") ) - SMART_ENUM(LogLevelOption, - (LOGGING_TRACE, "TRACE"), - (LOGGING_DEBUG, "DEBUG"), - (LOGGING_INFO, "INFO"), - (LOGGING_WARN, "WARN"), - (LOGGING_ERROR, "ERROR"), - (LOGGING_OFF, "OFF") - ) - explicit FetchFile(std::string name, const utils::Identifier& uuid = {}) : core::Processor(std::move(name), uuid) { } @@ -101,9 +93,6 @@ class FetchFile : public core::Processor { void onTrigger(const std::shared_ptr<core::ProcessContext> &context, const std::shared_ptr<core::ProcessSession> &session) override; private: - template<typename... Args> - void logWithLevel(LogLevelOption log_level, Args&&... args) const; - static std::filesystem::path getFileToFetch(core::ProcessContext& context, const std::shared_ptr<core::FlowFile>& flow_file); std::filesystem::path getMoveAbsolutePath(const std::filesystem::path& file_name) const; bool moveDestinationConflicts(const std::filesystem::path& file_name) const; @@ -115,8 +104,8 @@ class FetchFile : public core::Processor { std::filesystem::path move_destination_directory_; CompletionStrategyOption completion_strategy_; MoveConflictStrategyOption move_confict_strategy_; - LogLevelOption log_level_when_file_not_found_; - LogLevelOption log_level_when_permission_denied_; + utils::LogUtils::LogLevelOption log_level_when_file_not_found_; + utils::LogUtils::LogLevelOption log_level_when_permission_denied_; std::shared_ptr<core::logging::Logger> logger_ = core::logging::LoggerFactory<FetchFile>::getLogger(uuid_); }; diff --git a/libminifi/include/core/state/LogMetricsPublisher.h b/libminifi/include/core/state/LogMetricsPublisher.h new file mode 100644 index 000000000..04b28b172 --- /dev/null +++ b/libminifi/include/core/state/LogMetricsPublisher.h @@ -0,0 +1,57 @@ +/** + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +#pragma once + +#include <memory> +#include <vector> +#include <mutex> + +#include "core/logging/Logger.h" +#include "core/logging/LoggerConfiguration.h" +#include "MetricsPublisher.h" +#include "core/state/nodes/MetricsBase.h" +#include "utils/LogUtils.h" +#include "utils/StoppableThread.h" + +namespace org::apache::nifi::minifi::state { + +class LogMetricsPublisher : public MetricsPublisher { + public: + using MetricsPublisher::MetricsPublisher; + + MINIFIAPI static constexpr const char* Description = "Serializes all the configured metrics into a json output and writes the json to the MiNiFi logs periodically"; + + void initialize(const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader) override; + void clearMetricNodes() override; + void loadMetricNodes() override; + ~LogMetricsPublisher() override; + + private: + void readLoggingInterval(); + void readLogLevel(); + void logMetrics(); + + std::unique_ptr<utils::StoppableThread> metrics_logger_thread_; + utils::LogUtils::LogLevelOption log_level_ = utils::LogUtils::LogLevelOption::LOGGING_INFO; + std::chrono::milliseconds logging_interval_; + std::mutex response_nodes_mutex_; + std::vector<state::response::SharedResponseNode> response_nodes_; + std::shared_ptr<core::logging::Logger> logger_{core::logging::LoggerFactory<LogMetricsPublisher>::getLogger()}; +}; + +} // namespace org::apache::nifi::minifi::state diff --git a/libminifi/include/core/state/MetricsPublisherFactory.h b/libminifi/include/core/state/MetricsPublisherFactory.h index e483d36e6..79313eb0f 100644 --- a/libminifi/include/core/state/MetricsPublisherFactory.h +++ b/libminifi/include/core/state/MetricsPublisherFactory.h @@ -19,6 +19,7 @@ #include <memory> #include <string> +#include <vector> #include "MetricsPublisher.h" #include "utils/gsl.h" @@ -27,6 +28,7 @@ namespace org::apache::nifi::minifi::state { gsl::not_null<std::unique_ptr<MetricsPublisher>> createMetricsPublisher(const std::string& name, const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader); -std::unique_ptr<MetricsPublisher> createMetricsPublisher(const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader); +std::vector<gsl::not_null<std::unique_ptr<MetricsPublisher>>> createMetricsPublishers( + const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader); } // namespace org::apache::nifi::minifi::state diff --git a/libminifi/include/core/state/MetricsPublisherStore.h b/libminifi/include/core/state/MetricsPublisherStore.h index b4fafaa92..544e82445 100644 --- a/libminifi/include/core/state/MetricsPublisherStore.h +++ b/libminifi/include/core/state/MetricsPublisherStore.h @@ -40,12 +40,12 @@ class MetricsPublisherStore { std::weak_ptr<state::MetricsPublisher> getMetricsPublisher(const std::string& name) const; private: - void addMetricsPublisher(const std::string& name, std::shared_ptr<state::MetricsPublisher> publisher) { + void addMetricsPublisher(std::string name, std::shared_ptr<state::MetricsPublisher> publisher) { if (!publisher) { return; } - metrics_publishers_.emplace(name, gsl::make_not_null(std::move(publisher))); + metrics_publishers_.emplace(std::move(name), gsl::make_not_null(std::move(publisher))); } std::shared_ptr<Configure> configuration_; diff --git a/libminifi/include/core/state/Value.h b/libminifi/include/core/state/Value.h index 6d8b77219..e956fe2f3 100644 --- a/libminifi/include/core/state/Value.h +++ b/libminifi/include/core/state/Value.h @@ -15,8 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -#ifndef LIBMINIFI_INCLUDE_CORE_STATE_VALUE_H_ -#define LIBMINIFI_INCLUDE_CORE_STATE_VALUE_H_ +#pragma once #include <typeindex> #include <limits> @@ -31,6 +30,9 @@ #include "utils/ValueCaster.h" #include "utils/Export.h" #include "utils/meta/type_list.h" +#include "rapidjson/writer.h" +#include "rapidjson/document.h" +#include "rapidjson/stringbuffer.h" namespace org::apache::nifi::minifi::state::response { @@ -558,13 +560,26 @@ struct SerializedResponseNode { return value.empty() && children.empty(); } - [[nodiscard]] std::string to_string() const; + template<typename Writer = rapidjson::Writer<rapidjson::StringBuffer>> + [[nodiscard]] std::string to_string() const { + rapidjson::Document doc; + doc.SetObject(); + doc.AddMember(rapidjson::Value(name.c_str(), doc.GetAllocator()), nodeToJson(*this, doc.GetAllocator()), doc.GetAllocator()); + rapidjson::StringBuffer buf; + Writer writer{buf}; + doc.Accept(writer); + return buf.GetString(); + } + + [[nodiscard]] std::string to_pretty_string() const; + + private: + static rapidjson::Value nodeToJson(const SerializedResponseNode& node, rapidjson::MemoryPoolAllocator<rapidjson::CrtAllocator>& alloc); }; + inline std::string to_string(const SerializedResponseNode& node) { return node.to_string(); } std::string hashResponseNodes(const std::vector<SerializedResponseNode>& nodes); } // namespace org::apache::nifi::minifi::state::response - -#endif // LIBMINIFI_INCLUDE_CORE_STATE_VALUE_H_ diff --git a/libminifi/include/properties/Configuration.h b/libminifi/include/properties/Configuration.h index 00e581a52..3e8295510 100644 --- a/libminifi/include/properties/Configuration.h +++ b/libminifi/include/properties/Configuration.h @@ -183,6 +183,10 @@ class Configuration : public Properties { static constexpr const char *nifi_metrics_publisher_agent_identifier = "nifi.metrics.publisher.agent.identifier"; static constexpr const char *nifi_metrics_publisher_class = "nifi.metrics.publisher.class"; static constexpr const char *nifi_metrics_publisher_prometheus_metrics_publisher_port = "nifi.metrics.publisher.PrometheusMetricsPublisher.port"; + static constexpr const char *nifi_metrics_publisher_prometheus_metrics_publisher_metrics = "nifi.metrics.publisher.PrometheusMetricsPublisher.metrics"; + static constexpr const char *nifi_metrics_publisher_log_metrics_publisher_metrics = "nifi.metrics.publisher.LogMetricsPublisher.metrics"; + static constexpr const char *nifi_metrics_publisher_log_metrics_logging_interval = "nifi.metrics.publisher.LogMetricsPublisher.logging.interval"; + static constexpr const char *nifi_metrics_publisher_log_metrics_log_level = "nifi.metrics.publisher.LogMetricsPublisher.log.level"; static constexpr const char *nifi_metrics_publisher_metrics = "nifi.metrics.publisher.metrics"; // Controller socket options diff --git a/libminifi/include/utils/LogUtils.h b/libminifi/include/utils/LogUtils.h new file mode 100644 index 000000000..89191ce06 --- /dev/null +++ b/libminifi/include/utils/LogUtils.h @@ -0,0 +1,60 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +#pragma once + +#include <memory> +#include <utility> + +#include "utils/Enum.h" + +namespace org::apache::nifi::minifi::utils::LogUtils { + +SMART_ENUM(LogLevelOption, + (LOGGING_TRACE, "TRACE"), + (LOGGING_DEBUG, "DEBUG"), + (LOGGING_INFO, "INFO"), + (LOGGING_WARN, "WARN"), + (LOGGING_ERROR, "ERROR"), + (LOGGING_OFF, "OFF") +) + +template<typename... Args> +void logWithLevel(const std::shared_ptr<core::logging::Logger>& logger, LogLevelOption log_level, Args&&... args) { + switch (log_level.value()) { + case LogLevelOption::LOGGING_TRACE: + logger->log_trace(std::forward<Args>(args)...); + break; + case LogLevelOption::LOGGING_DEBUG: + logger->log_debug(std::forward<Args>(args)...); + break; + case LogLevelOption::LOGGING_INFO: + logger->log_info(std::forward<Args>(args)...); + break; + case LogLevelOption::LOGGING_WARN: + logger->log_warn(std::forward<Args>(args)...); + break; + case LogLevelOption::LOGGING_ERROR: + logger->log_error(std::forward<Args>(args)...); + break; + case LogLevelOption::LOGGING_OFF: + default: + break; + } +} + + +} // namespace org::apache::nifi::minifi::utils::LogUtils diff --git a/libminifi/src/Configuration.cpp b/libminifi/src/Configuration.cpp index 34bbaac40..658743b0d 100644 --- a/libminifi/src/Configuration.cpp +++ b/libminifi/src/Configuration.cpp @@ -145,6 +145,10 @@ const std::unordered_map<std::string_view, gsl::not_null<const core::PropertyVal {Configuration::nifi_metrics_publisher_agent_identifier, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, {Configuration::nifi_metrics_publisher_class, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, {Configuration::nifi_metrics_publisher_prometheus_metrics_publisher_port, gsl::make_not_null(&core::StandardValidators::PORT_VALIDATOR)}, + {Configuration::nifi_metrics_publisher_prometheus_metrics_publisher_metrics, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, + {Configuration::nifi_metrics_publisher_log_metrics_publisher_metrics, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, + {Configuration::nifi_metrics_publisher_log_metrics_logging_interval, gsl::make_not_null(&core::StandardValidators::TIME_PERIOD_VALIDATOR)}, + {Configuration::nifi_metrics_publisher_log_metrics_log_level, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, {Configuration::nifi_metrics_publisher_metrics, gsl::make_not_null(&core::StandardValidators::VALID_VALIDATOR)}, {Configuration::controller_socket_enable, gsl::make_not_null(&core::StandardValidators::BOOLEAN_VALIDATOR)}, {Configuration::controller_socket_local_any_interface, gsl::make_not_null(&core::StandardValidators::BOOLEAN_VALIDATOR)}, diff --git a/libminifi/src/core/state/LogMetricsPublisher.cpp b/libminifi/src/core/state/LogMetricsPublisher.cpp new file mode 100644 index 000000000..e6830af16 --- /dev/null +++ b/libminifi/src/core/state/LogMetricsPublisher.cpp @@ -0,0 +1,126 @@ +/** + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +#include "core/state/LogMetricsPublisher.h" + +#include "core/Resource.h" +#include "properties/Configuration.h" +#include "core/TypedValues.h" + +namespace org::apache::nifi::minifi::state { + +LogMetricsPublisher::~LogMetricsPublisher() { + if (metrics_logger_thread_) { + metrics_logger_thread_->stopAndJoin(); + } +} + +void LogMetricsPublisher::initialize(const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader) { + state::MetricsPublisher::initialize(configuration, response_node_loader); + readLoggingInterval(); + readLogLevel(); +} + +void LogMetricsPublisher::logMetrics() { + do { + response::SerializedResponseNode parent_node; + parent_node.name = "LogMetrics"; + { + std::lock_guard<std::mutex> lock(response_nodes_mutex_); + for (const auto& response_node : response_nodes_) { + response::SerializedResponseNode metric_response_node; + metric_response_node.name = response_node->getName(); + for (const auto& serialized_node : response_node->serialize()) { + metric_response_node.children.push_back(serialized_node); + } + parent_node.children.push_back(metric_response_node); + } + } + utils::LogUtils::logWithLevel(logger_, log_level_, parent_node.to_pretty_string().c_str()); + } while (!utils::StoppableThread::waitForStopRequest(logging_interval_)); +} + +void LogMetricsPublisher::readLoggingInterval() { + gsl_Expects(configuration_); + if (auto logging_interval_str = configuration_->get(Configure::nifi_metrics_publisher_log_metrics_logging_interval)) { + if (auto logging_interval = minifi::core::TimePeriodValue::fromString(logging_interval_str.value())) { + logging_interval_ = logging_interval->getMilliseconds(); + logger_->log_info("Metric logging interval is set to %" PRId64 " milliseconds", int64_t{logging_interval_.count()}); + return; + } else { + logger_->log_error("Configured logging interval '%s' is invalid!", logging_interval_str.value()); + } + } + + throw Exception(GENERAL_EXCEPTION, "Metrics logging interval not configured for log metrics publisher!"); +} + +void LogMetricsPublisher::readLogLevel() { + gsl_Expects(configuration_); + if (auto log_level_str = configuration_->get(Configure::nifi_metrics_publisher_log_metrics_log_level)) { + log_level_ = utils::LogUtils::LogLevelOption::parse(log_level_str->c_str(), utils::LogUtils::LogLevelOption::LOGGING_INFO, false); + logger_->log_info("Metric log level is set to %s", log_level_.toString()); + return; + } + + logger_->log_info("Metric log level is set to INFO by default"); +} + +void LogMetricsPublisher::clearMetricNodes() { + { + std::lock_guard<std::mutex> lock(response_nodes_mutex_); + logger_->log_debug("Clearing all metric nodes."); + response_nodes_.clear(); + } + if (metrics_logger_thread_) { + metrics_logger_thread_->stopAndJoin(); + metrics_logger_thread_.reset(); + } +} + +void LogMetricsPublisher::loadMetricNodes() { + gsl_Expects(response_node_loader_ && configuration_); + auto metric_classes_str = configuration_->get(minifi::Configuration::nifi_metrics_publisher_log_metrics_publisher_metrics); + if (!metric_classes_str || metric_classes_str->empty()) { + metric_classes_str = configuration_->get(minifi::Configuration::nifi_metrics_publisher_metrics); + } + if (metric_classes_str && !metric_classes_str->empty()) { + auto metric_classes = utils::StringUtils::split(*metric_classes_str, ","); + std::lock_guard<std::mutex> lock(response_nodes_mutex_); + for (const std::string& clazz : metric_classes) { + auto loaded_response_nodes = response_node_loader_->loadResponseNodes(clazz); + if (loaded_response_nodes.empty()) { + logger_->log_warn("Metric class '%s' could not be loaded.", clazz); + continue; + } + response_nodes_.insert(response_nodes_.end(), loaded_response_nodes.begin(), loaded_response_nodes.end()); + } + } + if (response_nodes_.empty()) { + logger_->log_warn("LogMetricsPublisher is configured without any valid metrics!"); + } + if (response_nodes_.empty() && metrics_logger_thread_) { + metrics_logger_thread_->stopAndJoin(); + metrics_logger_thread_.reset(); + } else if (!response_nodes_.empty() && !metrics_logger_thread_) { + metrics_logger_thread_ = std::make_unique<utils::StoppableThread>([this] { logMetrics(); }); + } +} + +REGISTER_RESOURCE(LogMetricsPublisher, DescriptionOnly); + +} // namespace org::apache::nifi::minifi::state diff --git a/libminifi/src/core/state/MetricsPublisherFactory.cpp b/libminifi/src/core/state/MetricsPublisherFactory.cpp index 7453ce67b..e869e0701 100644 --- a/libminifi/src/core/state/MetricsPublisherFactory.cpp +++ b/libminifi/src/core/state/MetricsPublisherFactory.cpp @@ -17,6 +17,8 @@ */ #include "core/state/MetricsPublisherFactory.h" +#include "utils/StringUtils.h" + namespace org::apache::nifi::minifi::state { gsl::not_null<std::unique_ptr<MetricsPublisher>> createMetricsPublisher(const std::string& name, const std::shared_ptr<Configure>& configuration, @@ -35,11 +37,18 @@ gsl::not_null<std::unique_ptr<MetricsPublisher>> createMetricsPublisher(const st return gsl::make_not_null(std::move(metrics_publisher)); } -std::unique_ptr<MetricsPublisher> createMetricsPublisher(const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader) { - if (auto metrics_publisher_class = configuration->get(minifi::Configure::nifi_metrics_publisher_class)) { - return createMetricsPublisher(*metrics_publisher_class, configuration, response_node_loader); +std::vector<gsl::not_null<std::unique_ptr<MetricsPublisher>>> createMetricsPublishers( + const std::shared_ptr<Configure>& configuration, const std::shared_ptr<state::response::ResponseNodeLoader>& response_node_loader) { + if (auto metrics_publisher_class_str = configuration->get(minifi::Configure::nifi_metrics_publisher_class)) { + std::vector<gsl::not_null<std::unique_ptr<MetricsPublisher>>> publishers; + auto publisher_classes = minifi::utils::StringUtils::split(*metrics_publisher_class_str, ","); + publishers.reserve(publisher_classes.size()); + for (const auto& publisher_class : publisher_classes) { + publishers.push_back(createMetricsPublisher(publisher_class, configuration, response_node_loader)); + } + return publishers; } - return nullptr; + return {}; } } // namespace org::apache::nifi::minifi::state diff --git a/libminifi/src/core/state/MetricsPublisherStore.cpp b/libminifi/src/core/state/MetricsPublisherStore.cpp index c837d264a..269cb6267 100644 --- a/libminifi/src/core/state/MetricsPublisherStore.cpp +++ b/libminifi/src/core/state/MetricsPublisherStore.cpp @@ -42,9 +42,9 @@ void MetricsPublisherStore::initialize(core::controller::ControllerServiceProvid addMetricsPublisher(c2::CONTROLLER_SOCKET_METRICS_PUBLISHER, std::move(controller_socket_metrics_publisher)); } - std::shared_ptr metrics_publisher = minifi::state::createMetricsPublisher(configuration_, response_node_loader_); - if (metrics_publisher) { - addMetricsPublisher(minifi::Configure::nifi_metrics_publisher_class, std::move(metrics_publisher)); + for (auto&& publisher : minifi::state::createMetricsPublishers(configuration_, response_node_loader_)) { + auto name = publisher->getName(); + addMetricsPublisher(std::move(name), std::move(publisher)); } loadMetricNodes(nullptr); diff --git a/libminifi/src/core/state/Value.cpp b/libminifi/src/core/state/Value.cpp index 06e8d0cb7..f34478fd1 100644 --- a/libminifi/src/core/state/Value.cpp +++ b/libminifi/src/core/state/Value.cpp @@ -20,9 +20,7 @@ #include <openssl/sha.h> #include <utility> #include <string> -#include "rapidjson/document.h" -#include "rapidjson/writer.h" -#include "rapidjson/stringbuffer.h" +#include "rapidjson/prettywriter.h" namespace org::apache::nifi::minifi::state::response { @@ -56,8 +54,7 @@ std::string hashResponseNodes(const std::vector<SerializedResponseNode>& nodes) return utils::StringUtils::to_hex(digest, true /*uppercase*/); } -namespace { -rapidjson::Value nodeToJson(const SerializedResponseNode& node, rapidjson::MemoryPoolAllocator<rapidjson::CrtAllocator>& alloc) { +rapidjson::Value SerializedResponseNode::nodeToJson(const SerializedResponseNode& node, rapidjson::MemoryPoolAllocator<rapidjson::CrtAllocator>& alloc) { if (node.value.empty()) { if (node.array) { rapidjson::Value result(rapidjson::kArrayType); @@ -76,16 +73,9 @@ rapidjson::Value nodeToJson(const SerializedResponseNode& node, rapidjson::Memor return {node.value.to_string().c_str(), alloc}; } } -} // namespace -std::string SerializedResponseNode::to_string() const { - rapidjson::Document doc; - doc.SetObject(); - doc.AddMember(rapidjson::Value(name.c_str(), doc.GetAllocator()), nodeToJson(*this, doc.GetAllocator()), doc.GetAllocator()); - rapidjson::StringBuffer buf; - rapidjson::Writer<rapidjson::StringBuffer> writer{buf}; - doc.Accept(writer); - return buf.GetString(); +std::string SerializedResponseNode::to_pretty_string() const { + return to_string<rapidjson::PrettyWriter<rapidjson::StringBuffer>>(); } } // namespace org::apache::nifi::minifi::state::response diff --git a/libminifi/test/unit/LogMetricsPublisherTests.cpp b/libminifi/test/unit/LogMetricsPublisherTests.cpp new file mode 100644 index 000000000..35d4ba4f2 --- /dev/null +++ b/libminifi/test/unit/LogMetricsPublisherTests.cpp @@ -0,0 +1,217 @@ +/** + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +#include <memory> +#include <thread> + +#include "../TestBase.h" +#include "../Catch.h" +#include "core/state/LogMetricsPublisher.h" +#include "core/state/nodes/ResponseNodeLoader.h" +#include "core/RepositoryFactory.h" +#include "utils/IntegrationTestUtils.h" + +using namespace std::literals::chrono_literals; + +namespace org::apache::nifi::minifi::test { + +class LogPublisherTestFixture { + public: + LogPublisherTestFixture() + : configuration_(std::make_shared<Configure>()), + provenance_repo_(core::createRepository("provenancerepository", "provenancerepository")), + flow_file_repo_(core::createRepository("flowfilerepository", "flowfilerepository")), + response_node_loader_(std::make_shared<state::response::ResponseNodeLoader>(configuration_, + std::vector<std::shared_ptr<core::RepositoryMetricsSource>>{provenance_repo_, flow_file_repo_}, nullptr)), + publisher_("LogMetricsPublisher") { + } + + protected: + std::shared_ptr<Configure> configuration_; + std::shared_ptr<core::Repository> provenance_repo_; + std::shared_ptr<core::Repository> flow_file_repo_; + std::shared_ptr<state::response::ResponseNodeLoader> response_node_loader_; + minifi::state::LogMetricsPublisher publisher_; +}; + +TEST_CASE_METHOD(LogPublisherTestFixture, "Logging interval property is mandatory", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + SECTION("No logging interval is set") { + REQUIRE_THROWS_WITH(publisher_.initialize(configuration_, response_node_loader_), "General Operation: Metrics logging interval not configured for log metrics publisher!"); + } + SECTION("Logging interval is set to 2 seconds") { + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "2s"); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + publisher_.initialize(configuration_, response_node_loader_); + REQUIRE(verifyLogLinePresenceInPollTime(5s, "Metric logging interval is set to 2000 milliseconds")); + } +} + +TEST_CASE_METHOD(LogPublisherTestFixture, "Verify empty metrics if no valid metrics are defined", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "100ms"); + SECTION("No metrics are defined") {} + SECTION("Only invalid metrics are defined") { + configuration_->set(Configure::nifi_metrics_publisher_metrics, "InvalidMetric,NotValidMetricNode"); + } + publisher_.initialize(configuration_, response_node_loader_); + publisher_.loadMetricNodes(); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + REQUIRE(verifyLogLinePresenceInPollTime(5s, "LogMetricsPublisher is configured without any valid metrics!")); +} + +TEST_CASE_METHOD(LogPublisherTestFixture, "Verify multiple metric nodes in logs", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "100ms"); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "RepositoryMetrics,DeviceInfoNode"); + publisher_.initialize(configuration_, response_node_loader_); + publisher_.loadMetricNodes(); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + std::string expected_log = R"([info] { + "LogMetrics": { + "RepositoryMetrics": { + "provenancerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + }, + "flowfilerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + } + }, + "deviceInfo": { + "identifier":)"; + REQUIRE(verifyLogLinePresenceInPollTime(5s, expected_log)); +} + +TEST_CASE_METHOD(LogPublisherTestFixture, "Verify reloading different metrics", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "100ms"); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "RepositoryMetrics"); + publisher_.initialize(configuration_, response_node_loader_); + publisher_.loadMetricNodes(); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + std::string expected_log = R"([info] { + "LogMetrics": { + "RepositoryMetrics": { + "provenancerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + }, + "flowfilerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + } + } + } +})"; + REQUIRE(verifyLogLinePresenceInPollTime(5s, expected_log)); + publisher_.clearMetricNodes(); + LogTestController::getInstance().reset(); + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "DeviceInfoNode"); + publisher_.loadMetricNodes(); + expected_log = R"([info] { + "LogMetrics": { + "deviceInfo": { + "identifier":)"; + REQUIRE(verifyLogLinePresenceInPollTime(5s, expected_log)); +} + +TEST_CASE_METHOD(LogPublisherTestFixture, "Verify generic and publisher specific metric properties", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "100ms"); + SECTION("Only generic metrics are defined") { + configuration_->set(Configure::nifi_metrics_publisher_metrics, "RepositoryMetrics"); + } + SECTION("Only publisher specific metrics are defined") { + configuration_->set(Configure::nifi_metrics_publisher_log_metrics_publisher_metrics, "RepositoryMetrics"); + } + SECTION("If both generic and publisher specific metrics are defined the publisher specific metrics are used") { + configuration_->set(Configure::nifi_metrics_publisher_log_metrics_publisher_metrics, "RepositoryMetrics"); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "DeviceInfoNode"); + } + publisher_.initialize(configuration_, response_node_loader_); + publisher_.loadMetricNodes(); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + std::string expected_log = R"([info] { + "LogMetrics": { + "RepositoryMetrics": { + "provenancerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + }, + "flowfilerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + } + } + } +})"; + REQUIRE(verifyLogLinePresenceInPollTime(5s, expected_log)); +} + +TEST_CASE_METHOD(LogPublisherTestFixture, "Verify changing log level property for logging", "[LogMetricsPublisher]") { + LogTestController::getInstance().setTrace<minifi::state::LogMetricsPublisher>(); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_logging_interval, "100ms"); + configuration_->set(minifi::Configuration::nifi_metrics_publisher_log_metrics_log_level, "dEbUg"); + configuration_->set(Configure::nifi_metrics_publisher_metrics, "RepositoryMetrics"); + publisher_.initialize(configuration_, response_node_loader_); + publisher_.loadMetricNodes(); + using org::apache::nifi::minifi::utils::verifyLogLinePresenceInPollTime; + std::string expected_log = R"([debug] { + "LogMetrics": { + "RepositoryMetrics": { + "provenancerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + }, + "flowfilerepository": { + "running": "false", + "full": "false", + "size": "0", + "maxSize": "0", + "entryCount": "0" + } + } + } +})"; + REQUIRE(verifyLogLinePresenceInPollTime(5s, expected_log)); +} + +} // namespace org::apache::nifi::minifi::test diff --git a/libminifi/test/unit/MetricsPublisherStoreTests.cpp b/libminifi/test/unit/MetricsPublisherStoreTests.cpp new file mode 100644 index 000000000..77d620453 --- /dev/null +++ b/libminifi/test/unit/MetricsPublisherStoreTests.cpp @@ -0,0 +1,72 @@ +/** + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +#include <memory> + +#include "../TestBase.h" +#include "../Catch.h" +#include "ProvenanceTestHelper.h" +#include "core/repository/VolatileContentRepository.h" +#include "properties/Configuration.h" +#include "core/Resource.h" + +namespace org::apache::nifi::minifi::test { + +class FirstDummyMetricsPublisher : public minifi::state::MetricsPublisher { + public: + using MetricsPublisher::MetricsPublisher; + + static constexpr const char* Description = "FirstDummyMetricsPublisher"; + + void clearMetricNodes() override {} + void loadMetricNodes() override {} +}; + +class SecondDummyMetricsPublisher : public minifi::state::MetricsPublisher { + public: + using MetricsPublisher::MetricsPublisher; + + static constexpr const char* Description = "SecondDummyMetricsPublisher"; + + void clearMetricNodes() override {} + void loadMetricNodes() override {} +}; + +REGISTER_RESOURCE(FirstDummyMetricsPublisher, DescriptionOnly); +REGISTER_RESOURCE(SecondDummyMetricsPublisher, DescriptionOnly); + +TEST_CASE("Test single metrics publisher store", "[MetricsPublisherStore]") { + auto configuration = std::make_shared<minifi::Configure>(); + configuration->set(Configuration::nifi_metrics_publisher_class, "FirstDummyMetricsPublisher"); + minifi::state::MetricsPublisherStore metrics_publisher_store(configuration, std::vector<std::shared_ptr<core::RepositoryMetricsSource>>{}, nullptr); + metrics_publisher_store.initialize(nullptr, nullptr); + auto publisher = metrics_publisher_store.getMetricsPublisher("FirstDummyMetricsPublisher"); + REQUIRE(publisher.lock()); +} + +TEST_CASE("Test multiple metrics publisher stores", "[MetricsPublisherStore]") { + auto configuration = std::make_shared<minifi::Configure>(); + configuration->set(Configuration::nifi_metrics_publisher_class, "FirstDummyMetricsPublisher,SecondDummyMetricsPublisher"); + minifi::state::MetricsPublisherStore metrics_publisher_store(configuration, std::vector<std::shared_ptr<core::RepositoryMetricsSource>>{}, nullptr); + metrics_publisher_store.initialize(nullptr, nullptr); + auto first_publisher = metrics_publisher_store.getMetricsPublisher("FirstDummyMetricsPublisher"); + REQUIRE(first_publisher.lock()); + auto second_publisher = metrics_publisher_store.getMetricsPublisher("SecondDummyMetricsPublisher"); + REQUIRE(second_publisher.lock()); +} + +} // namespace org::apache::nifi::minifi::test
