Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions src/core/cm/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,13 @@ set(INSTALL_HEADERS
aos_core_cm_iamclient
aos_core_cm_imagemanager
aos_core_cm_nodeinfoprovider
aos_core_cm_statushandler
aos_core_cm_smcontroller
aos_core_cm_storagestate
aos_core_cm_unitconfig
)

set(INSTALL_LIBRARIES aos_core_cm_alerts aos_core_cm_monitoring aos_core_cm_storagestate aos_core_cm_updatemanager)
set(INSTALL_LIBRARIES aos_core_cm_alerts aos_core_cm_monitoring aos_core_cm_storagestate)

if(WITH_TEST)
list(APPEND INSTALL_HEADERS aos_core_cm_tests_mocks)
Expand Down Expand Up @@ -60,10 +61,11 @@ add_subdirectory(imagemanager)
add_subdirectory(launcher)
add_subdirectory(monitoring)
add_subdirectory(nodeinfoprovider)
add_subdirectory(statushandler)
add_subdirectory(smcontroller)
add_subdirectory(storagestate)
add_subdirectory(unitconfig)
add_subdirectory(updatemanager)
# add_subdirectory(updatemanager)

if(WITH_TEST)
add_subdirectory(tests)
Expand Down
91 changes: 83 additions & 8 deletions src/core/cm/launcher/balancer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,16 +15,17 @@ namespace aos::cm::launcher {
**********************************************************************************************************************/

void Balancer::Init(InstanceManager& instanceManager, ImageInfoProvider& imageInfoProvider, NodeManager& nodeManager,
MonitoringProviderItf& monitorProvider, InstanceRunnerItf& runner)
MonitoringProviderItf& monitorProvider, InstanceRunnerItf& runner, statushandler::HandlerItf& statusHandler)
{
mInstanceManager = &instanceManager;
mImageInfoProvider = &imageInfoProvider;
mNodeManager = &nodeManager;
mMonitorProvider = &monitorProvider;
mRunner = &runner;
mStatusHandler = &statusHandler;
}

