5#include <experimental/memory>
10#include <shared_mutex>
17#include <Ice/Current.h>
18#include <Ice/LocalException.h>
19#include <IceUtil/Optional.h>
28#include <RobotAPI/interface/aron/Aron.h>
29#include <RobotAPI/interface/skills/SkillManagerInterface.h>
30#include <RobotAPI/interface/skills/SkillProviderInterface.h>
55 p.getProxy(myPrx, -1);
62 const std::string providerName = p.getName();
66 .providerInterface = myPrx,
69 ARMARX_INFO <<
"Adding provider to manager: " << i.providerId;
70 manager->addProvider(i.toIce());
77 std::string providerName = p.getName();
79 auto id = skills::manager::dto::ProviderID{providerName};
82 manager->removeProvider(
id);
84 catch (
const Ice::LocalException& e)
87 <<
"' from the skill manager: " << e.what();
94 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> drained;
96 const std::unique_lock l(skillExecutionsMutex);
97 for (
auto& [
id, runtime] : skillExecutions)
99 drained.push_back(std::move(runtime));
101 skillExecutions.clear();
104 <<
" skill executions to finish before disconnect.";
105 for (
auto& runtime : drained)
107 if (not runtime or not runtime->execution.joinable())
113 runtime->stopSkill();
118 << runtime->statusUpdate.executionId.skillId
122 const auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(5);
123 bool terminated =
false;
124 while (not terminated and std::chrono::steady_clock::now() < deadline)
127 std::scoped_lock l(runtime->skillStatusesMutex);
128 terminated = runtime->statusUpdate.hasBeenTerminated();
132 std::this_thread::sleep_for(std::chrono::milliseconds(20));
138 runtime->execution.join();
142 ARMARX_WARNING <<
"Skill execution '" << runtime->statusUpdate.executionId.skillId
143 <<
"' did not terminate within 5s during disconnect. Detaching.";
144 runtime->execution.detach();
151 skillFactories.clear();
157 std::string
prefix =
"skill.";
158 properties->component(
162 "The name of the SkillManager (or SkillMemory) proxy this provider belongs to.");
163 properties->topic<armarx::skills::SkillEventListenerInterface>(
164 "SkillEventListener",
prefix +
"tpc.sub.SkillEventListener");
176 const std::string componentName = p.getName();
181 const std::unique_lock l(skillFactoriesMutex);
182 auto skillId = fac->createSkillDescription(providerId).skillId;
184 if (skillFactories.find(skillId) != skillFactories.end())
186 ARMARX_WARNING <<
"Try to add a skill factory for skill '" + skillId.toString() +
187 "' which already exists in list. Ignoring this skill.";
191 ARMARX_INFO <<
"Adding skill `" << skillId <<
"` to component `" << componentName <<
"` .";
193 skillFactories.emplace(skillId, std::move(fac));
218 if (skillFactories.count(skillId) == 0)
220 ARMARX_INFO <<
"Could not find a skill factory for id: " << skillId;
224 auto* facPtr = skillFactories.at(skillId).get();
228 std::optional<skills::SkillStatusUpdate>
234 const std::shared_lock l(skillExecutionsMutex);
235 auto it = skillExecutions.find(execId);
236 if (it == skillExecutions.end())
243 std::scoped_lock l2{it->second->skillStatusesMutex};
244 return it->second->statusUpdate;
247 std::map<skills::SkillExecutionID, skills::SkillStatusUpdate>
250 std::map<skills::SkillExecutionID, skills::SkillStatusUpdate> skillUpdates;
252 const std::shared_lock l(skillExecutionsMutex);
253 for (
const auto& [key, impl] : skillExecutions)
255 const std::scoped_lock l2(impl->skillStatusesMutex);
256 skillUpdates.insert({key, impl->statusUpdate});
261 std::optional<skills::SkillDescription>
266 const std::shared_lock l(skillFactoriesMutex);
267 if (skillFactories.find(skillId) == skillFactories.end())
269 std::stringstream ss;
270 ss <<
"Skill description for skill '" + skillId.
toString() +
271 "' not found! Found instead: {"
273 for (
const auto& [k, _] : skillFactories)
275 ss <<
"\t" << k.toString() <<
"\n";
283 return skillFactories.at(skillId)->createSkillDescription(*skillId.
providerId);
286 std::map<skills::SkillID, skills::SkillDescription>
289 std::map<skills::SkillID, skills::SkillDescription> skillDesciptions;
290 const std::shared_lock l(skillFactoriesMutex);
291 for (
const auto& [key, fac] : skillFactories)
294 skillDesciptions.insert({key, fac->createSkillDescription(*key.providerId)});
296 return skillDesciptions;
307 .executionStartedTime =
313 std::shared_ptr<skills::detail::SkillRuntime> runtime;
314 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> finishedToJoin;
316 auto l1 = std::unique_lock{skillFactoriesMutex};
318 const auto& fac = getSkillFactory(executionId.skillId);
319 ARMARX_CHECK(fac) <<
"Could not find a factory for skill " << executionId.skillId;
322 const std::unique_lock l2{skillExecutionsMutex};
327 finishedToJoin = collectFinishedExecutions_locked();
329 runtime = std::make_shared<skills::detail::SkillRuntime>(
334 skillExecutions.emplace(executionId, runtime);
338 runtime->setLocalMinimumLoggingLevel(loggingLevel);
345 runtime->execution = std::thread(
350 auto x = runtime->executeSkill();
351 ret.result =
x.result;
354 catch (std::exception& e)
357 "skill. Exception was: "
365 for (
auto& finished : finishedToJoin)
367 if (finished && finished->execution.joinable())
369 finished->execution.join();
372 finishedToJoin.clear();
374 if (runtime && runtime->execution.joinable())
376 runtime->execution.join();
389 std::shared_ptr<skills::detail::SkillRuntime> runtime;
390 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> finishedToJoin;
392 auto l1 = std::unique_lock{skillFactoriesMutex};
394 const auto& fac = getSkillFactory(executionRequest.
skillId);
398 const std::unique_lock l2{skillExecutionsMutex};
401 finishedToJoin = collectFinishedExecutions_locked();
407 if (skillExecutions.count(executionId) > 0)
409 ARMARX_ERROR <<
"SkillsExecutionID already exists! This is undefined behaviour "
410 "and should not occur!";
413 runtime = std::make_shared<skills::detail::SkillRuntime>(
418 skillExecutions.emplace(executionId, runtime);
420 ARMARX_INFO <<
"Setting skill runtime's logging level to `"
422 runtime->setLocalMinimumLoggingLevel(loggingLevel);
428 runtime->execution = std::thread(
433 auto x = runtime->executeSkill();
436 catch (std::exception& e)
439 "skill. Exception was: "
447 for (
auto& finished : finishedToJoin)
449 if (finished && finished->execution.joinable())
451 finished->execution.join();
454 finishedToJoin.clear();
460 std::scoped_lock l(runtime->skillStatusesMutex);
462 if (runtime->statusUpdate.hasBeenConstructed())
468 std::this_thread::sleep_for(std::chrono::milliseconds(20));
480 std::shared_ptr<skills::detail::SkillRuntime> runtime;
482 std::shared_lock l{skillExecutionsMutex};
483 auto it = skillExecutions.find(executionId);
484 if (it == skillExecutions.end())
487 "' found! Ignoring prepareSkill request.";
490 runtime = it->second;
493 std::scoped_lock l2{runtime->skillStatusesMutex};
497 "' because its not in preparing phase.";
501 runtime->updateSkillParameters(input);
510 std::shared_ptr<skills::detail::SkillRuntime> runtime;
512 std::shared_lock l(skillExecutionsMutex);
513 auto it = skillExecutions.find(executionId);
514 if (it == skillExecutions.end())
517 "' found! Ignoring abortSkill request.";
520 runtime = it->second;
523 runtime->stopSkill();
528 std::scoped_lock l2(runtime->skillStatusesMutex);
529 auto status = runtime->statusUpdate;
531 if (
status.hasBeenTerminated())
536 std::this_thread::sleep_for(std::chrono::milliseconds(20));
547 std::shared_ptr<skills::detail::SkillRuntime> runtime;
549 std::shared_lock l(skillExecutionsMutex);
550 auto it = skillExecutions.find(executionId);
551 if (it == skillExecutions.end())
554 "' found! Ignoring abortSkill request.";
557 runtime = it->second;
560 runtime->stopSkill();
571 std::shared_ptr<skills::detail::SkillRuntime>>>
574 std::shared_lock l(skillExecutionsMutex);
575 snapshot.reserve(skillExecutions.size());
576 for (
const auto& [
id, runtime] : skillExecutions)
578 snapshot.emplace_back(
id, runtime);
581 for (
auto& [
id, runtime] : snapshot)
583 ARMARX_DEBUG <<
"updating subskill status for " <<
id.toString() <<
" with "
585 runtime->updateSubSkillStatus(statusUpdate);
589 const skills::manager::dti::SkillManagerInterfacePrx&
595 std::vector<std::shared_ptr<skills::detail::SkillRuntime>>
596 SkillProviderComponentPlugin::collectFinishedExecutions_locked()
604 std::vector<std::shared_ptr<skills::detail::SkillRuntime>> ret;
605 for (
auto it = skillExecutions.begin(); it != skillExecutions.end();)
607 const auto& runtime = it->second;
608 bool terminated =
false;
610 std::scoped_lock statusLock(runtime->skillStatusesMutex);
611 terminated = runtime->statusUpdate.hasBeenTerminated();
615 ret.push_back(std::move(it->second));
616 it = skillExecutions.erase(it);
634 IceUtil::Optional<skills::provider::dto::SkillDescription>
636 const skills::provider::dto::SkillID& skillId,
637 const Ice::Current& )
640 auto o = plugin->getSkillDescription(
id);
643 return o->toProviderIce();
648 skills::provider::dto::SkillDescriptionMap
651 skills::provider::dto::SkillDescriptionMap ret;
652 for (
const auto& [k, v] : plugin->getSkillDescriptions())
654 ret.insert({k.toProviderIce(), v.toProviderIce()});
659 IceUtil::Optional<skills::provider::dto::SkillStatusUpdate>
661 const skills::provider::dto::SkillExecutionID& executionId,
662 const Ice::Current& )
666 auto o = plugin->getSkillExecutionStatus(execId);
669 return o->toProviderIce();
674 skills::provider::dto::SkillStatusUpdateMap
677 skills::provider::dto::SkillStatusUpdateMap ret;
678 for (
const auto& [k, v] : plugin->getSkillExecutionStatuses())
680 ret.insert({k.toProviderIce(), v.toProviderIce()});
686 skills::provider::dto::SkillStatusUpdate
688 const skills::provider::dto::SkillExecutionRequest& info,
689 const Ice::Current& )
693 auto up = this->plugin->executeSkill(exec);
694 return up.toProviderIce();
697 skills::provider::dto::SkillExecutionID
699 const skills::provider::dto::SkillExecutionRequest& info,
700 const Ice::Current& current )
704 auto id = this->plugin->executeSkillAsync(exec);
708 skills::provider::dto::ParameterUpdateResult
710 const skills::provider::dto::SkillExecutionID&
id,
712 const Ice::Current& current )
714 skills::provider::dto::ParameterUpdateResult res;
719 res.success = this->plugin->updateSkillParameters(exec, prep);
723 skills::provider::dto::AbortSkillResult
725 const Ice::Current& )
727 skills::provider::dto::AbortSkillResult res;
730 res.success = this->plugin->abortSkill(exec);
734 skills::provider::dto::AbortSkillResult
736 const skills::provider::dto::SkillExecutionID&
id,
737 const Ice::Current& )
739 skills::provider::dto::AbortSkillResult res;
742 res.success = this->plugin->abortSkillAsync(exec);
748 const skills::provider::dto::SkillStatusUpdate& statusUpdate,
749 const std::string& providerName,
750 const Ice::Current& )
753 status.executionId.skillId.providerId.emplace().providerName = providerName;
754 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.
std::string GetHandledExceptionString()
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)