inprocess.go 2.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114
  1. package broker
  2. import (
  3. "context"
  4. "fmt"
  5. "git.nspix.com/golang/micro/log"
  6. "runtime"
  7. "sync"
  8. )
  9. type InProcBus struct {
  10. ctx context.Context
  11. once sync.Once
  12. eventChan chan *Event
  13. subscriberLocker sync.RWMutex
  14. subscribers map[string]Subscriber
  15. }
  16. func (bus *InProcBus) worker() {
  17. for {
  18. select {
  19. case e, ok := <-bus.eventChan:
  20. if ok {
  21. _ = bus.Dispatch(e)
  22. }
  23. case <-bus.ctx.Done():
  24. return
  25. }
  26. }
  27. }
  28. //Subscribers 获取所有订阅者的快照
  29. func (bus *InProcBus) Subscribers() []Subscriber {
  30. bus.subscriberLocker.RLock()
  31. defer bus.subscriberLocker.RUnlock()
  32. vs := make([]Subscriber, len(bus.subscribers))
  33. i := 0
  34. for _, v := range bus.subscribers {
  35. vs[i] = v
  36. i++
  37. }
  38. return vs
  39. }
  40. // Publish 发布一个事件
  41. func (bus *InProcBus) Publish(e *Event) {
  42. bus.once.Do(func() {
  43. for i := 0; i < runtime.NumCPU(); i++ {
  44. go bus.worker()
  45. }
  46. })
  47. select {
  48. case bus.eventChan <- e:
  49. default:
  50. }
  51. return
  52. }
  53. // Dispatch 分配事件,关心返回值
  54. func (bus *InProcBus) Dispatch(e *Event) (err error) {
  55. return bus.DispatchCtx(context.Background(), e)
  56. }
  57. //DispatchCtx 分配一个事件
  58. func (bus *InProcBus) DispatchCtx(ctx context.Context, e *Event) (err error) {
  59. bus.subscriberLocker.RLock()
  60. defer bus.subscriberLocker.RUnlock()
  61. for _, sub := range bus.subscribers {
  62. if sub.HasTopic(e.Name) {
  63. if err = sub.Process(ctx, e); err != nil {
  64. log.Warnf("subscriber %s process event %s error: %s", sub.ID(), e.Name, err.Error())
  65. }
  66. }
  67. }
  68. releaseEvent(e)
  69. return
  70. }
  71. //Subscribe 订阅一个事件
  72. func (bus *InProcBus) Subscribe(sub Subscriber) (err error) {
  73. bus.subscriberLocker.Lock()
  74. defer bus.subscriberLocker.Unlock()
  75. if _, ok := bus.subscribers[sub.ID()]; ok {
  76. err = fmt.Errorf("subscriber %s already exists", sub.ID())
  77. return
  78. }
  79. if err = sub.OnAttach(); err != nil {
  80. return
  81. }
  82. bus.subscribers[sub.ID()] = sub
  83. return
  84. }
  85. //UnSubscribe 取消一个订阅
  86. func (bus *InProcBus) UnSubscribe(sub Subscriber) (err error) {
  87. bus.subscriberLocker.Lock()
  88. defer bus.subscriberLocker.Unlock()
  89. if instance, ok := bus.subscribers[sub.ID()]; !ok {
  90. err = fmt.Errorf("subscriber %s not exists", sub.ID())
  91. return
  92. } else {
  93. delete(bus.subscribers, sub.ID())
  94. err = instance.OnDetach()
  95. }
  96. return
  97. }
  98. func NewInPrcBus(ctx context.Context) *InProcBus {
  99. return &InProcBus{
  100. ctx: ctx,
  101. eventChan: make(chan *Event, 100),
  102. subscribers: make(map[string]Subscriber),
  103. }
  104. }