From 5dc6c30cef14408c54fb8df14af92c0e570acec5 Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Fri, 6 Feb 2026 12:46:29 +0200 Subject: [PATCH 1/8] cm: launcher: don't return error if instance cache failed Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/instancemanager.cpp | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/src/core/cm/launcher/instancemanager.cpp b/src/core/cm/launcher/instancemanager.cpp index 4a9277f56..9f7ad3dff 100644 --- a/src/core/cm/launcher/instancemanager.cpp +++ b/src/core/cm/launcher/instancemanager.cpp @@ -203,7 +203,9 @@ Error InstanceManager::SubmitScheduledInstances() // Cache deleted instances if (!isStashed) { if (auto err = instance->Cache(); !err.IsNone()) { - return AOS_ERROR_WRAP(err); + const auto& id = instance->GetInfo().mInstanceIdent; + + LOG_ERR() << "Cache instance failed" << Log::Field("instanceID", id) << AOS_ERROR_WRAP(err); } mCachedInstances.PushBack(instance); @@ -219,7 +221,9 @@ Error InstanceManager::SubmitScheduledInstances() Error InstanceManager::DisableInstance(SharedPtr& instance) { if (auto err = instance->Cache(true); !err.IsNone()) { - return AOS_ERROR_WRAP(err); + const auto& id = instance->GetInfo().mInstanceIdent; + + LOG_ERR() << "Disable instance failed" << Log::Field("instanceID", id) << AOS_ERROR_WRAP(err); } mCachedInstances.PushBack(instance); From 04142e5516127472495393f04d14745d7d488e3c Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:36:13 +0200 Subject: [PATCH 2/8] cm: launcher: move aos::InstanceInfo into Instance class Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/instance.cpp | 103 +++++++++++++++++++----------- src/core/cm/launcher/instance.hpp | 63 ++++++++++++++++-- 2 files changed, 124 insertions(+), 42 deletions(-) diff --git a/src/core/cm/launcher/instance.cpp b/src/core/cm/launcher/instance.cpp index 4205bbb3e..1bd9b11e3 100644 --- a/src/core/cm/launcher/instance.cpp +++ b/src/core/cm/launcher/instance.cpp @@ -268,29 +268,25 @@ oci::BalancingPolicyEnum ComponentInstance::GetBalancingPolicy() return oci::BalancingPolicyEnum::eBalancingDisabled; } -Error ComponentInstance::Schedule(NodeItf& node, const String& runtimeID, aos::InstanceInfo& info) +Error ComponentInstance::Schedule(NodeItf& node, const String& runtimeID) { auto releaseConfig = DeferRelease(reinterpret_cast(1), [&](int*) { mImageConfig = nullptr; }); - static_cast(info) = mInfo.mInstanceIdent; - info.mVersion = mInfo.mVersion; - info.mManifestDigest = mInfo.mManifestDigest; - info.mRuntimeID = runtimeID; - info.mOwnerID = mInfo.mOwnerID; - info.mSubjectType = mInfo.mSubjectType; - info.mUID = mInfo.mUID; - info.mGID = mInfo.mGID; - info.mPriority = mInfo.mPriority; - - info.mStoragePath = ""; - info.mStatePath = ""; - info.mEnvVars.Clear(); - info.mNetworkParameters.Reset(); - info.mMonitoringParams.Reset(); - - if (auto err = node.ScheduleInstance(info); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } + static_cast(mSMInfo) = mInfo.mInstanceIdent; + mSMInfo.mVersion = mInfo.mVersion; + mSMInfo.mManifestDigest = mInfo.mManifestDigest; + mSMInfo.mRuntimeID = runtimeID; + mSMInfo.mOwnerID = mInfo.mOwnerID; + mSMInfo.mSubjectType = mInfo.mSubjectType; + mSMInfo.mUID = mInfo.mUID; + mSMInfo.mGID = mInfo.mGID; + mSMInfo.mPriority = mInfo.mPriority; + + mSMInfo.mStoragePath = ""; + mSMInfo.mStatePath = ""; + mSMInfo.mEnvVars.Clear(); + mSMInfo.mNetworkParameters.Reset(); + mSMInfo.mMonitoringParams.Reset(); if (auto err = SetActive(node.GetConfig().mNodeID, runtimeID); !err.IsNone()) { return AOS_ERROR_WRAP(err); @@ -299,6 +295,18 @@ Error ComponentInstance::Schedule(NodeItf& node, const String& runtimeID, aos::I return ErrorEnum::eNone; } +Error ComponentInstance::PrepareNetworkParams(bool onlyExposedPorts) +{ + (void)onlyExposedPorts; + + return ErrorEnum::eNone; +} + +Error ComponentInstance::RemoveNetworkParams() +{ + return ErrorEnum::eNone; +} + /*********************************************************************************************************************** * ServiceInstance implementation **********************************************************************************************************************/ @@ -434,7 +442,7 @@ oci::BalancingPolicyEnum ServiceInstance::GetBalancingPolicy() return mItemConfig->mBalancingPolicy; } -Error ServiceInstance::Schedule(NodeItf& node, const String& runtimeID, aos::InstanceInfo& info) +Error ServiceInstance::Schedule(NodeItf& node, const String& runtimeID) { assert(mItemConfig); @@ -443,46 +451,67 @@ Error ServiceInstance::Schedule(NodeItf& node, const String& runtimeID, aos::Ins mImageConfig.Reset(); }); - static_cast(info) = mInfo.mInstanceIdent; - info.mVersion = mInfo.mVersion; - info.mManifestDigest = mInfo.mManifestDigest; - info.mRuntimeID = runtimeID; - info.mOwnerID = mInfo.mOwnerID; - info.mSubjectType = mInfo.mSubjectType; - info.mUID = mInfo.mUID; - info.mGID = mInfo.mGID; - info.mPriority = mInfo.mPriority; - - if (auto err = SetupStateStorage(node.GetConfig(), info.mStoragePath, info.mStatePath); !err.IsNone()) { + static_cast(mSMInfo) = mInfo.mInstanceIdent; + mSMInfo.mVersion = mInfo.mVersion; + mSMInfo.mManifestDigest = mInfo.mManifestDigest; + mSMInfo.mRuntimeID = runtimeID; + mSMInfo.mOwnerID = mInfo.mOwnerID; + mSMInfo.mSubjectType = mInfo.mSubjectType; + mSMInfo.mUID = mInfo.mUID; + mSMInfo.mGID = mInfo.mGID; + mSMInfo.mPriority = mInfo.mPriority; + + if (auto err = SetupStateStorage(node.GetConfig(), mSMInfo.mStoragePath, mSMInfo.mStatePath); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - info.mEnvVars.Clear(); + mSMInfo.mEnvVars.Clear(); if (auto err = SetupNetworkServiceData(); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - info.mMonitoringParams.EmplaceValue(); + mSMInfo.mMonitoringParams.EmplaceValue(); if (mItemConfig->mAlertRules.HasValue()) { - info.mMonitoringParams.GetValue().mAlertRules = mItemConfig->mAlertRules.GetValue(); + mSMInfo.mMonitoringParams.GetValue().mAlertRules = mItemConfig->mAlertRules.GetValue(); } if (auto err = ReserveRuntimeResources(node); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - if (auto err = node.ScheduleInstance(info); !err.IsNone()) { + if (auto err = SetActive(node.GetConfig().mNodeID, runtimeID); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - if (auto err = SetActive(node.GetConfig().mNodeID, runtimeID); !err.IsNone()) { + return ErrorEnum::eNone; +} + +Error ServiceInstance::PrepareNetworkParams(bool onlyExposedPorts) +{ + auto err = mNetworkManager.PrepareInstanceNetworkParameters( + mInfo.mInstanceIdent, mInfo.mOwnerID, mInfo.mNodeID, onlyExposedPorts, mSMInfo.mNetworkParameters); + if (!err.IsNone()) { return AOS_ERROR_WRAP(err); } return ErrorEnum::eNone; } +Error ServiceInstance::RemoveNetworkParams() +{ + if (mSMInfo.mNetworkParameters.HasValue()) { + auto err = mNetworkManager.RemoveInstanceNetworkParameters(mInfo.mInstanceIdent, mInfo.mNodeID); + if (!err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + + mSMInfo.mNetworkParameters.Reset(); + } + + return ErrorEnum::eNone; +} + size_t ServiceInstance::GetRequestedCPU(const NodeConfig& nodeConfig, bool useMonitoringData) { assert(mItemConfig); diff --git a/src/core/cm/launcher/instance.hpp b/src/core/cm/launcher/instance.hpp index 84b053422..7b0b34038 100644 --- a/src/core/cm/launcher/instance.hpp +++ b/src/core/cm/launcher/instance.hpp @@ -96,6 +96,13 @@ class Instance { */ const InstanceInfo& GetInfo() const { return mInfo; } + /** + * Returns SM instance information. + * + * @return const aos::InstanceInfo&. + */ + const aos::InstanceInfo& GetSMInfo() const { return mSMInfo; } + /** * Returns instance status. * @@ -215,13 +222,29 @@ class Instance { * @param[out] info preallocate instance info used as temporary location. * @return Error. */ - virtual Error Schedule(NodeItf& node, const String& runtimeID, aos::InstanceInfo& info) = 0; + virtual Error Schedule(NodeItf& node, const String& runtimeID) = 0; + + /** + * Setups network parameters. + * + * @param onlyExposedPorts setup only for exposed ports. + * @return Error. + */ + virtual Error PrepareNetworkParams(bool onlyExposedPorts) = 0; + + /** + * Removes network parameters. + * + * @return Error. + */ + virtual Error RemoveNetworkParams() = 0; protected: Error SetActive(const String& nodeID, const String& runtimeID); - InstanceInfo mInfo; - InstanceStatus mStatus; + InstanceInfo mInfo; + aos::InstanceInfo mSMInfo; + InstanceStatus mStatus; StorageItf& mStorage; ImageInfoProvider& mImageInfoProvider; @@ -314,7 +337,22 @@ class ComponentInstance : public Instance { * @param[out] info preallocate instance info used as temporary location. * @return Error. */ - Error Schedule(NodeItf& node, const String& runtimeID, aos::InstanceInfo& info) override; + Error Schedule(NodeItf& node, const String& runtimeID) override; + + /** + * Prepares network parameters. + * + * @param onlyExposedPorts prepare only for exposed ports. + * @return Error. + */ + Error PrepareNetworkParams(bool onlyExposedPorts) override; + + /** + * Removes network parameters. + * + * @return Error. + */ + Error RemoveNetworkParams() override; }; /** @@ -401,7 +439,22 @@ class ServiceInstance : public Instance { * @param[out] info preallocate instance info used as temporary location. * @return Error. */ - Error Schedule(NodeItf& node, const String& runtimeID, aos::InstanceInfo& info) override; + Error Schedule(NodeItf& node, const String& runtimeID) override; + + /** + * Prepares network parameters. + * + * @param onlyExposedPorts prepare only for exposed ports. + * @return Error. + */ + Error PrepareNetworkParams(bool onlyExposedPorts) override; + + /** + * Removes network parameters. + * + * @return Error. + */ + Error RemoveNetworkParams() override; private: static constexpr auto cDefaultResourceRation = 50.0; From 51bbc8e19f3d4b2f97ce3d910c1efd2a03baae09 Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:38:40 +0200 Subject: [PATCH 3/8] cm: launcher: remove internal instance arrays from node Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/node.cpp | 265 +++++++++++++++---------------- src/core/cm/launcher/node.hpp | 74 ++------- src/core/cm/launcher/nodeitf.hpp | 8 - 3 files changed, 137 insertions(+), 210 deletions(-) diff --git a/src/core/cm/launcher/node.cpp b/src/core/cm/launcher/node.cpp index 6e4b2b5c4..7bb893683 100644 --- a/src/core/cm/launcher/node.cpp +++ b/src/core/cm/launcher/node.cpp @@ -10,15 +10,95 @@ namespace aos::cm::launcher { +template +class Filter { +public: + class Iterator { + public: + Iterator(typename Array::ConstIterator it, typename Array::ConstIterator end, Cmp cmp) + : mIt(it) + , mEnd(end) + , mCmp(cmp) + { + while (mIt != mEnd && !mCmp(*mIt)) { + ++mIt; + } + } + + Iterator& operator++() + { + assert(mIt != mEnd); + + ++mIt; + + while (mIt != mEnd && !mCmp(*mIt)) { + ++mIt; + } + + return *this; + } + + Iterator operator++(int) + { + assert(mIt != mEnd); + + Iterator tmp = *this; + + ++(*this); + + return tmp; + } + + bool operator==(const Iterator& other) const { return mIt == other.mIt; } + bool operator!=(const Iterator& other) const { return mIt != other.mIt; } + + const T& operator*() const { return *mIt; } + const T* operator->() const { return mIt; } + + private: + typename Array::ConstIterator mIt; + typename Array::ConstIterator mEnd; + Cmp mCmp; + }; + + Filter(const Array& array, Cmp cmp) + : mArray(&array) + , mCmp(cmp) + { + } + + Iterator begin() const { return Iterator(mArray->begin(), mArray->end(), mCmp); } + Iterator end() const { return Iterator(mArray->end(), mArray->end(), mCmp); } + +private: + const Array* mArray; + Cmp mCmp; +}; + +auto FilterByNode(const Array& array, const String& nodeID) +{ + auto cmp = [nodeID](const InstanceStatus& status) { return status.mNodeID == nodeID; }; + + return Filter(array, cmp); +} + +auto FilterByNode(const Array>& array, const String& nodeID) +{ + auto cmp = [nodeID](const SharedPtr& instance) { return instance->GetInfo().mNodeID == nodeID; }; + + return Filter, decltype(cmp)>(array, cmp); +} + /*********************************************************************************************************************** * Public **********************************************************************************************************************/ -void Node::Init( - const String& id, unitconfig::NodeConfigProviderItf& nodeConfigProvider, InstanceRunnerItf& instanceRunner) +void Node::Init(const String& id, unitconfig::NodeConfigProviderItf& nodeConfigProvider, + InstanceRunnerItf& instanceRunner, Allocator* allocator) { mNodeConfigProvider = &nodeConfigProvider; mInstanceRunner = &instanceRunner; + mAllocator = allocator; mInfo.mNodeID = id; mInfo.mState = NodeStateEnum::eUnprovisioned; @@ -32,8 +112,6 @@ void Node::PrepareForBalancing(bool rebalancing) mRuntimeAvailableRAM.Clear(); mMaxInstances.Clear(); - mScheduledInstances.Clear(); - mNeedBalancing = false; if (rebalancing && mConfig.mAlertRules.HasValue()) { @@ -112,41 +190,6 @@ bool Node::UpdateInfo(const UnitNodeInfo& info) return nodeChanged; } -Error Node::LoadInstances() -{ - mSentInstances = mScheduledInstances; - mRunningInstances.Clear(); - mScheduledInstances.Clear(); - - LOG_DBG() << "Load instances for node" << Log::Field("nodeID", mInfo.mNodeID) - << Log::Field("instances", mSentInstances.Size()); - - for (const auto& instance : mSentInstances) { - if (auto err = mRunningInstances.EmplaceBack(); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - - Convert(instance, mRunningInstances.Back()); - } - - return ErrorEnum::eNone; -} - -Error Node::UpdateRunningInstances(const Array& instances) -{ - mRunningInstances.Clear(); - - for (const auto& instance : instances) { - if (!instance.mPreinstalled) { - if (auto err = mRunningInstances.EmplaceBack(instance); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - } - } - - return ErrorEnum::eNone; -} - size_t Node::GetAvailableCPU() { return mAvailableCPU; @@ -260,61 +303,17 @@ Error Node::ReserveResources(const InstanceIdent& instanceIdent, const String& r return ErrorEnum::eNone; } -Error Node::ScheduleInstance(const aos::InstanceInfo& instance) -{ - if (auto err = mScheduledInstances.PushBack(instance); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - - return ErrorEnum::eNone; -} - -Error Node::SetupNetworkParams(const InstanceIdent& instanceID, bool onlyExposedPorts, NetworkManager& networkManager) -{ - auto instance = mScheduledInstances.FindIf( - [&instanceID](const aos::InstanceInfo& item) { return static_cast(item) == instanceID; }); - if (instance == mScheduledInstances.end()) { - return AOS_ERROR_WRAP(ErrorEnum::eNotFound); - } - - auto netErr = networkManager.PrepareInstanceNetworkParameters( - instanceID, instance->mOwnerID, mInfo.mNodeID, onlyExposedPorts, instance->mNetworkParameters); - if (!netErr.IsNone()) { - return AOS_ERROR_WRAP(netErr); - - mScheduledInstances.Erase(instance); - } - - return ErrorEnum::eNone; -} - -Error Node::RemoveNetworkParams(const InstanceIdent& instanceID, NetworkManager& networkManager) -{ - auto sentInstance = mSentInstances.FindIf( - [&instanceID](const aos::InstanceInfo& item) { return static_cast(item) == instanceID; }); - if (sentInstance == mSentInstances.end()) { - return ErrorEnum::eNotFound; - } - - if (sentInstance->mNetworkParameters.HasValue()) { - if (auto err = networkManager.RemoveInstanceNetworkParameters(instanceID, mInfo.mNodeID); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - - sentInstance->mNetworkParameters.Reset(); - } - - return ErrorEnum::eNone; -} - -Error Node::SendScheduledInstances() +Error Node::SendScheduledInstances( + const Array>& scheduledInstances, const Array& runningInstances) { - auto stopInstances = MakeUnique>(&mAllocator); - - for (const auto& instance : mSentInstances) { - auto isScheduled = mScheduledInstances.ContainsIf([&instance](const aos::InstanceInfo& item) { - return static_cast(instance) == static_cast(item) - && instance.mRuntimeID == item.mRuntimeID; + auto stopInstances = MakeUnique>(mAllocator); + auto startInstances = MakeUnique>(mAllocator); + + for (const auto& status : FilterByNode(runningInstances, mInfo.mNodeID)) { + // Check if the instance is scheduled on this node. + auto isScheduled = scheduledInstances.ContainsIf([&status, this](const SharedPtr& item) { + return static_cast(status) == item->GetInfo().mInstanceIdent + && status.mRuntimeID == item->GetInfo().mRuntimeID && item->GetInfo().mNodeID == mInfo.mNodeID; }); if (!isScheduled) { @@ -322,80 +321,74 @@ Error Node::SendScheduledInstances() return AOS_ERROR_WRAP(err); } - static_cast(stopInstances->Back()) = static_cast(instance); - stopInstances->Back().mRuntimeID = instance.mRuntimeID; + Convert(status, stopInstances->Back()); + } + } + + for (const auto& instance : FilterByNode(scheduledInstances, mInfo.mNodeID)) { + if (auto err = startInstances->PushBack(instance->GetSMInfo()); !err.IsNone()) { + return AOS_ERROR_WRAP(err); } } LOG_INF() << "Send scheduled instances" << Log::Field("nodeID", mInfo.mNodeID) << Log::Field("stopInstances", stopInstances->Size()) - << Log::Field("startInstances", mScheduledInstances.Size()); + << Log::Field("startInstances", startInstances->Size()); - if (auto err = mInstanceRunner->UpdateInstances(mInfo.mNodeID, *stopInstances, mScheduledInstances); - !err.IsNone()) { + if (auto err = mInstanceRunner->UpdateInstances(mInfo.mNodeID, *stopInstances, *startInstances); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - mSentInstances = mScheduledInstances; - mScheduledInstances.Clear(); - return ErrorEnum::eNone; } -RetWithError Node::ResendInstances() +RetWithError Node::ResendInstances( + const Array>& activeInstances, const Array& runningInstances) { - auto stopInstances = MakeUnique>(&mAllocator); + auto stopInstances = MakeUnique>(mAllocator); + auto startInstances = MakeUnique>(mAllocator); + size_t runningNodeInstances = 0; + + for (const auto& status : FilterByNode(runningInstances, mInfo.mNodeID)) { + runningNodeInstances++; - for (const auto& instance : mRunningInstances) { - auto isRunningInstanceSent = mSentInstances.ContainsIf([&instance](const aos::InstanceInfo& item) { - return static_cast(instance) == static_cast(item) - && instance.mRuntimeID == item.mRuntimeID; + auto isActive = activeInstances.ContainsIf([&status](const SharedPtr& item) { + return static_cast(status) == item->GetInfo().mInstanceIdent + && status.mRuntimeID == item->GetInfo().mRuntimeID; }); - if (!isRunningInstanceSent) { + if (!isActive) { if (auto err = stopInstances->EmplaceBack(); !err.IsNone()) { return {false, AOS_ERROR_WRAP(err)}; } - static_cast(stopInstances->Back()) = static_cast(instance); - stopInstances->Back().mRuntimeID = instance.mRuntimeID; + Convert(status, stopInstances->Back()); + } + } + + for (const auto& instance : FilterByNode(activeInstances, mInfo.mNodeID)) { + if (auto err = startInstances->PushBack(instance->GetSMInfo()); !err.IsNone()) { + return {false, AOS_ERROR_WRAP(err)}; } } // Instance list didn't change, skip update. - if (stopInstances->IsEmpty() && mSentInstances.Size() == mRunningInstances.Size()) { + if (stopInstances->IsEmpty() && startInstances->Size() == runningNodeInstances) { return {false, ErrorEnum::eNone}; } // Send request to node. LOG_INF() << "Resend instance update" << Log::Field("nodeID", mInfo.mNodeID) << Log::Field("stopInstances", stopInstances->Size()) - << Log::Field("startInstances", mSentInstances.Size()); + << Log::Field("startInstances", startInstances->Size()); - if (auto err = mInstanceRunner->UpdateInstances(mInfo.mNodeID, *stopInstances, mSentInstances); !err.IsNone()) { + if (auto err = mInstanceRunner->UpdateInstances(mInfo.mNodeID, *stopInstances, *startInstances); !err.IsNone()) { return {false, AOS_ERROR_WRAP(err)}; } - // Reset running instances to sent, in order to not duplicate requests. - mRunningInstances.Clear(); - - for (const auto& instance : mSentInstances) { - if (auto err = mRunningInstances.EmplaceBack(); !err.IsNone()) { - return {false, AOS_ERROR_WRAP(err)}; - } - static_cast(mRunningInstances.Back()) = static_cast(instance); - mRunningInstances.Back().mRuntimeID = instance.mRuntimeID; - } - return {true, ErrorEnum::eNone}; } -bool Node::IsScheduled(const InstanceIdent& id) const -{ - return mScheduledInstances.ContainsIf( - [&id](const aos::InstanceInfo& info) { return static_cast(info) == id; }); -} - /*********************************************************************************************************************** * Private **********************************************************************************************************************/ @@ -487,16 +480,10 @@ size_t* Node::GetPtrToMaxNumInstances(const String& runtimeID) return &mMaxInstances.Find(runtimeID)->mSecond; } -void Node::Convert(const aos::InstanceInfo& instance, InstanceStatus& status) +void Node::Convert(const InstanceStatus& status, aos::InstanceInfo& info) { - static_cast(status) = static_cast(instance); - status.mRuntimeID = instance.mRuntimeID; - status.mManifestDigest = instance.mManifestDigest; - status.mState = aos::InstanceStateEnum::eActivating; - status.mError = ErrorEnum::eNone; - status.mVersion = instance.mVersion; - status.mPreinstalled = false; - status.mNodeID = mInfo.mNodeID; + static_cast(info) = static_cast(status); + info.mRuntimeID = status.mRuntimeID; } } // namespace aos::cm::launcher diff --git a/src/core/cm/launcher/node.hpp b/src/core/cm/launcher/node.hpp index f6899378e..0f8f51dbe 100644 --- a/src/core/cm/launcher/node.hpp +++ b/src/core/cm/launcher/node.hpp @@ -35,10 +35,10 @@ class Node : public NodeItf { * @param info node information. * @param nodeConfigProvider node config provider. * @param instanceRunner instance runner interface. - * @return Error. + * @param allocator allocator. */ - void Init( - const String& id, unitconfig::NodeConfigProviderItf& nodeConfigProvider, InstanceRunnerItf& instanceRunner); + void Init(const String& id, unitconfig::NodeConfigProviderItf& nodeConfigProvider, + InstanceRunnerItf& instanceRunner, Allocator* allocator); /** * Prepares node for balancing. @@ -75,21 +75,6 @@ class Node : public NodeItf { */ bool UpdateInfo(const UnitNodeInfo& info); - /** - * Loads sent / running instance list from scheduled. - * - * @return Error. - */ - Error LoadInstances(); - - /** - * Updates running instances. - * - * @param statuses array of running instance statuses. - * @return Error. - */ - Error UpdateRunningInstances(const Array& statuses); - /** * Returns available CPU. * @@ -133,54 +118,23 @@ class Node : public NodeItf { Error ReserveResources(const InstanceIdent& instanceIdent, const String& runtimeID, size_t reqCPU, size_t reqRAM, const Array>& reqResources) override; - /** - * Adds instance to scheduled instances map. - * - * @param instance instance information. - * @return Error. - */ - Error ScheduleInstance(const aos::InstanceInfo& instance) override; - - /** - * Sets up network parameters. - * - * @param onlyExposedPorts flag for only exposed ports. - * @param netMgr network manager. - * @param instances all scheduled instances. - * @return Error. - */ - Error SetupNetworkParams(const InstanceIdent& instanceIdent, bool onlyExposedPorts, NetworkManager& networkManager); - - /** - * Removes network parameters. - * - * @param instanceIdent instance identifier. - * @param networkManager network manager. - * @return Error. - */ - Error RemoveNetworkParams(const InstanceIdent& instanceIdent, NetworkManager& networkManager); - /** * Sends scheduled instances to node. * + * @param scheduledInstances scheduled instances. + * @param runningInstances running instances. * @return Error. */ - Error SendScheduledInstances(); + Error SendScheduledInstances( + const Array>& scheduledInstances, const Array& runningInstances); /** * Resends instances to node. * * @return Error. */ - RetWithError ResendInstances(); - - /** - * Checks whether instance is scheduled. - * - * @param instance instance identifier. - * @return bool. - */ - bool IsScheduled(const InstanceIdent& instance) const; + RetWithError ResendInstances( + const Array>& scheduledInstances, const Array& runningInstances); /** * Checks whether max number of instances is reached. @@ -196,8 +150,6 @@ class Node : public NodeItf { void UpdateConfig(); private: - static constexpr auto cAllocatorSize = sizeof(StaticArray); - // Returns CPU usage without Aos service instances. size_t GetSystemCPUUsage(const monitoring::NodeMonitoringData& monitoringData) const; // Returns CPU usage without Aos service instances. @@ -207,7 +159,7 @@ class Node : public NodeItf { size_t* GetPtrToAvailableRAM(const String& runtimeID); size_t* GetPtrToMaxNumInstances(const String& runtimeID); - void Convert(const aos::InstanceInfo& instance, InstanceStatus& status); + void Convert(const InstanceStatus& status, aos::InstanceInfo& info); unitconfig::NodeConfigProviderItf* mNodeConfigProvider {}; InstanceRunnerItf* mInstanceRunner {}; @@ -215,10 +167,6 @@ class Node : public NodeItf { UnitNodeInfo mInfo {}; bool mNeedBalancing {}; - StaticArray mSentInstances; - StaticArray mScheduledInstances; - StaticArray mRunningInstances; - size_t mTotalCPUUsage {}; size_t mTotalRAMUsage {}; size_t mSystemCPUUsage {}; @@ -232,7 +180,7 @@ class Node : public NodeItf { StaticMap, size_t, cMaxNumNodeRuntimes> mRuntimeAvailableCPU; StaticMap, size_t, cMaxNumNodeResources> mMaxInstances; - StaticAllocator mAllocator; + Allocator* mAllocator {}; }; /** @}*/ diff --git a/src/core/cm/launcher/nodeitf.hpp b/src/core/cm/launcher/nodeitf.hpp index 1f7193fc8..18d45f732 100644 --- a/src/core/cm/launcher/nodeitf.hpp +++ b/src/core/cm/launcher/nodeitf.hpp @@ -41,14 +41,6 @@ class NodeItf { size_t reqRAM, const Array>& reqResources) = 0; - /** - * Schedules instance on node. - * - * @param instance instance information. - * @return Error. - */ - virtual Error ScheduleInstance(const aos::InstanceInfo& instance) = 0; - /** * Returns node configuration. * From bd18121c2c0fd98f543e5dd3a9da6130ae804b35 Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:40:27 +0200 Subject: [PATCH 4/8] cm: launcher: adapt nodemanager interface to instance arrays relocation Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/nodemanager.cpp | 44 ++++++++++------------------ src/core/cm/launcher/nodemanager.hpp | 44 +++++++++++++++++----------- 2 files changed, 43 insertions(+), 45 deletions(-) diff --git a/src/core/cm/launcher/nodemanager.cpp b/src/core/cm/launcher/nodemanager.cpp index bbd835252..9ca3aafdd 100644 --- a/src/core/cm/launcher/nodemanager.cpp +++ b/src/core/cm/launcher/nodemanager.cpp @@ -33,9 +33,9 @@ Error NodeManager::Start() LOG_DBG() << "Start node manager" << Log::Field("nodes", nodes->Size()); - for (const auto& nodeID : *nodes) { - auto nodeInfo = MakeUnique(&mAllocator); + auto nodeInfo = MakeUnique(&mAllocator); + for (const auto& nodeID : *nodes) { if (auto err = mNodeInfoProvider->GetNodeInfo(nodeID, *nodeInfo); !err.IsNone()) { return AOS_ERROR_WRAP(err); } @@ -49,7 +49,7 @@ Error NodeManager::Start() // Add online provisioned node mNodes.EmplaceBack(); - mNodes.Back().Init(nodeInfo->mNodeID, *mNodeConfigProvider, *mRunner); + mNodes.Back().Init(nodeInfo->mNodeID, *mNodeConfigProvider, *mRunner, &mNodeAllocator); mNodes.Back().UpdateInfo(*nodeInfo); } @@ -76,7 +76,8 @@ Error NodeManager::PrepareForBalancing(bool rebalancing) return ErrorEnum::eNone; } -Error NodeManager::LoadInstances(const Array>& instances, ImageInfoProvider& imageInfoProvider) +Error NodeManager::LoadSMDataForActiveInstances( + const Array>& instances, ImageInfoProvider& imageInfoProvider) { for (const auto& instance : instances) { const auto& nodeID = instance->GetInfo().mNodeID; @@ -113,8 +114,7 @@ Error NodeManager::LoadInstances(const Array>& instances, Im continue; } - auto instanceInfo = MakeUnique(&mAllocator); - if (auto err = instance->Schedule(*node, runtimeID, *instanceInfo); !err.IsNone()) { + if (auto err = instance->Schedule(*node, runtimeID); !err.IsNone()) { LOG_ERR() << "Can't load instance" << Log::Field("nodeID", nodeID) << Log::Field("instanceID", instanceID) << Log::Field(AOS_ERROR_WRAP(err)); @@ -122,16 +122,10 @@ Error NodeManager::LoadInstances(const Array>& instances, Im } } - for (auto& node : mNodes) { - if (auto err = node.LoadInstances(); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - } - return ErrorEnum::eNone; } -Error NodeManager::UpdateRunnigInstances(const String& nodeID, const Array& statuses) +Error NodeManager::NotifyNodeStatusReceived(const String& nodeID) { auto node = FindNode(nodeID); if (node == nullptr) { @@ -140,7 +134,7 @@ Error NodeManager::UpdateRunnigInstances(const String& nodeID, const ArrayUpdateRunningInstances(statuses); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - if (node->GetInfo().mIsConnected && node->GetInfo().mState == NodeStateEnum::eProvisioned) { if (mNodesExpectedToSendStatus.Remove(nodeID) != 0) { mStatusUpdateCondVar.NotifyAll(); @@ -195,17 +185,14 @@ Array& NodeManager::GetNodes() return mNodes; } -bool NodeManager::IsScheduled(const InstanceIdent& id) -{ - return mNodes.ContainsIf([&id](const Node& info) { return info.IsScheduled(id); }); -} - -Error NodeManager::SendScheduledInstances(UniqueLock& lock) +Error NodeManager::SendScheduledInstances(UniqueLock& lock, const Array>& scheduledInstances, + const Array& runningInstances) { Error firstErr = ErrorEnum::eNone; for (auto& node : mNodes) { - if (auto err = node.SendScheduledInstances(); !err.IsNone()) { + auto err = node.SendScheduledInstances(scheduledInstances, runningInstances); + if (!err.IsNone()) { LOG_ERR() << "Can't send instance update" << Log::Field("nodeID", node.GetInfo().mNodeID) << Log::Field(err); @@ -237,7 +224,8 @@ Error NodeManager::SendScheduledInstances(UniqueLock& lock) return ErrorEnum::eNone; } -Error NodeManager::ResendInstances(UniqueLock& lock, const Array>& updatedNodes) +Error NodeManager::ResendInstances(UniqueLock& lock, const Array>& updatedNodes, + const Array>& activeInstances, const Array& runningInstances) { Error firstErr = ErrorEnum::eNone; @@ -248,7 +236,7 @@ Error NodeManager::ResendInstances(UniqueLock& lock, const Array>& instances, ImageInfoProvider& imageInfoProvider); + Error LoadSMDataForActiveInstances( + const Array>& activeInstances, ImageInfoProvider& imageInfoProvider); /** - * Updates list of running instances for a node. + * Notifies that node status has been received. * * @param nodeID node identifier. - * @param statuses list of running instance statuses. * @return Error. */ - Error UpdateRunnigInstances(const String& nodeID, const Array& statuses); + Error NotifyNodeStatusReceived(const String& nodeID); /** * Updates node info. @@ -106,36 +114,37 @@ class NodeManager { */ Array& GetNodes(); - /** - * Checks whether instance is scheduled. - * - * @param instance instance identifier. - * @return bool. - */ - bool IsScheduled(const InstanceIdent& instance); - /** * Sends scheduled instances to nodes and waits for instance statuses from them. * * @param lock mutex lock. + * @param scheduledInstances scheduled instances. + * @param runningInstances running instances. * @return Error. */ - Error SendScheduledInstances(UniqueLock& lock); + Error SendScheduledInstances(UniqueLock& lock, const Array>& scheduledInstances, + const Array& runningInstances); /** * Resends instances to nodes and waits for instance statuses from them. * * @param lock mutex lock. * @param updatedNodes updated nodes. + * @param activeInstances active instances. + * @param runningInstances running instances. * @return Error. */ - Error ResendInstances(UniqueLock& lock, const Array>& updatedNodes); + Error ResendInstances(UniqueLock& lock, const Array>& updatedNodes, + const Array>& activeInstances, const Array& runningInstances); private: static constexpr auto cStatusUpdateTimeout = Time::cMinutes * 10; + static constexpr auto cAllocatorSize = sizeof(StaticArray, cMaxNumNodes>) + sizeof(UnitNodeInfo); + static constexpr auto cNodeAllocatorSize = sizeof(StaticArray) * 2; + Error FindImageDescriptor(const String& itemID, const String& version, const String& manifestDigest, ImageInfoProvider& imageInfoProvider, oci::IndexContentDescriptor& imageDescriptor); @@ -143,7 +152,8 @@ class NodeManager { unitconfig::NodeConfigProviderItf* mNodeConfigProvider {}; InstanceRunnerItf* mRunner {}; - StaticAllocator mAllocator; + StaticAllocator mAllocator; + StaticAllocator mNodeAllocator; StaticArray mNodes; From 9a4bb6b43dba781da8e26050a0e6a69806e6eb6e Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:42:28 +0200 Subject: [PATCH 5/8] cm: launcher: extend instancemanager interface Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/instancemanager.cpp | 137 +++++++++++++++-------- src/core/cm/launcher/instancemanager.hpp | 76 ++++++++++--- 2 files changed, 152 insertions(+), 61 deletions(-) diff --git a/src/core/cm/launcher/instancemanager.cpp b/src/core/cm/launcher/instancemanager.cpp index 9f7ad3dff..3ead20f0c 100644 --- a/src/core/cm/launcher/instancemanager.cpp +++ b/src/core/cm/launcher/instancemanager.cpp @@ -101,6 +101,7 @@ Error InstanceManager::Stop() mScheduledInstances.Clear(); mCachedInstances.Clear(); mPreinstalledComponents.Clear(); + mRunningInstances.Clear(); return ErrorEnum::eNone; } @@ -111,6 +112,8 @@ Error InstanceManager::PrepareForBalancing() return AOS_ERROR_WRAP(err); } + mScheduledInstances.Clear(); + return ErrorEnum::eNone; } @@ -134,6 +137,11 @@ Array& InstanceManager::GetPreinstalledComponents() return mPreinstalledComponents; } +Array& InstanceManager::GetRunningInstances() +{ + return mRunningInstances; +} + Error InstanceManager::UpdateStatus(const InstanceStatus& status) { if (status.mPreinstalled) { @@ -160,38 +168,21 @@ Error InstanceManager::UpdateStatus(const InstanceStatus& status) return instance->UpdateStatus(status); } -Error InstanceManager::ScheduleInstance(const InstanceIdent& id, const RunInstanceRequest& request) +RetWithError> InstanceManager::CreateInstance( + const InstanceIdent& id, const RunInstanceRequest& request) { + if (auto instance = FindReadyInstance(id); instance) { + return {instance, ErrorEnum::eNone}; + } + auto instanceInfo = MakeUnique(&mAllocator); CreateInfo(id, request, *instanceInfo); - auto instance = ScheduleReadyInstance(id); - if (!instance) { - if (auto err = mStorage->AddInstance(*instanceInfo); !err.IsNone()) { - return AOS_ERROR_WRAP(err); - } - - auto [newInstance, createErr] = CreateInstance(*instanceInfo); - if (!createErr.IsNone()) { - return createErr; - } - - instance = newInstance; - - if (auto err = mScheduledInstances.PushBack(instance); !err.IsNone()) { - return err; - } + if (auto err = mStorage->AddInstance(*instanceInfo); !err.IsNone()) { + return {nullptr, AOS_ERROR_WRAP(err)}; } - return CheckSubjectEnabled(instance); -} - -Error InstanceManager::ScheduleInstance(SharedPtr& instance) -{ - auto readyInstance = ScheduleReadyInstance(instance->GetInfo().mInstanceIdent); - assert(readyInstance); - - return CheckSubjectEnabled(readyInstance); + return CreateInstance(*instanceInfo); } Error InstanceManager::SubmitScheduledInstances() @@ -213,6 +204,11 @@ Error InstanceManager::SubmitScheduledInstances() } mActiveInstances = mScheduledInstances; + mCachedInstances.RemoveIf([this](const SharedPtr& instance) { + return mActiveInstances.ContainsIf( + [instance](const SharedPtr& item) { return instance.Get() == item.Get(); }); + }); + mScheduledInstances.Clear(); return ErrorEnum::eNone; @@ -298,6 +294,10 @@ Error InstanceManager::LoadInstancesFromStorage() } } + if (auto err = LoadInstanceStatuses(); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + return ErrorEnum::eNone; } @@ -327,6 +327,19 @@ Error InstanceManager::LoadInstanceFromStorage(const InstanceInfo& info) return ErrorEnum::eNone; } +Error InstanceManager::LoadInstanceStatuses() +{ + for (auto& instance : mActiveInstances) { + if (!instance->GetInfo().mNodeID.IsEmpty()) { + if (auto err = mRunningInstances.EmplaceBack(instance->GetStatus()); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + } + } + + return ErrorEnum::eNone; +} + Error InstanceManager::SetExpiredStatus() { for (auto& instance : mActiveInstances) { @@ -442,25 +455,71 @@ RetWithError InstanceManager::SetSubjects(const Array return {false, ErrorEnum::eNone}; } -Error InstanceManager::CheckSubjectEnabled(SharedPtr& instance) +bool InstanceManager::IsSubjectEnabled(const Instance& instance) { - if (IsSubjectEnabled(*instance)) { - return ErrorEnum::eNone; + return !instance.GetInfo().mIsUnitSubject || mSubjects.Contains(instance.GetInfo().mInstanceIdent.mSubjectID); +} + +bool InstanceManager::IsScheduled(const InstanceIdent& id) +{ + return FindScheduledInstance(id).Get() != nullptr; +} + +Error InstanceManager::UpdateRunningInstances(const String& nodeID, const Array& statuses) +{ + mRunningInstances.RemoveIf([&nodeID](const InstanceStatus& status) { return status.mNodeID == nodeID; }); + mPreinstalledComponents.RemoveIf([&nodeID](const InstanceStatus& status) { return status.mNodeID == nodeID; }); + + for (const auto& status : statuses) { + if (status.mNodeID == nodeID) { + if (status.mPreinstalled) { + if (auto err = mPreinstalledComponents.EmplaceBack(status); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + } else { + if (auto err = mRunningInstances.EmplaceBack(status); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + } + } + } + + Error firstErr = ErrorEnum::eNone; + + for (const auto& status : statuses) { + if (auto err = UpdateStatus(status); !err.IsNone() && firstErr.IsNone()) { + firstErr = err; + } + } + + return firstErr; +} + +Error InstanceManager::ScheduleInstance(SharedPtr& instance, NodeItf& node, const String& runtimeID) +{ + if (auto err = instance->Schedule(node, runtimeID); !err.IsNone()) { + return AOS_ERROR_WRAP(err); } - if (auto err = DisableInstance(instance); !err.IsNone()) { + if (auto err = mScheduledInstances.EmplaceBack(instance); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - return AOS_ERROR_WRAP(Error(ErrorEnum::eNotSupported, "subject disabled")); + return ErrorEnum::eNone; } -bool InstanceManager::IsSubjectEnabled(const Instance& instance) +Error InstanceManager::ScheduleInstance(SharedPtr& instance, const Error& error) { - return !instance.GetInfo().mIsUnitSubject || mSubjects.Contains(instance.GetInfo().mInstanceIdent.mSubjectID); + instance->SetError(error); + + if (auto err = mScheduledInstances.EmplaceBack(instance); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + + return ErrorEnum::eNone; } -SharedPtr InstanceManager::ScheduleReadyInstance(const InstanceIdent& id) +SharedPtr InstanceManager::FindReadyInstance(const InstanceIdent& id) { auto instance = FindScheduledInstance(id); if (instance) { @@ -469,21 +528,11 @@ SharedPtr InstanceManager::ScheduleReadyInstance(const InstanceIdent& instance = FindActiveInstance(id); if (instance) { - if (auto err = mScheduledInstances.PushBack(instance); !err.IsNone()) { - return nullptr; - } - return instance; } instance = FindCachedInstance(id); if (instance) { - if (auto err = mScheduledInstances.PushBack(instance); !err.IsNone()) { - return nullptr; - } - - mCachedInstances.Remove(instance); - return instance; } diff --git a/src/core/cm/launcher/instancemanager.hpp b/src/core/cm/launcher/instancemanager.hpp index 9f005e3e0..7f16ce605 100644 --- a/src/core/cm/launcher/instancemanager.hpp +++ b/src/core/cm/launcher/instancemanager.hpp @@ -97,6 +97,13 @@ class InstanceManager { */ Array& GetPreinstalledComponents(); + /** + * Returns the collection of running instances. + * + * @return Array&. + */ + Array& GetRunningInstances(); + /** * Finds an active instance by its identifier. * @@ -138,21 +145,13 @@ class InstanceManager { Error UpdateStatus(const InstanceStatus& status); /** - * Schedules instance. + * Creates new instance object or returns existing one. * * @param id instance identifier. * @param request run instance request. - * @return Error. - */ - Error ScheduleInstance(const InstanceIdent& id, const RunInstanceRequest& request); - - /** - * Schedules instance. - * - * @param instance instance. - * @return Error. + * @return SharedPtr. */ - Error ScheduleInstance(SharedPtr& instance); + RetWithError> CreateInstance(const InstanceIdent& id, const RunInstanceRequest& request); /** * Submits all stashed instances for execution. @@ -184,15 +183,59 @@ class InstanceManager { */ RetWithError SetSubjects(const Array>& subjects); + /** + * Checks if instance is subject enabled. + * @param instance instance. + * @return bool. + */ + bool IsSubjectEnabled(const Instance& instance); + + /** + * Checks if instance is scheduled. + * + * @param id instance identifier. + * @return bool. + */ + bool IsScheduled(const InstanceIdent& id); + + /** + * Updates running instances. + * + * @param statuses list of running instance statuses. + * @return Error. + */ + Error UpdateRunningInstances(const String& nodeID, const Array& statuses); + + /** + * Schedules instances on specified node and runtime. + * + * @param instance instance. + * @param node node. + * @param runtimeID runtime ID. + * @return Error. + */ + Error ScheduleInstance(SharedPtr& instance, NodeItf& node, const String& runtimeID); + + /** + * Adds instances with specified error to schedule list. + * It will not be sent for execution but will be available for future rebalancing. + * + * @param instance instance. + * @param error error. + * @return Error. + */ + Error ScheduleInstance(SharedPtr& instance, const Error& error); + private: static constexpr auto cRemovePeriod = Time::cDay; static constexpr auto cAllocatorSize = Max(sizeof(ComponentInstance), sizeof(ServiceInstance)) * cMaxNumInstances + sizeof(InstanceInfo) * cMaxNumInstances + sizeof(InstanceInfo) + sizeof(oci::ImageIndex); static constexpr auto cInstanceAllocatorSize - = Max(sizeof(oci::ImageConfig) + sizeof(oci::ItemConfig), sizeof(oci::ImageIndex)); + = sizeof(oci::ImageConfig) + sizeof(oci::ItemConfig) + sizeof(InstanceStatus) + sizeof(oci::ImageIndex); Error LoadInstancesFromStorage(); Error LoadInstanceFromStorage(const InstanceInfo& info); + Error LoadInstanceStatuses(); Error SetExpiredStatus(); Error RemoveOutdatedInstances(); @@ -200,10 +243,7 @@ class InstanceManager { RetWithError> CreateInstance(const InstanceInfo& info); - Error CheckSubjectEnabled(SharedPtr& instance); - bool IsSubjectEnabled(const Instance& instance); - - SharedPtr ScheduleReadyInstance(const InstanceIdent& id); + SharedPtr FindReadyInstance(const InstanceIdent& id); void CreateInfo(const InstanceIdent& id, const RunInstanceRequest& request, InstanceInfo& info); Config mConfig; @@ -224,7 +264,9 @@ class InstanceManager { StaticArray, cMaxNumInstances> mActiveInstances; StaticArray, cMaxNumInstances> mScheduledInstances; StaticArray, cMaxNumInstances> mCachedInstances; - StaticArray mPreinstalledComponents; + + StaticArray mPreinstalledComponents; + StaticArray mRunningInstances; SubjectArray mSubjects; }; From cfd5f744f90a4057061feb1ac094f979e2e59dfa Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:43:08 +0200 Subject: [PATCH 6/8] cm: launcher: adapt balancer to interface changes Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/balancer.cpp | 103 ++++++++++++------------------ src/core/cm/launcher/balancer.hpp | 25 ++++---- 2 files changed, 54 insertions(+), 74 deletions(-) diff --git a/src/core/cm/launcher/balancer.cpp b/src/core/cm/launcher/balancer.cpp index 4c2f25e96..7efd44fa2 100644 --- a/src/core/cm/launcher/balancer.cpp +++ b/src/core/cm/launcher/balancer.cpp @@ -26,19 +26,19 @@ void Balancer::Init(InstanceManager& instanceManager, imagemanager::ItemInfoProv mNetworkManager = &networkManager; } -Error Balancer::RunInstances(UniqueLock& lock, bool rebalancing) +Error Balancer::RunInstances(UniqueLock& lock, Array>& instances, bool rebalancing) { if (auto err = PrepareForBalancing(rebalancing); !err.IsNone()) { return AOS_ERROR_WRAP(err); } if (rebalancing) { - if (auto err = PerformPolicyBalancing(); !err.IsNone()) { + if (auto err = PerformPolicyBalancing(instances); !err.IsNone()) { return AOS_ERROR_WRAP(err); } } - if (auto err = PerformNodeBalancing(); !err.IsNone()) { + if (auto err = PerformNodeBalancing(instances); !err.IsNone()) { return AOS_ERROR_WRAP(err); } @@ -46,26 +46,33 @@ Error Balancer::RunInstances(UniqueLock& lock, bool rebalancing) return AOS_ERROR_WRAP(err); } + // Submit scheduled instances before sending instances to nodes. + // So following status updates will change active instances. if (auto err = mInstanceManager->SubmitScheduledInstances(); !err.IsNone()) { return AOS_ERROR_WRAP(err); } - if (auto err = mNodeManager->SendScheduledInstances(lock); !err.IsNone()) { + if (auto err = mNodeManager->SendScheduledInstances( + lock, mInstanceManager->GetActiveInstances(), mInstanceManager->GetRunningInstances()); + !err.IsNone()) { return AOS_ERROR_WRAP(err); } return ErrorEnum::eNone; } -Error Balancer::LoadInstances() +Error Balancer::LoadSMDataForActiveInstances() { - PrepareForBalancing(false); - - auto err = mNodeManager->LoadInstances(mInstanceManager->GetActiveInstances(), mImageInfoProvider); - if (!err.IsNone()) { + if (auto err = PrepareForBalancing(false); !err.IsNone()) { return AOS_ERROR_WRAP(err); } + auto loadErr + = mNodeManager->LoadSMDataForActiveInstances(mInstanceManager->GetActiveInstances(), mImageInfoProvider); + if (!loadErr.IsNone()) { + return AOS_ERROR_WRAP(loadErr); + } + return ErrorEnum::eNone; } @@ -73,18 +80,18 @@ Error Balancer::LoadInstances() * Private **********************************************************************************************************************/ -Error Balancer::PerformNodeBalancing() +Error Balancer::PerformNodeBalancing(Array>& instances) { LOG_DBG() << "Perform node balancing" << Log::Field("numNodes", mNodeManager->GetNodes().Size()) - << Log::Field("numInstances", mInstanceManager->GetScheduledInstances().Size()); + << Log::Field("numInstances", instances.Size()); - for (const auto& instance : mInstanceManager->GetScheduledInstances()) { + for (auto& instance : instances) { const auto& info = instance->GetInfo(); const auto& id = info.mInstanceIdent; LOG_DBG() << "Perform node balancing" << Log::Field("instance", id); - if (mNodeManager->IsScheduled(id)) { + if (mInstanceManager->IsScheduled(id)) { LOG_DBG() << "Instance aready scheduled" << Log::Field("instance", id); continue; @@ -95,6 +102,7 @@ Error Balancer::PerformNodeBalancing() if (auto err = mImageInfoProvider.GetImageIndex(id.mItemID, info.mVersion, *imageIndex); !err.IsNone()) { LOG_ERR() << "Can't get images" << Log::Field("instance", id) << Log::Field(err); + mInstanceManager->ScheduleInstance(instance, AOS_ERROR_WRAP(err)); continue; } @@ -104,7 +112,7 @@ Error Balancer::PerformNodeBalancing() LOG_DBG() << "Try to schedule instance" << Log::Field("instance", id) << Log::Field("manifest", manifest.mDigest); - if (auto err = ScheduleInstance(*instance, manifest); err.IsNone()) { + if (auto err = ScheduleInstance(instance, manifest); err.IsNone()) { LOG_DBG() << "Instance scheduled successfully" << Log::Field("nodeID", info.mNodeID); break; @@ -116,21 +124,20 @@ Error Balancer::PerformNodeBalancing() } if (!scheduleErr.IsNone()) { - instance->SetError(scheduleErr); + mInstanceManager->ScheduleInstance(instance, scheduleErr); } } return ErrorEnum::eNone; } -Error Balancer::ScheduleInstance(Instance& instance, const oci::IndexContentDescriptor& imageDescriptor) +Error Balancer::ScheduleInstance(SharedPtr& instance, const oci::IndexContentDescriptor& imageDescriptor) { - auto nodes = MakeUnique>(&mAllocator); - auto instanceInfo = MakeUnique(&mAllocator); + auto nodes = MakeUnique>(&mAllocator); - auto releaseConfigs = DeferRelease(reinterpret_cast(1), [&](int*) { instance.ResetConfigs(); }); + auto releaseConfigs = DeferRelease(reinterpret_cast(1), [&](int*) { instance->ResetConfigs(); }); - if (auto err = instance.LoadConfigs(imageDescriptor); !err.IsNone()) { + if (auto err = instance->LoadConfigs(imageDescriptor); !err.IsNone()) { return AOS_ERROR_WRAP(Error(err, "can't load instance configs")); } @@ -139,11 +146,11 @@ Error Balancer::ScheduleInstance(Instance& instance, const oci::IndexContentDesc return AOS_ERROR_WRAP(Error(err, "get connected nodes failed")); } - if (auto err = SelectNodes(instance, *nodes); !err.IsNone()) { + if (auto err = SelectNodes(*instance, *nodes); !err.IsNone()) { return AOS_ERROR_WRAP(Error(err, "can't find node for instance")); } - auto [nodeRuntime, selectErr] = SelectRuntime(instance, *nodes); + auto [nodeRuntime, selectErr] = SelectRuntime(*instance, *nodes); if (!selectErr.IsNone()) { return AOS_ERROR_WRAP(Error(selectErr, "can't find runtime for instance")); } @@ -152,7 +159,7 @@ Error Balancer::ScheduleInstance(Instance& instance, const oci::IndexContentDesc auto& node = nodeRuntime.mFirst; const auto& runtime = nodeRuntime.mSecond; - if (auto err = instance.Schedule(*node, runtime->mRuntimeID, *instanceInfo); !err.IsNone()) { + if (auto err = mInstanceManager->ScheduleInstance(instance, *node, runtime->mRuntimeID); !err.IsNone()) { return AOS_ERROR_WRAP(Error(err, "can't schedule instance")); } @@ -395,22 +402,10 @@ Error Balancer::RemoveNetworkForDeletedInstances() for (const auto& instance : mInstanceManager->GetActiveInstances()) { bool isScheduled = scheduledInstances.ContainsIf( - [&instance](const SharedPtr& stashInst) { return stashInst.Get() == instance.Get(); }); + [&instance](const SharedPtr& schedInst) { return schedInst.Get() == instance.Get(); }); if (!isScheduled) { - const auto& info = instance->GetInfo(); - - if (info.mNodeID.IsEmpty()) { - // Instance has not been sent to any node(failed to schedule) - continue; - } - - auto node = mNodeManager->FindNode(info.mNodeID); - if (!node) { - return AOS_ERROR_WRAP(ErrorEnum::eNotFound); - } - - if (auto err = node->RemoveNetworkParams(info.mInstanceIdent, *mNetworkManager); !err.IsNone()) { + if (auto err = instance->RemoveNetworkParams(); !err.IsNone()) { return AOS_ERROR_WRAP(err); } } @@ -422,21 +417,7 @@ Error Balancer::RemoveNetworkForDeletedInstances() Error Balancer::SetNetworkParams(bool onlyWithExposedPorts) { for (auto& instance : mInstanceManager->GetScheduledInstances()) { - const auto& id = instance->GetInfo().mInstanceIdent; - - if (instance->GetInfo().mNodeID.IsEmpty()) { - continue; - } - - auto node = mNodeManager->FindNode(instance->GetStatus().mNodeID); - if (!node) { - LOG_ERR() << "Can't find node for instance" << Log::Field("instance", id) - << Log::Field("nodeID", instance->GetStatus().mNodeID); - - continue; - } - - auto err = node->SetupNetworkParams(id, onlyWithExposedPorts, *mNetworkManager); + auto err = instance->PrepareNetworkParams(onlyWithExposedPorts); if (!err.IsNone()) { instance->SetError(AOS_ERROR_WRAP(Error(err, "can't setup network params"))); @@ -470,12 +451,11 @@ Error Balancer::SetupNetworkForNewInstances() return ErrorEnum::eNone; } -Error Balancer::PerformPolicyBalancing() +Error Balancer::PerformPolicyBalancing(Array>& instances) { - auto imageIndex = MakeUnique(&mAllocator); - auto instanceInfo = MakeUnique(&mAllocator); + auto imageIndex = MakeUnique(&mAllocator); - for (const auto& instance : mInstanceManager->GetScheduledInstances()) { + for (auto& instance : instances) { const auto& info = instance->GetInfo(); const auto& id = info.mInstanceIdent; const auto& version = info.mVersion; @@ -528,18 +508,15 @@ Error Balancer::PerformPolicyBalancing() auto node = mNodeManager->FindNode(info.mNodeID); if (!node) { - LOG_ERR() << "Can't find node for instance" << Log::Field("instance", id) + LOG_WRN() << "Can't find node for instance" << Log::Field("instance", id) << Log::Field("nodeID", info.mNodeID); - - instance->SetError(AOS_ERROR_WRAP(ErrorEnum::eFailed)); - continue; } - if (auto err = instance->Schedule(*node, info.mRuntimeID, *instanceInfo); !err.IsNone()) { - LOG_ERR() << "Can't schedule instance" << Log::Field("instance", id) << Log::Field(err); + if (auto err = mInstanceManager->ScheduleInstance(instance, *node, info.mRuntimeID); !err.IsNone()) { + LOG_WRN() << "Can't schedule instance" << Log::Field("instance", id) << Log::Field(AOS_ERROR_WRAP(err)); - instance->SetError(err); + continue; } } diff --git a/src/core/cm/launcher/balancer.hpp b/src/core/cm/launcher/balancer.hpp index 9d61c47c8..ae2627010 100644 --- a/src/core/cm/launcher/balancer.hpp +++ b/src/core/cm/launcher/balancer.hpp @@ -47,27 +47,30 @@ class Balancer { * @param rebalancing flag indicating rebalancing. * @return Error. */ - Error RunInstances(UniqueLock& lock, bool rebalancing); + Error RunInstances(UniqueLock& lock, Array>& instances, bool rebalancing); /** - * Loads instances from storage. + * Loads Service Manager (SM) data for active instances that were loaded from storage. * * @return Error. */ - Error LoadInstances(); + Error LoadSMDataForActiveInstances(); private: using NodeRuntimes = StaticMap, cMaxNumInstances>; - static constexpr auto cAllocatorSize = sizeof(StaticArray) - + sizeof(StaticArray) + sizeof(oci::ItemConfig) + sizeof(oci::ImageConfig) - + sizeof(oci::ImageIndex) + sizeof(aos::cm::networkmanager::NetworkServiceData) + sizeof(aos::InstanceInfo) - + sizeof(StaticMap, cMaxNumInstances>) - + sizeof(StaticArray, cMaxNumInstances>); + static constexpr size_t cScheduleInstanceSize + = sizeof(oci::ImageIndex) + sizeof(StaticArray) + sizeof(NodeRuntimes); + static constexpr size_t cPolicyBalancingSize = sizeof(oci::ImageIndex); + static constexpr size_t cNetworkSetupSize = sizeof(StaticArray, cMaxNumInstances>); + static constexpr size_t cMonitoringSize = sizeof(monitoring::NodeMonitoringData); - Error PerformNodeBalancing(); + static constexpr size_t cAllocatorSize + = Max(cScheduleInstanceSize, cPolicyBalancingSize, cNetworkSetupSize, cMonitoringSize); - Error ScheduleInstance(Instance& instance, const oci::IndexContentDescriptor& imageDescriptor); + Error PerformNodeBalancing(Array>& instances); + + Error ScheduleInstance(SharedPtr& instance, const oci::IndexContentDescriptor& imageDescriptor); // Selects nodes Error SelectNodes(Instance& instance, Array& nodes); @@ -94,7 +97,7 @@ class Balancer { Error SetNetworkParams(bool onlyWithExposedPorts); Error SetupNetworkForNewInstances(); - Error PerformPolicyBalancing(); + Error PerformPolicyBalancing(Array>& instances); Error PrepareForBalancing(bool rebalancing); void UpdateMonitoringData(); From 32486a3078f52a72ef8df2dd9de38a58d4eafed0 Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:44:11 +0200 Subject: [PATCH 7/8] cm: launcher: adapt launcher implementation Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/launcher.cpp | 117 +++++++++++++++++++----------- src/core/cm/launcher/launcher.hpp | 7 +- 2 files changed, 80 insertions(+), 44 deletions(-) diff --git a/src/core/cm/launcher/launcher.cpp b/src/core/cm/launcher/launcher.cpp index 3e3c6ee71..2186b2244 100644 --- a/src/core/cm/launcher/launcher.cpp +++ b/src/core/cm/launcher/launcher.cpp @@ -77,8 +77,8 @@ Error Launcher::Start() { LOG_DBG() << "Start Launcher"; - LockGuard updateLock {mUpdateMutex}; LockGuard balancingLock {mBalancingMutex}; + LockGuard updateLock {mUpdateMutex}; mIsRunning = true; @@ -91,10 +91,6 @@ Error Launcher::Start() return err; } - if (auto err = mBalancer.LoadInstances(); !err.IsNone()) { - return err; - } - // Subscribe to providers. if (auto err = mNodeInfoProvider->SubscribeListener(*this); !err.IsNone()) { return AOS_ERROR_WRAP(err); @@ -130,6 +126,10 @@ Error Launcher::Start() UpdateInstanceStatuses(); + if (auto err = mBalancer.LoadSMDataForActiveInstances(); !err.IsNone()) { + return AOS_ERROR_WRAP(err); + } + if (auto err = mWorkerThread.Run([this](void*) { ProcessUpdate(); }); !err.IsNone()) { return AOS_ERROR_WRAP(err); } @@ -189,8 +189,8 @@ Error Launcher::RunInstances(const Array& requests, Array& requests, ArraymProcessUpdatesCondVar.NotifyAll(); }); - ScheduleInstances(requests); + auto instances = MakeUnique, cMaxNumInstances>>(&mAllocator); + + CreateRequestedInstances(requests, *instances); - auto runErr = mBalancer.RunInstances(updateLock, false); + auto runErr = mBalancer.RunInstances(updateLock, *instances, false); FailActivatingInstances(); UpdateInstanceStatuses(); @@ -335,9 +337,11 @@ Error Launcher::Rebalance(UniqueLock& lock) self->mProcessUpdatesCondVar.NotifyAll(); }); - ScheduleInstances(); + // Get instances that need rebalancing (from active and cached) + auto instances = MakeUnique, cMaxNumInstances>>(&mAllocator); + CreateRequestedInstances(*instances); - auto runErr = mBalancer.RunInstances(lock, true); + auto runErr = mBalancer.RunInstances(lock, *instances, true); FailActivatingInstances(); @@ -353,16 +357,19 @@ Error Launcher::Rebalance(UniqueLock& lock) void Launcher::ProcessUpdate() { while (true) { - UniqueLock updateLock {mUpdateMutex}; + { + UniqueLock updateLock {mUpdateMutex}; - mProcessUpdatesCondVar.Wait(updateLock, [this]() { - return (!mUpdatedNodes.IsEmpty() || mNewSubjects.HasValue() || mAlertReceived || !mIsRunning) - && !mDisableProcessUpdates; - }); + mProcessUpdatesCondVar.Wait(updateLock, [this]() { + return (!mUpdatedNodes.IsEmpty() || mNewSubjects.HasValue() || mAlertReceived || !mIsRunning) + && !mDisableProcessUpdates; + }); - WaitAllNodesConnected(updateLock); + WaitAllNodesConnected(updateLock); + } UniqueLock balancingLock {mBalancingMutex}; + UniqueLock updateLock {mUpdateMutex}; if (!mIsRunning) { return; @@ -381,10 +388,18 @@ void Launcher::ProcessUpdate() mNewSubjects.Reset(); } + // Process received alert + if (mAlertReceived) { + mAlertReceived = false; + doRebalance = true; + } + // Resend instances. if (!mUpdatedNodes.IsEmpty()) { if (!doRebalance) { - if (err = mNodeManager.ResendInstances(updateLock, mUpdatedNodes); !err.IsNone()) { + err = mNodeManager.ResendInstances(updateLock, mUpdatedNodes, mInstanceManager.GetActiveInstances(), + mInstanceManager.GetRunningInstances()); + if (!err.IsNone()) { LOG_ERR() << "Failed to resend instances" << Log::Field(AOS_ERROR_WRAP(err)); } } else { @@ -394,12 +409,6 @@ void Launcher::ProcessUpdate() mUpdatedNodes.Clear(); } - // Process received alert - if (mAlertReceived) { - mAlertReceived = false; - doRebalance = true; - } - // Rebalance. if (doRebalance) { if (err = Rebalance(updateLock); !err.IsNone()) { @@ -420,13 +429,22 @@ void Launcher::WaitAllNodesConnected(UniqueLock& lock) mAllNodesConnectedCondVar.Wait(lock, allNodesConnected); } -void Launcher::ScheduleInstances() +void Launcher::CreateRequestedInstances(Array>& instances) { for (auto& instance : mInstanceManager.GetActiveInstances()) { auto instanceIdent = instance->GetInfo().mInstanceIdent; - if (auto err = mInstanceManager.ScheduleInstance(instance); !err.IsNone()) { - LOG_ERR() << "Can't schedule instance" << Log::Field("instance", instanceIdent) << Log::Field(err); + if (!mInstanceManager.IsSubjectEnabled(*instance)) { + LOG_WRN() << "Subject is not enabled for instance" << Log::Field("instance", instanceIdent); + + mInstanceManager.DisableInstance(instance); + + continue; + } + + if (auto err = instances.PushBack(instance); !err.IsNone()) { + LOG_ERR() << "Can't add instance to array" << Log::Field("instance", instanceIdent) + << Log::Field(AOS_ERROR_WRAP(err)); continue; } @@ -439,15 +457,21 @@ void Launcher::ScheduleInstances() continue; } - if (auto err = mInstanceManager.ScheduleInstance(instance); !err.IsNone()) { - LOG_DBG() << "Can't schedule disabled instance" << Log::Field("instance", instanceIdent) << Log::Field(err); + if (!mInstanceManager.IsSubjectEnabled(*instance)) { + continue; + } + + if (auto err = instances.PushBack(instance); !err.IsNone()) { + LOG_ERR() << "Can't add instance to array" << Log::Field("instance", instanceIdent) + << Log::Field(AOS_ERROR_WRAP(err)); continue; } } } -void Launcher::ScheduleInstances(const Array& requests) +void Launcher::CreateRequestedInstances( + const Array& requests, Array>& instances) { // Sort input requests by priority auto sortedRequests = MakeUnique>(&mAllocator); @@ -457,15 +481,30 @@ void Launcher::ScheduleInstances(const Array& requests) return left.mPriority > right.mPriority || (left.mPriority == right.mPriority && left.mItemID < right.mItemID); }); - // Schedule instances - for (const auto& request : requests) { + // Create instances + for (const auto& request : *sortedRequests) { for (size_t i = 0; i < request.mNumInstances; i++) { InstanceIdent instanceIdent {request.mItemID, request.mSubjectInfo.mSubjectID, i, request.mUpdateItemType}; - auto err = mInstanceManager.ScheduleInstance(instanceIdent, request); - if (!err.IsNone()) { - LOG_ERR() << "Can't schedule instance" << Log::Field("instance", instanceIdent) << Log::Field(err); + auto [instance, createErr] = mInstanceManager.CreateInstance(instanceIdent, request); + if (!createErr.IsNone()) { + LOG_ERR() << "Can't create instance" << Log::Field("instance", instanceIdent) << Log::Field(createErr); + + continue; + } + + if (!mInstanceManager.IsSubjectEnabled(*instance)) { + LOG_WRN() << "Subject is not enabled for instance" << Log::Field("instance", instanceIdent); + + mInstanceManager.DisableInstance(instance); + + continue; + } + + if (auto err = instances.PushBack(instance); !err.IsNone()) { + LOG_ERR() << "Can't add instance to array" << Log::Field("instance", instanceIdent) + << Log::Field(AOS_ERROR_WRAP(err)); continue; } @@ -500,15 +539,11 @@ Error Launcher::OnNodeInstancesStatusesReceived(const String& nodeID, const Arra Error firstErr; - // Update instance manager. - for (const auto& status : statuses) { - if (auto err = mInstanceManager.UpdateStatus(status); !err.IsNone() && firstErr.IsNone()) { - firstErr = err; - } + if (auto err = mInstanceManager.UpdateRunningInstances(nodeID, statuses); !err.IsNone() && firstErr.IsNone()) { + firstErr = err; } - // Update node manager. - if (auto err = mNodeManager.UpdateRunnigInstances(nodeID, statuses); !err.IsNone() && firstErr.IsNone()) { + if (auto err = mNodeManager.NotifyNodeStatusReceived(nodeID); !err.IsNone() && firstErr.IsNone()) { firstErr = err; } diff --git a/src/core/cm/launcher/launcher.hpp b/src/core/cm/launcher/launcher.hpp index 535833382..6723e8885 100644 --- a/src/core/cm/launcher/launcher.hpp +++ b/src/core/cm/launcher/launcher.hpp @@ -134,7 +134,8 @@ class Launcher : public LauncherItf, private: static constexpr auto cMaxNumInstanceStatusListeners = 8; - static constexpr auto cAllocatorSize = sizeof(StaticArray); + static constexpr auto cAllocatorSize = sizeof(StaticArray) + + sizeof(StaticArray, cMaxNumInstances>); void SendRunStatus(); @@ -146,8 +147,8 @@ class Launcher : public LauncherItf, void ProcessUpdate(); void WaitAllNodesConnected(UniqueLock& lock); - void ScheduleInstances(); - void ScheduleInstances(const Array& requests); + void CreateRequestedInstances(Array>& instances); + void CreateRequestedInstances(const Array& requests, Array>& instances); // InstanceStatusReceiverItf implementation Error OnInstanceStatusReceived(const InstanceStatus& status) override; From 628ce86866caf18f35590ba83ecdf0701cef9b17 Mon Sep 17 00:00:00 2001 From: Mykola Kobets Date: Tue, 10 Feb 2026 03:44:37 +0200 Subject: [PATCH 8/8] cm: launcher: update launcher unit tests Signed-off-by: Mykola Kobets Reviewed-by: Oleksandr Grytsov Reviewed-by: Mykola Solianko --- src/core/cm/launcher/tests/launcher.cpp | 150 ++++++++++-------- .../tests/stubs/instancerunnerstub.hpp | 2 + 2 files changed, 84 insertions(+), 68 deletions(-) diff --git a/src/core/cm/launcher/tests/launcher.cpp b/src/core/cm/launcher/tests/launcher.cpp index 38cc8ef9b..ddf059078 100644 --- a/src/core/cm/launcher/tests/launcher.cpp +++ b/src/core/cm/launcher/tests/launcher.cpp @@ -272,14 +272,9 @@ aos::InstanceInfo CreateServiceRunInfo(const InstanceIdent& id, const std::strin aos::InstanceInfo result; static_cast(result) = id; - StaticString itemID; - itemID = id.mItemID; - - StaticString image; - image = imageID.c_str(); result.mVersion = version.c_str(); - result.mManifestDigest = imagemanager::ImageStoreStub::BuildManifestDigest(itemID, image); + result.mManifestDigest = BuildManifestDigest(id.mItemID.CStr(), imageID); result.mRuntimeID = runtimeID.c_str(); result.mOwnerID = ownerID.c_str(); result.mUID = uid; @@ -310,14 +305,9 @@ aos::InstanceInfo CreateComponentRunInfo(const InstanceIdent& id, const std::str aos::InstanceInfo result; static_cast(result) = id; - StaticString itemID; - itemID = id.mItemID; - - StaticString image; - image = imageID.c_str(); result.mVersion = version.c_str(); - result.mManifestDigest = imagemanager::ImageStoreStub::BuildManifestDigest(itemID, image); + result.mManifestDigest = BuildManifestDigest(id.mItemID.CStr(), imageID); result.mRuntimeID = runtimeID.c_str(); result.mPriority = priority; result.mSubjectType = subjectType; @@ -805,9 +795,11 @@ TEST_F(CMLauncherTest, Components) // Check run status auto expectedRunStatus = std::make_unique>(); - expectedRunStatus->PushBack( - CreateInstanceStatus(CreateInstanceIdent(cComponent1, cSubject1, 0, UpdateItemTypeEnum::eComponent), - cNodeIDRemoteSM1, cRunnerRootfs, aos::InstanceStateEnum::eActivating, ErrorEnum::eNone)); + auto manifestDigest = BuildManifestDigest(cComponent1, cRootfsImageID); + + expectedRunStatus->PushBack(CreateInstanceStatus( + CreateInstanceIdent(cComponent1, cSubject1, 0, UpdateItemTypeEnum::eComponent), cNodeIDRemoteSM1, cRunnerRootfs, + aos::InstanceStateEnum::eActivating, ErrorEnum::eNone, "", false, manifestDigest.CStr())); Array componentStatuses(expectedRunStatus->begin(), expectedRunStatus->Size()); @@ -882,18 +874,22 @@ TestDataPtr TestItemNodePriority() testData->mExpectedRunRequests[cNodeIDRunxSM].mStartInstances = runxSMRequests; // Expected run status + auto digest1 = BuildManifestDigest(cService1, cImageID1); + auto digest2 = BuildManifestDigest(cService2, cImageID1); + auto digest3 = BuildManifestDigest(cService3, cImageID1); + testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 1), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 1), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), - cNodeIDRunxSM, cRunnerRunx, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRunxSM, cRunnerRunx, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 1), - cNodeIDRunxSM, cRunnerRunx, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRunxSM, cRunnerRunx, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); return testData; } @@ -943,14 +939,17 @@ TestDataPtr TestItemLabels() testData->mExpectedRunRequests[cNodeIDRunxSM].mStartInstances = std::vector(); // Expected run status + auto digest1 = BuildManifestDigest(cService1, cImageID1); + auto digest2 = BuildManifestDigest(cService2, cImageID1); + testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 1), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 1), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), "", "", aos::InstanceStateEnum::eFailed, Error(ErrorEnum::eNotFound, "no nodes with instance labels"))); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 1), "", "", @@ -1007,18 +1006,22 @@ TestDataPtr TestItemResources() testData->mExpectedRunRequests[cNodeIDRunxSM].mStartInstances = std::vector(); // Expected run status + auto digest1 = BuildManifestDigest(cService1, cImageID1); + auto digest2 = BuildManifestDigest(cService2, cImageID1); + auto digest3 = BuildManifestDigest(cService3, cImageID1); + testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 1), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 1), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 1), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); return testData; } @@ -1044,13 +1047,14 @@ TestDataPtr TestItemStorageRatio() // Instances 3 and 4 fail before being sent due to storage quota limits static const char* cIpSuffixes[] = {"2", "3", "4"}; auto& localSMRequests = testData->mExpectedRunRequests[cNodeIDLocalSM].mStartInstances; + auto digest1 = BuildManifestDigest(cService1, cImageID1); for (size_t i = 0; i < 3; ++i) { localSMRequests.push_back(CreateServiceRunInfo( CreateInstanceIdent(cService1, cSubject1, i), cImageID1, cRunnerRunc, 5000 + i, 5000, cIpSuffixes[i], 100)); testData->mExpectedRunStatus.PushBack( CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, static_cast(i)), cNodeIDLocalSM, - cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); } // Failed instances (3, 4) - not sent to nodes, fail immediately @@ -1089,13 +1093,14 @@ TestDataPtr TestItemStateRatio() // Instances 3 and 4 fail before being sent due to state quota limits static const char* cStateIpSuffixes[] = {"2", "3", "4"}; auto& stateLocalRequests = testData->mExpectedRunRequests[cNodeIDLocalSM].mStartInstances; + auto digest1 = BuildManifestDigest(cService1, cImageID1); for (size_t i = 0; i < 3; ++i) { stateLocalRequests.push_back(CreateServiceRunInfo(CreateInstanceIdent(cService1, cSubject1, i), cImageID1, cRunnerRunc, 5000 + i, 5000, cStateIpSuffixes[i], 100)); testData->mExpectedRunStatus.PushBack( CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, static_cast(i)), cNodeIDLocalSM, - cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); } // Failed instances (3, 4) - not sent to nodes, fail immediately @@ -1133,6 +1138,7 @@ TestDataPtr TestItemCpuRatio() // Expected run requests - all 5 instances are scheduled and distributed across nodes static const char* cCpuIpSuffixes[] = {"2", "3", "4", "5", "6"}; auto& cpuLocalRequests = testData->mExpectedRunRequests[cNodeIDLocalSM].mStartInstances; + auto digest1 = BuildManifestDigest(cService1, cImageID1); // Instances 0, 1, 2 on localSM for (size_t i = 0; i < 3; ++i) { @@ -1140,7 +1146,7 @@ TestDataPtr TestItemCpuRatio() cRunnerRunc, 5000 + i, 5000, cCpuIpSuffixes[i], 100)); testData->mExpectedRunStatus.PushBack( CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, static_cast(i)), cNodeIDLocalSM, - cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); } // Instance 3 on remoteSM1 @@ -1148,14 +1154,14 @@ TestDataPtr TestItemCpuRatio() cpuRemoteSM1Requests.push_back(CreateServiceRunInfo( CreateInstanceIdent(cService1, cSubject1, 3), cImageID1, cRunnerRunc, 5003, 5000, cCpuIpSuffixes[3], 100)); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 3), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); // Instance 4 on remoteSM2 auto& cpuRemoteSM2Requests = testData->mExpectedRunRequests[cNodeIDRemoteSM2].mStartInstances; cpuRemoteSM2Requests.push_back(CreateServiceRunInfo( CreateInstanceIdent(cService1, cSubject1, 4), cImageID1, cRunnerRunc, 5004, 5000, cCpuIpSuffixes[4], 100)); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 4), - cNodeIDRemoteSM2, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM2, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); // Initialize empty requests for runxSM testData->mExpectedRunRequests[cNodeIDRunxSM].mStartInstances = std::vector(); @@ -1183,6 +1189,7 @@ TestDataPtr TestItemRamRatio() // Expected run requests - instances distributed across nodes static const char* cRamIpSuffixes[] = {"2", "3", "4", "5", "6"}; auto& ramLocalRequests = testData->mExpectedRunRequests[cNodeIDLocalSM].mStartInstances; + auto digest1 = BuildManifestDigest(cService1, cImageID1); // Instances 0, 1, 2 on localSM for (size_t i = 0; i < 3; ++i) { @@ -1190,7 +1197,7 @@ TestDataPtr TestItemRamRatio() cRunnerRunc, 5000 + i, 5000, cRamIpSuffixes[i], 100)); testData->mExpectedRunStatus.PushBack( CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, static_cast(i)), cNodeIDLocalSM, - cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); } // Instance 3 on remoteSM1 @@ -1198,14 +1205,14 @@ TestDataPtr TestItemRamRatio() ramRemoteSM1Requests.push_back(CreateServiceRunInfo( CreateInstanceIdent(cService1, cSubject1, 3), cImageID1, cRunnerRunc, 5003, 5000, cRamIpSuffixes[3], 100)); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 3), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); // Instance 4 on remoteSM2 auto& ramRemoteSM2Requests = testData->mExpectedRunRequests[cNodeIDRemoteSM2].mStartInstances; ramRemoteSM2Requests.push_back(CreateServiceRunInfo( CreateInstanceIdent(cService1, cSubject1, 4), cImageID1, cRunnerRunc, 5004, 5000, cRamIpSuffixes[4], 100)); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 4), - cNodeIDRemoteSM2, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM2, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); // runxSM is present in the test environment; expect an empty request for it testData->mExpectedRunRequests[cNodeIDRunxSM] = {}; @@ -1267,12 +1274,16 @@ TestDataPtr TestItemRebalancing() testData->mExpectedRunRequests[cNodeIDRunxSM] = {}; // Expected run status + auto digest1 = BuildManifestDigest(cService1, cImageID1); + auto digest2 = BuildManifestDigest(cService2, cImageID1); + auto digest3 = BuildManifestDigest(cService3, cImageID1); + testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); // Monitoring data CreateNodeMonitoring(testData->mMonitoring[cNodeIDLocalSM], cNodeIDLocalSM, 1000, @@ -1315,29 +1326,33 @@ TestDataPtr TestItemRebalancingPolicy() // localSM: starts service1, stops service2 (which was initially scheduled there) auto& policyLocalRequests = testData->mExpectedRunRequests[cNodeIDLocalSM]; policyLocalRequests.mStartInstances.push_back(CreateServiceRunInfo( - CreateInstanceIdent(cService1, cSubject1, 0), cImageID1, cRunnerRunc, 5000, 5000, "5", 100)); + CreateInstanceIdent(cService1, cSubject1, 0), cImageID1, cRunnerRunc, 5000, 5000, "6", 100)); policyLocalRequests.mStopInstances.push_back( CreateAosStopInstanceInfo(CreateInstanceIdent(cService2, cSubject1, 0), cRunnerRunc)); // remoteSM1: starts service3 (stays there, no stops since service3 has BalancingDisabled) auto& policyRemoteSM1Requests = testData->mExpectedRunRequests[cNodeIDRemoteSM1]; policyRemoteSM1Requests.mStartInstances.push_back(CreateServiceRunInfo( - CreateInstanceIdent(cService3, cSubject1, 0), cImageID1, cRunnerRunc, 5002, 5002, "7", 50)); + CreateInstanceIdent(cService3, cSubject1, 0), cImageID1, cRunnerRunc, 5002, 5002, "5", 50)); // remoteSM2: starts service2, no stops auto& policyRemoteSM2Requests = testData->mExpectedRunRequests[cNodeIDRemoteSM2]; policyRemoteSM2Requests.mStartInstances.push_back(CreateServiceRunInfo( - CreateInstanceIdent(cService2, cSubject1, 0), cImageID1, cRunnerRunc, 5001, 5001, "6", 50)); + CreateInstanceIdent(cService2, cSubject1, 0), cImageID1, cRunnerRunc, 5001, 5001, "7", 50)); testData->mExpectedRunRequests[cNodeIDRunxSM] = {}; // Expected run status (sorted by priority desc, then itemID asc: service1(100), service2(50), service3(50)) // Initial state after RunInstances() (before rebalancing): service1 and service2 on localSM, service3 on remoteSM1 + auto digest1 = BuildManifestDigest(cService1, cImageID1); + auto digest2 = BuildManifestDigest(cService2, cImageID1); + auto digest3 = BuildManifestDigest(cService3, cImageID1); + testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest1.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest2.CStr())); testData->mExpectedRunStatus.PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), - cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone)); + cNodeIDRemoteSM1, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, digest3.CStr())); // Monitoring data CreateNodeMonitoring(testData->mMonitoring[cNodeIDLocalSM], cNodeIDLocalSM, 1000, @@ -1455,7 +1470,7 @@ TEST_F(CMLauncherTest, Balancing) // Wait for rebalancing to complete (expect at least 2 more notifications: // one from rebalancing and one from status updates after rebalancing) using namespace std::chrono_literals; - ASSERT_TRUE(instanceStatusListener.WaitForNotifyCount(currentNotifyCount + 2, 2000ms)); + ASSERT_TRUE(instanceStatusListener.WaitForNotifyCount(currentNotifyCount + 3, 2000ms)); } ASSERT_TRUE(mLauncher.Stop().IsNone()); @@ -1480,7 +1495,7 @@ TEST_F(CMLauncherTest, PlatformFiltering) mNodeInfoProvider.Init(); mImageStore.Init(); mNetworkManager.Init(); - mInstanceRunner.Init(mLauncher); + mInstanceRunner.Init(mLauncher, true, aos::InstanceStateEnum::eActive); mInstanceStatusProvider.Init(); mMonitoringProvider.Init(); mResourceManager.Init(); @@ -1560,12 +1575,14 @@ TEST_F(CMLauncherTest, PlatformFiltering) // Check run status - service1 and service2 should fail, service3 should succeed auto expectedRunStatus = std::make_unique>(); - expectedRunStatus->PushBack(CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), "", "", - aos::InstanceStateEnum::eFailed, Error(ErrorEnum::eNotFound))); - expectedRunStatus->PushBack(CreateInstanceStatus(CreateInstanceIdent(cService2, cSubject1, 0), "", "", - aos::InstanceStateEnum::eFailed, Error(ErrorEnum::eNotFound))); - expectedRunStatus->PushBack( - CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), cNodeIDRemoteSM2, cRunnerRunc)); + auto manifestDigest = BuildManifestDigest(cService3, cImageID1); + + expectedRunStatus->PushBack(CreateInstanceStatus( + CreateInstanceIdent(cService1, cSubject1, 0), "", "", aos::InstanceStateEnum::eFailed, ErrorEnum::eNotFound)); + expectedRunStatus->PushBack(CreateInstanceStatus( + CreateInstanceIdent(cService2, cSubject1, 0), "", "", aos::InstanceStateEnum::eFailed, ErrorEnum::eNotFound)); + expectedRunStatus->PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), cNodeIDRemoteSM2, + cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, manifestDigest.CStr())); ASSERT_EQ(instanceStatusListener.GetLastStatuses(), *expectedRunStatus); } @@ -1602,7 +1619,7 @@ TEST_F(CMLauncherTest, ResendInstancesOnMismatchedNodeStatus) CreateItemConfig(*itemConfig, {cRunnerRunc}); AddItem(cService1, cImageID1, *itemConfig, CreateImageConfig()); - mInstanceRunner.Init(mLauncher, false); + mInstanceRunner.Init(mLauncher, false, aos::InstanceStateEnum::eActive); // First request: send empty statuses (auto-update disabled => empty statuses). // After first request is prepared, enable auto-update so the next request sends correct statuses from @@ -1640,9 +1657,10 @@ TEST_F(CMLauncherTest, ResendInstancesOnMismatchedNodeStatus) ASSERT_TRUE(mLauncher.Stop().IsNone()); // Verify latest instance statuses are correct (after resend). + auto manifestDigest = BuildManifestDigest(cService1, cImageID1); std::vector expectedStatuses = { CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), cNodeIDLocalSM, cRunnerRunc, - aos::InstanceStateEnum::eActivating, ErrorEnum::eNone), + aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, manifestDigest.CStr()), }; ASSERT_EQ(instanceStatusListener.GetLastStatuses(), @@ -1702,20 +1720,14 @@ TEST_F(CMLauncherTest, SubjectChanged) auto runStatuses = std::make_unique>(); ASSERT_TRUE(mLauncher.RunInstances(*runRequest, *runStatuses).IsNone()); - // Wait until we have at least some statuses recorded. - // Expect 2 status notifications: - // - 1st: from Launcher::RunInstances() completion notification - // - 2nd: from OnNodeInstancesStatusesReceived() after instance runner sends status + // Wait until we receive notification after run instances. ASSERT_TRUE(instanceStatusListener.WaitForNotifyCount(2, 2000ms)); // 2) Change subjects (remove all of them). ASSERT_TRUE(mIdentProvider.SetSubjects({}).IsNone()); - // 3) Check that instance fails with BadSubject error. - // Expect one more notification caused by rebalance after subjects update. + // Wait until we receive notification caused by rebalance after subjects update. ASSERT_TRUE(instanceStatusListener.WaitForNotifyCount(3, 2000ms)); - - // Verify latest instance statuses are failed with BadSubject error. ASSERT_EQ(instanceStatusListener.GetLastStatuses(), Array()); ASSERT_TRUE(mLauncher.Stop().IsNone()); @@ -1929,8 +1941,10 @@ TEST_F(CMLauncherTest, PreinstalledComponents) auto statuses = std::make_unique>(); ASSERT_TRUE(mLauncher.GetInstancesStatuses(*statuses).IsNone()); - InstanceStatus expectedRegularStatus = CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), - cNodeIDLocalSM, cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false); + auto manifestDigest = BuildManifestDigest(cService1, cImageID1); + InstanceStatus expectedRegularStatus + = CreateInstanceStatus(CreateInstanceIdent(cService1, cSubject1, 0), cNodeIDLocalSM, cRunnerRunc, + aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, manifestDigest.CStr()); std::vector expectedStatuses = {expectedRegularStatus, preinstalledStatus}; diff --git a/src/core/cm/launcher/tests/stubs/instancerunnerstub.hpp b/src/core/cm/launcher/tests/stubs/instancerunnerstub.hpp index 992b1af66..ac98b09d6 100644 --- a/src/core/cm/launcher/tests/stubs/instancerunnerstub.hpp +++ b/src/core/cm/launcher/tests/stubs/instancerunnerstub.hpp @@ -105,6 +105,8 @@ class InstanceRunnerStub : public InstanceRunnerItf { static_cast(status) = static_cast(inst); status.mNodeID = nodeID; status.mRuntimeID = inst.mRuntimeID; + status.mManifestDigest = inst.mManifestDigest; + status.mVersion = inst.mVersion; status.mState = mInitialState; status.mError = ErrorEnum::eNone;