inprocess.go 2.3 KB

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