package exchange import ( "encoding/json" "log" "time" ) // BitgetWS connects to Bitget WebSocket for ticker channel. type BitgetWS struct { Conn *PriceConnector Tracked []string // Bitget symbols like BTCUSDT } type bitgetSubscribeMsg struct { Op string `json:"op"` Args []bitgetChannel `json:"args"` } type bitgetChannel struct { InstType string `json:"instType"` Channel string `json:"channel"` InstID string `json:"instId"` } type bitgetTickerMsg struct { Action string `json:"action"` Arg bitgetChannel `json:"arg"` Data []bitgetTickerData `json:"data"` } type bitgetTickerData struct { LastPr string `json:"lastPr"` BidPr string `json:"bidPr"` AskPr string `json:"askPr"` } func NewBitgetWS(tracked []string) *BitgetWS { return &BitgetWS{ Tracked: tracked, } } // Run connects to Bitget WS and streams ticker data. func (b *BitgetWS) Run(updateFn func(coin string, price, bid, ask float64)) error { url := "wss://ws.bitget.com/v2/ws/public" b.Conn = NewPriceConnector(url, "Bitget", 120*time.Second, 30*time.Second) b.Conn.PingInterval = 25 * time.Second // Bitget requires ping within 30s b.Conn.TextPing = true // Bitget v2 expects text "ping" message b.Conn.OnConnect = func() { log.Printf("[Bitget WS] Connected, subscribing (%d symbols)", len(b.Tracked)) // Batch subscriptions — Bitget WS has a limit per message batchSize := 20 for i := 0; i < len(b.Tracked); i += batchSize { end := i + batchSize if end > len(b.Tracked) { end = len(b.Tracked) } batch := b.Tracked[i:end] args := make([]map[string]string, len(batch)) for j, sym := range batch { args[j] = map[string]string{ "instType": "USDT-FUTURES", "channel": "ticker", "instId": sym, } } sub := map[string]interface{}{ "op": "subscribe", "args": args, } if err := b.Conn.SendJSON(sub); err != nil { log.Printf("[Bitget WS] Subscribe error (batch %d): %v", i/batchSize, err) } } } b.Conn.OnMessage = func(msg []byte) { // Handle Bitget text "pong" response if string(msg) == "pong" { return } // Check for Bitget subscription confirmation or error response var generic map[string]interface{} if err := json.Unmarshal(msg, &generic); err == nil { if evt, _ := generic["event"].(string); evt == "error" { log.Printf("[Bitget WS] Subscribe error response: %s", string(msg)) return } } var ticker bitgetTickerMsg if err := json.Unmarshal(msg, &ticker); err != nil { log.Printf("[Bitget WS] Unrecognized message: %s", string(msg)) return } if len(ticker.Data) == 0 || ticker.Data[0].LastPr == "" { return } // Convert BTCUSDT -> BTC coin := symbolToCoin(ticker.Arg.InstID, "USDT") if coin == "" { return } price := parseFloat(ticker.Data[0].LastPr) if price > 0 { bid := parseFloat(ticker.Data[0].BidPr) ask := parseFloat(ticker.Data[0].AskPr) updateFn(coin, price, bid, ask) } } return b.Conn.Run() }