Line data Source code
1 : /**
2 : * Copyright (c) 2025 Huawei Technologies Co., Ltd.
3 : * This program is free software, you can redistribute it and/or modify it under the terms and conditions of
4 : * CANN Open Software License Agreement Version 2.0 (the "License").
5 : * Please refer to the License for details. You may not use this file except in compliance with the License.
6 : * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, EITHER EXPRESS OR IMPLIED,
7 : * INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, MERCHANTABILITY, OR FITNESS FOR A PARTICULAR PURPOSE.
8 : * See LICENSE in the root of the software repository for the full text of the License.
9 : */
10 :
11 : #include "aicpu_kfc_deprecated_process.h"
12 :
13 : #include "dfx/dfx_extend_info.h"
14 : #include "aicpu_kfc_process.h"
15 : #include "algorithm/task_orchestrator.h"
16 : #include "common/aicpu_kfc_utils.h"
17 : #include "utils/aicpu_hdc_utils.h"
18 : #include "framework/aicpu_hccl_process.h"
19 :
20 : using namespace hccl;
21 :
22 : ANONYMOUS_NAMESPACE_BEGIN
23 42 : bool HcclOpCheckSupportRetry(HcclCMDType opType)
24 : {
25 : const std::set<HcclCMDType> HcclSupportRetryOpSet
26 : = {HcclCMDType::HCCL_CMD_BROADCAST, HcclCMDType::HCCL_CMD_ALLREDUCE, HcclCMDType::HCCL_CMD_REDUCE,
27 : HcclCMDType::HCCL_CMD_ALLGATHER, HcclCMDType::HCCL_CMD_REDUCE_SCATTER, HcclCMDType::HCCL_CMD_ALLTOALLV,
28 : HcclCMDType::HCCL_CMD_ALLTOALLVC, HcclCMDType::HCCL_CMD_ALLTOALL, HcclCMDType::HCCL_CMD_GATHER,
29 84 : HcclCMDType::HCCL_CMD_SCATTER};
30 84 : return (HcclSupportRetryOpSet.find(opType) != HcclSupportRetryOpSet.end());
31 42 : }
32 :
33 36 : void HcclUpdateOpIndex(HcclCMDType opType, AicpuComContext* ctx)
34 : {
35 36 : if (HcclOpCheckSupportRetry(opType)) {
36 35 : auto opIndex = ctx->opIndex + 1;
37 35 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, opIndex), opIndex);
38 : } else {
39 : // NOTE: send / recv / batchsendrecv 算子不是通信域内所有卡都参与,opIndex需要另行处理;重执行暂不支持该类算子
40 : }
41 36 : return;
42 : }
43 :
44 30 : HcclResult UpdateOpExecStatus(
45 : AicpuComContext* ctx, HcclOpExecFSM& fsmState, KfcStatus state, KfcError& errorCode, uint32_t retryCnt)
46 : {
47 30 : auto ret = AicpuHdcUtils::SetOpExecStatus(ctx->kfcStatusTransferD2H, state, errorCode, retryCnt);
48 30 : if (ret != HCCL_SUCCESS) {
49 1 : HCCL_ERROR("SetOpExecStatus failed, ret:%u", ret);
50 1 : errorCode = KfcError::kExec;
51 1 : fsmState = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
52 : }
53 30 : return ret;
54 : }
55 :
56 7 : bool HcclOpCheckInplace(const AivAicpuOpParam& opParams)
57 : {
58 7 : if (opParams.sendBuffer != opParams.recvBuffer) {
59 6 : return false;
60 : }
61 :
62 : const std::set<HcclCMDType> HcclInplaceOpSet
63 : = {HcclCMDType::HCCL_CMD_ALLREDUCE, HcclCMDType::HCCL_CMD_REDUCE, HcclCMDType::HCCL_CMD_ALLGATHER,
64 : HcclCMDType::HCCL_CMD_REDUCE_SCATTER, HcclCMDType::HCCL_CMD_ALLTOALLV, HcclCMDType::HCCL_CMD_ALLTOALLVC,
65 2 : HcclCMDType::HCCL_CMD_ALLTOALL, HcclCMDType::HCCL_CMD_GATHER, HcclCMDType::HCCL_CMD_SCATTER};
66 1 : if (HcclInplaceOpSet.find(opParams.commType) != HcclInplaceOpSet.end()) {
67 1 : return true;
68 : }
69 0 : return false;
70 1 : }
71 :
72 8 : bool HcclOpSupportRetry(AicpuComContext* ctx, AivAicpuOpParam& opParams)
73 : {
74 8 : if (!ctx->retryEnable) {
75 1 : HCCL_INFO("hccl aicpu can not retry, enable[%u].", ctx->retryEnable);
76 1 : return false;
77 : }
78 :
79 : // 不支持inplace的通信算子重执行
80 7 : if (HcclOpCheckInplace(opParams)) {
81 1 : HCCL_INFO(
82 : "hccl aicpu can not retry, opType[%u], sendBuffer[0x%016lx], recvBuffer[0x%016lx].", opParams.commType,
83 : opParams.sendBuffer, opParams.recvBuffer);
84 1 : return false;
85 : }
86 :
87 6 : if (HcclOpCheckSupportRetry(opParams.commType)) {
88 6 : return true;
89 : }
90 0 : return false;
91 : }
92 :
93 : #ifdef CCL_LLT
94 : static constexpr u32 HCCL_AICPU_WAIT_HOST_BASE_TIME_MS = 200U;
95 : #else
96 : static constexpr u32 HCCL_AICPU_WAIT_HOST_BASE_TIME_MS = 200000U;
97 : #endif
98 1035 : u32 HcclGetWaitRetryCmdTimeout(AicpuComContext* ctx, uint32_t retryCnt)
99 : {
100 1035 : if (retryCnt == 0) {
101 1035 : return HCCL_AICPU_WAIT_HOST_BASE_TIME_MS + ctx->retryHoldTime;
102 : } else {
103 0 : return HCCL_AICPU_WAIT_HOST_BASE_TIME_MS + ctx->retryIntervalTime;
104 : }
105 : }
106 :
107 15 : HcclResult HcclOpExecFsmInitProcess(
108 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode, AicpuKfcRpcServer& rpc, AivAicpuOpParam& opParams)
109 : {
110 15 : rpc.CheckRcvAddrMsg(&opParams, 0);
111 15 : ctx->directlySendMainSteramSqe = true;
112 :
113 15 : HcclUpdateOpIndex(opParams.commType, ctx);
114 15 : opParams.opId.index = ctx->opIndex;
115 15 : if (ctx->endStopLaunch) {
116 0 : HCCL_WARNING("[NsRecovery] Suspending status should not launch task");
117 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_STOP_LAUNCH;
118 0 : return HCCL_SUCCESS;
119 : }
120 15 : auto ret = AicpuHdcUtils::InitOpExecStatus(ctx->kfcStatusTransferD2H, opParams.opId);
121 15 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, isOpLaunch), true);
122 15 : if (ret == HCCL_SUCCESS) {
123 14 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_LAUNCH;
124 : } else {
125 1 : HCCL_ERROR("InitOpExecStatus failed, ret:%u", ret);
126 1 : errorCode = KfcError::kInner;
127 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
128 : }
129 15 : return ret;
130 : }
131 :
132 448367 : HcclResult HcclOpExecFsmStoppingProcess(AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode)
133 : {
134 448367 : HCCL_DEBUG("hccl aicpu stopping.");
135 448367 : if (TaskOrchestrator::IsTaskExceptionForHccs(ctx)) {
136 1 : HCCL_INFO("hccl aicpu recoverable task exception occurs.");
137 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPED;
138 1 : return HCCL_SUCCESS;
139 : }
140 :
141 448366 : KfcCommand cmd = KfcCommand::kNone;
142 448366 : auto ret = AicpuHdcUtils::GetOpExecCtrlCmd(ctx->kfcControlTransferH2D, cmd);
143 448366 : if (ret != HCCL_SUCCESS) {
144 1 : HCCL_ERROR("GetOpExecCtrlCmd failed, ret:%u", ret);
145 1 : errorCode = KfcError::kExec;
146 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
147 1 : return ret;
148 : }
149 448365 : if (cmd == KfcCommand::kExit) {
150 2 : HCCL_WARNING("hccl aicpu exec fsm stop by exit cmd.");
151 2 : errorCode = KfcError::kExit;
152 2 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
153 448363 : } else if ((cmd == KfcCommand::kStopExec)) {
154 7 : HCCL_INFO("hccl aicpu get stop exec cmd.");
155 7 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPED;
156 448356 : } else if ((cmd == KfcCommand::kNone) || (cmd == KfcCommand::kStopLaunch)) {
157 448356 : HCCL_DEBUG("hccl aicpu wait for stop exec cmd.");
158 : // do nothing
159 : } else {
160 0 : HCCL_ERROR("GetOpExecCtrlCmd failed, invalid cmd[%u]", cmd);
161 0 : errorCode = KfcError::kExec;
162 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
163 : }
164 448365 : return HCCL_SUCCESS;
165 : }
166 :
167 10 : HcclResult HcclOpExecFsmStoppedProcess(
168 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode, u32 retryCnt, AivAicpuOpParam& opParams,
169 : u32 beginSqePos, u32 endSqePos)
170 : {
171 10 : HCCL_DEBUG("hccl aicpu stop exec.");
172 10 : KfcCommand cmd = KfcCommand::kNone;
173 10 : auto ret = AicpuHdcUtils::GetOpExecCtrlCmd(ctx->kfcControlTransferH2D, cmd);
174 10 : if (ret != HCCL_SUCCESS) {
175 1 : HCCL_ERROR("GetOpExecCtrlCmd failed, ret:%u", ret);
176 1 : errorCode = KfcError::kExec;
177 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
178 1 : return ret;
179 : }
180 :
181 9 : if (cmd == KfcCommand::kExit) {
182 1 : HCCL_ERROR("hccl aicpu exec fsm stop by exit cmd.");
183 1 : errorCode = KfcError::kExit;
184 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
185 1 : return HCCL_SUCCESS;
186 : }
187 :
188 8 : if (!HcclOpSupportRetry(ctx, opParams)) {
189 2 : HCCL_ERROR("hccl aicpu not support retry, enable[%u], commType[%u].", ctx->retryEnable, opParams.commType);
190 2 : errorCode = KfcError::kExec;
191 2 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
192 2 : return HCCL_SUCCESS;
193 : }
194 :
195 6 : uint32_t sqHead = 0xFFFFFFFF;
196 6 : CHK_RET(QuerySqStatusByType(ctx->devId, ctx->streamInfo[ctx->rankId].sqId, DRV_SQCQ_PROP_SQ_HEAD, sqHead));
197 6 : if (sqHead == endSqePos) {
198 1 : HCCL_INFO(
199 : "hccl aicpu record complete task is complete, can not retry. params: sqHead %u, beginSqePos %u "
200 : "endSqePos %u",
201 : sqHead, beginSqePos, endSqePos);
202 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_END;
203 5 : } else if (sqHead == beginSqePos) {
204 0 : HCCL_ERROR(
205 : "hccl aicpu wait start task is not complete, can not retry. params: sqHead %u, beginSqePos %u "
206 : "endSqePos %u",
207 : sqHead, beginSqePos, endSqePos);
208 0 : errorCode = KfcError::kExec;
209 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
210 : } else {
211 5 : HCCL_INFO(
212 : "hccl aicpu op is running, can retry. params: sqHead %u, beginSqePos %u endSqePos %u", sqHead, beginSqePos,
213 : endSqePos);
214 5 : if (TaskOrchestrator::IsTaskExceptionForHccs(ctx)) {
215 1 : HCCL_INFO("hccl aicpu stop by sdma/write task exception, can retry.");
216 1 : errorCode = KfcError::kSdma;
217 : }
218 5 : CHK_RET(UpdateOpExecStatus(ctx, state, KfcStatus::kStopExec, errorCode, retryCnt));
219 5 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_RETRY;
220 : }
221 6 : return HCCL_SUCCESS;
222 : }
223 :
224 4 : HcclResult HcclOpExecFsmEndProcess(AicpuComContext* ctx, uint32_t retryCnt, AivAicpuOpParam& opParams)
225 : {
226 4 : auto ret = AicpuHdcUtils::SetOpExecStatus(ctx->kfcStatusTransferD2H, KfcStatus::kEnd, KfcError::kNone, retryCnt);
227 4 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, isOpLaunch), false);
228 4 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.kfcStatus), DfxKfcStatus::kOneFinished);
229 4 : AicpuKfcUtils::PrintBuffer(ctx, opParams);
230 4 : ctx->directlySendMainSteramSqe = false;
231 4 : return ret;
232 : }
233 : ANONYMOUS_NAMESPACE_END
234 :
235 14 : HcclResult AicpuKfcDeprecatedProcess::LaunchHcclOp(
236 : AicpuComContext* ctx, AivAicpuOpParam* commParam, uint32_t& beginSqePos, uint32_t& endSqePos)
237 : {
238 : // 获取通信stream上首次下发的notify wait
239 : // task的尾指针,已便重执行stop时判断是否已执行该task,如果该task已执行完成则可支持通信重执行
240 14 : CHK_RET(QuerySqStatusByType(ctx->devId, ctx->streamInfo[ctx->rankId].sqId, DRV_SQCQ_PROP_SQ_TAIL, beginSqePos));
241 :
242 : // STARS调度执行到该通信算子时,会触发一次本地notify record触发通信算子在AICPU上展开、执行
243 14 : CHK_RET(AicpuDispatcher::AicpuUnfoldSignalWait(ctx->rankId, 0, AicpuDispatcher::IPC));
244 14 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(commParam, nullptr, ctx));
245 :
246 : // AICPU上通信task下发完成后,在通信stream上紧跟着下发一个notify record,以通知通信主stream通信算子执行完成
247 13 : CHK_RET(AicpuDispatcher::AicpuUnfoldSignalRecord(ctx->rankId, 1, AicpuDispatcher::IPC));
248 :
249 13 : KfcCommand cmd = KfcCommand::kNone;
250 13 : if ((ctx->endStopLaunch == false) && (ctx->commOpenStatus == true)) {
251 13 : CHK_RET(AicpuHdcUtils::GetOpExecCtrlCmd(ctx->kfcControlTransferH2D, cmd));
252 13 : if (cmd == KfcCommand::NsStopLaunch) {
253 0 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, endStopLaunch), true);
254 0 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, isStopLaunch), true);
255 0 : return HCCL_E_SUSPENDING;
256 : }
257 : }
258 : // 启动通信task执行
259 13 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
260 13 : CHK_RET(QuerySqStatusByType(ctx->devId, ctx->streamInfo[ctx->rankId].sqId, DRV_SQCQ_PROP_SQ_TAIL, endSqePos));
261 13 : HCCL_INFO(
262 : "hccl aicpu launch hccl op task success. stream sqid:%d begin:%u end:%u", ctx->streamInfo[ctx->rankId].sqId,
263 : beginSqePos, endSqePos);
264 13 : return HCCL_SUCCESS;
265 : }
266 :
267 : HcclResult
268 2 : AicpuKfcDeprecatedProcess::RetryLaunchHcclOp(AicpuComContext* ctx, AivAicpuOpParam* commParam, uint32_t& endSqePos)
269 : {
270 2 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(commParam, nullptr, ctx));
271 :
272 : // AICPU上通信task下发完成后,在通信stream上紧跟着下发一个notify record,以通知通信主stream通信算子执行完成
273 2 : CHK_RET(AicpuDispatcher::AicpuUnfoldSignalRecord(ctx->rankId, 1, AicpuDispatcher::IPC));
274 :
275 : // 启动通信task执行
276 2 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
277 :
278 2 : CHK_RET(QuerySqStatusByType(ctx->devId, ctx->streamInfo[ctx->rankId].sqId, DRV_SQCQ_PROP_SQ_TAIL, endSqePos));
279 :
280 2 : HCCL_INFO(
281 : "hccl aicpu retry launch hccl op task success. stream sqid:%d end:%u", ctx->streamInfo[ctx->rankId].sqId,
282 : endSqePos);
283 2 : return HCCL_SUCCESS;
284 : }
285 :
286 21 : HcclResult AicpuKfcDeprecatedProcess::RunRpcServerOneStageWait(AicpuComContext* ctx, AicpuKfcRpcServer& rpc)
287 : {
288 84 : AivAicpuOpParam g_msg[3];
289 21 : AivAicpuOpParam* msg = &g_msg[0];
290 21 : AivAicpuOpParam* preMsg = &g_msg[1];
291 21 : AivAicpuOpParam* nextMsg = &g_msg[2];
292 21 : AivAicpuOpParam* tmpptr = nullptr;
293 21 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.kfcStatus), DfxKfcStatus::kOneStart);
294 21 : AicpuHcclProcess::CallMC2MaintenanceThread(ctx);
295 : // 读取首轮任务,并准备,直接先下主流(直接激活),后续还是先下从流,激活时下主流
296 21 : rpc.CheckRcvAddrMsg(msg, 0);
297 21 : if (!rpc.CheckAivIsEnd(0)) {
298 21 : rpc.ReadAddrMsg(nextMsg, 0);
299 21 : tmpptr = nextMsg;
300 : }
301 21 : HcclUpdateOpIndex(msg->commType, ctx);
302 21 : msg->opId.index = ctx->opIndex;
303 21 : if (ctx->endStopLaunch) {
304 0 : HCCL_WARNING("the op should not be launched in suspending status");
305 0 : return HCCL_E_SUSPENDING;
306 : }
307 21 : auto ret = AicpuHdcUtils::InitOpExecStatus(ctx->kfcStatusTransferD2H, msg->opId);
308 21 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, isOpLaunch), true);
309 21 : if (ret != HCCL_SUCCESS) {
310 0 : HCCL_ERROR("InitOpExecStatus failed, ret:%u", ret);
311 0 : return ret;
312 : }
313 21 : ctx->directlySendMainSteramSqe = true;
314 21 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(msg, tmpptr, ctx));
315 :
316 54 : while (!rpc.CheckAivIsEnd(0)) {
317 33 : tmpptr = msg;
318 33 : msg = preMsg;
319 33 : preMsg = tmpptr; // msg <-> preMsg
320 33 : tmpptr = nullptr;
321 :
322 : // 读取下一次任务,并编排
323 33 : rpc.CheckRcvAddrMsg(msg, 0);
324 33 : if (!rpc.CheckAivIsEnd(0)) {
325 12 : rpc.ReadAddrMsg(nextMsg, 0);
326 12 : tmpptr = nextMsg;
327 : }
328 33 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(msg, tmpptr, ctx));
329 :
330 : // 激活下一次任务执行
331 33 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
332 : }
333 :
334 : // 激活下一次任务执行
335 21 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
336 21 : ctx->directlySendMainSteramSqe = false;
337 21 : CHK_RET(AicpuKfcProcess::WaitTaskFinish(ctx));
338 20 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.kfcStatus), DfxKfcStatus::kOneFinished);
339 20 : AicpuKfcUtils::PrintBuffer(ctx, *msg);
340 20 : return HCCL_SUCCESS;
341 : }
342 :
343 15 : HcclResult AicpuKfcDeprecatedProcess::HcclOpExecFsmWaitEndProcess(
344 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode, u32 retryCnt)
345 : {
346 15 : bool isWaitTask = (ctx->debugMode == MC2_DEBUG_WAIT_COMM);
347 15 : auto ret = AicpuKfcProcess::WaitTaskFinish(ctx, isWaitTask);
348 15 : if (ret == HCCL_SUCCESS) {
349 3 : HCCL_DEBUG("hccl aicpu exec complete.");
350 3 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_END;
351 12 : } else if (ret == HCCL_E_SUSPENDING) {
352 11 : HCCL_RUN_INFO("[NsRecovery][AICPU]hccl aicpu force stop in launch loop");
353 11 : if (ctx->isStopLaunch == true) {
354 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_STOP_LAUNCH;
355 : } else {
356 11 : CHK_RET(UpdateOpExecStatus(ctx, state, KfcStatus::kStoplaunch, errorCode, retryCnt));
357 11 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPING;
358 : }
359 : } else {
360 1 : errorCode = KfcError::kExec;
361 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
362 : }
363 15 : return ret;
364 : }
365 :
366 1036 : HcclResult AicpuKfcDeprecatedProcess::HcclOpExecFsmWaitRetryProcess(
367 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode)
368 : {
369 1036 : KfcCommand cmd = KfcCommand::kNone;
370 1036 : auto ret = AicpuHdcUtils::GetOpExecCtrlCmd(ctx->kfcControlTransferH2D, cmd);
371 1036 : if (ret != HCCL_SUCCESS) {
372 1 : HCCL_ERROR("GetOpExecCtrlCmd failed, ret:%u", ret);
373 1 : errorCode = KfcError::kExec;
374 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
375 1 : return ret;
376 : }
377 1035 : if (cmd == KfcCommand::kRetry) {
378 4 : HCCL_INFO("hccl aicpu recv retry cmd from host.");
379 4 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.pollStatus), PollStatus::kDefault);
380 4 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.cqeStatus), dfx::CqeStatus::kDefault);
381 4 : ret = AicpuKfcProcess::ResetSqBuff(ctx);
382 4 : if (ret != HCCL_SUCCESS) {
383 1 : errorCode = KfcError::kInner;
384 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
385 1 : return ret;
386 : }
387 3 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_RETRY;
388 1031 : } else if (cmd == KfcCommand::kExit) {
389 1 : errorCode = KfcError::kExit;
390 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
391 : } else {
392 : // do nothing
393 : }
394 1034 : return HCCL_SUCCESS;
395 : }
396 :
397 15 : HcclResult AicpuKfcDeprecatedProcess::AICPU_RpcServerUnfoldStageWait(AicpuComContext* ctx, AicpuKfcRpcServer& rpc)
398 : {
399 15 : AivAicpuOpParam opParams;
400 15 : auto waitStopExecCmdTimeout = std::chrono::milliseconds(HCCL_AICPU_WAIT_HOST_BASE_TIME_MS);
401 15 : auto startTime = std::chrono::steady_clock::now();
402 :
403 15 : KfcError errorCode = KfcError::kNone;
404 15 : uint32_t retryCnt = 0;
405 15 : uint32_t beginSqePos = INVALID_UINT;
406 15 : uint32_t endSqePos = INVALID_UINT;
407 15 : HcclOpExecFSM state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_INIT;
408 15 : HcclResult ret = HCCL_SUCCESS;
409 15 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, dfxExtendInfo.kfcStatus), DfxKfcStatus::kOneStart);
410 15 : AicpuHcclProcess::CallMC2MaintenanceThread(ctx);
411 : while (true) {
412 449471 : switch (state) {
413 15 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_INIT:
414 15 : ret = HcclOpExecFsmInitProcess(ctx, state, errorCode, rpc, opParams);
415 15 : break;
416 14 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_LAUNCH:
417 14 : ret = HcclOpExecFsmLaunchProcess(ctx, state, errorCode, opParams, beginSqePos, endSqePos);
418 14 : break;
419 15 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_END:
420 15 : ret = HcclOpExecFsmWaitEndProcess(ctx, state, errorCode, retryCnt);
421 15 : if (state == HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPING) {
422 11 : startTime = std::chrono::steady_clock::now();
423 : }
424 15 : break;
425 448366 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPING:
426 448366 : if ((std::chrono::steady_clock::now() - startTime) >= waitStopExecCmdTimeout) {
427 1 : HCCL_ERROR("hccl aicpu wait stop exec timeout[%u ms].", HCCL_AICPU_WAIT_HOST_BASE_TIME_MS);
428 1 : errorCode = KfcError::kTimeout;
429 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
430 : } else {
431 448365 : ret = HcclOpExecFsmStoppingProcess(ctx, state, errorCode);
432 : }
433 448366 : break;
434 8 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_STOPPED:
435 8 : ret = HcclOpExecFsmStoppedProcess(ctx, state, errorCode, retryCnt, opParams, beginSqePos, endSqePos);
436 8 : if (state == HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_RETRY) {
437 5 : startTime = std::chrono::steady_clock::now();
438 : }
439 8 : break;
440 1035 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_RETRY:
441 1035 : if ((std::chrono::steady_clock::now() - startTime)
442 2070 : >= std::chrono::milliseconds(HcclGetWaitRetryCmdTimeout(ctx, retryCnt))) {
443 0 : HCCL_ERROR("hccl aicpu wait retry timeout[%u ms].", HcclGetWaitRetryCmdTimeout(ctx, retryCnt));
444 0 : errorCode = KfcError::kTimeout;
445 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
446 : } else {
447 1035 : ret = HcclOpExecFsmWaitRetryProcess(ctx, state, errorCode);
448 : }
449 1035 : break;
450 3 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_RETRY:
451 3 : ret = HcclOpExecFsmRetryProcess(ctx, state, errorCode, retryCnt, opParams, endSqePos);
452 3 : break;
453 4 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_END:
454 4 : return HcclOpExecFsmEndProcess(ctx, retryCnt, opParams);
455 0 : case HcclOpExecFSM::HCCL_OP_EXEC_STOP_LAUNCH:
456 0 : HCCL_DEBUG("[NsTest][AICPU] stop the kernel");
457 0 : if (!ctx->isStopLaunch) {
458 0 : return HCCL_E_SUSPENDING;
459 : } else {
460 0 : HCCL_RUN_INFO("[NsTest][AICPU] stop the kernel for stop command");
461 0 : AicpuHcclProcess::CopyCtxForBackGroundDfx(ctx);
462 0 : if (UpdateOpExecStatus(ctx, state, KfcStatus::kStoplaunch, errorCode, 0) == HCCL_SUCCESS) {
463 0 : return HCCL_E_SUSPENDING;
464 : } else {
465 0 : break;
466 : }
467 : }
468 11 : case HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR:
469 : default:
470 11 : UpdateOpExecStatus(ctx, state, KfcStatus::kError, errorCode, retryCnt);
471 11 : return (ret == HCCL_SUCCESS) ? HCCL_E_INTERNAL : ret;
472 : }
473 : }
474 : return HCCL_SUCCESS;
475 : }
476 :
477 1 : HcclResult AicpuKfcDeprecatedProcess::RunRpcServerTwoStageWait(AicpuComContext* ctx, AicpuKfcRpcServer& rpc)
478 : {
479 4 : AivAicpuOpParam gMsg[3];
480 1 : AivAicpuOpParam* msg = &gMsg[0];
481 1 : AivAicpuOpParam* msgWork = &gMsg[1];
482 1 : AivAicpuOpParam* nextMsg = &gMsg[2];
483 1 : AivAicpuOpParam* tmpptr = nullptr;
484 :
485 : // 读取首轮任务,并准备,直接先下主流(直接激活),后续还是先下从流,激活时下主流
486 : // 1.1、首轮读地址(需要自动产生)
487 1 : rpc.CheckRcvAddrMsg(msg, 0);
488 1 : if (!rpc.CheckAivIsEnd(0)) {
489 : // 读取下一轮地址
490 1 : rpc.ReadAddrMsg(nextMsg, 0);
491 1 : tmpptr = nextMsg;
492 : }
493 :
494 : // 1.2 首轮提前读看是否需要提前下主流,即判断sendcnt是否大于等于当前轮次
495 1 : if (rpc.ReadWorkMsg(msgWork, 0, (ctx->curTurnCnt + 1)) && rpc.GetWaitPolicy() != 0) {
496 0 : ctx->directlySendMainSteramSqe = true;
497 : }
498 :
499 : // 1.3 首轮编排开始
500 1 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(msg, tmpptr, ctx));
501 0 : ctx->directlySendMainSteramSqe = false;
502 : // 2、等待激活任务执行,如果前面已经激活,则ActiveRecordMain会空转一圈
503 0 : if (rpc.GetWaitPolicy() != 0) {
504 0 : rpc.CheckRcvWorkMsg(msgWork, 0, ctx->curTurnCnt);
505 : }
506 :
507 0 : AicpuKfcUtils::PrintBuffer(ctx, *msg);
508 0 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
509 0 : while (!rpc.CheckAivIsEnd(0)) {
510 0 : tmpptr = nullptr;
511 : // 3.1、读取下一轮任务
512 0 : rpc.CheckRcvAddrMsg(msg, 0);
513 0 : if (!rpc.CheckAivIsEnd(0)) {
514 0 : rpc.ReadAddrMsg(nextMsg, 0);
515 0 : tmpptr = nextMsg;
516 : }
517 :
518 : // 3.2 开始编排下一轮
519 0 : CHK_RET(AicpuKfcProcess::AicpuCcOpExe(msg, tmpptr, ctx));
520 :
521 : // 5.1 等待上一轮执行结束
522 0 : TaskOrchestrator::WaitMainStreamFinish(ctx);
523 :
524 : // 6.1 激活下一轮
525 0 : rpc.CheckRcvWorkMsg(msgWork, 0, ctx->curTurnCnt);
526 0 : CHK_RET(TaskOrchestrator::ActiveRecordMain(AicpuKfcProcess::GetActiveSqId(ctx)));
527 :
528 : // 7.1 发送上一轮消息
529 0 : rpc.PostMsg(ctx->curTurnCnt - 1);
530 : }
531 :
532 0 : if (rpc.GetRspPolicy() != 0) {
533 : // 8.1 等待执行结束
534 0 : TaskOrchestrator::WaitMainStreamFinish(ctx);
535 0 : rpc.ClearWorkMsg();
536 0 : HCCL_INFO(
537 : "[commType:%d, opType:%s, sendBuffer:%p, recvBuffer:%p, count:%d, data_type:%s, "
538 : "sendCnt:%d, rcvCnt:%d, funID:%d, valid:%d, everyTurnRsp:%d, strideLen:%d, isLast:%d",
539 : msg->commType, GetReduceOpEnumStr(msg->opType).c_str(), msg->sendBuffer, msg->recvBuffer, msg->count,
540 : GetDataTypeEnumStr(msg->hcclDataType).c_str(), msg->sendCnt, msg->rcvCnt, msg->funID, msg->valid,
541 : msg->everyTurnRsp, msg->strideLen, msg->isLast);
542 : // 9.1 发送最后一轮消息
543 0 : rpc.PostMsg(ctx->curTurnCnt);
544 : }
545 0 : AicpuKfcUtils::PrintBuffer(ctx, *msg);
546 0 : return HCCL_SUCCESS;
547 : }
548 :
549 21 : HcclResult AicpuKfcDeprecatedProcess::TryRunRpcServerOneStageWait(AicpuComContext* ctx, AicpuKfcRpcServer& rpc)
550 : {
551 21 : HCCL_INFO("Start to run, round %u", ctx->dfxExtendInfo.kfcRestartConfig.tryRestartTimes);
552 21 : if (dfx::DfxExtendInfoHelper::TryRestartTooManyTimes(ctx->dfxExtendInfo)) {
553 0 : HCCL_ERROR("Restart too many times, max try count is %u", ctx->dfxExtendInfo.kfcRestartConfig.maxRestartTimes);
554 0 : return HCCL_E_INTERNAL;
555 : }
556 21 : const auto ret = RunRpcServerOneStageWait(ctx, rpc);
557 21 : if (ret == HCCL_SUCCESS) {
558 20 : dfx::DfxExtendInfoHelper::ResetTryRestartTimes(ctx->dfxExtendInfo);
559 20 : CHK_RET(AicpuHdcUtils::SetOpExecStatus(ctx->kfcStatusTransferD2H, KfcStatus::kEnd, KfcError::kNone, 0));
560 20 : AicpuUpdatComContextMumber(offsetof(AicpuComContext, isOpLaunch), false);
561 20 : return HCCL_SUCCESS;
562 : }
563 1 : if (ctx->dfxExtendInfo.commandToKfc == CommandToKfc::kRestart) {
564 0 : dfx::DfxExtendInfoHelper::TryRestartOnceMore(ctx->dfxExtendInfo);
565 0 : return TryRunRpcServerOneStageWait(ctx, rpc);
566 : }
567 1 : dfx::DfxExtendInfoHelper::ResetTryRestartTimes(ctx->dfxExtendInfo);
568 1 : if (ctx->isStopLaunch) {
569 0 : AicpuHcclProcess::CopyCtxForBackGroundDfx(ctx);
570 0 : CHK_RET(AicpuHdcUtils::SetOpExecStatus(ctx->kfcStatusTransferD2H, KfcStatus::kStoplaunch, KfcError::kNone, 0));
571 : } else {
572 1 : CHK_RET(AicpuHdcUtils::SetOpExecStatus(ctx->kfcStatusTransferD2H, KfcStatus::kError, KfcError::kInner, 0));
573 : }
574 1 : return ret;
575 : }
576 :
577 15 : HcclResult AicpuKfcDeprecatedProcess::HcclOpExecFsmLaunchProcess(
578 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode, AivAicpuOpParam& opParams, uint32_t& beginSqePos,
579 : uint32_t& endSqePos)
580 : {
581 15 : HCCL_DEBUG("hccl aicpu start launch task");
582 15 : auto ret = AicpuKfcDeprecatedProcess::LaunchHcclOp(ctx, &opParams, beginSqePos, endSqePos);
583 15 : if (ret == HCCL_SUCCESS) {
584 13 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_END;
585 2 : } else if (ret == HCCL_E_SUSPENDING) {
586 0 : HCCL_RUN_INFO("[NsRecovery][AICPU]hccl aicpu force stop in launch process");
587 0 : state = HcclOpExecFSM::HCCL_OP_EXEC_STOP_LAUNCH;
588 : } else {
589 2 : HCCL_ERROR("Failed to launch hccl op, ret:%u", ret);
590 2 : errorCode = KfcError::kInner;
591 2 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
592 : }
593 15 : return ret;
594 : }
595 :
596 3 : HcclResult AicpuKfcDeprecatedProcess::HcclOpExecFsmRetryProcess(
597 : AicpuComContext* ctx, HcclOpExecFSM& state, KfcError& errorCode, uint32_t& retryCnt, AivAicpuOpParam& opParams,
598 : uint32_t& endSqePos)
599 : {
600 3 : HCCL_DEBUG("hccl retry launch task");
601 3 : retryCnt++;
602 3 : auto ret = AicpuKfcDeprecatedProcess::RetryLaunchHcclOp(ctx, &opParams, endSqePos);
603 3 : if (ret != HCCL_SUCCESS) {
604 1 : errorCode = KfcError::kInner;
605 1 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_ERROR;
606 1 : return ret;
607 : }
608 2 : errorCode = KfcError::kNone;
609 2 : CHK_RET(UpdateOpExecStatus(ctx, state, KfcStatus::kRuning, errorCode, retryCnt));
610 2 : state = HcclOpExecFSM::HCCL_OP_EXEC_FSM_WAIT_END;
611 2 : return HCCL_SUCCESS;
612 : }
|