package exchange import ( "encoding/json" "log" "strings" "time" ) // OKXWS connects to OKX WebSocket for tickers channel (perpetual swaps). type OKXWS struct { Tracked []string // OKX symbols like BTC-USDT-SWAP } type okxSubscribeMsg struct { Op string `json:"op"` Args []okxChannel `json:"args"` } type okxChannel struct { Channel string `json:"channel"` InstID string `json:"instId"` } type okxTickerMsg struct { Arg okxChannel `json:"arg"` Data []okxTickerData `json:"data"` } type okxTickerData struct { Last string `json:"last"` BidPx string `json:"bidPx"` AskPx string `json:"askPx"` } func NewOKXWS(tracked []string) *OKXWS { return &OKXWS{ Tracked: tracked, } } // Run connects to OKX WS and streams ticker data. func (o *OKXWS) Run(updateFn func(coin string, price, bid, ask float64)) error { url := "wss://ws.okx.com:8443/ws/v5/public" conn := NewPriceConnector(url, "OKX", 120*time.Second, 30*time.Second) conn.PingInterval = 20 * time.Second // OKX requires ping within 30s conn.TextPing = true // OKX expects text "ping" message conn.OnConnect = func() { log.Printf("[OKX WS] Connected, subscribing (%d symbols)", len(o.Tracked)) // Batch subscriptions — OKX has rate limits (3 req/s, 480/hr) batchSize := 20 for i := 0; i < len(o.Tracked); i += batchSize { end := i + batchSize if end > len(o.Tracked) { end = len(o.Tracked) } batch := o.Tracked[i:end] args := make([]okxChannel, len(batch)) for j, sym := range batch { args[j] = okxChannel{ Channel: "tickers", InstID: sym, } } sub := okxSubscribeMsg{ Op: "subscribe", Args: args, } if err := conn.SendJSON(sub); err != nil { log.Printf("[OKX WS] Subscribe error (batch %d): %v", i/batchSize, err) } } } conn.OnMessage = func(msg []byte) { // Handle OKX text "pong" response if string(msg) == "pong" { return } // Check for 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" || evt == "subscribe" { return } } var ticker okxTickerMsg if err := json.Unmarshal(msg, &ticker); err != nil { return } if len(ticker.Data) == 0 || ticker.Data[0].Last == "" { return } // Convert BTC-USDT-SWAP -> BTC coin := okxSymbolToCoin(ticker.Arg.InstID) if coin == "" { return } price := parseFloat(ticker.Data[0].Last) if price > 0 { bid := parseFloat(ticker.Data[0].BidPx) ask := parseFloat(ticker.Data[0].AskPx) updateFn(coin, price, bid, ask) } } return conn.Run() } // okxSymbolToCoin converts "BTC-USDT-SWAP" to "BTC". func okxSymbolToCoin(symbol string) string { // Strip "-USDT-SWAP" suffix const suffix = "-USDT-SWAP" if !strings.HasSuffix(symbol, suffix) { return "" } return symbol[:len(symbol)-len(suffix)] }