| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113 |
- // 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()
- }
|