impl.go 107.1 KB
Newer Older
1 2 3 4 5 6
// 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
7 8
// with the License. You may obtain a copy of the License at
//
9
//     http://www.apache.org/licenses/LICENSE-2.0
10
//
11 12 13 14 15
// 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.
16

C
Cai Yudong 已提交
17
package proxy
18 19 20

import (
	"context"
21
	"errors"
22
	"fmt"
C
cai.zhang 已提交
23
	"os"
24 25
	"strconv"

26 27
	"github.com/milvus-io/milvus/internal/common"

28 29
	"github.com/milvus-io/milvus/internal/util/funcutil"

C
cai.zhang 已提交
30 31
	"github.com/milvus-io/milvus/internal/util/trace"

32 33
	"github.com/milvus-io/milvus/internal/util/metricsinfo"

34
	"go.uber.org/zap"
S
sunby 已提交
35

X
Xiangyu Wang 已提交
36 37 38 39 40 41 42 43
	"github.com/milvus-io/milvus/internal/log"
	"github.com/milvus-io/milvus/internal/msgstream"
	"github.com/milvus-io/milvus/internal/proto/commonpb"
	"github.com/milvus-io/milvus/internal/proto/datapb"
	"github.com/milvus-io/milvus/internal/proto/internalpb"
	"github.com/milvus-io/milvus/internal/proto/milvuspb"
	"github.com/milvus-io/milvus/internal/proto/proxypb"
	"github.com/milvus-io/milvus/internal/proto/querypb"
44
	"github.com/milvus-io/milvus/internal/proto/schemapb"
45
	"github.com/milvus-io/milvus/internal/util/distance"
X
Xiangyu Wang 已提交
46
	"github.com/milvus-io/milvus/internal/util/typeutil"
47 48
)

49
// UpdateStateCode updates the state code of Proxy.
C
Cai Yudong 已提交
50
func (node *Proxy) UpdateStateCode(code internalpb.StateCode) {
51
	node.stateCode.Store(code)
Z
zhenshan.cao 已提交
52 53
}

54
// GetComponentStates get state of Proxy.
C
Cai Yudong 已提交
55
func (node *Proxy) GetComponentStates(ctx context.Context) (*internalpb.ComponentStates, error) {
G
godchen 已提交
56 57 58 59 60 61 62 63 64 65 66 67
	stats := &internalpb.ComponentStates{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
	}
	code, ok := node.stateCode.Load().(internalpb.StateCode)
	if !ok {
		errMsg := "unexpected error in type assertion"
		stats.Status = &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    errMsg,
		}
G
godchen 已提交
68
		return stats, nil
G
godchen 已提交
69
	}
70 71 72 73
	nodeID := common.NotRegisteredID
	if node.session != nil && node.session.Registered() {
		nodeID = node.session.ServerID
	}
G
godchen 已提交
74
	info := &internalpb.ComponentInfo{
75 76
		// NodeID:    Params.ProxyID, // will race with Proxy.Register()
		NodeID:    nodeID,
C
Cai Yudong 已提交
77
		Role:      typeutil.ProxyRole,
G
godchen 已提交
78 79 80 81 82 83
		StateCode: code,
	}
	stats.State = info
	return stats, nil
}

84
// GetStatisticsChannel get statistics channel of Proxy.
C
Cai Yudong 已提交
85
func (node *Proxy) GetStatisticsChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
G
godchen 已提交
86 87 88 89 90 91 92 93 94
	return &milvuspb.StringResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
		Value: "",
	}, nil
}