Error Balancer::RunInstances(UniqueLock<Mutex>& lock, Array<SharedPtr<Instance>>& instances, bool rebalancing)
Error Balancer::RunInstances(Array<SharedPtr<Instance>>& instances, bool rebalancing)
{
if (auto err = PrepareForBalancing(rebalancing); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
Expand All @@ -41,18 +42,26 @@ Error Balancer::RunInstances(UniqueLock<Mutex>& lock, Array<SharedPtr<Instance>>
}

// Submit scheduled instances before sending them to nodes.
// So status updates expecting by SendScheduledInstances will be assigned to active instances.
if (auto err = mInstanceManager->SubmitScheduledInstances(); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

if (auto err = mNodeManager->SendScheduledInstances(
lock, mInstanceManager->GetActiveInstances(), mInstanceManager->GetRunningInstances());
!err.IsNone()) {
return AOS_ERROR_WRAP(err);
if (mStatusHandler == nullptr) {
return AOS_ERROR_WRAP(Error(ErrorEnum::eFailed, "status handler is not initialized"));
}

return ErrorEnum::eNone;
const auto& activeInstances = mInstanceManager->GetActiveInstances();

Error sendErr;
if (auto err = mNodeManager->SendScheduledInstances(activeInstances); !err.IsNone()) {
sendErr = AOS_ERROR_WRAP(err);
}

if (auto err = SendInstanceStatuses(); !err.IsNone() && sendErr.IsNone()) {
sendErr = AOS_ERROR_WRAP(err);
}

return sendErr;
}

Error Balancer::LoadSMDataForActiveInstances()
Expand All @@ -70,6 +79,43 @@ Error Balancer::LoadSMDataForActiveInstances()
return ErrorEnum::eNone;
}

Error Balancer::SendInstanceStatuses()
{
const auto& instances = mInstanceManager->GetActiveInstances();

Error firstErr;
if (auto err = SendFailedInstanceStatuses(); !err.IsNone()) {
firstErr = AOS_ERROR_WRAP(err);
}

auto statuses = MakeUnique<StaticArray<InstanceStatus, cMaxNumInstances>>(&mAllocator);

for (const auto& node : mNodeManager->GetNodes()) {
const auto& nodeID = node.GetInfo().mNodeID;

statuses->Clear();

for (const auto& instance : instances) {
const auto& status = instance->GetStatus();

if (status.mNodeID != nodeID) {
continue;
}

if (auto err = statuses->PushBack(status); !err.IsNone() && firstErr.IsNone()) {
firstErr = AOS_ERROR_WRAP(err);
}
}

if (auto err = mStatusHandler->SetNodeInstancesStatuses(nodeID, *statuses);
!err.IsNone() && firstErr.IsNone()) {
firstErr = AOS_ERROR_WRAP(err);
}
}

return firstErr;
}

/***********************************************************************************************************************
* Private
**********************************************************************************************************************/
Expand Down Expand Up @@ -491,4 +537,33 @@ Error Balancer::PrepareForBalancing(bool rebalancing, bool isInitialUpdate)
return ErrorEnum::eNone;
}

Error Balancer::SendFailedInstanceStatuses()
{
Error firstErr = ErrorEnum::eNone;

for (const auto& instance : mInstanceManager->GetActiveInstances()) {
const auto& status = instance->GetStatus();

if (status.mState != aos::InstanceStateEnum::eFailed) {
continue;
}

if (auto err = mStatusHandler->SetInstanceStatus(status); !err.IsNone() && firstErr.IsNone()) {
firstErr = AOS_ERROR_WRAP(err);
}
}

for (const auto& status : mInstanceManager->GetPreinstalledComponents()) {
if (status.mState != aos::InstanceStateEnum::eFailed) {
continue;
}

if (auto err = mStatusHandler->SetInstanceStatus(status); !err.IsNone() && firstErr.IsNone()) {
firstErr = AOS_ERROR_WRAP(err);
}
}

return firstErr;
}

} // namespace aos::cm::launcher
31 changes: 21 additions & 10 deletions src/core/cm/launcher/balancer.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
#include "itf/instancerunner.hpp"
#include "itf/launcher.hpp"
#include "itf/monitoringprovider.hpp"
#include <core/common/statushandler/itf/handler.hpp>

#include "imageinfoprovider.hpp"
#include "instancemanager.hpp"
Expand All @@ -34,18 +35,19 @@ class Balancer {
* @param nodeManager node manager.
* @param monitorProvider monitoring provider.
* @param runner instance runner interface.
* @param statusHandler status handler interface.
*/
void Init(InstanceManager& instanceManager, ImageInfoProvider& imageInfoProvider, NodeManager& nodeManager,
MonitoringProviderItf& monitorProvider, InstanceRunnerItf& runner);
MonitoringProviderItf& monitorProvider, InstanceRunnerItf& runner, statushandler::HandlerItf& statusHandler);

/**
* Runs instances.
*
* @param lock lock on the balancing mutex.
* @param instances instances to run.
* @param rebalancing flag indicating rebalancing.
* @return Error.
*/
Error RunInstances(UniqueLock<Mutex>& lock, Array<SharedPtr<Instance>>& instances, bool rebalancing);
Error RunInstances(Array<SharedPtr<Instance>>& instances, bool rebalancing);

/**
* Loads Service Manager (SM) data for active instances that were loaded from storage.
Expand All @@ -54,6 +56,13 @@ class Balancer {
*/
Error LoadSMDataForActiveInstances();

/**
* Sends instance statuses to nodes.
*
* @return Error.
*/
Error SendInstanceStatuses();

private:
using NodeRuntimes = StaticMap<Node*, StaticArray<const RuntimeInfo*, cMaxNumNodeRuntimes>, cMaxNumInstances>;

