handleWebsocket starts two message readers (consumers), one consuming a Rabbit queue and another reading from a Websocket connection. Each consumer receives a message handler to relay messages – from Rabbit to Websocket and vice versa.
(w http.ResponseWriter, r *http.Request, conn *rabbit.Conn)
| 109 | // Each consumer receives a message handler to relay messages – |
| 110 | // from Rabbit to Websocket and vice versa. |
| 111 | func handleWebsocket(w http.ResponseWriter, r *http.Request, conn *rabbit.Conn) { |
| 112 | ws, err := upgrader.Upgrade(w, r, nil) |
| 113 | if err != nil { |
| 114 | log.Printf("upgrade websocket: %s", err) |
| 115 | return |
| 116 | } |
| 117 | defer ws.Close() |
| 118 | |
| 119 | // A separate channel for a publisher in a go routine. |
| 120 | ch, err := conn.Connection.Channel() |
| 121 | if err != nil { |
| 122 | log.Printf("open channel: %s", err) |
| 123 | return |
| 124 | } |
| 125 | defer ch.Close() |
| 126 | |
| 127 | // done and cancel() makes sure all spawned go routines are |
| 128 | // terminated if any one of them is finished. |
| 129 | done := make(chan bool) |
| 130 | ctx, cancel := context.WithCancel(context.Background()) |
| 131 | defer cancel() |
| 132 | |
| 133 | // Start a Rabbit consumer |
| 134 | err = conn.StartConsumerTemp(ctx, done, conf.Exchange, conf.KeyFront, handleWriteWebsocket(ws)) |
| 135 | if err != nil { |
| 136 | log.Printf("start temp consumer: %s", err) |
| 137 | return |
| 138 | } |
| 139 | |
| 140 | // Start a websocket reader (consumer) |
| 141 | err = iwebsocket.StartReader(ctx, done, ws, handlePublishRabbit(ch)) |
| 142 | if err != nil { |
| 143 | log.Printf("start websocket reader: %s", err) |
| 144 | return |
| 145 | } |
| 146 | |
| 147 | <-done |
| 148 | } |
| 149 | |
| 150 | // handleWriteWebsocket writes a Rabbit message to Websocket. |
| 151 | // A Rabbit consumer only passes a message. So, a Websocket connection is |
no test coverage detected