dict_topic_sub.go 2.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778
  1. // package dict_topic_sub -- потокобезопасный словарь подписчиков локальной шины.
  2. package dict_topic_sub
  3. import (
  4. "sync"
  5. "gitp78su.ipnodns.ru/svi/kern/v4/d0"
  6. "gitp78su.ipnodns.ru/svi/kern/v4/d0/helpers"
  7. "gitp78su.ipnodns.ru/svi/kern/v4/d2/kern_ctx"
  8. "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_ent"
  9. "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_ent/topic"
  10. "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/bus_spec"
  11. "gitp78su.ipnodns.ru/svi/kern/v4/d2/lti/kbus/dict_sub_hook"
  12. )
  13. type tReadReq struct {
  14. topic *topic.LTopic
  15. binMsg []byte
  16. }
  17. // dictTopicSub -- потокобезопасный словарь подписчиков.
  18. type dictTopicSub struct {
  19. sync.RWMutex
  20. kCtx *kern_ctx.KernCtx
  21. dictTopicHook map[bus_ent.ATopic]bus_spec.IDictSubHook
  22. }
  23. // NewDictTopicSub -- возвращает потокобезопасный словарь подписчиков.
  24. func NewDictTopicSub() *dictTopicSub {
  25. sf := &dictTopicSub{
  26. kCtx: kern_ctx.GetKernCtx(),
  27. dictTopicHook: map[bus_ent.ATopic]bus_spec.IDictSubHook{},
  28. }
  29. return sf
  30. }
  31. // Read -- вызывает обработчики при поступлении события.
  32. func (sf *dictTopicSub) Read(topic *topic.LTopic, binMsg []byte) {
  33. sf.RLock()
  34. defer sf.RUnlock()
  35. msg := &tReadReq{
  36. topic: topic,
  37. binMsg: binMsg,
  38. }
  39. dictHook := sf.dictTopicHook[msg.topic.Get()]
  40. if dictHook == nil {
  41. return
  42. }
  43. dictHook.Read(msg.binMsg)
  44. }
  45. // Subscribe -- подписывает обработчик на топик.
  46. func (sf *dictTopicSub) Subscribe(handler bus_spec.IBusHandlerSubscribe) {
  47. sf.Lock()
  48. defer sf.Unlock()
  49. d0.If(handler == nil).
  50. Hassert("dictTopicSub.Subscribe(): handler==nil")
  51. topic := handler.Topic()
  52. dictSubHook := sf.dictTopicHook[topic.Get()]
  53. if dictSubHook == nil {
  54. dictSubHook := dict_sub_hook.NewDictSubHook()
  55. sf.dictTopicHook[topic.Get()] = dictSubHook
  56. }
  57. dictSubHook.Subscribe(handler)
  58. }
  59. // Unsubscribe -- отписывает обработчик.
  60. func (sf *dictTopicSub) Unsubscribe(handler bus_spec.IBusHandlerSubscribe) {
  61. sf.Lock()
  62. defer sf.Unlock()
  63. helpers.If(handler == nil).Hassert("dictTopicSub.Unsubscribe(): handler==nil", []any{})
  64. topic := handler.Topic()
  65. dictSubHook := sf.dictTopicHook[topic.Get()]
  66. if dictSubHook == nil {
  67. return
  68. }
  69. dictSubHook.Unsubscribe(handler)
  70. }