Files
exchange-monitor-go/exchange/dydx.go
T
2026-05-03 16:54:36 +08:00

175 lines
4.0 KiB
Go

package exchange
import (
"encoding/json"
"log"
"time"
)
// DydxWS connects to dYdX v4 WebSocket for market data (oracle prices).
type DydxWS struct {
Conn *PriceConnector
Tracked []string // coin names like ["BTC", "ETH", ...]
}
type dydxSubscribeMsg struct {
Type string `json:"type"`
Channel string `json:"channel"`
ID string `json:"id,omitempty"`
}
type dydxMarketMsg struct {
Type string `json:"type"`
ID string `json:"id"`
Contents json.RawMessage `json:"contents"`
}
type dydxMarketContents struct {
OraclePrice string `json:"oraclePrice"`
MarkPrice string `json:"markPrice"`
NextFundingRate string `json:"nextFundingRate"`
}
func NewDydxWS(tracked []string) *DydxWS {
return &DydxWS{Tracked: tracked}
}
// Run connects to dYdX v4 WS and streams oracle/market prices.
func (d *DydxWS) Run(updateFn func(coin string, price, bid, ask float64)) error {
url := "wss://indexer.dydx.trade/v4/ws"
d.Conn = NewPriceConnector(url, "dYdX", 120*time.Second, 30*time.Second)
// dYdX v4 requires JSON {"type":"ping"} heartbeat
d.Conn.OnConnect = func() {
log.Printf("[dYdX WS] Connected, subscribing")
// Subscribe to all markets (gets all coins in one stream)
sub := dydxSubscribeMsg{
Type: "subscribe",
Channel: "v4_markets",
}
if err := d.Conn.SendJSON(sub); err != nil {
log.Printf("[dYdX WS] Subscribe error: %v", err)
}
// dYdX requires JSON {"type":"ping"} every ~30s
go func() {
heartbeat := time.NewTicker(15 * time.Second)
defer heartbeat.Stop()
// Send first ping after 10s (let subscription settle)
time.Sleep(10 * time.Second)
for {
select {
case <-heartbeat.C:
if err := d.Conn.SendJSON(map[string]string{"type": "ping"}); err != nil {
log.Printf("[dYdX WS] Heartbeat send error: %v", err)
}
case <-d.Conn.Done():
return
}
}
}()
}
d.Conn.OnMessage = func(msg []byte) {
var raw map[string]json.RawMessage
if err := json.Unmarshal(msg, &raw); err != nil {
return
}
// Check type
var msgType string
if err := json.Unmarshal(raw["type"], &msgType); err != nil {
return
}
if msgType != "channel_data" {
// Handle initial subscription response with all markets
if msgType == "subscribed" {
var contents struct {
Markets map[string]struct {
OraclePrice string `json:"oraclePrice"`
} `json:"markets"`
}
contentsRaw, ok := raw["contents"]
if !ok {
return
}
if err := json.Unmarshal(contentsRaw, &contents); err != nil {
return
}
for marketID, market := range contents.Markets {
if market.OraclePrice == "" {
continue
}
coin := dydxSymbolToCoin(marketID)
if coin == "" {
continue
}
if !isTracked(d.Tracked, coin) {
continue
}
price := parseFloat(market.OraclePrice)
if price > 0 {
updateFn(coin, price, 0, 0)
}
}
}
return
}
// Live updates - type "channel_data" with oraclePrices
contentsRaw, ok := raw["contents"]
if !ok {
return
}
var contents struct {
OraclePrices map[string]struct {
OraclePrice string `json:"oraclePrice"`
} `json:"oraclePrices"`
}
if err := json.Unmarshal(contentsRaw, &contents); err != nil {
return
}
for marketID, data := range contents.OraclePrices {
if data.OraclePrice == "" {
continue
}
coin := dydxSymbolToCoin(marketID)
if coin == "" {
continue
}
if !isTracked(d.Tracked, coin) {
continue
}
price := parseFloat(data.OraclePrice)
if price > 0 {
updateFn(coin, price, 0, 0)
}
}
}
return d.Conn.Run()
}
// dydxSymbolToCoin converts "BTC-USD" -> "BTC", "ETH-USD" -> "ETH"
func dydxSymbolToCoin(symbol string) string {
if len(symbol) < 4 {
return ""
}
// Remove "-USD" suffix
if len(symbol) > 4 && symbol[len(symbol)-4:] == "-USD" {
return symbol[:len(symbol)-4]
}
return symbol
}
func isTracked(list []string, coin string) bool {
for _, t := range list {
if t == coin {
return true
}
}
return false
}