From 5c0d82b6ecf49b4026acd6e2f11f8db70d36641e Mon Sep 17 00:00:00 2001 From: XuPeng-SH Date: Thu, 23 Jul 2026 07:15:29 +0800 Subject: [PATCH 1/3] fix(tae): tolerate pruned catalog entries during replay --- pkg/vm/engine/tae/txn/txnimpl/replaystore.go | 11 +++++++ pkg/vm/engine/tae/txn/txnimpl/txn_test.go | 30 ++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go index e94079d60db5b..ca303e03501a2 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go +++ b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go @@ -109,10 +109,21 @@ func (store *replayTxnStore) registerPreparedDMLTables(txn txnif.AsyncTxn) { for _, record := range dirty.Tables { db, err := store.catalog.GetDatabaseByID(record.DbID) if err != nil { + // The checkpoint catalog is an end-of-checkpoint view. A WAL txn + // can therefore still refer to a database that was dropped before + // the checkpoint ended. There is no live table to fence in that case. + if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + continue + } panic(err) } table, err := db.GetTableEntryByID(record.ID) if err != nil { + // As with databases, checkpoint pruning can legitimately remove a + // dropped table while its transaction is still present in the WAL. + if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + continue + } panic(err) } table.RegisterReplayedPreparedDML(store.preparedTxnID) diff --git a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go index b2ddaba4aa158..45d572621ab96 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go +++ b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go @@ -141,6 +141,36 @@ func TestReplayPreparedRollbackReleasesAutoIncrementFence(t *testing.T) { assert.True(t, watermark.IsEmpty()) } +func TestReplayPreparedDMLSkipsCheckpointPrunedTables(t *testing.T) { + c := catalog.MockCatalog(nil) + defer c.Close() + mgr := txnbase.NewTxnManager(catalog.MockTxnStoreFactory(c), catalog.MockTxnFactory(c), types.NewMockHLCClock(1)) + mgr.Start(context.Background()) + defer mgr.Stop() + + setupTxn, err := mgr.StartTxn(nil) + assert.NoError(t, err) + dbEntry, err := c.CreateDBEntry("replay_pruned", "", "", setupTxn) + assert.NoError(t, err) + tableEntry, err := dbEntry.CreateTableEntry(catalog.MockSchemaAll(3, 1), setupTxn, nil) + assert.NoError(t, err) + assert.NoError(t, setupTxn.Commit(context.Background())) + + startTS := types.BuildTS(10, 0) + replayTxn := newPreparingEpochTestTxn(t, "replay-pruned", startTS, types.BuildTS(11, 0)) + replayTxn.GetMemo().AddTable(dbEntry.ID, tableEntry.ID) + replayTxn.GetMemo().AddTable(dbEntry.ID, tableEntry.ID+1) + replayTxn.GetMemo().AddTable(dbEntry.ID+1, tableEntry.ID+2) + store := &replayTxnStore{Cmd: &txnbase.TxnCmd{ComposedCmd: txnbase.NewComposedCmd()}, Observer: noopReplayObserver{}, catalog: c} + + assert.NoError(t, store.prepareCommit(replayTxn)) + assert.Len(t, store.preparedTables, 1) + assert.Same(t, tableEntry, store.preparedTables[tableEntry.ID]) + assert.True(t, tableEntry.ShouldRetryAutoIncrementAlter(startTS)) + assert.NoError(t, store.applyRollback(replayTxn)) + assert.False(t, tableEntry.ShouldRetryAutoIncrementAlter(startTS)) +} + type waitingSchemaTxn struct { txnif.TxnReader prepareTS types.TS From 8f9aae5086e0353d1f7a46eb2945eef0ae42c101 Mon Sep 17 00:00:00 2001 From: XuPeng-SH Date: Thu, 23 Jul 2026 07:39:20 +0800 Subject: [PATCH 2/3] fix(tae): use replayable catalog DDL in append GC test --- pkg/vm/engine/tae/db/test/db_test.go | 29 ++++++++++++++----- pkg/vm/engine/tae/txn/txnimpl/replaystore.go | 11 ------- pkg/vm/engine/tae/txn/txnimpl/txn_test.go | 30 -------------------- 3 files changed, 22 insertions(+), 48 deletions(-) diff --git a/pkg/vm/engine/tae/db/test/db_test.go b/pkg/vm/engine/tae/db/test/db_test.go index 5ab8f1b266ea0..256accccceb73 100644 --- a/pkg/vm/engine/tae/db/test/db_test.go +++ b/pkg/vm/engine/tae/db/test/db_test.go @@ -7077,13 +7077,13 @@ func TestAppendAndGC2(t *testing.T) { schema2.Extra.ObjectMaxBlocks = 2 { txn, _ := db.StartTxn(nil) - database, err := txn.CreateDatabase("db", "", "") - assert.Nil(t, err) - _, err = database.CreateRelation(schema1) - assert.Nil(t, err) - _, err = database.CreateRelation(schema2) - assert.Nil(t, err) - assert.Nil(t, txn.Commit(context.Background())) + database, err := testutil.CreateDatabase2(ctx, txn, "db") + require.NoError(t, err) + _, err = testutil.CreateRelation2(ctx, txn, database, schema1) + require.NoError(t, err) + _, err = testutil.CreateRelation2(ctx, txn, database, schema2) + require.NoError(t, err) + require.NoError(t, txn.Commit(ctx)) } bat := catalog.MockBatch(schema1, int(schema1.Extra.BlockMaxRows*10-1)) defer bat.Close() @@ -7110,6 +7110,21 @@ func TestAppendAndGC2(t *testing.T) { metaFile := db.BGCheckpointRunner.GetCheckpointMetaFiles() tae.Restart(ctx) db = tae.DB + func() { + replayTxn, err := db.StartTxn(nil) + require.NoError(t, err) + defer func() { + require.NoError(t, replayTxn.Rollback(ctx)) + }() + replayDB, err := replayTxn.GetDatabase("db") + require.NoError(t, err) + replayRel1, err := replayDB.GetRelationByName(schema1.Name) + require.NoError(t, err) + replayRel2, err := replayDB.GetRelationByName(schema2.Name) + require.NoError(t, err) + testutil.CheckAllColRowsByScan(t, replayRel1, bat.Length(), false) + testutil.CheckAllColRowsByScan(t, replayRel2, bat.Length(), false) + }() files := make(map[string]struct{}, 0) loadFiles := func(group uint32, lsn uint64, payload []byte, typ uint16, info any) driver.ReplayEntryState { if group != wal.GroupFiles { diff --git a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go index ca303e03501a2..e94079d60db5b 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go +++ b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go @@ -109,21 +109,10 @@ func (store *replayTxnStore) registerPreparedDMLTables(txn txnif.AsyncTxn) { for _, record := range dirty.Tables { db, err := store.catalog.GetDatabaseByID(record.DbID) if err != nil { - // The checkpoint catalog is an end-of-checkpoint view. A WAL txn - // can therefore still refer to a database that was dropped before - // the checkpoint ended. There is no live table to fence in that case. - if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { - continue - } panic(err) } table, err := db.GetTableEntryByID(record.ID) if err != nil { - // As with databases, checkpoint pruning can legitimately remove a - // dropped table while its transaction is still present in the WAL. - if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { - continue - } panic(err) } table.RegisterReplayedPreparedDML(store.preparedTxnID) diff --git a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go index 45d572621ab96..b2ddaba4aa158 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go +++ b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go @@ -141,36 +141,6 @@ func TestReplayPreparedRollbackReleasesAutoIncrementFence(t *testing.T) { assert.True(t, watermark.IsEmpty()) } -func TestReplayPreparedDMLSkipsCheckpointPrunedTables(t *testing.T) { - c := catalog.MockCatalog(nil) - defer c.Close() - mgr := txnbase.NewTxnManager(catalog.MockTxnStoreFactory(c), catalog.MockTxnFactory(c), types.NewMockHLCClock(1)) - mgr.Start(context.Background()) - defer mgr.Stop() - - setupTxn, err := mgr.StartTxn(nil) - assert.NoError(t, err) - dbEntry, err := c.CreateDBEntry("replay_pruned", "", "", setupTxn) - assert.NoError(t, err) - tableEntry, err := dbEntry.CreateTableEntry(catalog.MockSchemaAll(3, 1), setupTxn, nil) - assert.NoError(t, err) - assert.NoError(t, setupTxn.Commit(context.Background())) - - startTS := types.BuildTS(10, 0) - replayTxn := newPreparingEpochTestTxn(t, "replay-pruned", startTS, types.BuildTS(11, 0)) - replayTxn.GetMemo().AddTable(dbEntry.ID, tableEntry.ID) - replayTxn.GetMemo().AddTable(dbEntry.ID, tableEntry.ID+1) - replayTxn.GetMemo().AddTable(dbEntry.ID+1, tableEntry.ID+2) - store := &replayTxnStore{Cmd: &txnbase.TxnCmd{ComposedCmd: txnbase.NewComposedCmd()}, Observer: noopReplayObserver{}, catalog: c} - - assert.NoError(t, store.prepareCommit(replayTxn)) - assert.Len(t, store.preparedTables, 1) - assert.Same(t, tableEntry, store.preparedTables[tableEntry.ID]) - assert.True(t, tableEntry.ShouldRetryAutoIncrementAlter(startTS)) - assert.NoError(t, store.applyRollback(replayTxn)) - assert.False(t, tableEntry.ShouldRetryAutoIncrementAlter(startTS)) -} - type waitingSchemaTxn struct { txnif.TxnReader prepareTS types.TS From 670b1d66a139af342578467b3d139ad1070e75db Mon Sep 17 00:00:00 2001 From: XuPeng-SH Date: Thu, 23 Jul 2026 10:01:56 +0800 Subject: [PATCH 3/3] test(shardservice): wait for CN readiness before shard reads --- pkg/shardservice/service_read_test.go | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/pkg/shardservice/service_read_test.go b/pkg/shardservice/service_read_test.go index bf792768bd19b..13eed607d4e17 100644 --- a/pkg/shardservice/service_read_test.go +++ b/pkg/shardservice/service_read_test.go @@ -208,12 +208,23 @@ func TestReadValidatesAllRemoteShardsBeforeSending(t *testing.T) { runServicesTest( t, "cn1,cn2,cn3", - func(ctx context.Context, _ *server, services []*service) { + func(_ context.Context, server *server, services []*service) { s := services[0] table := uint64(1) - mustAddTestShards(t, ctx, s, table, 3, 1, services[1:]...) + require.Eventually(t, func() bool { + server.r.RLock() + defer server.r.RUnlock() + return len(server.r.cns) == len(services) + }, 10*time.Second, 10*time.Millisecond) + func() { + setupCtx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + mustAddTestShards(t, setupCtx, s, table, 3, 1, services[1:]...) + }() for _, service := range services { - waitReplicaCount(table, service, 1) + require.Eventually(t, func() bool { + return service.TableReplicaCount(table) == 1 + }, 10*time.Second, 10*time.Millisecond) } cache, err := s.getShards(table) @@ -253,8 +264,10 @@ func TestReadValidatesAllRemoteShardsBeforeSending(t *testing.T) { client := &countingMethodBasedClient{MethodBasedClient: s.remote.client} s.remote.client = client adjustCalls := make(map[uint64]int) + readCtx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() err = s.Read( - ctx, + readCtx, ReadRequest{ TableID: table, Param: shard.ReadParam{Process: pipeline.ProcessInfo{ @@ -266,6 +279,7 @@ func TestReadValidatesAllRemoteShardsBeforeSending(t *testing.T) { }), ) require.Error(t, err) + require.ErrorContains(t, err, "incompatible or unknown commit") require.Zero(t, client.asyncCalls.Load()) require.Len(t, adjustCalls, len(remoteTargets)) for _, calls := range adjustCalls {