| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161 |
- // package kbus_base -- базовая часть шины данных.
- package kbus_base
- import (
- "sync"
- mL0 "gitp78su.ipnodns.ru/svi/kern/v4/lev0"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/core_spec"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/defs"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/defs/stream_name"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/etypes/ebool"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev0/helpers"
- mL1 "gitp78su.ipnodns.ru/svi/kern/v4/lev1"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev1/comp_spec"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kctx"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/kspec"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_ent"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/bus_spec"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus/dict_topic_serve"
- "gitp78su.ipnodns.ru/svi/kern/v4/lev2/lti/kbus/dict_topic_sub"
- )
- var (
- busBaseStreamName = stream_name.NewAStreamName(defs.Txt("bus_base"))
- )
- // KBusBase -- базовая часть шины данных.
- type KBusBase struct {
- KCtx_ kspec.IKernelCtx
- IsWork_ core_spec.ISafeBool
- lCtx comp_spec.ILocalCtx
- log comp_spec.ILogBuf
- dictSub bus_spec.IDictTopicSub
- dictServe bus_spec.IDictTopicServe
- }
- var (
- Bus_ *KBusBase
- block sync.Mutex
- )
- // GetKernelBusBase -- возвращает базовую шину сообщений.
- func GetKernelBusBase() *KBusBase {
- block.Lock()
- defer block.Unlock()
- if Bus_ != nil {
- return Bus_
- }
- kCtx := kctx.GetKernelCtx()
- lCtx := mL1.NewLocalCtx(kCtx.Ctx())
- Bus_ = &KBusBase{
- KCtx_: kCtx,
- IsWork_: mL1.NewSafeBool(false),
- dictSub: dict_topic_sub.NewDictTopicSub(),
- dictServe: dict_topic_serve.NewDictServe(),
- lCtx: lCtx,
- }
- Bus_.log = Bus_.lCtx.Log()
- go Bus_.close()
- go Bus_.run()
- Bus_.IsWork_.Set()
- Bus_.KCtx_.Wg().Add(busBaseStreamName)
- Bus_.KCtx_.Set(defs.Txt("kernBusBase"), Bus_, "base of data bus")
- _ = bus_spec.IKernelBus(Bus_)
- return Bus_
- }
- // Log -- возвращает лог шины.
- func (sf *KBusBase) Log() comp_spec.ILogBuf {
- return sf.log
- }
- func (sf *KBusBase) run() {
- sf.log.Debug(defs.Txt("KBusBase.run()"))
- for {
- break
- }
- }
- // Unsubscribe -- отписывает обработчик от топика.
- func (sf *KBusBase) Unsubscribe(handler bus_spec.IBusHandlerSubscribe) {
- msg := defs.Txt("KBusBase.Unsubscribe(): handler='%v'", handler.Name())
- sf.log.Debug(msg)
- sf.dictSub.Unsubscribe(handler)
- }
- // Subscribe -- подписывает обработчик на топик.
- func (sf *KBusBase) Subscribe(handler bus_spec.IBusHandlerSubscribe) core_spec.IResult[core_spec.EBool] {
- msg := defs.Txt("KBusBase.Subscribe(): handler='%v'", handler.Name())
- sf.log.Debug(msg)
- helpers.Hassert(!sf.IsWork_.Get(), "KBusBase.Subscribe(): handler='%v', bus already closed", []any{handler.Name()})
- sf.dictSub.Subscribe(handler)
- return mL0.NewOk(ebool.NewEBool(true))
- }
- // SendRequest -- отправляет запрос в шину данных.
- func (sf *KBusBase) SendRequest(topic bus_ent.ITopicFix, binReq []byte) mL0.IResult[[]byte] {
- msg := defs.Txt("KBusBase.SendRequest(): topic='%v'", topic.Get())
- sf.log.Debug(msg)
- if !sf.IsWork_.Get() {
- err := defs.Err("KBusBase.SendRequest(): topic='%v', bus already closed", topic.Get())
- sf.log.Err(err)
- return mL0.NewErr[[]byte](err)
- }
- res := sf.dictServe.SendRequest(topic, binReq)
- if res.IsErr() {
- err := defs.Err("KBusBase.SendRequest(): topic='%v', err=\n\t%w", topic.Get(), res.Err())
- sf.log.Err(err)
- return mL0.NewErr[[]byte](err)
- }
- return res
- }
- // RegisterServe -- регистрирует обработчики входящих запросов.
- func (sf *KBusBase) RegisterServe(handler bus_spec.IBusHandlerServe) mL0.IResult[core_spec.EBool] {
- if handler == nil {
- err := defs.Err("KBusBase.RegisterServe(): IBusHandlerServe==nil")
- return mL0.NewErr[core_spec.EBool](err)
- }
- msg := defs.Txt("KBusBase.RegisterServe(): handler='%v'", handler.Name())
- sf.log.Debug(msg)
- res := sf.dictServe.Register(handler)
- if res.IsErr() {
- err := defs.Err("KBusBase.RegisterServe(): handler='%v', err=\n\t%w", handler.Name(), res.Err())
- sf.log.Err(err)
- return mL0.NewErr[core_spec.EBool](err)
- }
- return mL0.NewOk(ebool.NewEBool(true))
- }
- // Publish -- публикует сообщение в шину.
- func (sf *KBusBase) Publish(topic bus_ent.ITopicFix, binMsg []byte) mL0.IResult[core_spec.EBool] {
- msg := defs.Txt("KBusBase.Publish(): topic='%v'", topic)
- sf.log.Debug(msg)
- if !sf.IsWork_.Get() {
- err := defs.Err("KBusBase.Publish(): topic='%v',bus already closed", topic)
- sf.log.Err(err)
- return mL0.NewErr[core_spec.EBool](err)
- }
- // Асинхронный запуск чтения
- go sf.dictSub.Read(topic, binMsg)
- return mL0.NewOk(ebool.NewEBool(true))
- }
- // IsWork -- возвращает признак работы шины.
- func (sf *KBusBase) IsWork() core_spec.EBool {
- return ebool.NewEBool(sf.IsWork_.Get())
- }
- // Ожидает закрытия шины в отдельном потоке.
- func (sf *KBusBase) close() {
- sf.KCtx_.Wait()
- sf.KCtx_.Lock()
- defer sf.KCtx_.Unlock()
- if !sf.IsWork_.Get() {
- return
- }
- sf.IsWork_.Reset()
- sf.KCtx_.Wg().Done(busBaseStreamName)
- sf.log.Debug(defs.Txt("KBusBase.close(): done"))
- }
|