| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778 |
- // package dict_topic_sub -- потокобезопасный словарь подписчиков локальной шины.
- package dict_topic_sub
- import (
- "sync"
- "gitp78su.ipnodns.ru/svi/kern/v4/d0"
- "gitp78su.ipnodns.ru/svi/kern/v4/d0/helpers"
- "gitp78su.ipnodns.ru/svi/kern/v4/d2/kern_ctx"
- "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_ent"
- "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_ent/topic"
- "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_spec"
- "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/kbus/dict_sub_hook"
- )
- type tReadReq struct {
- topic *topic.LTopic
- binMsg []byte
- }
- // dictTopicSub -- потокобезопасный словарь подписчиков.
- type dictTopicSub struct {
- sync.RWMutex
- kCtx *kern_ctx.KernCtx
- dictTopicHook map[bus_ent.ATopic]bus_spec.IDictSubHook
- }
- // NewDictTopicSub -- возвращает потокобезопасный словарь подписчиков.
- func NewDictTopicSub() *dictTopicSub {
- sf := &dictTopicSub{
- kCtx: kern_ctx.GetKernCtx(),
- dictTopicHook: map[bus_ent.ATopic]bus_spec.IDictSubHook{},
- }
- return sf
- }
- // Read -- вызывает обработчики при поступлении события.
- func (sf *dictTopicSub) Read(topic *topic.LTopic, 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()
- d0.If(handler == nil).
- Hassert("dictTopicSub.Subscribe(): handler==nil")
- 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.If(handler == nil).Hassert("dictTopicSub.Unsubscribe(): handler==nil", []any{})
- topic := handler.Topic()
- dictSubHook := sf.dictTopicHook[topic.Get()]
- if dictSubHook == nil {
- return
- }
- dictSubHook.Unsubscribe(handler)
- }
|