// package kchan -- умный потокобезопасный канал package kchan import ( "sync" "gitp78su.ipnodns.ru/svi/kern/v4/lev0" "gitp78su.ipnodns.ru/svi/kern/v4/lev0/defs" "gitp78su.ipnodns.ru/svi/kern/v4/lev1/comp_spec" "gitp78su.ipnodns.ru/svi/kern/v4/lev1/safe_bool" ) // KChanParam -- параметры умного канала. type KChanParam[T any] struct { Limit_ lev0.Num // лимит на размер канала, не может быть пустым OnWrite_ func(T) // Обратный вызов при записи в канал OnRead_ func(T) // Обратный вызов при чтении из канала OnClose_ func() // Обратный вызов при закрытии канала OnLimit_ func() // Обратный вызов при достижении лимита Ctx_ comp_spec.ILocalCtx // Локальный контекст для контроля создателя } var msg1 = lev0.ATxt("KChanParam[T]: ctx==nil") // SelfCheck -- проверка корректности правильности параметров умного канала. func (sf *KChanParam[T]) SelfCheck() { lev0.If(sf.Limit_ <= 0).Hassert("KChanParam[T].SelfCheck(): limit=%v, канал должен иметь положительный лимит", sf.Limit_) lev0.If(sf.Ctx_ == nil).Hassert(msg1) } // KChan -- умный потокобезопасный канал // // Канал является резиновым, но лимит можно поставить сверху. // При необходимости лимит ёмкости можно поднять, но не снизить. // Умный канал можно безопасно закрывать многократно. // Умный канал можно использовать неблокирующим способом. // Умный канал точно знает своё состояние (длина, закрыт и т.п.) // При передаче инстанса канала -- из него можно только читать. // Умный канал строго типизирован и его тип видно из сигнатуры. // Кроме того, можно повесить хуки на события записи, чтения и закрытия // (например, для целей валидации). type KChan[T any] struct { *KChanParam[T] block sync.RWMutex isClosed lev0.ISafeBool lstMsg []T } var msg2 = lev0.ATxt("NewKChan: param==nil") // NewKChan -- создаёт новый безопасный умный канал. func NewKChan[T any](param *KChanParam[T]) lev0.IChan[T] { lev0.If(param == nil).Hassert(msg2) param.SelfCheck() sf := &KChan[T]{ KChanParam: param, isClosed: safe_bool.NewSafeBool(true), } go sf.close() return sf } // Read -- возвращает первый элемент очереди (если есть). func (sf *KChan[T]) Read() *lev0.Result[T] { sf.block.Lock() defer sf.block.Unlock() if len(sf.lstMsg) == 0 { err := lev0.LetErr("KChan[T].Read(): empty list msg") return lev0.ResErr[T](err) } msg := sf.lstMsg[0] sf.lstMsg = sf.lstMsg[1:] return lev0.ResOk(msg) } // Limit -- ограничение размера канала. func (sf *KChan[T]) Limit() defs.Num { sf.block.RLock() defer sf.block.RUnlock() return defs.Num(len(sf.lstMsg)) } // Len -- возвращает количество элементов в канале. func (sf *KChan[T]) Len() defs.Num { sf.block.RLock() defer sf.block.RUnlock() return defs.Num(len(sf.lstMsg)) } // IsClosed -- возвращает признак закрытия канала. func (sf *KChan[T]) IsClosed() defs.Bool { return sf.isClosed.Get() } // Close -- закрывает канал. func (sf *KChan[T]) Close() { sf.block.Lock() defer sf.block.Unlock() sf.isClosed.Set() } // Ожидает закрытия канала по контексту. func (sf *KChan[T]) close() { sf.Ctx_.Wait() sf.isClosed.Set() }