You can not select more than 25 topics Topics must start with a chinese character,a letter or number, can include dashes ('-') and can be up to 35 characters long.

queryresourceslogic.go 4.8 kB

11 months ago
11 months ago

  1. package schedule
  2. import (
  3. "context"
  4. "errors"
  5. "github.com/zeromicro/go-zero/core/logx"
  6. "gitlink.org.cn/JointCloud/pcm-coordinator/internal/scheduler/service/collector"
  7. "gitlink.org.cn/JointCloud/pcm-coordinator/internal/storeLink"
  8. "gitlink.org.cn/JointCloud/pcm-coordinator/internal/svc"
  9. "gitlink.org.cn/JointCloud/pcm-coordinator/internal/types"
  10. "strconv"
  11. "strings"
  12. "sync"
  13. "time"
  14. )
  15. const (
  16. ADAPTERID = "1777144940459986944" // 异构适配器id
  17. QUERY_RESOURCES = "query_resources"
  18. )
  19. type QueryResourcesLogic struct {
  20. logx.Logger
  21. ctx context.Context
  22. svcCtx *svc.ServiceContext
  23. }
  24. func NewQueryResourcesLogic(ctx context.Context, svcCtx *svc.ServiceContext) *QueryResourcesLogic {
  25. return &QueryResourcesLogic{
  26. Logger: logx.WithContext(ctx),
  27. ctx: ctx,
  28. svcCtx: svcCtx,
  29. }
  30. }
  31. func (l *QueryResourcesLogic) QueryResources(req *types.QueryResourcesReq) (resp *types.QueryResourcesResp, err error) {
  32. resp = &types.QueryResourcesResp{}
  33. if len(req.ClusterIDs) == 0 {
  34. cs, err := l.svcCtx.Scheduler.AiStorages.GetClustersByAdapterId(ADAPTERID)
  35. if err != nil {
  36. return nil, err
  37. }
  38. resources, ok := l.svcCtx.Scheduler.AiService.LocalCache[QUERY_RESOURCES]
  39. if ok {
  40. specs, ok := resources.([]*collector.ResourceSpec)
  41. if ok {
  42. results := handleEmptyResourceUsage(cs.List, specs)
  43. resp.Data = results
  44. return resp, nil
  45. }
  46. }
  47. rus, err := l.QueryResourcesByClusterId(cs.List)
  48. if err != nil {
  49. return nil, err
  50. }
  51. if checkCachingCondition(cs.List, rus) {
  52. l.svcCtx.Scheduler.AiService.LocalCache[QUERY_RESOURCES] = rus
  53. }
  54. results := handleEmptyResourceUsage(cs.List, rus)
  55. resp.Data = results
  56. } else {
  57. var clusters []types.ClusterInfo
  58. for _, id := range req.ClusterIDs {
  59. cluster, err := l.svcCtx.Scheduler.AiStorages.GetClustersById(id)
  60. if err != nil {
  61. return nil, err
  62. }
  63. clusters = append(clusters, *cluster)
  64. }
  65. if len(clusters) == 0 {
  66. return nil, errors.New("no clusters found ")
  67. }
  68. rus, err := l.QueryResourcesByClusterId(clusters)
  69. if err != nil {
  70. return nil, err
  71. }
  72. results := handleEmptyResourceUsage(clusters, rus)
  73. resp.Data = results
  74. }
  75. return resp, nil
  76. }
  77. func (l *QueryResourcesLogic) QueryResourcesByClusterId(clusterinfos []types.ClusterInfo) ([]*collector.ResourceSpec, error) {
  78. var clusters []types.ClusterInfo
  79. if len(clusterinfos) == 0 {
  80. cs, err := l.svcCtx.Scheduler.AiStorages.GetClustersByAdapterId(ADAPTERID)
  81. if err != nil {
  82. return nil, err
  83. }
  84. clusters = cs.List
  85. } else {
  86. clusters = clusterinfos
  87. }
  88. var ulist []*collector.ResourceSpec
  89. var ch = make(chan *collector.ResourceSpec, len(clusters))
  90. var wg sync.WaitGroup
  91. for _, cluster := range clusters {
  92. wg.Add(1)
  93. c := cluster
  94. go func() {
  95. defer wg.Done()
  96. done := make(chan bool)
  97. var u *collector.ResourceSpec
  98. var err error
  99. go func() {
  100. col, found := l.svcCtx.Scheduler.AiService.AiCollectorAdapterMap[strconv.FormatInt(c.AdapterId, 10)][c.Id]
  101. if !found {
  102. done <- true
  103. return
  104. }
  105. u, err = col.GetResourceSpecs(l.ctx)
  106. if err != nil {
  107. done <- true
  108. return
  109. }
  110. done <- true
  111. }()
  112. select {
  113. case <-done:
  114. if u != nil {
  115. ch <- u
  116. }
  117. case <-time.After(10 * time.Second):
  118. return
  119. }
  120. }()
  121. }
  122. wg.Wait()
  123. close(ch)
  124. for v := range ch {
  125. ulist = append(ulist, v)
  126. }
  127. return ulist, nil
  128. }
  129. func handleEmptyResourceUsage(list []types.ClusterInfo, ulist []*collector.ResourceSpec) []*collector.ResourceSpec {
  130. var rus []*collector.ResourceSpec
  131. m := make(map[string]interface{})
  132. for _, u := range ulist {
  133. if u == nil {
  134. continue
  135. }
  136. m[u.ClusterId] = u
  137. }
  138. for _, l := range list {
  139. s, ok := m[l.Id]
  140. if !ok {
  141. ru := &collector.ResourceSpec{
  142. ClusterId: l.Id,
  143. Resources: nil,
  144. Msg: "resources unavailable, please retry later",
  145. }
  146. rus = append(rus, ru)
  147. } else {
  148. if s == nil {
  149. ru := &collector.ResourceSpec{
  150. ClusterId: l.Id,
  151. Resources: nil,
  152. Msg: "resources unavailable, please retry later",
  153. }
  154. rus = append(rus, ru)
  155. } else {
  156. r, ok := s.(*collector.ResourceSpec)
  157. if ok {
  158. if r.Resources == nil || len(r.Resources) == 0 {
  159. ru := &collector.ResourceSpec{
  160. ClusterId: r.ClusterId,
  161. Resources: nil,
  162. Msg: "resources unavailable, please retry later",
  163. }
  164. rus = append(rus, ru)
  165. } else {
  166. // add cluster type
  167. t, ok := storeLink.ClusterTypeMap[strings.Title(l.Name)]
  168. if ok {
  169. r.ClusterType = t
  170. }
  171. rus = append(rus, r)
  172. }
  173. }
  174. }
  175. }
  176. }
  177. return rus
  178. }
  179. func checkCachingCondition(clusters []types.ClusterInfo, specs []*collector.ResourceSpec) bool {
  180. var count int
  181. for _, spec := range specs {
  182. if spec.Resources != nil {
  183. count++
  184. }
  185. }
  186. if count == len(clusters) {
  187. return true
  188. }
  189. return false
  190. }

PCM is positioned as Software stack over Cloud, aiming to build the standards and ecology of heterogeneous cloud collaboration for JCC in a non intrusive and autonomous peer-to-peer manner.