Expand Down Expand Up @@ -91,13 +100,15 @@ class Balancer {
Error PerformPolicyBalancing(Array<SharedPtr<Instance>>& instances);
Error PrepareForBalancing(bool rebalancing, bool isInitialUpdate = false);
Error UpdateMonitoringData(bool isInitialUpdate = false);

ImageInfoProvider* mImageInfoProvider {};
InstanceManager* mInstanceManager {};
NodeManager* mNodeManager {};
MonitoringProviderItf* mMonitorProvider {};
InstanceRunnerItf* mRunner {};
SubjectArray mSubjects;
Error SendFailedInstanceStatuses();

ImageInfoProvider* mImageInfoProvider {};
InstanceManager* mInstanceManager {};
NodeManager* mNodeManager {};
MonitoringProviderItf* mMonitorProvider {};
InstanceRunnerItf* mRunner {};
statushandler::HandlerItf* mStatusHandler {};
SubjectArray mSubjects;

StaticAllocator<cAllocatorSize> mAllocator;
};
Expand Down
3 changes: 1 addition & 2 deletions src/core/cm/launcher/itf/launcher.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
#ifndef AOS_CORE_CM_LAUNCHER_ITF_LAUNCHER_HPP_
#define AOS_CORE_CM_LAUNCHER_ITF_LAUNCHER_HPP_

#include <core/common/instancestatusprovider/itf/instancestatusprovider.hpp>
#include <core/common/types/instance.hpp>

