logic: improved indexer parallelization
* added task that injects into PersistentStorage parallel to indexer tasks. * added Blackboard for Task classes
This commit is contained in:
@@ -0,0 +1,31 @@
|
||||
#include "utility/scheduling/Blackboard.h"
|
||||
|
||||
Blackboard::Blackboard()
|
||||
{
|
||||
}
|
||||
|
||||
Blackboard::Blackboard(std::shared_ptr<Blackboard> parent)
|
||||
: m_parent(parent)
|
||||
{
|
||||
}
|
||||
|
||||
Blackboard::~Blackboard()
|
||||
{
|
||||
}
|
||||
|
||||
bool Blackboard::exists(const std::string& key)
|
||||
{
|
||||
ItemMap::const_iterator it = m_values.find(key);
|
||||
return (it != m_values.end());
|
||||
}
|
||||
|
||||
bool Blackboard::clear(const std::string& key)
|
||||
{
|
||||
ItemMap::const_iterator it = m_values.find(key);
|
||||
if (it != m_values.end())
|
||||
{
|
||||
m_values.erase(it);
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
#ifndef BLACKBOARD_H
|
||||
#define BLACKBOARD_H
|
||||
|
||||
#include <map>
|
||||
#include <memory>
|
||||
#include <string>
|
||||
|
||||
#include "Utility/UtilityString.h"
|
||||
#include "utility/logging/logging.h"
|
||||
|
||||
struct BlackboardItemBase
|
||||
{
|
||||
virtual ~BlackboardItemBase()
|
||||
{
|
||||
}
|
||||
};
|
||||
|
||||
template <typename T>
|
||||
struct BlackboardItem: public BlackboardItemBase
|
||||
{
|
||||
BlackboardItem(const T &v)
|
||||
: value(v)
|
||||
{
|
||||
}
|
||||
|
||||
virtual ~BlackboardItem()
|
||||
{
|
||||
}
|
||||
|
||||
T value;
|
||||
};
|
||||
|
||||
class Blackboard
|
||||
{
|
||||
public:
|
||||
Blackboard();
|
||||
Blackboard(std::shared_ptr<Blackboard> parent);
|
||||
~Blackboard();
|
||||
|
||||
template <typename T>
|
||||
void set(const std::string& key, const T& value);
|
||||
|
||||
template <typename T>
|
||||
bool get(const std::string& key, T& value);
|
||||
|
||||
bool exists(const std::string& key);
|
||||
bool clear(const std::string& key);
|
||||
|
||||
private:
|
||||
typedef std::map<std::string, std::shared_ptr<BlackboardItemBase>> ItemMap;
|
||||
|
||||
std::shared_ptr<Blackboard> m_parent;
|
||||
ItemMap m_values;
|
||||
};
|
||||
|
||||
|
||||
template <typename T>
|
||||
void Blackboard::set(const std::string& key, const T& value)
|
||||
{
|
||||
m_values[key] = std::make_shared<BlackboardItem<T>>(value);
|
||||
}
|
||||
|
||||
template <typename T>
|
||||
bool Blackboard::get(const std::string& key, T& value)
|
||||
{
|
||||
ItemMap::const_iterator it = m_values.find(key);
|
||||
if (it != m_values.end())
|
||||
{
|
||||
std::shared_ptr<BlackboardItem<T>> item = std::dynamic_pointer_cast<BlackboardItem<T>>(it->second);
|
||||
if (item)
|
||||
{
|
||||
value = item->value;
|
||||
return true;
|
||||
}
|
||||
}
|
||||
if (m_parent)
|
||||
{
|
||||
return m_parent->get(key, value);
|
||||
}
|
||||
|
||||
LOG_WARNING("Entry for \"" + key + "\" not found on blackboard.");
|
||||
return false;
|
||||
}
|
||||
|
||||
#endif // BLACKBOARD_H
|
||||
@@ -23,28 +23,28 @@ Task::~Task()
|
||||
{
|
||||
}
|
||||
|
||||
Task::TaskState Task::update()
|
||||
Task::TaskState Task::update(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
if (!m_enterCalled)
|
||||
{
|
||||
doEnter();
|
||||
doEnter(blackboard);
|
||||
m_enterCalled = true;
|
||||
}
|
||||
|
||||
TaskState state = doUpdate();
|
||||
TaskState state = doUpdate(blackboard);
|
||||
|
||||
if (state != STATE_RUNNING && !m_exitCalled)
|
||||
{
|
||||
doExit();
|
||||
doExit(blackboard);
|
||||
m_exitCalled = true;
|
||||
}
|
||||
|
||||
return state;
|
||||
}
|
||||
|
||||
void Task::reset()
|
||||
void Task::reset(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
doReset();
|
||||
doReset(blackboard);
|
||||
m_enterCalled = false;
|
||||
m_exitCalled = false;
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
|
||||
#include <memory>
|
||||
|
||||
class Blackboard;
|
||||
|
||||
class Task
|
||||
{
|
||||
public:
|
||||
@@ -19,16 +21,14 @@ public:
|
||||
Task();
|
||||
virtual ~Task();
|
||||
|
||||
// virtual TaskState getState() const = 0;
|
||||
|
||||
TaskState update();
|
||||
void reset();
|
||||
TaskState update(std::shared_ptr<Blackboard> blackboard);
|
||||
void reset(std::shared_ptr<Blackboard> blackboard);
|
||||
|
||||
private:
|
||||
virtual void doEnter() = 0;
|
||||
virtual Task::TaskState doUpdate() = 0;
|
||||
virtual void doExit() = 0;
|
||||
virtual void doReset() = 0;
|
||||
virtual void doEnter(std::shared_ptr<Blackboard> blackboard) = 0;
|
||||
virtual Task::TaskState doUpdate(std::shared_ptr<Blackboard> blackboard) = 0;
|
||||
virtual void doExit(std::shared_ptr<Blackboard> blackboard) = 0;
|
||||
virtual void doReset(std::shared_ptr<Blackboard> blackboard) = 0;
|
||||
|
||||
bool m_enterCalled;
|
||||
bool m_exitCalled;
|
||||
|
||||
@@ -16,7 +16,7 @@ void TaskGroupParallel::addTask(std::shared_ptr<Task> task)
|
||||
m_tasks.push_back(std::make_shared<TaskInfo>(std::make_shared<TaskRunner>(task)));
|
||||
}
|
||||
|
||||
void TaskGroupParallel::doEnter()
|
||||
void TaskGroupParallel::doEnter(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
m_taskFailed = false;
|
||||
|
||||
@@ -26,14 +26,14 @@ void TaskGroupParallel::doEnter()
|
||||
m_activeTaskCount = 0;
|
||||
for (size_t i = 0; i < m_tasks.size(); i++)
|
||||
{
|
||||
m_tasks[i]->thread = std::make_shared<std::thread>(&TaskGroupParallel::processTaskThreaded, this, m_tasks[i]);
|
||||
m_tasks[i]->thread = std::make_shared<std::thread>(&TaskGroupParallel::processTaskThreaded, this, m_tasks[i], blackboard);
|
||||
m_tasks[i]->active = true;
|
||||
m_activeTaskCount++;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Task::TaskState TaskGroupParallel::doUpdate()
|
||||
Task::TaskState TaskGroupParallel::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
if (m_tasks.size() != 0 && getActveTaskCount() > 0)
|
||||
{
|
||||
@@ -43,7 +43,7 @@ Task::TaskState TaskGroupParallel::doUpdate()
|
||||
return (m_taskFailed ? STATE_FAILURE : STATE_SUCCESS);
|
||||
}
|
||||
|
||||
void TaskGroupParallel::doExit()
|
||||
void TaskGroupParallel::doExit(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
for (size_t i = 0; i < m_tasks.size(); i++)
|
||||
{
|
||||
@@ -52,7 +52,7 @@ void TaskGroupParallel::doExit()
|
||||
}
|
||||
}
|
||||
|
||||
void TaskGroupParallel::doReset()
|
||||
void TaskGroupParallel::doReset(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
for (size_t i = 0; i < m_tasks.size(); i++)
|
||||
{
|
||||
@@ -60,14 +60,14 @@ void TaskGroupParallel::doReset()
|
||||
if (!m_tasks[i]->active)
|
||||
{
|
||||
m_tasks[i]->thread->join();
|
||||
m_tasks[i]->thread = std::make_shared<std::thread>(&TaskGroupParallel::processTaskThreaded, this, m_tasks[i]);
|
||||
m_tasks[i]->thread = std::make_shared<std::thread>(&TaskGroupParallel::processTaskThreaded, this, m_tasks[i], blackboard);
|
||||
m_tasks[i]->active = true;
|
||||
m_activeTaskCount++;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void TaskGroupParallel::processTaskThreaded(std::shared_ptr<TaskInfo> taskInfo)
|
||||
void TaskGroupParallel::processTaskThreaded(std::shared_ptr<TaskInfo> taskInfo, std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
ScopedFunctor functor([&](){
|
||||
std::lock_guard<std::mutex> lock(m_activeTaskCountMutex);
|
||||
@@ -77,7 +77,7 @@ void TaskGroupParallel::processTaskThreaded(std::shared_ptr<TaskInfo> taskInfo)
|
||||
|
||||
while (true)
|
||||
{
|
||||
TaskState state = taskInfo->taskRunner->update();
|
||||
TaskState state = taskInfo->taskRunner->update(blackboard);
|
||||
|
||||
if (state != STATE_RUNNING)
|
||||
{
|
||||
|
||||
@@ -29,12 +29,12 @@ private:
|
||||
volatile bool active;
|
||||
};
|
||||
|
||||
virtual void doEnter();
|
||||
virtual TaskState doUpdate();
|
||||
virtual void doExit();
|
||||
virtual void doReset();
|
||||
virtual void doEnter(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual TaskState doUpdate(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doExit(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doReset(std::shared_ptr<Blackboard> blackboard);
|
||||
|
||||
void processTaskThreaded(std::shared_ptr<TaskInfo> taskInfo);
|
||||
void processTaskThreaded(std::shared_ptr<TaskInfo> taskInfo, std::shared_ptr<Blackboard> blackboard);
|
||||
int getActveTaskCount() const;
|
||||
|
||||
std::vector<std::shared_ptr<TaskInfo>> m_tasks;
|
||||
|
||||
@@ -14,12 +14,12 @@ void TaskGroupSequential::addTask(std::shared_ptr<Task> task)
|
||||
m_taskRunners.push_back(std::make_shared<TaskRunner>(task));
|
||||
}
|
||||
|
||||
void TaskGroupSequential::doEnter()
|
||||
void TaskGroupSequential::doEnter(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
m_taskIndex = 0;
|
||||
}
|
||||
|
||||
Task::TaskState TaskGroupSequential::doUpdate()
|
||||
Task::TaskState TaskGroupSequential::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
if (m_taskIndex >= int(m_taskRunners.size()))
|
||||
{
|
||||
@@ -30,7 +30,7 @@ Task::TaskState TaskGroupSequential::doUpdate()
|
||||
return STATE_FAILURE;
|
||||
}
|
||||
|
||||
TaskState state = m_taskRunners[m_taskIndex]->update();
|
||||
TaskState state = m_taskRunners[m_taskIndex]->update(blackboard);
|
||||
|
||||
if (state == STATE_SUCCESS)
|
||||
{
|
||||
@@ -44,11 +44,11 @@ Task::TaskState TaskGroupSequential::doUpdate()
|
||||
return STATE_RUNNING;
|
||||
}
|
||||
|
||||
void TaskGroupSequential::doExit()
|
||||
void TaskGroupSequential::doExit(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
void TaskGroupSequential::doReset()
|
||||
void TaskGroupSequential::doReset(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
for (size_t i = 0; i < m_taskRunners.size(); i++)
|
||||
{
|
||||
|
||||
@@ -14,10 +14,10 @@ public:
|
||||
virtual void addTask(std::shared_ptr<Task> task);
|
||||
|
||||
private:
|
||||
virtual void doEnter();
|
||||
virtual TaskState doUpdate();
|
||||
virtual void doExit();
|
||||
virtual void doReset();
|
||||
virtual void doEnter(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual TaskState doUpdate(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doExit(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doReset(std::shared_ptr<Blackboard> blackboard);
|
||||
|
||||
std::vector<std::shared_ptr<TaskRunner>> m_taskRunners;
|
||||
int m_taskIndex;
|
||||
|
||||
@@ -9,20 +9,20 @@ TaskLambda::~TaskLambda()
|
||||
{
|
||||
}
|
||||
|
||||
void TaskLambda::doEnter()
|
||||
void TaskLambda::doEnter(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
Task::TaskState TaskLambda::doUpdate()
|
||||
Task::TaskState TaskLambda::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
m_func();
|
||||
return STATE_SUCCESS;
|
||||
}
|
||||
|
||||
void TaskLambda::doExit()
|
||||
void TaskLambda::doExit(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
void TaskLambda::doReset()
|
||||
void TaskLambda::doReset(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
@@ -13,10 +13,10 @@ public:
|
||||
virtual ~TaskLambda();
|
||||
|
||||
private:
|
||||
virtual void doEnter();
|
||||
virtual TaskState doUpdate();
|
||||
virtual void doExit();
|
||||
virtual void doReset();
|
||||
virtual void doEnter(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual TaskState doUpdate(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doExit(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doReset(std::shared_ptr<Blackboard> blackboard);
|
||||
|
||||
std::function<void()> m_func;
|
||||
};
|
||||
|
||||
@@ -0,0 +1,39 @@
|
||||
#include "utility/scheduling/TaskRepeatWhileSuccess.h"
|
||||
|
||||
TaskRepeatWhileSuccess::TaskRepeatWhileSuccess()
|
||||
{
|
||||
}
|
||||
|
||||
void TaskRepeatWhileSuccess::setTask(std::shared_ptr<Task> task)
|
||||
{
|
||||
if (task)
|
||||
{
|
||||
m_taskRunner = std::make_shared<TaskRunner>(task);
|
||||
}
|
||||
}
|
||||
|
||||
void TaskRepeatWhileSuccess::doEnter(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
Task::TaskState TaskRepeatWhileSuccess::doUpdate(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
TaskState state = m_taskRunner->update(blackboard);
|
||||
|
||||
if (state == Task::STATE_SUCCESS)
|
||||
{
|
||||
state = Task::STATE_RUNNING;
|
||||
m_taskRunner->reset();
|
||||
}
|
||||
|
||||
return state;
|
||||
}
|
||||
|
||||
void TaskRepeatWhileSuccess::doExit(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
}
|
||||
|
||||
void TaskRepeatWhileSuccess::doReset(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
m_taskRunner->reset();
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
#ifndef TASK_REPEAT_WHILE_SUCCESS_H
|
||||
#define TASK_REPEAT_WHILE_SUCCESS_H
|
||||
|
||||
#include <vector>
|
||||
|
||||
#include "utility/scheduling/TaskDecorator.h"
|
||||
#include "utility/scheduling/TaskRunner.h"
|
||||
|
||||
class TaskRepeatWhileSuccess
|
||||
: public TaskDecorator
|
||||
{
|
||||
public:
|
||||
TaskRepeatWhileSuccess();
|
||||
|
||||
virtual void setTask(std::shared_ptr<Task> task);
|
||||
|
||||
private:
|
||||
virtual void doEnter(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual TaskState doUpdate(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doExit(std::shared_ptr<Blackboard> blackboard);
|
||||
virtual void doReset(std::shared_ptr<Blackboard> blackboard);
|
||||
|
||||
std::shared_ptr<TaskRunner> m_taskRunner;
|
||||
};
|
||||
|
||||
#endif // TASK_REPEAT_WHILE_SUCCESS_H
|
||||
@@ -10,20 +10,15 @@ TaskRunner::~TaskRunner()
|
||||
{
|
||||
}
|
||||
|
||||
//Task::TaskState TaskRunner::getState() const
|
||||
//{
|
||||
// return m_task->getState();
|
||||
//}
|
||||
|
||||
Task::TaskState TaskRunner::update()
|
||||
Task::TaskState TaskRunner::update(std::shared_ptr<Blackboard> blackboard)
|
||||
{
|
||||
if (m_reset)
|
||||
{
|
||||
m_task->reset();
|
||||
m_task->reset(blackboard);
|
||||
m_reset = false;
|
||||
}
|
||||
|
||||
return m_task->update();
|
||||
return m_task->update(blackboard);
|
||||
}
|
||||
|
||||
void TaskRunner::reset()
|
||||
|
||||
@@ -11,9 +11,7 @@ public:
|
||||
TaskRunner(std::shared_ptr<Task> task);
|
||||
~TaskRunner();
|
||||
|
||||
//Task::TaskState getState() const;
|
||||
|
||||
Task::TaskState update();
|
||||
Task::TaskState update(std::shared_ptr<Blackboard> blackboard);
|
||||
void reset();
|
||||
|
||||
private:
|
||||
|
||||
@@ -5,6 +5,8 @@
|
||||
|
||||
#include "utility/logging/logging.h"
|
||||
#include "utility/messaging/type/MessageStatus.h"
|
||||
#include "utility/scheduling/Blackboard.h"
|
||||
#include "utility/ScopedFunctor.h"
|
||||
|
||||
std::shared_ptr<TaskScheduler> TaskScheduler::getInstance()
|
||||
{
|
||||
@@ -140,15 +142,22 @@ void TaskScheduler::processTasks()
|
||||
{
|
||||
std::shared_ptr<TaskRunner> runner = m_taskRunners.front();
|
||||
|
||||
m_tasksMutex.unlock();
|
||||
|
||||
Task::TaskState state = runner->update();
|
||||
|
||||
m_tasksMutex.lock();
|
||||
|
||||
if (state != Task::STATE_RUNNING)
|
||||
{
|
||||
m_taskRunners.pop_front();
|
||||
m_tasksMutex.unlock();
|
||||
ScopedFunctor functor([this](){
|
||||
m_tasksMutex.lock();
|
||||
});
|
||||
|
||||
std::shared_ptr<Blackboard> blackboard = std::make_shared<Blackboard>();
|
||||
while (true)
|
||||
{
|
||||
if (runner->update(blackboard) != Task::STATE_RUNNING)
|
||||
{
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
m_taskRunners.pop_front();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user