dict_sub_hook.go 2.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566
  1. // package dict_sub_hook -- словарь потребителей топика по подписке.
  2. package dict_sub_hook
  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. )
  11. // dictSubHook -- словарь потребителей топика по подписке.
  12. type dictSubHook struct {
  13. kCtx kspec.IKernelCtx
  14. dict map[bus_ent.AHandlerName]struct{} // В качестве ключа -- URL веб-хука
  15. block sync.RWMutex
  16. }
  17. // NewDictSubHook -- возвращает новый словарь веб-хуков одного топика.
  18. func NewDictSubHook() bus_spec.IDictSubHook {
  19. sf := &dictSubHook{
  20. kCtx: kctx.GetKernelCtx(),
  21. dict: map[bus_ent.AHandlerName]struct{}{},
  22. }
  23. return sf
  24. }
  25. // Unsubscribe -- удаляет из словаря подписки обработчик.
  26. func (sf *dictSubHook) Unsubscribe(handler bus_spec.IBusHandlerSubscribe) {
  27. sf.block.Lock()
  28. defer sf.block.Unlock()
  29. if handler == nil {
  30. sf.kCtx.Log().Err("dictSubHook.Unsubscribe(): handler==nil")
  31. return
  32. }
  33. handlerName := handler.Name()
  34. delete(sf.dict, handlerName.Get())
  35. sf.kCtx.Del(handlerName.String())
  36. }
  37. // Subscribe -- добавляет в словарь подписки новый обработчик.
  38. func (sf *dictSubHook) Subscribe(handler bus_spec.IBusHandlerSubscribe) {
  39. sf.block.Lock()
  40. defer sf.block.Unlock()
  41. helpers.Hassert(handler != nil, "dictSubHook.Subscribe(): handler==nil", []any{})
  42. handlerName := handler.Name()
  43. sf.dict[handlerName.Get()] = struct{}{}
  44. sf.kCtx.Set(handlerName.String(), handler, "subscribe handler")
  45. }
  46. // Read -- вызывает все обработчики словаря подписок.
  47. func (sf *dictSubHook) Read(binMsg []byte) {
  48. sf.block.RLock()
  49. defer sf.block.RUnlock()
  50. for handlerName := range sf.dict {
  51. optHand := sf.kCtx.Get(string(handlerName))
  52. if optHand.IsNone() {
  53. sf.kCtx.Del(string(handlerName))
  54. continue
  55. }
  56. handler := optHand.Some().Val().(bus_spec.IBusHandlerSubscribe)
  57. go handler.FnBack(binMsg)
  58. }
  59. }