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
   }
 }
 

Reply via email to