5#include <experimental/memory>
10#include <shared_mutex>
17#include <Ice/Current.h>
18#include <IceUtil/Optional.h>
26#include <RobotAPI/interface/aron/Aron.h>
27#include <RobotAPI/interface/skills/SkillManagerInterface.h>
28#include <RobotAPI/interface/skills/SkillProviderInterface.h>
53 p.getProxy(myPrx, -1);
60 const std::string providerName = p.getName();
64 .providerInterface = myPrx,
67 ARMARX_INFO <<
"Adding provider to manager: " << i.providerId;
68 manager->addProvider(i.toIce());
75 std::string providerName = p.getName();
77 auto id = skills::manager::dto::ProviderID{providerName};
78 manager->removeProvider(
id);
84 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> drained;
86 const std::unique_lock l(skillExecutionsMutex);
87 for (
auto& [
id, runtime] : skillExecutions)
89 drained.push_back(std::move(runtime));
91 skillExecutions.clear();
94 <<
" skill executions to finish before disconnect.";
95 for (
auto& runtime : drained)
97 if (not runtime or not runtime->execution.joinable())
101 runtime->stopSkill();
103 const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(5);
104 bool terminated =
false;
105 while (not terminated and std::chrono::steady_clock::now() < deadline)
108 std::scoped_lock l(runtime->skillStatusesMutex);
109 terminated = runtime->statusUpdate.hasBeenTerminated();
113 std::this_thread::sleep_for(std::chrono::milliseconds(20));
119 runtime->execution.join();
123 ARMARX_WARNING <<
"Skill execution '" << runtime->statusUpdate.executionId.skillId
124 <<
"' did not terminate within 5s during disconnect. Detaching.";
125 runtime->execution.detach();
132 skillFactories.clear();
138 std::string
prefix =
"skill.";
139 properties->component(
143 "The name of the SkillManager (or SkillMemory) proxy this provider belongs to.");
144 properties->topic<armarx::skills::SkillEventListenerInterface>(
145 "SkillEventListener",
prefix +
"tpc.sub.SkillEventListener");
157 const std::string componentName = p.getName();
162 const std::unique_lock l(skillFactoriesMutex);
163 auto skillId = fac->createSkillDescription(providerId).skillId;
165 if (skillFactories.find(skillId) != skillFactories.end())
167 ARMARX_WARNING <<
"Try to add a skill factory for skill '" + skillId.toString() +
168 "' which already exists in list. Ignoring this skill.";
172 ARMARX_INFO <<
"Adding skill `" << skillId <<
"` to component `" << componentName <<
"` .";
174 skillFactories.emplace(skillId, std::move(fac));
199 if (skillFactories.count(skillId) == 0)
201 ARMARX_INFO <<
"Could not find a skill factory for id: " << skillId;
205 auto* facPtr = skillFactories.at(skillId).get();
209 std::optional<skills::SkillStatusUpdate>
215 const std::shared_lock l(skillExecutionsMutex);
216 auto it = skillExecutions.find(execId);
217 if (it == skillExecutions.end())
224 std::scoped_lock l2{it->second->skillStatusesMutex};
225 return it->second->statusUpdate;
228 std::map<skills::SkillExecutionID, skills::SkillStatusUpdate>
231 std::map<skills::SkillExecutionID, skills::SkillStatusUpdate> skillUpdates;
233 const std::shared_lock l(skillExecutionsMutex);
234 for (
const auto& [key, impl] : skillExecutions)
236 const std::scoped_lock l2(impl->skillStatusesMutex);
237 skillUpdates.insert({key, impl->statusUpdate});
242 std::optional<skills::SkillDescription>
247 const std::shared_lock l(skillFactoriesMutex);
248 if (skillFactories.find(skillId) == skillFactories.end())
250 std::stringstream ss;
251 ss <<
"Skill description for skill '" + skillId.
toString() +
252 "' not found! Found instead: {"
254 for (
const auto& [k, _] : skillFactories)
256 ss <<
"\t" << k.toString() <<
"\n";
264 return skillFactories.at(skillId)->createSkillDescription(*skillId.
providerId);
267 std::map<skills::SkillID, skills::SkillDescription>
270 std::map<skills::SkillID, skills::SkillDescription> skillDesciptions;
271 const std::shared_lock l(skillFactoriesMutex);
272 for (
const auto& [key, fac] : skillFactories)
275 skillDesciptions.insert({key, fac->createSkillDescription(*key.providerId)});
277 return skillDesciptions;
288 .executionStartedTime =
294 std::shared_ptr<skills::detail::SkillRuntime> runtime;
295 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> finishedToJoin;
297 auto l1 = std::unique_lock{skillFactoriesMutex};
299 const auto& fac = getSkillFactory(executionId.skillId);
300 ARMARX_CHECK(fac) <<
"Could not find a factory for skill " << executionId.skillId;
303 const std::unique_lock l2{skillExecutionsMutex};
308 finishedToJoin = collectFinishedExecutions_locked();
310 runtime = std::make_shared<skills::detail::SkillRuntime>(
315 skillExecutions.emplace(executionId, runtime);
319 runtime->setLocalMinimumLoggingLevel(loggingLevel);
326 runtime->execution = std::thread(
331 auto x = runtime->executeSkill();
332 ret.result =
x.result;
335 catch (std::exception& e)
338 "skill. Exception was: "
346 for (
auto& finished : finishedToJoin)
348 if (finished && finished->execution.joinable())
350 finished->execution.join();
353 finishedToJoin.clear();
355 if (runtime && runtime->execution.joinable())
357 runtime->execution.join();
370 std::shared_ptr<skills::detail::SkillRuntime> runtime;
371 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> finishedToJoin;
373 auto l1 = std::unique_lock{skillFactoriesMutex};
375 const auto& fac = getSkillFactory(executionRequest.
skillId);
379 const std::unique_lock l2{skillExecutionsMutex};
382 finishedToJoin = collectFinishedExecutions_locked();
388 if (skillExecutions.count(executionId) > 0)
390 ARMARX_ERROR <<
"SkillsExecutionID already exists! This is undefined behaviour "
391 "and should not occur!";
394 runtime = std::make_shared<skills::detail::SkillRuntime>(
399 skillExecutions.emplace(executionId, runtime);
401 ARMARX_INFO <<
"Setting skill runtime's logging level to `"
403 runtime->setLocalMinimumLoggingLevel(loggingLevel);
409 runtime->execution = std::thread(
414 auto x = runtime->executeSkill();
417 catch (std::exception& e)
420 "skill. Exception was: "
428 for (
auto& finished : finishedToJoin)
430 if (finished && finished->execution.joinable())
432 finished->execution.join();
435 finishedToJoin.clear();
441 std::scoped_lock l(runtime->skillStatusesMutex);
443 if (runtime->statusUpdate.hasBeenConstructed())
449 std::this_thread::sleep_for(std::chrono::milliseconds(20));
461 std::shared_ptr<skills::detail::SkillRuntime> runtime;
463 std::shared_lock l{skillExecutionsMutex};
464 auto it = skillExecutions.find(executionId);
465 if (it == skillExecutions.end())
468 "' found! Ignoring prepareSkill request.";
471 runtime = it->second;
474 std::scoped_lock l2{runtime->skillStatusesMutex};
478 "' because its not in preparing phase.";
482 runtime->updateSkillParameters(input);
491 std::shared_ptr<skills::detail::SkillRuntime> runtime;
493 std::shared_lock l(skillExecutionsMutex);
494 auto it = skillExecutions.find(executionId);
495 if (it == skillExecutions.end())
498 "' found! Ignoring abortSkill request.";
501 runtime = it->second;
504 runtime->stopSkill();
509 std::scoped_lock l2(runtime->skillStatusesMutex);
510 auto status = runtime->statusUpdate;
512 if (
status.hasBeenTerminated())
517 std::this_thread::sleep_for(std::chrono::milliseconds(20));
528 std::shared_ptr<skills::detail::SkillRuntime> runtime;
530 std::shared_lock l(skillExecutionsMutex);
531 auto it = skillExecutions.find(executionId);
532 if (it == skillExecutions.end())
535 "' found! Ignoring abortSkill request.";
538 runtime = it->second;
541 runtime->stopSkill();
552 std::shared_ptr<skills::detail::SkillRuntime>>>
555 std::shared_lock l(skillExecutionsMutex);
556 snapshot.reserve(skillExecutions.size());
557 for (
const auto& [
id, runtime] : skillExecutions)
559 snapshot.emplace_back(
id, runtime);
562 for (
auto& [
id, runtime] : snapshot)
564 ARMARX_DEBUG <<
"updating subskill status for " <<
id.toString() <<
" with "
566 runtime->updateSubSkillStatus(statusUpdate);
570 const skills::manager::dti::SkillManagerInterfacePrx&
576 std::vector<std::shared_ptr<skills::detail::SkillRuntime>>
577 SkillProviderComponentPlugin::collectFinishedExecutions_locked()
585 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> ret;
586 for (
auto it = skillExecutions.begin(); it != skillExecutions.end();)
588 const auto& runtime = it->second;
589 bool terminated =
false;
591 std::scoped_lock statusLock(runtime->skillStatusesMutex);
592 terminated = runtime->statusUpdate.hasBeenTerminated();
596 ret.push_back(std::move(it->second));
597 it = skillExecutions.erase(it);
615 IceUtil::Optional<skills::provider::dto::SkillDescription>
617 const skills::provider::dto::SkillID& skillId,
618 const Ice::Current& )
621 auto o = plugin->getSkillDescription(
id);
624 return o->toProviderIce();
629 skills::provider::dto::SkillDescriptionMap
632 skills::provider::dto::SkillDescriptionMap ret;
633 for (
const auto& [k, v] : plugin->getSkillDescriptions())
635 ret.insert({k.toProviderIce(), v.toProviderIce()});
640 IceUtil::Optional<skills::provider::dto::SkillStatusUpdate>
642 const skills::provider::dto::SkillExecutionID& executionId,
643 const Ice::Current& )
647 auto o = plugin->getSkillExecutionStatus(execId);
650 return o->toProviderIce();
655 skills::provider::dto::SkillStatusUpdateMap
658 skills::provider::dto::SkillStatusUpdateMap ret;
659 for (
const auto& [k, v] : plugin->getSkillExecutionStatuses())
661 ret.insert({k.toProviderIce(), v.toProviderIce()});
667 skills::provider::dto::SkillStatusUpdate
669 const skills::provider::dto::SkillExecutionRequest& info,
670 const Ice::Current& )
674 auto up = this->plugin->executeSkill(exec);
675 return up.toProviderIce();
678 skills::provider::dto::SkillExecutionID
680 const skills::provider::dto::SkillExecutionRequest& info,
681 const Ice::Current& current )
685 auto id = this->plugin->executeSkillAsync(exec);
689 skills::provider::dto::ParameterUpdateResult
691 const skills::provider::dto::SkillExecutionID&
id,
693 const Ice::Current& current )
695 skills::provider::dto::ParameterUpdateResult res;
700 res.success = this->plugin->updateSkillParameters(exec, prep);
704 skills::provider::dto::AbortSkillResult
706 const Ice::Current& )
708 skills::provider::dto::AbortSkillResult res;
711 res.success = this->plugin->abortSkill(exec);
715 skills::provider::dto::AbortSkillResult
717 const skills::provider::dto::SkillExecutionID&
id,
718 const Ice::Current& )
720 skills::provider::dto::AbortSkillResult res;
723 res.success = this->plugin->abortSkillAsync(exec);
729 const skills::provider::dto::SkillStatusUpdate& statusUpdate,
730 const std::string& providerName,
731 const Ice::Current& )
734 status.executionId.skillId.providerId.emplace().providerName = providerName;
735 plugin->updateSubSkillStatus(
status);
static std::string levelToString(MessageTypeT type)
MessageTypeT getEffectiveLoggingLevel() const
ManagedIceObject & parent()
const std::string & prefix() const
PluginT * addPlugin(const std::string prefix="", ParamsT &&... params)
std::string getName() const
Retrieve name of object.
IceUtil::Optional< skills::provider::dto::SkillStatusUpdate > getSkillExecutionStatus(const skills::provider::dto::SkillExecutionID &executionId, const Ice::Current ¤t=Ice::Current()) override
skills::provider::dto::SkillStatusUpdateMap getSkillExecutionStatuses(const Ice::Current ¤t=Ice::Current()) override
SkillProviderComponentPluginUser()
skills::provider::dto::AbortSkillResult abortSkill(const skills::provider::dto::SkillExecutionID &skill, const Ice::Current ¤t=Ice::Current()) override
IceUtil::Optional< skills::provider::dto::SkillDescription > getSkillDescription(const skills::provider::dto::SkillID &skill, const Ice::Current ¤t=Ice::Current()) override
skills::provider::dto::SkillStatusUpdate executeSkill(const skills::provider::dto::SkillExecutionRequest &executionInfo, const Ice::Current ¤t=Ice::Current()) override
skills::provider::dto::SkillExecutionID executeSkillAsync(const skills::provider::dto::SkillExecutionRequest &executionInfo, const Ice::Current ¤t=Ice::Current()) override
skills::provider::dto::AbortSkillResult abortSkillAsync(const skills::provider::dto::SkillExecutionID &skill, const Ice::Current ¤t=Ice::Current()) override
const std::experimental::observer_ptr< plugins::SkillProviderComponentPlugin > & getSkillProviderPlugin() const
void reportSkillEvent(const skills::provider::dto::SkillStatusUpdate &statusUpdate, const std::string &providerName, const Ice::Current ¤t) override
skills::provider::dto::ParameterUpdateResult updateSkillParameters(const skills::provider::dto::SkillExecutionID &executionId, const armarx::aron::data::dto::DictPtr ¶meters, const Ice::Current ¤t=Ice::Current()) override
skills::provider::dto::SkillDescriptionMap getSkillDescriptions(const Ice::Current ¤t=Ice::Current()) override
static PointerType FromAronDictDTO(const data::dto::DictPtr &aron)
skills::SkillExecutionID executeSkillAsync(const skills::SkillExecutionRequest &executionInfo)
bool abortSkill(const skills::SkillExecutionID &execId)
skills::SkillStatusUpdate executeSkill(const skills::SkillExecutionRequest &executionInfo)
std::optional< skills::SkillStatusUpdate > getSkillExecutionStatus(const skills::SkillExecutionID &) const
void preOnInitComponent() override
std::optional< skills::SkillDescription > getSkillDescription(const skills::SkillID &) const
void preOnConnectComponent() override
void postOnConnectComponent() override
void postCreatePropertyDefinitions(PropertyDefinitionsPtr &properties) override
const skills::manager::dti::SkillManagerInterfacePrx & skillManager() const
std::map< skills::SkillExecutionID, skills::SkillStatusUpdate > getSkillExecutionStatuses() const
void addSkillFactory(std::unique_ptr< skills::SkillBlueprint > &&)
bool updateSkillParameters(const skills::SkillExecutionID &id, const armarx::aron::data::DictPtr ¶ms)
std::map< skills::SkillID, skills::SkillDescription > getSkillDescriptions() const
void preOnDisconnectComponent() override
void updateSubSkillStatus(const skills::SkillStatusUpdate &statusUpdate)
bool abortSkillAsync(const skills::SkillExecutionID &execId)
std::function< TerminatedSkillStatus()> FunctionType
callback::dti::SkillProviderCallbackInterfacePrx callbackInterface
static SkillExecutionRequest FromIce(const manager::dto::SkillExecutionRequest &)
armarx::aron::data::DictPtr parameters
provider::dto::SkillExecutionRequest toProviderIce() const
std::string toString() const
std::optional< ProviderID > providerId
bool isFullySpecified() const
bool isSkillSpecified() const
static SkillID FromIce(const manager::dto::SkillID &)
#define ARMARX_CHECK(expression)
Shortcut for ARMARX_CHECK_EXPRESSION.
#define ARMARX_INFO
The normal logging level.
#define ARMARX_ERROR
The logging level for unexpected behaviour, that must be fixed.
#define ARMARX_DEBUG
The logging level for output that is only interesting while debugging.
#define ARMARX_WARNING
The logging level for unexpected behaviour, but not a serious problem.
#define ARMARX_VERBOSE
The logging level for verbose information.
::IceInternal::Handle< Dict > DictPtr
std::shared_ptr< Dict > DictPtr
This file is part of ArmarX.
SkillStatus toSkillStatus(const ActiveOrTerminatedSkillStatus &d)
This file offers overloads of toIce() and fromIce() functions for STL container types.
IceUtil::Handle< class PropertyDefinitionContainer > PropertyDefinitionsPtr
PropertyDefinitions smart pointer type.
std::string toString() const
static SkillExecutionID FromIce(const skills::manager::dto::SkillExecutionID &)
SkillExecutionID executionId
static SkillStatusUpdate FromIce(const provider::dto::SkillStatusUpdate &update, const std::optional< skills::ProviderID > &providerId=std::nullopt)