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
	method := "CreateCollection"
159 160 161
	// avoid data race
	lenOfSchema := len(request.Schema)

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

171 172 173
	if err := node.sched.ddQueue.Enqueue(cct); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
174 175
			zap.Error(err),
			zap.String("traceID", traceID),
176
			zap.String("role", typeutil.ProxyRole),
177 178 179 180 181
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Int("len(schema)", lenOfSchema),
			zap.Int32("shards_num", request.ShardsNum))

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

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

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

215
		return &commonpb.Status{
216
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
217 218 219 220
			Reason:    err.Error(),
		}, nil
	}

221 222
	log.Debug(
		rpcDone(method),
223
		zap.String("traceID", traceID),
224
		zap.String("role", typeutil.ProxyRole),
225 226 227 228 229 230 231 232
		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))

233 234 235
	return cct.result, nil
}

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

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

246
	dct := &dropCollectionTask{
S
sunby 已提交
247
		ctx:                   ctx,
248 249
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
250
		rootCoord:             node.rootCoord,
251
		chMgr:                 node.chMgr,
S
sunby 已提交
252
		chTicker:              node.chTicker,
253 254
	}

255 256
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
257
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
258 259
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
260 261 262 263 264

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

269
		return &commonpb.Status{
270
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
271 272 273 274
			Reason:    err.Error(),
		}, nil
	}

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

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

295
		return &commonpb.Status{
296
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
297 298 299 300
			Reason:    err.Error(),
		}, nil
	}

301 302
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
303
		zap.String("role", typeutil.ProxyRole),
304 305 306 307 308 309
		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))

310 311 312
	return dct.result, nil
}

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

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

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
327
		zap.String("role", typeutil.ProxyRole),
328 329 330
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

331
	hct := &hasCollectionTask{
S
sunby 已提交
332
		ctx:                  ctx,
333 334
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
335
		rootCoord:            node.rootCoord,
336 337
	}

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

346 347
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
348
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
349 350 351 352 353
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

374 375
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
376
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
377 378 379 380 381
				Reason:    err.Error(),
			},
		}, nil
	}

382 383
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
384
		zap.String("role", typeutil.ProxyRole),
385 386 387 388 389 390
		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))

391 392 393
	return hct.result, nil
}

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

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

404
	lct := &loadCollectionTask{
S
sunby 已提交
405
		ctx:                   ctx,
406 407
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
408
		queryCoord:            node.queryCoord,
409 410
	}

411 412
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
413
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
414 415
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
416 417 418 419 420

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

425
		return &commonpb.Status{
426
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
427 428 429
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
430

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

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

451
		return &commonpb.Status{
452
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
453 454 455 456
			Reason:    err.Error(),
		}, nil
	}

457 458
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
459
		zap.String("role", typeutil.ProxyRole),
460 461 462 463 464 465
		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))

466
	return lct.result, nil
467 468
}

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

475
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
476 477 478
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

479
	rct := &releaseCollectionTask{
S
sunby 已提交
480
		ctx:                      ctx,
481 482
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
483
		queryCoord:               node.queryCoord,
484
		chMgr:                    node.chMgr,
485 486
	}

487 488 489 490
	method := "ReleaseCollection"

	log.Debug(
		rpcReceived(method),
491
		zap.String("traceID", traceID),
492
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
493 494
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
495 496

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

505
		return &commonpb.Status{
506
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
507 508 509 510
			Reason:    err.Error(),
		}, nil
	}

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

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

533
		return &commonpb.Status{
534
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
535 536 537 538
			Reason:    err.Error(),
		}, nil
	}

539 540
	log.Debug(
		rpcDone(method),
541
		zap.String("traceID", traceID),
542
		zap.String("role", typeutil.ProxyRole),
543 544 545 546 547 548
		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))

549
	return rct.result, nil
550 551
}

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

560
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
561 562 563
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

564
	dct := &describeCollectionTask{
S
sunby 已提交
565
		ctx:                       ctx,
566 567
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
568
		rootCoord:                 node.rootCoord,
569 570
	}

571 572
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
573
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
574 575
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
576 577 578 579 580

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

585 586
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
587
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
588 589 590 591 592
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

613 614
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
615
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
616 617 618 619 620
				Reason:    err.Error(),
			},
		}, nil
	}

621 622
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
623
		zap.String("role", typeutil.ProxyRole),
