Skip to content
Merged
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
4 changes: 2 additions & 2 deletions src/core/cm/launcher/balancer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -46,8 +46,8 @@ Error Balancer::RunInstances(UniqueLock<Mutex>& lock, Array<SharedPtr<Instance>>
return AOS_ERROR_WRAP(err);
}

// Submit scheduled instances before sending instances to nodes.
// So following status updates will change active instances.
// 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);
}
Expand Down
7 changes: 6 additions & 1 deletion src/core/cm/launcher/instancemanager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,12 @@ Error InstanceManager::UpdateStatus(const InstanceStatus& status)

auto instance = FindActiveInstance(status);
if (!instance) {
// not expected instance received from SM.
// Ignore inactive instance, SM sometimes sends inactive status for stopped instances.
if (status.mState == aos::InstanceStateEnum::eInactive) {
return ErrorEnum::eNone;
}

// Not expected instance received from SM.
return AOS_ERROR_WRAP(ErrorEnum::eNotFound);
}

Expand Down
51 changes: 40 additions & 11 deletions src/core/cm/launcher/launcher.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,15 @@ void Launcher::UpdateInstanceStatuses()
const auto& preinstalledComponents = mInstanceManager.GetPreinstalledComponents();
const auto totalSize = activeInstances.Size() + preinstalledComponents.Size();

// Copy old statuses.
auto oldInstanceStatuses = MakeUnique<StaticArray<InstanceStatus, cMaxNumInstances>>(&mAllocator);
if (auto err = oldInstanceStatuses->Assign(mInstanceStatuses); !err.IsNone()) {
LOG_ERR() << "Failed to copy old instance statuses" << Log::Field(AOS_ERROR_WRAP(err));

return;
}

// Keep current statuses.
bool changed = mInstanceStatuses.Size() != totalSize;

if (auto err = mInstanceStatuses.Resize(totalSize); !err.IsNone()) {
Expand Down Expand Up @@ -316,15 +325,29 @@ void Launcher::UpdateInstanceStatuses()
return;
}

for (const auto& status : mInstanceStatuses) {
// Find new statuses.
auto changedStatuses = MakeUnique<StaticArray<InstanceStatus, cMaxNumInstances>>(&mAllocator);

for (size_t i = 0; i < mInstanceStatuses.Size(); ++i) {
auto newStatus = !oldInstanceStatuses->Contains(mInstanceStatuses[i]);
if (newStatus) {
if (auto err = changedStatuses->PushBack(mInstanceStatuses[i]); !err.IsNone()) {
LOG_ERR() << "Failed to add changed status" << Log::Field(AOS_ERROR_WRAP(err));

return;
}
}
}

for (const auto& status : *changedStatuses) {
LOG_INF() << "Instance status changed" << Log::Field("instance", static_cast<const InstanceIdent&>(status))
<< Log::Field("version", status.mVersion) << Log::Field("nodeID", status.mNodeID)
<< Log::Field("runtimeID", status.mRuntimeID) << Log::Field("manifestDigest", status.mManifestDigest)
<< Log::Field("state", status.mState) << Log::Field(status.mError);
}

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

Expand Down Expand Up @@ -374,7 +397,8 @@ void Launcher::ProcessUpdate()
UniqueLock updateLock {mUpdateMutex};

mProcessUpdatesCondVar.Wait(updateLock, [this]() {
return (!mUpdatedNodes.IsEmpty() || mNewSubjects.HasValue() || mAlertReceived || !mIsRunning)
return (!mUpdatedNodes.IsEmpty() || mNewSubjects.HasValue() || mAlertReceived || !mIsRunning
|| mIsNodeInfoChanged)
&& !mDisableProcessUpdates;
});

Expand Down Expand Up @@ -407,6 +431,12 @@ void Launcher::ProcessUpdate()
doRebalance = true;
}

// Process node info changed.
if (mIsNodeInfoChanged) {
mIsNodeInfoChanged = false;
doRebalance = true;
}

