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
The following commit(s) were added to refs/heads/main by this push:
new f2fd698 MINIFICPP-1449 Add pause and resume message to C2
functionality
f2fd698 is described below
commit f2fd698504c0cba9e024d4a3a5868a374ba18b5e
Author: Gabor Gyimesi <[email protected]>
AuthorDate: Tue Jan 26 18:31:35 2021 +0100
MINIFICPP-1449 Add pause and resume message to C2 functionality
This closes #974
Signed-off-by: Marton Szasz <[email protected]>
---
extensions/coap/protocols/CoapC2Protocol.cpp | 4 +
extensions/http-curl/tests/C2PauseResumeTest.cpp | 147 +++++++++++++++++++++
extensions/http-curl/tests/CMakeLists.txt | 1 +
libminifi/include/FlowController.h | 5 +-
libminifi/include/c2/C2Agent.h | 7 +-
libminifi/include/c2/C2Client.h | 2 +-
libminifi/include/c2/C2Payload.h | 4 +-
libminifi/include/c2/PayloadSerializer.h | 10 ++
libminifi/include/core/state/ProcessorController.h | 14 +-
libminifi/include/core/state/UpdateController.h | 13 +-
libminifi/include/utils/ThreadPool.h | 11 ++
libminifi/src/FlowController.cpp | 28 +++-
libminifi/src/c2/C2Agent.cpp | 20 ++-
libminifi/src/c2/C2Client.cpp | 4 +-
libminifi/src/c2/protocols/RESTProtocol.cpp | 8 ++
libminifi/src/core/state/ProcessorController.cpp | 7 +-
libminifi/src/utils/ThreadPool.cpp | 19 ++-
libminifi/test/resources/C2PauseResumeTest.yml | 77 +++++++++++
libminifi/test/unit/ControllerTests.cpp | 8 ++
libminifi/test/unit/ProvenanceTestHelper.h | 4 +
nanofi/include/cxx/C2CallbackAgent.h | 6 +-
nanofi/include/cxx/Instance.h | 2 +-
nanofi/src/cxx/C2CallbackAgent.cpp | 9 +-
23 files changed, 379 insertions(+), 31 deletions(-)
diff --git a/extensions/coap/protocols/CoapC2Protocol.cpp
b/extensions/coap/protocols/CoapC2Protocol.cpp
index 1cf7003..aa2da16 100644
--- a/extensions/coap/protocols/CoapC2Protocol.cpp
+++ b/extensions/coap/protocols/CoapC2Protocol.cpp
@@ -182,6 +182,10 @@ minifi::c2::Operation CoapProtocol::getOperation(int type)
const {
return minifi::c2::UPDATE;
case 7:
return minifi::c2::STOP;
+ case 8:
+ return minifi::c2::PAUSE;
+ case 9:
+ return minifi::c2::RESUME;
}
return minifi::c2::ACKNOWLEDGE;
}
diff --git a/extensions/http-curl/tests/C2PauseResumeTest.cpp
b/extensions/http-curl/tests/C2PauseResumeTest.cpp
new file mode 100644
index 0000000..c462965
--- /dev/null
+++ b/extensions/http-curl/tests/C2PauseResumeTest.cpp
@@ -0,0 +1,147 @@
+/**
+ *
+ * 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.
+ */
+
+#undef NDEBUG
+#include "TestBase.h"
+#include "HTTPIntegrationBase.h"
+#include "HTTPHandlers.h"
+#include "InvokeHTTP.h"
+#include "TestServer.h"
+#include "core/yaml/YamlConfiguration.h"
+#include "FlowController.h"
+#include "properties/Configure.h"
+#include "io/StreamFactory.h"
+#include "integration/IntegrationBase.h"
+#include "utils/GeneralUtils.h"
+
+class VerifyC2PauseResume : public VerifyC2Base {
+ public:
+ explicit VerifyC2PauseResume(const std::atomic_bool&
flow_resumed_successfully) : VerifyC2Base(),
flow_resumed_successfully_(flow_resumed_successfully) {}
+
+ void configureC2() override {
+ VerifyC2Base::configureC2();
+ configuration->set("nifi.c2.agent.heartbeat.period", "500");
+ }
+
+ void runAssertions() override {
+ using org::apache::nifi::minifi::utils::verifyEventHappenedInPollTime;
+ assert(verifyEventHappenedInPollTime(std::chrono::seconds(20), [&] {
return flow_resumed_successfully_.load(); }));
+ }
+
+ private:
+ const std::atomic_bool& flow_resumed_successfully_;
+};
+
+class PauseResumeHandler: public HeartbeatHandler {
+ public:
+ static const uint32_t PAUSE_SECONDS = 3;
+ static const uint32_t INITIAL_GET_INVOKE_COUNT = 2;
+
+ explicit PauseResumeHandler(std::atomic_bool& flow_resumed_successfully) :
HeartbeatHandler(), flow_resumed_successfully_(flow_resumed_successfully) {}
+ bool handleGet(CivetServer *server, struct mg_connection *conn) override {
+ assert(flow_state_ != FlowState::PAUSED);
+ ++get_invoke_count_;
+ if (flow_state_ == FlowState::RESUMED) {
+ flow_resumed_successfully_ = true;
+ }
+
+ mg_printf(conn, "HTTP/1.1 200 OK\r\n");
+ return true;
+ }
+
+ void handleHeartbeat(const rapidjson::Document&, struct mg_connection *
conn) override {
+ std::string operation = "resume";
+ if (flow_state_ == FlowState::PAUSE_INITIATED) {
+ pause_start_time_ = std::chrono::system_clock::now();
+ flow_state_ = FlowState::PAUSED;
+ operation = "pause";
+ } else if (get_invoke_count_ == INITIAL_GET_INVOKE_COUNT && flow_state_ ==
FlowState::STARTED) {
+ flow_state_ = FlowState::PAUSE_INITIATED;
+ operation = "pause";
+ } else if (flow_state_ == FlowState::PAUSED) {
+ operation = "pause";
+ }
+
+ std::string heartbeat_response = "{\"operation\" :
\"heartbeat\",\"requested_operations\": [ {"
+ "\"operation\" : \"" + operation + "\","
+ "\"operationid\" : \"8675309\"}]}";
+
+ if (flow_state_ == FlowState::PAUSED &&
std::chrono::duration_cast<std::chrono::seconds>(std::chrono::system_clock::now()
- pause_start_time_).count() > PAUSE_SECONDS) {
+ flow_state_ = FlowState::RESUMED;
+ }
+
+ mg_printf(conn, "HTTP/1.1 200 OK\r\nContent-Type: "
+ "text/plain\r\nContent-Length: %lu\r\nConnection: close\r\n\r\n",
+ heartbeat_response.length());
+ mg_printf(conn, "%s", heartbeat_response.c_str());
+ }
+
+ private:
+ enum class FlowState {
+ STARTED,
+ PAUSE_INITIATED,
+ PAUSED,
+ RESUMED
+ };
+
+ std::atomic<uint32_t> get_invoke_count_{0};
+ std::chrono::time_point<std::chrono::system_clock> pause_start_time_;
+ std::atomic<FlowState> flow_state_{FlowState::STARTED};
+ std::atomic_bool& flow_resumed_successfully_;
+};
+
+int main(int argc, char **argv) {
+ const cmd_args args = parse_cmdline_args(argc, argv, "heartbeat");
+ std::atomic_bool flow_resumed_successfully{false};
+ VerifyC2PauseResume harness{flow_resumed_successfully};
+ harness.setKeyDir(args.key_dir);
+ PauseResumeHandler responder{flow_resumed_successfully};
+
+ std::shared_ptr<core::Repository> test_repo =
std::make_shared<TestRepository>();
+ std::shared_ptr<core::Repository> test_flow_repo =
std::make_shared<TestFlowRepository>();
+ std::shared_ptr<minifi::Configure> configuration =
std::make_shared<minifi::Configure>();
+ configuration->set(minifi::Configure::nifi_default_directory, args.key_dir);
+ configuration->set(minifi::Configure::nifi_flow_configuration_file,
args.test_file);
+
+ std::shared_ptr<minifi::io::StreamFactory> stream_factory =
minifi::io::StreamFactory::getInstance(configuration);
+ std::shared_ptr<core::ContentRepository> content_repo =
std::make_shared<core::repository::VolatileContentRepository>();
+ content_repo->initialize(configuration);
+
+ std::unique_ptr<core::FlowConfiguration> yaml_ptr =
utils::make_unique<core::YamlConfiguration>(
+ test_repo, test_repo, content_repo, stream_factory, configuration,
args.test_file);
+
+ std::shared_ptr<minifi::FlowController> controller =
std::make_shared<minifi::FlowController>(
+ test_repo, test_flow_repo, configuration, std::move(yaml_ptr),
content_repo, DEFAULT_ROOT_GROUP_NAME, true);
+
+ core::YamlConfiguration yaml_config(test_repo, test_repo, content_repo,
stream_factory, configuration, args.test_file);
+
+ std::shared_ptr<core::Processor> proc =
yaml_config.getRoot()->findProcessorByName("invoke");
+ assert(proc != nullptr);
+
+ const auto inv =
std::dynamic_pointer_cast<minifi::processors::InvokeHTTP>(proc);
+ assert(inv != nullptr);
+ std::string url;
+ inv->getProperty(minifi::processors::InvokeHTTP::URL.getName(), url);
+ std::string port, scheme, path;
+ std::unique_ptr<TestServer> server;
+ parse_http_components(url, port, scheme, path);
+ server = utils::make_unique<TestServer>(port, path, &responder);
+
+ harness.setUrl(args.url, &responder);
+ harness.run(args.test_file);
+}
diff --git a/extensions/http-curl/tests/CMakeLists.txt
b/extensions/http-curl/tests/CMakeLists.txt
index e6ff766..0cbb237 100644
--- a/extensions/http-curl/tests/CMakeLists.txt
+++ b/extensions/http-curl/tests/CMakeLists.txt
@@ -98,3 +98,4 @@ add_test(NAME ControllerServiceIntegrationTests COMMAND
ControllerServiceIntegra
add_test(NAME ThreadPoolAdjust COMMAND ThreadPoolAdjust
"${TEST_RESOURCES}/ThreadPoolAdjust.yml" "${TEST_RESOURCES}/")
add_test(NAME VerifyInvokeHTTPTest COMMAND VerifyInvokeHTTPTest
"${TEST_RESOURCES}/TestInvokeHTTPPost.yml")
add_test(NAME AbsoluteTimeoutTest COMMAND AbsoluteTimeoutTest)
+add_test(NAME C2PauseResumeTest COMMAND C2PauseResumeTest
"${TEST_RESOURCES}/C2PauseResumeTest.yml" "${TEST_RESOURCES}/")
diff --git a/libminifi/include/FlowController.h
b/libminifi/include/FlowController.h
index 9ed88a3..932fbbb 100644
--- a/libminifi/include/FlowController.h
+++ b/libminifi/include/FlowController.h
@@ -108,9 +108,8 @@ class FlowController : public
core::controller::ForwardingControllerServiceProvi
}
// Start to run the Flow Controller which internally start the root process
group and all its children
int16_t start() override;
- int16_t pause() override {
- return -1;
- }
+ int16_t pause() override;
+ int16_t resume() override;
// Unload the current flow YAML, clean the root process group and all its
children
int16_t stop() override;
int16_t applyUpdate(const std::string &source, const std::string
&configuration, bool persist) override;
diff --git a/libminifi/include/c2/C2Agent.h b/libminifi/include/c2/C2Agent.h
index 9cc6c28..e3cd68f 100644
--- a/libminifi/include/c2/C2Agent.h
+++ b/libminifi/include/c2/C2Agent.h
@@ -64,10 +64,11 @@ class C2Agent : public state::UpdateController {
public:
static constexpr const char* UPDATE_NAME = "C2UpdatePolicy";
- C2Agent(core::controller::ControllerServiceProvider* controller,
+ C2Agent(core::controller::ControllerServiceProvider *controller,
+ state::Pausable *pause_handler,
const std::shared_ptr<state::StateMonitor> &updateSink,
const std::shared_ptr<Configure> &configure,
- const std::shared_ptr<utils::file::FileSystem>& filesystem =
std::make_shared<utils::file::FileSystem>());
+ const std::shared_ptr<utils::file::FileSystem> &filesystem =
std::make_shared<utils::file::FileSystem>());
virtual ~C2Agent() noexcept {
delete protocol_.load();
}
@@ -206,6 +207,8 @@ class C2Agent : public state::UpdateController {
// controller service provider reference.
core::controller::ControllerServiceProvider* controller_;
+ state::Pausable* pause_handler_;
+
// shared pointer to the configuration of this agent
std::shared_ptr<Configure> configuration_;
diff --git a/libminifi/include/c2/C2Client.h b/libminifi/include/c2/C2Client.h
index 982d6cf..ab912e1 100644
--- a/libminifi/include/c2/C2Client.h
+++ b/libminifi/include/c2/C2Client.h
@@ -48,7 +48,7 @@ class C2Client : public core::Flow, public
state::response::NodeReporter {
std::unique_ptr<core::FlowConfiguration> flow_configuration,
std::shared_ptr<utils::file::FileSystem> filesystem,
std::shared_ptr<logging::Logger> logger =
logging::LoggerFactory<C2Client>::getLogger());
- void initialize(core::controller::ControllerServiceProvider* controller,
const std::shared_ptr<state::StateMonitor> &update_sink);
+ void initialize(core::controller::ControllerServiceProvider *controller,
state::Pausable *pause_handler, const std::shared_ptr<state::StateMonitor>
&update_sink);
std::shared_ptr<state::response::ResponseNode> getMetricsNode(const
std::string& metrics_class) const override;
diff --git a/libminifi/include/c2/C2Payload.h b/libminifi/include/c2/C2Payload.h
index 4be8ae6..1b7875b 100644
--- a/libminifi/include/c2/C2Payload.h
+++ b/libminifi/include/c2/C2Payload.h
@@ -44,7 +44,9 @@ enum Operation {
UPDATE,
VALIDATE,
CLEAR,
- TRANSFER
+ TRANSFER,
+ PAUSE,
+ RESUME
};
#define PAYLOAD_NO_STATUS 0
diff --git a/libminifi/include/c2/PayloadSerializer.h
b/libminifi/include/c2/PayloadSerializer.h
index 89042c9..3c71d09 100644
--- a/libminifi/include/c2/PayloadSerializer.h
+++ b/libminifi/include/c2/PayloadSerializer.h
@@ -130,6 +130,12 @@ class PayloadSerializer {
case Operation::UPDATE:
op = 7;
break;
+ case Operation::PAUSE:
+ op = 8;
+ break;
+ case Operation::RESUME:
+ op = 9;
+ break;
default:
op = 2;
break;
@@ -309,6 +315,10 @@ class PayloadSerializer {
return Operation::START;
case 7:
return Operation::UPDATE;
+ case 8:
+ return Operation::PAUSE;
+ case 9:
+ return Operation::RESUME;
default:
return Operation::HEARTBEAT;
}
diff --git a/libminifi/include/core/state/ProcessorController.h
b/libminifi/include/core/state/ProcessorController.h
index abf01e6..5b76121 100644
--- a/libminifi/include/core/state/ProcessorController.h
+++ b/libminifi/include/core/state/ProcessorController.h
@@ -42,11 +42,11 @@ class ProcessorController : public StateController {
virtual ~ProcessorController();
- virtual std::string getComponentName() const {
+ std::string getComponentName() const override {
return processor_->getName();
}
- virtual utils::Identifier getComponentUUID() const {
+ utils::Identifier getComponentUUID() const override {
return processor_->getUUID();
}
@@ -56,15 +56,17 @@ class ProcessorController : public StateController {
/**
* Start the client
*/
- virtual int16_t start();
+ int16_t start() override;
/**
* Stop the client
*/
- virtual int16_t stop();
+ int16_t stop() override;
- virtual bool isRunning();
+ bool isRunning() override;
- virtual int16_t pause();
+ int16_t pause() override;
+
+ int16_t resume() override;
protected:
std::shared_ptr<core::Processor> processor_;
diff --git a/libminifi/include/core/state/UpdateController.h
b/libminifi/include/core/state/UpdateController.h
index 512fc9f..1d4c96e 100644
--- a/libminifi/include/core/state/UpdateController.h
+++ b/libminifi/include/core/state/UpdateController.h
@@ -148,7 +148,16 @@ class UpdateRunner : public utils::AfterExecute<Update> {
std::chrono::milliseconds delay_;
};
-class StateController {
+class Pausable {
+ public:
+ virtual ~Pausable() = default;
+
+ virtual int16_t pause() = 0;
+
+ virtual int16_t resume() = 0;
+};
+
+class StateController : public Pausable {
public:
virtual ~StateController() = default;
@@ -165,8 +174,6 @@ class StateController {
virtual int16_t stop() = 0;
virtual bool isRunning() = 0;
-
- virtual int16_t pause() = 0;
};
/**
diff --git a/libminifi/include/utils/ThreadPool.h
b/libminifi/include/utils/ThreadPool.h
index 9607533..e206ac9 100644
--- a/libminifi/include/utils/ThreadPool.h
+++ b/libminifi/include/utils/ThreadPool.h
@@ -27,6 +27,7 @@
#include <atomic>
#include <mutex>
#include <map>
+#include <unordered_map>
#include <vector>
#include <queue>
#include <future>
@@ -227,6 +228,16 @@ class ThreadPool {
void stopTasks(const TaskId &identifier);
/**
+ * resumes work queue processing.
+ */
+ void resume();
+
+ /**
+ * pauses work queue processing
+ */
+ void pause();
+
+ /**
* Returns true if a task is running.
*/
bool isTaskRunning(const TaskId &identifier) {
diff --git a/libminifi/src/FlowController.cpp b/libminifi/src/FlowController.cpp
index 25c7036..458a155 100644
--- a/libminifi/src/FlowController.cpp
+++ b/libminifi/src/FlowController.cpp
@@ -268,7 +268,7 @@ std::unique_ptr<core::ProcessGroup>
FlowController::loadInitialFlow() {
// since we don't have access to the flow definition, the C2 communication
// won't be able to use the services defined there, e.g. SSLContextService
controller_service_provider_impl_ =
flow_configuration_->getControllerServiceProvider();
- C2Client::initialize(this, shared_from_this());
+ C2Client::initialize(this, this, shared_from_this());
auto opt_source = fetchFlow(*opt_flow_url);
if (!opt_source) {
logger_->log_error("Couldn't fetch flow configuration from C2 server");
@@ -379,7 +379,7 @@ int16_t FlowController::start() {
// as the thread_pool_ is started in load()
this->root_->startProcessing(timer_scheduler_, event_scheduler_,
cron_scheduler_);
}
- C2Client::initialize(this, shared_from_this());
+ C2Client::initialize(this, this, shared_from_this());
running_ = true;
this->protocol_->start();
this->provenance_repo_->start();
@@ -391,6 +391,30 @@ int16_t FlowController::start() {
}
}
+int16_t FlowController::pause() {
+ std::lock_guard<std::recursive_mutex> flow_lock(mutex_);
+ if (!running_) {
+ logger_->log_warn("Can not pause flow controller that is not running");
+ return 0;
+ }
+
+ logger_->log_info("Pausing Flow Controller");
+ thread_pool_.pause();
+ return 0;
+}
+
+int16_t FlowController::resume() {
+ std::lock_guard<std::recursive_mutex> flow_lock(mutex_);
+ if (!running_) {
+ logger_->log_warn("Can not resume flow controller tasks because the flow
controller is not running");
+ return 0;
+ }
+
+ logger_->log_info("Resuming Flow Controller");
+ thread_pool_.resume();
+ return 0;
+}
+
int16_t FlowController::applyUpdate(const std::string &source, const
std::string &configuration, bool persist) {
if (applyConfiguration(source, configuration)) {
if (persist) {
diff --git a/libminifi/src/c2/C2Agent.cpp b/libminifi/src/c2/C2Agent.cpp
index 8f25c07..0914ed5 100644
--- a/libminifi/src/c2/C2Agent.cpp
+++ b/libminifi/src/c2/C2Agent.cpp
@@ -51,15 +51,17 @@ namespace nifi {
namespace minifi {
namespace c2 {
-C2Agent::C2Agent(core::controller::ControllerServiceProvider* controller,
+C2Agent::C2Agent(core::controller::ControllerServiceProvider *controller,
+ state::Pausable *pause_handler,
const std::shared_ptr<state::StateMonitor> &updateSink,
const std::shared_ptr<Configure> &configuration,
- const std::shared_ptr<utils::file::FileSystem>& filesystem)
+ const std::shared_ptr<utils::file::FileSystem> &filesystem)
: heart_beat_period_(3000),
max_c2_responses(5),
update_sink_(updateSink),
update_service_(nullptr),
controller_(controller),
+ pause_handler_(pause_handler),
configuration_(configuration),
filesystem_(filesystem),
protocol_(nullptr),
@@ -429,6 +431,20 @@ void C2Agent::handle_c2_server_response(const
C2ContentResponse &resp) {
}
//
break;
+ case Operation::PAUSE:
+ if (pause_handler_ != nullptr) {
+ pause_handler_->pause();
+ } else {
+ logger_->log_warn("Pause functionality is not supported!");
+ }
+ break;
+ case Operation::RESUME:
+ if (pause_handler_ != nullptr) {
+ pause_handler_->resume();
+ } else {
+ logger_->log_warn("Resume functionality is not supported!");
+ }
+ break;
default:
break;
// do nothing
diff --git a/libminifi/src/c2/C2Client.cpp b/libminifi/src/c2/C2Client.cpp
index d24fc2a..7547ba8 100644
--- a/libminifi/src/c2/C2Client.cpp
+++ b/libminifi/src/c2/C2Client.cpp
@@ -58,7 +58,7 @@ bool C2Client::isC2Enabled() const {
return utils::StringUtils::toBool(c2_enable_str).value_or(false);
}
-void C2Client::initialize(core::controller::ControllerServiceProvider
*controller, const std::shared_ptr<state::StateMonitor> &update_sink) {
+void C2Client::initialize(core::controller::ControllerServiceProvider
*controller, state::Pausable *pause_handler, const
std::shared_ptr<state::StateMonitor> &update_sink) {
if (!isC2Enabled()) {
return;
}
@@ -120,7 +120,7 @@ void
C2Client::initialize(core::controller::ControllerServiceProvider *controlle
if (!initialized_) {
// C2Agent is initialized once, meaning that a C2-triggered
flow/configuration update
// might not be equal to a fresh restart
- c2_agent_ = std::unique_ptr<c2::C2Agent>(new c2::C2Agent(controller,
update_sink, configuration_, filesystem_));
+ c2_agent_ = std::unique_ptr<c2::C2Agent>(new c2::C2Agent(controller,
pause_handler, update_sink, configuration_, filesystem_));
c2_agent_->start();
initialized_ = true;
}
diff --git a/libminifi/src/c2/protocols/RESTProtocol.cpp
b/libminifi/src/c2/protocols/RESTProtocol.cpp
index 8903c71..ddd0012 100644
--- a/libminifi/src/c2/protocols/RESTProtocol.cpp
+++ b/libminifi/src/c2/protocols/RESTProtocol.cpp
@@ -431,6 +431,10 @@ std::string RESTProtocol::getOperation(const C2Payload
&payload) {
return "start";
case Operation::UPDATE:
return "update";
+ case Operation::PAUSE:
+ return "pause";
+ case Operation::RESUME:
+ return "resume";
default:
return "heartbeat";
}
@@ -455,6 +459,10 @@ Operation RESTProtocol::stringToOperation(const
std::string str) {
return Operation::STOP;
} else if (op == "start") {
return Operation::START;
+ } else if (op == "pause") {
+ return Operation::PAUSE;
+ } else if (op == "resume") {
+ return Operation::RESUME;
}
return Operation::HEARTBEAT;
}
diff --git a/libminifi/src/core/state/ProcessorController.cpp
b/libminifi/src/core/state/ProcessorController.cpp
index edc766a..dcf235a 100644
--- a/libminifi/src/core/state/ProcessorController.cpp
+++ b/libminifi/src/core/state/ProcessorController.cpp
@@ -52,8 +52,11 @@ bool ProcessorController::isRunning() {
}
int16_t ProcessorController::pause() {
- scheduler_->unschedule(processor_);
- return 0;
+ return stop();
+}
+
+int16_t ProcessorController::resume() {
+ return start();
}
} /* namespace state */
diff --git a/libminifi/src/utils/ThreadPool.cpp
b/libminifi/src/utils/ThreadPool.cpp
index 3c5d3bf..a74313d 100644
--- a/libminifi/src/utils/ThreadPool.cpp
+++ b/libminifi/src/utils/ThreadPool.cpp
@@ -44,6 +44,9 @@ void ThreadPool<T>::run_tasks(std::shared_ptr<WorkerThread>
thread) {
std::unique_lock<std::mutex> lock(worker_queue_mutex_);
if (!task_status_[task.getIdentifier()]) {
continue;
+ } else if (!worker_queue_.isRunning()) {
+ worker_queue_.enqueue(std::move(task));
+ continue;
}
}
if (task.run()) {
@@ -64,7 +67,7 @@ void ThreadPool<T>::run_tasks(std::shared_ptr<WorkerThread>
thread) {
}
}
} else {
- // This means that the threadpool is running, but the ConcurrentQueue is
stopped -> shouldn't happen during normal conditions
+ // The threadpool is running, but the ConcurrentQueue is stopped ->
shouldn't happen during normal conditions
// Might happen during startup or shutdown for a very short time
if (running_.load()) {
std::this_thread::sleep_for(std::chrono::milliseconds(1));
@@ -200,6 +203,20 @@ void ThreadPool<T>::stopTasks(const TaskId &identifier) {
}
template<typename T>
+void ThreadPool<T>::resume() {
+ if (!worker_queue_.isRunning()) {
+ worker_queue_.start();
+ }
+}
+
+template<typename T>
+void ThreadPool<T>::pause() {
+ if (worker_queue_.isRunning()) {
+ worker_queue_.stop();
+ }
+}
+
+template<typename T>
void ThreadPool<T>::shutdown() {
if (running_.load()) {
std::lock_guard<std::recursive_mutex> lock(manager_mutex_);
diff --git a/libminifi/test/resources/C2PauseResumeTest.yml
b/libminifi/test/resources/C2PauseResumeTest.yml
new file mode 100644
index 0000000..4755c0b
--- /dev/null
+++ b/libminifi/test/resources/C2PauseResumeTest.yml
@@ -0,0 +1,77 @@
+#
+# 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.
+#
+Flow Controller:
+ name: MiNiFi Flow
+ id: 2438e3c8-015a-1000-79ca-83af40ec1990
+Processors:
+ - name: invoke
+ id: 2438e3c8-015a-1000-79ca-83af40ec1991
+ class: org.apache.nifi.processors.standard.InvokeHTTP
+ max concurrent tasks: 1
+ scheduling strategy: TIMER_DRIVEN
+ scheduling period: 1 sec
+ penalization period: 30 sec
+ yield period: 1 sec
+ run duration nanos: 0
+ auto-terminated relationships list:
+ - retry
+ - no retry
+ - response
+ - failure
+ Properties:
+ HTTP Method: GET
+ Remote URL: http://localhost:10008/geturl
+ - name: LogAttribute
+ id: 2438e3c8-015a-1000-79ca-83af40ec1992
+ class: org.apache.nifi.processors.standard.LogAttribute
+ max concurrent tasks: 1
+ scheduling strategy: TIMER_DRIVEN
+ scheduling period: 1 sec
+ penalization period: 30 sec
+ yield period: 1 sec
+ run duration nanos: 0
+ auto-terminated relationships list: response
+ Properties:
+ Log Level: info
+ Log Payload: true
+
+Connections:
+ - name: TransferFilesToRPG
+ id: 2438e3c8-015a-1000-79ca-83af40ec1997
+ source name: invoke
+ source id: 2438e3c8-015a-1000-79ca-83af40ec1991
+ source relationship name: success
+ destination name: LogAttribute
+ destination id: 2438e3c8-015a-1000-79ca-83af40ec1992
+ max work queue size: 0
+ max work queue data size: 1 MB
+ flowfile expiration: 60 sec
+ - name: TransferFilesToRPG2
+ id: 2438e3c8-015a-1000-79ca-83af40ec1917
+ source name: LogAttribute
+ source id: 2438e3c8-015a-1000-79ca-83af40ec1992
+ destination name: LogAttribute
+ destination id: 2438e3c8-015a-1000-79ca-83af40ec1992
+ source relationship name: success
+ max work queue size: 0
+ max work queue data size: 1 MB
+ flowfile expiration: 60 sec
+
+Remote Processing Groups:
+
diff --git a/libminifi/test/unit/ControllerTests.cpp
b/libminifi/test/unit/ControllerTests.cpp
index 6ea798d..de14f60 100644
--- a/libminifi/test/unit/ControllerTests.cpp
+++ b/libminifi/test/unit/ControllerTests.cpp
@@ -66,6 +66,10 @@ class TestStateController : public
minifi::state::StateController {
return 0;
}
+ virtual int16_t resume() {
+ return 0;
+ }
+
std::atomic<bool> is_running;
};
@@ -119,6 +123,10 @@ class TestUpdateSink : public minifi::state::StateMonitor {
int16_t pause() override {
return 0;
}
+
+ int16_t resume() override {
+ return 0;
+ }
std::vector<BackTrace> getTraces() override {
std::vector<BackTrace> traces;
return traces;
diff --git a/libminifi/test/unit/ProvenanceTestHelper.h
b/libminifi/test/unit/ProvenanceTestHelper.h
index 0460eef..82f7978 100644
--- a/libminifi/test/unit/ProvenanceTestHelper.h
+++ b/libminifi/test/unit/ProvenanceTestHelper.h
@@ -266,6 +266,10 @@ class TestFlowController : public minifi::FlowController {
return -1;
}
+ int16_t resume() override {
+ return -1;
+ }
+
void unload() override {
stop();
}
diff --git a/nanofi/include/cxx/C2CallbackAgent.h
b/nanofi/include/cxx/C2CallbackAgent.h
index 81198bb..287ba73 100644
--- a/nanofi/include/cxx/C2CallbackAgent.h
+++ b/nanofi/include/cxx/C2CallbackAgent.h
@@ -47,7 +47,11 @@ class C2CallbackAgent : public c2::C2Agent {
public:
- explicit C2CallbackAgent(core::controller::ControllerServiceProvider*
controller, const std::shared_ptr<state::StateMonitor> &updateSink, const
std::shared_ptr<Configure> &configure);
+ explicit C2CallbackAgent(
+ core::controller::ControllerServiceProvider* controller,
+ state::Pausable* pause_handler,
+ const std::shared_ptr<state::StateMonitor> &updateSink,
+ const std::shared_ptr<Configure> &configure);
virtual ~C2CallbackAgent() = default;
diff --git a/nanofi/include/cxx/Instance.h b/nanofi/include/cxx/Instance.h
index 8eab940..3261914 100644
--- a/nanofi/include/cxx/Instance.h
+++ b/nanofi/include/cxx/Instance.h
@@ -106,7 +106,7 @@ class Instance {
configure_->set("c2.rest.url", server->url);
configure_->set("c2.rest.url.ack", server->ack_url);
}
- agent_ = std::make_shared<c2::C2CallbackAgent>(nullptr, nullptr,
configure_);
+ agent_ = std::make_shared<c2::C2CallbackAgent>(nullptr, nullptr, nullptr,
configure_);
listener_thread_pool_.start();
registerUpdateListener(agent_, 1000);
agent_->setStopCallback(c1);
diff --git a/nanofi/src/cxx/C2CallbackAgent.cpp
b/nanofi/src/cxx/C2CallbackAgent.cpp
index b41a870..3498aeb 100644
--- a/nanofi/src/cxx/C2CallbackAgent.cpp
+++ b/nanofi/src/cxx/C2CallbackAgent.cpp
@@ -34,9 +34,9 @@ namespace nifi {
namespace minifi {
namespace c2 {
-C2CallbackAgent::C2CallbackAgent(core::controller::ControllerServiceProvider*
controller, const std::shared_ptr<state::StateMonitor> &updateSink,
+C2CallbackAgent::C2CallbackAgent(core::controller::ControllerServiceProvider*
controller, state::Pausable* pause_handler, const
std::shared_ptr<state::StateMonitor> &updateSink,
const std::shared_ptr<Configure>
&configuration)
- : C2Agent(controller, updateSink, configuration),
+ : C2Agent(controller, pause_handler, updateSink, configuration),
stop(nullptr),
logger_(logging::LoggerFactory<C2CallbackAgent>::getLogger()) {
}
@@ -64,11 +64,12 @@ void C2CallbackAgent::handle_c2_server_response(const
C2ContentResponse &resp) {
break;
}
- //
+ case Operation::PAUSE:
+ break;
+ case Operation::RESUME:
break;
default:
break;
- // do nothing
}
}