Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 12 additions & 14 deletions packages/sdk-go/websocket/client_dispatch.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,23 +52,21 @@ func (c *Client) dispatchFrame(stop <-chan struct{}, payload []byte, generation
// the explicit event discriminator before the heuristic market-data probes
// below; otherwise an order update can be silently delivered as a trade.
if bytes.Contains(payload, bOrder) || bytes.Contains(payload, bOrderUpdate) {
var eventEnvelope map[string]json.RawMessage
if err := json.Unmarshal(payload, &eventEnvelope); err == nil {
var explicitType string
if rawType, ok := eventEnvelope["e"]; ok {
_ = json.Unmarshal(rawType, &explicitType)
}
if explicitType == "order" || explicitType == "orderUpdate" {
var order OrderEvent
if err := json.Unmarshal(payload, &order); err == nil {
for _, subs := range tables.orderSubs {
for _, sub := range subs {
sub.send(stop, &order)
}
var eventEnvelope struct {
EventType string `json:"e"`
EventTime int64 `json:"E"`
}
if err := json.Unmarshal(payload, &eventEnvelope); err == nil &&
(eventEnvelope.EventType == "order" || eventEnvelope.EventType == "orderUpdate") {
var order OrderEvent
if err := json.Unmarshal(payload, &order); err == nil {
for _, subs := range tables.orderSubs {
for _, sub := range subs {
sub.send(stop, &order)
}
}
return nil
}
return nil
}
}

Expand Down
43 changes: 43 additions & 0 deletions packages/sdk-go/websocket/client_dispatch_bench_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
package websocket

import "testing"

var dispatchBenchPayloads = map[string][]byte{
"depth": []byte(
`{"e":"depthUpdate","E":1720000000000,"s":"BTCUSD","U":100,"u":101,"b":[["65000.00","1.25"],["64999.50","2.50"]],"a":[["65001.00","0.75"],["65002.00","3.00"]]}`,
),
"trade": []byte(
`{"e":"trade","E":1720000000000,"s":"BTCUSD","t":123456,"p":"65000.50","q":"0.125","T":1720000000000,"m":false}`,
),
"ticker": []byte(
`{"e":"bookTicker","E":1720000000000,"s":"BTCUSD","b":"65000.00","B":"1.25","a":"65001.00","A":"0.75"}`,
),
"order": []byte(
`{"e":"orderUpdate","E":1720000000000,"s":"BTCUSD","i":12345,"t":777,"S":"BUY","X":"NEW","p":"65000.00","q":"0.10"}`,
),
}

func newDispatchBenchClient() *Client {
client := NewClient("wss://ws.gemini.com")
client.state.Store(int32(StateConnected))
return client
}

func BenchmarkDispatchFrame(b *testing.B) {
for name, payload := range dispatchBenchPayloads {
b.Run(name, func(b *testing.B) {
client := newDispatchBenchClient()
stop := make(chan struct{})

b.ReportAllocs()
b.SetBytes(int64(len(payload)))
b.ResetTimer()

for i := 0; i < b.N; i++ {
if err := client.dispatchFrame(stop, payload, 0); err != nil {
b.Fatal(err)
}
}
})
}
}
7 changes: 5 additions & 2 deletions packages/sdk-go/websocket/client_dispatch_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,13 @@ func TestDispatchFrame_OrderUpdateWithTradeIDRemainsOrderEvent(t *testing.T) {
client.subTables.Store(tables)
client.subsMu.Unlock()

client.dispatchFrame(make(chan struct{}), []byte(`{"e":"orderUpdate","s":"BTCUSD","t":777,"i":12345,"S":"BUY","X":"NEW"}`), 0)
client.dispatchFrame(make(chan struct{}), []byte(`{"e":"orderUpdate","E":1720000000000,"s":"BTCUSD","t":777,"i":12345,"S":"BUY","X":"NEW"}`), 0)
select {
case event := <-orderSub.ch:
if event.EventType != "orderUpdate" || event.TradeID != 777 || event.OrderID != 12345 {
if event.EventType != "orderUpdate" ||
event.EventTime != 1720000000000 ||
event.TradeID != 777 ||
event.OrderID != 12345 {
t.Fatalf("unexpected order event: %+v", event)
}
case <-time.After(time.Second):
Expand Down