kwg.go 4.0 KB

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