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