mirror of
https://github.com/c9s/bbgo.git
synced 2024-11-22 06:53:52 +00:00
add public only mode to stream
This commit is contained in:
parent
ce0e28708a
commit
f56318c9b6
|
@ -45,6 +45,7 @@ var rootCmd = &cobra.Command{
|
|||
var exchange = binance.New(key, secret)
|
||||
|
||||
stream := exchange.NewStream()
|
||||
stream.SetPublicOnly()
|
||||
stream.Subscribe(types.BookChannel, symbol, types.SubscribeOptions{})
|
||||
|
||||
stream.OnBookSnapshot(func(book types.OrderBook) {
|
||||
|
|
|
@ -186,7 +186,13 @@ func (s *Stream) SetPublicOnly() {
|
|||
}
|
||||
|
||||
func (s *Stream) dial(listenKey string) (*websocket.Conn, error) {
|
||||
url := "wss://stream.binance.com:9443/ws/" + listenKey
|
||||
var url string
|
||||
if s.publicOnly {
|
||||
url = "wss://stream.binance.com:9443/ws"
|
||||
} else {
|
||||
url = "wss://stream.binance.com:9443/ws/" + listenKey
|
||||
}
|
||||
|
||||
conn, _, err := websocket.DefaultDialer.Dial(url, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
|
@ -6,7 +6,6 @@ import (
|
|||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
@ -82,13 +81,6 @@ func (s *WebSocketService) Connect(ctx context.Context) error {
|
|||
}
|
||||
})
|
||||
|
||||
s.OnConnect(func(conn *websocket.Conn) {
|
||||
if err := s.Auth(); err != nil {
|
||||
s.EmitError(err)
|
||||
logger.WithError(err).Error("failed to send auth request")
|
||||
}
|
||||
})
|
||||
|
||||
// pre-allocate the websocket client, the websocket client can be used for reconnecting.
|
||||
if err := s.connect(ctx); err != nil {
|
||||
return err
|
||||
|
|
|
@ -19,6 +19,8 @@ type Stream struct {
|
|||
types.StandardStream
|
||||
|
||||
websocketService *max.WebSocketService
|
||||
|
||||
publicOnly bool
|
||||
}
|
||||
|
||||
func NewStream(key, secret string) *Stream {
|
||||
|
@ -28,6 +30,17 @@ func NewStream(key, secret string) *Stream {
|
|||
websocketService: wss,
|
||||
}
|
||||
|
||||
wss.OnConnect(func(conn *websocket.Conn) {
|
||||
if key == "" || secret == "" {
|
||||
log.Warn("MAX API key or secret is empty, will not send authentication command")
|
||||
} else {
|
||||
if err := wss.Auth(); err != nil {
|
||||
wss.EmitError(err)
|
||||
logger.WithError(err).Error("failed to send auth request")
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
wss.OnMessage(func(message []byte) {
|
||||
logger.Debugf("M: %s", message)
|
||||
})
|
||||
|
@ -77,19 +90,19 @@ func NewStream(key, secret string) *Stream {
|
|||
})
|
||||
|
||||
wss.OnBookEvent(func(e max.BookEvent) {
|
||||
newbook, err := e.OrderBook()
|
||||
newBook, err := e.OrderBook()
|
||||
if err != nil {
|
||||
logger.WithError(err).Error("book convert error")
|
||||
return
|
||||
}
|
||||
|
||||
newbook.Symbol = toGlobalSymbol(e.Market)
|
||||
newBook.Symbol = toGlobalSymbol(e.Market)
|
||||
|
||||
switch e.Event {
|
||||
case "snapshot":
|
||||
stream.EmitBookSnapshot(newbook)
|
||||
stream.EmitBookSnapshot(newBook)
|
||||
case "update":
|
||||
stream.EmitBookUpdate(newbook)
|
||||
stream.EmitBookUpdate(newBook)
|
||||
}
|
||||
})
|
||||
|
||||
|
@ -128,6 +141,10 @@ func NewStream(key, secret string) *Stream {
|
|||
return stream
|
||||
}
|
||||
|
||||
func (s *Stream) SetPublicOnly() {
|
||||
s.publicOnly = true
|
||||
}
|
||||
|
||||
func (s *Stream) Subscribe(channel types.Channel, symbol string, options types.SubscribeOptions) {
|
||||
s.websocketService.Subscribe(string(channel), toLocalSymbol(symbol))
|
||||
}
|
||||
|
|
|
@ -8,6 +8,7 @@ type Stream interface {
|
|||
StandardStreamEventHub
|
||||
|
||||
Subscribe(channel Channel, symbol string, options SubscribeOptions)
|
||||
SetPublicOnly()
|
||||
Connect(ctx context.Context) error
|
||||
Close() error
|
||||
}
|
||||
|
|
Loading…
Reference in New Issue
Block a user