()
| 83 | } |
| 84 | |
| 85 | func (sess *Session) writeLoop() { |
| 86 | ticker := time.NewTicker(pingPeriod) |
| 87 | |
| 88 | defer func() { |
| 89 | ticker.Stop() |
| 90 | sess.closeWS() // break readLoop |
| 91 | }() |
| 92 | |
| 93 | for { |
| 94 | select { |
| 95 | case msg, ok := <-sess.send: |
| 96 | if !ok { |
| 97 | // channel closed |
| 98 | return |
| 99 | } |
| 100 | if err := ws_write(sess.ws, websocket.TextMessage, []byte(msg)); err != nil { |
| 101 | log.Println("sess.writeLoop: " + err.Error()) |
| 102 | return |
| 103 | } |
| 104 | case <-sess.stop: |
| 105 | // shutdown requested |
| 106 | return |
| 107 | |
| 108 | case topic := <-sess.detach: |
| 109 | delete(sess.subs, topic) |
| 110 | |
| 111 | case <-ticker.C: |
| 112 | if err := ws_write(sess.ws, websocket.PingMessage, []byte{}); err != nil { |
| 113 | log.Println("sess.writeLoop: ping/" + err.Error()) |
| 114 | return |
| 115 | } |
| 116 | } |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | // Writes a message with the given message type (mt) and payload. |
| 121 | func ws_write(ws *websocket.Conn, mt int, payload []byte) error { |
no test coverage detected