package market import ( "context" "fmt" "sync" "time" "github.com/alpacahq/alpaca-trade-api-go/v3/marketdata" "github.com/alpacahq/alpaca-trade-api-go/v3/marketdata/stream" ) // AlpacaClient wraps the Alpaca market data API for ingesting OHLCV bars. // It handles both historical (REST) and real-time (WebSocket) bar data, // converting everything into the organism's native Bar type. type AlpacaClient struct { mu sync.Mutex dataClient *marketdata.Client streamClient *stream.CryptoClient Symbols []string // e.g., ["BTC/USD", "ETH/USD"] Bars chan *Bar Errors chan error connected bool } // AlpacaConfig holds the configuration for connecting to Alpaca. type AlpacaConfig struct { APIKey string // APCA_API_KEY_ID APISecret string // APCA_API_SECRET_KEY Symbols []string // Paper controls whether trading goes to paper-api or live api. // Market data always goes to data.alpaca.markets regardless. Paper bool } // NewAlpacaClient creates a client connected to Alpaca's market data API. func NewAlpacaClient(cfg AlpacaConfig) *AlpacaClient { opts := marketdata.ClientOpts{ APIKey: cfg.APIKey, APISecret: cfg.APISecret, } return &AlpacaClient{ dataClient: marketdata.NewClient(opts), Symbols: cfg.Symbols, Bars: make(chan *Bar, 256), Errors: make(chan error, 16), } } // GetHistoricalBars fetches OHLCV bars for a symbol between start and end. // TimeFrame controls the bar resolution (1 minute, 1 hour, 1 day, etc.). func (ac *AlpacaClient) GetHistoricalBars( symbol string, tf marketdata.TimeFrame, start, end time.Time, ) ([]*Bar, error) { req := marketdata.GetCryptoBarsRequest{ TimeFrame: tf, Start: start, End: end, } cryptoBars, err := ac.dataClient.GetCryptoBars(symbol, req) if err != nil { return nil, fmt.Errorf("get crypto bars %s: %w", symbol, err) } bars := make([]*Bar, len(cryptoBars)) for i, cb := range cryptoBars { bars[i] = &Bar{ Symbol: symbol, Open: cb.Open, High: cb.High, Low: cb.Low, Close: cb.Close, Volume: cb.Volume, VWAP: cb.VWAP, TradeCount: cb.TradeCount, Timestamp: cb.Timestamp, } } return bars, nil } // GetMultiHistoricalBars fetches bars for multiple symbols at once. func (ac *AlpacaClient) GetMultiHistoricalBars( tf marketdata.TimeFrame, start, end time.Time, ) (map[string][]*Bar, error) { req := marketdata.GetCryptoBarsRequest{ TimeFrame: tf, Start: start, End: end, } cryptoBars, err := ac.dataClient.GetCryptoMultiBars(ac.Symbols, req) if err != nil { return nil, fmt.Errorf("get crypto multi bars: %w", err) } result := make(map[string][]*Bar, len(cryptoBars)) for sym, cbs := range cryptoBars { bars := make([]*Bar, len(cbs)) for i, cb := range cbs { bars[i] = &Bar{ Symbol: sym, Open: cb.Open, High: cb.High, Low: cb.Low, Close: cb.Close, Volume: cb.Volume, VWAP: cb.VWAP, TradeCount: cb.TradeCount, Timestamp: cb.Timestamp, } } result[sym] = bars } return result, nil } // StreamBars connects to Alpaca's WebSocket and streams real-time bars // for all configured symbols. Bars are sent to the Bars channel. // Call this in a goroutine. Cancel the context to stop streaming. func (ac *AlpacaClient) StreamBars(ctx context.Context, apiKey, apiSecret string) error { ac.mu.Lock() if ac.connected { ac.mu.Unlock() return fmt.Errorf("already streaming") } ac.mu.Unlock() cryptoClient := stream.NewCryptoClient( marketdata.US, stream.WithCredentials(apiKey, apiSecret), ) if err := cryptoClient.Connect(ctx); err != nil { return fmt.Errorf("stream connect: %w", err) } ac.mu.Lock() ac.streamClient = cryptoClient ac.connected = true ac.mu.Unlock() err := cryptoClient.SubscribeToBars(func(cb stream.CryptoBar) { bar := &Bar{ Symbol: cb.Symbol, Open: cb.Open, High: cb.High, Low: cb.Low, Close: cb.Close, Volume: cb.Volume, VWAP: cb.VWAP, TradeCount: cb.TradeCount, Timestamp: cb.Timestamp, } select { case ac.Bars <- bar: default: // drop if consumer is slow } }, ac.Symbols...) if err != nil { return fmt.Errorf("subscribe bars: %w", err) } // Wait for termination or context cancellation. select { case err := <-cryptoClient.Terminated(): ac.mu.Lock() ac.connected = false ac.mu.Unlock() if err != nil { return fmt.Errorf("stream terminated: %w", err) } return nil case <-ctx.Done(): ac.mu.Lock() ac.connected = false ac.mu.Unlock() return ctx.Err() } } // IsConnected returns whether the streaming connection is active. func (ac *AlpacaClient) IsConnected() bool { ac.mu.Lock() defer ac.mu.Unlock() return ac.connected }