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 "offload_stream_manager.h"
12 :
13 : namespace hccl {
14 533 : OffloadStreamManager::OffloadStreamManager() = default;
15 532 : OffloadStreamManager::~OffloadStreamManager() = default;
16 :
17 3 : HcclResult OffloadStreamManager::RegisterMaster(const std::string& tag, Stream& stream)
18 : {
19 3 : std::unique_lock<std::mutex> lock(masterMapMutex_);
20 3 : masterMap_[tag] = stream;
21 3 : HCCL_DEBUG(
22 : "[OffloadStreamManager][RegisterMaster]register master stream[%p] success, tag[%s].", stream.ptr(),
23 : tag.c_str());
24 3 : return HCCL_SUCCESS;
25 3 : }
26 :
27 131 : HcclResult OffloadStreamManager::RegisterSlaves(const std::string& tag, std::vector<Stream>& stream)
28 : {
29 131 : HCCL_DEBUG(
30 : "[OffloadStreamManager][RegisterSlaves]start register slaves stream, tag[%s], size[%u].", tag.c_str(),
31 : stream.size());
32 :
33 131 : std::unique_lock<std::mutex> lock(slavesMapMutex_);
34 131 : auto iter = slavesMap_.find(tag);
35 131 : if (iter != slavesMap_.end()) {
36 1 : HCCL_ERROR(
37 : "[OffloadStreamManager][RegisterSlaves]in offload stream manager, register slaves fail,"
38 : "tag[%s] has existed",
39 : tag.c_str());
40 1 : return HCCL_E_PARA;
41 : }
42 130 : slavesMap_.insert(std::make_pair(tag, stream));
43 130 : HCCL_INFO(
44 : "[OffloadStreamManager][RegisterSlaves]register slaves stream success, tag[%s], size[%u].", tag.c_str(),
45 : stream.size());
46 130 : return HCCL_SUCCESS;
47 131 : }
48 :
49 2 : Stream OffloadStreamManager::GetMaster(const std::string& tag)
50 : {
51 2 : std::unique_lock<std::mutex> lock(masterMapMutex_);
52 2 : auto iter = masterMap_.find(tag);
53 2 : if (iter == masterMap_.end()) {
54 1 : HCCL_ERROR("[OffloadStreamManager][GetMaster]can't find tag[%s]", tag.c_str());
55 1 : return Stream();
56 : }
57 1 : return iter->second;
58 2 : }
59 :
60 22 : std::vector<Stream> OffloadStreamManager::GetSlaves(const std::string& tag, u32 num)
61 : {
62 22 : HCCL_DEBUG("[OffloadStreamManager][GetSlaves]requesting for [%u] slaves, tag[%s].", num, tag.c_str());
63 23 : if (num == 0) {
64 2 : HCCL_WARNING("[OffloadStreamManager][GetSlaves]requesting for 0 slaves, return empty vector.");
65 4 : return std::vector<Stream>(0);
66 : }
67 :
68 21 : std::unique_lock<std::mutex> lock(slavesMapMutex_);
69 21 : auto iter = slavesMap_.find(tag);
70 19 : if (iter == slavesMap_.end()) {
71 1 : HCCL_ERROR("[OffloadStreamManager][GetSlaves]can't find tag[%s]", tag.c_str());
72 1 : return std::vector<Stream>();
73 : }
74 :
75 18 : if (iter->second.size() < num) {
76 1 : HCCL_ERROR(
77 : "[OffloadStreamManager][GetSlaves]"
78 : "trying to get [%u] slaves fail, only [%u] slaves available, tag[%s].",
79 : num, iter->second.size(), tag.c_str());
80 1 : return std::vector<Stream>();
81 : }
82 :
83 17 : std::vector<Stream> res(iter->second.begin(), iter->second.begin() + num);
84 19 : iter->second.erase(iter->second.begin(), iter->second.begin() + num);
85 17 : HCCL_INFO("[OffloadStreamManager][GetSlaves]get [%u] slaves success, returning.", res.size());
86 19 : return res;
87 21 : }
88 :
89 18 : HcclResult OffloadStreamManager::ClearSlaves(const std::string& tag)
90 : {
91 18 : std::unique_lock<std::mutex> lock(slavesMapMutex_);
92 17 : auto iter = slavesMap_.find((tag));
93 17 : if (iter != slavesMap_.end()) {
94 17 : slavesMap_.erase(tag);
95 : }
96 18 : HCCL_DEBUG("[OffloadStreamManager][ClearSlaves]Destroy slaves stream success, tag[%s]", tag.c_str());
97 18 : return HCCL_SUCCESS;
98 18 : }
99 :
100 6 : HcclResult OffloadStreamManager::ClearSlaves()
101 : {
102 6 : std::unique_lock<std::mutex> lock(slavesMapMutex_);
103 6 : slavesMap_.clear();
104 6 : HCCL_DEBUG("[OffloadStreamManager][ClearSlaves]clear all slave streams.");
105 6 : return HCCL_SUCCESS;
106 6 : }
107 :
108 : } // namespace hccl
|