624 625 626 627 628 629
		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))

630 631 632
	return dct.result, nil
}

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

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

645
	g := &getCollectionStatisticsTask{
G
godchen 已提交
646 647 648
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
649
		dataCoord:                      node.dataCoord,
650 651
	}

652 653
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
654
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
655 656
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
657 658 659 660 661

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

G
godchen 已提交
666
		return &milvuspb.GetCollectionStatisticsResponse{
667
			Status: &commonpb.Status{
668
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
669 670 671 672 673
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

G
godchen 已提交
694
		return &milvuspb.GetCollectionStatisticsResponse{
695
			Status: &commonpb.Status{
696
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
697 698 699 700 701
				Reason:    err.Error(),
			},
		}, nil
	}

702 703
	log.Debug("GetCollectionStatistics done",
		zap.String("traceID", traceID),
704
		zap.String("role", typeutil.ProxyRole),
705 706 707 708 709 710
		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))

711
	return g.result, nil
712 713
}

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

729
	log.Debug("ShowCollections received",
730
		zap.String("role", typeutil.ProxyRole),
731 732 733 734 735 736
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

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

G
godchen 已提交
748
		return &milvuspb.ShowCollectionsResponse{
749
			Status: &commonpb.Status{
750
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
751 752 753 754 755
				Reason:    err.Error(),
			},
		}, nil
	}

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

765 766
	err = sct.WaitToFinish()
	if err != nil {
767 768
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
769
			zap.String("role", typeutil.ProxyRole),
770 771 772 773 774 775 776
			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 已提交
777
		return &milvuspb.ShowCollectionsResponse{
778
			Status: &commonpb.Status{
779
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
780 781 782 783 784
				Reason:    err.Error(),
			},
		}, nil
	}

785
	log.Debug("ShowCollections Done",
786
		zap.String("role", typeutil.ProxyRole),
787 788 789 790 791 792 793 794
		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),
	)

795 796 797
	return sct.result, nil
}

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

804
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
805 806 807
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

808
	cpt := &createPartitionTask{
S
sunby 已提交
809
		ctx:                    ctx,
810 811
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
812
		rootCoord:              node.rootCoord,
813 814 815
		result:                 nil,
	}

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

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

834
		return &commonpb.Status{
835
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
836 837 838
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
839

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

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

864
		return &commonpb.Status{
865
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
866 867 868
			Reason:    err.Error(),
		}, nil
	}
869 870 871 872

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
873
		zap.String("role", typeutil.ProxyRole),
874 875 876 877 878 879 880
		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))

881 882 883
	return cpt.result, nil
}

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

890
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
891 892 893
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

894
	dpt := &dropPartitionTask{
S
sunby 已提交
895
		ctx:                  ctx,
896 897
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
898
		rootCoord:            node.rootCoord,
899 900 901
		result:               nil,
	}

902 903 904 905 906
	method := "DropPartition"

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

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

922
		return &commonpb.Status{
923
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
924 925 926
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
927

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

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

952
		return &commonpb.Status{
953
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
954 955 956
			Reason:    err.Error(),
		}, nil
	}
957 958 959 960

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
961
		zap.String("role", typeutil.ProxyRole),
962 963 964 965 966 967 968
		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))

969 970 971
	return dpt.result, nil
}

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

D
dragondriver 已提交
980
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
981 982 983
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

984
	hpt := &hasPartitionTask{
S
sunby 已提交
985
		ctx:                 ctx,
986 987
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
988
		rootCoord:           node.rootCoord,
989 990 991
		result:              nil,
	}

D
dragondriver 已提交
992 993 994 995 996
	method := "HasPartition"

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

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

1012 1013
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1014
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1015 1016 1017 1018 1019
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1020

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

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

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

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1057
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1058 1059 1060 1061 1062 1063 1064
		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))

1065 1066 1067
	return hpt.result, nil
}

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

D
dragondriver 已提交
1074
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1075 1076 1077
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1078
	lpt := &loadPartitionsTask{
G
godchen 已提交
1079 1080 1081
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1082
		queryCoord:            node.queryCoord,
1083 1084
	}

1085 1086 1087 1088 1089
	method := "LoadPartitions"

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

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

1105
		return &commonpb.Status{
1106
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1107 1108 1109 1110
			Reason:    err.Error(),
		}, nil
	}

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

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

