package market import ( "context" "encoding/json" "fmt" "strconv" "time" "github.com/coder/websocket" ) // Feed connects to an exchange WebSocket API and streams order book // updates as OrderBook snapshots. Each exchange has a different wire // format; the Feed normalizes them into the common OrderBook type. type Feed struct { URL string Venue Venue Asset string conn *websocket.Conn Books chan *OrderBook Errors chan error } // NewFeed creates a feed for a specific asset on a specific venue. func NewFeed(url string, venue Venue, asset string) *Feed { return &Feed{ URL: url, Venue: venue, Asset: asset, Books: make(chan *OrderBook, 64), Errors: make(chan error, 8), } } // Connect establishes the WebSocket connection and sends the // subscription message. func (f *Feed) Connect(ctx context.Context) error { conn, _, err := websocket.Dial(ctx, f.URL, nil) if err != nil { return fmt.Errorf("dial %s: %w", f.URL, err) } conn.SetReadLimit(1 << 20) // 1MB f.conn = conn return nil } // Listen reads messages from the exchange and parses them into // OrderBook snapshots. Call this in a goroutine. func (f *Feed) Listen(ctx context.Context) { defer close(f.Books) for { _, data, err := f.conn.Read(ctx) if err != nil { select { case f.Errors <- err: default: } return } ob, err := f.parseMessage(data) if err != nil { continue // skip unparseable messages } if ob != nil { select { case f.Books <- ob: default: // drop if consumer is slow } } } } // Close disconnects the feed. func (f *Feed) Close() { if f.conn != nil { f.conn.Close(websocket.StatusNormalClosure, "done") } } // parseMessage tries to parse an exchange message into an OrderBook. // Supports a generic JSON format that covers multiple exchanges. func (f *Feed) parseMessage(data []byte) (*OrderBook, error) { // Try generic order book format: {"bids": [[price, qty], ...], "asks": [[price, qty], ...]} var msg struct { Bids [][]json.Number `json:"bids"` Asks [][]json.Number `json:"asks"` } if err := json.Unmarshal(data, &msg); err != nil { return nil, err } if len(msg.Bids) == 0 && len(msg.Asks) == 0 { return nil, nil // not an order book message } ob := &OrderBook{ Asset: f.Asset, Venue: f.Venue, Time: time.Now(), } for _, level := range msg.Bids { if len(level) < 2 { continue } price, _ := strconv.ParseFloat(string(level[0]), 64) volume, _ := strconv.ParseFloat(string(level[1]), 64) if price > 0 { ob.Bids = append(ob.Bids, OrderEntry{ Price: price, Volume: volume, Side: Bid, Venue: f.Venue, Time: ob.Time, }) } } for _, level := range msg.Asks { if len(level) < 2 { continue } price, _ := strconv.ParseFloat(string(level[0]), 64) volume, _ := strconv.ParseFloat(string(level[1]), 64) if price > 0 { ob.Asks = append(ob.Asks, OrderEntry{ Price: price, Volume: volume, Side: Ask, Venue: f.Venue, Time: ob.Time, }) } } return ob, nil }