diff --git a/server/socket.go b/server/socket.go index 7c818e4..4132706 100644 --- a/server/socket.go +++ b/server/socket.go @@ -6,7 +6,7 @@ import ( "log" "net" "os" - + jsoniter "github.com/json-iterator/go" ) diff --git a/server/websocket.go b/server/websocket.go index 0f2226b..26503a7 100644 --- a/server/websocket.go +++ b/server/websocket.go @@ -15,6 +15,11 @@ type Client struct { Conn *websocket.Conn } +type wsMsg struct { + Mes []byte + client *Client +} + var ( wsClients = make(map[*Client]bool) @@ -22,10 +27,12 @@ var ( ws = struct { Broadcast chan []byte + Send chan *wsMsg Add chan *Client Del chan *Client }{ Broadcast: make(chan []byte, 1), + Send: make(chan *wsMsg, 1), Add: make(chan *Client), Del: make(chan *Client), } @@ -50,12 +57,15 @@ func broadcast() { case client := <-ws.Add: wsClients[client] = true - case conn := <-ws.Del: - delete(wsClients, conn) + case client := <-ws.Del: + delete(wsClients, client) - case r := <-ws.Broadcast: + case msg := <-ws.Send: + sendMessage(msg) + + case b := <-ws.Broadcast: //fmt.Println(len(ws.Broadcast), cap(ws.Broadcast)) - sendBroadcast(r) + sendBroadcast(b) case <-ticker.C: fmt.Println(len(wsClients), runtime.NumGoroutine(), len(ws.Broadcast), cap(ws.Broadcast)) @@ -63,6 +73,15 @@ func broadcast() { } } +func sendMessage(x *wsMsg) { + if _, ok := wsClients[x.client]; ok { + if err := x.client.Conn.WriteMessage(1, x.Mes); err != nil { + delete(wsClients, x.client) + x.client.Conn.Close() + } + } +} + func sendBroadcast(message []byte) { input := struct { Chanel string `json:"chanel"` @@ -126,11 +145,16 @@ func readWS(conn *websocket.Conn) { ws.Del <- client }() + ping := time.Now().Unix() for { conn.SetReadDeadline(time.Now().Add(30 * time.Minute)) - _, _, err := conn.ReadMessage() + _, message, err := conn.ReadMessage() if err != nil { return } + if string(message) == "ping" && time.Now().Unix() > ping { + ws.Send <- &wsMsg{client: client, Mes: []byte("pong")} + ping = time.Now().Unix() + 15 + } } }