inprocess.go 2.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106
  1. package broker
  2. import (
  3. "context"
  4. "fmt"
  5. "runtime"
  6. "sync"
  7. "git.nspix.com/golang/micro/log"
  8. )
  9. type InProcBus struct {
  10. ctx context.Context
  11. once sync.Once
  12. eventChan chan *Event
  13. subscribers sync.Map
  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. vs := make([]Subscriber, 0)
  30. bus.subscribers.Range(func(key, value interface{}) bool {
  31. vs = append(vs, value.(Subscriber))
  32. return true
  33. })
  34. return vs
  35. }
  36. // Publish 发布一个事件
  37. func (bus *InProcBus) Publish(e *Event) {
  38. select {
  39. case bus.eventChan <- e:
  40. default:
  41. log.Warnf("event queue is full, event %s@%s has been discarded", e.Name, e.Namespace)
  42. }
  43. return
  44. }
  45. // Dispatch 分配事件,关心返回值
  46. func (bus *InProcBus) Dispatch(e *Event) (err error) {
  47. return bus.DispatchCtx(bus.ctx, e)
  48. }
  49. //DispatchCtx 分配一个事件
  50. func (bus *InProcBus) DispatchCtx(ctx context.Context, e *Event) (err error) {
  51. bus.subscribers.Range(func(key, value interface{}) bool {
  52. sub := value.(Subscriber)
  53. if sub.HasTopic(e.Name) {
  54. if err = sub.Process(ctx, e); err != nil {
  55. log.Warnf("subscriber %s process event %s error: %s", sub.ID(), e.Name, err.Error())
  56. }
  57. }
  58. return true
  59. })
  60. releaseEvent(e)
  61. return
  62. }
  63. //Subscribe 订阅一个事件
  64. func (bus *InProcBus) Subscribe(sub Subscriber) (err error) {
  65. if _, ok := bus.subscribers.Load(sub.ID()); ok {
  66. err = fmt.Errorf("subscriber %s already exists", sub.ID())
  67. return
  68. }
  69. if err = sub.OnAttach(); err != nil {
  70. return
  71. }
  72. bus.subscribers.Store(sub.ID(), sub)
  73. return
  74. }
  75. //UnSubscribe 取消一个订阅
  76. func (bus *InProcBus) UnSubscribe(sub Subscriber) (err error) {
  77. bus.subscribers.Delete(sub.ID())
  78. return
  79. }
  80. func (bus *InProcBus) WithContext(ctx context.Context) {
  81. bus.ctx = ctx
  82. }
  83. func NewInPrcBus(ctx context.Context) *InProcBus {
  84. bus := &InProcBus{
  85. ctx: ctx,
  86. eventChan: make(chan *Event, 1024),
  87. }
  88. bus.once.Do(func() {
  89. for i := 0; i < runtime.NumCPU()*2; i++ {
  90. go bus.worker()
  91. }
  92. })
  93. return bus
  94. }