Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 1 addition & 11 deletions core-framework/include/core/ProcessorImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -42,19 +42,11 @@
#include "minifi-cpp/core/ProcessorMetadata.h"
#include "minifi-cpp/Exception.h"

#define ADD_GET_PROCESSOR_NAME \
std::string getProcessorType() const override { \
auto class_name = org::apache::nifi::minifi::core::className<decltype(*this)>(); \
auto splitted = org::apache::nifi::minifi::utils::string::split(class_name, "::"); \
return splitted[splitted.size() - 1]; \
}

#define ADD_COMMON_VIRTUAL_FUNCTIONS_FOR_PROCESSORS \
bool supportsDynamicProperties() const override { return SupportsDynamicProperties; } \
bool supportsDynamicRelationships() const override { return SupportsDynamicRelationships; } \
minifi::core::annotation::Input getInputRequirement() const override { return InputRequirement; } \
bool isSingleThreaded() const override { return IsSingleThreaded; } \
ADD_GET_PROCESSOR_NAME
bool isSingleThreaded() const override { return IsSingleThreaded; }

namespace org::apache::nifi::minifi {

Expand Down Expand Up @@ -85,8 +77,6 @@ class ProcessorImpl : public virtual ProcessorApi {

[[nodiscard]] bool supportsDynamicRelationships() const override = 0;

std::string getProcessorType() const override = 0;

void initialize(ProcessorDescriptor& self) final;

void setSupportedRelationships(std::span<const RelationshipDefinition> relationships);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,19 +25,22 @@
#include "utils/CControllerService.h"
#include "utils/CProcessor.h"
#include "minifi-cpp/agent/agent_docs.h"
#include "utils/StringUtils.h"

namespace org::apache::nifi::minifi::test::utils {

template<typename T, typename... Args>
std::unique_ptr<minifi::core::Processor> make_custom_c_processor(minifi::core::ProcessorMetadata metadata,
Args&&... args) { // NOLINT(cppcoreguidelines-missing-std-forward)
std::string type;
std::unique_ptr<minifi::core::ProcessorApi> processor_impl;
minifi::api::core::useProcessorClassDefinition<T>([&](const minifi_processor_definition& definition) {
minifi::utils::useCProcessorClassDescription(definition, [&](const auto&, auto c_description) {
minifi::utils::useCProcessorClassDescription(definition, [&](const auto& description, auto c_description) {
type = description.short_name_;
processor_impl = std::make_unique<minifi::utils::CProcessor>(std::move(c_description), metadata, new T(metadata, std::forward<Args>(args)...));
});
});
return std::make_unique<minifi::core::Processor>(metadata.name, metadata.uuid, std::move(processor_impl));
return std::make_unique<minifi::core::Processor>(std::move(type), metadata.name, metadata.uuid, std::move(processor_impl));
}

template<typename T, typename... Args>
Expand Down
3 changes: 0 additions & 3 deletions extension-framework/include/core/AbstractProcessor.h
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,5 @@ class AbstractProcessor : public ProcessorImpl {
bool supportsDynamicRelationships() const noexcept final { return ProcessorT::SupportsDynamicRelationships; }
minifi::core::annotation::Input getInputRequirement() const noexcept final { return ProcessorT::InputRequirement; }
bool isSingleThreaded() const noexcept final { return ProcessorT::IsSingleThreaded; }
std::string getProcessorType() const final {
return utils::string::partAfterLastOccurrenceOf(className<ProcessorT>(), ':');
}
};
} // namespace org::apache::nifi::minifi::core
4 changes: 2 additions & 2 deletions extensions/aws/tests/S3TestsFixture.h
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ class FlowProcessorS3TestsFixture : public S3TestsFixture<T> {
this->mock_s3_request_sender->setUseVirtualAddressing(use_virtual_addressing);
return std::make_unique<minifi::aws::s3::S3Wrapper>(std::move(this->mock_s3_request_sender));
}));
auto s3_processor_unique_ptr = std::make_unique<core::Processor>("S3Processor", uuid, std::move(impl));
auto s3_processor_unique_ptr = std::make_unique<core::Processor>(utils::string::partAfterLastOccurrenceOf(core::className<T>(), ':'), "S3Processor", uuid, std::move(impl));
this->s3_processor = s3_processor_unique_ptr.get();

auto input_dir = this->test_controller.createTempDirectory();
Expand Down Expand Up @@ -226,7 +226,7 @@ class FlowProducerS3TestsFixture : public S3TestsFixture<T> {
this->mock_s3_request_sender->setUseVirtualAddressing(use_virtual_addressing);
return std::make_unique<minifi::aws::s3::S3Wrapper>(std::move(this->mock_s3_request_sender));
}));
auto s3_processor_unique_ptr = std::make_unique<core::Processor>("S3Processor", uuid, std::move(impl));
auto s3_processor_unique_ptr = std::make_unique<core::Processor>(utils::string::partAfterLastOccurrenceOf(core::className<T>(), ':'), "S3Processor", uuid, std::move(impl));
this->s3_processor = s3_processor_unique_ptr.get();

