martinzink commented on code in PR #1335:
URL: https://github.com/apache/nifi-minifi-cpp/pull/1335#discussion_r882738381


##########
libminifi/test/unit/SchedulingAgentTests.cpp:
##########
@@ -0,0 +1,138 @@
+/**
+ * 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 "../Catch.h"
+#include "../TestBase.h"
+#include "ProvenanceTestHelper.h"
+#include "utils/TestUtils.h"
+
+using namespace std::literals::chrono_literals;
+
+namespace org::apache::nifi::minifi::testing {
+
+class CountOnTriggersProcessor : public minifi::core::Processor {
+ public:
+  using minifi::core::Processor::Processor;
+
+  void onTrigger(core::ProcessContext*, core::ProcessSession*) override {
+    ++number_of_triggers;
+  }
+
+  size_t getNumberOfTriggers() const { return number_of_triggers; }
+
+ private:
+  std::atomic<size_t> number_of_triggers = 0;
+};
+
+
+TEST_CASE("SchedulingAgentTests", "[SchedulingAgent]") {
+  std::shared_ptr<core::Repository> test_repo = 
std::make_shared<TestRepository>();
+  std::shared_ptr<core::ContentRepository> content_repo = 
std::make_shared<core::repository::VolatileContentRepository>();
+  std::shared_ptr<TestRepository> repo = 
std::static_pointer_cast<TestRepository>(test_repo);
+  std::shared_ptr<minifi::FlowController> controller =
+      std::make_shared<TestFlowController>(test_repo, test_repo, content_repo);
+
+  TestController testController;
+  auto test_plan = testController.createPlan();
+  auto controller_services_ = 
std::make_shared<minifi::core::controller::ControllerServiceMap>();
+  auto configuration = std::make_shared<minifi::Configure>();
+  auto controller_services_provider_ = 
std::make_shared<minifi::core::controller::StandardControllerServiceProvider>(controller_services_,
 nullptr, configuration);
+  utils::ThreadPool<utils::TaskRescheduleInfo> thread_pool;
+  auto count_proc = std::make_shared<CountOnTriggersProcessor>("count_proc");
+  count_proc->incrementActiveTasks();
+  count_proc->setScheduledState(core::RUNNING);
+  auto node = std::make_shared<core::ProcessorNode>(count_proc.get());
+  auto context = std::make_shared<core::ProcessContext>(node, nullptr, repo, 
repo, content_repo);
+  std::shared_ptr<core::ProcessSessionFactory> factory = 
std::make_shared<core::ProcessSessionFactory>(context);
+  count_proc->setSchedulingPeriodNano(1250ms);
+#ifdef WIN32
+  utils::dateSetInstall(TZ_DATA_DIR);
+#endif
+
+  SECTION("Timer Driven") {
+    auto timer_driven_agent = 
std::make_shared<TimerDrivenSchedulingAgent>(gsl::make_not_null(controller_services_provider_.get()),
 test_repo, test_repo, content_repo, configuration, thread_pool);
+    timer_driven_agent->start();
+    auto first_task_reschedule_info = 
timer_driven_agent->run(count_proc.get(), context, factory);
+    CHECK(!first_task_reschedule_info.finished_);
+    CHECK(first_task_reschedule_info.wait_time_ == 1250ms);
+    CHECK(count_proc->getNumberOfTriggers() == 1);
+
+    auto second_task_reschedule_info = 
timer_driven_agent->run(count_proc.get(), context, factory);
+
+    CHECK(!second_task_reschedule_info.finished_);
+    CHECK(second_task_reschedule_info.wait_time_ == 1250ms);
+    CHECK(count_proc->getNumberOfTriggers() == 2);
+  }
+
+  SECTION("Event Driven") {
+    auto event_driven_agent = 
std::make_shared<EventDrivenSchedulingAgent>(gsl::make_not_null(controller_services_provider_.get()),
 test_repo, test_repo, content_repo, configuration, thread_pool);
+    event_driven_agent->start();
+    auto first_task_reschedule_info = 
event_driven_agent->run(count_proc.get(), context, factory);
+    CHECK(!first_task_reschedule_info.finished_);
+    CHECK(first_task_reschedule_info.wait_time_ == 0ms);
+    auto count_num_after_one_schedule = count_proc->getNumberOfTriggers();
+    CHECK(count_num_after_one_schedule > 100);
+
+    auto second_task_reschedule_info = 
event_driven_agent->run(count_proc.get(), context, factory);
+    CHECK(!second_task_reschedule_info.finished_);
+    CHECK(second_task_reschedule_info.wait_time_ == 0ms);
+    auto count_num_after_two_schedule = count_proc->getNumberOfTriggers();
+    CHECK(count_num_after_two_schedule > count_num_after_one_schedule+100);
+  }
+
+  SECTION("Cron Driven every year") {
+    count_proc->setCronPeriod("0 0 0 1 1/12 ?");
+    auto cron_driven_agent = 
std::make_shared<CronDrivenSchedulingAgent>(gsl::make_not_null(controller_services_provider_.get()),
 test_repo, test_repo, content_repo, configuration, thread_pool);
+    cron_driven_agent->start();
+    auto first_task_reschedule_info = cron_driven_agent->run(count_proc.get(), 
context, factory);
+    CHECK(!first_task_reschedule_info.finished_);
+    auto next_run_time_point = 
std::chrono::round<std::chrono::years>(std::chrono::system_clock::now() + 
first_task_reschedule_info.wait_time_);
+    CHECK(next_run_time_point == 
std::chrono::ceil<std::chrono::years>(std::chrono::system_clock::now()));
+    CHECK(count_proc->getNumberOfTriggers() == 0);
+
+    auto second_task_reschedule_info = 
cron_driven_agent->run(count_proc.get(), context, factory);
+    CHECK(!second_task_reschedule_info.finished_);
+    next_run_time_point = 
std::chrono::round<std::chrono::years>(std::chrono::system_clock::now() + 
first_task_reschedule_info.wait_time_);
+    CHECK(next_run_time_point == 
std::chrono::ceil<std::chrono::years>(std::chrono::system_clock::now()));
+    CHECK(count_proc->getNumberOfTriggers() == 0);
+  }
+
+  SECTION("Cron Driven every sec") {
+    count_proc->setCronPeriod("* * * * * *");
+    auto cron_driven_agent = 
std::make_shared<CronDrivenSchedulingAgent>(gsl::make_not_null(controller_services_provider_.get()),
 test_repo, test_repo, content_repo, configuration, thread_pool);
+    cron_driven_agent->start();
+    auto first_task_reschedule_info = cron_driven_agent->run(count_proc.get(), 
context, factory);
+    CHECK(!first_task_reschedule_info.finished_);
+    CHECK(first_task_reschedule_info.wait_time_ <= 1s);
+    CHECK(count_proc->getNumberOfTriggers() == 0);

Review Comment:
   I don't think so. The processor can only be triggered once per 
`cron_driven_agent->run` and the first run will never trigger it.
   
   We use local_seconds so the precision is 1 second.
   During `cron_driven_agent->run` we calculate the current_time and as there 
is no previous execution time_point (first run) we set this as the last 
execution time.
   Then we calculate the next execution time starting from 
last_execution_time+1s (due to roundToNextSecond and second precision)
   So the next execution time will be current_time+1s at the earliest, no 
matter what the cronexpr or the current time is.
   and we compare this the current_time (this is the same current_time (we only 
call system_clock::now once))
   ```
       if (*next_to_last_trigger > current_time.get_local_time())
         return 
utils::TaskRescheduleInfo::RetryIn(ceil<milliseconds>(*next_to_last_trigger-current_time.get_local_time()));
   ```
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to