// Resend instances.
if (!mUpdatedNodes.IsEmpty()) {
if (!doRebalance) {
Expand Down Expand Up @@ -545,12 +575,15 @@ Error Launcher::OnInstanceStatusReceived(const InstanceStatus& status)

Error Launcher::OnNodeInstancesStatusesReceived(const String& nodeID, const Array<InstanceStatus>& statuses)
{
LOG_INF() << "Node instances statuses received" << Log::Field("nodeID", nodeID)
<< Log::Field("numStatuses", statuses.Size());

for (const auto& status : statuses) {
LOG_INF() << "Node instance status received"
<< Log::Field("instance", static_cast<const InstanceIdent&>(status))
<< Log::Field("version", status.mVersion) << Log::Field("nodeID", status.mNodeID)
<< Log::Field("runtimeID", status.mRuntimeID) << Log::Field("manifestDigest", status.mManifestDigest)
<< Log::Field("state", status.mState) << Log::Field(status.mError);
<< Log::Field("version", status.mVersion) << Log::Field("runtimeID", status.mRuntimeID)
<< Log::Field("manifestDigest", status.mManifestDigest) << Log::Field("state", status.mState)
<< Log::Field(status.mError);
}

LockGuard updateLock {mUpdateMutex};
Expand Down Expand Up @@ -593,11 +626,7 @@ void Launcher::OnNodeInfoChanged(const UnitNodeInfo& info)
UniqueLock updateLock {mUpdateMutex};

if (mNodeManager.UpdateNodeInfo(info)) {
if (auto err = PushUnique(mUpdatedNodes, info.mNodeID); !err.IsNone()) {
LOG_ERR() << "Failed to add node ID to updated nodes" << Log::Field(AOS_ERROR_WRAP(err));

return;
}
mIsNodeInfoChanged = true;

mProcessUpdatesCondVar.NotifyAll();
mAllNodesConnectedCondVar.NotifyAll();
Expand Down
3 changes: 2 additions & 1 deletion src/core/cm/launcher/launcher.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,7 @@ class Launcher : public LauncherItf,

private:
static constexpr auto cMaxNumInstanceStatusListeners = 8;
static constexpr auto cAllocatorSize = sizeof(StaticArray<InstanceStatus, cMaxNumInstances>)
static constexpr auto cAllocatorSize = 2 * sizeof(StaticArray<InstanceStatus, cMaxNumInstances>)
+ sizeof(StaticArray<SharedPtr<Instance>, cMaxNumInstances>);

void SendRunStatus();
Expand Down Expand Up @@ -188,6 +188,7 @@ class Launcher : public LauncherItf,
bool mDisableProcessUpdates {};
StaticArray<StaticString<cIDLen>, cMaxNumNodes> mUpdatedNodes;
bool mAlertReceived {};
bool mIsNodeInfoChanged {};
Optional<SubjectArray> mNewSubjects;

// Misc
Expand Down
10 changes: 6 additions & 4 deletions src/core/cm/launcher/node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -75,9 +75,11 @@ class Filter {
Cmp mCmp;
};

auto FilterByNode(const Array<InstanceStatus>& array, const String& nodeID)
auto FilterActiveNodeInstances(const Array<InstanceStatus>& array, const String& nodeID)
{
auto cmp = [nodeID](const InstanceStatus& status) { return status.mNodeID == nodeID; };
auto cmp = [nodeID](const InstanceStatus& status) {
return status.mNodeID == nodeID && status.mState != aos::InstanceStateEnum::eInactive;
};

return Filter<InstanceStatus, decltype(cmp)>(array, cmp);
}
Expand Down Expand Up @@ -309,7 +311,7 @@ Error Node::SendScheduledInstances(
auto stopInstances = MakeUnique<StaticArray<aos::InstanceInfo, cMaxNumInstances>>(mAllocator);
auto startInstances = MakeUnique<StaticArray<aos::InstanceInfo, cMaxNumInstances>>(mAllocator);

for (const auto& status : FilterByNode(runningInstances, mInfo.mNodeID)) {
for (const auto& status : FilterActiveNodeInstances(runningInstances, mInfo.mNodeID)) {
// Check if the instance is scheduled on this node.
auto isScheduled = scheduledInstances.ContainsIf([&status, this](const SharedPtr<Instance>& item) {
return static_cast<const InstanceIdent&>(status) == item->GetInfo().mInstanceIdent
Expand Down Expand Up @@ -359,7 +361,7 @@ RetWithError<bool> Node::ResendInstances(
auto startInstances = MakeUnique<StaticArray<aos::InstanceInfo, cMaxNumInstances>>(mAllocator);
size_t runningNodeInstances = 0;

for (const auto& status : FilterByNode(runningInstances, mInfo.mNodeID)) {
for (const auto& status : FilterActiveNodeInstances(runningInstances, mInfo.mNodeID)) {
runningNodeInstances++;

auto isActive = activeInstances.ContainsIf([&status](const SharedPtr<Instance>& item) {
Expand Down
3 changes: 1 addition & 2 deletions src/core/cm/launcher/tests/launcher.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1454,7 +1454,6 @@ TEST_F(CMLauncherTest, Balancing)
ASSERT_TRUE(mLauncher.RunInstances(testItem.mRunRequests, *runStatuses).IsNone());

ASSERT_EQ(*runStatuses, testItem.mExpectedRunStatus);
ASSERT_EQ(instanceStatusListener.GetLastStatuses(), testItem.mExpectedRunStatus);

// Rebalance
if (testItem.mRebalancing) {
Expand Down Expand Up @@ -1584,7 +1583,7 @@ TEST_F(CMLauncherTest, PlatformFiltering)
expectedRunStatus->PushBack(CreateInstanceStatus(CreateInstanceIdent(cService3, cSubject1, 0), cNodeIDRemoteSM2,
cRunnerRunc, aos::InstanceStateEnum::eActive, ErrorEnum::eNone, "", false, manifestDigest.CStr()));

ASSERT_EQ(instanceStatusListener.GetLastStatuses(), *expectedRunStatus);
ASSERT_EQ(*runStatuses, *expectedRunStatus);
}

TEST_F(CMLauncherTest, ResendInstancesOnMismatchedNodeStatus)
Expand Down
Loading