kwg.go 3.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141
  1. // package kwg -- именованный ожидатель потоков ядра.
  2. //
  3. // Не позволяет завершиться ядру, если есть хоть один работающий поток
  4. package kwg
  5. import (
  6. "context"
  7. "sync"
  8. "time"
  9. "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
  10. "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
  11. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kern_ent/stream_name"
  12. )
  13. // KernWg -- именованный ожидатель потоков ядра.
  14. type KernWg struct {
  15. sync.RWMutex
  16. ctx context.Context
  17. dictStream map[*stream_name.AStreamName]lev0.Bool // Словарь имён потоков с признаком работы
  18. isWork *lev0.MSafeBool
  19. log *lev1.LogBuf
  20. }
  21. var (
  22. kernWg *KernWg // Глобальный объект
  23. block sync.Mutex
  24. )
  25. // GetKernelWg -- возвращает новый именованный ожидатель потоков ядра.
  26. func GetKernelWg(ctx context.Context) *KernWg {
  27. block.Lock()
  28. defer block.Unlock()
  29. if kernWg != nil {
  30. kernWg.log.Debug("GetKernelWg()")
  31. return kernWg
  32. }
  33. lev0.If(ctx == nil).Hassert("GetKernelWg(): ctx==nil")
  34. paramLogBuf := &lev1.LogBufParam{
  35. IsTerm_: lev0.MutSafeBool(true),
  36. Prefix_: lev0.LetTxt("kernelWg"),
  37. }
  38. sf := &KernWg{
  39. ctx: ctx,
  40. dictStream: map[*stream_name.AStreamName]lev0.Bool{},
  41. isWork: lev0.MutSafeBool(false),
  42. log: lev1.NewLogBuf(paramLogBuf),
  43. }
  44. sf.log.Debug("GetKernelWg(): run")
  45. go sf.close()
  46. sf.isWork.Set()
  47. kernWg = sf
  48. return kernWg
  49. }
  50. // Log -- возвращает лог ожидателя потоков.
  51. func (sf *KernWg) Log() *lev1.LogBuf {
  52. return sf.log
  53. }
  54. // Len -- возвращает размер списка ожидания потоков.
  55. func (sf *KernWg) Len() lev0.Num {
  56. sf.RLock()
  57. defer sf.RUnlock()
  58. return lev0.Num(len(sf.dictStream))
  59. }
  60. // IsWork -- возвращает признак работы ядра.
  61. func (sf *KernWg) IsWork() lev0.Bool {
  62. return sf.isWork.Get()
  63. }
  64. // IsStop -- возвращает признак остановки ядра.
  65. func (sf *KernWg) IsStop() lev0.Bool {
  66. return !sf.isWork.Get()
  67. }
  68. // List -- возвращает список имён потоков на ожидании.
  69. func (sf *KernWg) List() []*stream_name.AStreamName {
  70. sf.RLock()
  71. defer sf.RUnlock()
  72. lst := make([]*stream_name.AStreamName, 0, len(sf.dictStream))
  73. for name := range sf.dictStream {
  74. lst = append(lst, name)
  75. }
  76. return lst
  77. }
  78. // Done -- удаляет поток из ожидания.
  79. func (sf *KernWg) Done(name *stream_name.AStreamName) {
  80. sf.Lock()
  81. defer sf.Unlock()
  82. delete(sf.dictStream, name)
  83. sf.log.Debug("Done(): stream(%v) done", name)
  84. }
  85. // Wait -- блокирующий вызов; возвращает управление, только когда все потоки завершили работу.
  86. func (sf *KernWg) Wait() {
  87. for {
  88. time.Sleep(time.Millisecond * 1)
  89. if !sf.isWork.Get() {
  90. break
  91. }
  92. }
  93. sf.log.Debug("Wait(): done")
  94. }
  95. // Add -- добавляет поток в ожидание.
  96. func (sf *KernWg) Add(name *stream_name.AStreamName) {
  97. sf.Lock()
  98. defer sf.Unlock()
  99. sf.log.Debug("Add(): stream('%v') done", name)
  100. lev0.If(!sf.isWork.Get()).Hassert("Add(): stream=%v, work end", []any{name})
  101. lev0.If(name.Get() == "").Hassert("Add(): name stream is empty")
  102. _, isOk := sf.dictStream[name]
  103. lev0.If(lev0.Bool(isOk)).Hassert("Add(): stream '%v' already exists", name)
  104. sf.dictStream[name] = true
  105. }
  106. // Ожидает окончания работы ожидателя групп.
  107. func (sf *KernWg) close() {
  108. <-sf.ctx.Done()
  109. fnDone := func() bool {
  110. sf.Lock()
  111. defer sf.Unlock()
  112. return len(sf.dictStream) == 0
  113. }
  114. for {
  115. time.Sleep(time.Millisecond * 1)
  116. if fnDone() {
  117. break
  118. }
  119. }
  120. sf.Lock()
  121. defer sf.Unlock()
  122. if !sf.isWork.Get() {
  123. return
  124. }
  125. sf.isWork.Reset()
  126. sf.log.Debug("close(): end")
  127. }