websocket_hub.go 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225
  1. // Copyright 2016 Martin Hebnes Pedersen (LA5NTA). All rights reserved.
  2. // Use of this source code is governed by the MIT-license that can be
  3. // found in the LICENSE file.
  4. package main
  5. import (
  6. "bufio"
  7. "encoding/json"
  8. "io"
  9. "log"
  10. "os"
  11. "path"
  12. "sync"
  13. "time"
  14. "github.com/fsnotify/fsnotify"
  15. "github.com/gorilla/websocket"
  16. "github.com/la5nta/wl2k-go/mailbox"
  17. )
  18. // WSConn represent one connection in the WSHub pool
  19. type WSConn struct {
  20. conn *websocket.Conn
  21. out chan interface{}
  22. }
  23. // WSHub is a hub for broadcasting data to several websocket connections
  24. type WSHub struct {
  25. mu sync.Mutex
  26. pool map[*WSConn]struct{}
  27. }
  28. func NewWSHub() *WSHub {
  29. w := &WSHub{pool: map[*WSConn]struct{}{}}
  30. go w.watchMBox()
  31. return w
  32. }
  33. func (w *WSHub) UpdateStatus() { w.WriteJSON(struct{ Status Status }{getStatus()}) }
  34. func (w *WSHub) WriteProgress(p Progress) { w.WriteJSON(struct{ Progress Progress }{p}) }
  35. func (w *WSHub) WriteNotification(n Notification) { w.WriteJSON(struct{ Notification Notification }{n}) }
  36. func (w *WSHub) Prompt(p Prompt) {
  37. w.WriteJSON(struct{ Prompt Prompt }{p})
  38. go func() { <-p.cancel; w.WriteJSON(struct{ PromptAbort Prompt }{p}) }()
  39. }
  40. func (w *WSHub) WriteJSON(v interface{}) {
  41. if w == nil {
  42. return
  43. }
  44. w.mu.Lock()
  45. for c, _ := range w.pool {
  46. select {
  47. case c.out <- v:
  48. default:
  49. log.Println("Closing one unresponsive web socket")
  50. c.conn.Close()
  51. delete(w.pool, c)
  52. }
  53. }
  54. w.mu.Unlock()
  55. }
  56. func (w *WSHub) ClientAddrs() []string {
  57. if w == nil {
  58. return nil
  59. }
  60. w.mu.Lock()
  61. defer w.mu.Unlock()
  62. addrs := make([]string, 0, len(w.pool))
  63. for c, _ := range w.pool {
  64. addrs = append(addrs, c.conn.RemoteAddr().String())
  65. }
  66. return addrs
  67. }
  68. func (w *WSHub) watchMBox() {
  69. fsWatcher, err := fsnotify.NewWatcher()
  70. if err != nil {
  71. log.Println("Unable to start fs watcher: ", err)
  72. } else {
  73. p := path.Join(mbox.MBoxPath, mailbox.DIR_INBOX)
  74. if err := fsWatcher.Add(p); err != nil {
  75. log.Printf("Unable to add path '%s' to fs watcher: %s", p, err)
  76. }
  77. // These will probably fail if the first failed, but it's not important to log all.
  78. fsWatcher.Add(path.Join(mbox.MBoxPath, mailbox.DIR_OUTBOX))
  79. fsWatcher.Add(path.Join(mbox.MBoxPath, mailbox.DIR_SENT))
  80. fsWatcher.Add(path.Join(mbox.MBoxPath, mailbox.DIR_ARCHIVE))
  81. defer fsWatcher.Close()
  82. }
  83. for {
  84. select {
  85. // Filesystem events
  86. case <-fsWatcher.Events:
  87. drainEvents(fsWatcher)
  88. websocketHub.WriteJSON(struct {
  89. UpdateMailbox bool
  90. }{true})
  91. case err := <-fsWatcher.Errors:
  92. log.Println(err)
  93. }
  94. }
  95. }
  96. // Handle adds a new websocket to the hub
  97. //
  98. // It will block until the client either stops responding or closes the connection.
  99. func (w *WSHub) Handle(conn *websocket.Conn) {
  100. c := &WSConn{
  101. conn: conn,
  102. out: make(chan interface{}, 1),
  103. }
  104. w.mu.Lock()
  105. w.pool[c] = struct{}{}
  106. w.mu.Unlock()
  107. // Initial status update
  108. // (broadcasted as it includes info to other clients about this new one)
  109. w.UpdateStatus()
  110. quit := wsReadLoop(conn)
  111. lines, done, err := tailFile(fOptions.LogPath)
  112. if err != nil {
  113. log.Println(err)
  114. return
  115. }
  116. defer close(done)
  117. for {
  118. select {
  119. case line := <-lines:
  120. c.conn.WriteJSON(struct {
  121. LogLine string
  122. }{string(line)})
  123. case v := <-c.out:
  124. err := c.conn.WriteJSON(v)
  125. if err != nil {
  126. log.Println(err)
  127. }
  128. case <-quit:
  129. // The read loop failed/disconnected. Remove from hub.
  130. c.conn.Close()
  131. w.mu.Lock()
  132. delete(w.pool, c)
  133. defer w.UpdateStatus()
  134. w.mu.Unlock()
  135. return
  136. }
  137. }
  138. }
  139. func drainEvents(w *fsnotify.Watcher) {
  140. for {
  141. select {
  142. case <-w.Events:
  143. default:
  144. return
  145. }
  146. }
  147. }
  148. // Expects the file to never get renamed/truncated or deleted
  149. func tailFile(path string) (<-chan []byte, chan<- struct{}, error) {
  150. lines := make(chan []byte)
  151. done := make(chan struct{})
  152. file, err := os.Open(path)
  153. if err != nil {
  154. return nil, nil, err
  155. }
  156. go func() {
  157. rd := bufio.NewReader(file)
  158. for {
  159. data, _, err := rd.ReadLine()
  160. if err == io.EOF {
  161. time.Sleep(time.Millisecond * 100)
  162. continue
  163. }
  164. select {
  165. case <-done:
  166. file.Close()
  167. return
  168. case lines <- data:
  169. }
  170. }
  171. }()
  172. return (<-chan []byte)(lines), (chan<- struct{})(done), nil
  173. }
  174. func handleWSMessage(v map[string]json.RawMessage) {
  175. raw, ok := v["prompt_response"]
  176. if !ok {
  177. return
  178. }
  179. var resp PromptResponse
  180. json.Unmarshal(raw, &resp)
  181. promptHub.Respond(resp.ID, resp.Value, resp.Err)
  182. }
  183. func wsReadLoop(c *websocket.Conn) <-chan struct{} {
  184. quit := make(chan struct{})
  185. go func() {
  186. for {
  187. v := map[string]json.RawMessage{}
  188. err := c.ReadJSON(&v)
  189. if err != nil {
  190. close(quit)
  191. return
  192. }
  193. go handleWSMessage(v)
  194. }
  195. }()
  196. return quit
  197. }