kbus_base.go 5.0 KB

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