498b235461
Build and test / Build and test AMD64 Ubuntu 22.04 (push) Failing after 0s
Publish Builder / amazonlinux2023 (push) Failing after 1s
Build and test / UT for Go (push) Has been skipped
Publish KRTE Images / KRTE (push) Failing after 1s
Build and test / Integration Test (push) Has been skipped
Build and test / Upload Code Coverage (push) Has been skipped
Publish Builder / rockylinux9 (push) Failing after 1s
Publish Builder / ubuntu22.04 (push) Failing after 0s
Publish Builder / ubuntu24.04 (push) Failing after 0s
Publish Gpu Builder / publish-gpu-builder (push) Failing after 1s
Publish Test Images / PyTest (push) Failing after 0s
Build and test / UT for Cpp (push) Has been cancelled
941 lines
33 KiB
Go
941 lines
33 KiB
Go
// Licensed to the LF AI & Data foundation under one
|
|
// or more contributor license agreements. See the NOTICE file
|
|
// distributed with this work for additional information
|
|
// regarding copyright ownership. The ASF licenses this file
|
|
// to you under the Apache License, Version 2.0 (the
|
|
// "License"); you may not use this file except in compliance
|
|
// with the License. You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
package datacoord
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/cockroachdb/errors"
|
|
"github.com/stretchr/testify/mock"
|
|
"github.com/stretchr/testify/suite"
|
|
"go.uber.org/atomic"
|
|
|
|
"github.com/milvus-io/milvus-proto/go-api/v3/commonpb"
|
|
"github.com/milvus-io/milvus-proto/go-api/v3/schemapb"
|
|
"github.com/milvus-io/milvus/internal/datacoord/allocator"
|
|
"github.com/milvus-io/milvus/internal/datacoord/broker"
|
|
"github.com/milvus-io/milvus/internal/datacoord/session"
|
|
"github.com/milvus-io/milvus/internal/metastore/kv/datacoord"
|
|
"github.com/milvus-io/milvus/internal/metastore/mocks"
|
|
"github.com/milvus-io/milvus/internal/storage"
|
|
"github.com/milvus-io/milvus/internal/storagev2/packed"
|
|
"github.com/milvus-io/milvus/pkg/v3/objectstorage"
|
|
"github.com/milvus-io/milvus/pkg/v3/proto/datapb"
|
|
taskcommon "github.com/milvus-io/milvus/pkg/v3/taskcommon"
|
|
"github.com/milvus-io/milvus/pkg/v3/util/merr"
|
|
)
|
|
|
|
func TestBumpSchemaVersionCompactionTaskSuite(t *testing.T) {
|
|
suite.Run(t, new(BumpSchemaVersionCompactionTaskSuite))
|
|
}
|
|
|
|
type BumpSchemaVersionCompactionTaskSuite struct {
|
|
suite.Suite
|
|
|
|
mockID atomic.Int64
|
|
mockAlloc *allocator.MockAllocator
|
|
meta *meta
|
|
handler *NMockHandler
|
|
ievm IndexEngineVersionManager
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) SetupTest() {
|
|
ctx := context.Background()
|
|
cm := storage.NewLocalChunkManager(objectstorage.RootPath(""))
|
|
catalog := datacoord.NewCatalog(NewMetaMemoryKV(), "", "")
|
|
broker := broker.NewMockBroker(s.T())
|
|
broker.EXPECT().ShowCollectionIDs(mock.Anything).Return(nil, nil)
|
|
meta, err := newMeta(ctx, catalog, cm, broker)
|
|
s.NoError(err)
|
|
s.meta = meta
|
|
|
|
s.mockID.Store(time.Now().UnixMilli())
|
|
s.mockAlloc = allocator.NewMockAllocator(s.T())
|
|
s.mockAlloc.EXPECT().AllocN(mock.Anything).RunAndReturn(func(x int64) (int64, int64, error) {
|
|
start := s.mockID.Load()
|
|
end := s.mockID.Add(x)
|
|
return start, end, nil
|
|
}).Maybe()
|
|
s.mockAlloc.EXPECT().AllocID(mock.Anything).RunAndReturn(func(ctx context.Context) (int64, error) {
|
|
end := s.mockID.Add(1)
|
|
return end, nil
|
|
}).Maybe()
|
|
|
|
s.handler = NewNMockHandler(s.T())
|
|
s.handler.EXPECT().GetCollection(mock.Anything, mock.Anything).Return(&collectionInfo{}, nil).Maybe()
|
|
s.ievm = newIndexEngineVersionManager()
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) SetupSubTest() {
|
|
s.SetupTest()
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) generateBasicTask() *bumpSchemaVersionTask {
|
|
schema := &schemapb.CollectionSchema{
|
|
Name: "test_schema_bump_collection",
|
|
Description: "test collection for schema bump compaction",
|
|
Version: 2,
|
|
Fields: []*schemapb.FieldSchema{
|
|
{
|
|
FieldID: 100,
|
|
Name: "pk",
|
|
IsPrimaryKey: true,
|
|
DataType: schemapb.DataType_Int64,
|
|
AutoID: true,
|
|
},
|
|
{
|
|
FieldID: 101,
|
|
Name: "text",
|
|
DataType: schemapb.DataType_VarChar,
|
|
},
|
|
{
|
|
FieldID: 102,
|
|
Name: "sparse_vector",
|
|
DataType: schemapb.DataType_SparseFloatVector,
|
|
},
|
|
},
|
|
}
|
|
|
|
compactionTask := &datapb.CompactionTask{
|
|
PlanID: 1,
|
|
TriggerID: 19530,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
Type: datapb.CompactionType_BumpSchemaVersionCompaction,
|
|
NodeID: 1,
|
|
State: datapb.CompactionTaskState_pipelining,
|
|
Schema: schema,
|
|
InputSegments: []int64{101},
|
|
ResultSegments: []int64{1000},
|
|
PreAllocatedSegmentIDs: &datapb.IDRange{
|
|
Begin: 1000,
|
|
End: 2000,
|
|
},
|
|
Channel: "ch-1",
|
|
}
|
|
|
|
task := newBumpSchemaVersionTask(compactionTask, s.mockAlloc, s.meta, s.ievm)
|
|
return task
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestBumpSchemaVersionCompactionTaskBasic() {
|
|
task := s.generateBasicTask()
|
|
|
|
// Test basic getters
|
|
s.Equal(int64(1), task.GetTaskID())
|
|
s.Equal(taskcommon.Compaction, task.GetTaskType())
|
|
s.Equal(int64(1), task.GetSlotUsage())
|
|
s.Equal("10-ch-1", task.GetLabel())
|
|
s.Equal(int64(0), task.GetTaskVersion())
|
|
|
|
// Test task proto
|
|
taskProto := task.GetTaskProto()
|
|
s.NotNil(taskProto)
|
|
s.Equal(int64(1), taskProto.GetPlanID())
|
|
s.Equal(int64(19530), taskProto.GetTriggerID())
|
|
s.Equal(int64(1), taskProto.GetCollectionID())
|
|
s.Equal(int64(10), taskProto.GetPartitionID())
|
|
s.Equal(datapb.CompactionType_BumpSchemaVersionCompaction, taskProto.GetType())
|
|
s.Equal(datapb.CompactionTaskState_pipelining, taskProto.GetState())
|
|
s.Equal([]int64{101}, taskProto.GetInputSegments())
|
|
s.Equal([]int64{1000}, taskProto.GetResultSegments())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestBuildCompactionRequest() {
|
|
// Add a segment to meta
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
Binlogs: []*datapb.FieldBinlog{
|
|
{
|
|
FieldID: 101,
|
|
Binlogs: []*datapb.Binlog{
|
|
{LogID: 1000, EntriesNum: 1000},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
|
|
// Build compaction request
|
|
plan, err := task.BuildCompactionRequest()
|
|
s.NoError(err)
|
|
s.NotNil(plan)
|
|
|
|
// Verify plan
|
|
s.Equal(int64(1), plan.GetPlanID())
|
|
s.Equal(datapb.CompactionType_BumpSchemaVersionCompaction, plan.GetType())
|
|
s.Equal("ch-1", plan.GetChannel())
|
|
s.Equal(1, len(plan.GetSegmentBinlogs()))
|
|
s.Equal(segmentID, plan.GetSegmentBinlogs()[0].GetSegmentID())
|
|
s.Equal(int64(1), plan.GetSegmentBinlogs()[0].GetCollectionID())
|
|
s.Equal(int64(10), plan.GetSegmentBinlogs()[0].GetPartitionID())
|
|
s.Require().NotNil(plan.GetSchema())
|
|
s.Equal(task.GetTaskProto().GetSchema().GetVersion(), plan.GetSchema().GetVersion())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestBuildCompactionRequestCarriesV3ManifestAndPreAllocatedLogs() {
|
|
segmentID := int64(101)
|
|
manifest := "manifest-v3"
|
|
commitTimestamp := uint64(5000)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
StorageVersion: storage.StorageV3,
|
|
ManifestPath: manifest,
|
|
CommitTimestamp: commitTimestamp,
|
|
Binlogs: []*datapb.FieldBinlog{
|
|
{
|
|
FieldID: 101,
|
|
Binlogs: []*datapb.Binlog{
|
|
{LogID: 1000, EntriesNum: 1000},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
plan, err := task.BuildCompactionRequest()
|
|
s.NoError(err)
|
|
s.Require().Len(plan.GetSegmentBinlogs(), 1)
|
|
s.EqualValues(storage.StorageV3, plan.GetSegmentBinlogs()[0].GetStorageVersion())
|
|
s.Equal(manifest, plan.GetSegmentBinlogs()[0].GetManifest())
|
|
s.Equal(commitTimestamp, plan.GetSegmentBinlogs()[0].GetCommitTimestamp())
|
|
s.Require().NotNil(plan.GetSchema())
|
|
s.Equal(task.GetTaskProto().GetSchema().GetVersion(), plan.GetSchema().GetVersion())
|
|
s.Equal(task.GetTaskProto().GetPreAllocatedSegmentIDs(), plan.GetPreAllocatedSegmentIDs())
|
|
s.Require().NotNil(plan.GetPreAllocatedLogIDs())
|
|
s.Greater(plan.GetPreAllocatedLogIDs().GetEnd(), plan.GetPreAllocatedLogIDs().GetBegin())
|
|
s.Equal(plan.GetPreAllocatedLogIDs().GetBegin(), plan.GetBeginLogID())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestBuildCompactionRequestPreAllocateLogIDsError() {
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
StorageVersion: storage.StorageV3,
|
|
ManifestPath: "manifest-v3",
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
allocErr := errors.New("alloc failed")
|
|
mockAlloc := allocator.NewMockAllocator(s.T())
|
|
mockAlloc.EXPECT().AllocN(mock.Anything).Return(int64(0), int64(0), allocErr).Once()
|
|
task := newBumpSchemaVersionTask(s.generateBasicTask().GetTaskProto(), mockAlloc, s.meta, s.ievm)
|
|
|
|
plan, err := task.BuildCompactionRequest()
|
|
s.ErrorIs(err, allocErr)
|
|
s.Nil(plan)
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestBuildCompactionRequestSegmentNotFound() {
|
|
task := s.generateBasicTask()
|
|
|
|
// Try to build compaction request without adding segment to meta
|
|
plan, err := task.BuildCompactionRequest()
|
|
s.Error(err)
|
|
s.Nil(plan)
|
|
s.Contains(err.Error(), "segment not found")
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestCreateTaskOnWorker() {
|
|
s.Run("CreateTaskOnWorker fail, segment not found", func() {
|
|
task := s.generateBasicTask()
|
|
cluster := session.NewMockCluster(s.T())
|
|
task.CreateTaskOnWorker(1, cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("CreateTaskOnWorker fail, CreateCompaction error", func() {
|
|
// Add segment to meta
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
Binlogs: []*datapb.FieldBinlog{
|
|
{
|
|
FieldID: 101,
|
|
Binlogs: []*datapb.Binlog{
|
|
{LogID: 1000, EntriesNum: 1000},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().CreateCompaction(mock.Anything, mock.Anything, mock.Anything).Return(merr.WrapErrNodeNotFound(1))
|
|
task.CreateTaskOnWorker(1, cluster)
|
|
// Should remain in pipelining state when CreateCompaction fails
|
|
s.Equal(datapb.CompactionTaskState_pipelining, task.GetTaskProto().GetState())
|
|
// NodeID should be set to NullNodeID (-1) when CreateCompaction fails
|
|
s.Equal(int64(-1), task.GetTaskProto().GetNodeID())
|
|
})
|
|
|
|
s.Run("CreateTaskOnWorker succeed", func() {
|
|
// Add segment to meta
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
Binlogs: []*datapb.FieldBinlog{
|
|
{
|
|
FieldID: 101,
|
|
Binlogs: []*datapb.Binlog{
|
|
{LogID: 1000, EntriesNum: 1000},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().CreateCompaction(mock.Anything, mock.Anything, mock.Anything).Return(nil)
|
|
task.CreateTaskOnWorker(1, cluster)
|
|
s.Equal(datapb.CompactionTaskState_executing, task.GetTaskProto().GetState())
|
|
s.Equal(int64(1), task.GetTaskProto().GetNodeID())
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestQueryTaskOnWorker() {
|
|
s.Run("QueryTaskOnWorker, node not found", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(nil, merr.WrapErrNodeNotFound(1)).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_pipelining, task.GetTaskProto().GetState())
|
|
s.Equal(int64(-1), task.GetTaskProto().GetNodeID())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, result is nil", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(nil, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
// State should remain unchanged when result is nil
|
|
s.Equal(datapb.CompactionTaskState_executing, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed with empty segments", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{},
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, pipelining state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_pipelining,
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
// State should remain unchanged
|
|
s.Equal(datapb.CompactionTaskState_executing, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, executing state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_executing,
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
// State should remain unchanged
|
|
s.Equal(datapb.CompactionTaskState_executing, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, timeout state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_timeout,
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_timeout, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, failed state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_failed,
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed with ValidateSegmentState error (segment not in meta)", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
// segment 101 is NOT in meta → ValidateSegmentStateBeforeCompleteCompactionMutation returns error
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{{SegmentID: 101}},
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed with ErrIllegalCompactionPlan from saveSegmentMeta", func() {
|
|
// Add segment 101 so ValidateSegmentState passes, but result segmentID mismatches → ErrIllegalCompactionPlan
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{{SegmentID: 999}}, // mismatched → ErrIllegalCompactionPlan
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed success path", func() {
|
|
segmentID := int64(101)
|
|
manifest := "manifest-v3"
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
SchemaVersion: 1,
|
|
StorageVersion: storage.StorageV3,
|
|
ManifestPath: manifest,
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{{
|
|
SegmentID: segmentID,
|
|
InsertLogs: []*datapb.FieldBinlog{},
|
|
Manifest: manifest,
|
|
StorageVersion: storage.StorageV3,
|
|
}},
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_completed, task.GetTaskProto().GetState())
|
|
s.Equal([]int64{segmentID}, task.GetTaskProto().GetResultSegments())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed replacement success path", func() {
|
|
segmentID := int64(101)
|
|
newSegmentID := int64(1000)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
SchemaVersion: 1,
|
|
StorageVersion: storage.StorageV3,
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(
|
|
setState(datapb.CompactionTaskState_executing),
|
|
setNodeID(1),
|
|
func(t *datapb.CompactionTask) {
|
|
t.PreAllocatedSegmentIDs = &datapb.IDRange{Begin: newSegmentID, End: newSegmentID + 1}
|
|
},
|
|
))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{{
|
|
SegmentID: newSegmentID,
|
|
NumOfRows: 5,
|
|
InsertLogs: []*datapb.FieldBinlog{{FieldID: 101, Binlogs: []*datapb.Binlog{{LogID: 1001}}}},
|
|
Manifest: "replacement-manifest-v3",
|
|
StorageVersion: storage.StorageV3,
|
|
}},
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_completed, task.GetTaskProto().GetState())
|
|
s.Equal([]int64{newSegmentID}, task.GetTaskProto().GetResultSegments())
|
|
s.Equal(commonpb.SegmentState_Dropped, s.meta.GetSegment(context.TODO(), segmentID).GetState())
|
|
s.Equal(commonpb.SegmentState_Flushed, s.meta.GetSegment(context.TODO(), newSegmentID).GetState())
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, completed invalid manifest marks failed", func() {
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
SchemaVersion: 1,
|
|
StorageVersion: storage.StorageV3,
|
|
ManifestPath: "manifest-v3",
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_completed,
|
|
Segments: []*datapb.CompactionSegment{{
|
|
SegmentID: segmentID,
|
|
StorageVersion: storage.StorageV3,
|
|
}},
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
s.Contains(task.GetTaskProto().GetFailReason(), "StorageV3 manifest")
|
|
})
|
|
|
|
s.Run("QueryTaskOnWorker, default unknown state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().QueryCompaction(mock.Anything, mock.Anything).Return(&datapb.CompactionPlanResult{
|
|
State: datapb.CompactionTaskState_unknown,
|
|
}, nil).Once()
|
|
task.QueryTaskOnWorker(cluster)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestProcess() {
|
|
s.Run("Process meta_saved state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_meta_saved)))
|
|
result := task.Process()
|
|
// processMetaSaved should transition to completed and return true
|
|
s.True(result)
|
|
s.Equal(datapb.CompactionTaskState_completed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("Process completed state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_completed)))
|
|
result := task.Process()
|
|
s.True(result)
|
|
s.Equal(datapb.CompactionTaskState_completed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("Process failed state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_failed)))
|
|
result := task.Process()
|
|
s.True(result)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("Process timeout state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_timeout)))
|
|
result := task.Process()
|
|
s.True(result)
|
|
s.Equal(datapb.CompactionTaskState_timeout, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("Process other states return false", func() {
|
|
testStates := []datapb.CompactionTaskState{
|
|
datapb.CompactionTaskState_pipelining,
|
|
datapb.CompactionTaskState_executing,
|
|
datapb.CompactionTaskState_unknown,
|
|
}
|
|
for _, state := range testStates {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(state)))
|
|
result := task.Process()
|
|
s.False(result, "state %s should return false", state.String())
|
|
}
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestClean() {
|
|
task := s.generateBasicTask()
|
|
// Mark segment as compacting
|
|
s.meta.SetSegmentsCompacting(context.TODO(), []int64{101}, true)
|
|
|
|
// Add segment to meta
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
|
|
result := task.Clean()
|
|
s.True(result)
|
|
s.Equal(datapb.CompactionTaskState_cleaned, task.GetTaskProto().GetState())
|
|
|
|
// Verify segment compacting flag is reset
|
|
seg := s.meta.GetSegment(context.TODO(), segmentID)
|
|
s.NotNil(seg)
|
|
s.False(seg.isCompacting)
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestNeedReAssignNodeID() {
|
|
s.Run("NeedReAssignNodeID, pipelining with nodeID 0", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_pipelining), setNodeID(0)))
|
|
s.True(task.NeedReAssignNodeID())
|
|
})
|
|
|
|
s.Run("NeedReAssignNodeID, pipelining with NullNodeID", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_pipelining), setNodeID(-1)))
|
|
s.True(task.NeedReAssignNodeID())
|
|
})
|
|
|
|
s.Run("NeedReAssignNodeID, pipelining with valid nodeID", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_pipelining), setNodeID(1)))
|
|
s.False(task.NeedReAssignNodeID())
|
|
})
|
|
|
|
s.Run("NeedReAssignNodeID, non-pipelining state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(0)))
|
|
s.False(task.NeedReAssignNodeID())
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestDropTaskOnWorker() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing), setNodeID(1)))
|
|
cluster := session.NewMockCluster(s.T())
|
|
cluster.EXPECT().DropCompaction(mock.Anything, mock.Anything).Return(nil).Once()
|
|
task.DropTaskOnWorker(cluster)
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestSetNodeID() {
|
|
task := s.generateBasicTask()
|
|
err := task.SetNodeID(100)
|
|
s.NoError(err)
|
|
s.Equal(int64(100), task.GetTaskProto().GetNodeID())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestSaveSegmentMeta() {
|
|
s.Run("success", func() {
|
|
segmentID := int64(101)
|
|
currentManifest := packed.MarshalManifestPath("/data/segments/101", 1)
|
|
resultManifest := packed.MarshalManifestPath("/data/segments/101", 2)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
NumOfRows: 1000,
|
|
StorageVersion: storage.StorageV3,
|
|
ManifestPath: currentManifest,
|
|
Binlogs: []*datapb.FieldBinlog{
|
|
{FieldID: 101, Binlogs: []*datapb.Binlog{{LogID: 1000, EntriesNum: 1000}}},
|
|
},
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
result := &datapb.CompactionPlanResult{
|
|
PlanID: 1,
|
|
State: datapb.CompactionTaskState_completed,
|
|
Type: datapb.CompactionType_BumpSchemaVersionCompaction,
|
|
Segments: []*datapb.CompactionSegment{
|
|
{
|
|
SegmentID: segmentID,
|
|
StorageVersion: storage.StorageV3,
|
|
Manifest: resultManifest,
|
|
InsertLogs: []*datapb.FieldBinlog{
|
|
{FieldID: 101, Binlogs: []*datapb.Binlog{{LogID: 1000, EntriesNum: 1000}}},
|
|
{FieldID: 102, Binlogs: []*datapb.Binlog{{LogID: 2000, EntriesNum: 1000}}},
|
|
},
|
|
},
|
|
},
|
|
}
|
|
err = task.saveSegmentMeta(result)
|
|
s.NoError(err)
|
|
s.Equal(datapb.CompactionTaskState_meta_saved, task.GetTaskProto().GetState())
|
|
})
|
|
|
|
s.Run("CompleteCompactionMutation error", func() {
|
|
// Without adding a segment to meta, CompleteCompactionMutation should fail
|
|
task := s.generateBasicTask()
|
|
result := &datapb.CompactionPlanResult{
|
|
PlanID: 1,
|
|
State: datapb.CompactionTaskState_completed,
|
|
Type: datapb.CompactionType_BumpSchemaVersionCompaction,
|
|
Segments: []*datapb.CompactionSegment{
|
|
{SegmentID: 101},
|
|
},
|
|
}
|
|
err := task.saveSegmentMeta(result)
|
|
s.Error(err)
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestProcessCompleted() {
|
|
segmentID := int64(101)
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segmentID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
},
|
|
isCompacting: true,
|
|
})
|
|
s.NoError(err)
|
|
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_completed)))
|
|
|
|
result := task.processCompleted()
|
|
s.True(result)
|
|
// processCompleted() does NOT reset the compacting flag — that is done by Clean().
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestUpdateAndSaveTaskMeta() {
|
|
s.Run("normal state update", func() {
|
|
task := s.generateBasicTask()
|
|
err := task.updateAndSaveTaskMeta(setState(datapb.CompactionTaskState_executing))
|
|
s.NoError(err)
|
|
s.Equal(datapb.CompactionTaskState_executing, task.GetTaskProto().GetState())
|
|
// EndTime should not be set for non-terminal states
|
|
s.Equal(int64(0), task.GetTaskProto().GetEndTime())
|
|
})
|
|
|
|
s.Run("terminal state sets end time", func() {
|
|
task := s.generateBasicTask()
|
|
// First set to completed state
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_completed)))
|
|
// Then call updateAndSaveTaskMeta which checks current state
|
|
err := task.updateAndSaveTaskMeta()
|
|
s.NoError(err)
|
|
s.Equal(datapb.CompactionTaskState_completed, task.GetTaskProto().GetState())
|
|
s.NotZero(task.GetTaskProto().GetEndTime())
|
|
})
|
|
|
|
s.Run("failed state sets end time", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_failed)))
|
|
err := task.updateAndSaveTaskMeta()
|
|
s.NoError(err)
|
|
s.Equal(datapb.CompactionTaskState_failed, task.GetTaskProto().GetState())
|
|
s.NotZero(task.GetTaskProto().GetEndTime())
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestProcessFailed() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_failed)))
|
|
s.True(task.processFailed())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestGetSlotUsage() {
|
|
task := s.generateBasicTask()
|
|
s.Equal(int64(1), task.GetSlotUsage())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestSetTaskTime() {
|
|
task := s.generateBasicTask()
|
|
now := time.Now()
|
|
task.SetTaskTime(taskcommon.TimeQueue, now)
|
|
s.False(task.GetTaskTime(taskcommon.TimeQueue).IsZero())
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestGetTaskState() {
|
|
s.Run("pipelining state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_pipelining)))
|
|
state := task.GetTaskState()
|
|
s.Equal(taskcommon.Init, state)
|
|
})
|
|
|
|
s.Run("executing state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_executing)))
|
|
state := task.GetTaskState()
|
|
s.Equal(taskcommon.InProgress, state)
|
|
})
|
|
|
|
s.Run("completed state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_completed)))
|
|
state := task.GetTaskState()
|
|
s.Equal(taskcommon.Finished, state)
|
|
})
|
|
|
|
s.Run("failed state", func() {
|
|
task := s.generateBasicTask()
|
|
task.SetTask(task.ShadowClone(setState(datapb.CompactionTaskState_failed)))
|
|
state := task.GetTaskState()
|
|
s.Equal(taskcommon.Failed, state)
|
|
})
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestGetTaskSlot() {
|
|
task := s.generateBasicTask()
|
|
slot := task.GetTaskSlot()
|
|
// GetTaskSlot reads from paramtable; default is 1
|
|
s.GreaterOrEqual(slot, int64(1))
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestCleanError() {
|
|
// Make the compactionTaskMeta catalog fail on SaveCompactionTask so that doClean returns an error.
|
|
// meta.compactionTaskMeta.catalog is the catalog used by SaveCompactionTask,
|
|
// separate from meta.catalog which is used by segment operations.
|
|
mockCatalog := mocks.NewDataCoordCatalog(s.T())
|
|
mockCatalog.EXPECT().SaveCompactionTask(mock.Anything, mock.Anything).Return(errors.New("catalog write error"))
|
|
s.meta.compactionTaskMeta.catalog = mockCatalog
|
|
|
|
task := s.generateBasicTask()
|
|
result := task.Clean()
|
|
s.False(result, "Clean() must return false when doClean fails")
|
|
}
|
|
|
|
func (s *BumpSchemaVersionCompactionTaskSuite) TestResetSegmentCompacting() {
|
|
// Add two segments and mark them as compacting.
|
|
for _, segID := range []int64{101, 102} {
|
|
err := s.meta.AddSegment(context.TODO(), &SegmentInfo{
|
|
SegmentInfo: &datapb.SegmentInfo{
|
|
ID: segID,
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
InsertChannel: "ch-1",
|
|
Level: datapb.SegmentLevel_L1,
|
|
State: commonpb.SegmentState_Flushed,
|
|
},
|
|
})
|
|
s.NoError(err)
|
|
}
|
|
s.meta.SetSegmentsCompacting(context.TODO(), []int64{101, 102}, true)
|
|
|
|
task := s.generateBasicTask()
|
|
// Override input segments to include both IDs.
|
|
task.SetTask(task.ShadowClone(func(t *datapb.CompactionTask) {
|
|
t.InputSegments = []int64{101, 102}
|
|
}))
|
|
|
|
task.resetSegmentCompacting()
|
|
|
|
for _, segID := range []int64{101, 102} {
|
|
seg := s.meta.GetSegment(context.TODO(), segID)
|
|
s.Require().NotNil(seg)
|
|
s.False(seg.isCompacting, "segment %d should no longer be compacting after reset", segID)
|
|
}
|
|
}
|