95
// InvalidateCollectionMetaCache invalidate the meta cache of specific collection.
C
Cai Yudong 已提交
96
func (node *Proxy) InvalidateCollectionMetaCache(ctx context.Context, request *proxypb.InvalidateCollMetaCacheRequest) (*commonpb.Status, error) {
D
dragondriver 已提交
97
	log.Debug("InvalidateCollectionMetaCache",
98
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
99 100 101
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

102
	collectionName := request.CollectionName
N
neza2017 已提交
103 104 105
	if globalMetaCache != nil {
		globalMetaCache.RemoveCollection(ctx, collectionName) // no need to return error, though collection may be not cached
	}
D
dragondriver 已提交
106
	log.Debug("InvalidateCollectionMetaCache Done",
107
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
108 109 110
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

111
	return &commonpb.Status{
112
		ErrorCode: commonpb.ErrorCode_Success,
113 114
		Reason:    "",
	}, nil
115 116
}

117
// ReleaseDQLMessageStream release the query message stream of specific collection.
C
Cai Yudong 已提交
118
func (node *Proxy) ReleaseDQLMessageStream(ctx context.Context, request *proxypb.ReleaseDQLMessageStreamRequest) (*commonpb.Status, error) {
119
	log.Debug("ReleaseDQLMessageStream",
120
		zap.Any("role", typeutil.ProxyRole),
121 122 123
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

124 125 126 127
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

128 129 130
	_ = node.chMgr.removeDQLStream(request.CollectionID)

	log.Debug("ReleaseDQLMessageStream Done",
131
		zap.Any("role", typeutil.ProxyRole),
132 133 134 135 136 137 138 139 140
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_Success,
		Reason:    "",
	}, nil
}

141
// CreateCollection create a collection by the schema.
C
Cai Yudong 已提交
142
func (node *Proxy) CreateCollection(ctx context.Context, request *milvuspb.CreateCollectionRequest) (*commonpb.Status, error) {
143 144 145
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
146 147 148 149 150

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

151
	cct := &createCollectionTask{
S
sunby 已提交
152
		ctx:                     ctx,
153 154
		Condition:               NewTaskCondition(ctx),
		CreateCollectionRequest: request,
155
		rootCoord:               node.rootCoord,
156 157
	}

158 159 160 161 162
	// avoid data race
	lenOfSchema := len(request.Schema)

	log.Debug("CreateCollection received",
		zap.String("traceID", traceID),
163
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
164 165
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
166 167 168
		zap.Int("len(schema)", lenOfSchema),
		zap.Int32("shards_num", request.ShardsNum))

169
	err := node.sched.ddQueue.Enqueue(cct)
170
	if err != nil {
171 172 173
		log.Debug("CreateCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
174
			zap.String("role", typeutil.ProxyRole),
175 176 177 178 179
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Int("len(schema)", lenOfSchema),
			zap.Int32("shards_num", request.ShardsNum))

180
		return &commonpb.Status{
181
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
182 183 184 185
			Reason:    err.Error(),
		}, nil
	}

186 187
	log.Debug("CreateCollection enqueued",
		zap.String("traceID", traceID),
188
		zap.String("role", typeutil.ProxyRole),
189 190 191
		zap.Int64("MsgID", cct.ID()),
		zap.Uint64("BeginTs", cct.BeginTs()),
		zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
192 193
		zap.Uint64("timestamp", request.Base.Timestamp),
		zap.String("db", request.DbName),
194 195 196 197 198 199 200
		zap.String("collection", request.CollectionName),
		zap.Int("len(schema)", lenOfSchema),
		zap.Int32("shards_num", request.ShardsNum))

	err = cct.WaitToFinish()
	if err != nil {
		log.Debug("CreateCollection failed to WaitToFinish",
D
dragondriver 已提交
201
			zap.Error(err),
202
			zap.String("traceID", traceID),
203
			zap.String("role", typeutil.ProxyRole),
204 205 206
			zap.Int64("MsgID", cct.ID()),
			zap.Uint64("BeginTs", cct.BeginTs()),
			zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
207 208
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
209 210
			zap.Int("len(schema)", lenOfSchema),
			zap.Int32("shards_num", request.ShardsNum))
D
dragondriver 已提交
211

212
		return &commonpb.Status{
213
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
214 215 216 217
			Reason:    err.Error(),
		}, nil
	}

218 219
	log.Debug("CreateCollection done",
		zap.String("traceID", traceID),
220
		zap.String("role", typeutil.ProxyRole),
221 222 223 224 225 226 227 228
		zap.Int64("MsgID", cct.ID()),
		zap.Uint64("BeginTs", cct.BeginTs()),
		zap.Uint64("EndTs", cct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Int("len(schema)", lenOfSchema),
		zap.Int32("shards_num", request.ShardsNum))

229 230 231
	return cct.result, nil
}

232
// DropCollection drop a collection.
C
Cai Yudong 已提交
233
func (node *Proxy) DropCollection(ctx context.Context, request *milvuspb.DropCollectionRequest) (*commonpb.Status, error) {
234 235 236
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
237 238 239 240 241

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

242
	dct := &dropCollectionTask{
S
sunby 已提交
243
		ctx:                   ctx,
244 245
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
246
		rootCoord:             node.rootCoord,
247
		chMgr:                 node.chMgr,
S
sunby 已提交
248
		chTicker:              node.chTicker,
249 250
	}

251 252
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
253
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
254 255
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
256 257 258 259 260

	if err := node.sched.ddQueue.Enqueue(dct); err != nil {
		log.Warn("DropCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
261
			zap.String("role", typeutil.ProxyRole),
262 263 264
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

265
		return &commonpb.Status{
266
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
267 268 269 270
			Reason:    err.Error(),
		}, nil
	}

271 272
	log.Debug("DropCollection enqueued",
		zap.String("traceID", traceID),
273
		zap.String("role", typeutil.ProxyRole),
274 275 276
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
277 278
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
279 280 281

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DropCollection failed to WaitToFinish",
D
dragondriver 已提交
282
			zap.Error(err),
283
			zap.String("traceID", traceID),
284
			zap.String("role", typeutil.ProxyRole),
285 286 287
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTs", dct.BeginTs()),
			zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
288 289 290
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

291
		return &commonpb.Status{
292
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
293 294 295 296
			Reason:    err.Error(),
		}, nil
	}

297 298
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
299
		zap.String("role", typeutil.ProxyRole),
300 301 302 303 304 305
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

306 307 308
	return dct.result, nil
}

309
// HasCollection check if the specific collection exists in Milvus.
C
Cai Yudong 已提交
310
func (node *Proxy) HasCollection(ctx context.Context, request *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
311 312 313 314 315
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
316 317 318 319 320 321 322

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
323
		zap.String("role", typeutil.ProxyRole),
324 325 326
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

327
	hct := &hasCollectionTask{
S
sunby 已提交
328
		ctx:                  ctx,
329 330
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
331
		rootCoord:            node.rootCoord,
332 333
	}

334 335 336 337
	if err := node.sched.ddQueue.Enqueue(hct); err != nil {
		log.Warn("HasCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
338
			zap.String("role", typeutil.ProxyRole),
339 340 341
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

342 343
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
344
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
345 346 347 348 349
				Reason:    err.Error(),
			},
		}, nil
	}

350 351
	log.Debug("HasCollection enqueued",
		zap.String("traceID", traceID),
352
		zap.String("role", typeutil.ProxyRole),
353 354 355
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
356 357
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
358 359 360

	if err := hct.WaitToFinish(); err != nil {
		log.Warn("HasCollection failed to WaitToFinish",
D
dragondriver 已提交
361
			zap.Error(err),
362
			zap.String("traceID", traceID),
363
			zap.String("role", typeutil.ProxyRole),
364 365 366
			zap.Int64("MsgID", hct.ID()),
			zap.Uint64("BeginTS", hct.BeginTs()),
			zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
367 368 369
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

370 371
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
372
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
373 374 375 376 377
				Reason:    err.Error(),
			},
		}, nil
	}

378 379
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
380
		zap.String("role", typeutil.ProxyRole),
381 382 383 384 385 386
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

387 388 389
	return hct.result, nil
}

390
// LoadCollection load a collection into query nodes.
C
Cai Yudong 已提交
391
func (node *Proxy) LoadCollection(ctx context.Context, request *milvuspb.LoadCollectionRequest) (*commonpb.Status, error) {
392 393 394
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
395 396 397 398 399

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

400
	lct := &loadCollectionTask{
S
sunby 已提交
401
		ctx:                   ctx,
402 403
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
404
		queryCoord:            node.queryCoord,
405 406
	}

407 408
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
409
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
410 411
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
412 413 414 415 416

	if err := node.sched.ddQueue.Enqueue(lct); err != nil {
		log.Warn("LoadCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
417
			zap.String("role", typeutil.ProxyRole),
418 419 420
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

421
		return &commonpb.Status{
422
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
423 424 425
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
426

427 428
	log.Debug("LoadCollection enqueued",
		zap.String("traceID", traceID),
429
		zap.String("role", typeutil.ProxyRole),
430 431 432
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
433 434
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
435 436 437

	if err := lct.WaitToFinish(); err != nil {
		log.Warn("LoadCollection failed to WaitToFinish",
D
dragondriver 已提交
438
			zap.Error(err),
439
			zap.String("traceID", traceID),
440
			zap.String("role", typeutil.ProxyRole),
441 442 443
			zap.Int64("MsgID", lct.ID()),
			zap.Uint64("BeginTS", lct.BeginTs()),
			zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
444 445 446
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

447
		return &commonpb.Status{
448
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
449 450 451 452
			Reason:    err.Error(),
		}, nil
	}

453 454
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
455
		zap.String("role", typeutil.ProxyRole),
456 457 458 459 460 461
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

462
	return lct.result, nil
463 464
}

465
// ReleaseCollection remove the loaded collection from query nodes.
C
Cai Yudong 已提交
466
func (node *Proxy) ReleaseCollection(ctx context.Context, request *milvuspb.ReleaseCollectionRequest) (*commonpb.Status, error) {
467 468 469
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
470

471
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
472 473 474
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

475
	rct := &releaseCollectionTask{
S
sunby 已提交
476
		ctx:                      ctx,
477 478
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
479
		queryCoord:               node.queryCoord,
480
		chMgr:                    node.chMgr,
481 482
	}

483 484 485 486
	method := "ReleaseCollection"

	log.Debug(
		rpcReceived(method),
487
		zap.String("traceID", traceID),
488
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
489 490
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
491 492

	if err := node.sched.ddQueue.Enqueue(rct); err != nil {
493 494
		log.Warn(
			rpcFailedToEnqueue(method),
495 496
			zap.Error(err),
			zap.String("traceID", traceID),
497
			zap.String("role", typeutil.ProxyRole),
498 499 500
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

501
		return &commonpb.Status{
502
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
503 504 505 506
			Reason:    err.Error(),
		}, nil
	}

507 508
	log.Debug(
		rpcEnqueued(method),
509
		zap.String("traceID", traceID),
510
		zap.String("role", typeutil.ProxyRole),
511 512 513
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
514 515
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
516 517

	if err := rct.WaitToFinish(); err != nil {
518 519
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
520
			zap.Error(err),
521
			zap.String("traceID", traceID),
522
			zap.String("role", typeutil.ProxyRole),
523 524 525
			zap.Int64("MsgID", rct.ID()),
			zap.Uint64("BeginTS", rct.BeginTs()),
			zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
526 527 528
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

529
		return &commonpb.Status{
530
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
531 532 533 534
			Reason:    err.Error(),
		}, nil
	}

535 536
	log.Debug(
		rpcDone(method),
537
		zap.String("traceID", traceID),
538
		zap.String("role", typeutil.ProxyRole),
539 540 541 542 543 544
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

545
	return rct.result, nil
546 547
}

548
// DescribeCollection get the meta information of specific collection, such as schema, created timestamp and etc.
C
Cai Yudong 已提交
549
func (node *Proxy) DescribeCollection(ctx context.Context, request *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
550 551 552 553 554
	if !node.checkHealthy() {
		return &milvuspb.DescribeCollectionResponse{
			Status: unhealthyStatus(),
		}, nil
	}
555

556
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
557 558 559
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

560
	dct := &describeCollectionTask{
S
sunby 已提交
561
		ctx:                       ctx,
562 563
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
564
		rootCoord:                 node.rootCoord,
565 566
	}

567 568
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
569
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
570 571
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
572 573 574 575 576

	if err := node.sched.ddQueue.Enqueue(dct); err != nil {
		log.Warn("DescribeCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
577
			zap.String("role", typeutil.ProxyRole),
578 579 580
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

581 582
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
583
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
584 585 586 587 588
				Reason:    err.Error(),
			},
		}, nil
	}

589 590
	log.Debug("DescribeCollection enqueued",
		zap.String("traceID", traceID),
591
		zap.String("role", typeutil.ProxyRole),
592 593 594
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
595 596
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
597 598 599

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DescribeCollection failed to WaitToFinish",
D
dragondriver 已提交
600
			zap.Error(err),
601
			zap.String("traceID", traceID),
602
			zap.String("role", typeutil.ProxyRole),
603 604 605
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTS", dct.BeginTs()),
			zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
606 607 608
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

609 610
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
611
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
612 613 614 615 616
				Reason:    err.Error(),
			},
		}, nil
	}

617 618
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
619
		zap.String("role", typeutil.ProxyRole),
620 621 622 623 624 625
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

626 627 628
	return dct.result, nil
}

629
// GetCollectionStatistics get the collection statistics, such as `num_rows`.
C
Cai Yudong 已提交
630
func (node *Proxy) GetCollectionStatistics(ctx context.Context, request *milvuspb.GetCollectionStatisticsRequest) (*milvuspb.GetCollectionStatisticsResponse, error) {
631 632 633 634 635
	if !node.checkHealthy() {
		return &milvuspb.GetCollectionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
636 637 638 639 640

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

641
	g := &getCollectionStatisticsTask{
G
godchen 已提交
642 643 644
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
645
		dataCoord:                      node.dataCoord,
646 647
	}

648 649
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
650
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
651 652
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
653 654 655 656 657

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn("GetCollectionStatistics failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
658
			zap.String("role", typeutil.ProxyRole),
659 660 661
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

G
godchen 已提交
662
		return &milvuspb.GetCollectionStatisticsResponse{
663
			Status: &commonpb.Status{
664
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
665 666 667 668 669
				Reason:    err.Error(),
			},
		}, nil
	}

670 671
	log.Debug("GetCollectionStatistics enqueued",
		zap.String("traceID", traceID),
672
		zap.String("role", typeutil.ProxyRole),
673 674 675
		zap.Int64("MsgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
676 677
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
678 679 680

	if err := g.WaitToFinish(); err != nil {
		log.Warn("GetCollectionStatistics failed to WaitToFinish",
D
dragondriver 已提交
681
			zap.Error(err),
682
			zap.String("traceID", traceID),
683
			zap.String("role", typeutil.ProxyRole),
684 685 686
			zap.Int64("MsgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
687 688 689
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

G
godchen 已提交
690
		return &milvuspb.GetCollectionStatisticsResponse{
691
			Status: &commonpb.Status{
692
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
693 694 695 696 697
				Reason:    err.Error(),
			},
		}, nil
	}

698 699
	log.Debug("GetCollectionStatistics done",
		zap.String("traceID", traceID),
700
		zap.String("role", typeutil.ProxyRole),
701 702 703 704 705 706
		zap.Int64("MsgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

707
	return g.result, nil
708 709
}

710
// ShowCollections list all collections in Milvus.
C
Cai Yudong 已提交
711
func (node *Proxy) ShowCollections(ctx context.Context, request *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
712 713 714 715 716
	if !node.checkHealthy() {
		return &milvuspb.ShowCollectionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
717
	sct := &showCollectionsTask{
G
godchen 已提交
718 719 720
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
721
		queryCoord:             node.queryCoord,
722
		rootCoord:              node.rootCoord,
723 724
	}

725
	log.Debug("ShowCollections received",
726
		zap.String("role", typeutil.ProxyRole),
727 728 729 730 731 732
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

733
	err := node.sched.ddQueue.Enqueue(sct)
734
	if err != nil {
735 736
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
737
			zap.String("role", typeutil.ProxyRole),
738 739 740 741 742 743
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

G
godchen 已提交
744
		return &milvuspb.ShowCollectionsResponse{
745
			Status: &commonpb.Status{
746
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
747 748 749 750 751
				Reason:    err.Error(),
			},
		}, nil
	}

752
	log.Debug("ShowCollections enqueued",
753
		zap.String("role", typeutil.ProxyRole),
754
		zap.Int64("MsgID", sct.ID()),
755
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
756
		zap.Uint64("TimeStamp", request.TimeStamp),
757 758 759
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
760

761 762
	err = sct.WaitToFinish()
	if err != nil {
763 764
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
765
			zap.String("role", typeutil.ProxyRole),
766 767 768 769 770 771 772
			zap.Int64("MsgID", sct.ID()),
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

G
godchen 已提交
773
		return &milvuspb.ShowCollectionsResponse{
774
			Status: &commonpb.Status{
775
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
776 777 778 779 780
				Reason:    err.Error(),
			},
		}, nil
	}

781
	log.Debug("ShowCollections Done",
782
		zap.String("role", typeutil.ProxyRole),
783 784 785 786 787 788 789 790
		zap.Int64("MsgID", sct.ID()),
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
		zap.Any("result", sct.result),
	)

791 792 793
	return sct.result, nil
}

794
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
795
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
796 797 798
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
799

800
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
801 802 803
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

804
	cpt := &createPartitionTask{
S
sunby 已提交
805
		ctx:                    ctx,
806 807
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
808
		rootCoord:              node.rootCoord,
809 810 811
		result:                 nil,
	}

812 813 814
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
815
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
816 817 818
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
819 820 821 822 823 824

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
825
			zap.String("role", typeutil.ProxyRole),
826 827 828 829
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

830
		return &commonpb.Status{
831
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
832 833 834
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
835

836 837 838
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
839
		zap.String("role", typeutil.ProxyRole),
840 841 842
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
843 844 845
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
846 847 848 849

	if err := cpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish("CreatePartition"),
D
dragondriver 已提交
850
			zap.Error(err),
851
			zap.String("traceID", traceID),
852
			zap.String("role", typeutil.ProxyRole),
853 854 855
			zap.Int64("MsgID", cpt.ID()),
			zap.Uint64("BeginTS", cpt.BeginTs()),
			zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
856 857 858 859
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

860
		return &commonpb.Status{
861
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
862 863 864
			Reason:    err.Error(),
		}, nil
	}
865 866 867 868

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
869
		zap.String("role", typeutil.ProxyRole),
870 871 872 873 874 875 876
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

877 878 879
	return cpt.result, nil
}

880
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
881
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
882 883 884
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
885

886
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
887 888 889
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

890
	dpt := &dropPartitionTask{
S
sunby 已提交
891
		ctx:                  ctx,
892 893
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
894
		rootCoord:            node.rootCoord,
895 896 897
		result:               nil,
	}

898 899 900 901 902
	method := "DropPartition"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
903
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
904 905 906
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
907 908 909 910 911 912

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
913
			zap.String("role", typeutil.ProxyRole),
914 915 916 917
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

918
		return &commonpb.Status{
919
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
920 921 922
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
923

924 925 926
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
927
		zap.String("role", typeutil.ProxyRole),
928 929 930
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
931 932 933
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
934 935 936 937

	if err := dpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
938
			zap.Error(err),
939
			zap.String("traceID", traceID),
940
			zap.String("role", typeutil.ProxyRole),
941 942 943
			zap.Int64("MsgID", dpt.ID()),
			zap.Uint64("BeginTS", dpt.BeginTs()),
			zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
944 945 946 947
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

948
		return &commonpb.Status{
949
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
950 951 952
			Reason:    err.Error(),
		}, nil
	}
953 954 955 956

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
957
		zap.String("role", typeutil.ProxyRole),
958 959 960 961 962 963 964
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

965 966 967
	return dpt.result, nil
}

968
// HasPartition check if partition exist.
C
Cai Yudong 已提交
969
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
970 971 972 973 974
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
975

D
dragondriver 已提交
976
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
977 978 979
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

980
	hpt := &hasPartitionTask{
S
sunby 已提交
981
		ctx:                 ctx,
982 983
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
984
		rootCoord:           node.rootCoord,
985 986 987
		result:              nil,
	}

D
dragondriver 已提交
988 989 990 991 992
	method := "HasPartition"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
993
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
994 995 996
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
997 998 999 1000 1001 1002

	if err := node.sched.ddQueue.Enqueue(hpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1003
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1004 1005 1006 1007
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1008 1009
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1010
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1011 1012 1013 1014 1015
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1016

D
dragondriver 已提交
1017 1018 1019
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1020
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1021 1022 1023
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1024 1025 1026
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1027 1028 1029 1030

	if err := hpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1031
			zap.Error(err),
D
dragondriver 已提交
1032
			zap.String("traceID", traceID),
1033
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1034 1035 1036
			zap.Int64("MsgID", hpt.ID()),
			zap.Uint64("BeginTS", hpt.BeginTs()),
			zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1037 1038 1039 1040
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1041 1042
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1043
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1044 1045 1046 1047 1048
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1049 1050 1051 1052

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1053
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1054 1055 1056 1057 1058 1059 1060
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

1061 1062 1063
	return hpt.result, nil
}

1064
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1065
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1066 1067 1068
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1069

D
dragondriver 已提交
1070
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1071 1072 1073
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1074
	lpt := &loadPartitionsTask{
G
godchen 已提交
1075 1076 1077
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1078
		queryCoord:            node.queryCoord,
1079 1080
	}

1081 1082 1083 1084 1085
	method := "LoadPartitions"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1086
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1087 1088 1089
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1090 1091 1092 1093 1094 1095

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1096
			zap.String("role", typeutil.ProxyRole),
1097 1098 1099 1100
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1101
		return &commonpb.Status{
1102
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1103 1104 1105 1106
			Reason:    err.Error(),
		}, nil
	}

1107 1108 1109
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1110
		zap.String("role", typeutil.ProxyRole),
1111 1112 1113
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1114 1115 1116
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1117 1118 1119 1120

	if err := lpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1121
			zap.Error(err),
1122
			zap.String("traceID", traceID),
1123
			zap.String("role", typeutil.ProxyRole),
1124 1125 1126
			zap.Int64("MsgID", lpt.ID()),
			zap.Uint64("BeginTS", lpt.BeginTs()),
			zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1127 1128 1129 1130
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1131
		return &commonpb.Status{
1132
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1133 1134 1135 1136
			Reason:    err.Error(),
		}, nil
	}

1137 1138 1139
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1140
		zap.String("role", typeutil.ProxyRole),
1141 1142 1143 1144 1145 1146 1147
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))

1148
	return lpt.result, nil
1149 1150
}

1151
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1152
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1153 1154 1155
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1156 1157 1158 1159 1160

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleasePartitions")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1161
	rpt := &releasePartitionsTask{
G
godchen 已提交
1162 1163 1164
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1165
		queryCoord:               node.queryCoord,
1166 1167
	}

1168 1169 1170 1171 1172
	method := "ReleasePartitions"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1173
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1174 1175 1176
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1177 1178 1179 1180 1181 1182

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1183
			zap.String("role", typeutil.ProxyRole),
1184 1185 1186 1187
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1188
		return &commonpb.Status{
1189
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1190 1191 1192 1193
			Reason:    err.Error(),
		}, nil
	}

1194 1195 1196
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1197
		zap.String("role", typeutil.ProxyRole),
1198 1199 1200
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1201 1202 1203
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1204 1205 1206 1207

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1208
			zap.Error(err),
1209
			zap.String("traceID", traceID),
1210
			zap.String("role", typeutil.ProxyRole),
1211 1212 1213
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1214 1215 1216 1217
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1218
		return &commonpb.Status{
1219
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1220 1221 1222 1223
			Reason:    err.Error(),
		}, nil
	}

1224 1225 1226
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1227
		zap.String("role", typeutil.ProxyRole),
1228 1229 1230 1231 1232 1233 1234
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))

1235
	return rpt.result, nil
1236 1237
}

1238
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1239
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1240 1241 1242 1243 1244
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1245 1246 1247 1248 1249

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetPartitionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1250
	g := &getPartitionStatisticsTask{
1251 1252 1253
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1254
		dataCoord:                     node.dataCoord,
1255 1256
	}

1257 1258 1259 1260 1261
	method := "GetPartitionStatistics"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1262
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1263 1264 1265
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1266 1267 1268 1269 1270 1271

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1272
			zap.String("role", typeutil.ProxyRole),
1273 1274 1275 1276
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1277 1278 1279 1280 1281 1282 1283 1284
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1285 1286 1287
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1288
		zap.String("role", typeutil.ProxyRole),
1289 1290 1291
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1292 1293 1294
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1295 1296 1297 1298

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1299
			zap.Error(err),
1300
			zap.String("traceID", traceID),
1301
			zap.String("role", typeutil.ProxyRole),
1302 1303 1304
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1317 1318 1319
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1320
		zap.String("role", typeutil.ProxyRole),
1321 1322 1323 1324 1325 1326 1327
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

1328
	return g.result, nil
1329 1330
}

1331
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1332
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1333 1334 1335 1336 1337
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1338 1339 1340 1341 1342

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ShowPartitions")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1343
	spt := &showPartitionsTask{
G
godchen 已提交
1344 1345 1346
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1347
		rootCoord:             node.rootCoord,
1348
		queryCoord:            node.queryCoord,
G
godchen 已提交
1349
		result:                nil,
1350 1351
	}

1352 1353 1354 1355 1356
	method := "ShowPartitions"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1357
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1358
		zap.Any("request", request))
1359 1360 1361 1362 1363 1364

	if err := node.sched.ddQueue.Enqueue(spt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1365
			zap.String("role", typeutil.ProxyRole),
1366 1367
			zap.Any("request", request))

G
godchen 已提交
1368
		return &milvuspb.ShowPartitionsResponse{
1369
			Status: &commonpb.Status{
1370
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1371 1372 1373 1374 1375
				Reason:    err.Error(),
			},
		}, nil
	}

1376 1377 1378
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1379
		zap.String("role", typeutil.ProxyRole),
1380 1381 1382
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1383 1384
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1385 1386 1387 1388 1389
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1390
			zap.Error(err),
1391
			zap.String("traceID", traceID),
1392
			zap.String("role", typeutil.ProxyRole),
1393 1394 1395 1396 1397 1398
			zap.Int64("msgID", spt.ID()),
			zap.Uint64("BeginTS", spt.BeginTs()),
			zap.Uint64("EndTS", spt.EndTs()),
			zap.String("db", spt.ShowPartitionsRequest.DbName),
			zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
			zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))
D
dragondriver 已提交
1399

G
godchen 已提交
1400
		return &milvuspb.ShowPartitionsResponse{
1401
			Status: &commonpb.Status{
1402
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1403 1404 1405 1406
				Reason:    err.Error(),
			},
		}, nil
	}
1407 1408 1409 1410

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1411
		zap.String("role", typeutil.ProxyRole),
1412 1413 1414 1415 1416 1417 1418
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

1419 1420 1421
	return spt.result, nil
}

1422
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1423
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1424 1425 1426
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1427 1428 1429 1430 1431

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ShowPartitions")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1432
	cit := &createIndexTask{
S
sunby 已提交
1433
		ctx:                ctx,
1434 1435
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1436
		rootCoord:          node.rootCoord,
1437 1438
	}

D
dragondriver 已提交
1439 1440 1441 1442 1443
	method := "CreateIndex"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1444
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1445 1446 1447 1448
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1449 1450 1451 1452 1453 1454

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1455
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1456 1457 1458 1459 1460
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1461
		return &commonpb.Status{
1462
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1463 1464 1465 1466
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1467 1468 1469
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1470
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1471 1472 1473
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1474 1475 1476 1477
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1478 1479 1480 1481

	if err := cit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1482
			zap.Error(err),
D
dragondriver 已提交
1483
			zap.String("traceID", traceID),
1484
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1485 1486 1487
			zap.Int64("MsgID", cit.ID()),
			zap.Uint64("BeginTs", cit.BeginTs()),
			zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1488 1489 1490 1491 1492
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1493
		return &commonpb.Status{
1494
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1495 1496 1497 1498
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1499 1500 1501
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1502
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1503 1504 1505 1506 1507 1508 1509 1510
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))

1511 1512 1513
	return cit.result, nil
}

