package exchange import ( "encoding/json" "fmt" "log" "math" "strings" "time" ) // BinanceWS connects to Binance WS for ticker data. // Splits symbols across multiple combined-stream connections. type BinanceWS struct { Tracked []string connections int } func NewBinanceWS(tracked []string) *BinanceWS { conns := int(math.Ceil(float64(len(tracked)) / 60)) if conns < 1 { conns = 1 } if conns > 10 { conns = 10 } return &BinanceWS{Tracked: tracked, connections: conns} } func (b *BinanceWS) runSingle(symbols []string, connIdx int, updateFn func(coin string, price, bid, ask float64)) error { streams := "" for i, sym := range symbols { if i > 0 { streams += "/" } streams += fmt.Sprintf("%s@bookTicker", strings.ToLower(sym)) } url := fmt.Sprintf("wss://fstream.binance.com/stream?streams=%s", streams) name := fmt.Sprintf("Binance-%d", connIdx) conn := NewPriceConnector(url, name, 60*time.Second, 15*time.Second) // No client-side pings — let the proxy handle keepalive conn.PingInterval = 0 conn.OnConnect = func() { log.Printf("[%s] Connected (%d symbols)", name, len(symbols)) } conn.OnMessage = func(msg []byte) { var raw map[string]json.RawMessage if err := json.Unmarshal(msg, &raw); err != nil { return } dataRaw, ok := raw["data"] if !ok { return } // Parse data as a generic map to avoid field name conflicts // (bookTicker has both "b" bid price and "B" bid quantity) var dataMap map[string]interface{} if err := json.Unmarshal(dataRaw, &dataMap); err != nil { return } symbol, _ := dataMap["s"].(string) bidStr, _ := dataMap["b"].(string) askStr, _ := dataMap["a"].(string) if symbol == "" || bidStr == "" || askStr == "" { return } bid := parseFloat(bidStr) ask := parseFloat(askStr) if bid <= 0 || ask <= 0 { return } coin := symbolToCoin(symbol, "USDT") if coin == "" { return } mid := (bid + ask) / 2.0 updateFn(coin, mid, bid, ask) } return conn.Run() } func (b *BinanceWS) Run(updateFn func(coin string, price, bid, ask float64)) error { if len(b.Tracked) == 0 { return nil } n := b.connections perConn := (len(b.Tracked) + n - 1) / n errCh := make(chan error, n) for i := 0; i < n; i++ { start := i * perConn end := start + perConn if end > len(b.Tracked) { end = len(b.Tracked) } if start >= end { errCh <- nil continue } batch := b.Tracked[start:end] go func(idx int, syms []string) { errCh <- b.runSingle(syms, idx, updateFn) }(i+1, batch) } return <-errCh }