Files
jackyu66gitandClaude Opus 4.6 b7767c95ae feat: 添加OKX行情接入+趋势检测+累积变动系统+界面重构
- 新增OKX WebSocket行情连接器,扩展4交易所价格监控
- 新增z-score趋势检测引擎(TrendDetector),识别价格异动/趋势启动
- 新增累积变动跟踪(CumulativeTracker),基于1min/5min多交易所共识
- 趋势事件和累积变动事件持久化到SQLite
- 新增Binance/OKX动量检测字段,扩展前端动量卡片至15列
- 迁移至macOS(darwin-arm64),更新前端依赖
- Dashboard网格重构:非交易卡片置顶,交易卡片置底
- TrackedCoin添加OK字段,添加ExBinance/ExOKX常量
- 前端新增趋势检测卡片、趋势历史卡片、累积变动卡片

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-05-06 13:26:05 +08:00

126 lines
2.9 KiB
Go

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)]
}