// package dict_topic_sub -- потокобезопасный словарь подписчиков локальной шины. package dict_topic_sub import ( "sync" "gitp78su.ipnodns.ru/svi/kern/v4/lev0/helpers" "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_sub_hook" ) type tReadReq struct { topic bus_ent.ITopicFix binMsg []byte } // dictTopicSub -- потокобезопасный словарь подписчиков. type dictTopicSub struct { sync.RWMutex kCtx kspec.IKernelCtx dictTopicHook map[bus_ent.ATopic]bus_spec.IDictSubHook } // NewDictTopicSub -- возвращает потокобезопасный словарь подписчиков. func NewDictTopicSub() *dictTopicSub { sf := &dictTopicSub{ kCtx: kctx.GetKernelCtx(), dictTopicHook: map[bus_ent.ATopic]bus_spec.IDictSubHook{}, } return sf } // Read -- вызывает обработчики при поступлении события. func (sf *dictTopicSub) Read(topic bus_ent.ITopicFix, binMsg []byte) { sf.RLock() defer sf.RUnlock() msg := &tReadReq{ topic: topic, binMsg: binMsg, } dictHook := sf.dictTopicHook[msg.topic.Get()] if dictHook == nil { return } dictHook.Read(msg.binMsg) } // Subscribe -- подписывает обработчик на топик. func (sf *dictTopicSub) Subscribe(handler bus_spec.IBusHandlerSubscribe) { sf.Lock() defer sf.Unlock() helpers.Hassert(handler != nil, "dictTopicSub.Subscribe(): handler==nil", []any{}) topic := handler.Topic() dictSubHook := sf.dictTopicHook[topic.Get()] if dictSubHook == nil { dictSubHook := dict_sub_hook.NewDictSubHook() sf.dictTopicHook[topic.Get()] = dictSubHook } dictSubHook.Subscribe(handler) } // Unsubscribe -- отписывает обработчик. func (sf *dictTopicSub) Unsubscribe(handler bus_spec.IBusHandlerSubscribe) { sf.Lock() defer sf.Unlock() helpers.Hassert(handler != nil, "dictTopicSub.Unsubscribe(): handler==nil", []any{}) topic := handler.Topic() dictSubHook := sf.dictTopicHook[topic.Get()] if dictSubHook == nil { return } dictSubHook.Unsubscribe(handler) }