gateway.go 2.3 KB

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