kchan.go 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108
  1. // package kchan -- умный потокобезопасный канал
  2. package kchan
  3. import (
  4. "sync"
  5. "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
  6. "gitp78su.ipnodns.ru/svi/kern/v4/lev0/defs"
  7. "gitp78su.ipnodns.ru/svi/kern/v4/lev1/comp_spec"
  8. "gitp78su.ipnodns.ru/svi/kern/v4/lev1/safe_bool"
  9. )
  10. // KChanParam -- параметры умного канала.
  11. type KChanParam[T any] struct {
  12. Limit_ lev0.Num // лимит на размер канала, не может быть пустым
  13. OnWrite_ func(T) // Обратный вызов при записи в канал
  14. OnRead_ func(T) // Обратный вызов при чтении из канала
  15. OnClose_ func() // Обратный вызов при закрытии канала
  16. OnLimit_ func() // Обратный вызов при достижении лимита
  17. Ctx_ comp_spec.ILocalCtx // Локальный контекст для контроля создателя
  18. }
  19. var msg1 = lev0.ATxt("KChanParam[T]: ctx==nil")
  20. // SelfCheck -- проверка корректности правильности параметров умного канала.
  21. func (sf *KChanParam[T]) SelfCheck() {
  22. lev0.If(sf.Limit_ <= 0).Hassert("KChanParam[T].SelfCheck(): limit=%v, канал должен иметь положительный лимит",
  23. sf.Limit_)
  24. lev0.If(sf.Ctx_ == nil).Hassert(msg1)
  25. }
  26. // KChan -- умный потокобезопасный канал
  27. //
  28. // Канал является резиновым, но лимит можно поставить сверху.
  29. // При необходимости лимит ёмкости можно поднять, но не снизить.
  30. // Умный канал можно безопасно закрывать многократно.
  31. // Умный канал можно использовать неблокирующим способом.
  32. // Умный канал точно знает своё состояние (длина, закрыт и т.п.)
  33. // При передаче инстанса канала -- из него можно только читать.
  34. // Умный канал строго типизирован и его тип видно из сигнатуры.
  35. // Кроме того, можно повесить хуки на события записи, чтения и закрытия
  36. // (например, для целей валидации).
  37. type KChan[T any] struct {
  38. *KChanParam[T]
  39. block sync.RWMutex
  40. isClosed lev0.ISafeBool
  41. lstMsg []T
  42. }
  43. var msg2 = lev0.ATxt("NewKChan: param==nil")
  44. // NewKChan -- создаёт новый безопасный умный канал.
  45. func NewKChan[T any](param *KChanParam[T]) lev0.IChan[T] {
  46. lev0.If(param == nil).Hassert(msg2)
  47. param.SelfCheck()
  48. sf := &KChan[T]{
  49. KChanParam: param,
  50. isClosed: safe_bool.NewSafeBool(true),
  51. }
  52. go sf.close()
  53. return sf
  54. }
  55. // Read -- возвращает первый элемент очереди (если есть).
  56. func (sf *KChan[T]) Read() *lev0.Result[T] {
  57. sf.block.Lock()
  58. defer sf.block.Unlock()
  59. if len(sf.lstMsg) == 0 {
  60. err := lev0.LetErr("KChan[T].Read(): empty list msg")
  61. return lev0.ResErr[T](err)
  62. }
  63. msg := sf.lstMsg[0]
  64. sf.lstMsg = sf.lstMsg[1:]
  65. return lev0.ResOk(msg)
  66. }
  67. // Limit -- ограничение размера канала.
  68. func (sf *KChan[T]) Limit() defs.Num {
  69. sf.block.RLock()
  70. defer sf.block.RUnlock()
  71. return defs.Num(len(sf.lstMsg))
  72. }
  73. // Len -- возвращает количество элементов в канале.
  74. func (sf *KChan[T]) Len() defs.Num {
  75. sf.block.RLock()
  76. defer sf.block.RUnlock()
  77. return defs.Num(len(sf.lstMsg))
  78. }
  79. // IsClosed -- возвращает признак закрытия канала.
  80. func (sf *KChan[T]) IsClosed() defs.Bool {
  81. return sf.isClosed.Get()
  82. }
  83. // Close -- закрывает канал.
  84. func (sf *KChan[T]) Close() {
  85. sf.block.Lock()
  86. defer sf.block.Unlock()
  87. sf.isClosed.Set()
  88. }
  89. // Ожидает закрытия канала по контексту.
  90. func (sf *KChan[T]) close() {
  91. sf.Ctx_.Wait()
  92. sf.isClosed.Set()
  93. }