package main import ( "flag" "log" "net/url" "os" "os/signal" "time" "github.com/gorilla/websocket" ) var ( addr = flag.String("addr", "localhost:8080", "http service address") streams = flag.String("streams", "eurusd.ob-inc,eurusd.kline-12h", "streams to connect to") wait = flag.Float64("wait", 2, "Time to wait between submit batch of messages") ) func main() { flag.Parse() interrupt := make(chan os.Signal, 1) signal.Notify(interrupt, os.Interrupt) u := url.URL{ Scheme: "ws", Host: *addr, Path: "", RawQuery: "stream=" + *streams, } log.Printf("connecting to %s", u.String()) dialer := websocket.DefaultDialer dialer.ReadBufferSize = 1 dialer.WriteBufferSize = 1 c, _, err := dialer.Dial(u.String(), nil) if err != nil { log.Fatal("dial:", err) } defer c.Close() done := make(chan bool) go func() { defer close(done) for { _, message, err := c.ReadMessage() if err != nil { log.Println("read:", err) return } log.Printf("recv: %s", message) } }() ticker := time.NewTicker(time.Second) defer ticker.Stop() for { select { case <-done: return // case t := <-ticker.C: // err := c.WriteMessage(websocket.TextMessage, []byte(t.String())) // if err != nil { // log.Println("write:", err) // return // } case <-interrupt: log.Println("interrupt") // Cleanly close the connection by sending a close message and then // waiting (with timeout) for the server to close the connection. err := c.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "")) if err != nil { log.Println("write close:", err) return } select { case <-done: case <-time.After(time.Second): } return } } }