1514
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
1515
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
1516 1517 1518 1519 1520
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1521 1522 1523 1524 1525

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeIndex")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1526
	dit := &describeIndexTask{
S
sunby 已提交
1527
		ctx:                  ctx,
1528 1529
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1530
		rootCoord:            node.rootCoord,
1531 1532
	}

1533 1534 1535 1536 1537 1538 1539
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1540
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1541 1542 1543
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1544 1545 1546 1547 1548 1549 1550
		zap.String("index name", indexName))

	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1551
			zap.String("role", typeutil.ProxyRole),
1552 1553 1554 1555 1556
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

1557 1558
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1559
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1560 1561 1562 1563 1564
				Reason:    err.Error(),
			},
		}, nil
	}

1565 1566 1567
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1568
		zap.String("role", typeutil.ProxyRole),
1569 1570 1571
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1572 1573 1574
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1575 1576 1577 1578 1579
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1580
			zap.Error(err),
1581
			zap.String("traceID", traceID),
1582
			zap.String("role", typeutil.ProxyRole),
1583 1584 1585
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1586 1587 1588
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
1589
			zap.String("index name", indexName))
D
dragondriver 已提交
1590

Z
zhenshan.cao 已提交
1591 1592 1593 1594
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
1595 1596
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
1597
				ErrorCode: errCode,
