fix(monitor/collectors/coinbase): connection-global sequencer; gap is non-fatal log

Two related issues caught in live verification:

1. Coinbase's sequence_num is connection-global (not per-product). Previous
   per-product tracking flagged every cross-product frame as a gap.

2. Even after switching to global tracking, gaps still occur regularly under
   normal operation (Coinbase's exact sequencing pattern across channels is
   not strictly seq+1 frame-to-frame). Treating gaps as fatal kills the
   ReadLoop in seconds, never giving the minute aggregator a chance to cross
   a boundary, so btc_eth_minute stays empty.

Resolution: rewrite Sequencer for a single global counter and downgrade
gap-detection to a non-fatal Warn. ReadLoop continues; (symbol, ts) PK on
btc_eth_ticks is the safety net for any duplicate replay. Includes a
diagnostic log on loop exit (handleConn) for ops visibility.
This commit is contained in:
CarterPerez-dev 2026-05-02 04:03:09 -04:00
parent 989e53c958
commit d4995c1258
5 changed files with 62 additions and 47 deletions

View File

@ -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 "<nil>"
}
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()

View File

@ -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,
)
}
}

View File

@ -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) {

View File

@ -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
}

View File

@ -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))
}