From 23a4e2ded4ed8e474f1f9688809bb44432bd5bc2 Mon Sep 17 00:00:00 2001 From: CarterPerez-dev Date: Sat, 2 May 2026 03:29:43 -0400 Subject: [PATCH] =?UTF-8?q?feat(monitor/collectors/coinbase):=20gap-aware?= =?UTF-8?q?=20ReadLoop=20(snapshot=20resets=20sequencer;=20gap=20=E2=86=92?= =?UTF-8?q?=20ErrSequenceGap)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../internal/collectors/coinbase/readloop.go | 46 ++++++++ .../collectors/coinbase/readloop_test.go | 106 ++++++++++++++++++ 2 files changed, 152 insertions(+) create mode 100644 PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop.go create mode 100644 PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop_test.go diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop.go new file mode 100644 index 00000000..a3b404cf --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop.go @@ -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 + } + } +} diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop_test.go new file mode 100644 index 00000000..b1af6e69 --- /dev/null +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/readloop_test.go @@ -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) +}