kwg.go 4.0 KB

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