|
|
@@ -2,12 +2,34 @@
|
|
|
package kchan
|
|
|
|
|
|
import (
|
|
|
+ "fmt"
|
|
|
"sync"
|
|
|
|
|
|
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/etypes/ebool"
|
|
|
+ mL0 "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
|
|
|
"gitp78su.ipnodns.ru/svi/kern/v4/lev0/kspec"
|
|
|
+ "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_ kspec.ILocalCtx // Локальный контекст для контроля создателя
|
|
|
+}
|
|
|
+
|
|
|
+// SelfCheck -- проверка корректности правильности параметров умного канала
|
|
|
+func (sf *KChanParam[T]) SelfCheck() {
|
|
|
+ hassert(sf.Limit_ > 0, "KChanParam[T].SelfCheck(): limit=%v, канал должен иметь положительный лимит")
|
|
|
+ hassert(sf.Ctx_ != nil, "KChanParam[T]: ctx==nil")
|
|
|
+}
|
|
|
+
|
|
|
// KChan -- умный потокобезопасный канал
|
|
|
//
|
|
|
// Канал является резиновым, но лимит можно поставить сверху.
|
|
|
@@ -20,23 +42,66 @@ import (
|
|
|
// Кроме того, можно повесить хуки на события записи, чтения и закрытия
|
|
|
// (например, для целей валидации).
|
|
|
type KChan[T any] struct {
|
|
|
- block sync.RWMutex
|
|
|
+ *KChanParam[T]
|
|
|
+ block sync.RWMutex
|
|
|
isClosed kspec.ISafeBool
|
|
|
- lstMsg []T
|
|
|
- onWrite func(T)
|
|
|
- onRead func(T)
|
|
|
- onClose func()
|
|
|
+ lstMsg []T
|
|
|
}
|
|
|
|
|
|
// NewKChan -- создаёт новый безопасный умный канал
|
|
|
-func NewKChan[T any](fnWrite func()T) *KChan[T] {
|
|
|
- sf := &KChan[T]{}
|
|
|
+func NewKChan[T any](param *KChanParam[T]) kspec.IChan[T] {
|
|
|
+ hassert(param != nil, "NewKChan: param==nil")
|
|
|
+ param.SelfCheck()
|
|
|
+ sf := &KChan[T]{
|
|
|
+ KChanParam: param,
|
|
|
+ isClosed: safe_bool.NewSafeBool(true),
|
|
|
+ }
|
|
|
+ go sf.close()
|
|
|
return sf
|
|
|
}
|
|
|
|
|
|
-// Close -- закрывает канал
|
|
|
-func (sf *KChan[T]) Close() error {
|
|
|
- if sf.block.TryLock() == false {
|
|
|
- return nil
|
|
|
+// Read -- возвращает первый элемент очереди (если есть)
|
|
|
+func (sf *KChan[T]) Read() kspec.IResult[T] {
|
|
|
+ sf.block.Lock()
|
|
|
+ defer sf.block.Unlock()
|
|
|
+ if len(sf.lstMsg) == 0 {
|
|
|
+ err := fmt.Errorf("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() int {
|
|
|
+ sf.block.RLock()
|
|
|
+ defer sf.block.RUnlock()
|
|
|
+ return len(sf.lstMsg)
|
|
|
+}
|
|
|
+
|
|
|
+// Len -- возвращает количество элементов в канале
|
|
|
+func (sf *KChan[T]) Len() int {
|
|
|
+ sf.block.RLock()
|
|
|
+ defer sf.block.RUnlock()
|
|
|
+ return 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()
|
|
|
+}
|