kbus_base.go 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161
  1. // package kbus_base -- базовая часть шины данных.
  2. package kbus_base
  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/defs/stream_name"
  9. "gitp78su.ipnodns.ru/svi/kern/v4/lev0/etypes/ebool"
  10. "gitp78su.ipnodns.ru/svi/kern/v4/lev0/helpers"
  11. mL1 "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
  12. "gitp78su.ipnodns.ru/svi/kern/v4/lev1/comp_spec"
  13. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kctx"
  14. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kspec"
  15. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_ent"
  16. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_spec"
  17. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus/dict_topic_serve"
  18. "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus/dict_topic_sub"
  19. )
  20. var (
  21. busBaseStreamName = stream_name.NewAStreamName(defs.Txt("bus_base"))
  22. )
  23. // KBusBase -- базовая часть шины данных.
  24. type KBusBase struct {
  25. KCtx_ kspec.IKernelCtx
  26. IsWork_ core_spec.ISafeBool
  27. lCtx comp_spec.ILocalCtx
  28. log comp_spec.ILogBuf
  29. dictSub bus_spec.IDictTopicSub
  30. dictServe bus_spec.IDictTopicServe
  31. }
  32. var (
  33. Bus_ *KBusBase
  34. block sync.Mutex
  35. )
  36. // GetKernelBusBase -- возвращает базовую шину сообщений.
  37. func GetKernelBusBase() *KBusBase {
  38. block.Lock()
  39. defer block.Unlock()
  40. if Bus_ != nil {
  41. return Bus_
  42. }
  43. kCtx := kctx.GetKernelCtx()
  44. lCtx := mL1.NewLocalCtx(kCtx.Ctx())
  45. Bus_ = &KBusBase{
  46. KCtx_: kCtx,
  47. IsWork_: mL1.NewSafeBool(false),
  48. dictSub: dict_topic_sub.NewDictTopicSub(),
  49. dictServe: dict_topic_serve.NewDictServe(),
  50. lCtx: lCtx,
  51. }
  52. Bus_.log = Bus_.lCtx.Log()
  53. go Bus_.close()
  54. go Bus_.run()
  55. Bus_.IsWork_.Set()
  56. Bus_.KCtx_.Wg().Add(busBaseStreamName)
  57. Bus_.KCtx_.Set(defs.Txt("kernBusBase"), Bus_, "base of data bus")
  58. _ = bus_spec.IKernelBus(Bus_)
  59. return Bus_
  60. }
  61. // Log -- возвращает лог шины.
  62. func (sf *KBusBase) Log() comp_spec.ILogBuf {
  63. return sf.log
  64. }
  65. func (sf *KBusBase) run() {
  66. sf.log.Debug(defs.Txt("KBusBase.run()"))
  67. for {
  68. break
  69. }
  70. }
  71. // Unsubscribe -- отписывает обработчик от топика.
  72. func (sf *KBusBase) Unsubscribe(handler bus_spec.IBusHandlerSubscribe) {
  73. msg := defs.Txt("KBusBase.Unsubscribe(): handler='%v'", handler.Name())
  74. sf.log.Debug(msg)
  75. sf.dictSub.Unsubscribe(handler)
  76. }
  77. // Subscribe -- подписывает обработчик на топик.
  78. func (sf *KBusBase) Subscribe(handler bus_spec.IBusHandlerSubscribe) core_spec.IResult[core_spec.EBool] {
  79. msg := defs.Txt("KBusBase.Subscribe(): handler='%v'", handler.Name())
  80. sf.log.Debug(msg)
  81. helpers.Hassert(!sf.IsWork_.Get(), "KBusBase.Subscribe(): handler='%v', bus already closed", []any{handler.Name()})
  82. sf.dictSub.Subscribe(handler)
  83. return mL0.NewOk(ebool.NewEBool(true))
  84. }
  85. // SendRequest -- отправляет запрос в шину данных.
  86. func (sf *KBusBase) SendRequest(topic bus_ent.ITopicFix, binReq []byte) mL0.IResult[[]byte] {
  87. msg := defs.Txt("KBusBase.SendRequest(): topic='%v'", topic.Get())
  88. sf.log.Debug(msg)
  89. if !sf.IsWork_.Get() {
  90. err := defs.Err("KBusBase.SendRequest(): topic='%v', bus already closed", topic.Get())
  91. sf.log.Err(err)
  92. return mL0.NewErr[[]byte](err)
  93. }
  94. res := sf.dictServe.SendRequest(topic, binReq)
  95. if res.IsErr() {
  96. err := defs.Err("KBusBase.SendRequest(): topic='%v', err=\n\t%w", topic.Get(), res.Err())
  97. sf.log.Err(err)
  98. return mL0.NewErr[[]byte](err)
  99. }
  100. return res
  101. }
  102. // RegisterServe -- регистрирует обработчики входящих запросов.
  103. func (sf *KBusBase) RegisterServe(handler bus_spec.IBusHandlerServe) mL0.IResult[core_spec.EBool] {
  104. if handler == nil {
  105. err := defs.Err("KBusBase.RegisterServe(): IBusHandlerServe==nil")
  106. return mL0.NewErr[core_spec.EBool](err)
  107. }
  108. msg := defs.Txt("KBusBase.RegisterServe(): handler='%v'", handler.Name())
  109. sf.log.Debug(msg)
  110. res := sf.dictServe.Register(handler)
  111. if res.IsErr() {
  112. err := defs.Err("KBusBase.RegisterServe(): handler='%v', err=\n\t%w", handler.Name(), res.Err())
  113. sf.log.Err(err)
  114. return mL0.NewErr[core_spec.EBool](err)
  115. }
  116. return mL0.NewOk(ebool.NewEBool(true))
  117. }
  118. // Publish -- публикует сообщение в шину.
  119. func (sf *KBusBase) Publish(topic bus_ent.ITopicFix, binMsg []byte) mL0.IResult[core_spec.EBool] {
  120. msg := defs.Txt("KBusBase.Publish(): topic='%v'", topic)
  121. sf.log.Debug(msg)
  122. if !sf.IsWork_.Get() {
  123. err := defs.Err("KBusBase.Publish(): topic='%v',bus already closed", topic)
  124. sf.log.Err(err)
  125. return mL0.NewErr[core_spec.EBool](err)
  126. }
  127. // Асинхронный запуск чтения
  128. go sf.dictSub.Read(topic, binMsg)
  129. return mL0.NewOk(ebool.NewEBool(true))
  130. }
  131. // IsWork -- возвращает признак работы шины.
  132. func (sf *KBusBase) IsWork() core_spec.EBool {
  133. return ebool.NewEBool(sf.IsWork_.Get())
  134. }
  135. // Ожидает закрытия шины в отдельном потоке.
  136. func (sf *KBusBase) close() {
  137. sf.KCtx_.Wait()
  138. sf.KCtx_.Lock()
  139. defer sf.KCtx_.Unlock()
  140. if !sf.IsWork_.Get() {
  141. return
  142. }
  143. sf.IsWork_.Reset()
  144. sf.KCtx_.Wg().Done(busBaseStreamName)
  145. sf.log.Debug(defs.Txt("KBusBase.close(): done"))
  146. }