logic: Fixed termination issues in task scheduling

* don't use thread::detach() in TaskGroupParallel::terminate to make it useful for tabs
* properly terminate TaskFillIndexerCommandQueue
* fixed wrong microsecond delays
* moved some task delays to TaskDecoratorRepeat for better transparency
This commit is contained in:
Eberhard Graether
2018-10-28 02:36:41 +01:00
parent aef4970aa2
commit 27627b6d05
10 changed files with 40 additions and 29 deletions
-7
View File
@@ -1,8 +1,5 @@
#include "TaskInjectStorage.h"
#include <chrono>
#include <thread>
#include "Storage.h"
#include "StorageProvider.h"
@@ -33,10 +30,6 @@ Task::TaskState TaskInjectStorage::doUpdate(std::shared_ptr<Blackboard> blackboa
}
}
}
else
{
std::this_thread::sleep_for(std::chrono::milliseconds(25));
}
return STATE_FAILURE;
}
-7
View File
@@ -1,8 +1,5 @@
#include "TaskMergeStorages.h"
#include <chrono>
#include <thread>
#include "StorageProvider.h"
TaskMergeStorages::TaskMergeStorages(
@@ -40,10 +37,6 @@ Task::TaskState TaskMergeStorages::doUpdate(std::shared_ptr<Blackboard> blackboa
}
}
}
else
{
std::this_thread::sleep_for(std::chrono::milliseconds(25));
}
return STATE_FAILURE;
}
@@ -71,6 +71,11 @@ void TaskFillIndexerCommandsQueue::doEnter(std::shared_ptr<Blackboard> blackboar
Task::TaskState TaskFillIndexerCommandsQueue::doUpdate(std::shared_ptr<Blackboard> blackboard)
{
if (m_interrupted)
{
return STATE_FAILURE;
}
if (!fillCommandQueue())
{
std::lock_guard<std::mutex> lock(m_commandsMutex);
@@ -93,6 +98,12 @@ void TaskFillIndexerCommandsQueue::doExit(std::shared_ptr<Blackboard> blackboard
void TaskFillIndexerCommandsQueue::doReset(std::shared_ptr<Blackboard> blackboard)
{
m_interrupted = false;
}
void TaskFillIndexerCommandsQueue::terminate()
{
m_interrupted = true;
}
void TaskFillIndexerCommandsQueue::handleMessage(MessageInterruptTasks* message)
@@ -27,6 +27,7 @@ protected:
TaskState doUpdate(std::shared_ptr<Blackboard> blackboard) override;
void doExit(std::shared_ptr<Blackboard> blackboard) override;
void doReset(std::shared_ptr<Blackboard> blackboard) override;
void terminate() override;
void handleMessage(MessageInterruptTasks* message) override;
@@ -35,9 +36,13 @@ protected:
private:
std::unique_ptr<IndexerCommandProvider> m_indexerCommandProvider;
InterprocessIndexerCommandManager m_indexerCommandManager;
const size_t m_maximumQueueSize;
std::queue<FilePath> m_filePathQueue;
std::mutex m_commandsMutex;
bool m_interrupted = false;
};
#endif // TASK_FILL_INDEXER_COMMAND_QUEUE_H
+6 -6
View File
@@ -488,7 +488,7 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
taskParallelIndexing->addChildTasks(
std::make_shared<TaskGroupSequence>()->addChildTasks(
// block until there are indexer commands to process
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask(
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_command_queue_started", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
),
std::make_shared<TaskBuildIndex>(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess)
@@ -499,11 +499,11 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
taskParallelIndexing->addTask(
std::make_shared<TaskGroupSequence>()->addChildTasks(
// block until there are indexers running
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask(
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_threads_started", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
),
// merge until all indexers stopped and nothing left to merge
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 250)->addChildTask(
std::make_shared<TaskGroupSelector>()->addChildTasks(
std::make_shared<TaskMergeStorages>(storageProvider),
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_threads_stopped", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
@@ -516,10 +516,10 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
taskParallelIndexing->addTask(
std::make_shared<TaskGroupSequence>()->addChildTasks(
// block until there are indexers running
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask(
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_threads_started", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
),
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask(
std::make_shared<TaskGroupSelector>()->addChildTasks(
std::make_shared<TaskInjectStorage>(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<DialogView> di
// add task that injects the remaining intermediate storages into the persistent storage
taskSequential->addTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS, 25)->addChildTask(
std::make_shared<TaskInjectStorage>(storageProvider, tempStorage)
)
);
@@ -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> blackboard)
{
m_delayComplete = (m_delayMS == 0);
m_start = TimeStamp::now();
}
@@ -20,8 +21,8 @@ Task::TaskState TaskDecoratorDelay::doUpdate(std::shared_ptr<Blackboard> 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);
@@ -1,8 +1,12 @@
#include "TaskDecoratorRepeat.h"
TaskDecoratorRepeat::TaskDecoratorRepeat(ConditionType condition, TaskState exitState)
#include <chrono>
#include <thread>
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<Blackboard> blackb
break;
}
std::this_thread::sleep_for(std::chrono::milliseconds(m_delayMS));
return state;
}
@@ -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> 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
@@ -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();
}
}
@@ -44,9 +44,6 @@ void TaskReturnSuccessIf<T>::doEnter(std::shared_ptr<Blackboard> blackboard)
template <typename T>
Task::TaskState TaskReturnSuccessIf<T>::doUpdate(std::shared_ptr<Blackboard> blackboard)
{
const int SLEEP_TIME_MS = 25;
std::this_thread::sleep_for(std::chrono::microseconds(SLEEP_TIME_MS));
T lhsValue = 0;
blackboard->get<T>(m_lhsValueName, lhsValue);