logic: Fixed race conditions in indexing tasks
* Fixed TaskFillIndexerCommandQueue exiting right away when all commands fitted into queue on first call * Fixed TaskBuildIndex exiting right away when starting of indexer processes takes longer than the update cycle * Refactored Indexing Tasks Behavior Tree
This commit is contained in:
@@ -41,13 +41,13 @@ TaskBuildIndex::TaskBuildIndex(
|
||||
|
||||
void TaskBuildIndex::doEnter(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
blackboard->set<bool>("indexer_threads_started", true);
|
||||
|
||||
m_interprocessIndexingStatusManager.setIndexingInterrupted(false);
|
||||
|
||||
m_indexingFileCount = 0;
|
||||
updateIndexingDialog(blackboard, std::vector<FilePath>());
|
||||
|
||||
blackboard->set("indexer_count", (int)m_processCount);
|
||||
|
||||
std::wstring logFilePath;
|
||||
Logger* logger = LogManager::getInstance()->getLoggerByType("FileLogger");
|
||||
if (logger)
|
||||
@@ -91,14 +91,14 @@ Task::TaskState TaskBuildIndex::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
updateIndexingDialog(blackboard, indexingFiles);
|
||||
}
|
||||
|
||||
if (runningThreadCount == 0)
|
||||
if (m_indexerCommandQueueStopped && runningThreadCount == 0)
|
||||
{
|
||||
return STATE_FAILURE;
|
||||
return STATE_SUCCESS;
|
||||
}
|
||||
else if (m_interrupted)
|
||||
{
|
||||
blackboard->set("interrupted_indexing", true);
|
||||
return STATE_FAILURE;
|
||||
return STATE_SUCCESS;
|
||||
}
|
||||
|
||||
if (fetchIntermediateStorages(blackboard))
|
||||
@@ -138,7 +138,7 @@ void TaskBuildIndex::doExit(std::shared_ptr<Blackboard> blackboard)
|
||||
m_storageProvider->insert(is);
|
||||
}
|
||||
|
||||
blackboard->set("indexer_count", 0);
|
||||
blackboard->set<bool>("indexer_threads_stopped", true);
|
||||
}
|
||||
|
||||
void TaskBuildIndex::doReset(std::shared_ptr<Blackboard> blackboard)
|
||||
|
||||
@@ -64,15 +64,12 @@ void TaskFillIndexerCommandsQueue::doEnter(std::shared_ptr<Blackboard> blackboar
|
||||
}
|
||||
}
|
||||
|
||||
fillCommandQueue();
|
||||
|
||||
blackboard->set<bool>("indexer_command_queue_started", true);
|
||||
}
|
||||
|
||||
Task::TaskState TaskFillIndexerCommandsQueue::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
fillCommandQueue();
|
||||
|
||||
if (!fillCommandQueue())
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(m_commandsMutex);
|
||||
|
||||
@@ -99,17 +96,23 @@ void TaskFillIndexerCommandsQueue::doReset(std::shared_ptr<Blackboard> blackboar
|
||||
void TaskFillIndexerCommandsQueue::handleMessage(MessageInterruptTasks* message)
|
||||
{
|
||||
std::lock_guard<std::mutex> lock(m_commandsMutex);
|
||||
|
||||
LOG_INFO("Discarding remaining " + std::to_string(m_indexerCommandProvider->size() + m_indexerCommandManager.indexerCommandCount()) + " indexer commands.");
|
||||
|
||||
std::queue<FilePath> empty;
|
||||
std::swap(m_filePathQueue, empty);
|
||||
|
||||
m_indexerCommandProvider->clear();
|
||||
m_indexerCommandManager.clearIndexerCommands();
|
||||
|
||||
LOG_INFO("Remaining: " + std::to_string(m_indexerCommandProvider->size() + m_indexerCommandManager.indexerCommandCount()) + ".");
|
||||
}
|
||||
|
||||
void TaskFillIndexerCommandsQueue::fillCommandQueue()
|
||||
bool TaskFillIndexerCommandsQueue::fillCommandQueue()
|
||||
{
|
||||
bool filled = false;
|
||||
std::lock_guard<std::mutex> lock(m_commandsMutex);
|
||||
|
||||
while (!m_indexerCommandProvider->empty() && m_indexerCommandManager.indexerCommandCount() < m_maximumQueueSize)
|
||||
{
|
||||
if (!m_filePathQueue.empty())
|
||||
@@ -121,5 +124,9 @@ void TaskFillIndexerCommandsQueue::fillCommandQueue()
|
||||
{
|
||||
m_indexerCommandManager.pushIndexerCommands({ m_indexerCommandProvider->consumeCommand() });
|
||||
}
|
||||
|
||||
filled = true;
|
||||
}
|
||||
|
||||
return filled;
|
||||
}
|
||||
|
||||
@@ -30,7 +30,7 @@ protected:
|
||||
|
||||
void handleMessage(MessageInterruptTasks* message) override;
|
||||
|
||||
void fillCommandQueue();
|
||||
bool fillCommandQueue();
|
||||
|
||||
private:
|
||||
std::unique_ptr<IndexerCommandProvider> m_indexerCommandProvider;
|
||||
|
||||
+21
-30
@@ -469,7 +469,8 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
|
||||
// add tasks for setting some variables on the blackboard that are used during indexing
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<int>>("source_file_count", indexerCommandProvider->size()));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<int>>("indexed_source_file_count", 0));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<int>>("indexer_count", 0));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<bool>>("indexer_threads_started", false));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<bool>>("indexer_threads_stopped", false));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<bool>>("indexer_command_queue_started", false));
|
||||
taskSequential->addTask(std::make_shared<TaskSetValue<bool>>("indexer_command_queue_stopped", false));
|
||||
|
||||
@@ -479,41 +480,35 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
|
||||
std::shared_ptr<TaskGroupParallel> taskParallelIndexing = std::make_shared<TaskGroupParallel>();
|
||||
taskParserWrapper->setTask(taskParallelIndexing);
|
||||
|
||||
// add task for indexing
|
||||
if (indexerThreadCount > 0)
|
||||
{
|
||||
bool multiProcess = ApplicationSettings::getInstance()->getMultiProcessIndexingEnabled() && hasCxxSourceGroup();
|
||||
|
||||
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<TaskReturnSuccessIf<bool>>("indexer_command_queue_started", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
|
||||
),
|
||||
std::make_shared<TaskDecoratorRepeat>(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask(
|
||||
std::make_shared<TaskBuildIndex>(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess)
|
||||
)
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
// add task for refilling the indexer command queue
|
||||
taskParallelIndexing->addTask(
|
||||
std::make_shared<TaskFillIndexerCommandsQueue>(m_appUUID, std::move(indexerCommandProvider), 20)
|
||||
);
|
||||
|
||||
// add task for indexing
|
||||
bool multiProcess = ApplicationSettings::getInstance()->getMultiProcessIndexingEnabled() && hasCxxSourceGroup();
|
||||
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<TaskReturnSuccessIf<bool>>("indexer_command_queue_started", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
|
||||
),
|
||||
std::make_shared<TaskBuildIndex>(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess)
|
||||
)
|
||||
);
|
||||
|
||||
// add task for merging the intermediate storages
|
||||
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<TaskReturnSuccessIf<int>>("indexer_count", TaskReturnSuccessIf<int>::CONDITION_EQUALS, 0)
|
||||
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<TaskGroupSelector>()->addChildTasks(
|
||||
std::make_shared<TaskMergeStorages>(storageProvider),
|
||||
std::make_shared<TaskReturnSuccessIf<int>>("indexer_count", TaskReturnSuccessIf<int>::CONDITION_GREATER_THAN, 0)
|
||||
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_threads_stopped", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
|
||||
)
|
||||
)
|
||||
)
|
||||
@@ -524,17 +519,13 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr<DialogView> di
|
||||
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<TaskReturnSuccessIf<int>>("indexer_count", TaskReturnSuccessIf<int>::CONDITION_EQUALS, 0)
|
||||
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<TaskGroupSequence>()->addChildTasks(
|
||||
// stopping when indexer count is zero, regardless wether there are still storages left to insert.
|
||||
std::make_shared<TaskReturnSuccessIf<int>>("indexer_count", TaskReturnSuccessIf<int>::CONDITION_GREATER_THAN, 0),
|
||||
std::make_shared<TaskGroupSelector>()->addChildTasks(
|
||||
std::make_shared<TaskInjectStorage>(storageProvider, tempStorage),
|
||||
// continuing when indexer count is greater than zero, even if there are no storages right now.
|
||||
std::make_shared<TaskReturnSuccessIf<int>>("indexer_count", TaskReturnSuccessIf<int>::CONDITION_GREATER_THAN, 0)
|
||||
)
|
||||
std::make_shared<TaskGroupSelector>()->addChildTasks(
|
||||
std::make_shared<TaskInjectStorage>(storageProvider, tempStorage),
|
||||
// continuing when indexers still running, even if there are no storages right now.
|
||||
std::make_shared<TaskReturnSuccessIf<bool>>("indexer_threads_stopped", TaskReturnSuccessIf<bool>::CONDITION_EQUALS, false)
|
||||
)
|
||||
)
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user