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