1135
		return &commonpb.Status{
1136
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1137 1138 1139 1140
			Reason:    err.Error(),
		}, nil
	}

1141 1142 1143
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1144
		zap.String("role", typeutil.ProxyRole),
1145 1146 1147 1148 1149 1150 1151
		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))

1152
	return lpt.result, nil
1153 1154
}

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

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

1165
	rpt := &releasePartitionsTask{
G
godchen 已提交
1166 1167 1168
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1169
		queryCoord:               node.queryCoord,
1170 1171
	}

1172 1173 1174 1175 1176
	method := "ReleasePartitions"

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

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

1192
		return &commonpb.Status{
1193
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1194 1195 1196 1197
			Reason:    err.Error(),
		}, nil
	}

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

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

1222
		return &commonpb.Status{
1223
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1224 1225 1226 1227
			Reason:    err.Error(),
		}, nil
	}

1228 1229 1230
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1231
		zap.String("role", typeutil.ProxyRole),
1232 1233 1234 1235 1236 1237 1238
		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))

1239
	return rpt.result, nil
1240 1241
}

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

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

1254
	g := &getPartitionStatisticsTask{
1255 1256 1257
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1258
		dataCoord:                     node.dataCoord,
1259 1260
	}

1261 1262 1263 1264 1265
	method := "GetPartitionStatistics"

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

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

1281 1282 1283 1284 1285 1286 1287 1288
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

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

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1303
			zap.Error(err),
1304
			zap.String("traceID", traceID),
1305
			zap.String("role", typeutil.ProxyRole),
1306 1307 1308
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320
			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
	}

1321 1322 1323
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1324
		zap.String("role", typeutil.ProxyRole),
1325 1326 1327 1328 1329 1330 1331
		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))

1332
	return g.result, nil
1333 1334
}

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

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

1347
	spt := &showPartitionsTask{
G
godchen 已提交
1348 1349 1350
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1351
		rootCoord:             node.rootCoord,
1352
		queryCoord:            node.queryCoord,
G
godchen 已提交
1353
		result:                nil,
1354 1355
	}

1356 1357 1358 1359 1360
	method := "ShowPartitions"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1361
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1362
		zap.Any("request", request))
1363 1364 1365 1366 1367 1368

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

G
godchen 已提交
1372
		return &milvuspb.ShowPartitionsResponse{
1373
			Status: &commonpb.Status{
1374
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1375 1376 1377 1378 1379
				Reason:    err.Error(),
			},
		}, nil
	}

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

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1394
			zap.Error(err),
1395
			zap.String("traceID", traceID),
1396
			zap.String("role", typeutil.ProxyRole),
1397 1398 1399 1400 1401 1402
			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 已提交
1403

G
godchen 已提交
1404
		return &milvuspb.ShowPartitionsResponse{
1405
			Status: &commonpb.Status{
1406
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1407 1408 1409 1410
				Reason:    err.Error(),
			},
		}, nil
	}
1411 1412 1413 1414

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1415
		zap.String("role", typeutil.ProxyRole),
1416 1417 1418 1419 1420 1421 1422
		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))

1423 1424 1425
	return spt.result, nil
}

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

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

1436
	cit := &createIndexTask{
S
sunby 已提交
1437
		ctx:                ctx,
1438 1439
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1440
		rootCoord:          node.rootCoord,
1441 1442
	}

D
dragondriver 已提交
1443 1444 1445 1446 1447
	method := "CreateIndex"

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

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

1465
		return &commonpb.Status{
1466
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1467 1468 1469 1470
			Reason:    err.Error(),
		}, nil
	}

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

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

1497
		return &commonpb.Status{
1498
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1499 1500 1501 1502
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1503 1504 1505
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1506
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1507 1508 1509 1510 1511 1512 1513 1514
		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))

1515 1516 1517
	return cit.result, nil
}

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

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

1530
	dit := &describeIndexTask{
S
sunby 已提交
1531
		ctx:                  ctx,
1532 1533
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1534
		rootCoord:            node.rootCoord,
1535 1536
	}

1537 1538 1539 1540 1541 1542 1543
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1544
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1545 1546 1547
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1548 1549 1550 1551 1552 1553 1554
		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),
1555
			zap.String("role", typeutil.ProxyRole),
1556 1557 1558 1559 1560
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

1561 1562
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1563
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1564 1565 1566 1567 1568
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

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

