diff --git a/src/lib/data/TaskInjectStorage.cpp b/src/lib/data/TaskInjectStorage.cpp index 5e79db28..3cae84d6 100644 --- a/src/lib/data/TaskInjectStorage.cpp +++ b/src/lib/data/TaskInjectStorage.cpp @@ -1,8 +1,5 @@ #include "TaskInjectStorage.h" -#include -#include - #include "Storage.h" #include "StorageProvider.h" @@ -33,10 +30,6 @@ Task::TaskState TaskInjectStorage::doUpdate(std::shared_ptr blackboa } } } - else - { - std::this_thread::sleep_for(std::chrono::milliseconds(25)); - } return STATE_FAILURE; } diff --git a/src/lib/data/TaskMergeStorages.cpp b/src/lib/data/TaskMergeStorages.cpp index 8b40118b..f612de79 100644 --- a/src/lib/data/TaskMergeStorages.cpp +++ b/src/lib/data/TaskMergeStorages.cpp @@ -1,8 +1,5 @@ #include "TaskMergeStorages.h" -#include -#include - #include "StorageProvider.h" TaskMergeStorages::TaskMergeStorages( @@ -40,10 +37,6 @@ Task::TaskState TaskMergeStorages::doUpdate(std::shared_ptr blackboa } } } - else - { - std::this_thread::sleep_for(std::chrono::milliseconds(25)); - } return STATE_FAILURE; } diff --git a/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp b/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp index 4d837450..633dc36e 100644 --- a/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp +++ b/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp @@ -71,6 +71,11 @@ void TaskFillIndexerCommandsQueue::doEnter(std::shared_ptr blackboar Task::TaskState TaskFillIndexerCommandsQueue::doUpdate(std::shared_ptr blackboard) { + if (m_interrupted) + { + return STATE_FAILURE; + } + if (!fillCommandQueue()) { std::lock_guard lock(m_commandsMutex); @@ -93,6 +98,12 @@ void TaskFillIndexerCommandsQueue::doExit(std::shared_ptr blackboard void TaskFillIndexerCommandsQueue::doReset(std::shared_ptr blackboard) { + m_interrupted = false; +} + +void TaskFillIndexerCommandsQueue::terminate() +{ + m_interrupted = true; } void TaskFillIndexerCommandsQueue::handleMessage(MessageInterruptTasks* message) diff --git a/src/lib/data/indexer/TaskFillIndexerCommandQueue.h b/src/lib/data/indexer/TaskFillIndexerCommandQueue.h index dd613e11..2caddc77 100644 --- a/src/lib/data/indexer/TaskFillIndexerCommandQueue.h +++ b/src/lib/data/indexer/TaskFillIndexerCommandQueue.h @@ -27,6 +27,7 @@ protected: TaskState doUpdate(std::shared_ptr blackboard) override; void doExit(std::shared_ptr blackboard) override; void doReset(std::shared_ptr blackboard) override; + void terminate() override; void handleMessage(MessageInterruptTasks* message) override; @@ -35,9 +36,13 @@ protected: private: std::unique_ptr m_indexerCommandProvider; InterprocessIndexerCommandManager m_indexerCommandManager; + const size_t m_maximumQueueSize; + std::queue m_filePathQueue; std::mutex m_commandsMutex; + + bool m_interrupted = false; }; #endif // TASK_FILL_INDEXER_COMMAND_QUEUE_H diff --git a/src/lib/project/Project.cpp b/src/lib/project/Project.cpp index 7fd549e1..b2ea8590 100644 --- a/src/lib/project/Project.cpp +++ b/src/lib/project/Project.cpp @@ -488,7 +488,7 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di taskParallelIndexing->addChildTasks( std::make_shared()->addChildTasks( // block until there are indexer commands to process - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask( std::make_shared>("indexer_command_queue_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) ), std::make_shared(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess) @@ -499,11 +499,11 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di taskParallelIndexing->addTask( std::make_shared()->addChildTasks( // block until there are indexers running - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask( std::make_shared>("indexer_threads_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) ), // merge until all indexers stopped and nothing left to merge - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 250)->addChildTask( std::make_shared()->addChildTasks( std::make_shared(storageProvider), std::make_shared>("indexer_threads_stopped", TaskReturnSuccessIf::CONDITION_EQUALS, false) @@ -516,10 +516,10 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di taskParallelIndexing->addTask( std::make_shared()->addChildTasks( // block until there are indexers running - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask( std::make_shared>("indexer_threads_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) ), - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask( std::make_shared()->addChildTasks( std::make_shared(storageProvider, tempStorage), // continuing when indexers still running, even if there are no storages right now. @@ -538,7 +538,7 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di // add task that injects the remaining intermediate storages into the persistent storage taskSequential->addTask( - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( + std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask( std::make_shared(storageProvider, tempStorage) ) ); diff --git a/src/lib/utility/scheduling/TaskDecoratorDelay.cpp b/src/lib/utility/scheduling/TaskDecoratorDelay.cpp index 97ae7a96..4a3b5327 100644 --- a/src/lib/utility/scheduling/TaskDecoratorDelay.cpp +++ b/src/lib/utility/scheduling/TaskDecoratorDelay.cpp @@ -4,12 +4,13 @@ TaskDecoratorDelay::TaskDecoratorDelay(size_t delayMS) : m_delayMS(delayMS) - , m_delayComplete(delayMS == 0) + , m_delayComplete(false) { } void TaskDecoratorDelay::doEnter(std::shared_ptr blackboard) { + m_delayComplete = (m_delayMS == 0); m_start = TimeStamp::now(); } @@ -20,8 +21,8 @@ Task::TaskState TaskDecoratorDelay::doUpdate(std::shared_ptr blackbo return m_taskRunner->update(blackboard); } - const int SLEEP_TIME_MS = 25; - std::this_thread::sleep_for(std::chrono::microseconds(SLEEP_TIME_MS)); + const int SLEEP_TIME_MS = (m_delayMS / 3) + 1; + std::this_thread::sleep_for(std::chrono::milliseconds(SLEEP_TIME_MS)); m_delayComplete = (TimeStamp::now().deltaMS(m_start) >= m_delayMS); diff --git a/src/lib/utility/scheduling/TaskDecoratorRepeat.cpp b/src/lib/utility/scheduling/TaskDecoratorRepeat.cpp index 000d17fa..64d63517 100644 --- a/src/lib/utility/scheduling/TaskDecoratorRepeat.cpp +++ b/src/lib/utility/scheduling/TaskDecoratorRepeat.cpp @@ -1,8 +1,12 @@ #include "TaskDecoratorRepeat.h" -TaskDecoratorRepeat::TaskDecoratorRepeat(ConditionType condition, TaskState exitState) +#include +#include + +TaskDecoratorRepeat::TaskDecoratorRepeat(ConditionType condition, TaskState exitState, size_t delayMS) : m_condition(condition) , m_exitState(exitState) + , m_delayMS(delayMS) { } @@ -29,6 +33,8 @@ Task::TaskState TaskDecoratorRepeat::doUpdate(std::shared_ptr blackb break; } + std::this_thread::sleep_for(std::chrono::milliseconds(m_delayMS)); + return state; } diff --git a/src/lib/utility/scheduling/TaskDecoratorRepeat.h b/src/lib/utility/scheduling/TaskDecoratorRepeat.h index 3c7256f0..4907f8e3 100644 --- a/src/lib/utility/scheduling/TaskDecoratorRepeat.h +++ b/src/lib/utility/scheduling/TaskDecoratorRepeat.h @@ -15,7 +15,7 @@ public: CONDITION_WHILE_SUCCESS }; - TaskDecoratorRepeat(ConditionType condition, TaskState exitState); + TaskDecoratorRepeat(ConditionType condition, TaskState exitState, size_t delayMS); private: void doEnter(std::shared_ptr blackboard) override; @@ -25,6 +25,7 @@ private: const ConditionType m_condition; const TaskState m_exitState; + const size_t m_delayMS; }; #endif // TASK_DECORATOR_REPEAT_H diff --git a/src/lib/utility/scheduling/TaskGroupParallel.cpp b/src/lib/utility/scheduling/TaskGroupParallel.cpp index f810131a..575aed5a 100644 --- a/src/lib/utility/scheduling/TaskGroupParallel.cpp +++ b/src/lib/utility/scheduling/TaskGroupParallel.cpp @@ -78,9 +78,13 @@ void TaskGroupParallel::doTerminate() for (size_t i = 0; i < m_tasks.size(); i++) { m_tasks[i]->taskRunner->terminate(); + } + + for (size_t i = 0; i < m_tasks.size(); i++) + { if (m_tasks[i]->thread) { - m_tasks[i]->thread->detach(); + m_tasks[i]->thread->join(); m_tasks[i]->thread.reset(); } } diff --git a/src/lib/utility/scheduling/TaskReturnSuccessIf.h b/src/lib/utility/scheduling/TaskReturnSuccessIf.h index 876f3c14..8c19f4d8 100644 --- a/src/lib/utility/scheduling/TaskReturnSuccessIf.h +++ b/src/lib/utility/scheduling/TaskReturnSuccessIf.h @@ -44,9 +44,6 @@ void TaskReturnSuccessIf::doEnter(std::shared_ptr blackboard) template Task::TaskState TaskReturnSuccessIf::doUpdate(std::shared_ptr blackboard) { - const int SLEEP_TIME_MS = 25; - std::this_thread::sleep_for(std::chrono::microseconds(SLEEP_TIME_MS)); - T lhsValue = 0; blackboard->get(m_lhsValueName, lhsValue);