Files
milvus-io--milvus/internal/datacoord/compaction_task_bump_schema_version_test.go
wehub-resource-sync 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
chore: import upstream snapshot with attribution
2026-07-13 12:31:17 +08:00

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)
}
}