1607 1608 1609
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1610
		zap.String("role", typeutil.ProxyRole),
1611 1612 1613 1614 1615 1616 1617 1618
		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))

1619 1620 1621
	return dit.result, nil
}

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

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

1632
	dit := &dropIndexTask{
S
sunby 已提交
1633
		ctx:              ctx,
B
BossZou 已提交
1634 1635
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1636
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1637
	}
G
godchen 已提交
1638

D
dragondriver 已提交
1639 1640 1641 1642 1643
	method := "DropIndex"

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

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

B
BossZou 已提交
1661
		return &commonpb.Status{
1662
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1663 1664 1665
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1666

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

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

B
BossZou 已提交
1693
		return &commonpb.Status{
1694
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1695 1696 1697
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1698 1699 1700 1701

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1702
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1703 1704 1705 1706 1707 1708 1709 1710
		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 已提交
1711 1712 1713
	return dit.result, nil
}

1714 1715
// 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 已提交
1716
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1717 1718 1719 1720 1721
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1722 1723 1724 1725 1726

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

1727
	gibpt := &getIndexBuildProgressTask{
1728 1729 1730
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1731 1732
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1733
		dataCoord:                    node.dataCoord,
1734 1735
	}

1736 1737 1738 1739 1740
	method := "GetIndexBuildProgress"

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

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

1758 1759 1760 1761 1762 1763 1764 1765
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

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

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1781
			zap.Error(err),
1782
			zap.String("traceID", traceID),
1783
			zap.String("role", typeutil.ProxyRole),
1784 1785 1786
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798
			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
	}
1799 1800 1801 1802

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1803
		zap.String("role", typeutil.ProxyRole),
1804 1805 1806 1807 1808 1809 1810 1811
		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))
1812 1813

	return gibpt.result, nil
1814 1815
}

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

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

1828
	dipt := &getIndexStateTask{
G
godchen 已提交
1829 1830 1831
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
1832 1833
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
1834 1835
	}

1836 1837 1838 1839 1840
	method := "GetIndexState"

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

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

G
godchen 已提交
1858
		return &milvuspb.GetIndexStateResponse{
1859
			Status: &commonpb.Status{
1860
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1861 1862 1863 1864 1865
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

G
godchen 已提交
1892
		return &milvuspb.GetIndexStateResponse{
1893
			Status: &commonpb.Status{
1894
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1895 1896 1897 1898 1899
				Reason:    err.Error(),
			},
		}, nil
	}

1900 1901 1902
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1903
		zap.String("role", typeutil.ProxyRole),
1904 1905 1906 1907 1908 1909 1910 1911
		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))

1912 1913 1914
	return dipt.result, nil
}

1915
// Insert insert records into collection.
C
Cai Yudong 已提交
1916
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
1917 1918 1919 1920 1921 1922
	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))

1923 1924 1925 1926 1927
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1928

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

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

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

X
Xiangyu Wang 已提交
1972
	log.Debug("Enqueue insert request in Proxy",
1973
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1974 1975 1976 1977 1978
		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)),
1979 1980
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
1981

X
Xiangyu Wang 已提交
1982 1983 1984
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
		return constructFailedResponse(err), nil
1985
	}
D
dragondriver 已提交
1986

X
Xiangyu Wang 已提交
1987
	log.Debug("Detail of insert request in Proxy",
1988
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
1989
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
1990 1991 1992 1993 1994
		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),
X
Xiangyu Wang 已提交
1995 1996 1997 1998 1999 2000 2001 2002 2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016 2017
		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 已提交
2018

2019 2020 2021
	return it.result, nil
}

2022
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2023
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2024 2025 2026
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2027 2028
	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))
2029

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

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

2043
	dt := &deleteTask{
C
Cai Yudong 已提交
2044 2045 2046
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       deleteReq,
G
godchen 已提交
2047
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2048 2049 2050
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2051 2052 2053 2054 2055 2056 2057 2058
			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 已提交
2059 2060 2061 2062
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2063 2064
	}

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

	// 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 已提交
2075 2076 2077 2078 2079 2080 2081 2082
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2083
	log.Debug("Detail of delete request in Proxy",
2084
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2085 2086 2087 2088 2089
		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),
2090 2091
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2092

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

	return dt.result, nil
}

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

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

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