this->plan->addProcessor(
Expand Down
3 changes: 2 additions & 1 deletion extensions/azure/tests/AzureBlobStorageTestsFixture.h
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,8 @@ class AzureBlobStorageTestsFixture {
auto uuid = utils::IdGenerator::getIdGenerator()->generate();
auto impl = std::unique_ptr<ProcessorType>(
new ProcessorType({.uuid = uuid, .name = "AzureBlobStorageProcessor", .logger = logging::LoggerFactory<ProcessorType>::getLogger(uuid)}, std::move(mock_blob_storage)));
auto azure_blob_storage_processor_unique_ptr = std::make_unique<core::Processor>(impl->getName(), impl->getUUID(), std::move(impl));
auto azure_blob_storage_processor_unique_ptr = std::make_unique<core::Processor>(
utils::string::partAfterLastOccurrenceOf(core::className<ProcessorType>(), ':'), impl->getName(), impl->getUUID(), std::move(impl));
azure_blob_storage_processor_ = azure_blob_storage_processor_unique_ptr.get();
auto input_dir = test_controller_.createTempDirectory();
std::ofstream input_file_stream(input_dir / GET_FILE_NAME);
Expand Down
3 changes: 2 additions & 1 deletion extensions/azure/tests/AzureDataLakeStorageTestsFixture.h
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,8 @@ class AzureDataLakeStorageTestsFixture {
new AzureDataLakeStorageProcessor({
.uuid = uuid, .name = "AzureDataLakeStorageProcessor",
.logger = logging::LoggerFactory<AzureDataLakeStorageProcessor>::getLogger(uuid)}, std::move(mock_data_lake_storage_client)));
auto azure_data_lake_storage_unique_ptr = std::make_unique<core::Processor>(impl->getName(), impl->getUUID(), std::move(impl));
auto azure_data_lake_storage_unique_ptr = std::make_unique<core::Processor>(
utils::string::partAfterLastOccurrenceOf(core::className<AzureDataLakeStorageProcessor>(), ':'), impl->getName(), impl->getUUID(), std::move(impl));
azure_data_lake_storage_ = azure_data_lake_storage_unique_ptr.get();
auto input_dir = test_controller_.createTempDirectory();
minifi::test::utils::putFileToDir(input_dir, GETFILE_FILE_NAME, TEST_DATA);
Expand Down
2 changes: 1 addition & 1 deletion extensions/azure/tests/ListAzureBlobStorageTests.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ class ListAzureBlobStorageTestsFixture {
core::ProcessorMetadata{
.uuid = uuid, .name = "ListAzureBlobStorage",
.logger = logging::LoggerFactory<minifi::azure::processors::ListAzureBlobStorage>::getLogger(uuid)}, std::move(mock_blob_storage));
auto list_azure_blob_storage_unique_ptr = std::make_unique<core::Processor>(impl->getName(), impl->getUUID(), std::move(impl));
auto list_azure_blob_storage_unique_ptr = std::make_unique<core::Processor>("ListAzureBlobStorage", impl->getName(), impl->getUUID(), std::move(impl));
list_azure_blob_storage_ = list_azure_blob_storage_unique_ptr.get();

plan_->addProcessor(std::move(list_azure_blob_storage_unique_ptr), "ListAzureBlobStorage", { {"success", "d"} });
Expand Down
2 changes: 1 addition & 1 deletion extensions/azure/tests/ListAzureDataLakeStorageTests.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ class ListAzureDataLakeStorageTestsFixture {
new minifi::azure::processors::ListAzureDataLakeStorage({
.uuid = uuid, .name = "ListAzureDataLakeStorage",
.logger = logging::LoggerFactory<minifi::azure::processors::ListAzureDataLakeStorage>::getLogger(uuid)}, std::move(mock_data_lake_storage_client)));
auto list_azure_data_lake_storage_unique_ptr = std::make_unique<core::Processor>(impl->getName(), impl->getUUID(), std::move(impl));
auto list_azure_data_lake_storage_unique_ptr = std::make_unique<core::Processor>("ListAzureDataLakeStorage", impl->getName(), impl->getUUID(), std::move(impl));
list_azure_data_lake_storage_ = list_azure_data_lake_storage_unique_ptr.get();

plan_->addProcessor(std::move(list_azure_data_lake_storage_unique_ptr), "ListAzureDataLakeStorage", { {"success", "d"} });
Expand Down
1 change: 0 additions & 1 deletion extensions/python/ExecutePythonProcessor.h
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,6 @@ class ExecutePythonProcessor : public core::ProcessorImpl {
bool supportsDynamicRelationships() const override { return SupportsDynamicRelationships; }
minifi::core::annotation::Input getInputRequirement() const override { return InputRequirement; }
bool isSingleThreaded() const override { return python_single_threaded_; }
ADD_GET_PROCESSOR_NAME

void initialize() override;
void onSchedule(core::ProcessContext& context, core::ProcessSessionFactory&) override;
Expand Down
3 changes: 2 additions & 1 deletion extensions/python/tests/ExecutePythonProcessorTests.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -152,7 +152,8 @@ class SimplePythonFlowFileTransferTest : public ExecutePythonProcessorTestBase {
.logger = logging::LoggerFactory<minifi::extensions::python::processors::ExecutePythonProcessor>::getLogger(uuid)
});
execute_python_processor->setScriptFilePath(getScriptFullPath(used_as_script_file).string());
auto execute_python_processor_unique_ptr = std::make_unique<core::Processor>(execute_python_processor->getName(), execute_python_processor->getUUID(), std::move(execute_python_processor));
auto execute_python_processor_unique_ptr = std::make_unique<core::Processor>(
"ExecutePythonProcessor", execute_python_processor->getName(), execute_python_processor->getUUID(), std::move(execute_python_processor));
auto processor = plan_->addProcessor(std::move(execute_python_processor_unique_ptr), "executePythonProcessor", core::Relationship("success", "description"), link_to_previous);
return processor;
}
Expand Down
4 changes: 2 additions & 2 deletions libminifi/include/Port.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,9 +50,9 @@ class PortImpl final : public ForwardingNode {
PortType port_type_;
};

class Port : public core::Processor {
class Port final : public core::Processor {
public:
Port(std::string_view name, const utils::Identifier& uuid, std::unique_ptr<PortImpl> impl): Processor(name, uuid, std::move(impl)) {}
Port(std::string_view name, const utils::Identifier& uuid, std::unique_ptr<PortImpl> impl): Processor("Port", name, uuid, std::move(impl)) {}

PortType getPortType() const {
auto* port_impl = dynamic_cast<const PortImpl*>(impl_.get());
Expand Down
4 changes: 2 additions & 2 deletions libminifi/include/core/Processor.h
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,7 @@ class ProcessSessionFactory;

class Processor : public ConnectableImpl, public ConfigurableComponentImpl, public state::response::ResponseNodeSource {
public:
Processor(std::string_view name, const utils::Identifier& uuid, std::unique_ptr<ProcessorApi> impl);
explicit Processor(std::string_view name, std::unique_ptr<ProcessorApi> impl);
Processor(std::string type, std::string_view name, const utils::Identifier& uuid, std::unique_ptr<ProcessorApi> impl);

Processor(const Processor& parent) = delete;
Processor& operator=(const Processor& parent) = delete;
Expand Down Expand Up @@ -141,6 +140,7 @@ class Processor : public ConnectableImpl, public ConfigurableComponentImpl, publ
}

protected:
std::string type_;
std::atomic<ScheduledState> state_;

std::atomic<std::chrono::steady_clock::duration> scheduling_period_;
Expand Down
4 changes: 0 additions & 4 deletions libminifi/include/utils/CProcessor.h
Original file line number Diff line number Diff line change
Expand Up @@ -110,10 +110,6 @@ class CProcessor : public minifi::core::ProcessorApi {
return class_description_.is_single_threaded;
}

std::string getProcessorType() const override {
return class_description_.name;
}

void onTrigger(minifi::core::ProcessContext& process_context, minifi::core::ProcessSession& process_session) override {
std::optional<std::string> error;
auto status = class_description_.callbacks.trigger(impl_, reinterpret_cast<minifi_process_context*>(&process_context), reinterpret_cast<minifi_process_session*>(&process_session));
Expand Down
3 changes: 2 additions & 1 deletion libminifi/src/core/ClassLoader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,8 @@ class ProcessorFactoryWrapper : public ObjectFactoryImpl {

[[nodiscard]] gsl::owner<CoreComponent*> createRaw(const std::string &name, const utils::Identifier &uuid) override {
auto logger = logging::LoggerFactoryBase::getAliasedLogger(getClassName(), uuid);
return gsl::owner<CoreComponent*>{new Processor(name, uuid, factory_->create({.uuid = uuid, .name = name, .logger = std::move(logger)}))};
return gsl::owner<CoreComponent*>{new Processor(
utils::string::partAfterLastOccurrenceOf(factory_->getClassName(), ':'), name, uuid, factory_->create({.uuid = uuid, .name = name, .logger = std::move(logger)}))};
}

[[nodiscard]] std::string getModuleName() const override {
Expand Down
5 changes: 4 additions & 1 deletion libminifi/src/core/FlowConfiguration.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,10 @@ std::unique_ptr<core::Processor> FlowConfiguration::createProcessor(const std::s
}

std::unique_ptr<core::Processor> FlowConfiguration::createProvenanceReportTask() {
auto processor = std::make_unique<core::Processor>("", std::make_unique<core::reporting::SiteToSiteProvenanceReportingTask>(this->configuration_));
auto impl = std::make_unique<core::reporting::SiteToSiteProvenanceReportingTask>(this->configuration_);
auto uuid = impl->getUUID();
auto name = impl->getName();
auto processor = std::make_unique<core::Processor>("SiteToSiteProvenanceReportingTask", name, uuid, std::move(impl));
processor->initialize();
return processor;
}
Expand Down
25 changes: 3 additions & 22 deletions libminifi/src/core/Processor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -52,28 +52,9 @@ static std::mutex& getGraphMutex() {
return mutex;
}

Processor::Processor(std::string_view name, std::unique_ptr<ProcessorApi> impl)
: ConnectableImpl(name),
state_(DISABLED),
scheduling_period_(MINIMUM_SCHEDULING_PERIOD),
run_duration_(DEFAULT_RUN_DURATION),
yield_period_(DEFAULT_YIELD_PERIOD_SECONDS),
active_tasks_(0),
trigger_when_empty_(false),
logger_(logging::LoggerFactory<Processor>::getLogger(uuid_)),
metrics_(gsl::make_not_null(std::make_shared<ProcessorMetrics>(*this))),
impl_(std::move(impl)) {
has_work_.store(false);
// Setup the default values
strategy_ = TIMER_DRIVEN;
penalization_period_ = DEFAULT_PENALIZATION_PERIOD;
max_concurrent_tasks_ = DEFAULT_MAX_CONCURRENT_TASKS;
incoming_connections_Iter = this->incoming_connections_.begin();
logger_->log_debug("Processor {} created UUID {}", name_, getUUIDStr());
}

Processor::Processor(std::string_view name, const utils::Identifier& uuid, std::unique_ptr<ProcessorApi> impl)
Processor::Processor(std::string type, std::string_view name, const utils::Identifier& uuid, std::unique_ptr<ProcessorApi> impl)
: ConnectableImpl(name, uuid),
type_(std::move(type)),
state_(DISABLED),
scheduling_period_(MINIMUM_SCHEDULING_PERIOD),
run_duration_(DEFAULT_RUN_DURATION),
Expand Down Expand Up @@ -447,7 +428,7 @@ bool Processor::isSingleThreaded() const {
}

std::string Processor::getProcessorType() const {
return impl_->getProcessorType();
return type_;
}

void Processor::setTriggerWhenEmpty(bool trigger_when_empty) {
Expand Down
4 changes: 2 additions & 2 deletions libminifi/src/core/flow/StructuredConfiguration.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -741,7 +741,7 @@ void StructuredConfiguration::parseRPGPort(const Node& port_node, core::ProcessG
auto port_impl = std::make_unique<minifi::RemoteProcessGroupPort>(
nameStr, parent->getURL(), this->configuration_, uuid, direction);
auto* port = port_impl.get();
auto port_wrapper = std::make_unique<core::Processor>(nameStr, uuid, std::move(port_impl));
auto port_wrapper = std::make_unique<core::Processor>("RemoteProcessGroupPort", nameStr, uuid, std::move(port_impl));
port->setTimeout(std::chrono::milliseconds(parent->getTimeout()));
port->setTransmitting(true);
port_wrapper->setYieldPeriodMsec(parent->getYieldPeriodMsec());
Expand Down Expand Up @@ -986,7 +986,7 @@ void StructuredConfiguration::parseFunnels(const Node& node, core::ProcessGroup*
throw Exception(ExceptionType::GENERAL_EXCEPTION, "Incorrect funnel UUID format.");
});

auto funnel = std::make_unique<core::Processor>(name, uuid.value(), std::make_unique<minifi::Funnel>(name, uuid.value()));
auto funnel = std::make_unique<core::Processor>("Funnel", name, uuid.value(), std::make_unique<minifi::Funnel>(name, uuid.value()));
logger_->log_debug("Created funnel with UUID {} and name {}", id, name);
funnel->setScheduledState(core::RUNNING);
funnel->setSchedulingStrategy(core::EVENT_DRIVEN);
Expand Down
5 changes: 3 additions & 2 deletions libminifi/test/libtest/unit/ProcessorUtils.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

#include "core/Processor.h"
#include "minifi-cpp/core/ProcessorMetadata.h"
#include "utils/StringUtils.h"

namespace org::apache::nifi::minifi::test::utils {

Expand All @@ -32,15 +33,15 @@ std::unique_ptr<core::Processor> make_processor(std::string_view name, std::opti
.name = std::string{name},
.logger = minifi::core::logging::LoggerFactory<T>::getLogger(uuid.value())
});
return std::make_unique<core::Processor>(name, uuid.value(), std::move(processor_impl));
return std::make_unique<core::Processor>(minifi::utils::string::partAfterLastOccurrenceOf(core::className<T>(), ':'), name, uuid.value(), std::move(processor_impl));
}

template<typename T, typename ...Args>
std::unique_ptr<core::Processor> make_custom_processor(Args&&... args) {
auto processor_impl = std::make_unique<T>(std::forward<Args>(args)...);
auto name = processor_impl->getName();
auto uuid = processor_impl->getUUID();
return std::make_unique<core::Processor>(name, uuid, std::move(processor_impl));
return std::make_unique<core::Processor>(minifi::utils::string::partAfterLastOccurrenceOf(core::className<T>(), ':'), name, uuid, std::move(processor_impl));
}

} // namespace org::apache::nifi::minifi::test::utils
1 change: 0 additions & 1 deletion minifi-api/include/minifi-cpp/core/ProcessorApi.h
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,6 @@ class ProcessorApi {

virtual void initialize(ProcessorDescriptor& descriptor) = 0;
virtual bool isSingleThreaded() const = 0;
virtual std::string getProcessorType() const = 0;
virtual void onTrigger(ProcessContext&, ProcessSession&) = 0;
virtual void onSchedule(ProcessContext&, ProcessSessionFactory&) = 0;
virtual void onUnSchedule() = 0;
Expand Down
Loading