Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 18 additions & 4 deletions pkg/shardservice/service_read_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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{
Expand All @@ -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 {
Expand Down
29 changes: 22 additions & 7 deletions pkg/vm/engine/tae/db/test/db_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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 {
Expand Down
Loading