completed_queue.go 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108
  1. // Copyright 2019 Yunion
  2. //
  3. // Licensed under the Apache License, Version 2.0 (the "License");
  4. // you may not use this file except in compliance with the License.
  5. // You may obtain a copy of the License at
  6. //
  7. // http://www.apache.org/licenses/LICENSE-2.0
  8. //
  9. // Unless required by applicable law or agreed to in writing, software
  10. // distributed under the License is distributed on an "AS IS" BASIS,
  11. // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  12. // See the License for the specific language governing permissions and
  13. // limitations under the License.
  14. package manager
  15. import (
  16. "time"
  17. "yunion.io/x/log"
  18. "yunion.io/x/pkg/utils"
  19. "yunion.io/x/onecloud/pkg/scheduler/api"
  20. o "yunion.io/x/onecloud/pkg/scheduler/options"
  21. )
  22. type CompletedManager struct {
  23. completedChannel chan *api.CompletedNotifyArgs
  24. stopCh <-chan struct{}
  25. }
  26. func NewCompletedManager(stopCh <-chan struct{}) *CompletedManager {
  27. return &CompletedManager{
  28. completedChannel: make(chan *api.CompletedNotifyArgs, o.Options.CompletedQueueMaxLength),
  29. stopCh: stopCh,
  30. }
  31. }
  32. func (c *CompletedManager) Add(completedNotifyArgs *api.CompletedNotifyArgs) {
  33. c.completedChannel <- completedNotifyArgs
  34. }
  35. func (c *CompletedManager) Run() {
  36. t := time.Tick(utils.ToDuration(o.Options.CompletedQueueConsumptionPeriod))
  37. removeSession := func() {
  38. //completedNotifyArgs := <-c.completedChannel
  39. //pool, err := schedManager.ReservedPoolManager.SearchReservedPoolBySessionID(completedNotifyArgs.SessionID)
  40. //if err != nil {
  41. //log.Errorln(err)
  42. //return
  43. //}
  44. //sessionItem := pool.GetSessionItem(completedNotifyArgs.SessionID)
  45. //if sessionItem == nil {
  46. //log.Errorln(fmt.Errorf("session %v not found\n", completedNotifyArgs.SessionID))
  47. //return
  48. //}
  49. //candidateIds := sessionItem.AllCandidateIDs()
  50. // load candidates
  51. //if len(candidateIds) > 0 {
  52. //schedManager.CandidateManager.Reload(pool.Name, candidateIds)
  53. //}
  54. // remove session
  55. //pool.RemoveSession(completedNotifyArgs.SessionID)
  56. }
  57. reloadAndRemoveSessions := func() {
  58. completedRequestNumber := len(c.completedChannel)
  59. // If the completedRequestNumber then return right now.
  60. if completedRequestNumber <= 0 {
  61. return
  62. }
  63. wg := &utils.WaitGroupWrapper{}
  64. for i := 0; i < completedRequestNumber; i++ {
  65. wg.Wrap(removeSession)
  66. }
  67. if ok := utils.WaitTimeOut(wg, time.Duration(completedRequestNumber)*utils.ToDuration(o.Options.CompletedQueueConsumptionTimeout)); !ok {
  68. log.Errorln("time out reload data in completed when remove sessions.")
  69. }
  70. }
  71. // Watching the completed sessions.
  72. for {
  73. select {
  74. case <-t:
  75. reloadAndRemoveSessions()
  76. case <-c.stopCh:
  77. // update all the sessions before return.
  78. reloadAndRemoveSessions()
  79. close(c.completedChannel)
  80. c.completedChannel = nil
  81. log.Errorln("completed manager EXIT!")
  82. return
  83. default:
  84. // if sessions' number is bigger then 10 then reload and remove.
  85. if len(c.completedChannel) >= o.Options.CompletedQueueDealLength {
  86. reloadAndRemoveSessions()
  87. } else {
  88. time.Sleep(10 * time.Second)
  89. }
  90. }
  91. }
  92. }