2021-11-10 23:56:35 +08:00
|
|
|
// 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
|
2021-04-19 10:09:43 +08:00
|
|
|
// with the License. You may obtain a copy of the License at
|
|
|
|
//
|
2021-11-10 23:56:35 +08:00
|
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
2021-04-19 10:09:43 +08:00
|
|
|
//
|
2021-11-10 23:56:35 +08:00
|
|
|
// 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.
|
2021-04-19 10:09:43 +08:00
|
|
|
|
2021-06-22 14:40:07 +08:00
|
|
|
package proxy
|
2021-01-31 14:55:36 +08:00
|
|
|
|
|
|
|
import (
|
2021-02-26 17:44:24 +08:00
|
|
|
"context"
|
2021-09-12 17:54:01 +08:00
|
|
|
"errors"
|
|
|
|
"fmt"
|
2021-01-31 14:55:36 +08:00
|
|
|
"testing"
|
|
|
|
|
2022-08-04 11:04:34 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/util/funcutil"
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
"github.com/stretchr/testify/assert"
|
|
|
|
"github.com/stretchr/testify/require"
|
|
|
|
|
2022-10-16 20:49:27 +08:00
|
|
|
"github.com/milvus-io/milvus-proto/go-api/commonpb"
|
|
|
|
"github.com/milvus-io/milvus-proto/go-api/milvuspb"
|
|
|
|
"github.com/milvus-io/milvus-proto/go-api/schemapb"
|
2021-09-12 17:54:01 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/log"
|
2022-08-04 11:04:34 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/proto/internalpb"
|
2022-05-19 10:13:56 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/proto/querypb"
|
2022-04-20 16:15:41 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/proto/rootcoordpb"
|
2021-09-12 17:54:01 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/types"
|
2022-04-20 16:15:41 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/util/crypto"
|
2021-04-22 14:45:57 +08:00
|
|
|
"github.com/milvus-io/milvus/internal/util/typeutil"
|
2021-01-31 14:55:36 +08:00
|
|
|
)
|
|
|
|
|
2021-06-21 17:28:03 +08:00
|
|
|
type MockRootCoordClientInterface struct {
|
|
|
|
types.RootCoord
|
2021-09-12 17:54:01 +08:00
|
|
|
Error bool
|
|
|
|
AccessCount int
|
2022-08-04 11:04:34 +08:00
|
|
|
|
|
|
|
listPolicy func(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error)
|
2021-01-31 14:55:36 +08:00
|
|
|
}
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
type MockQueryCoordClientInterface struct {
|
|
|
|
types.QueryCoord
|
|
|
|
Error bool
|
|
|
|
AccessCount int
|
|
|
|
}
|
|
|
|
|
2021-06-21 17:28:03 +08:00
|
|
|
func (m *MockRootCoordClientInterface) ShowPartitions(ctx context.Context, in *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
|
2021-09-12 17:54:01 +08:00
|
|
|
if m.Error {
|
|
|
|
return nil, errors.New("mocked error")
|
|
|
|
}
|
2021-01-31 14:55:36 +08:00
|
|
|
if in.CollectionName == "collection1" {
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
2021-01-31 14:55:36 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-31 14:55:36 +08:00
|
|
|
},
|
2021-09-12 17:54:01 +08:00
|
|
|
PartitionIDs: []typeutil.UniqueID{1, 2},
|
|
|
|
CreatedTimestamps: []uint64{100, 200},
|
|
|
|
CreatedUtcTimestamps: []uint64{100, 200},
|
|
|
|
PartitionNames: []string{"par1", "par2"},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
if in.CollectionName == "collection2" {
|
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
PartitionIDs: []typeutil.UniqueID{3, 4},
|
|
|
|
CreatedTimestamps: []uint64{201, 202},
|
|
|
|
CreatedUtcTimestamps: []uint64{201, 202},
|
|
|
|
PartitionNames: []string{"par1", "par2"},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
if in.CollectionName == "errorCollection" {
|
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
PartitionIDs: []typeutil.UniqueID{5, 6},
|
|
|
|
CreatedTimestamps: []uint64{201},
|
|
|
|
CreatedUtcTimestamps: []uint64{201},
|
|
|
|
PartitionNames: []string{"par1", "par2"},
|
2021-01-31 14:55:36 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
2021-01-31 14:55:36 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-09-12 17:54:01 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-31 14:55:36 +08:00
|
|
|
},
|
2021-09-12 17:54:01 +08:00
|
|
|
PartitionIDs: []typeutil.UniqueID{},
|
|
|
|
CreatedTimestamps: []uint64{},
|
|
|
|
CreatedUtcTimestamps: []uint64{},
|
|
|
|
PartitionNames: []string{},
|
2021-01-31 14:55:36 +08:00
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-06-21 17:28:03 +08:00
|
|
|
func (m *MockRootCoordClientInterface) DescribeCollection(ctx context.Context, in *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
|
2021-09-12 17:54:01 +08:00
|
|
|
if m.Error {
|
|
|
|
return nil, errors.New("mocked error")
|
|
|
|
}
|
|
|
|
m.AccessCount++
|
2021-01-31 14:55:36 +08:00
|
|
|
if in.CollectionName == "collection1" {
|
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-31 14:55:36 +08:00
|
|
|
},
|
|
|
|
CollectionID: typeutil.UniqueID(1),
|
|
|
|
Schema: &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
2021-09-12 17:54:01 +08:00
|
|
|
if in.CollectionName == "collection2" {
|
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
CollectionID: typeutil.UniqueID(2),
|
|
|
|
Schema: &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
if in.CollectionName == "errorCollection" {
|
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
CollectionID: typeutil.UniqueID(3),
|
|
|
|
Schema: &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
err := fmt.Errorf("can't find collection: " + in.CollectionName)
|
2021-01-31 14:55:36 +08:00
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-09-12 17:54:01 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_CollectionNotExists,
|
|
|
|
Reason: "describe collection failed: " + err.Error(),
|
2021-01-31 14:55:36 +08:00
|
|
|
},
|
2021-09-12 17:54:01 +08:00
|
|
|
Schema: nil,
|
2021-01-31 14:55:36 +08:00
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2022-04-11 19:49:34 +08:00
|
|
|
func (m *MockRootCoordClientInterface) GetCredential(ctx context.Context, req *rootcoordpb.GetCredentialRequest) (*rootcoordpb.GetCredentialResponse, error) {
|
|
|
|
if m.Error {
|
|
|
|
return nil, errors.New("mocked error")
|
|
|
|
}
|
|
|
|
m.AccessCount++
|
|
|
|
if req.Username == "mockUser" {
|
|
|
|
encryptedPassword, _ := crypto.PasswordEncrypt("mockPass")
|
|
|
|
return &rootcoordpb.GetCredentialResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
Username: "mockUser",
|
|
|
|
Password: encryptedPassword,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
err := fmt.Errorf("can't find credential: " + req.Username)
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *MockRootCoordClientInterface) ListCredUsers(ctx context.Context, req *milvuspb.ListCredUsersRequest) (*milvuspb.ListCredUsersResponse, error) {
|
|
|
|
if m.Error {
|
|
|
|
return nil, errors.New("mocked error")
|
|
|
|
}
|
|
|
|
|
|
|
|
return &milvuspb.ListCredUsersResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
Usernames: []string{"mockUser"},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2022-08-04 11:04:34 +08:00
|
|
|
func (m *MockRootCoordClientInterface) ListPolicy(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error) {
|
|
|
|
if m.listPolicy != nil {
|
|
|
|
return m.listPolicy(ctx, in)
|
|
|
|
}
|
|
|
|
return &internalpb.ListPolicyResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
func (m *MockQueryCoordClientInterface) ShowCollections(ctx context.Context, req *querypb.ShowCollectionsRequest) (*querypb.ShowCollectionsResponse, error) {
|
|
|
|
if m.Error {
|
|
|
|
return nil, errors.New("mocked error")
|
|
|
|
}
|
|
|
|
m.AccessCount++
|
|
|
|
rsp := &querypb.ShowCollectionsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
CollectionIDs: []UniqueID{1, 2},
|
|
|
|
InMemoryPercentages: []int64{100, 50},
|
|
|
|
}
|
|
|
|
return rsp, nil
|
|
|
|
}
|
|
|
|
|
2021-09-12 17:54:01 +08:00
|
|
|
//Simulate the cache path and the
|
2021-01-31 14:55:36 +08:00
|
|
|
func TestMetaCache_GetCollection(t *testing.T) {
|
2021-02-26 17:44:24 +08:00
|
|
|
ctx := context.Background()
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
id, err := globalMetaCache.GetCollectionID(ctx, "collection1")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(1))
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
2021-09-12 17:54:01 +08:00
|
|
|
|
|
|
|
// should'nt be accessed to remote root coord.
|
2021-02-26 17:44:24 +08:00
|
|
|
schema, err := globalMetaCache.GetCollectionSchema(ctx, "collection1")
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, schema, &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
2021-09-12 17:54:01 +08:00
|
|
|
Fields: []*schemapb.FieldSchema{},
|
2021-01-31 14:55:36 +08:00
|
|
|
})
|
2021-02-26 17:44:24 +08:00
|
|
|
id, err = globalMetaCache.GetCollectionID(ctx, "collection2")
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(2))
|
2021-02-26 17:44:24 +08:00
|
|
|
schema, err = globalMetaCache.GetCollectionSchema(ctx, "collection2")
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, schema, &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
Fields: []*schemapb.FieldSchema{},
|
|
|
|
})
|
|
|
|
|
|
|
|
// test to get from cache, this should trigger root request
|
|
|
|
id, err = globalMetaCache.GetCollectionID(ctx, "collection1")
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(1))
|
|
|
|
schema, err = globalMetaCache.GetCollectionSchema(ctx, "collection1")
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, schema, &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
Fields: []*schemapb.FieldSchema{},
|
|
|
|
})
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_GetCollectionFailure(t *testing.T) {
|
|
|
|
ctx := context.Background()
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord.Error = true
|
2021-09-12 17:54:01 +08:00
|
|
|
|
|
|
|
schema, err := globalMetaCache.GetCollectionSchema(ctx, "collection1")
|
|
|
|
assert.NotNil(t, err)
|
|
|
|
assert.Nil(t, schema)
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord.Error = false
|
2021-09-12 17:54:01 +08:00
|
|
|
|
|
|
|
schema, err = globalMetaCache.GetCollectionSchema(ctx, "collection1")
|
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, schema, &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
Fields: []*schemapb.FieldSchema{},
|
|
|
|
})
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord.Error = true
|
2021-09-12 17:54:01 +08:00
|
|
|
// should be cached with no error
|
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, schema, &schemapb.CollectionSchema{
|
|
|
|
AutoID: true,
|
|
|
|
Fields: []*schemapb.FieldSchema{},
|
|
|
|
})
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_GetNonExistCollection(t *testing.T) {
|
|
|
|
ctx := context.Background()
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
|
|
|
id, err := globalMetaCache.GetCollectionID(ctx, "collection3")
|
|
|
|
assert.NotNil(t, err)
|
|
|
|
assert.Equal(t, id, int64(0))
|
|
|
|
schema, err := globalMetaCache.GetCollectionSchema(ctx, "collection3")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.NotNil(t, err)
|
|
|
|
assert.Nil(t, schema)
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_GetPartitionID(t *testing.T) {
|
2021-02-26 17:44:24 +08:00
|
|
|
ctx := context.Background()
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
id, err := globalMetaCache.GetPartitionID(ctx, "collection1", "par1")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(1))
|
2021-02-26 17:44:24 +08:00
|
|
|
id, err = globalMetaCache.GetPartitionID(ctx, "collection1", "par2")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(2))
|
2021-09-12 17:54:01 +08:00
|
|
|
id, err = globalMetaCache.GetPartitionID(ctx, "collection2", "par1")
|
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(3))
|
|
|
|
id, err = globalMetaCache.GetPartitionID(ctx, "collection2", "par2")
|
|
|
|
assert.Nil(t, err)
|
|
|
|
assert.Equal(t, id, typeutil.UniqueID(4))
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_GetPartitionError(t *testing.T) {
|
|
|
|
ctx := context.Background()
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2021-09-12 17:54:01 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
|
|
|
// Test the case where ShowPartitionsResponse is not aligned
|
|
|
|
id, err := globalMetaCache.GetPartitionID(ctx, "errorCollection", "par1")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.NotNil(t, err)
|
2021-09-12 17:54:01 +08:00
|
|
|
log.Debug(err.Error())
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Equal(t, id, typeutil.UniqueID(0))
|
2021-09-12 17:54:01 +08:00
|
|
|
|
|
|
|
partitions, err2 := globalMetaCache.GetPartitions(ctx, "errorCollection")
|
|
|
|
assert.NotNil(t, err2)
|
|
|
|
log.Debug(err.Error())
|
|
|
|
assert.Equal(t, len(partitions), 0)
|
|
|
|
|
|
|
|
// Test non existed tables
|
|
|
|
id, err = globalMetaCache.GetPartitionID(ctx, "nonExisted", "par1")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.NotNil(t, err)
|
2021-09-12 17:54:01 +08:00
|
|
|
log.Debug(err.Error())
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Equal(t, id, typeutil.UniqueID(0))
|
2021-09-12 17:54:01 +08:00
|
|
|
|
|
|
|
// Test non existed partition
|
|
|
|
id, err = globalMetaCache.GetPartitionID(ctx, "collection1", "par3")
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.NotNil(t, err)
|
2021-09-12 17:54:01 +08:00
|
|
|
log.Debug(err.Error())
|
2021-01-31 14:55:36 +08:00
|
|
|
assert.Equal(t, id, typeutil.UniqueID(0))
|
|
|
|
}
|
2022-04-20 16:15:41 +08:00
|
|
|
|
|
|
|
func TestMetaCache_GetShards(t *testing.T) {
|
2022-08-04 11:04:34 +08:00
|
|
|
var (
|
|
|
|
ctx = context.Background()
|
|
|
|
collectionName = "collection1"
|
|
|
|
)
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
qc := NewQueryCoordMock()
|
|
|
|
shardMgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, qc, shardMgr)
|
2022-04-20 16:15:41 +08:00
|
|
|
require.Nil(t, err)
|
|
|
|
|
|
|
|
qc.Init()
|
|
|
|
qc.Start()
|
|
|
|
defer qc.Stop()
|
|
|
|
|
|
|
|
t.Run("No collection in meta cache", func(t *testing.T) {
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err := globalMetaCache.GetShards(ctx, true, "non-exists")
|
2022-04-20 16:15:41 +08:00
|
|
|
assert.Error(t, err)
|
|
|
|
assert.Empty(t, shards)
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("without shardLeaders in collection info invalid shardLeaders", func(t *testing.T) {
|
|
|
|
qc.validShardLeaders = false
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err := globalMetaCache.GetShards(ctx, false, collectionName)
|
2022-04-20 16:15:41 +08:00
|
|
|
assert.Error(t, err)
|
|
|
|
assert.Empty(t, shards)
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("without shardLeaders in collection info", func(t *testing.T) {
|
|
|
|
qc.validShardLeaders = true
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err := globalMetaCache.GetShards(ctx, true, collectionName)
|
2022-04-20 16:15:41 +08:00
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.NotEmpty(t, shards)
|
|
|
|
assert.Equal(t, 1, len(shards))
|
2022-05-17 22:35:57 +08:00
|
|
|
assert.Equal(t, 3, len(shards["channel-1"]))
|
2022-04-20 16:15:41 +08:00
|
|
|
|
|
|
|
// get from cache
|
|
|
|
qc.validShardLeaders = false
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err = globalMetaCache.GetShards(ctx, true, collectionName)
|
2022-08-23 10:44:52 +08:00
|
|
|
|
2022-04-20 16:15:41 +08:00
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.NotEmpty(t, shards)
|
|
|
|
assert.Equal(t, 1, len(shards))
|
2022-05-17 22:35:57 +08:00
|
|
|
assert.Equal(t, 3, len(shards["channel-1"]))
|
2022-04-20 16:15:41 +08:00
|
|
|
})
|
2022-05-17 11:11:56 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_ClearShards(t *testing.T) {
|
2022-08-04 11:04:34 +08:00
|
|
|
var (
|
|
|
|
ctx = context.TODO()
|
|
|
|
collectionName = "collection1"
|
|
|
|
)
|
|
|
|
|
2022-05-19 10:13:56 +08:00
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
qc := NewQueryCoordMock()
|
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, qc, mgr)
|
2022-05-17 11:11:56 +08:00
|
|
|
require.Nil(t, err)
|
|
|
|
|
|
|
|
qc.Init()
|
|
|
|
qc.Start()
|
|
|
|
defer qc.Stop()
|
|
|
|
|
|
|
|
t.Run("Clear with no collection info", func(t *testing.T) {
|
|
|
|
globalMetaCache.ClearShards("collection_not_exist")
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("Clear valid collection empty cache", func(t *testing.T) {
|
|
|
|
globalMetaCache.ClearShards(collectionName)
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("Clear valid collection valid cache", func(t *testing.T) {
|
|
|
|
|
|
|
|
qc.validShardLeaders = true
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err := globalMetaCache.GetShards(ctx, true, collectionName)
|
2022-05-17 11:11:56 +08:00
|
|
|
require.NoError(t, err)
|
|
|
|
require.NotEmpty(t, shards)
|
|
|
|
require.Equal(t, 1, len(shards))
|
2022-05-17 22:35:57 +08:00
|
|
|
require.Equal(t, 3, len(shards["channel-1"]))
|
2022-05-17 11:11:56 +08:00
|
|
|
|
|
|
|
globalMetaCache.ClearShards(collectionName)
|
|
|
|
|
|
|
|
qc.validShardLeaders = false
|
2022-06-02 12:16:03 +08:00
|
|
|
shards, err = globalMetaCache.GetShards(ctx, true, collectionName)
|
2022-05-17 11:11:56 +08:00
|
|
|
assert.Error(t, err)
|
|
|
|
assert.Empty(t, shards)
|
|
|
|
})
|
2022-08-04 11:04:34 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_PolicyInfo(t *testing.T) {
|
|
|
|
client := &MockRootCoordClientInterface{}
|
|
|
|
qc := &MockQueryCoordClientInterface{}
|
|
|
|
mgr := newShardClientMgr()
|
|
|
|
|
|
|
|
t.Run("InitMetaCache", func(t *testing.T) {
|
|
|
|
client.listPolicy = func(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error) {
|
|
|
|
return nil, fmt.Errorf("mock error")
|
|
|
|
}
|
|
|
|
err := InitMetaCache(context.Background(), client, qc, mgr)
|
|
|
|
assert.NotNil(t, err)
|
|
|
|
|
|
|
|
client.listPolicy = func(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error) {
|
|
|
|
return &internalpb.ListPolicyResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
PolicyInfos: []string{"policy1", "policy2", "policy3"},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
err = InitMetaCache(context.Background(), client, qc, mgr)
|
|
|
|
assert.Nil(t, err)
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("GetPrivilegeInfo", func(t *testing.T) {
|
|
|
|
client.listPolicy = func(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error) {
|
|
|
|
return &internalpb.ListPolicyResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
PolicyInfos: []string{"policy1", "policy2", "policy3"},
|
|
|
|
UserRoles: []string{funcutil.EncodeUserRoleCache("foo", "role1"), funcutil.EncodeUserRoleCache("foo", "role2"), funcutil.EncodeUserRoleCache("foo2", "role2")},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
err := InitMetaCache(context.Background(), client, qc, mgr)
|
|
|
|
assert.Nil(t, err)
|
|
|
|
policyInfos := globalMetaCache.GetPrivilegeInfo(context.Background())
|
|
|
|
assert.Equal(t, 3, len(policyInfos))
|
|
|
|
roles := globalMetaCache.GetUserRole("foo")
|
|
|
|
assert.Equal(t, 2, len(roles))
|
|
|
|
})
|
2022-04-20 16:15:41 +08:00
|
|
|
|
2022-08-04 11:04:34 +08:00
|
|
|
t.Run("GetPrivilegeInfo", func(t *testing.T) {
|
|
|
|
client.listPolicy = func(ctx context.Context, in *internalpb.ListPolicyRequest) (*internalpb.ListPolicyResponse, error) {
|
|
|
|
return &internalpb.ListPolicyResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
},
|
|
|
|
PolicyInfos: []string{"policy1", "policy2", "policy3"},
|
|
|
|
UserRoles: []string{funcutil.EncodeUserRoleCache("foo", "role1"), funcutil.EncodeUserRoleCache("foo", "role2"), funcutil.EncodeUserRoleCache("foo2", "role2")},
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
err := InitMetaCache(context.Background(), client, qc, mgr)
|
|
|
|
assert.Nil(t, err)
|
|
|
|
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: typeutil.CacheGrantPrivilege, OpKey: "policyX"})
|
|
|
|
assert.Nil(t, err)
|
|
|
|
policyInfos := globalMetaCache.GetPrivilegeInfo(context.Background())
|
|
|
|
assert.Equal(t, 4, len(policyInfos))
|
|
|
|
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: typeutil.CacheRevokePrivilege, OpKey: "policyX"})
|
|
|
|
assert.Nil(t, err)
|
|
|
|
policyInfos = globalMetaCache.GetPrivilegeInfo(context.Background())
|
|
|
|
assert.Equal(t, 3, len(policyInfos))
|
|
|
|
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: typeutil.CacheAddUserToRole, OpKey: funcutil.EncodeUserRoleCache("foo", "role3")})
|
|
|
|
assert.Nil(t, err)
|
|
|
|
roles := globalMetaCache.GetUserRole("foo")
|
|
|
|
assert.Equal(t, 3, len(roles))
|
|
|
|
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: typeutil.CacheRemoveUserFromRole, OpKey: funcutil.EncodeUserRoleCache("foo", "role3")})
|
|
|
|
assert.Nil(t, err)
|
|
|
|
roles = globalMetaCache.GetUserRole("foo")
|
|
|
|
assert.Equal(t, 2, len(roles))
|
|
|
|
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: typeutil.CacheGrantPrivilege, OpKey: ""})
|
|
|
|
assert.NotNil(t, err)
|
|
|
|
err = globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{OpType: 100, OpKey: "policyX"})
|
|
|
|
assert.NotNil(t, err)
|
|
|
|
})
|
2022-04-20 16:15:41 +08:00
|
|
|
}
|
2022-05-19 10:13:56 +08:00
|
|
|
|
|
|
|
func TestMetaCache_LoadCache(t *testing.T) {
|
|
|
|
ctx := context.Background()
|
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
mgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, mgr)
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
|
|
|
t.Run("test IsCollectionLoaded", func(t *testing.T) {
|
|
|
|
info, err := globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// no collectionInfo of collection1, should access RootCoord
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
|
|
|
// not loaded, should access QueryCoord
|
|
|
|
assert.Equal(t, queryCoord.AccessCount, 1)
|
|
|
|
|
|
|
|
info, err = globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// shouldn't access QueryCoord or RootCoord again
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
|
|
|
assert.Equal(t, queryCoord.AccessCount, 1)
|
|
|
|
|
|
|
|
// test collection2 not fully loaded
|
|
|
|
info, err = globalMetaCache.GetCollectionInfo(ctx, "collection2")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.False(t, info.isLoaded)
|
|
|
|
// no collectionInfo of collection2, should access RootCoord
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
|
|
|
// not loaded, should access QueryCoord
|
|
|
|
assert.Equal(t, queryCoord.AccessCount, 2)
|
|
|
|
})
|
|
|
|
|
|
|
|
t.Run("test RemoveCollectionLoadCache", func(t *testing.T) {
|
|
|
|
globalMetaCache.RemoveCollection(ctx, "collection1")
|
|
|
|
info, err := globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// should access QueryCoord
|
|
|
|
assert.Equal(t, queryCoord.AccessCount, 3)
|
|
|
|
})
|
|
|
|
}
|
|
|
|
|
|
|
|
func TestMetaCache_RemoveCollection(t *testing.T) {
|
|
|
|
ctx := context.Background()
|
|
|
|
rootCoord := &MockRootCoordClientInterface{}
|
|
|
|
queryCoord := &MockQueryCoordClientInterface{}
|
2022-06-02 12:16:03 +08:00
|
|
|
shardMgr := newShardClientMgr()
|
2022-08-04 11:04:34 +08:00
|
|
|
err := InitMetaCache(ctx, rootCoord, queryCoord, shardMgr)
|
2022-05-19 10:13:56 +08:00
|
|
|
assert.Nil(t, err)
|
|
|
|
|
|
|
|
info, err := globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// no collectionInfo of collection1, should access RootCoord
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
|
|
|
|
|
|
|
info, err = globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// shouldn't access RootCoord again
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 1)
|
|
|
|
|
|
|
|
globalMetaCache.RemoveCollection(ctx, "collection1")
|
|
|
|
// no collectionInfo of collection2, should access RootCoord
|
|
|
|
info, err = globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// shouldn't access RootCoord again
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 2)
|
|
|
|
|
|
|
|
globalMetaCache.RemoveCollectionsByID(ctx, UniqueID(1))
|
|
|
|
// no collectionInfo of collection2, should access RootCoord
|
|
|
|
info, err = globalMetaCache.GetCollectionInfo(ctx, "collection1")
|
|
|
|
assert.NoError(t, err)
|
|
|
|
assert.True(t, info.isLoaded)
|
|
|
|
// shouldn't access RootCoord again
|
|
|
|
assert.Equal(t, rootCoord.AccessCount, 3)
|
|
|
|
}
|