inprocess.go 2.1 KB

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