1598 1599 1600 1601 1602
				Reason:    err.Error(),
			},
		}, nil
	}

1603 1604 1605
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1606
		zap.String("role", typeutil.ProxyRole),
1607 1608 1609 1610 1611 1612 1613 1614
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", indexName))

1615 1616 1617
	return dit.result, nil
}

1618
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
1619
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
1620 1621 1622
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1623 1624 1625 1626 1627

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropIndex")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1628
	dit := &dropIndexTask{
S
sunby 已提交
1629
		ctx:              ctx,
B
BossZou 已提交
1630 1631
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1632
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1633
	}
G
godchen 已提交
1634

D
dragondriver 已提交
1635 1636 1637 1638 1639
	method := "DropIndex"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1640
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1641 1642 1643 1644 1645
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
1646 1647 1648 1649 1650
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1651
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1652 1653 1654 1655 1656
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

B
BossZou 已提交
1657
		return &commonpb.Status{
1658
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1659 1660 1661
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1662

D
dragondriver 已提交
1663 1664 1665
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1666
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1667 1668 1669
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1670 1671 1672 1673
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
1674 1675 1676 1677

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1678
			zap.Error(err),
D
dragondriver 已提交
1679
			zap.String("traceID", traceID),
1680
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1681 1682 1683
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1684 1685 1686 1687 1688
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

B
BossZou 已提交
1689
		return &commonpb.Status{
1690
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1691 1692 1693
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1694 1695 1696 1697

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1698
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1699 1700 1701 1702 1703 1704 1705 1706
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

B
BossZou 已提交
1707 1708 1709
	return dit.result, nil
}

1710 1711
// GetIndexBuildProgress gets index build progress with filed_name and index_name.
// IndexRows is the num of indexed rows. And TotalRows is the total number of segment rows.
C
Cai Yudong 已提交
1712
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1713 1714 1715 1716 1717
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1718 1719 1720 1721 1722

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetIndexBuildProgress")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1723
	gibpt := &getIndexBuildProgressTask{
1724 1725 1726
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1727 1728
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1729
		dataCoord:                    node.dataCoord,
1730 1731
	}

1732 1733 1734 1735 1736
	method := "GetIndexBuildProgress"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1737
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1738 1739 1740 1741
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1742 1743 1744 1745 1746 1747

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1748
			zap.String("role", typeutil.ProxyRole),
1749 1750 1751 1752 1753
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

1754 1755 1756 1757 1758 1759 1760 1761
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1762 1763 1764
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1765
		zap.String("role", typeutil.ProxyRole),
1766 1767 1768
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
1769 1770 1771 1772
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1773 1774 1775 1776

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1777
			zap.Error(err),
1778
			zap.String("traceID", traceID),
1779
			zap.String("role", typeutil.ProxyRole),
1780 1781 1782
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
1795 1796 1797 1798

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1799
		zap.String("role", typeutil.ProxyRole),
1800 1801 1802 1803 1804 1805 1806 1807
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName),
		zap.Any("result", gibpt.result))
1808 1809

	return gibpt.result, nil
1810 1811
}

1812
// GetIndexState get the build-state of index.
C
Cai Yudong 已提交
1813
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
1814 1815 1816 1817 1818
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1819 1820 1821 1822 1823

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Insert")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1824
	dipt := &getIndexStateTask{
G
godchen 已提交
1825 1826 1827
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
1828 1829
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
1830 1831
	}

1832 1833 1834 1835 1836
	method := "GetIndexState"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1837
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1838 1839 1840 1841
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1842 1843 1844 1845 1846 1847

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1848
			zap.String("role", typeutil.ProxyRole),
1849 1850 1851 1852 1853
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

G
godchen 已提交
1854
		return &milvuspb.GetIndexStateResponse{
1855
			Status: &commonpb.Status{
1856
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1857 1858 1859 1860 1861
				Reason:    err.Error(),
			},
		}, nil
	}

1862 1863 1864
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1865
		zap.String("role", typeutil.ProxyRole),
1866 1867 1868
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
1869 1870 1871 1872
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1873 1874 1875 1876

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1877
			zap.Error(err),
1878
			zap.String("traceID", traceID),
1879
			zap.String("role", typeutil.ProxyRole),
1880 1881 1882
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
1883 1884 1885 1886 1887
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

G
godchen 已提交
1888
		return &milvuspb.GetIndexStateResponse{
1889
			Status: &commonpb.Status{
1890
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1891 1892 1893 1894 1895
				Reason:    err.Error(),
			},
		}, nil
	}

1896 1897 1898
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1899
		zap.String("role", typeutil.ProxyRole),
1900 1901 1902 1903 1904 1905 1906 1907
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

1908 1909 1910
	return dipt.result, nil
}

1911
// Insert insert records into collection.
C
Cai Yudong 已提交
1912
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
1913 1914 1915 1916 1917 1918
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Insert")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
	log.Info("Start processing insert request in Proxy", zap.String("traceID", traceID))
	defer log.Info("Finish processing insert request in Proxy", zap.String("traceID", traceID))

1919 1920 1921 1922 1923
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1924

1925
	it := &insertTask{
1926 1927 1928
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       request,
1929 1930 1931 1932
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
1933
			InsertRequest: internalpb.InsertRequest{
1934
				Base: &commonpb.MsgBase{
1935
					MsgType: commonpb.MsgType_Insert,
1936 1937 1938 1939
					MsgID:   0,
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
1940
				// RowData: transfer column based request to this
1941 1942
			},
		},
1943
		rowIDAllocator: node.idAllocator,
1944
		segIDAssigner:  node.segAssigner,
1945
		chMgr:          node.chMgr,
1946
		chTicker:       node.chTicker,
1947
	}
1948 1949 1950 1951 1952

	if len(it.PartitionName) <= 0 {
		it.PartitionName = Params.DefaultPartitionName
	}

X
Xiangyu Wang 已提交
1953
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
1954 1955 1956 1957 1958
		numRows := it.req.NumRows
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
X
Xiangyu Wang 已提交
1959 1960 1961 1962 1963 1964 1965
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
1966 1967
	}

X
Xiangyu Wang 已提交
1968
	log.Debug("Enqueue insert request in Proxy",
1969
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1970 1971 1972 1973 1974 1975 1976
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.Int("len(FieldsData)", len(request.FieldsData)),
		zap.Int("len(HashKeys)", len(request.HashKeys)),
		zap.Uint32("NumRows", request.NumRows))

X
Xiangyu Wang 已提交
1977 1978 1979
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
		return constructFailedResponse(err), nil
