diff --git a/src/lib/data/indexer/TaskBuildIndex.cpp b/src/lib/data/indexer/TaskBuildIndex.cpp index 020f29f5..5d6fbb51 100644 --- a/src/lib/data/indexer/TaskBuildIndex.cpp +++ b/src/lib/data/indexer/TaskBuildIndex.cpp @@ -41,13 +41,13 @@ TaskBuildIndex::TaskBuildIndex( void TaskBuildIndex::doEnter(std::shared_ptr blackboard) { + blackboard->set("indexer_threads_started", true); + m_interprocessIndexingStatusManager.setIndexingInterrupted(false); m_indexingFileCount = 0; updateIndexingDialog(blackboard, std::vector()); - 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) 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) m_storageProvider->insert(is); } - blackboard->set("indexer_count", 0); + blackboard->set("indexer_threads_stopped", true); } void TaskBuildIndex::doReset(std::shared_ptr blackboard) diff --git a/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp b/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp index 9f432320..01ffded1 100644 --- a/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp +++ b/src/lib/data/indexer/TaskFillIndexerCommandQueue.cpp @@ -64,15 +64,12 @@ void TaskFillIndexerCommandsQueue::doEnter(std::shared_ptr blackboar } } - fillCommandQueue(); - blackboard->set("indexer_command_queue_started", true); } Task::TaskState TaskFillIndexerCommandsQueue::doUpdate(std::shared_ptr blackboard) { - fillCommandQueue(); - + if (!fillCommandQueue()) { std::lock_guard lock(m_commandsMutex); @@ -99,17 +96,23 @@ void TaskFillIndexerCommandsQueue::doReset(std::shared_ptr blackboar void TaskFillIndexerCommandsQueue::handleMessage(MessageInterruptTasks* message) { std::lock_guard lock(m_commandsMutex); + LOG_INFO("Discarding remaining " + std::to_string(m_indexerCommandProvider->size() + m_indexerCommandManager.indexerCommandCount()) + " indexer commands."); + std::queue 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 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; } diff --git a/src/lib/data/indexer/TaskFillIndexerCommandQueue.h b/src/lib/data/indexer/TaskFillIndexerCommandQueue.h index e7e19e94..dd613e11 100644 --- a/src/lib/data/indexer/TaskFillIndexerCommandQueue.h +++ b/src/lib/data/indexer/TaskFillIndexerCommandQueue.h @@ -30,7 +30,7 @@ protected: void handleMessage(MessageInterruptTasks* message) override; - void fillCommandQueue(); + bool fillCommandQueue(); private: std::unique_ptr m_indexerCommandProvider; diff --git a/src/lib/project/Project.cpp b/src/lib/project/Project.cpp index 7a2248c7..393f4217 100644 --- a/src/lib/project/Project.cpp +++ b/src/lib/project/Project.cpp @@ -469,7 +469,8 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di // add tasks for setting some variables on the blackboard that are used during indexing taskSequential->addTask(std::make_shared>("source_file_count", indexerCommandProvider->size())); taskSequential->addTask(std::make_shared>("indexed_source_file_count", 0)); - taskSequential->addTask(std::make_shared>("indexer_count", 0)); + taskSequential->addTask(std::make_shared>("indexer_threads_started", false)); + taskSequential->addTask(std::make_shared>("indexer_threads_stopped", false)); taskSequential->addTask(std::make_shared>("indexer_command_queue_started", false)); taskSequential->addTask(std::make_shared>("indexer_command_queue_stopped", false)); @@ -479,41 +480,35 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di std::shared_ptr taskParallelIndexing = std::make_shared(); taskParserWrapper->setTask(taskParallelIndexing); - // add task for indexing - if (indexerThreadCount > 0) - { - bool multiProcess = ApplicationSettings::getInstance()->getMultiProcessIndexingEnabled() && hasCxxSourceGroup(); - - 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>("indexer_command_queue_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) - ), - std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( - std::make_shared(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess) - ) - ) - ); - } - // add task for refilling the indexer command queue taskParallelIndexing->addTask( std::make_shared(m_appUUID, std::move(indexerCommandProvider), 20) ); + // add task for indexing + bool multiProcess = ApplicationSettings::getInstance()->getMultiProcessIndexingEnabled() && hasCxxSourceGroup(); + 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>("indexer_command_queue_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) + ), + std::make_shared(indexerThreadCount, storageProvider, dialogView, m_appUUID, multiProcess) + ) + ); + // add task for merging the intermediate storages 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>("indexer_count", TaskReturnSuccessIf::CONDITION_EQUALS, 0) + 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()->addChildTasks( std::make_shared(storageProvider), - std::make_shared>("indexer_count", TaskReturnSuccessIf::CONDITION_GREATER_THAN, 0) + std::make_shared>("indexer_threads_stopped", TaskReturnSuccessIf::CONDITION_EQUALS, false) ) ) ) @@ -524,17 +519,13 @@ void Project::buildIndex(const RefreshInfo& info, std::shared_ptr di std::make_shared()->addChildTasks( // block until there are indexers running std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( - std::make_shared>("indexer_count", TaskReturnSuccessIf::CONDITION_EQUALS, 0) + std::make_shared>("indexer_threads_started", TaskReturnSuccessIf::CONDITION_EQUALS, false) ), std::make_shared(TaskDecoratorRepeat::CONDITION_WHILE_SUCCESS, Task::STATE_SUCCESS)->addChildTask( - std::make_shared()->addChildTasks( - // stopping when indexer count is zero, regardless wether there are still storages left to insert. - std::make_shared>("indexer_count", TaskReturnSuccessIf::CONDITION_GREATER_THAN, 0), - std::make_shared()->addChildTasks( - std::make_shared(storageProvider, tempStorage), - // continuing when indexer count is greater than zero, even if there are no storages right now. - std::make_shared>("indexer_count", TaskReturnSuccessIf::CONDITION_GREATER_THAN, 0) - ) + std::make_shared()->addChildTasks( + std::make_shared(storageProvider, tempStorage), + // continuing when indexers still running, even if there are no storages right now. + std::make_shared>("indexer_threads_stopped", TaskReturnSuccessIf::CONDITION_EQUALS, false) ) ) )