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 "externalinput_pub.h"
12 : #include "callback_thread_manager.h"
13 : #include "device_capacity.h"
14 : #include "sal_pub.h"
15 : #include "hccl_callback_task.h"
16 :
17 : namespace hccl {
18 : constexpr u32 HOST_NIC_THREAD_WAIT = 10000; // 等待hostnic监控线程时间1s(10000 * 100us);
19 :
20 521 : HcclCallbackTask::HcclCallbackTask(u32 devicePhyId, u32 deviceLogicId,
21 521 : HcclDispatcher dispatcher, NICDeployment nicDeployment)
22 521 : : devicePhyId_(devicePhyId), deviceLogicId_(deviceLogicId),
23 521 : dispatcher_(dispatcher), nicDeployment_(nicDeployment),
24 521 : callbackThread_(nullptr), callbackThreadId_(INVALID_U64),
25 521 : callbackThreadShutDown_(false)
26 : {
27 521 : }
28 :
29 521 : HcclCallbackTask::~HcclCallbackTask()
30 : {
31 : // 析构时自动停止线程
32 521 : CloseCallbackThread();
33 521 : }
34 :
35 0 : void HcclCallbackTask::CallbackThread()
36 : {
37 : //给当前线程添加名字
38 0 : SetThreadName("Hccl_Callback");
39 :
40 0 : CHK_PRT(hrtSetDevice(deviceLogicId_));
41 0 : callbackThreadId_ = pthread_self();
42 0 : while (!callbackThreadShutDown_) {
43 : // 等待1000ms,等待callback函数的返回
44 0 : HcclResult ret = hrtProcessReport(1000);
45 0 : if (ret != HCCL_SUCCESS) {
46 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
47 : }
48 : }
49 0 : CHK_PRT(hrtResetDevice(deviceLogicId_));
50 0 : }
51 :
52 183 : HcclResult HcclCallbackTask::CallbackRegStream(rtStream_t stream)
53 : {
54 : // callback线程绑定stream场景:
55 : // 1、910A host RDMA、host tcp
56 : // 2、310P 2pg device侧下沉 tcp、标准roce(当前device SOC下沉场景都需要callback task,若之后有改动,判断条件需要更新)
57 356 : if ((static_cast<s32>(devicePhyId_) == HOST_DEVICE_ID) || (stream == nullptr) ||
58 183 : ((nicDeployment_ == NICDeployment::NIC_DEPLOYMENT_DEVICE && !GetExternalInputHcclIsTcpMode() &&
59 176 : !Is310PDevice()))) {
60 172 : return HCCL_SUCCESS;
61 : }
62 :
63 1 : if (HcclGetCallbackResult(dispatcher_) != HCCL_SUCCESS) {
64 0 : HCCL_ERROR("[HcclCallbackTask][CallbackRegStream]errNo[0x%016llx] callback func err",
65 : HCCL_ERROR_CODE(HcclGetCallbackResult(dispatcher_)));
66 0 : return HcclGetCallbackResult(dispatcher_);
67 : }
68 :
69 0 : if (callbackThread_ == nullptr) {
70 0 : callbackThread_.reset(new (std::nothrow) std::thread(&HcclCallbackTask::CallbackThread, this));
71 0 : CHK_SMART_PTR_NULL(callbackThread_);
72 : }
73 :
74 0 : u32 countTime = 0;
75 0 : while (callbackThreadId_ == INVALID_U64) {
76 0 : SaluSleep(ONE_MILLISECOND_OF_USLEEP);
77 0 : countTime++;
78 0 : CHK_PRT_RET(countTime >= HOST_NIC_THREAD_WAIT,
79 : HCCL_ERROR("[HcclCallbackTask][CallbackRegStream]errNo[0x%016llx] waiting Callback Thread time out",
80 : HCCL_ERROR_CODE(HCCL_E_INTERNAL)), HCCL_E_INTERNAL);
81 : }
82 :
83 : // 如果当前stream已经被注册,则直接返回成功
84 0 : if (ThreadStreamManager::Instance().StreamHasBeenReged(stream)) {
85 0 : HCCL_INFO("[HcclCallbackTask][CallbackRegStream]Cur stream Already registered, stream[%p] tid:[%llu] ",
86 : stream, callbackThreadId_);
87 0 : return HCCL_SUCCESS;
88 : } else {
89 0 : CHK_RET(hrtSubscribeReport(callbackThreadId_, stream));
90 : // ReleaseTidAndStream 暂无相应调用
91 0 : CHK_RET(ThreadStreamManager::Instance().RegTidAndStream(callbackThreadId_, stream));
92 0 : HCCL_INFO("[HcclCallbackTask][CallbackRegStream]rt Subscribe Report success[%llu]", callbackThreadId_);
93 : }
94 0 : return HCCL_SUCCESS;
95 : }
96 :
97 521 : HcclResult HcclCallbackTask::CloseCallbackThread()
98 : {
99 1042 : if (nicDeployment_ == NICDeployment::NIC_DEPLOYMENT_DEVICE && !GetExternalInputHcclIsTcpMode() &&
100 521 : !Is310PDevice()) {
101 521 : return HCCL_SUCCESS;
102 : }
103 :
104 0 : callbackThreadShutDown_ = true;
105 0 : if (callbackThread_ != nullptr && callbackThread_->joinable()) {
106 0 : callbackThread_->join();
107 0 : callbackThread_ = nullptr;
108 : }
109 :
110 0 : callbackThreadId_ = INVALID_U64;
111 0 : return HCCL_SUCCESS;
112 : }
113 : } // namespace hccl
|