1980
	}
D
dragondriver 已提交
1981

X
Xiangyu Wang 已提交
1982
	log.Debug("Detail of insert request in Proxy",
1983
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
1984
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
1985 1986 1987 1988 1989 1990 1991
		zap.Uint64("BeginTS", it.BeginTs()),
		zap.Uint64("EndTS", it.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.Int("len(FieldsData)", len(request.FieldsData)),
		zap.Int("len(HashKeys)", len(request.HashKeys)),
X
Xiangyu Wang 已提交
1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))

	if err := it.WaitToFinish(); err != nil {
		log.Debug("Failed to execute insert task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
			numRows := it.req.NumRows
			errIndex := make([]uint32, numRows)
			for i := uint32(0); i < numRows; i++ {
				errIndex[i] = i
			}
			it.result.ErrIndex = errIndex
		}

		setErrorIndex()
	}

	// InsertCnt always equals to the number of entities in the request
	it.result.InsertCnt = int64(it.req.NumRows)
D
dragondriver 已提交
2015

2016 2017 2018
	return it.result, nil
}

2019
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2020
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2021 2022 2023
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2024 2025
	log.Info("Start processing delete request in Proxy", zap.String("traceID", traceID))
	defer log.Info("Finish processing delete request in Proxy", zap.String("traceID", traceID))
2026

G
groot 已提交
2027 2028 2029 2030 2031 2032
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

C
Cai Yudong 已提交
2033 2034 2035 2036 2037 2038 2039
	deleteReq := &milvuspb.DeleteRequest{
		DbName:         request.DbName,
		CollectionName: request.CollectionName,
		PartitionName:  request.PartitionName,
		Expr:           request.Expr,
	}

2040
	dt := &deleteTask{
C
Cai Yudong 已提交
2041 2042 2043
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       deleteReq,
G
godchen 已提交
2044
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2045 2046 2047
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2048 2049 2050 2051 2052 2053 2054 2055
			DeleteRequest: internalpb.DeleteRequest{
				Base: &commonpb.MsgBase{
					MsgType: commonpb.MsgType_Delete,
					MsgID:   0,
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2056 2057 2058 2059
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2060 2061
	}

2062
	log.Debug("Enqueue delete request in Proxy",
2063
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2064 2065 2066 2067
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2068 2069 2070 2071

	// MsgID will be set by Enqueue()
	if err := node.sched.dmQueue.Enqueue(dt); err != nil {
		log.Error("Failed to enqueue delete task: "+err.Error(), zap.String("traceID", traceID))
G
groot 已提交
2072 2073 2074 2075 2076 2077 2078 2079
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2080
	log.Debug("Detail of delete request in Proxy",
2081
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2082 2083 2084 2085 2086
		zap.Int64("msgID", dt.Base.MsgID),
		zap.Uint64("timestamp", dt.Base.Timestamp),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
2087 2088
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2089

2090 2091
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
G
groot 已提交
2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

	return dt.result, nil
}

2103
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2104
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2105 2106 2107 2108 2109
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
2110

C
cai.zhang 已提交
2111 2112
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2113 2114
	traceID, _, _ := trace.InfoFromSpan(sp)

2115
	qt := &searchTask{
S
sunby 已提交
2116
		ctx:       ctx,
2117
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2118
		SearchRequest: &internalpb.SearchRequest{
2119
			Base: &commonpb.MsgBase{
2120
				MsgType:  commonpb.MsgType_Search,
2121
				SourceID: Params.ProxyID,
2122
			},
2123
			ResultChannelID: strconv.FormatInt(Params.ProxyID, 10),
2124
		},
2125 2126 2127
		resultBuf: make(chan []*internalpb.SearchResults),
		query:     request,
		chMgr:     node.chMgr,
2128
		qc:        node.queryCoord,
2129 2130
	}

D
dragondriver 已提交
2131 2132
	log.Debug("Search received",
		zap.String("traceID", traceID),
2133
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2134 2135 2136 2137 2138 2139
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames),
		zap.Any("dsl", request.Dsl),
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
		zap.Any("OutputFields", request.OutputFields))
D
dragondriver 已提交
2140

2141
	err := node.sched.dqQueue.Enqueue(qt)
2142
	if err != nil {
D
dragondriver 已提交
2143 2144 2145
		log.Debug("Search failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
2146
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2147 2148 2149 2150 2151 2152 2153 2154
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
			zap.Any("OutputFields", request.OutputFields),
		)

2155 2156
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2157
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2158 2159 2160 2161 2162
				Reason:    err.Error(),
			},
		}, nil
	}

D
dragondriver 已提交
2163 2164
	log.Debug("Search enqueued",
		zap.String("traceID", traceID),
2165
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2166
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2167 2168 2169 2170 2171
		zap.Uint64("timestamp", qt.Base.Timestamp),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames),
		zap.Any("dsl", request.Dsl),
2172 2173
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
		zap.Any("OutputFields", request.OutputFields))
D
dragondriver 已提交
2174 2175 2176 2177 2178

	err = qt.WaitToFinish()

	if err != nil {
		log.Debug("Search failed to WaitToFinish",
D
dragondriver 已提交
2179
			zap.Error(err),
D
dragondriver 已提交
2180
			zap.String("traceID", traceID),
2181
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2182
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2183 2184 2185 2186
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2187 2188
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
			zap.Any("OutputFields", request.OutputFields))
2189

2190 2191
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2192
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2193 2194 2195 2196 2197
				Reason:    err.Error(),
			},
		}, nil
	}

D
dragondriver 已提交
2198 2199
	log.Debug("Search Done",
		zap.String("traceID", traceID),
2200
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2201 2202 2203 2204 2205 2206 2207 2208
		zap.Int64("msgID", qt.ID()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames),
		zap.Any("dsl", request.Dsl),
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
		zap.Any("OutputFields", request.OutputFields))

2209 2210 2211
	return qt.result, nil
}

2212
// Flush notify data nodes to persist the data of collection.
2213 2214 2215 2216 2217 2218 2219
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2220
	if !node.checkHealthy() {
2221 2222
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2223
	}
D
dragondriver 已提交
2224 2225 2226 2227 2228

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Flush")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

2229
	ft := &flushTask{
T
ThreadDao 已提交
2230 2231 2232
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2233
		dataCoord:    node.dataCoord,
2234 2235
	}

D
dragondriver 已提交
2236 2237 2238 2239 2240
	method := "Flush"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2241
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2242 2243
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2244 2245 2246 2247 2248 2249

	if err := node.sched.ddQueue.Enqueue(ft); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2250
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2251 2252 2253
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

2254 2255
		resp.Status.Reason = err.Error()
		return resp, nil
2256 2257
	}

D
dragondriver 已提交
2258 2259 2260
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2261
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2262 2263 2264
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2265 2266
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2267 2268 2269 2270

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2271
			zap.Error(err),
D
dragondriver 已提交
2272
			zap.String("traceID", traceID),
2273
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2274 2275 2276
			zap.Int64("MsgID", ft.ID()),
			zap.Uint64("BeginTs", ft.BeginTs()),
			zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2277 2278 2279
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

D
dragondriver 已提交
2280
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2281 2282
		resp.Status.Reason = err.Error()
		return resp, nil
2283 2284
	}

D
dragondriver 已提交
2285 2286 2287
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2288
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2289 2290 2291 2292 2293 2294
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))

2295
	return ft.result, nil
2296 2297
}

2298
// Query get the records by primary keys.
C
Cai Yudong 已提交
2299
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2300 2301 2302 2303 2304
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2305

D
dragondriver 已提交
2306 2307 2308 2309
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

2310 2311 2312 2313 2314 2315
	queryRequest := &milvuspb.QueryRequest{
		DbName:         request.DbName,
		CollectionName: request.CollectionName,
		PartitionNames: request.PartitionNames,
		Expr:           request.Expr,
		OutputFields:   request.OutputFields,
2316
	}
2317

2318
	qt := &queryTask{
2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
				SourceID: Params.ProxyID,
			},
			ResultChannelID: strconv.FormatInt(Params.ProxyID, 10),
		},
		resultBuf: make(chan []*internalpb.RetrieveResults),
		query:     queryRequest,
		chMgr:     node.chMgr,
		qc:        node.queryCoord,
2332 2333
	}

D
dragondriver 已提交
2334 2335 2336 2337 2338
	method := "Query"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2339
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2340 2341 2342 2343
		zap.String("db", queryRequest.DbName),
		zap.String("collection", queryRequest.CollectionName),
		zap.Any("partitions", queryRequest.PartitionNames))

D
dragondriver 已提交
2344 2345 2346 2347 2348 2349 2350 2351 2352 2353
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames))

2354 2355 2356 2357 2358 2359
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2360 2361
	}

D
dragondriver 已提交
2362 2363 2364
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2365
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2366 2367 2368
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
2369 2370 2371
		zap.String("db", queryRequest.DbName),
		zap.String("collection", queryRequest.CollectionName),
		zap.Any("partitions", queryRequest.PartitionNames))
D
dragondriver 已提交
2372 2373 2374 2375 2376 2377

	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2378
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2379 2380 2381
			zap.Int64("MsgID", qt.ID()),
			zap.Uint64("BeginTs", qt.BeginTs()),
			zap.Uint64("EndTs", qt.EndTs()),
2382 2383 2384
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames))
2385

2386 2387 2388 2389 2390 2391 2392
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2393

D
dragondriver 已提交
2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
		zap.String("db", queryRequest.DbName),
		zap.String("collection", queryRequest.CollectionName),
		zap.Any("partitions", queryRequest.PartitionNames))

2405 2406 2407 2408 2409
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2410

2411
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
2412 2413 2414 2415
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2416 2417 2418 2419 2420

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateAlias")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