#include <core/cm/launcher/itf/types.hpp>
Expand All @@ -21,7 +20,7 @@ namespace aos::cm::launcher {
/**
* Instance launcher interface.
*/
class LauncherItf : public instancestatusprovider::ProviderItf {
class LauncherItf {
public:
/**
* Destructor.
Expand Down
77 changes: 13 additions & 64 deletions src/core/cm/launcher/launcher.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
unitconfig::NodeConfigProviderItf& nodeConfigProvider, storagestate::StorageStateItf& storageState,
MonitoringProviderItf& monitorProvider, alerts::AlertsProviderItf& alertsProvider,
iamclient::IdentProviderItf& identProvider, IdentifierPoolValidator gidValidator,
IdentifierPoolValidator uidValidator, StorageItf& storage)
IdentifierPoolValidator uidValidator, StorageItf& storage, statushandler::HandlerItf& statusHandler)
{
LOG_DBG() << "Init Launcher";

Expand All @@ -48,6 +48,7 @@
mMonitorProvider = &monitorProvider;
mAlertsProvider = &alertsProvider;
mIdentProvider = &identProvider;
mStatusHandler = &statusHandler;

mImageInfoProvider.Init(itemInfoProvider, ociSpec);

Expand All @@ -58,7 +59,7 @@

mRunRequestsLoader.Init(storage, mInstanceManager, mImageInfoProvider);
mNodeManager.Init(*mNodeInfoProvider, *mNodeConfigProvider, *mRunner);
mBalancer.Init(mInstanceManager, mImageInfoProvider, mNodeManager, *mMonitorProvider, *mRunner);
mBalancer.Init(mInstanceManager, mImageInfoProvider, mNodeManager, *mMonitorProvider, *mRunner, statusHandler);

return ErrorEnum::eNone;
}
Expand Down Expand Up @@ -133,6 +134,11 @@
return AOS_ERROR_WRAP(err);
}

// Send instance statuses to statushandler.
if (auto err = mBalancer.SendInstanceStatuses(); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

// Load SM data for active instances.
if (auto err = mBalancer.LoadSMDataForActiveInstances(); !err.IsNone()) {
LOG_ERR() << "Can't load SM data for active instances" << Log::Field(err);
Expand Down Expand Up @@ -241,52 +247,17 @@
return AOS_ERROR_WRAP(ErrorEnum::eCanceled);
}

if (auto err = BalanceInstances(updateLock, false); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

if (auto err = statuses.Assign(mInstanceStatuses); !err.IsNone()) {
if (auto err = BalanceInstances(false); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

return ErrorEnum::eNone;
}

Error Launcher::GetInstancesStatuses(Array<InstanceStatus>& statuses)
{
LockGuard updateLock {mUpdateMutex};

if (auto err = statuses.Assign(mInstanceStatuses); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

return ErrorEnum::eNone;
}

Error Launcher::SubscribeListener(instancestatusprovider::ListenerItf& listener)
{
LockGuard updateLock {mUpdateMutex};

LOG_DBG() << "Subscribe instance status listener";

if (auto err = mInstanceStatusListeners.PushBack(&listener); !err.IsNone()) {
return AOS_ERROR_WRAP(err);
}

return ErrorEnum::eNone;
}

Error Launcher::UnsubscribeListener(instancestatusprovider::ListenerItf& listener)
{
LockGuard updateLock {mUpdateMutex};

LOG_DBG() << "Unsubscribe instance status listener";

auto count = mInstanceStatusListeners.Remove(&listener);

return count == 0 ? AOS_ERROR_WRAP(ErrorEnum::eNotFound) : ErrorEnum::eNone;
}

Error Launcher::OverrideEnvVars(const OverrideEnvVarsRequest& envVars)
{
LOG_DBG() << "Override env vars";
Expand Down Expand Up @@ -366,40 +337,18 @@
<< Log::Field("manifestDigest", status.mManifestDigest) << Log::Field("state", status.mState)
<< Log::Field(status.mError);
}

for (auto& listener : mInstanceStatusListeners) {
listener->OnInstancesStatusesChanged(*changedStatuses);
}
}

void Launcher::FailActivatingInstances()
{
for (auto& instance : mInstanceManager.GetActiveInstances()) {
if (instance->GetStatus().mState == aos::InstanceStateEnum::eActivating
&& instance->GetStatus().mType != UpdateItemTypeEnum::eComponent) {
const auto& instanceInfo = instance->GetInfo();

// Keep node ID, because instance still scheduled, but node didn't send activating status.
LOG_ERR() << "Instance failed to activate" << Log::Field("instance", instanceInfo.mInstanceIdent)
<< Log::Field("runtimeID", instanceInfo.mRuntimeID)
<< Log::Field("manifestDigest", instanceInfo.mManifestDigest);

instance->SetError(AOS_ERROR_WRAP(ErrorEnum::eTimeout), false);
}
}
}

Error Launcher::BalanceInstances(UniqueLock<Mutex>& lock, bool rebalance)
Error Launcher::BalanceInstances(bool rebalance)
{
LOG_DBG() << "Balance instances" << Log::Field("rebalance", rebalance);

// Create instances from run requests.
auto instances = MakeUnique<StaticArray<SharedPtr<Instance>, cMaxNumInstances>>(&mAllocator);
mRunRequestsLoader.CreateInstances(mNodeManager.GetNodes(), *instances);

auto runErr = mBalancer.RunInstances(lock, *instances, rebalance);
auto runErr = mBalancer.RunInstances(*instances, rebalance);

FailActivatingInstances();
UpdateInstanceStatuses();

if (!runErr.IsNone()) {
Expand Down Expand Up @@ -486,7 +435,7 @@
// Resend instances.
if (!mUpdatedNodes.IsEmpty()) {
if (!doRebalance) {
err = mNodeManager.ResendInstances(updateLock, mUpdatedNodes, mInstanceManager.GetActiveInstances(),
err = mNodeManager.ResendInstances(mUpdatedNodes, mInstanceManager.GetActiveInstances(),
mInstanceManager.GetRunningInstances(), forceRestart);
if (!err.IsNone()) {
LOG_ERR() << "Failed to resend instances" << Log::Field(AOS_ERROR_WRAP(err));
Expand All @@ -502,7 +451,7 @@
if (doRebalance) {
mForceRebalance = false;

if (err = BalanceInstances(updateLock, true); !err.IsNone()) {
if (err = BalanceInstances(true); !err.IsNone()) {

Check warning on line 454 in src/core/cm/launcher/launcher.cpp

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

This initialization is not a variable declaration; move it out of the "if".

See more on https://sonarcloud.io/project/issues?id=aosedge_aos_core_lib_cpp&issues=AZ5y0bimYLTB-VuxFrvZ&open=AZ5y0bimYLTB-VuxFrvZ&pullRequest=573
LOG_ERR() << "Rebalancing failed" << Log::Field(AOS_ERROR_WRAP(err));
}
}
Expand Down
Loading
Loading