safe_chan.go 4.0 KB

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