ws_log_distr.go 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. package wslogdist
  2. import (
  3. "bytes"
  4. "fmt"
  5. "git.swzry.com/zry/gorillaws2netconn"
  6. "github.com/GUAIK-ORG/go-snowflake/snowflake"
  7. "github.com/gin-gonic/gin"
  8. "github.com/gorilla/websocket"
  9. "net/http"
  10. "sync"
  11. "time"
  12. )
  13. type WsSessions struct {
  14. Conn *gorillaws2netconn.NetConn4Gorilla
  15. BeginTime time.Time
  16. }
  17. func NewWsSessions(conn *gorillaws2netconn.NetConn4Gorilla) *WsSessions {
  18. s := &WsSessions{
  19. Conn: conn,
  20. BeginTime: time.Now(),
  21. }
  22. return s
  23. }
  24. type WsLogDistributor struct {
  25. wsList map[string]*WsSessions
  26. lock sync.RWMutex
  27. upgrader *websocket.Upgrader
  28. startLogBuffer *bytes.Buffer
  29. startLogBufferSize int
  30. snowFlaker *snowflake.Snowflake
  31. }
  32. func NewWsLogDistributor(startLogBufferSize int) *WsLogDistributor {
  33. d := &WsLogDistributor{
  34. wsList: make(map[string]*WsSessions),
  35. startLogBufferSize: startLogBufferSize,
  36. }
  37. d.upgrader = &websocket.Upgrader{
  38. CheckOrigin: func(r *http.Request) bool {
  39. return true
  40. },
  41. }
  42. bbuf := make([]byte, 0, startLogBufferSize)
  43. d.startLogBuffer = bytes.NewBuffer(bbuf)
  44. flake, err := snowflake.NewSnowflake(int64(0), int64(0))
  45. if err != nil {
  46. panic(err)
  47. }
  48. d.snowFlaker = flake
  49. return d
  50. }
  51. func (w *WsLogDistributor) Write(p []byte) (n int, err error) {
  52. if w.startLogBuffer.Len() < w.startLogBufferSize {
  53. w.startLogBuffer.Write(p)
  54. }
  55. defer w.lock.Unlock()
  56. w.lock.Lock()
  57. for k, v := range w.wsList {
  58. _, err := v.Conn.Write(p)
  59. if err != nil {
  60. _ = v.Conn.WS.Close()
  61. delete(w.wsList, k)
  62. }
  63. }
  64. return len(p), nil
  65. }
  66. func (w *WsLogDistributor) HandleNewConnections(ctx *gin.Context) {
  67. ustr := fmt.Sprintf("%16X", w.snowFlaker.NextVal())
  68. ws, err := w.upgrader.Upgrade(ctx.Writer, ctx.Request, nil)
  69. if err != nil {
  70. ctx.JSON(500, gin.H{
  71. "suc": false,
  72. "err_msg": "failed upgrading to websocket",
  73. })
  74. return
  75. }
  76. nconn := &gorillaws2netconn.NetConn4Gorilla{WS: ws, WriteMessageType: websocket.TextMessage}
  77. w.lock.Lock()
  78. sess := NewWsSessions(nconn)
  79. w.wsList[ustr] = sess
  80. w.lock.Unlock()
  81. buf := make([]byte, 1024)
  82. kc := make(chan int)
  83. go func() {
  84. for {
  85. _, err := sess.Conn.Read(buf)
  86. if err != nil {
  87. kc <- 0
  88. return
  89. }
  90. }
  91. }()
  92. <-kc
  93. _ = sess.Conn.Close()
  94. w.lock.Lock()
  95. delete(w.wsList, ustr)
  96. w.lock.Unlock()
  97. }
  98. func (w *WsLogDistributor) ClearConnections() {
  99. defer w.lock.Unlock()
  100. w.lock.Lock()
  101. for k, v := range w.wsList {
  102. _ = v.Conn.Close()
  103. delete(w.wsList, k)
  104. }
  105. }
  106. func (w *WsLogDistributor) GetClientsInfo() map[string]map[string]interface{} {
  107. defer w.lock.RUnlock()
  108. w.lock.RLock()
  109. now := time.Now()
  110. m := map[string]map[string]interface{}{}
  111. for k, v := range w.wsList {
  112. m[k] = map[string]interface{}{
  113. "begin_time": v.BeginTime.Unix(),
  114. "begin_time_in_rfc3399": v.BeginTime.Format(time.RFC3339),
  115. "begin_to_now_nanoseconds": now.Sub(v.BeginTime),
  116. "begin_to_now_str": now.Sub(v.BeginTime).String(),
  117. "remote_addr": v.Conn.RemoteAddr().String(),
  118. "local_addr": v.Conn.LocalAddr().String(),
  119. }
  120. }
  121. return m
  122. }
  123. func (w *WsLogDistributor) GetStartLog() string {
  124. return w.startLogBuffer.String()
  125. }
  126. func (w *WsLogDistributor) ClearStartLog() {
  127. w.startLogBuffer.Reset()
  128. }