feat(monitor/collectors/coinbase): gap-aware ReadLoop (snapshot resets sequencer; gap → ErrSequenceGap)
This commit is contained in:
parent
3be4c9d815
commit
23a4e2ded4
|
|
@ -0,0 +1,46 @@
|
|||
// ©AngelaMos | 2026
|
||||
// readloop.go
|
||||
|
||||
package coinbase
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
)
|
||||
|
||||
var ErrSequenceGap = errors.New("coinbase: sequence gap detected")
|
||||
|
||||
type FrameHandler func(ctx context.Context, f Frame) error
|
||||
|
||||
func ReadLoop(ctx context.Context, conn *Conn, seq *Sequencer, handler FrameHandler) error {
|
||||
for {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
frame, err := conn.ReadFrame(ctx)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
switch frame.Kind {
|
||||
case FrameTypeSnapshot:
|
||||
seq.Reset()
|
||||
for _, t := range frame.Tickers {
|
||||
_ = seq.Observe(t.ProductID, frame.SequenceNum)
|
||||
}
|
||||
case FrameTypeTicker:
|
||||
for _, t := range frame.Tickers {
|
||||
if seq.Observe(t.ProductID, frame.SequenceNum) {
|
||||
return ErrSequenceGap
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err := handler(ctx, frame); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,106 @@
|
|||
// ©AngelaMos | 2026
|
||||
// readloop_test.go
|
||||
|
||||
package coinbase_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/coinbase"
|
||||
)
|
||||
|
||||
func TestReadLoop_DeliversTickerFrames(t *testing.T) {
|
||||
fs := newFakeServer(t,
|
||||
loadFixture(t, "subscriptions.json"),
|
||||
[]byte(`{"channel":"ticker","sequence_num":1000,"timestamp":"2026-05-01T22:30:00Z","events":[{"type":"update","tickers":[{"product_id":"BTC-USD","price":"42000.00","volume_24_h":"1.0","time":"2026-05-01T22:30:00Z"}]}]}`),
|
||||
[]byte(`{"channel":"ticker","sequence_num":1001,"timestamp":"2026-05-01T22:30:01Z","events":[{"type":"update","tickers":[{"product_id":"BTC-USD","price":"42001.00","volume_24_h":"1.0","time":"2026-05-01T22:30:01Z"}]}]}`),
|
||||
)
|
||||
|
||||
d := coinbase.NewWSDialer(coinbase.DialerConfig{URL: fs.URL(), ProductIDs: []string{"BTC-USD"}})
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
conn, err := d.Dial(ctx)
|
||||
require.NoError(t, err)
|
||||
defer conn.Close()
|
||||
|
||||
seq := coinbase.NewSequencer()
|
||||
mu := sync.Mutex{}
|
||||
tickerFrames := 0
|
||||
|
||||
loopCtx, loopCancel := context.WithTimeout(ctx, 500*time.Millisecond)
|
||||
defer loopCancel()
|
||||
|
||||
err = coinbase.ReadLoop(loopCtx, conn, seq, func(_ context.Context, f coinbase.Frame) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if f.Kind == coinbase.FrameTypeTicker {
|
||||
tickerFrames++
|
||||
}
|
||||
return nil
|
||||
})
|
||||
require.True(t, errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) || err == nil)
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
require.GreaterOrEqual(t, tickerFrames, 2)
|
||||
}
|
||||
|
||||
func TestReadLoop_GapTriggersErrSequenceGap(t *testing.T) {
|
||||
fs := newFakeServer(t,
|
||||
loadFixture(t, "subscriptions.json"),
|
||||
[]byte(`{"channel":"ticker","sequence_num":100,"timestamp":"2026-05-01T22:30:00Z","events":[{"type":"update","tickers":[{"product_id":"BTC-USD","price":"42000.00","volume_24_h":"1.0","time":"2026-05-01T22:30:00Z"}]}]}`),
|
||||
[]byte(`{"channel":"ticker","sequence_num":250,"timestamp":"2026-05-01T22:30:01Z","events":[{"type":"update","tickers":[{"product_id":"BTC-USD","price":"42001.00","volume_24_h":"1.0","time":"2026-05-01T22:30:01Z"}]}]}`),
|
||||
)
|
||||
|
||||
d := coinbase.NewWSDialer(coinbase.DialerConfig{URL: fs.URL(), ProductIDs: []string{"BTC-USD"}})
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
conn, err := d.Dial(ctx)
|
||||
require.NoError(t, err)
|
||||
defer conn.Close()
|
||||
|
||||
seq := coinbase.NewSequencer()
|
||||
loopErr := coinbase.ReadLoop(ctx, conn, seq, func(context.Context, coinbase.Frame) error {
|
||||
return nil
|
||||
})
|
||||
require.ErrorIs(t, loopErr, coinbase.ErrSequenceGap)
|
||||
}
|
||||
|
||||
func TestReadLoop_SnapshotResetsSequencer(t *testing.T) {
|
||||
fs := newFakeServer(t,
|
||||
loadFixture(t, "subscriptions.json"),
|
||||
loadFixture(t, "snapshot.json"),
|
||||
[]byte(`{"channel":"ticker","sequence_num":2,"timestamp":"2026-05-01T22:30:02Z","events":[{"type":"update","tickers":[{"product_id":"BTC-USD","price":"42164.00","volume_24_h":"1.0","time":"2026-05-01T22:30:02Z"}]}]}`),
|
||||
)
|
||||
|
||||
d := coinbase.NewWSDialer(coinbase.DialerConfig{URL: fs.URL(), ProductIDs: []string{"BTC-USD"}})
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
conn, err := d.Dial(ctx)
|
||||
require.NoError(t, err)
|
||||
defer conn.Close()
|
||||
|
||||
seq := coinbase.NewSequencer()
|
||||
mu := sync.Mutex{}
|
||||
kinds := []coinbase.FrameType{}
|
||||
|
||||
loopCtx, loopCancel := context.WithTimeout(ctx, 800*time.Millisecond)
|
||||
defer loopCancel()
|
||||
err = coinbase.ReadLoop(loopCtx, conn, seq, func(_ context.Context, f coinbase.Frame) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
kinds = append(kinds, f.Kind)
|
||||
return nil
|
||||
})
|
||||
require.True(t, errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) || err == nil)
|
||||
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
require.Contains(t, kinds, coinbase.FrameTypeSnapshot)
|
||||
require.Contains(t, kinds, coinbase.FrameTypeTicker)
|
||||
}
|
||||
Loading…
Reference in New Issue