Y
Yusup 已提交
2421 2422 2423 2424 2425 2426 2427
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
2428 2429 2430 2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 2442 2443 2444 2445 2446 2447
	method := "CreateAlias"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

	if err := node.sched.ddQueue.Enqueue(cat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

Y
Yusup 已提交
2448 2449 2450 2451 2452 2453
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2454 2455 2456
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2457
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2458 2459 2460 2461
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2462 2463
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2464 2465 2466 2467

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2468
			zap.Error(err),
D
dragondriver 已提交
2469
			zap.String("traceID", traceID),
2470
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2471 2472 2473 2474
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2475 2476 2477 2478 2479 2480 2481 2482 2483
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2484 2485 2486 2487 2488 2489 2490 2491 2492 2493 2494
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

Y
Yusup 已提交
2495 2496 2497
	return cat.result, nil
}

2498
// DropAlias alter the alias of collection.
Y
Yusup 已提交
2499 2500 2501 2502
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2503 2504 2505 2506 2507

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropAlias")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

Y
Yusup 已提交
2508 2509 2510 2511 2512 2513 2514
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
2515 2516 2517 2518 2519 2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532
	method := "DropAlias"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias))

	if err := node.sched.ddQueue.Enqueue(dat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias))

Y
Yusup 已提交
2533 2534 2535 2536 2537 2538
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2539 2540 2541
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2542
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2543 2544 2545 2546
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2547
		zap.String("alias", request.Alias))
D
dragondriver 已提交
2548 2549 2550 2551

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2552
			zap.Error(err),
D
dragondriver 已提交
2553
			zap.String("traceID", traceID),
2554
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2555 2556 2557 2558
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2559 2560 2561 2562 2563 2564 2565 2566
			zap.String("alias", request.Alias))

		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2567 2568 2569 2570 2571 2572 2573 2574 2575 2576
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias))

Y
Yusup 已提交
2577 2578 2579
	return dat.result, nil
}

2580
// AlterAlias alter alias of collection.
Y
Yusup 已提交
2581 2582 2583 2584
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2585 2586 2587 2588 2589

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-AlterAlias")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

Y
Yusup 已提交
2590 2591 2592 2593 2594 2595 2596
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
2597 2598 2599 2600 2601 2602 2603 2604 2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616
	method := "AlterAlias"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

	if err := node.sched.ddQueue.Enqueue(aat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

Y
Yusup 已提交
2617 2618 2619 2620 2621 2622
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2623 2624 2625
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2626
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2627 2628 2629 2630
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2631 2632
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2633 2634 2635 2636

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2637
			zap.Error(err),
D
dragondriver 已提交
2638
			zap.String("traceID", traceID),
2639
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2640 2641 2642 2643
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2644 2645 2646 2647 2648 2649 2650 2651 2652
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2653 2654 2655 2656 2657 2658 2659 2660 2661 2662 2663
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

Y
Yusup 已提交
2664 2665 2666
	return aat.result, nil
}

2667
// CalcDistance calculates the distances between vectors.
2668
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
2669 2670 2671 2672 2673 2674
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}

2675
	param, _ := funcutil.GetAttrByKeyFromRepeatedKV("metric", request.GetParams())
2676 2677 2678 2679 2680 2681 2682 2683
	metric, err := distance.ValidateMetricType(param)
	if err != nil {
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2684 2685
	}

2686 2687 2688 2689
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

2690 2691
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
2692

2693 2694 2695 2696 2697
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
2698 2699
		}

2700
		qt := &queryTask{
2701 2702 2703 2704 2705 2706 2707 2708 2709
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
					SourceID: Params.ProxyID,
				},
				ResultChannelID: strconv.FormatInt(Params.ProxyID, 10),
			},
2710
			resultBuf: make(chan []*internalpb.RetrieveResults),
2711
			query:     queryRequest,
2712
			chMgr:     node.chMgr,
2713
			qc:        node.queryCoord,
Y
yukun 已提交
2714
			ids:       ids.IdArray,
2715 2716
		}

2717
		err := node.sched.dqQueue.Enqueue(qt)
2718
		if err != nil {
2719 2720 2721
			log.Debug("CalcDistance queryTask failed to enqueue",
				zap.Error(err),
				zap.String("traceID", traceID),
2722
				zap.String("role", typeutil.ProxyRole),
2723 2724 2725 2726
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames))

2727 2728 2729 2730 2731
			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
2732
			}, err
2733
		}
2734 2735 2736

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
2737
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2738 2739 2740 2741 2742 2743
			zap.Int64("msgID", qt.Base.MsgID),
			zap.Uint64("timestamp", qt.Base.Timestamp),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields))
2744 2745 2746 2747

		err = qt.WaitToFinish()
		if err != nil {
			log.Debug("CalcDistance queryTask failed to WaitToFinish",
G
godchen 已提交
2748
				zap.Error(err),
2749
				zap.String("traceID", traceID),
2750
				zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2751 2752 2753 2754 2755 2756
				zap.Int64("msgID", qt.Base.MsgID),
				zap.Uint64("timestamp", qt.Base.Timestamp),
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames),
				zap.Any("OutputFields", queryRequest.OutputFields))
2757 2758 2759 2760 2761 2762

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
2763
			}, err
2764
		}
2765 2766 2767

		log.Debug("CalcDistance queryTask Done",
			zap.String("traceID", traceID),
2768
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2769 2770 2771 2772 2773 2774
			zap.Int64("msgID", qt.Base.MsgID),
			zap.Uint64("timestamp", qt.Base.Timestamp),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields))
2775 2776

		return &milvuspb.QueryResults{
2777 2778
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
2779 2780 2781
		}, nil
	}

2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792 2793 2794 2795
	// the vectors retrieved are random order, we need re-arrange the vectors by the order of input ids
	arrangeFunc := func(ids *milvuspb.VectorIDs, retrievedFields []*schemapb.FieldData) (*schemapb.VectorField, error) {
		var retrievedIds *schemapb.ScalarField
		var retrievedVectors *schemapb.VectorField
		for _, fieldData := range retrievedFields {
			if fieldData.FieldName == ids.FieldName {
				retrievedVectors = fieldData.GetVectors()
			}
			if fieldData.Type == schemapb.DataType_Int64 {
				retrievedIds = fieldData.GetScalars()
			}
		}

		if retrievedIds == nil || retrievedVectors == nil {
2796
			return nil, errors.New("failed to fetch vectors")
2797 2798 2799 2800 2801 2802 2803 2804 2805 2806 2807 2808 2809 2810 2811 2812
		}

		dict := make(map[int64]int)
		for index, id := range retrievedIds.GetLongData().Data {
			dict[id] = index
		}

		inputIds := ids.IdArray.GetIntId().Data
		if retrievedVectors.GetFloatVector() != nil {
			floatArr := retrievedVectors.GetFloatVector().Data
			element := retrievedVectors.GetDim()
			result := make([]float32, 0, int64(len(inputIds))*element)
			for _, id := range inputIds {
				index, ok := dict[id]
				if !ok {
					log.Error("id not found in CalcDistance", zap.Int64("id", id))
2813
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
2814 2815 2816 2817 2818 2819 2820 2821 2822 2823 2824 2825 2826 2827 2828 2829 2830 2831 2832 2833 2834 2835 2836 2837 2838 2839
				}
				result = append(result, floatArr[int64(index)*element:int64(index+1)*element]...)
			}

			return &schemapb.VectorField{
				Dim: element,
				Data: &schemapb.VectorField_FloatVector{
					FloatVector: &schemapb.FloatArray{
						Data: result,
					},
				},
			}, nil
		}

		if retrievedVectors.GetBinaryVector() != nil {
			binaryArr := retrievedVectors.GetBinaryVector()
			element := retrievedVectors.GetDim()
			if element%8 != 0 {
				element = element + 8 - element%8
			}

			result := make([]byte, 0, int64(len(inputIds))*element)
			for _, id := range inputIds {
				index, ok := dict[id]
				if !ok {
					log.Error("id not found in CalcDistance", zap.Int64("id", id))
2840
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
2841 2842 2843 2844 2845 2846 2847 2848 2849 2850 2851 2852
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

			return &schemapb.VectorField{
				Dim: element * 8,
				Data: &schemapb.VectorField_BinaryVector{
					BinaryVector: result,
				},
			}, nil
		}

2853
		return nil, errors.New("failed to fetch vectors")
2854 2855
	}

2856 2857
	log.Debug("CalcDistance received",
		zap.String("traceID", traceID),
2858
		zap.String("role", typeutil.ProxyRole),
2859
		zap.String("metric", metric))
G
godchen 已提交
2860

2861 2862 2863
	vectorsLeft := request.GetOpLeft().GetDataArray()
	opLeft := request.GetOpLeft().GetIdArray()
	if opLeft != nil {
2864 2865
		log.Debug("OpLeft IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
2866
			zap.String("role", typeutil.ProxyRole))
2867

2868
		result, err := query(opLeft)
2869
		if err != nil {
2870 2871 2872
			log.Debug("Failed to get left vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
2873
				zap.String("role", typeutil.ProxyRole))
2874

2875 2876 2877 2878 2879 2880 2881 2882
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

2883 2884
		log.Debug("OpLeft IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
2885
			zap.String("role", typeutil.ProxyRole))
2886

2887 2888
		vectorsLeft, err = arrangeFunc(opLeft, result.FieldsData)
		if err != nil {
2889 2890 2891
			log.Debug("Failed to re-arrange left vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
2892
				zap.String("role", typeutil.ProxyRole))
2893

2894 2895 2896 2897 2898 2899
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
2900
		}
2901 2902 2903

		log.Debug("Re-arrange left vectors done",
			zap.String("traceID", traceID),
2904
			zap.String("role", typeutil.ProxyRole))
2905 2906
	}

G
groot 已提交
2907
	if vectorsLeft == nil {
2908 2909 2910
		msg := "Left vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
2911
			zap.String("role", typeutil.ProxyRole))
2912

G
groot 已提交
2913 2914 2915
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2916
				Reason:    msg,
G
groot 已提交
2917 2918 2919 2920
			},
		}, nil
	}

