| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141 |
- // package kwg -- именованный ожидатель потоков ядра.
- //
- // Не позволяет завершиться ядру, если есть хоть один работающий поток
- package kwg
- import (
- "context"
- "sync"
- "time"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kern_ent/stream_name"
- )
- // KernWg -- именованный ожидатель потоков ядра.
- type KernWg struct {
- sync.RWMutex
- ctx context.Context
- dictStream map[*stream_name.AStreamName]lev0.Bool // Словарь имён потоков с признаком работы
- isWork *lev0.MSafeBool
- log *lev1.LogBuf
- }
- var (
- kernWg *KernWg // Глобальный объект
- block sync.Mutex
- )
- // GetKernelWg -- возвращает новый именованный ожидатель потоков ядра.
- func GetKernelWg(ctx context.Context) *KernWg {
- block.Lock()
- defer block.Unlock()
- if kernWg != nil {
- kernWg.log.Debug("GetKernelWg()")
- return kernWg
- }
- lev0.If(ctx == nil).Hassert("GetKernelWg(): ctx==nil")
- paramLogBuf := &lev1.LogBufParam{
- IsTerm_: lev0.MutSafeBool(true),
- Prefix_: lev0.LetTxt("kernelWg"),
- }
- sf := &KernWg{
- ctx: ctx,
- dictStream: map[*stream_name.AStreamName]lev0.Bool{},
- isWork: lev0.MutSafeBool(false),
- log: lev1.NewLogBuf(paramLogBuf),
- }
- sf.log.Debug("GetKernelWg(): run")
- go sf.close()
- sf.isWork.Set()
- kernWg = sf
- return kernWg
- }
- // Log -- возвращает лог ожидателя потоков.
- func (sf *KernWg) Log() *lev1.LogBuf {
- return sf.log
- }
- // Len -- возвращает размер списка ожидания потоков.
- func (sf *KernWg) Len() lev0.Num {
- sf.RLock()
- defer sf.RUnlock()
- return lev0.Num(len(sf.dictStream))
- }
- // IsWork -- возвращает признак работы ядра.
- func (sf *KernWg) IsWork() lev0.Bool {
- return sf.isWork.Get()
- }
- // IsStop -- возвращает признак остановки ядра.
- func (sf *KernWg) IsStop() lev0.Bool {
- return !sf.isWork.Get()
- }
- // List -- возвращает список имён потоков на ожидании.
- func (sf *KernWg) List() []*stream_name.AStreamName {
- sf.RLock()
- defer sf.RUnlock()
- lst := make([]*stream_name.AStreamName, 0, len(sf.dictStream))
- for name := range sf.dictStream {
- lst = append(lst, name)
- }
- return lst
- }
- // Done -- удаляет поток из ожидания.
- func (sf *KernWg) Done(name *stream_name.AStreamName) {
- sf.Lock()
- defer sf.Unlock()
- delete(sf.dictStream, name)
- sf.log.Debug("Done(): stream(%v) done", name)
- }
- // Wait -- блокирующий вызов; возвращает управление, только когда все потоки завершили работу.
- func (sf *KernWg) Wait() {
- for {
- time.Sleep(time.Millisecond * 1)
- if !sf.isWork.Get() {
- break
- }
- }
- sf.log.Debug("Wait(): done")
- }
- // Add -- добавляет поток в ожидание.
- func (sf *KernWg) Add(name *stream_name.AStreamName) {
- sf.Lock()
- defer sf.Unlock()
- sf.log.Debug("Add(): stream('%v') done", name)
- lev0.If(!sf.isWork.Get()).Hassert("Add(): stream=%v, work end", []any{name})
- lev0.If(name.Get() == "").Hassert("Add(): name stream is empty")
- _, isOk := sf.dictStream[name]
- lev0.If(lev0.Bool(isOk)).Hassert("Add(): stream '%v' already exists", name)
- sf.dictStream[name] = true
- }
- // Ожидает окончания работы ожидателя групп.
- func (sf *KernWg) close() {
- <-sf.ctx.Done()
- fnDone := func() bool {
- sf.Lock()
- defer sf.Unlock()
- return len(sf.dictStream) == 0
- }
- for {
- time.Sleep(time.Millisecond * 1)
- if fnDone() {
- break
- }
- }
- sf.Lock()
- defer sf.Unlock()
- if !sf.isWork.Get() {
- return
- }
- sf.isWork.Reset()
- sf.log.Debug("close(): end")
- }
|