From 640ca41b316c1d3f5508da7bafe4bcf31de05041 Mon Sep 17 00:00:00 2001 From: nhsmw Date: Thu, 10 Sep 2026 17:48:27 +0800 Subject: [PATCH 1/4] Update event_store.go --- logservice/eventstore/event_store.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index d85319b97c..75c773302a 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -1521,6 +1521,9 @@ func (e *eventStore) writeEvents( } metrics.EventStoreWritePrepareDurationHistogram.Observe(time.Since(prepareStart).Seconds()) start := time.Now() + // Simulate slow EventStore storage so write workers remain occupied and + // incoming events queue up behind them. + failpoint.Inject("SlowEventStoreWrite", nil) err := batch.Commit(pebble.NoSync) metrics.EventStoreWriteDurationHistogram.Observe(time.Since(start).Seconds()) return err From 60e8f565aa3a6c3b1162eb5b1b1c619f4f71d28b Mon Sep 17 00:00:00 2001 From: wk989898 Date: Thu, 10 Sep 2026 10:28:35 +0000 Subject: [PATCH 2/4] chore Signed-off-by: wk989898 --- logservice/eventstore/event_store.go | 1 + 1 file changed, 1 insertion(+) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 75c773302a..0a171de62d 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -28,6 +28,7 @@ import ( "github.com/cockroachdb/pebble" "github.com/klauspost/compress/zstd" "github.com/pingcap/errors" + "github.com/pingcap/failpoint" "github.com/pingcap/log" "github.com/pingcap/ticdc/heartbeatpb" "github.com/pingcap/ticdc/logservice/logpuller" From fc5ed2c724042f91b20ffa509e08c5862333a7fa Mon Sep 17 00:00:00 2001 From: wk989898 Date: Fri, 11 Sep 2026 04:46:59 +0000 Subject: [PATCH 3/4] update Signed-off-by: wk989898 --- logservice/eventstore/event_store.go | 53 ++++++++++++- .../eventstore/event_store_bench_test.go | 4 +- logservice/eventstore/event_store_test.go | 77 ++++++++++++++++--- 3 files changed, 118 insertions(+), 16 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 0a171de62d..79a338928f 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -48,6 +48,7 @@ import ( "github.com/tikv/client-go/v2/oracle" "go.uber.org/zap" "golang.org/x/sync/errgroup" + "golang.org/x/time/rate" ) var ( @@ -241,6 +242,9 @@ type eventStore struct { tableCache *pebble.TableCache chs []*chann.UnlimitedChannel[eventWithCallback, uint64] writeTaskPools []*writeTaskPool + // Used only by SlowEventStoreWrite; shared by every DB and write worker. + slowWriteLimiterOnce sync.Once + slowWriteLimiter *rate.Limiter gcManager *gcManager @@ -388,7 +392,10 @@ func (p *writeTaskPool) run(ctx context.Context) { queueDuration.Observe(float64(time.Now().UnixNano()-events[0].enqueueTimeNano) / float64(time.Second)) } start := time.Now() - if err = p.store.writeEvents(p.db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if err = p.store.writeEvents(ctx, p.db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if ctx.Err() != nil { + return + } log.Panic("write events failed", zap.Error(err)) } ioWriteDuration.Observe(time.Since(start).Seconds()) @@ -1388,6 +1395,7 @@ func (e *eventStore) collectAndReportStoreMetrics() { } func (e *eventStore) writeEvents( + ctx context.Context, db *pebble.DB, events []eventWithCallback, encoder *zstd.Encoder, @@ -1522,14 +1530,51 @@ func (e *eventStore) writeEvents( } metrics.EventStoreWritePrepareDurationHistogram.Observe(time.Since(prepareStart).Seconds()) start := time.Now() - // Simulate slow EventStore storage so write workers remain occupied and - // incoming events queue up behind them. - failpoint.Inject("SlowEventStoreWrite", nil) + // return(bytesPerSecond) limits aggregate encoded batch throughput across + // all DBs and workers, allowing incoming events to build up naturally. + failpoint.Inject("SlowEventStoreWrite", func(val failpoint.Value) { + if bytesPerSecond, ok := val.(int); ok && bytesPerSecond > 0 && batch.Count() > 0 { + if err := e.waitForWriteBandwidth(ctx, batch.Len(), bytesPerSecond); err != nil { + failpoint.Return(err) + } + } + }) err := batch.Commit(pebble.NoSync) metrics.EventStoreWriteDurationHistogram.Observe(time.Since(start).Seconds()) return err } +func (e *eventStore) waitForWriteBandwidth(ctx context.Context, size, bytesPerSecond int) error { + const burst = 1 << 20 + e.slowWriteLimiterOnce.Do(func() { + e.slowWriteLimiter = rate.NewLimiter(rate.Limit(bytesPerSecond), burst) + // Charge the first batch too; idle time can later accumulate up to 1 MiB. + e.slowWriteLimiter.AllowN(time.Now(), burst) + }) + limiter := e.slowWriteLimiter + if limiter.Limit() != rate.Limit(bytesPerSecond) { + limiter.SetLimit(rate.Limit(bytesPerSecond)) + } + for size > 0 { + if err := ctx.Err(); err != nil { + return err + } + // Split large batches so reservations never exceed the limiter's burst. + n := min(size, burst, bytesPerSecond) + reservation := limiter.ReserveN(time.Now(), n) + timer := time.NewTimer(reservation.Delay()) + select { + case <-ctx.Done(): + timer.Stop() + reservation.Cancel() + return ctx.Err() + case <-timer.C: + } + size -= n + } + return nil +} + func encodeAndMaybeCompressValue( kv *common.RawKVEntry, encoder *zstd.Encoder, diff --git a/logservice/eventstore/event_store_bench_test.go b/logservice/eventstore/event_store_bench_test.go index c32a1ef884..2d696e25ce 100644 --- a/logservice/eventstore/event_store_bench_test.go +++ b/logservice/eventstore/event_store_bench_test.go @@ -108,7 +108,7 @@ func BenchmarkEventStoreWriteEvents(b *testing.B) { b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { - if err := store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if err := store.writeEvents(b.Context(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf); err != nil { b.Fatal(err) } } @@ -172,7 +172,7 @@ func BenchmarkEventStoreIteratorNext(b *testing.B) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - if err := store.writeEvents(db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if err := store.writeEvents(b.Context(), db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { b.Fatal(err) } diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index 8e5380b3d2..f53303d5de 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -26,6 +26,7 @@ import ( "github.com/cockroachdb/pebble" "github.com/klauspost/compress/zstd" + "github.com/pingcap/failpoint" "github.com/pingcap/ticdc/heartbeatpb" "github.com/pingcap/ticdc/logservice/logpuller" "github.com/pingcap/ticdc/pkg/common" @@ -353,7 +354,7 @@ func TestEventStoreUsesKeyspaceIDForEncryption(t *testing.T) { } var compressionBuf []byte var rawValueBuf []byte - err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = es.writeEvents(context.Background(), es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) require.Equal(t, uint32(42), spy.encryptKeyspaceID) require.Equal(t, 2, spy.encryptCalls) @@ -447,7 +448,7 @@ func TestEventStoreHandlesUnencryptedValuesFromEncryptionLayer(t *testing.T) { } var compressionBuf []byte var rawValueBuf []byte - err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = es.writeEvents(context.Background(), es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) subStat.resolvedTs.Store(largeKV.CRTs) @@ -1298,7 +1299,7 @@ func TestEventStoreRowLevelScanPositionSurvivesSubStatSwitch(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte writeRows := func(subStat *subscriptionStat) { - err := store.writeEvents(store.dbs[subStat.dbIndex], []eventWithCallback{{ + err := store.writeEvents(context.Background(), store.dbs[subStat.dbIndex], []eventWithCallback{{ subID: subStat.subID, tableID: tableID, kvs: rows, @@ -1422,7 +1423,7 @@ func TestWriteToEventStore(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) // Read events back and verify. @@ -1504,7 +1505,7 @@ func TestWriteToEventStoreZstdCompressionDisabled(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) iter, err := store.dbs[0].NewIter(&pebble.IterOptions{}) @@ -1612,7 +1613,7 @@ func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) afterMetric := testutil.ToFloat64(metrics.EventStoreCompressedRowsCount) require.InDelta(t, float64(len(expectedValues)), afterMetric-beforeMetric, 1e-9) @@ -1661,6 +1662,62 @@ func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { require.Equal(t, int64(len(expectedValues)), rowCount) } +func TestEventStoreSharedWriteBandwidth(t *testing.T) { + store := &eventStore{} + const bytesPerSecond = 10 << 20 + start := time.Now() + // Initialize without consuming bandwidth, so all workers share the same + // empty limiter before starting their reservations. + require.NoError(t, store.waitForWriteBandwidth(context.Background(), 0, bytesPerSecond)) + var wg sync.WaitGroup + for _, size := range []int{1 << 19, 1 << 20, 3 << 19} { + wg.Add(1) + go func() { + defer wg.Done() + assertErr := store.waitForWriteBandwidth(context.Background(), size, bytesPerSecond) + if assertErr != nil { + t.Errorf("wait for write bandwidth: %v", assertErr) + } + }() + } + wg.Wait() + // Three concurrent batches total 3 MiB, including one larger than burst. + // Independent per-worker limits would incorrectly finish sooner. + require.GreaterOrEqual(t, time.Since(start), 300*time.Millisecond) +} + +func TestSlowEventStoreWrite(t *testing.T) { + const failpointName = "github.com/pingcap/ticdc/logservice/eventstore/SlowEventStoreWrite" + require.NoError(t, failpoint.Enable(failpointName, "return(1)")) + t.Cleanup(func() { require.NoError(t, failpoint.Disable(failpointName)) }) + db, err := pebble.Open(t.TempDir(), &pebble.Options{}) + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, db.Close()) }) + store := &eventStore{} + events := []eventWithCallback{{ + subID: 1, + tableID: 1, + kvs: []common.RawKVEntry{{ + OpType: common.OpTypePut, StartTs: 1, CRTs: 2, + Key: []byte("key"), Value: []byte("value"), + }}, + }} + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + require.ErrorIs(t, store.writeEvents(ctx, db, events, nil, nil, nil), context.DeadlineExceeded) + iter, err := db.NewIter(nil) + require.NoError(t, err) + require.False(t, iter.First(), "canceled writes must not commit") + require.NoError(t, iter.Close()) + + require.NoError(t, failpoint.Disable(failpointName)) + require.NoError(t, store.writeEvents(context.Background(), db, events, nil, nil, nil)) + iter, err = db.NewIter(nil) + require.NoError(t, err) + require.True(t, iter.First(), "writes should resume after disabling the failpoint") + require.NoError(t, iter.Close()) +} + func TestEventStoreKVEntryCount(t *testing.T) { dir := t.TempDir() _, storeInt := newEventStoreForTest(dir) @@ -1691,7 +1748,7 @@ func TestEventStoreKVEntryCount(t *testing.T) { encoder, err := zstd.NewWriter(nil) require.NoError(t, err) defer encoder.Close() - require.NoError(t, store.writeEvents(store.dbs[0], events, encoder, nil, nil)) + require.NoError(t, store.writeEvents(context.Background(), store.dbs[0], events, encoder, nil, nil)) for i, metric := range entryMetrics { require.Equal(t, before[i]+1, testutil.ToFloat64(metric)) @@ -1727,7 +1784,7 @@ func TestEventStoreIterReadsLegacyCompressedValuesWithEncryptionManager(t *testi var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) innerIter, err := store.dbs[0].NewIter(&pebble.IterOptions{}) @@ -1802,7 +1859,7 @@ func TestEventStoreGetIteratorConcurrently(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.(*eventStore).writeEvents(store.(*eventStore).dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.(*eventStore).writeEvents(context.Background(), store.(*eventStore).dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) // 3. Advance resolved ts for the subscription. @@ -1890,7 +1947,7 @@ func TestEventStoreResumeTokenSupportsRowLevelResume(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(store.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(context.Background(), store.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) subStat.resolvedTs.Store(nextCommitTs) From 03c7a651961f3746f3665bc69a3ad8ca85d4306a Mon Sep 17 00:00:00 2001 From: wk989898 Date: Fri, 11 Sep 2026 06:41:00 +0000 Subject: [PATCH 4/4] Revert "update" This reverts commit fc5ed2c724042f91b20ffa509e08c5862333a7fa. --- logservice/eventstore/event_store.go | 53 +------------ .../eventstore/event_store_bench_test.go | 4 +- logservice/eventstore/event_store_test.go | 77 +++---------------- 3 files changed, 16 insertions(+), 118 deletions(-) diff --git a/logservice/eventstore/event_store.go b/logservice/eventstore/event_store.go index 79a338928f..0a171de62d 100644 --- a/logservice/eventstore/event_store.go +++ b/logservice/eventstore/event_store.go @@ -48,7 +48,6 @@ import ( "github.com/tikv/client-go/v2/oracle" "go.uber.org/zap" "golang.org/x/sync/errgroup" - "golang.org/x/time/rate" ) var ( @@ -242,9 +241,6 @@ type eventStore struct { tableCache *pebble.TableCache chs []*chann.UnlimitedChannel[eventWithCallback, uint64] writeTaskPools []*writeTaskPool - // Used only by SlowEventStoreWrite; shared by every DB and write worker. - slowWriteLimiterOnce sync.Once - slowWriteLimiter *rate.Limiter gcManager *gcManager @@ -392,10 +388,7 @@ func (p *writeTaskPool) run(ctx context.Context) { queueDuration.Observe(float64(time.Now().UnixNano()-events[0].enqueueTimeNano) / float64(time.Second)) } start := time.Now() - if err = p.store.writeEvents(ctx, p.db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { - if ctx.Err() != nil { - return - } + if err = p.store.writeEvents(p.db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { log.Panic("write events failed", zap.Error(err)) } ioWriteDuration.Observe(time.Since(start).Seconds()) @@ -1395,7 +1388,6 @@ func (e *eventStore) collectAndReportStoreMetrics() { } func (e *eventStore) writeEvents( - ctx context.Context, db *pebble.DB, events []eventWithCallback, encoder *zstd.Encoder, @@ -1530,51 +1522,14 @@ func (e *eventStore) writeEvents( } metrics.EventStoreWritePrepareDurationHistogram.Observe(time.Since(prepareStart).Seconds()) start := time.Now() - // return(bytesPerSecond) limits aggregate encoded batch throughput across - // all DBs and workers, allowing incoming events to build up naturally. - failpoint.Inject("SlowEventStoreWrite", func(val failpoint.Value) { - if bytesPerSecond, ok := val.(int); ok && bytesPerSecond > 0 && batch.Count() > 0 { - if err := e.waitForWriteBandwidth(ctx, batch.Len(), bytesPerSecond); err != nil { - failpoint.Return(err) - } - } - }) + // Simulate slow EventStore storage so write workers remain occupied and + // incoming events queue up behind them. + failpoint.Inject("SlowEventStoreWrite", nil) err := batch.Commit(pebble.NoSync) metrics.EventStoreWriteDurationHistogram.Observe(time.Since(start).Seconds()) return err } -func (e *eventStore) waitForWriteBandwidth(ctx context.Context, size, bytesPerSecond int) error { - const burst = 1 << 20 - e.slowWriteLimiterOnce.Do(func() { - e.slowWriteLimiter = rate.NewLimiter(rate.Limit(bytesPerSecond), burst) - // Charge the first batch too; idle time can later accumulate up to 1 MiB. - e.slowWriteLimiter.AllowN(time.Now(), burst) - }) - limiter := e.slowWriteLimiter - if limiter.Limit() != rate.Limit(bytesPerSecond) { - limiter.SetLimit(rate.Limit(bytesPerSecond)) - } - for size > 0 { - if err := ctx.Err(); err != nil { - return err - } - // Split large batches so reservations never exceed the limiter's burst. - n := min(size, burst, bytesPerSecond) - reservation := limiter.ReserveN(time.Now(), n) - timer := time.NewTimer(reservation.Delay()) - select { - case <-ctx.Done(): - timer.Stop() - reservation.Cancel() - return ctx.Err() - case <-timer.C: - } - size -= n - } - return nil -} - func encodeAndMaybeCompressValue( kv *common.RawKVEntry, encoder *zstd.Encoder, diff --git a/logservice/eventstore/event_store_bench_test.go b/logservice/eventstore/event_store_bench_test.go index 2d696e25ce..c32a1ef884 100644 --- a/logservice/eventstore/event_store_bench_test.go +++ b/logservice/eventstore/event_store_bench_test.go @@ -108,7 +108,7 @@ func BenchmarkEventStoreWriteEvents(b *testing.B) { b.ReportAllocs() b.ResetTimer() for i := 0; i < b.N; i++ { - if err := store.writeEvents(b.Context(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if err := store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf); err != nil { b.Fatal(err) } } @@ -172,7 +172,7 @@ func BenchmarkEventStoreIteratorNext(b *testing.B) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - if err := store.writeEvents(b.Context(), db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { + if err := store.writeEvents(db, events, encoder, &compressionBuf, &rawValueBuf); err != nil { b.Fatal(err) } diff --git a/logservice/eventstore/event_store_test.go b/logservice/eventstore/event_store_test.go index f53303d5de..8e5380b3d2 100644 --- a/logservice/eventstore/event_store_test.go +++ b/logservice/eventstore/event_store_test.go @@ -26,7 +26,6 @@ import ( "github.com/cockroachdb/pebble" "github.com/klauspost/compress/zstd" - "github.com/pingcap/failpoint" "github.com/pingcap/ticdc/heartbeatpb" "github.com/pingcap/ticdc/logservice/logpuller" "github.com/pingcap/ticdc/pkg/common" @@ -354,7 +353,7 @@ func TestEventStoreUsesKeyspaceIDForEncryption(t *testing.T) { } var compressionBuf []byte var rawValueBuf []byte - err = es.writeEvents(context.Background(), es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) require.Equal(t, uint32(42), spy.encryptKeyspaceID) require.Equal(t, 2, spy.encryptCalls) @@ -448,7 +447,7 @@ func TestEventStoreHandlesUnencryptedValuesFromEncryptionLayer(t *testing.T) { } var compressionBuf []byte var rawValueBuf []byte - err = es.writeEvents(context.Background(), es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = es.writeEvents(es.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) subStat.resolvedTs.Store(largeKV.CRTs) @@ -1299,7 +1298,7 @@ func TestEventStoreRowLevelScanPositionSurvivesSubStatSwitch(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte writeRows := func(subStat *subscriptionStat) { - err := store.writeEvents(context.Background(), store.dbs[subStat.dbIndex], []eventWithCallback{{ + err := store.writeEvents(store.dbs[subStat.dbIndex], []eventWithCallback{{ subID: subStat.subID, tableID: tableID, kvs: rows, @@ -1423,7 +1422,7 @@ func TestWriteToEventStore(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) // Read events back and verify. @@ -1505,7 +1504,7 @@ func TestWriteToEventStoreZstdCompressionDisabled(t *testing.T) { var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) iter, err := store.dbs[0].NewIter(&pebble.IterOptions{}) @@ -1613,7 +1612,7 @@ func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) afterMetric := testutil.ToFloat64(metrics.EventStoreCompressedRowsCount) require.InDelta(t, float64(len(expectedValues)), afterMetric-beforeMetric, 1e-9) @@ -1662,62 +1661,6 @@ func TestEventStoreCompressionAndIterDecodeBufferReuse(t *testing.T) { require.Equal(t, int64(len(expectedValues)), rowCount) } -func TestEventStoreSharedWriteBandwidth(t *testing.T) { - store := &eventStore{} - const bytesPerSecond = 10 << 20 - start := time.Now() - // Initialize without consuming bandwidth, so all workers share the same - // empty limiter before starting their reservations. - require.NoError(t, store.waitForWriteBandwidth(context.Background(), 0, bytesPerSecond)) - var wg sync.WaitGroup - for _, size := range []int{1 << 19, 1 << 20, 3 << 19} { - wg.Add(1) - go func() { - defer wg.Done() - assertErr := store.waitForWriteBandwidth(context.Background(), size, bytesPerSecond) - if assertErr != nil { - t.Errorf("wait for write bandwidth: %v", assertErr) - } - }() - } - wg.Wait() - // Three concurrent batches total 3 MiB, including one larger than burst. - // Independent per-worker limits would incorrectly finish sooner. - require.GreaterOrEqual(t, time.Since(start), 300*time.Millisecond) -} - -func TestSlowEventStoreWrite(t *testing.T) { - const failpointName = "github.com/pingcap/ticdc/logservice/eventstore/SlowEventStoreWrite" - require.NoError(t, failpoint.Enable(failpointName, "return(1)")) - t.Cleanup(func() { require.NoError(t, failpoint.Disable(failpointName)) }) - db, err := pebble.Open(t.TempDir(), &pebble.Options{}) - require.NoError(t, err) - t.Cleanup(func() { require.NoError(t, db.Close()) }) - store := &eventStore{} - events := []eventWithCallback{{ - subID: 1, - tableID: 1, - kvs: []common.RawKVEntry{{ - OpType: common.OpTypePut, StartTs: 1, CRTs: 2, - Key: []byte("key"), Value: []byte("value"), - }}, - }} - ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) - defer cancel() - require.ErrorIs(t, store.writeEvents(ctx, db, events, nil, nil, nil), context.DeadlineExceeded) - iter, err := db.NewIter(nil) - require.NoError(t, err) - require.False(t, iter.First(), "canceled writes must not commit") - require.NoError(t, iter.Close()) - - require.NoError(t, failpoint.Disable(failpointName)) - require.NoError(t, store.writeEvents(context.Background(), db, events, nil, nil, nil)) - iter, err = db.NewIter(nil) - require.NoError(t, err) - require.True(t, iter.First(), "writes should resume after disabling the failpoint") - require.NoError(t, iter.Close()) -} - func TestEventStoreKVEntryCount(t *testing.T) { dir := t.TempDir() _, storeInt := newEventStoreForTest(dir) @@ -1748,7 +1691,7 @@ func TestEventStoreKVEntryCount(t *testing.T) { encoder, err := zstd.NewWriter(nil) require.NoError(t, err) defer encoder.Close() - require.NoError(t, store.writeEvents(context.Background(), store.dbs[0], events, encoder, nil, nil)) + require.NoError(t, store.writeEvents(store.dbs[0], events, encoder, nil, nil)) for i, metric := range entryMetrics { require.Equal(t, before[i]+1, testutil.ToFloat64(metric)) @@ -1784,7 +1727,7 @@ func TestEventStoreIterReadsLegacyCompressedValuesWithEncryptionManager(t *testi var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(context.Background(), store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(store.dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) innerIter, err := store.dbs[0].NewIter(&pebble.IterOptions{}) @@ -1859,7 +1802,7 @@ func TestEventStoreGetIteratorConcurrently(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.(*eventStore).writeEvents(context.Background(), store.(*eventStore).dbs[0], events, encoder, &compressionBuf, &rawValueBuf) + err = store.(*eventStore).writeEvents(store.(*eventStore).dbs[0], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) // 3. Advance resolved ts for the subscription. @@ -1947,7 +1890,7 @@ func TestEventStoreResumeTokenSupportsRowLevelResume(t *testing.T) { defer encoder.Close() var compressionBuf []byte var rawValueBuf []byte - err = store.writeEvents(context.Background(), store.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) + err = store.writeEvents(store.dbs[subStat.dbIndex], events, encoder, &compressionBuf, &rawValueBuf) require.NoError(t, err) subStat.resolvedTs.Store(nextCommitTs)