kchan.go 4.1 KB

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