This is an automated email from the ASF dual-hosted git repository. martinzink pushed a commit to branch api_1_1_enrichment in repository https://gitbox.apache.org/repos/asf/nifi-minifi-cpp.git
commit 9ba3599257a847afa88f0723e194d84791f4a681 Author: Martin Zink <[email protected]> AuthorDate: Fri Sep 18 15:30:02 2026 +0200 Revert "cpp stable api enrichment" This reverts commit 57669e0a7ee79613fde39ec7374ab947a017daba. --- extensions/enrichment/CMakeLists.txt | 35 ------ extensions/enrichment/EnrichmentUtils.h | 23 ---- extensions/enrichment/ExtensionInitializer.cpp | 41 ------ extensions/enrichment/ForkEnrichment.cpp | 56 --------- extensions/enrichment/ForkEnrichment.h | 78 ------------ extensions/enrichment/JoinEnrichmentAttributes.cpp | 124 ------------------ extensions/enrichment/JoinEnrichmentAttributes.h | 138 -------------------- extensions/enrichment/tests/CMakeLists.txt | 40 ------ .../enrichment/tests/ForkEnrichmentTests.cpp | 70 ----------- .../tests/JoinEnrichmentAttributesTests.cpp | 140 --------------------- .../enrichment/tests/features/enrichment.feature | 47 ------- .../tests/features/enrichment_restart.feature | 71 ----------- .../enrichment/tests/features/environment.py | 28 ----- .../enrichment/tests/features/steps/steps.py | 21 ---- 14 files changed, 912 deletions(-) diff --git a/extensions/enrichment/CMakeLists.txt b/extensions/enrichment/CMakeLists.txt deleted file mode 100644 index 3588011c0..000000000 --- a/extensions/enrichment/CMakeLists.txt +++ /dev/null @@ -1,35 +0,0 @@ -# -# 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. -# - -if (NOT (ENABLE_ALL OR MINIFI_EXTENSION_ENRICHMENT)) - return() -endif() - -include(${CMAKE_SOURCE_DIR}/extensions/ExtensionHeader.txt) - -file(GLOB SOURCES "*.cpp") - -add_minifi_library(minifi-enrichment SHARED ${SOURCES}) -target_include_directories(minifi-enrichment PUBLIC "${CMAKE_SOURCE_DIR}/extensions/enrichment") -target_include_directories(minifi-enrichment PUBLIC "${ENRICHMENT_INCLUDE_DIRS}") - -include(Fetchstduuid) -target_link_libraries(minifi-enrichment minifi-cpp-extension-lib stduuid::stduuid) - -register_c_api_extension(minifi-enrichment "ENRICHMENT EXTENSION" ENRICHMENT-EXTENSION "Provides Fork/Join enrichment processors" "extensions/enrichment/tests") diff --git a/extensions/enrichment/EnrichmentUtils.h b/extensions/enrichment/EnrichmentUtils.h deleted file mode 100644 index 76dae0052..000000000 --- a/extensions/enrichment/EnrichmentUtils.h +++ /dev/null @@ -1,23 +0,0 @@ -/** - * 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 <string_view> - -constexpr std::string_view ENRICHMENT_ROLE = "enrichment.role"; -constexpr std::string_view ENRICHMENT_GROUP_ID = "enrichment.group.id"; diff --git a/extensions/enrichment/ExtensionInitializer.cpp b/extensions/enrichment/ExtensionInitializer.cpp deleted file mode 100644 index 19c26c071..000000000 --- a/extensions/enrichment/ExtensionInitializer.cpp +++ /dev/null @@ -1,41 +0,0 @@ -/** -* 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 "api/core/Resource.h" -#include "api/utils/minifi-c-utils.h" -#include "JoinEnrichmentAttributes.h" -#include "ForkEnrichment.h" - -#define MKSOC(x) #x -#define MAKESTRING(x) MKSOC(x) // NOLINT(cppcoreguidelines-macro-usage) - -namespace minifi = org::apache::nifi::minifi; - -CEXTENSIONAPI const uint32_t minifi_api_version = MINIFI_API_VERSION; - -CEXTENSIONAPI void minifi_init_extension(minifi_extension_context* extension_context) { - minifi_extension_definition extension_definition{ - .name = minifi::api::utils::minifiStringView(MAKESTRING(EXTENSION_NAME)), - .version = minifi::api::utils::minifiStringView(MAKESTRING(EXTENSION_VERSION)), - .group_name = minifi::api::utils::minifiStringView(MAKESTRING(MINIFI_EXTENSION_GROUP_NAME)), - .deinit = nullptr, - .user_data = nullptr - }; - auto* extension = minifi_register_extension(extension_context, &extension_definition); - minifi::api::core::registerProcessors<minifi::enrichment::ForkEnrichment>(extension); - minifi::api::core::registerProcessors<minifi::enrichment::JoinEnrichmentAttributes>(extension); -} diff --git a/extensions/enrichment/ForkEnrichment.cpp b/extensions/enrichment/ForkEnrichment.cpp deleted file mode 100644 index 52e35f460..000000000 --- a/extensions/enrichment/ForkEnrichment.cpp +++ /dev/null @@ -1,56 +0,0 @@ -/** - * 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 "ForkEnrichment.h" - -#include "api/core/ProcessSession.h" -#include "api/utils/ProcessorConfigUtils.h" -#include "stduuid/uuid.hpp" - -namespace org::apache::nifi::minifi::enrichment { - -minifi_status ForkEnrichment::onScheduleImpl(api::core::ProcessContext& context) { - max_batch_size_ = api::utils::parseOptionalU64Property(context, MaxBatchSize); - if (max_batch_size_ && *max_batch_size_ == 0) { - return MINIFI_STATUS_VALIDATION_FAILED; - } - - return MINIFI_STATUS_SUCCESS; -} - -minifi_status ForkEnrichment::onTriggerImpl(api::core::ProcessContext&, api::core::ProcessSession& session) { - uint64_t processed = 0; - while (api::core::FlowFile original = session.get()) { - api::core::FlowFile enrichment = session.clone(original); - - session.setAttribute(original, ENRICHMENT_ROLE, "ORIGINAL"); - session.setAttribute(enrichment, ENRICHMENT_ROLE, "ENRICHMENT"); - - const std::string group_id = uuids::to_string(uuid_gen()); - session.setAttribute(original, ENRICHMENT_GROUP_ID, group_id); - session.setAttribute(enrichment, ENRICHMENT_GROUP_ID, group_id); - - session.transfer(std::move(original), Original); - session.transfer(std::move(enrichment), Enrichment); - if (max_batch_size_ && ++processed >= max_batch_size_) { - break; - } - } - return MINIFI_STATUS_SUCCESS; -} - -} // namespace org::apache::nifi::minifi::enrichment diff --git a/extensions/enrichment/ForkEnrichment.h b/extensions/enrichment/ForkEnrichment.h deleted file mode 100644 index 312fb1a5d..000000000 --- a/extensions/enrichment/ForkEnrichment.h +++ /dev/null @@ -1,78 +0,0 @@ -/** - * 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 <random> - -#include "EnrichmentUtils.h" -#include "api/core/ProcessorImpl.h" -#include "api/utils/Export.h" -#include "core/PropertyDefinitionBuilder.h" -#include "minifi-cpp/core/Annotation.h" -#include "minifi-cpp/core/PropertyDefinition.h" -#include "stduuid/uuid.hpp" - -namespace org::apache::nifi::minifi::enrichment { - -class ForkEnrichment : public api::core::ProcessorImpl { - public: - using ProcessorImpl::ProcessorImpl; - - EXTENSIONAPI static constexpr const char* Description = - "Used in conjunction with the JoinEnrichmentAttributes processor, this processor is responsible for adding the attributes that are necessary " - "for the JoinEnrichmentAttributes processor to perform its function. Each incoming FlowFile will be cloned. The original FlowFile will have " - "appropriate attributes added and then be transferred to the 'original' relationship. The clone will have appropriate attributes added and " - "then be routed to the 'enrichment' relationship."; - - EXTENSIONAPI static constexpr auto MaxBatchSize = - core::PropertyDefinitionBuilder<>::createProperty("Max Batch Size") - .withDescription("The maximum number of flow files to process at a time. If unset, all FlowFiles will be processed at once.") - .withValidator(core::StandardPropertyValidators::UNSIGNED_INTEGER_VALIDATOR) - .build(); - - EXTENSIONAPI static constexpr auto Enrichment = core::RelationshipDefinition{"enrichment", - "A clone of the incoming FlowFile will be routed to this relationship, after adding appropriate attributes."}; - EXTENSIONAPI static constexpr auto Original = core::RelationshipDefinition{"original", - "The incoming FlowFile will be routed to this relationship, after adding appropriate attributes."}; - EXTENSIONAPI static constexpr auto Properties = std::array<core::PropertyReference, 1>{MaxBatchSize}; - EXTENSIONAPI static constexpr auto Relationships = std::array{Enrichment, Original}; - - EXTENSIONAPI static constexpr bool SupportsDynamicProperties = false; - EXTENSIONAPI static constexpr bool SupportsDynamicRelationships = false; - EXTENSIONAPI static constexpr core::annotation::Input InputRequirement = core::annotation::Input::INPUT_REQUIRED; - EXTENSIONAPI static constexpr bool IsSingleThreaded = false; - - EXTENSIONAPI static constexpr auto EnrichmentRole = core::OutputAttributeDefinition<2>{ - ENRICHMENT_ROLE, {Enrichment, Original}, "The role to use for enrichment. This will either be ORIGINAL or ENRICHMENT."}; - EXTENSIONAPI static constexpr auto EnrichmentGroupId = core::OutputAttributeDefinition<2>{ENRICHMENT_GROUP_ID, - {Enrichment, Original}, - "The Group ID to use in order to correlate the 'original' FlowFile with the 'enrichment' FlowFile."}; - - EXTENSIONAPI static constexpr auto OutputAttributes = std::array<core::OutputAttributeReference, 2>{EnrichmentRole, EnrichmentGroupId}; - - protected: - minifi_status onScheduleImpl(api::core::ProcessContext& context) override; - minifi_status onTriggerImpl(api::core::ProcessContext& context, api::core::ProcessSession& session) override; - - private: - std::optional<uint64_t> max_batch_size_; - std::random_device rd; - std::mt19937 rng{rd()}; - uuids::uuid_random_generator uuid_gen{rng}; -}; -} // namespace org::apache::nifi::minifi::enrichment diff --git a/extensions/enrichment/JoinEnrichmentAttributes.cpp b/extensions/enrichment/JoinEnrichmentAttributes.cpp deleted file mode 100644 index bed6e0237..000000000 --- a/extensions/enrichment/JoinEnrichmentAttributes.cpp +++ /dev/null @@ -1,124 +0,0 @@ -/** - * 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 "JoinEnrichmentAttributes.h" - -#include "EnrichmentUtils.h" -#include "api/core/Resource.h" -#include "api/utils/ProcessorConfigUtils.h" -#include "utils/AttributeErrors.h" - -namespace org::apache::nifi::minifi::enrichment { - -minifi_status JoinEnrichmentAttributes::onScheduleImpl(api::core::ProcessContext& context) { - using namespace std::literals::chrono_literals; - if (const auto timeout = api::utils::parseOptionalDurationProperty(context, TimeoutProperty); timeout && *timeout > 0ms) { - time_out_tracker_.emplace(*timeout); - } - max_batch_size_ = api::utils::parseOptionalU64Property(context, MaxBatchSize); - if (max_batch_size_ && *max_batch_size_ == 0) { - return MINIFI_STATUS_VALIDATION_FAILED; - } - return MINIFI_STATUS_SUCCESS; -} - -namespace { -bool checkRequiredAttributes(api::core::ProcessSession& session, api::core::FlowFile& flow_file) { - return session.getAttribute(flow_file, ENRICHMENT_ROLE).has_value() && session.getAttribute(flow_file, ENRICHMENT_GROUP_ID).has_value(); -} -} // namespace - -void JoinEnrichmentAttributes::join(api::core::FlowFile& original, api::core::FlowFile& enrichment, api::core::ProcessSession& session) { - api::core::FlowFile joined = session.clone(original); - for (const auto& [k, v] : session.getAttributes(enrichment)) { - if (k != ENRICHMENT_ROLE) { - session.setAttribute(joined, k, v); - } - } - session.setAttribute(joined, ENRICHMENT_ROLE, "JOINED"); - session.transfer(std::move(original), Original); - session.transfer(std::move(enrichment), Original); - session.transfer(std::move(joined), Joined); -} - -void JoinEnrichmentAttributes::handleFlowFile(api::core::FlowFile flow_file, api::core::ProcessSession& session, - const std::chrono::steady_clock::time_point current_time) { - if (!checkRequiredAttributes(session, flow_file)) { - logger_->log_warn("{} is missing enrichment.group.id and/or enrichment.role, routing it to Invalid", session.getFlowFileId(flow_file)); - session.transfer(std::move(flow_file), Invalid); - return; - } - - // SAFETY: checkRequiredAttributes already checks for ENRICHMENT_ROLE - const auto role = parsing::parseEnum<EnrichmentRole>(*session.getAttribute(flow_file, ENRICHMENT_ROLE)); - if (!role) { - logger_->log_warn("{} has invalid role due to {}", session.getFlowFileId(flow_file), role.error()); - session.transfer(std::move(flow_file), Invalid); - return; - } - - // SAFETY: checkRequiredAttributes already checks for ENRICHMENT_GROUP_ID - const std::string group_id = *session.getAttribute(flow_file, ENRICHMENT_GROUP_ID); - - auto& my_map = role == EnrichmentRole::ENRICHMENT ? enrichments_ : originals_; - auto& pair_map = role == EnrichmentRole::ENRICHMENT ? originals_ : enrichments_; - - if (const auto previous_node = my_map.extract(group_id)) { - logger_->log_warn("Encountered duplicate {} for {}, routing both to Invalid", magic_enum::enum_name(*role), group_id); - session.transfer(std::move(flow_file), Invalid); - session.transfer(session.unstash(std::move(previous_node.mapped())), Invalid); - return; - } - - if (const auto pair_node = pair_map.extract(group_id)) { - logger_->log_trace("Match found for {}", group_id); - api::core::FlowFile pair = session.unstash(std::move(pair_node.mapped())); - auto [original, enrichment] = - (*role == EnrichmentRole::ORIGINAL) ? std::tie(flow_file, pair) : std::tie(pair, flow_file); - join(original, enrichment, session); - } else { - my_map.insert({group_id, session.stash(std::move(flow_file))}); - if (time_out_tracker_) { - time_out_tracker_->track(std::move(group_id), current_time); - } - } -} - -minifi_status JoinEnrichmentAttributes::onTriggerImpl(api::core::ProcessContext&, api::core::ProcessSession& session) { - const auto current_time = std::chrono::steady_clock::now(); - uint64_t processed = 0; - while (auto flow_file = session.get()) { - handleFlowFile(std::move(flow_file), session, current_time); - if (max_batch_size_ && ++processed >= max_batch_size_) { - break; - } - } - - if (time_out_tracker_) { - for (auto timed_out_group : time_out_tracker_->getTimedOutFlowFiles(current_time)) { - const auto removeFromMap = [&](StoredFlowFileMap& map) { - if (auto timed_out_node = map.extract(timed_out_group)) { - session.transfer(session.unstash(std::move(timed_out_node.mapped())), TimeoutRelationship); - } - }; - removeFromMap(originals_); - removeFromMap(enrichments_); - } - } - return MINIFI_STATUS_SUCCESS; -} - -} // namespace org::apache::nifi::minifi::enrichment diff --git a/extensions/enrichment/JoinEnrichmentAttributes.h b/extensions/enrichment/JoinEnrichmentAttributes.h deleted file mode 100644 index e2c457407..000000000 --- a/extensions/enrichment/JoinEnrichmentAttributes.h +++ /dev/null @@ -1,138 +0,0 @@ -/** - * 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 <array> -#include <chrono> -#include <deque> -#include <functional> -#include <memory> -#include <optional> -#include <string> -#include <unordered_map> -#include <utility> -#include <vector> - -#include "EnrichmentUtils.h" -#include "api/core/ProcessorImpl.h" -#include "api/utils/Export.h" -#include "core/PropertyDefinitionBuilder.h" -#include "minifi-cpp/core/Annotation.h" -#include "minifi-cpp/core/PropertyDefinition.h" - -namespace org::apache::nifi::minifi::enrichment { - -namespace join_enrichment_attributes { -class TimeOutTracker { - public: - explicit TimeOutTracker(std::chrono::steady_clock::duration timeout) : time_out_(timeout) { - } - TimeOutTracker(const TimeOutTracker&) = delete; - TimeOutTracker& operator=(const TimeOutTracker&) = delete; - TimeOutTracker(TimeOutTracker&&) = delete; - TimeOutTracker& operator=(TimeOutTracker&&) = delete; - ~TimeOutTracker() = default; - - void track(std::string id, std::chrono::steady_clock::time_point timestamp) { - queue_.emplace_back(timestamp, std::move(id)); - } - - std::vector<std::string> getTimedOutFlowFiles(std::chrono::steady_clock::time_point current_time) { - std::vector<std::string> result; - // Even with 0 time_out_, we won't return just added FlowFiles - while (!queue_.empty() && queue_.front().timestamp + time_out_ < current_time) { - result.push_back(std::move(queue_.front().group_name)); - queue_.pop_front(); - } - return result; - } - - private: - struct TimeStampedGroup { - std::chrono::steady_clock::time_point timestamp; - std::string group_name; - }; - - std::chrono::steady_clock::duration time_out_; - std::deque<TimeStampedGroup> queue_; -}; -} // namespace join_enrichment_attributes - -using StoredFlowFileMap = std::unordered_map<std::string, api::core::StashedFlowFile, utils::string::transparent_string_hash, std::equal_to<>>; - -class JoinEnrichmentAttributes : public api::core::ProcessorImpl { - public: - using ProcessorImpl::ProcessorImpl; - - EXTENSIONAPI static constexpr const char* Description = - "Rejoins the forked FlowFiles coming from ForkEnrichment processor, the resulting FlowFile will have the Original's content and all attributes " - "from both of them (prioritizing Enrichment's)."; - - EXTENSIONAPI static constexpr auto Invalid = core::RelationshipDefinition{"invalid", - "Any FlowFiles without the requisite attributes will be routed here"}; - EXTENSIONAPI static constexpr auto Joined = core::RelationshipDefinition{"joined", - "The resultant FlowFile with Records joined together from both the original and enrichment FlowFiles will be routed to this relationship"}; - EXTENSIONAPI static constexpr auto Original = core::RelationshipDefinition{"original", - "Both of the incoming FlowFiles ('original' and 'enrichment') will be routed to this Relationship. I.e., this is the 'original' version of " - "both of these FlowFiles."}; - EXTENSIONAPI static constexpr auto TimeoutRelationship = core::RelationshipDefinition{"timeout", - "If one of the incoming FlowFiles (i.e., the 'original' FlowFile or the 'enrichment' FlowFile) arrives to this Processor but the other does " - "not arrive within the configured Timeout period, the FlowFile that did arrive is routed to this relationship."}; - - EXTENSIONAPI static constexpr auto MaxBatchSize = - core::PropertyDefinitionBuilder<>::createProperty("Max Batch Size") - .withDescription("The maximum number of flow files to process at a time. If unset, all FlowFiles will be processed at once.") - .withValidator(core::StandardPropertyValidators::UNSIGNED_INTEGER_VALIDATOR) - .build(); - - EXTENSIONAPI static constexpr auto TimeoutProperty = - core::PropertyDefinitionBuilder<>::createProperty("Timeout") - .withDescription( - "Specifies the maximum amount of time to wait for the second FlowFile once the first arrives at the processor, after which point the " - "first FlowFile will be routed to the 'timeout' relationship.") - .withValidator(core::StandardPropertyValidators::TIME_PERIOD_VALIDATOR) - .isRequired(false) - .build(); - - EXTENSIONAPI static constexpr auto Properties = std::array<core::PropertyReference, 2>{TimeoutProperty, MaxBatchSize}; - EXTENSIONAPI static constexpr auto Relationships = std::array{Invalid, Joined, Original, TimeoutRelationship}; - - EXTENSIONAPI static constexpr bool SupportsDynamicProperties = false; - EXTENSIONAPI static constexpr bool SupportsDynamicRelationships = false; - EXTENSIONAPI static constexpr auto InputRequirement = core::annotation::Input::INPUT_REQUIRED; - EXTENSIONAPI static constexpr bool IsSingleThreaded = true; - - protected: - minifi_status onScheduleImpl(api::core::ProcessContext& context) override; - minifi_status onTriggerImpl(api::core::ProcessContext& context, api::core::ProcessSession& session) override; - - private: - enum class EnrichmentRole { - ORIGINAL, - ENRICHMENT, - }; - - void handleFlowFile(api::core::FlowFile flow_file, api::core::ProcessSession& session, std::chrono::steady_clock::time_point current_time); - static void join(api::core::FlowFile& original, api::core::FlowFile& enrichment, api::core::ProcessSession& session); - - std::optional<join_enrichment_attributes::TimeOutTracker> time_out_tracker_; - StoredFlowFileMap originals_; - StoredFlowFileMap enrichments_; - std::optional<uint64_t> max_batch_size_; -}; -} // namespace org::apache::nifi::minifi::enrichment diff --git a/extensions/enrichment/tests/CMakeLists.txt b/extensions/enrichment/tests/CMakeLists.txt deleted file mode 100644 index 4629067d5..000000000 --- a/extensions/enrichment/tests/CMakeLists.txt +++ /dev/null @@ -1,40 +0,0 @@ -# -# 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. -# - -file(GLOB ENRICHMENT_TESTS "*.cpp") - -SET(EXTENSIONS_TEST_COUNT 0) -FOREACH(testfile ${ENRICHMENT_TESTS}) - get_filename_component(testfilename "${testfile}" NAME_WE) - add_minifi_executable(${testfilename} "${testfile}") - target_include_directories(${testfilename} BEFORE PRIVATE "${CMAKE_SOURCE_DIR}/libminifi/include") - target_include_directories(${testfilename} BEFORE PRIVATE "${CMAKE_SOURCE_DIR}/extensions/enrichment/processors") - createTests(${testfilename}) - target_link_libraries(${testfilename} Catch2WithMain) - target_link_libraries(${testfilename} minifi-enrichment) - target_link_libraries(${testfilename} minifi-standard-processors) - target_link_libraries(${testfilename} libminifi-c-unittest) - target_compile_definitions("${testfilename}" PRIVATE TZ_DATA_DIR="${CMAKE_BINARY_DIR}/tzdata") - target_compile_definitions(${testfilename} PRIVATE "MINIFI_EXTENSION_GROUP_NAME=org.apache.nifi.minifi.test") - - MATH(EXPR EXTENSIONS_TEST_COUNT "${EXTENSIONS_TEST_COUNT}+1") - add_test(NAME ${testfilename} COMMAND ${testfilename} WORKING_DIRECTORY ${TEST_DIR}) - set_tests_properties("${testfilename}" PROPERTIES LABELS "enrichment;memchecked") -ENDFOREACH() -message("-- Finished building ${EXTENSIONS_TEST_COUNT} enrichment related test file(s)...") diff --git a/extensions/enrichment/tests/ForkEnrichmentTests.cpp b/extensions/enrichment/tests/ForkEnrichmentTests.cpp deleted file mode 100644 index 320b15ec1..000000000 --- a/extensions/enrichment/tests/ForkEnrichmentTests.cpp +++ /dev/null @@ -1,70 +0,0 @@ -/** - * - * 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 "CProcessorTestUtils.h" -#include "ForkEnrichment.h" -#include "unit/Catch.h" -#include "unit/ProcessorUtils.h" -#include "unit/SingleProcessorTestController.h" - -namespace org::apache::nifi::minifi::enrichment::test { -TEST_CASE("Fork Enrichment processor works") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<ForkEnrichment>( - core::ProcessorMetadata{utils::Identifier{}, "ForkEnrichment", logging::LoggerFactory<ForkEnrichment>::getLogger()})); - const auto trigger_result = test_controller.trigger("test_content"); - REQUIRE(trigger_result.contains(ForkEnrichment::Original)); - REQUIRE(trigger_result.contains(ForkEnrichment::Enrichment)); - const auto original_results = trigger_result.at(ForkEnrichment::Original); - const auto enrichment_results = trigger_result.at(ForkEnrichment::Enrichment); - REQUIRE(original_results.size() == 1); - REQUIRE(enrichment_results.size() == 1); - - const auto original_content = test_controller.plan->getContent(original_results.at(0)); - const auto enrichment_content = test_controller.plan->getContent(enrichment_results.at(0)); - - CHECK(original_content == enrichment_content); - CHECK(original_content == "test_content"); - - CHECK(original_results.at(0)->getAttribute(ForkEnrichment::EnrichmentRole.name) == "ORIGINAL"); - CHECK(enrichment_results.at(0)->getAttribute(ForkEnrichment::EnrichmentRole.name) == "ENRICHMENT"); - - CHECK(original_results.at(0)->getAttribute(ForkEnrichment::EnrichmentGroupId.name) == - enrichment_results.at(0)->getAttribute(ForkEnrichment::EnrichmentGroupId.name)); -} - -TEST_CASE("ForkEnrichment no max batch size") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<ForkEnrichment>( - core::ProcessorMetadata{utils::Identifier{}, "ForkEnrichment", logging::LoggerFactory<ForkEnrichment>::getLogger()})); - const auto trigger_result = test_controller.trigger({{.content = "one"}, {.content = "two"}, {.content = "three"}}); - const auto original_results = trigger_result.at(ForkEnrichment::Original); - const auto enrichment_results = trigger_result.at(ForkEnrichment::Enrichment); - REQUIRE(original_results.size() == 3); - REQUIRE(enrichment_results.size() == 3); -} - -TEST_CASE("ForkEnrichment max batch size 2") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<ForkEnrichment>( - core::ProcessorMetadata{utils::Identifier{}, "ForkEnrichment", logging::LoggerFactory<ForkEnrichment>::getLogger()})); - const auto proc = test_controller.getProcessor(); - CHECK(test_controller.plan->setProperty(proc, ForkEnrichment::MaxBatchSize.name, "2")); - const auto trigger_result = test_controller.trigger({{.content = "one"}, {.content = "two"}, {.content = "three"}}); - const auto original_results = trigger_result.at(ForkEnrichment::Original); - const auto enrichment_results = trigger_result.at(ForkEnrichment::Enrichment); - REQUIRE(original_results.size() == 2); - REQUIRE(enrichment_results.size() == 2); -} -} // namespace org::apache::nifi::minifi::enrichment::test diff --git a/extensions/enrichment/tests/JoinEnrichmentAttributesTests.cpp b/extensions/enrichment/tests/JoinEnrichmentAttributesTests.cpp deleted file mode 100644 index 3ec640abd..000000000 --- a/extensions/enrichment/tests/JoinEnrichmentAttributesTests.cpp +++ /dev/null @@ -1,140 +0,0 @@ -/** - * - * 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 <thread> - -#include "CProcessorTestUtils.h" -#include "EnrichmentUtils.h" -#include "JoinEnrichmentAttributes.h" -#include "unit/Catch.h" -#include "unit/ProcessorUtils.h" -#include "unit/SingleProcessorTestController.h" - -namespace org::apache::nifi::minifi::enrichment::test { -TEST_CASE("JoinEnrichmentAttributes input without appropriate attributes") { - minifi::test::SingleProcessorTestController controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto trigger_result = controller.trigger("test_content"); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Invalid).size() == 1); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Original).empty()); - CHECK(trigger_result.at(JoinEnrichmentAttributes::TimeoutRelationship).empty()); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Joined).empty()); -} - -TEST_CASE("JoinEnrichmentAttributes invalid role") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto trigger_result = test_controller.trigger(minifi::test::InputFlowFileData{.content = "first", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "FIREBIRD"}}}); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Invalid).size() == 1); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Original).empty()); - CHECK(trigger_result.at(JoinEnrichmentAttributes::TimeoutRelationship).empty()); - CHECK(trigger_result.at(JoinEnrichmentAttributes::Joined).empty()); -} - -TEST_CASE("JoinEnrichmentAttributes same id same role same session") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto trigger = test_controller.trigger( - {minifi::test::InputFlowFileData{.content = "first", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}}}, - minifi::test::InputFlowFileData{.content = "second", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}}}}); - CHECK(trigger.at(JoinEnrichmentAttributes::Invalid).size() == 2); - CHECK(trigger.at(JoinEnrichmentAttributes::Original).empty()); - CHECK(trigger.at(JoinEnrichmentAttributes::TimeoutRelationship).empty()); - CHECK(trigger.at(JoinEnrichmentAttributes::Joined).empty()); -} - -TEST_CASE("JoinEnrichmentAttributes same id same role different session") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto first_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "first", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}}}); - // First trigger no output (it holds the original waiting for its pair) - CHECK(std::ranges::all_of(first_trigger, [](const auto& res) -> bool { return res.second.empty(); })); - - const auto second_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "second", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}}}); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Invalid).size() == 2); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Original).empty()); - CHECK(second_trigger.at(JoinEnrichmentAttributes::TimeoutRelationship).empty()); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Joined).empty()); -} - -TEST_CASE("JoinEnrichmentAttributes same id diff role different session") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto first_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "first", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}, {"first_attr", "1"}}}); - // First trigger no output (it holds the original waiting for its pair) - CHECK(std::ranges::all_of(first_trigger, [](const auto& res) -> bool { return res.second.empty(); })); - - const auto second_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "second", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ENRICHMENT"}, {"second_attr", "2"}}}); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Original).size() == 2); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Invalid).empty()); - CHECK(second_trigger.at(JoinEnrichmentAttributes::TimeoutRelationship).empty()); - REQUIRE(second_trigger.at(JoinEnrichmentAttributes::Joined).size() == 1); - - const auto joined_content = test_controller.plan->getContent(second_trigger.at(JoinEnrichmentAttributes::Joined).at(0)); - const auto joined_attrs = second_trigger.at(JoinEnrichmentAttributes::Joined).at(0)->getAttributes(); - - CHECK(joined_content == "first"); - CHECK(joined_attrs.at(std::string{ENRICHMENT_GROUP_ID}) == "foo"); - CHECK(joined_attrs.at(std::string{ENRICHMENT_ROLE}) == "JOINED"); - CHECK(joined_attrs.at("first_attr") == "1"); - CHECK(joined_attrs.at("second_attr") == "2"); -} - -TEST_CASE("JoinEnrichmentAttributes test timeout") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto proc = test_controller.getProcessor(); - CHECK(test_controller.plan->setProperty(proc, JoinEnrichmentAttributes::TimeoutProperty.name, "1 ms")); - - const auto first_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "first", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "foo"}, {std::string{ENRICHMENT_ROLE}, "ORIGINAL"}, {"first_attr", "1"}}}); - // First trigger no output (it holds the original waiting for its pair) - CHECK(std::ranges::all_of(first_trigger, [](const auto& res) -> bool { return res.second.empty(); })); - - std::this_thread::sleep_for(1ms); - const auto second_trigger = test_controller.trigger(minifi::test::InputFlowFileData{.content = "second", - .attributes = {{std::string{ENRICHMENT_GROUP_ID}, "bar"}, {std::string{ENRICHMENT_ROLE}, "ENRICHMENT"}, {"second_attr", "2"}}}); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Original).empty()); - CHECK(second_trigger.at(JoinEnrichmentAttributes::Invalid).empty()); - CHECK(second_trigger.at(JoinEnrichmentAttributes::TimeoutRelationship).size() == 1); - REQUIRE(second_trigger.at(JoinEnrichmentAttributes::Joined).empty()); -} - -TEST_CASE("JoinEnrichmentAttributes no max batch size") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto trigger_result = test_controller.trigger({{.content = "one"}, {.content = "two"}, {.content = "three"}}); - REQUIRE(trigger_result.at(JoinEnrichmentAttributes::Invalid).size() == 3); -} - -TEST_CASE("JoinEnrichmentAttributes max batch size 2") { - minifi::test::SingleProcessorTestController test_controller(minifi::test::utils::make_custom_c_processor<JoinEnrichmentAttributes>( - core::ProcessorMetadata{utils::Identifier{}, "JoinEnrichmentAttributes", logging::LoggerFactory<JoinEnrichmentAttributes>::getLogger()})); - const auto proc = test_controller.getProcessor(); - CHECK(test_controller.plan->setProperty(proc, JoinEnrichmentAttributes::MaxBatchSize.name, "2")); - const auto trigger_result = test_controller.trigger({{.content = "one"}, {.content = "two"}, {.content = "three"}}); - REQUIRE(trigger_result.at(JoinEnrichmentAttributes::Invalid).size() == 2); -} -} // namespace org::apache::nifi::minifi::enrichment::test diff --git a/extensions/enrichment/tests/features/enrichment.feature b/extensions/enrichment/tests/features/enrichment.feature deleted file mode 100644 index 1401abc46..000000000 --- a/extensions/enrichment/tests/features/enrichment.feature +++ /dev/null @@ -1,47 +0,0 @@ -# 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. - -@MINIFI_EXTENSION_ENRICHMENT @SUPPORTS_WINDOWS -Feature: ForkEnrichment and JoinEnrichmentAttributes - - Scenario: Merges correctly - Given a GenerateFlowFile processor with the "Custom Text" property set to "original_${literal("content")}" - And the scheduling period of the GenerateFlowFile processor is set to "1 hour" - And the "Data Format" property of the GenerateFlowFile processor is set to "Text" - And the "Unique FlowFiles" property of the GenerateFlowFile processor is set to "false" - - And a ForkEnrichment processor - And a JoinEnrichmentAttributes processor - - And a ReplaceText processor with the "Evaluation Mode" property set to "Entire text" - And the "Replacement Strategy" property of the ReplaceText processor is set to "Always Replace" - And the "Replacement Value" property of the ReplaceText processor is set to "replaced_content" - - And an UpdateAttribute processor with the "extra_prop" property set to "foo" - - And a LogAttribute processor with the "Log Payload" property set to "true" - - And the "success" relationship of the GenerateFlowFile processor is connected to the ForkEnrichment - And the "original" relationship of the ForkEnrichment processor is connected to the JoinEnrichmentAttributes - And the "enrichment" relationship of the ForkEnrichment processor is connected to the ReplaceText - And the "success" relationship of the ReplaceText processor is connected to the UpdateAttribute - And the "success" relationship of the UpdateAttribute processor is connected to the JoinEnrichmentAttributes - And the "joined" relationship of the JoinEnrichmentAttributes processor is connected to the LogAttribute - And JoinEnrichmentAttributes's original relationship is auto-terminated - And LogAttribute's success relationship is auto-terminated - When the MiNiFi instance starts up - Then the Minifi logs contain the following message: "key:enrichment.role value:JOINED" in less than 10 seconds - And the Minifi logs contain the following message: "key:extra_prop value:foo" in less than 1 second - And the Minifi logs contain the following message: "original_content" in less than 1 second diff --git a/extensions/enrichment/tests/features/enrichment_restart.feature b/extensions/enrichment/tests/features/enrichment_restart.feature deleted file mode 100644 index 36d329101..000000000 --- a/extensions/enrichment/tests/features/enrichment_restart.feature +++ /dev/null @@ -1,71 +0,0 @@ -# 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. - -@MINIFI_EXTENSION_ENRICHMENT -Feature: JoinEnrichmentAttributes recovers released FlowFiles after an agent restart - - # JoinEnrichmentAttributes releases the first half of a pair and holds it in memory until the - # other half arrives. The released FlowFile is not routed, persisted-as-moved, nor deleted on - # commit, so its record stays in the (default, RocksDB-backed) flow file repository. When the - # agent is restarted the held FlowFile must be recovered from disk and re-enqueued, so that the - # join still completes once its pair finally arrives. - # - # The two halves are produced independently with GetFile + UpdateAttribute (instead of - # ForkEnrichment) so their arrival can be separated across the restart: the ORIGINAL is present - # at startup, and its ENRICHMENT pair is only delivered after the restart. - - Scenario: A released FlowFile held across an agent restart is recovered and joined with its late-arriving pair - Given a GetFile processor with the name "GetOriginal" and the "Input Directory" property set to "/tmp/original_input" - And the scheduling period of the GetOriginal processor is set to "100 ms" - And a UpdateAttribute processor with the name "TagOriginal" and the "enrichment.role" property set to "ORIGINAL" - And the "enrichment.group.id" property of the TagOriginal processor is set to "group-1" - - And a GetFile processor with the name "GetEnrichment" and the "Input Directory" property set to "/tmp/enrichment_input" - And the scheduling period of the GetEnrichment processor is set to "100 ms" - And a UpdateAttribute processor with the name "TagEnrichment" and the "enrichment.role" property set to "ENRICHMENT" - And the "enrichment.group.id" property of the TagEnrichment processor is set to "group-1" - - And a JoinEnrichmentAttributes processor - And a PutFile processor with the "Directory" property set to "/tmp/output" - - And the "success" relationship of the GetOriginal processor is connected to the TagOriginal - And the "success" relationship of the TagOriginal processor is connected to the JoinEnrichmentAttributes - And the "success" relationship of the GetEnrichment processor is connected to the TagEnrichment - And the "success" relationship of the TagEnrichment processor is connected to the JoinEnrichmentAttributes - And the "joined" relationship of the JoinEnrichmentAttributes processor is connected to the PutFile - - And JoinEnrichmentAttributes's original relationship is auto-terminated - And JoinEnrichmentAttributes's invalid relationship is auto-terminated - And JoinEnrichmentAttributes's timeout relationship is auto-terminated - And PutFile's success relationship is auto-terminated - And PutFile's failure relationship is auto-terminated - - # Only the ORIGINAL half exists at startup; the ENRICHMENT half arrives after the restart. - And a directory at "/tmp/original_input" has a file with the content "original_content" - - When the MiNiFi instance starts up - # JoinEnrichmentAttributes gets the ORIGINAL, releases it and holds it waiting for its pair, - # so nothing is joined yet. - Then no files are placed in the "/tmp/output" directory in 5 seconds of running time - - # Graceful stop destroys JoinEnrichmentAttributes while it still holds the released ORIGINAL, - # then restart brings the agent back with the persistent repositories intact. - When MiNiFi is stopped - And MiNiFi is restarted - - # The released ORIGINAL is recovered from the flow file repository and re-enqueued. Once its - # ENRICHMENT pair arrives, the join produces a single FlowFile with the ORIGINAL's content. - And a file with the content "enrichment_content" is placed in "/tmp/enrichment_input" - Then a single file with the content "original_content" is placed in the "/tmp/output" directory in less than 60 seconds diff --git a/extensions/enrichment/tests/features/environment.py b/extensions/enrichment/tests/features/environment.py deleted file mode 100644 index 079c95b50..000000000 --- a/extensions/enrichment/tests/features/environment.py +++ /dev/null @@ -1,28 +0,0 @@ -# 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. - -from minifi_behave.core.hooks import ( - common_after_scenario, - common_before_scenario, -) -from minifi_behave.core.minifi_test_context import MinifiTestContext - - -def before_scenario(context: MinifiTestContext, scenario): - common_before_scenario(context, scenario) - - -def after_scenario(context, scenario): - common_after_scenario(context, scenario) diff --git a/extensions/enrichment/tests/features/steps/steps.py b/extensions/enrichment/tests/features/steps/steps.py deleted file mode 100644 index ff5817598..000000000 --- a/extensions/enrichment/tests/features/steps/steps.py +++ /dev/null @@ -1,21 +0,0 @@ -# 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. - -from minifi_behave.steps import ( - checking_steps, # noqa: F401 - configuration_steps, # noqa: F401 - core_steps, # noqa: F401 - flow_building_steps, # noqa: F401 -)