D
dragondriver 已提交
2134 2135
	log.Debug("Search received",
		zap.String("traceID", traceID),
2136
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2137 2138 2139 2140 2141 2142
		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 已提交
2143

2144
	err := node.sched.dqQueue.Enqueue(qt)
2145
	if err != nil {
D
dragondriver 已提交
2146 2147 2148
		log.Debug("Search failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
2149
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2150 2151 2152 2153 2154 2155 2156 2157
			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),
		)

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

D
dragondriver 已提交
2166 2167
	log.Debug("Search enqueued",
		zap.String("traceID", traceID),
2168
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2169
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2170 2171 2172 2173 2174
		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),
2175 2176
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
		zap.Any("OutputFields", request.OutputFields))
D
dragondriver 已提交
2177 2178 2179 2180 2181

	err = qt.WaitToFinish()

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

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

D
dragondriver 已提交
2201 2202
	log.Debug("Search Done",
		zap.String("traceID", traceID),
2203
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2204 2205 2206 2207 2208 2209 2210 2211
		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))

2212 2213 2214
	return qt.result, nil
}

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

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

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

D
dragondriver 已提交
2239 2240 2241 2242 2243
	method := "Flush"

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

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

2257 2258
		resp.Status.Reason = err.Error()
		return resp, nil
2259 2260
	}

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

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

D
dragondriver 已提交
2283
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2284 2285
		resp.Status.Reason = err.Error()
		return resp, nil
2286 2287
	}

D
dragondriver 已提交
2288 2289 2290
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2291
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2292 2293 2294 2295 2296 2297
		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))

2298
	return ft.result, nil
2299 2300
}

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

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

2313 2314 2315 2316 2317 2318
	queryRequest := &milvuspb.QueryRequest{
		DbName:         request.DbName,
		CollectionName: request.CollectionName,
		PartitionNames: request.PartitionNames,
		Expr:           request.Expr,
		OutputFields:   request.OutputFields,
2319
	}
2320

2321
	qt := &queryTask{
2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334
		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,
2335 2336
	}

D
dragondriver 已提交
2337 2338 2339 2340 2341
	method := "Query"

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

D
dragondriver 已提交
2347 2348 2349 2350 2351 2352 2353 2354 2355 2356
	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))

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

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

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

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

D
dragondriver 已提交
2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407
	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))

2408 2409 2410 2411 2412
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2413

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

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

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

D
dragondriver 已提交
2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 2442 2443 2444 2445 2446 2447 2448 2449 2450
	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 已提交
2451 2452 2453 2454 2455 2456
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

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

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

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

D
dragondriver 已提交
2487 2488 2489 2490 2491 2492 2493 2494 2495 2496 2497
	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 已提交
2498 2499 2500
	return cat.result, nil
}

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

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

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

D
dragondriver 已提交
2518 2519 2520 2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535
	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 已提交
2536 2537 2538 2539 2540 2541
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

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

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

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

D
dragondriver 已提交
2570 2571 2572 2573 2574 2575 2576 2577 2578 2579
	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 已提交
2580 2581 2582
	return dat.result, nil
}

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

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

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

D
dragondriver 已提交
2600 2601 2602 2603 2604 2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616 2617 2618 2619
	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 已提交
2620 2621 2622 2623 2624 2625
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

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

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

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

D
dragondriver 已提交
2656 2657 2658 2659 2660 2661 2662 2663 2664 2665 2666
	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 已提交
2667 2668 2669
	return aat.result, nil
}

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

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

2689 2690 2691 2692
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

2693 2694
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
2695

2696 2697 2698 2699 2700
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
2701 2702
		}

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

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

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

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
2740
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2741 2742 2743 2744 2745 2746
			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))
2747 2748 2749 2750

		err = qt.WaitToFinish()
		if err != nil {
			log.Debug("CalcDistance queryTask failed to WaitToFinish",
G
godchen 已提交
2751
				zap.Error(err),
2752
				zap.String("traceID", traceID),
2753
				zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2754 2755 2756 2757 2758 2759
				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))
2760 2761 2762 2763 2764 2765

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
2766
			}, err
2767
		}
2768 2769 2770

		log.Debug("CalcDistance queryTask Done",
			zap.String("traceID", traceID),
2771
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2772 2773 2774 2775 2776 2777
			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))
2778 2779

		return &milvuspb.QueryResults{
2780 2781
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
2782 2783 2784
		}, nil
	}