2921 2922 2923
	vectorsRight := request.GetOpRight().GetDataArray()
	opRight := request.GetOpRight().GetIdArray()
	if opRight != nil {
2924 2925
		log.Debug("OpRight IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
2926
			zap.String("role", typeutil.ProxyRole))
2927

2928
		result, err := query(opRight)
2929
		if err != nil {
2930 2931 2932
			log.Debug("Failed to get right vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
2933
				zap.String("role", typeutil.ProxyRole))
2934

2935 2936 2937 2938 2939 2940 2941 2942
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

2943 2944
		log.Debug("OpRight IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
2945
			zap.String("role", typeutil.ProxyRole))
2946

2947 2948
		vectorsRight, err = arrangeFunc(opRight, result.FieldsData)
		if err != nil {
2949 2950 2951
			log.Debug("Failed to re-arrange right vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
2952
				zap.String("role", typeutil.ProxyRole))
2953

2954 2955 2956 2957 2958 2959
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
2960
		}
2961 2962 2963

		log.Debug("Re-arrange right vectors done",
			zap.String("traceID", traceID),
2964
			zap.String("role", typeutil.ProxyRole))
2965 2966
	}

G
groot 已提交
2967
	if vectorsRight == nil {
2968 2969 2970
		msg := "Right vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
2971
			zap.String("role", typeutil.ProxyRole))
2972

G
groot 已提交
2973 2974 2975
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2976
				Reason:    msg,
G
groot 已提交
2977 2978 2979 2980
			},
		}, nil
	}

2981
	if vectorsLeft.Dim != vectorsRight.Dim {
2982 2983 2984
		msg := "Vectors dimension is not equal"
		log.Debug(msg,
			zap.String("traceID", traceID),
2985
			zap.String("role", typeutil.ProxyRole))
2986

2987 2988 2989
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2990
				Reason:    msg,
2991 2992 2993 2994 2995
			},
		}, nil
	}

	if vectorsLeft.GetFloatVector() != nil && vectorsRight.GetFloatVector() != nil {
2996 2997
		distances, err := distance.CalcFloatDistance(vectorsLeft.Dim, vectorsLeft.GetFloatVector().Data, vectorsRight.GetFloatVector().Data, metric)
		if err != nil {
2998 2999 3000
			log.Debug("Failed to CalcFloatDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3001
				zap.String("role", typeutil.ProxyRole))
3002

3003 3004 3005 3006 3007 3008 3009 3010
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3011 3012 3013
		log.Debug("CalcFloatDistance done",
			zap.Error(err),
			zap.String("traceID", traceID),
3014
			zap.String("role", typeutil.ProxyRole))
3015

3016 3017 3018 3019 3020 3021 3022 3023 3024 3025
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
			Array: &milvuspb.CalcDistanceResults_FloatDist{
				FloatDist: &schemapb.FloatArray{
					Data: distances,
				},
			},
		}, nil
	}

3026
	if vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetBinaryVector() != nil {
G
groot 已提交
3027
		hamming, err := distance.CalcHammingDistance(vectorsLeft.Dim, vectorsLeft.GetBinaryVector(), vectorsRight.GetBinaryVector())
3028
		if err != nil {
3029 3030 3031
			log.Debug("Failed to CalcHammingDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3032
				zap.String("role", typeutil.ProxyRole))
3033

3034 3035 3036 3037 3038 3039 3040 3041 3042
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

		if metric == distance.HAMMING {
3043 3044
			log.Debug("CalcHammingDistance done",
				zap.String("traceID", traceID),
3045
				zap.String("role", typeutil.ProxyRole))
3046

3047 3048 3049 3050
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_IntDist{
					IntDist: &schemapb.IntArray{
G
groot 已提交
3051
						Data: hamming,
3052 3053 3054 3055 3056 3057
					},
				},
			}, nil
		}

		if metric == distance.TANIMOTO {
G
groot 已提交
3058
			tanimoto, err := distance.CalcTanimotoCoefficient(vectorsLeft.Dim, hamming)
3059
			if err != nil {
3060 3061 3062
				log.Debug("Failed to CalcTanimotoCoefficient",
					zap.Error(err),
					zap.String("traceID", traceID),
3063
					zap.String("role", typeutil.ProxyRole))
3064

3065 3066 3067 3068 3069 3070 3071 3072
				return &milvuspb.CalcDistanceResults{
					Status: &commonpb.Status{
						ErrorCode: commonpb.ErrorCode_UnexpectedError,
						Reason:    err.Error(),
					},
				}, nil
			}

3073 3074
			log.Debug("CalcTanimotoCoefficient done",
				zap.String("traceID", traceID),
3075
				zap.String("role", typeutil.ProxyRole))
3076

3077 3078 3079 3080 3081 3082 3083 3084 3085 3086 3087
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_FloatDist{
					FloatDist: &schemapb.FloatArray{
						Data: tanimoto,
					},
				},
			}, nil
		}
	}

3088
	err = errors.New("unexpected error")
3089
	if (vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetFloatVector() != nil) || (vectorsLeft.GetFloatVector() != nil && vectorsRight.GetBinaryVector() != nil) {
3090
		err = errors.New("cannot calculate distance between binary vectors and float vectors")
3091 3092
	}

3093 3094 3095
	log.Debug("Failed to CalcDistance",
		zap.Error(err),
		zap.String("traceID", traceID),
3096
		zap.String("role", typeutil.ProxyRole))
3097

3098 3099 3100 3101 3102 3103 3104 3105
	return &milvuspb.CalcDistanceResults{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		},
	}, nil
}

3106
// GetDdChannel returns the used channel for dd operations.
C
Cai Yudong 已提交
3107
func (node *Proxy) GetDdChannel(ctx context.Context, request *internalpb.GetDdChannelRequest) (*milvuspb.StringResponse, error) {
3108 3109
	panic("implement me")
}
X
XuanYang-cn 已提交
3110

3111
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3112
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3113
	log.Debug("GetPersistentSegmentInfo",
3114
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3115 3116 3117
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3118
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3119
		Status: &commonpb.Status{
3120
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3121 3122
		},
	}
3123 3124 3125 3126
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
G
godchen 已提交
3127
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
X
XuanYang-cn 已提交
3128
	if err != nil {
3129
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3130 3131
		return resp, nil
	}
3132
	infoResp, err := node.dataCoord.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
X
XuanYang-cn 已提交
3133
		Base: &commonpb.MsgBase{
3134
			MsgType:   commonpb.MsgType_SegmentInfo,
X
XuanYang-cn 已提交
3135 3136 3137 3138 3139 3140 3141
			MsgID:     0,
			Timestamp: 0,
			SourceID:  Params.ProxyID,
		},
		SegmentIDs: segments,
	})
	if err != nil {
3142
		log.Debug("GetPersistentSegmentInfo fail", zap.Error(err))
3143
		resp.Status.Reason = fmt.Errorf("dataCoord:GetSegmentInfo, err:%w", err).Error()
X
XuanYang-cn 已提交
3144 3145
		return resp, nil
	}
3146
	log.Debug("GetPersistentSegmentInfo ", zap.Int("len(infos)", len(infoResp.Infos)), zap.Any("status", infoResp.Status))
3147
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3148 3149 3150 3151 3152 3153
		resp.Status.Reason = infoResp.Status.Reason
		return resp, nil
	}
	persistentInfos := make([]*milvuspb.PersistentSegmentInfo, len(infoResp.Infos))
	for i, info := range infoResp.Infos {
		persistentInfos[i] = &milvuspb.PersistentSegmentInfo{
S
sunby 已提交
3154
			SegmentID:    info.ID,
X
XuanYang-cn 已提交
3155 3156
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
S
sunby 已提交
3157
			NumRows:      info.NumOfRows,
X
XuanYang-cn 已提交
3158 3159 3160
			State:        info.State,
		}
	}
3161
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
X
XuanYang-cn 已提交
3162 3163 3164 3165
	resp.Infos = persistentInfos
	return resp, nil
}

J
jingkl 已提交
3166
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3167
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3168
	log.Debug("GetQuerySegmentInfo",
3169
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3170 3171 3172
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3173
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3174
		Status: &commonpb.Status{
3175
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3176 3177
		},
	}
3178 3179 3180 3181
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
G
godchen 已提交
3182
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
Z
zhenshan.cao 已提交
3183 3184 3185 3186
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3187 3188 3189 3190 3191
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3192
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
Z
zhenshan.cao 已提交
3193
		Base: &commonpb.MsgBase{
3194
			MsgType:   commonpb.MsgType_SegmentInfo,
Z
zhenshan.cao 已提交
3195 3196 3197 3198
			MsgID:     0,
			Timestamp: 0,
			SourceID:  Params.ProxyID,
		},
3199 3200
		CollectionID: collID,
		SegmentIDs:   segments,
Z
zhenshan.cao 已提交
3201 3202
	})
	if err != nil {
3203
		log.Error("Failed to get segment info from QueryCoord",
3204
			zap.Int64s("segmentIDs", segments), zap.Error(err))
Z
zhenshan.cao 已提交
3205 3206 3207
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3208
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3209
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3210
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3211 3212 3213 3214 3215 3216 3217 3218 3219 3220 3221 3222 3223
		resp.Status.Reason = infoResp.Status.Reason
		return resp, nil
	}
	queryInfos := make([]*milvuspb.QuerySegmentInfo, len(infoResp.Infos))
	for i, info := range infoResp.Infos {
		queryInfos[i] = &milvuspb.QuerySegmentInfo{
			SegmentID:    info.SegmentID,
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
			NumRows:      info.NumRows,
			MemSize:      info.MemSize,
			IndexName:    info.IndexName,
			IndexID:      info.IndexID,
3224
			NodeID:       info.NodeID,
X
xige-16 已提交
3225
			State:        info.SegmentState,
Z
zhenshan.cao 已提交
3226 3227
		}
	}
3228
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3229 3230 3231 3232
	resp.Infos = queryInfos
	return resp, nil
}

