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