kwg.go 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146
  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"
  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 core_spec.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(defs.Txt("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_: defs.Txt("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(defs.Txt("GetKernelWg(): run"))
  52. go sf.close()
  53. sf.isWork.Set()
  54. kernWg = sf
  55. return kernWg
  56. }
  57. // Log -- возвращает лог ожидателя потоков.
  58. func (sf *kernelWg) Log() comp_spec.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. msg := defs.Txt("Done(): stream(%v) done", name)
  88. sf.log.Debug(msg)
  89. }
  90. // Wait -- блокирующий вызов; возвращает управление, только когда все потоки завершили работу.
  91. func (sf *kernelWg) Wait() {
  92. for {
  93. time.Sleep(time.Millisecond * 1)
  94. if !sf.isWork.Get() {
  95. break
  96. }
  97. }
  98. sf.log.Debug(defs.Txt("Wait(): done"))
  99. }
  100. // Add -- добавляет поток в ожидание.
  101. func (sf *kernelWg) Add(name *stream_name.AStreamName) {
  102. sf.Lock()
  103. defer sf.Unlock()
  104. msg := defs.Txt("Add(): stream('%v') done", name)
  105. sf.log.Debug(msg)
  106. mKh.Hassert(sf.isWork.Get(), "Add(): stream=%v, work end", []any{name})
  107. mKh.Hassert(name.Name().Get() != "", "Add(): name stream is empty")
  108. _, isOk := sf.dictStream[name]
  109. mKh.Hassert(!isOk, "Add(): stream '%v' already exists", name)
  110. sf.dictStream[name] = true
  111. }
  112. // Ожидает окончания работы ожидателя групп.
  113. func (sf *kernelWg) close() {
  114. <-sf.ctx.Done()
  115. fnDone := func() bool {
  116. sf.Lock()
  117. defer sf.Unlock()
  118. return len(sf.dictStream) == 0
  119. }
  120. for {
  121. time.Sleep(time.Millisecond * 1)
  122. if fnDone() {
  123. break
  124. }
  125. }
  126. sf.Lock()
  127. defer sf.Unlock()
  128. if !sf.isWork.Get() {
  129. return
  130. }
  131. sf.isWork.Reset()
  132. sf.log.Debug(defs.Txt("close(): end"))
  133. }