C
Cai Yudong 已提交
3233
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3234
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3235
		Base: &commonpb.MsgBase{
3236
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3237 3238 3239 3240 3241 3242 3243 3244 3245 3246
			MsgID:     0,
			Timestamp: 0,
			SourceID:  Params.ProxyID,
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3247
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3248 3249 3250
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3251
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3252
		Base: &commonpb.MsgBase{
3253
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3254 3255 3256 3257 3258 3259 3260 3261 3262 3263 3264
			MsgID:     0,
			Timestamp: 0,
			SourceID:  Params.ProxyID,
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3265
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3266 3267 3268 3269 3270
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3271
		showSegmentResponse, err := node.rootCoord.ShowSegments(ctx, &milvuspb.ShowSegmentsRequest{
X
XuanYang-cn 已提交
3272
			Base: &commonpb.MsgBase{
3273
				MsgType:   commonpb.MsgType_ShowSegments,
X
XuanYang-cn 已提交
3274 3275 3276 3277 3278 3279 3280 3281 3282 3283
				MsgID:     0,
				Timestamp: 0,
				SourceID:  Params.ProxyID,
			},
			CollectionID: collectionID,
			PartitionID:  partitionID,
		})
		if err != nil {
			return nil, err
		}
3284
		if showSegmentResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3285 3286 3287 3288 3289 3290
			return nil, errors.New(showSegmentResponse.Status.Reason)
		}
		ret = append(ret, showSegmentResponse.SegmentIDs...)
	}
	return ret, nil
}
3291

J
jingkl 已提交
3292
// Dummy handles dummy request
C
Cai Yudong 已提交
3293
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3294 3295 3296 3297 3298 3299 3300 3301 3302 3303 3304
	failedResponse := &milvuspb.DummyResponse{
		Response: `{"status": "fail"}`,
	}

	// TODO(wxyu): change name RequestType to Request
	drt, err := parseDummyRequestType(req.RequestType)
	if err != nil {
		log.Debug("Failed to parse dummy request type")
		return failedResponse, nil
	}

3305 3306
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3307
		if err != nil {
3308
			log.Debug("Failed to parse dummy query request")
3309 3310 3311
			return failedResponse, nil
		}

3312
		request := &milvuspb.QueryRequest{
3313 3314 3315
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3316
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3317 3318
		}

3319
		_, err = node.Query(ctx, request)
3320
		if err != nil {
3321
			log.Debug("Failed to execute dummy query")
3322 3323
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3324 3325 3326 3327 3328 3329

		return &milvuspb.DummyResponse{
			Response: `{"status": "success"}`,
		}, nil
	}

3330 3331
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3332 3333
}

J
jingkl 已提交
3334
// RegisterLink registers a link
C
Cai Yudong 已提交
3335
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
G
godchen 已提交
3336
	code := node.stateCode.Load().(internalpb.StateCode)
D
dragondriver 已提交
3337
	log.Debug("RegisterLink",
3338
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3339
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3340

G
godchen 已提交
3341
	if code != internalpb.StateCode_Healthy {
3342 3343 3344
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3345
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3346
				Reason:    "proxy not healthy",
3347 3348 3349 3350 3351 3352
			},
		}, nil
	}
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3353
			ErrorCode: commonpb.ErrorCode_Success,
3354
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3355 3356 3357
		},
	}, nil
}
3358

3359 3360 3361 3362 3363 3364 3365 3366 3367 3368 3369 3370 3371 3372 3373 3374 3375 3376 3377 3378 3379 3380 3381 3382 3383 3384 3385 3386 3387 3388 3389 3390 3391 3392 3393 3394 3395 3396 3397 3398
// TODO(dragondriver): cache the Metrics and set a retention to the cache
func (node *Proxy) GetMetrics(ctx context.Context, req *milvuspb.GetMetricsRequest) (*milvuspb.GetMetricsResponse, error) {
	log.Debug("Proxy.GetMetrics",
		zap.Int64("node_id", Params.ProxyID),
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
			zap.Int64("node_id", Params.ProxyID),
			zap.String("req", req.Request),
			zap.Error(errProxyIsUnhealthy(Params.ProxyID)))

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    msgProxyIsUnhealthy(Params.ProxyID),
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
			zap.Int64("node_id", Params.ProxyID),
			zap.String("req", req.Request),
			zap.Error(err))

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			Response: "",
		}, nil
	}

	log.Debug("Proxy.GetMetrics",
		zap.String("metric_type", metricType))

D
dragondriver 已提交
3399 3400 3401 3402 3403 3404 3405 3406 3407 3408 3409 3410 3411
	msgID := UniqueID(0)
	msgID, err = node.idAllocator.AllocOne()
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to allocate id",
			zap.Error(err))
	}
	req.Base = &commonpb.MsgBase{
		MsgType:   commonpb.MsgType_SystemInfo,
		MsgID:     msgID,
		Timestamp: 0,
		SourceID:  Params.ProxyID,
	}

3412
	if metricType == metricsinfo.SystemInfoMetrics {
3413 3414 3415 3416 3417 3418 3419
		ret, err := node.metricsCacheManager.GetSystemInfoMetrics()
		if err == nil && ret != nil {
			return ret, nil
		}
		log.Debug("failed to get system info metrics from cache, recompute instead",
			zap.Error(err))

3420
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3421 3422 3423 3424 3425 3426 3427 3428

		log.Debug("Proxy.GetMetrics",
			zap.Int64("node_id", Params.ProxyID),
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3429 3430
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3431
		return metrics, nil
3432 3433 3434 3435 3436 3437 3438 3439 3440 3441 3442 3443 3444 3445 3446 3447
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
		zap.Int64("node_id", Params.ProxyID),
		zap.String("req", req.Request),
		zap.String("metric_type", metricType))

	return &milvuspb.GetMetricsResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    metricsinfo.MsgUnimplementedMetric,
		},
		Response: "",
	}, nil
}

B
bigsheeper 已提交
3448 3449 3450 3451 3452 3453 3454 3455 3456 3457 3458 3459 3460 3461 3462 3463 3464 3465 3466 3467 3468 3469
// LoadBalance would do a load balancing operation between query nodes
func (node *Proxy) LoadBalance(ctx context.Context, req *milvuspb.LoadBalanceRequest) (*commonpb.Status, error) {
	log.Debug("Proxy.LoadBalance",
		zap.Int64("proxy_id", Params.ProxyID),
		zap.Any("req", req))

	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
	infoResp, err := node.queryCoord.LoadBalance(ctx, &querypb.LoadBalanceRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_LoadBalanceSegments,
			MsgID:     0,
			Timestamp: 0,
			SourceID:  Params.ProxyID,
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3470
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3471 3472 3473 3474 3475 3476 3477 3478 3479 3480 3481 3482 3483 3484 3485 3486 3487 3488
		SealedSegmentIDs: req.SealedSegmentIDs,
	})
	if err != nil {
		log.Error("Failed to LoadBalance from Query Coordinator",
			zap.Any("req", req), zap.Error(err))
		status.Reason = err.Error()
		return status, nil
	}
	if infoResp.ErrorCode != commonpb.ErrorCode_Success {
		log.Error("Failed to LoadBalance from Query Coordinator", zap.String("errMsg", infoResp.Reason))
		status.Reason = infoResp.Reason
		return status, nil
	}
	log.Debug("LoadBalance Done", zap.Any("req", req), zap.Any("status", infoResp))
	status.ErrorCode = commonpb.ErrorCode_Success
	return status, nil
}

J
jingkl 已提交
3489
//GetCompactionState gets the compaction state of multiple segments
3490 3491 3492 3493 3494 3495 3496 3497 3498 3499 3500 3501 3502 3503 3504 3505 3506 3507 3508 3509 3510 3511 3512 3513 3514 3515 3516 3517 3518 3519 3520 3521 3522 3523 3524 3525 3526 3527 3528
func (node *Proxy) GetCompactionState(ctx context.Context, req *milvuspb.GetCompactionStateRequest) (*milvuspb.GetCompactionStateResponse, error) {
	log.Info("received GetCompactionState request", zap.Int64("compactionID", req.GetCompactionID()))
	resp := &milvuspb.GetCompactionStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.dataCoord.GetCompactionState(ctx, req)
	log.Info("received GetCompactionState response", zap.Int64("compactionID", req.GetCompactionID()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

func (node *Proxy) ManualCompaction(ctx context.Context, req *milvuspb.ManualCompactionRequest) (*milvuspb.ManualCompactionResponse, error) {
	log.Info("received ManualCompaction request", zap.Int64("collectionID", req.GetCollectionID()))
	resp := &milvuspb.ManualCompactionResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.dataCoord.ManualCompaction(ctx, req)
	log.Info("received ManualCompaction response", zap.Int64("collectionID", req.GetCollectionID()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

func (node *Proxy) GetCompactionStateWithPlans(ctx context.Context, req *milvuspb.GetCompactionPlansRequest) (*milvuspb.GetCompactionPlansResponse, error) {
	log.Info("received GetCompactionStateWithPlans request", zap.Int64("compactionID", req.GetCompactionID()))
	resp := &milvuspb.GetCompactionPlansResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.dataCoord.GetCompactionStateWithPlans(ctx, req)
	log.Info("received GetCompactionStateWithPlans response", zap.Int64("compactionID", req.GetCompactionID()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

B
Bingyi Sun 已提交
3529 3530 3531
// GetFlushState gets the flush state of multiple segments
func (node *Proxy) GetFlushState(ctx context.Context, req *milvuspb.GetFlushStateRequest) (*milvuspb.GetFlushStateResponse, error) {
	log.Info("received get flush state request", zap.Any("request", req))
3532
	var err error
B
Bingyi Sun 已提交
3533 3534 3535 3536 3537 3538 3539
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3540
	resp, err = node.dataCoord.GetFlushState(ctx, req)
B
Bingyi Sun 已提交
3541 3542 3543 3544
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3545 3546
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3547 3548 3549 3550
	code := node.stateCode.Load().(internalpb.StateCode)
	return code == internalpb.StateCode_Healthy
}

J
jingkl 已提交
3551
//unhealthyStatus returns the proxy not healthy status
3552 3553 3554
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3555
		Reason:    "proxy not healthy",
3556 3557
	}
}