gateway.go 2.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134
  1. package gateway
  2. import (
  3. "bytes"
  4. "context"
  5. "errors"
  6. "io"
  7. "net"
  8. "sync"
  9. )
  10. const (
  11. MinFeatureLength = 3
  12. )
  13. var (
  14. ErrShortFeature = errors.New("short feature")
  15. ErrFeatureExists = errors.New("feature already exists")
  16. )
  17. type (
  18. listener struct {
  19. feature []byte
  20. l *Listener
  21. }
  22. Gateway struct {
  23. listeners []*listener
  24. l net.Listener
  25. ch chan net.Conn
  26. }
  27. )
  28. func (g *Gateway) Attach(feature []byte, l *Listener) (err error) {
  29. //特征量不够大
  30. if len(feature) < MinFeatureLength {
  31. return ErrShortFeature
  32. }
  33. //判断重复
  34. for _, v := range g.listeners {
  35. if bytes.Equal(v.feature, feature) {
  36. return ErrFeatureExists
  37. }
  38. }
  39. g.listeners = append(g.listeners, &listener{
  40. feature: feature,
  41. l: l,
  42. })
  43. return
  44. }
  45. func (g *Gateway) Attaches(features [][]byte, l *Listener) (err error) {
  46. for _, b := range features {
  47. if err = g.Attach(b, l); err != nil {
  48. break
  49. }
  50. }
  51. return
  52. }
  53. func (g *Gateway) Detach(feature []byte) (err error) {
  54. for i, l := range g.listeners {
  55. if bytes.Equal(l.feature, feature) {
  56. g.listeners = append(g.listeners[:i], g.listeners[i+1:]...)
  57. break
  58. }
  59. }
  60. return
  61. }
  62. func (g *Gateway) loop(ctx context.Context) {
  63. var (
  64. n int
  65. ok bool
  66. err error
  67. conn net.Conn
  68. feature = make([]byte, MinFeatureLength)
  69. )
  70. for {
  71. select {
  72. case conn, ok = <-g.ch:
  73. if ok {
  74. if n, err = io.ReadFull(conn, feature); err != nil {
  75. continue
  76. }
  77. for _, l := range g.listeners {
  78. if bytes.Equal(feature[:n], l.feature[:n]) {
  79. l.l.Receive(wrapConn(conn, feature[:n]))
  80. break
  81. }
  82. }
  83. }
  84. case <-ctx.Done():
  85. return
  86. }
  87. }
  88. }
  89. func (g *Gateway) schedule(ctx context.Context) {
  90. for {
  91. conn, err := g.l.Accept()
  92. if err != nil {
  93. return
  94. }
  95. select {
  96. case g.ch <- conn:
  97. case <-ctx.Done():
  98. return
  99. }
  100. }
  101. }
  102. //运行项目
  103. func (g *Gateway) Run(ctx context.Context) {
  104. var wg sync.WaitGroup
  105. wg.Add(2)
  106. go func() {
  107. g.loop(ctx)
  108. wg.Done()
  109. }()
  110. go func() {
  111. g.schedule(ctx)
  112. wg.Done()
  113. }()
  114. wg.Wait()
  115. }
  116. func New(l net.Listener) *Gateway {
  117. return &Gateway{
  118. l: l,
  119. ch: make(chan net.Conn, 10),
  120. listeners: make([]*listener, 0),
  121. }
  122. }