diff --git a/packages/api/internal/ingest/ingest.go b/packages/api/internal/ingest/ingest.go index c4bc789..e12bec3 100644 --- a/packages/api/internal/ingest/ingest.go +++ b/packages/api/internal/ingest/ingest.go @@ -90,21 +90,6 @@ func (r *Runner) Run(ctx context.Context) error { continue } - if r.shouldSkipEvent(event.ID) { - if err := r.tap.AckEvent(ctx, event.ID); err != nil { - r.log.Warn("tap ack skipped event failed", - slog.Int64("event_id", event.ID), - slog.String("error", err.Error()), - ) - continue - } - r.statusMu.Lock() - highWaterMark := r.highWaterMark - r.statusMu.Unlock() - r.log.Debug("skipped previously-processed event", slog.Int64("event_id", event.ID), slog.Int64("resume_cursor", highWaterMark)) - continue - } - if err := r.processWithRetry(ctx, event); err != nil { if ctx.Err() != nil { return nil @@ -137,16 +122,10 @@ func (r *Runner) initializeCursor(ctx context.Context) error { r.highWaterMark = cursor r.lastCursor = state.Cursor r.statusMu.Unlock() - r.log.Info("indexer cursor resume enabled", slog.Int64("resume_cursor", cursor)) + r.log.Info("indexer cursor state loaded", slog.Int64("cursor", cursor)) return nil } -func (r *Runner) shouldSkipEvent(eventID int64) bool { - r.statusMu.Lock() - defer r.statusMu.Unlock() - return r.highWaterMark > 0 && eventID <= r.highWaterMark -} - func (r *Runner) processWithRetry(ctx context.Context, event normalize.TapRecordEvent) error { attempt := 0 for { diff --git a/packages/api/internal/ingest/ingest_test.go b/packages/api/internal/ingest/ingest_test.go index ab5f420..d7f1ee0 100644 --- a/packages/api/internal/ingest/ingest_test.go +++ b/packages/api/internal/ingest/ingest_test.go @@ -305,8 +305,8 @@ func TestRunner_CursorHighWaterMarkDoesNotRegress(t *testing.T) { if st.syncCursor != "200" { t.Fatalf("cursor regressed: got %q want 200", st.syncCursor) } - if !r.shouldSkipEvent(150) { - t.Fatal("expected older event id to be skipped once a newer cursor is recorded") + if r.highWaterMark != 200 { + t.Fatalf("high-water mark: got %d want 200", r.highWaterMark) } } @@ -320,11 +320,8 @@ func TestRunner_InitializeCursorUsesHighWaterMark(t *testing.T) { t.Fatalf("initialize cursor: %v", err) } - if !r.shouldSkipEvent(150) { - t.Fatal("expected stored cursor to act as skip high-water mark") - } - if r.shouldSkipEvent(151) { - t.Fatal("did not expect events above the high-water mark to be skipped") + if r.highWaterMark != 150 { + t.Fatalf("high-water mark: got %d want 150", r.highWaterMark) } } @@ -401,29 +398,6 @@ func TestAllowlistMatching(t *testing.T) { } } -func TestRunner_InitializeCursorResume(t *testing.T) { - st := newFakeStore() - st.initialSync = &store.SyncState{ConsumerName: "indexer-tap-v1", Cursor: "150"} - tap := &fakeTapClient{} - r := newRunnerForTest(st, tap, "sh.tangled.*") - - if err := r.initializeCursor(context.Background()); err != nil { - t.Fatalf("initialize cursor: %v", err) - } - if r.highWaterMark != 150 { - t.Fatalf("high-water mark: got %d want 150", r.highWaterMark) - } - if !r.shouldSkipEvent(149) { - t.Fatalf("expected event 149 to be skipped") - } - if !r.shouldSkipEvent(150) { - t.Fatalf("expected event 150 to be skipped") - } - if r.shouldSkipEvent(151) { - t.Fatalf("expected event 151 to be processed") - } -} - func TestRunner_PersistCursorBeforeAck(t *testing.T) { st := newFakeStore() tap := &fakeTapClient{}