2785 2786 2787 2788 2789 2790 2791 2792 2793 2794 2795 2796 2797 2798
	// 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 {
2799
			return nil, errors.New("failed to fetch vectors")
2800 2801 2802 2803 2804 2805 2806 2807 2808 2809 2810 2811 2812 2813 2814 2815
		}

		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))
2816
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
2817 2818 2819 2820 2821 2822 2823 2824 2825 2826 2827 2828 2829 2830 2831 2832 2833 2834 2835 2836 2837 2838 2839 2840 2841 2842
				}
				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))
2843
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
2844 2845 2846 2847 2848 2849 2850 2851 2852 2853 2854 2855
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

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

2856
		return nil, errors.New("failed to fetch vectors")
2857 2858
	}

2859 2860
	log.Debug("CalcDistance received",
		zap.String("traceID", traceID),
2861
		zap.String("role", typeutil.ProxyRole),
2862
		zap.String("metric", metric))
G
godchen 已提交
2863

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

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

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

2886 2887
		log.Debug("OpLeft IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
2888
			zap.String("role", typeutil.ProxyRole))
2889

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

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

		log.Debug("Re-arrange left vectors done",
			zap.String("traceID", traceID),
2907
			zap.String("role", typeutil.ProxyRole))
2908 2909
	}

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

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

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

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

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

2946 2947
		log.Debug("OpRight IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
2948
			zap.String("role", typeutil.ProxyRole))
2949

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

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

		log.Debug("Re-arrange right vectors done",
			zap.String("traceID", traceID),
2967
			zap.String("role", typeutil.ProxyRole))
2968 2969
	}

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

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

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

2990 2991 2992
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2993
				Reason:    msg,
2994 2995 2996 2997 2998
			},
		}, nil
	}

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

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

3014 3015 3016
		log.Debug("CalcFloatDistance done",
			zap.Error(err),
			zap.String("traceID", traceID),
3017
			zap.String("role", typeutil.ProxyRole))
3018

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

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

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

		if metric == distance.HAMMING {
3046 3047
			log.Debug("CalcHammingDistance done",
				zap.String("traceID", traceID),
3048
				zap.String("role", typeutil.ProxyRole))
3049

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

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

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

3076 3077
			log.Debug("CalcTanimotoCoefficient done",
				zap.String("traceID", traceID),
3078
				zap.String("role", typeutil.ProxyRole))
3079

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

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

3096 3097 3098
	log.Debug("Failed to CalcDistance",
		zap.Error(err),
		zap.String("traceID", traceID),
3099
		zap.String("role", typeutil.ProxyRole))
3100

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

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

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

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

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

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

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

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

J
jingkl 已提交
3295
// Dummy handles dummy request
C
Cai Yudong 已提交
3296
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3297 3298 3299 3300 3301 3302 3303 3304 3305 3306 3307
	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
	}

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

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

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

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

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

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

G
godchen 已提交
3344
	if code != internalpb.StateCode_Healthy {
3345 3346 3347
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3348
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3349
				Reason:    "proxy not healthy",
3350 3351 3352 3353 3354 3355
			},
		}, nil
	}
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3356
			ErrorCode: commonpb.ErrorCode_Success,
3357
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3358 3359 3360
		},
	}, nil
}
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 3399 3400 3401
// 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 已提交
3402 3403 3404 3405 3406 3407 3408 3409 3410 3411 3412 3413 3414
	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,
	}

3415
	if metricType == metricsinfo.SystemInfoMetrics {
3416 3417 3418 3419 3420 3421 3422
		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))

3423
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3424 3425 3426 3427 3428 3429 3430 3431

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

3432 3433
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

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

	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 已提交
3451 3452 3453 3454 3455 3456 3457 3458 3459 3460 3461 3462 3463 3464 3465 3466 3467 3468 3469 3470 3471 3472
// 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 已提交
3473
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3474 3475 3476 3477 3478 3479 3480 3481 3482 3483 3484 3485 3486 3487 3488 3489 3490 3491
		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 已提交
3492
//GetCompactionState gets the compaction state of multiple segments
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 3529 3530 3531
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 已提交
3532 3533 3534
// 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))
3535
	var err error
B
Bingyi Sun 已提交
3536 3537 3538 3539 3540 3541 3542
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

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

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

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