dict_topic_sub.go 2.3 KB

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