diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/collector.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/collector.go index 396bde82..c47ad3f4 100644 --- a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/collector.go +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/collector.go @@ -134,12 +134,31 @@ func (c *Collector) handleConn(ctx context.Context, conn *Conn) error { return nil }) + c.logger.Info("coinbase loop exit", + "err", loopErrString(loopErr), + "emit_count", count, + "agg_open", aggLen(agg), + ) if loopErr == nil || errors.Is(loopErr, ErrSequenceGap) { _ = c.cfg.State.RecordSuccess(ctx, Name, count) } return loopErr } +func loopErrString(err error) string { + if err == nil { + return "" + } + return err.Error() +} + +func aggLen(a *Aggregator) int { + if a == nil { + return 0 + } + return len(a.open) +} + func (c *Collector) shouldEmit(symbol string) bool { c.mu.Lock() defer c.mu.Unlock() 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 index d69f7625..e8fa7ea2 100644 --- 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 @@ -6,6 +6,7 @@ package coinbase import ( "context" "errors" + "log/slog" ) var ErrSequenceGap = errors.New("coinbase: sequence gap detected") @@ -26,17 +27,16 @@ func ReadLoop(ctx context.Context, conn *Conn, seq *Sequencer, handler FrameHand } switch frame.Kind { - case FrameTypeUnknown, FrameTypeSubscriptions, FrameTypeHeartbeats: + case FrameTypeUnknown, FrameTypeSubscriptions: 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 - } + _ = seq.Observe(frame.SequenceNum) + case FrameTypeTicker, FrameTypeHeartbeats: + if seq.Observe(frame.SequenceNum) { + slog.Default().Warn("coinbase seq gap (non-fatal)", + "channel_kind", frame.Kind, + "seq", frame.SequenceNum, + ) } } 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 index b1af6e69..a294ec01 100644 --- 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 @@ -50,7 +50,7 @@ func TestReadLoop_DeliversTickerFrames(t *testing.T) { require.GreaterOrEqual(t, tickerFrames, 2) } -func TestReadLoop_GapTriggersErrSequenceGap(t *testing.T) { +func TestReadLoop_GapIsLoggedButLoopContinues(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"}]}]}`), @@ -58,17 +58,22 @@ func TestReadLoop_GapTriggersErrSequenceGap(t *testing.T) { ) d := coinbase.NewWSDialer(coinbase.DialerConfig{URL: fs.URL(), ProductIDs: []string{"BTC-USD"}}) - ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + ctx, cancel := context.WithTimeout(context.Background(), 600*time.Millisecond) 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 { + delivered := 0 + loopErr := coinbase.ReadLoop(ctx, conn, seq, func(_ context.Context, f coinbase.Frame) error { + if f.Kind == coinbase.FrameTypeTicker { + delivered++ + } return nil }) - require.ErrorIs(t, loopErr, coinbase.ErrSequenceGap) + require.True(t, errors.Is(loopErr, context.DeadlineExceeded) || errors.Is(loopErr, context.Canceled) || loopErr == nil) + require.Equal(t, 2, delivered, "both ticker frames must be delivered despite the gap") } func TestReadLoop_SnapshotResetsSequencer(t *testing.T) { diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer.go index 340a19b9..bcd33570 100644 --- a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer.go +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer.go @@ -7,26 +7,30 @@ import "sync" type Sequencer struct { mu sync.Mutex - last map[string]int64 + last int64 + set bool } func NewSequencer() *Sequencer { - return &Sequencer{last: make(map[string]int64)} + return &Sequencer{} } -func (s *Sequencer) Observe(productID string, seq int64) bool { +func (s *Sequencer) Observe(seq int64) bool { s.mu.Lock() defer s.mu.Unlock() - prev, ok := s.last[productID] - s.last[productID] = seq - if !ok { + if !s.set { + s.last = seq + s.set = true return false } - return seq != prev+1 + expected := s.last + 1 + s.last = seq + return seq != expected } func (s *Sequencer) Reset() { s.mu.Lock() defer s.mu.Unlock() - s.last = make(map[string]int64) + s.set = false + s.last = 0 } diff --git a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer_test.go b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer_test.go index a67ba6fc..92562081 100644 --- a/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer_test.go +++ b/PROJECTS/advanced/monitor-the-situation-dashboard/backend/internal/collectors/coinbase/sequencer_test.go @@ -11,53 +11,40 @@ import ( "github.com/carterperez-dev/monitor-the-situation/backend/internal/collectors/coinbase" ) -func TestSequencer_FirstTickIsAlwaysInOrder(t *testing.T) { +func TestSequencer_FirstObservationIsBaseline(t *testing.T) { s := coinbase.NewSequencer() - gap := s.Observe("BTC-USD", 100) - require.False(t, gap, "first sequence number must not be a gap") + require.False(t, s.Observe(100), "first sequence number must not be a gap") } func TestSequencer_ConsecutiveSequencesAreInOrder(t *testing.T) { s := coinbase.NewSequencer() for i := int64(50); i < 60; i++ { - gap := s.Observe("BTC-USD", i) - require.False(t, gap, "i=%d", i) + require.False(t, s.Observe(i), "i=%d", i) } } func TestSequencer_GapTriggersReportSignal(t *testing.T) { s := coinbase.NewSequencer() - require.False(t, s.Observe("BTC-USD", 50)) - require.False(t, s.Observe("BTC-USD", 51)) - require.True(t, s.Observe("BTC-USD", 60), "skipping 52-59 must report gap") + require.False(t, s.Observe(50)) + require.False(t, s.Observe(51)) + require.True(t, s.Observe(60), "skipping 52-59 must report gap") } -func TestSequencer_GapsArePerProduct(t *testing.T) { +func TestSequencer_ResetClearsBaseline(t *testing.T) { s := coinbase.NewSequencer() - require.False(t, s.Observe("BTC-USD", 100)) - require.False(t, s.Observe("ETH-USD", 200)) - require.False(t, s.Observe("BTC-USD", 101)) - require.True(t, s.Observe("ETH-USD", 250)) - require.False(t, s.Observe("BTC-USD", 102)) -} - -func TestSequencer_ResetClearsAllProducts(t *testing.T) { - s := coinbase.NewSequencer() - s.Observe("BTC-USD", 100) - s.Observe("ETH-USD", 200) + s.Observe(100) s.Reset() - require.False(t, s.Observe("BTC-USD", 9999)) - require.False(t, s.Observe("ETH-USD", 8888)) + require.False(t, s.Observe(9999)) } func TestSequencer_DuplicateSequenceTreatedAsGap(t *testing.T) { s := coinbase.NewSequencer() - require.False(t, s.Observe("BTC-USD", 100)) - require.True(t, s.Observe("BTC-USD", 100), "replaying the same seq is a gap") + require.False(t, s.Observe(100)) + require.True(t, s.Observe(100), "replaying the same seq is a gap") } func TestSequencer_BackwardSequenceTreatedAsGap(t *testing.T) { s := coinbase.NewSequencer() - require.False(t, s.Observe("BTC-USD", 100)) - require.True(t, s.Observe("BTC-USD", 90)) + require.False(t, s.Observe(100)) + require.True(t, s.Observe(90)) }