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