impl.go 133.3 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
	"github.com/milvus-io/milvus/internal/util/timerecord"
27

28
	"go.uber.org/zap"
S
sunby 已提交
29

30
	"github.com/milvus-io/milvus/internal/common"
X
Xiangyu Wang 已提交
31
	"github.com/milvus-io/milvus/internal/log"
32
	"github.com/milvus-io/milvus/internal/metrics"
J
jaime 已提交
33
	"github.com/milvus-io/milvus/internal/mq/msgstream"
X
Xiangyu Wang 已提交
34 35 36 37 38 39
	"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"
40
	"github.com/milvus-io/milvus/internal/proto/schemapb"
41
	"github.com/milvus-io/milvus/internal/util/distance"
42 43 44 45
	"github.com/milvus-io/milvus/internal/util/funcutil"
	"github.com/milvus-io/milvus/internal/util/logutil"
	"github.com/milvus-io/milvus/internal/util/metricsinfo"
	"github.com/milvus-io/milvus/internal/util/trace"
X
Xiangyu Wang 已提交
46
	"github.com/milvus-io/milvus/internal/util/typeutil"
47 48
)

49 50
const moduleName = "Proxy"

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

56
// GetComponentStates get state of Proxy.
C
Cai Yudong 已提交
57
func (node *Proxy) GetComponentStates(ctx context.Context) (*internalpb.ComponentStates, error) {
G
godchen 已提交
58 59 60 61 62 63 64 65 66 67 68 69
	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 已提交
70
		return stats, nil
G
godchen 已提交
71
	}
72 73 74 75
	nodeID := common.NotRegisteredID
	if node.session != nil && node.session.Registered() {
		nodeID = node.session.ServerID
	}
G
godchen 已提交
76
	info := &internalpb.ComponentInfo{
77 78
		// NodeID:    Params.ProxyID, // will race with Proxy.Register()
		NodeID:    nodeID,
C
Cai Yudong 已提交
79
		Role:      typeutil.ProxyRole,
G
godchen 已提交
80 81 82 83 84 85
		StateCode: code,
	}
	stats.State = info
	return stats, nil
}

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

97
// InvalidateCollectionMetaCache invalidate the meta cache of specific collection.
C
Cai Yudong 已提交
98
func (node *Proxy) InvalidateCollectionMetaCache(ctx context.Context, request *proxypb.InvalidateCollMetaCacheRequest) (*commonpb.Status, error) {
99 100
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to invalidate collection meta cache",
101
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
102 103 104
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

105
	collectionName := request.CollectionName
N
neza2017 已提交
106 107 108
	if globalMetaCache != nil {
		globalMetaCache.RemoveCollection(ctx, collectionName) // no need to return error, though collection may be not cached
	}
109
	logutil.Logger(ctx).Debug("complete to invalidate collection meta cache",
110
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
111 112 113
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

114
	return &commonpb.Status{
115
		ErrorCode: commonpb.ErrorCode_Success,
116 117
		Reason:    "",
	}, nil
118 119
}

120
// ReleaseDQLMessageStream release the query message stream of specific collection.
C
Cai Yudong 已提交
121
func (node *Proxy) ReleaseDQLMessageStream(ctx context.Context, request *proxypb.ReleaseDQLMessageStreamRequest) (*commonpb.Status, error) {
122 123
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to release DQL message strem",
124
		zap.Any("role", typeutil.ProxyRole),
125 126 127
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

128 129 130 131
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

132 133
	_ = node.chMgr.removeDQLStream(request.CollectionID)

134
	logutil.Logger(ctx).Debug("complete to release DQL message stream",
135
		zap.Any("role", typeutil.ProxyRole),
136 137 138 139 140 141 142 143 144
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

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

145
// TODO(dragondriver): add more detailed ut for ConsistencyLevel, should we support multiple consistency level in Proxy?
146
// CreateCollection create a collection by the schema.
C
Cai Yudong 已提交
147
func (node *Proxy) CreateCollection(ctx context.Context, request *milvuspb.CreateCollectionRequest) (*commonpb.Status, error) {
148 149 150
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
151 152 153 154

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
155 156 157 158
	method := "CreateCollection"
	tr := timerecord.NewTimeRecorder(method)

	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
159

160
	cct := &createCollectionTask{
S
sunby 已提交
161
		ctx:                     ctx,
162 163
		Condition:               NewTaskCondition(ctx),
		CreateCollectionRequest: request,
164
		rootCoord:               node.rootCoord,
165 166
	}

167 168 169
	// avoid data race
	lenOfSchema := len(request.Schema)

170 171
	log.Debug(
		rpcReceived(method),
172
		zap.String("traceID", traceID),
173
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
174 175
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
176
		zap.Int("len(schema)", lenOfSchema),
177 178
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
179

180 181 182
	if err := node.sched.ddQueue.Enqueue(cct); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
183 184
			zap.Error(err),
			zap.String("traceID", traceID),
185
			zap.String("role", typeutil.ProxyRole),
186 187 188
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Int("len(schema)", lenOfSchema),
189 190
			zap.Int32("shards_num", request.ShardsNum),
			zap.String("consistency_level", request.ConsistencyLevel.String()))
191

192
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
193
		return &commonpb.Status{
194
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
195 196 197 198
			Reason:    err.Error(),
		}, nil
	}

199 200
	log.Debug(
		rpcEnqueued(method),
201
		zap.String("traceID", traceID),
202
		zap.String("role", typeutil.ProxyRole),
203 204 205
		zap.Int64("MsgID", cct.ID()),
		zap.Uint64("BeginTs", cct.BeginTs()),
		zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
206 207
		zap.Uint64("timestamp", request.Base.Timestamp),
		zap.String("db", request.DbName),
208 209
		zap.String("collection", request.CollectionName),
		zap.Int("len(schema)", lenOfSchema),
210 211
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
212

213 214 215
	if err := cct.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
216
			zap.Error(err),
217
			zap.String("traceID", traceID),
218
			zap.String("role", typeutil.ProxyRole),
219 220 221
			zap.Int64("MsgID", cct.ID()),
			zap.Uint64("BeginTs", cct.BeginTs()),
			zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
222 223
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
224
			zap.Int("len(schema)", lenOfSchema),
225 226
			zap.Int32("shards_num", request.ShardsNum),
			zap.String("consistency_level", request.ConsistencyLevel.String()))
D
dragondriver 已提交
227

228
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
229
		return &commonpb.Status{
230
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
231 232 233 234
			Reason:    err.Error(),
		}, nil
	}

235 236
	log.Debug(
		rpcDone(method),
237
		zap.String("traceID", traceID),
238
		zap.String("role", typeutil.ProxyRole),
239 240 241 242 243 244
		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),
245 246
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
247

248 249
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
250 251 252
	return cct.result, nil
}

253
// DropCollection drop a collection.
C
Cai Yudong 已提交
254
func (node *Proxy) DropCollection(ctx context.Context, request *milvuspb.DropCollectionRequest) (*commonpb.Status, error) {
255 256 257
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
258 259 260 261

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
262 263 264
	method := "DropCollection"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
265

266
	dct := &dropCollectionTask{
S
sunby 已提交
267
		ctx:                   ctx,
268 269
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
270
		rootCoord:             node.rootCoord,
271
		chMgr:                 node.chMgr,
S
sunby 已提交
272
		chTicker:              node.chTicker,
273 274
	}

275 276
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
277
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
278 279
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
280 281 282 283 284

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

289
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
290
		return &commonpb.Status{
291
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
292 293 294 295
			Reason:    err.Error(),
		}, nil
	}

296 297
	log.Debug("DropCollection enqueued",
		zap.String("traceID", traceID),
298
		zap.String("role", typeutil.ProxyRole),
299 300 301
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
302 303
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
304 305 306

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DropCollection failed to WaitToFinish",
D
dragondriver 已提交
307
			zap.Error(err),
308
			zap.String("traceID", traceID),
309
			zap.String("role", typeutil.ProxyRole),
310 311 312
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTs", dct.BeginTs()),
			zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
313 314 315
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

316
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
317
		return &commonpb.Status{
318
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
319 320 321 322
			Reason:    err.Error(),
		}, nil
	}

323 324
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
325
		zap.String("role", typeutil.ProxyRole),
326 327 328 329 330 331
		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))

332 333
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
334 335 336
	return dct.result, nil
}

337
// HasCollection check if the specific collection exists in Milvus.
C
Cai Yudong 已提交
338
func (node *Proxy) HasCollection(ctx context.Context, request *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
339 340 341 342 343
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
344 345 346 347

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
348 349 350
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
351
		metrics.TotalLabel).Inc()
352 353 354

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
355
		zap.String("role", typeutil.ProxyRole),
356 357 358
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

359
	hct := &hasCollectionTask{
S
sunby 已提交
360
		ctx:                  ctx,
361 362
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
363
		rootCoord:            node.rootCoord,
364 365
	}

366 367 368 369
	if err := node.sched.ddQueue.Enqueue(hct); err != nil {
		log.Warn("HasCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
370
			zap.String("role", typeutil.ProxyRole),
371 372 373
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

374
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
375
			metrics.AbandonLabel).Inc()
376 377
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
378
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
379 380 381 382 383
				Reason:    err.Error(),
			},
		}, nil
	}

384 385
	log.Debug("HasCollection enqueued",
		zap.String("traceID", traceID),
386
		zap.String("role", typeutil.ProxyRole),
387 388 389
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
390 391
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
392 393 394

	if err := hct.WaitToFinish(); err != nil {
		log.Warn("HasCollection failed to WaitToFinish",
D
dragondriver 已提交
395
			zap.Error(err),
396
			zap.String("traceID", traceID),
397
			zap.String("role", typeutil.ProxyRole),
398 399 400
			zap.Int64("MsgID", hct.ID()),
			zap.Uint64("BeginTS", hct.BeginTs()),
			zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
401 402 403
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

404
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
405
			metrics.FailLabel).Inc()
406 407
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
408
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
409 410 411 412 413
				Reason:    err.Error(),
			},
		}, nil
	}

414 415
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
416
		zap.String("role", typeutil.ProxyRole),
417 418 419 420 421 422
		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))

423
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
424 425
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
426 427 428
	return hct.result, nil
}

429
// LoadCollection load a collection into query nodes.
C
Cai Yudong 已提交
430
func (node *Proxy) LoadCollection(ctx context.Context, request *milvuspb.LoadCollectionRequest) (*commonpb.Status, error) {
431 432 433
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
434 435 436 437

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
438 439
	method := "LoadCollection"
	tr := timerecord.NewTimeRecorder(method)
440

441
	lct := &loadCollectionTask{
S
sunby 已提交
442
		ctx:                   ctx,
443 444
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
445
		queryCoord:            node.queryCoord,
446 447
	}

448 449
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
450
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
451 452
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
453 454 455 456 457

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

462
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
463
			metrics.AbandonLabel).Inc()
464
		return &commonpb.Status{
465
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
466 467 468
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
469

470 471
	log.Debug("LoadCollection enqueued",
		zap.String("traceID", traceID),
472
		zap.String("role", typeutil.ProxyRole),
473 474 475
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
476 477
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
478 479 480

	if err := lct.WaitToFinish(); err != nil {
		log.Warn("LoadCollection failed to WaitToFinish",
D
dragondriver 已提交
481
			zap.Error(err),
482
			zap.String("traceID", traceID),
483
			zap.String("role", typeutil.ProxyRole),
484 485 486
			zap.Int64("MsgID", lct.ID()),
			zap.Uint64("BeginTS", lct.BeginTs()),
			zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
487 488 489
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

490
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
491
			metrics.TotalLabel).Inc()
492
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
493
			metrics.FailLabel).Inc()
494
		return &commonpb.Status{
495
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
496 497 498 499
			Reason:    err.Error(),
		}, nil
	}

500 501
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
502
		zap.String("role", typeutil.ProxyRole),
503 504 505 506 507 508
		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))

509
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
510
		metrics.TotalLabel).Inc()
511
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
512 513
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
514
	return lct.result, nil
515 516
}

517
// ReleaseCollection remove the loaded collection from query nodes.
C
Cai Yudong 已提交
518
func (node *Proxy) ReleaseCollection(ctx context.Context, request *milvuspb.ReleaseCollectionRequest) (*commonpb.Status, error) {
519 520 521
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
522

523
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
524 525
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
526 527
	method := "ReleaseCollection"
	tr := timerecord.NewTimeRecorder(method)
528

529
	rct := &releaseCollectionTask{
S
sunby 已提交
530
		ctx:                      ctx,
531 532
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
533
		queryCoord:               node.queryCoord,
534
		chMgr:                    node.chMgr,
535 536
	}

537 538
	log.Debug(
		rpcReceived(method),
539
		zap.String("traceID", traceID),
540
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
541 542
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
543 544

	if err := node.sched.ddQueue.Enqueue(rct); err != nil {
545 546
		log.Warn(
			rpcFailedToEnqueue(method),
547 548
			zap.Error(err),
			zap.String("traceID", traceID),
549
			zap.String("role", typeutil.ProxyRole),
550 551 552
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

553
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
554
			metrics.AbandonLabel).Inc()
555
		return &commonpb.Status{
556
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
557 558 559 560
			Reason:    err.Error(),
		}, nil
	}

561 562
	log.Debug(
		rpcEnqueued(method),
563
		zap.String("traceID", traceID),
564
		zap.String("role", typeutil.ProxyRole),
565 566 567
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
568 569
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
570 571

	if err := rct.WaitToFinish(); err != nil {
572 573
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
574
			zap.Error(err),
575
			zap.String("traceID", traceID),
576
			zap.String("role", typeutil.ProxyRole),
577 578 579
			zap.Int64("MsgID", rct.ID()),
			zap.Uint64("BeginTS", rct.BeginTs()),
			zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
580 581 582
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

583
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
584
			metrics.TotalLabel).Inc()
585
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
586
			metrics.FailLabel).Inc()
587
		return &commonpb.Status{
588
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
589 590 591 592
			Reason:    err.Error(),
		}, nil
	}

593 594
	log.Debug(
		rpcDone(method),
595
		zap.String("traceID", traceID),
596
		zap.String("role", typeutil.ProxyRole),
597 598 599 600 601 602
		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))

603
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
604
		metrics.TotalLabel).Inc()
605
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
606 607
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
608
	return rct.result, nil
609 610
}

611
// DescribeCollection get the meta information of specific collection, such as schema, created timestamp and etc.
C
Cai Yudong 已提交
612
func (node *Proxy) DescribeCollection(ctx context.Context, request *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
613 614 615 616 617
	if !node.checkHealthy() {
		return &milvuspb.DescribeCollectionResponse{
			Status: unhealthyStatus(),
		}, nil
	}
618

619
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
620 621
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
622 623
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
624

625
	dct := &describeCollectionTask{
S
sunby 已提交
626
		ctx:                       ctx,
627 628
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
629
		rootCoord:                 node.rootCoord,
630 631
	}

632 633
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
634
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
635 636
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
637 638 639 640 641

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

646
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
647
			metrics.AbandonLabel).Inc()
648 649
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
650
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
651 652 653 654 655
				Reason:    err.Error(),
			},
		}, nil
	}

656 657
	log.Debug("DescribeCollection enqueued",
		zap.String("traceID", traceID),
658
		zap.String("role", typeutil.ProxyRole),
659 660 661
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
662 663
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
664 665 666

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DescribeCollection failed to WaitToFinish",
D
dragondriver 已提交
667
			zap.Error(err),
668
			zap.String("traceID", traceID),
669
			zap.String("role", typeutil.ProxyRole),
670 671 672
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTS", dct.BeginTs()),
			zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
673 674 675
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

676
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
677
			metrics.TotalLabel).Inc()
678
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
679
			metrics.FailLabel).Inc()
680

681 682
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
683
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
684 685 686 687 688
				Reason:    err.Error(),
			},
		}, nil
	}

689 690
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
691
		zap.String("role", typeutil.ProxyRole),
692 693 694 695 696 697
		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))

698
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
699
		metrics.TotalLabel).Inc()
700
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
701 702
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
703 704 705
	return dct.result, nil
}

706
// GetCollectionStatistics get the collection statistics, such as `num_rows`.
C
Cai Yudong 已提交
707
func (node *Proxy) GetCollectionStatistics(ctx context.Context, request *milvuspb.GetCollectionStatisticsRequest) (*milvuspb.GetCollectionStatisticsResponse, error) {
708 709 710 711 712
	if !node.checkHealthy() {
		return &milvuspb.GetCollectionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
713 714 715 716

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
717 718
	method := "GetCollectionStatistics"
	tr := timerecord.NewTimeRecorder(method)
719

720
	g := &getCollectionStatisticsTask{
G
godchen 已提交
721 722 723
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
724
		dataCoord:                      node.dataCoord,
725 726
	}

727 728
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
729
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
730 731
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
732 733 734 735 736

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

741
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
742
			metrics.AbandonLabel).Inc()
743

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

752 753
	log.Debug("GetCollectionStatistics enqueued",
		zap.String("traceID", traceID),
754
		zap.String("role", typeutil.ProxyRole),
755 756 757
		zap.Int64("MsgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
758 759
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
760 761 762

	if err := g.WaitToFinish(); err != nil {
		log.Warn("GetCollectionStatistics failed to WaitToFinish",
D
dragondriver 已提交
763
			zap.Error(err),
764
			zap.String("traceID", traceID),
765
			zap.String("role", typeutil.ProxyRole),
766 767 768
			zap.Int64("MsgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
769 770 771
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

772
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
773
			metrics.TotalLabel).Inc()
774
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
775
			metrics.FailLabel).Inc()
776

G
godchen 已提交
777
		return &milvuspb.GetCollectionStatisticsResponse{
778
			Status: &commonpb.Status{
779
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
780 781 782 783 784
				Reason:    err.Error(),
			},
		}, nil
	}

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

794
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
795
		metrics.TotalLabel).Inc()
796
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
797 798
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
799
	return g.result, nil
800 801
}

802
// ShowCollections list all collections in Milvus.
C
Cai Yudong 已提交
803
func (node *Proxy) ShowCollections(ctx context.Context, request *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
804 805 806 807 808
	if !node.checkHealthy() {
		return &milvuspb.ShowCollectionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
809 810 811 812
	method := "ShowCollections"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()

813
	sct := &showCollectionsTask{
G
godchen 已提交
814 815 816
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
817
		queryCoord:             node.queryCoord,
818
		rootCoord:              node.rootCoord,
819 820
	}

821
	log.Debug("ShowCollections received",
822
		zap.String("role", typeutil.ProxyRole),
823 824 825 826 827 828
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

829
	err := node.sched.ddQueue.Enqueue(sct)
830
	if err != nil {
831 832
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
833
			zap.String("role", typeutil.ProxyRole),
834 835 836 837 838 839
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

840
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
G
godchen 已提交
841
		return &milvuspb.ShowCollectionsResponse{
842
			Status: &commonpb.Status{
843
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
844 845 846 847 848
				Reason:    err.Error(),
			},
		}, nil
	}

849
	log.Debug("ShowCollections enqueued",
850
		zap.String("role", typeutil.ProxyRole),
851
		zap.Int64("MsgID", sct.ID()),
852
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
853
		zap.Uint64("TimeStamp", request.TimeStamp),
854 855 856
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
857

858 859
	err = sct.WaitToFinish()
	if err != nil {
860 861
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
862
			zap.String("role", typeutil.ProxyRole),
863 864 865 866 867 868 869
			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),
		)

870 871
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

G
godchen 已提交
872
		return &milvuspb.ShowCollectionsResponse{
873
			Status: &commonpb.Status{
874
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
875 876 877 878 879
				Reason:    err.Error(),
			},
		}, nil
	}

880
	log.Debug("ShowCollections Done",
881
		zap.String("role", typeutil.ProxyRole),
882 883 884 885 886 887 888 889
		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),
	)

890 891
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
892 893 894
	return sct.result, nil
}

895
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
896
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
897 898 899
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
900

901
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
902 903
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
904 905 906
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
907

908
	cpt := &createPartitionTask{
S
sunby 已提交
909
		ctx:                    ctx,
910 911
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
912
		rootCoord:              node.rootCoord,
913 914 915
		result:                 nil,
	}

916 917 918
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
919
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
920 921 922
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
923 924 925 926 927 928

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
929
			zap.String("role", typeutil.ProxyRole),
930 931 932 933
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

934 935
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

936
		return &commonpb.Status{
937
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
938 939 940
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
941

942 943 944
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
945
		zap.String("role", typeutil.ProxyRole),
946 947 948
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
949 950 951
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
952 953 954 955

	if err := cpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish("CreatePartition"),
D
dragondriver 已提交
956
			zap.Error(err),
957
			zap.String("traceID", traceID),
958
			zap.String("role", typeutil.ProxyRole),
959 960 961
			zap.Int64("MsgID", cpt.ID()),
			zap.Uint64("BeginTS", cpt.BeginTs()),
			zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
962 963 964 965
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

966 967
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

968
		return &commonpb.Status{
969
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
970 971 972
			Reason:    err.Error(),
		}, nil
	}
973 974 975 976

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
977
		zap.String("role", typeutil.ProxyRole),
978 979 980 981 982 983 984
		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))

985 986
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
987 988 989
	return cpt.result, nil
}

990
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
991
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
992 993 994
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
995

996
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
997 998
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
999 1000 1001
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
1002

1003
	dpt := &dropPartitionTask{
S
sunby 已提交
1004
		ctx:                  ctx,
1005 1006
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1007
		rootCoord:            node.rootCoord,
1008 1009 1010
		result:               nil,
	}

1011 1012 1013
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1014
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1015 1016 1017
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1018 1019 1020 1021 1022 1023

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1024
			zap.String("role", typeutil.ProxyRole),
1025 1026 1027 1028
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1029 1030
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

1031
		return &commonpb.Status{
1032
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1033 1034 1035
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1036

1037 1038 1039
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1040
		zap.String("role", typeutil.ProxyRole),
1041 1042 1043
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1044 1045 1046
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1047 1048 1049 1050

	if err := dpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1051
			zap.Error(err),
1052
			zap.String("traceID", traceID),
1053
			zap.String("role", typeutil.ProxyRole),
1054 1055 1056
			zap.Int64("MsgID", dpt.ID()),
			zap.Uint64("BeginTS", dpt.BeginTs()),
			zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1057 1058 1059 1060
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1061 1062
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

1063
		return &commonpb.Status{
1064
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1065 1066 1067
			Reason:    err.Error(),
		}, nil
	}
1068 1069 1070 1071

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1072
		zap.String("role", typeutil.ProxyRole),
1073 1074 1075 1076 1077 1078 1079
		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))

1080 1081
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1082 1083 1084
	return dpt.result, nil
}

1085
// HasPartition check if partition exist.
C
Cai Yudong 已提交
1086
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
1087 1088 1089 1090 1091
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1092

D
dragondriver 已提交
1093
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
1094 1095
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1096 1097 1098 1099
	method := "HasPartition"
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1100
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1101

1102
	hpt := &hasPartitionTask{
S
sunby 已提交
1103
		ctx:                 ctx,
1104 1105
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1106
		rootCoord:           node.rootCoord,
1107 1108 1109
		result:              nil,
	}

D
dragondriver 已提交
1110 1111 1112
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1113
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1114 1115 1116
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1117 1118 1119 1120 1121 1122

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

1128
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1129
			metrics.AbandonLabel).Inc()
1130

1131 1132
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1133
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1134 1135 1136 1137 1138
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1139

D
dragondriver 已提交
1140 1141 1142
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1143
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1144 1145 1146
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1147 1148 1149
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1150 1151 1152 1153

	if err := hpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1154
			zap.Error(err),
D
dragondriver 已提交
1155
			zap.String("traceID", traceID),
1156
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1157 1158 1159
			zap.Int64("MsgID", hpt.ID()),
			zap.Uint64("BeginTS", hpt.BeginTs()),
			zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1160 1161 1162 1163
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1164
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1165
			metrics.FailLabel).Inc()
1166

1167 1168
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1169
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1170 1171 1172 1173 1174
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1175 1176 1177 1178

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1179
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1180 1181 1182 1183 1184 1185 1186
		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))

1187
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1188 1189
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1190 1191 1192
	return hpt.result, nil
}

1193
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1194
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1195 1196 1197
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1198

D
dragondriver 已提交
1199
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1200 1201
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1202 1203
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1204

1205
	lpt := &loadPartitionsTask{
G
godchen 已提交
1206 1207 1208
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1209
		queryCoord:            node.queryCoord,
1210 1211
	}

1212 1213 1214
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1215
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1216 1217 1218
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1219 1220 1221 1222 1223 1224

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1225
			zap.String("role", typeutil.ProxyRole),
1226 1227 1228 1229
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1230
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1231
			metrics.AbandonLabel).Inc()
1232

1233
		return &commonpb.Status{
1234
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1235 1236 1237 1238
			Reason:    err.Error(),
		}, nil
	}

1239 1240 1241
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1242
		zap.String("role", typeutil.ProxyRole),
1243 1244 1245
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1246 1247 1248
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1249 1250 1251 1252

	if err := lpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1253
			zap.Error(err),
1254
			zap.String("traceID", traceID),
1255
			zap.String("role", typeutil.ProxyRole),
1256 1257 1258
			zap.Int64("MsgID", lpt.ID()),
			zap.Uint64("BeginTS", lpt.BeginTs()),
			zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1259 1260 1261 1262
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1263
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1264
			metrics.TotalLabel).Inc()
1265
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1266
			metrics.FailLabel).Inc()
1267

1268
		return &commonpb.Status{
1269
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1270 1271 1272 1273
			Reason:    err.Error(),
		}, nil
	}

1274 1275 1276
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1277
		zap.String("role", typeutil.ProxyRole),
1278 1279 1280 1281 1282 1283 1284
		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))

1285
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1286
		metrics.TotalLabel).Inc()
1287
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1288 1289
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1290
	return lpt.result, nil
1291 1292
}

1293
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1294
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1295 1296 1297
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1298 1299 1300 1301 1302

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

1303
	rpt := &releasePartitionsTask{
G
godchen 已提交
1304 1305 1306
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1307
		queryCoord:               node.queryCoord,
1308 1309
	}

1310
	method := "ReleasePartitions"
1311
	tr := timerecord.NewTimeRecorder(method)
1312 1313 1314 1315

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1316
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1317 1318 1319
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1320 1321 1322 1323 1324 1325

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1326
			zap.String("role", typeutil.ProxyRole),
1327 1328 1329 1330
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1331
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1332
			metrics.AbandonLabel).Inc()
1333

1334
		return &commonpb.Status{
1335
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1336 1337 1338 1339
			Reason:    err.Error(),
		}, nil
	}

1340 1341 1342
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1343
		zap.String("role", typeutil.ProxyRole),
1344 1345 1346
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1347 1348 1349
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1350 1351 1352 1353

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1354
			zap.Error(err),
1355
			zap.String("traceID", traceID),
1356
			zap.String("role", typeutil.ProxyRole),
1357 1358 1359
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1360 1361 1362 1363
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1364
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1365
			metrics.TotalLabel).Inc()
1366
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1367
			metrics.FailLabel).Inc()
1368

1369
		return &commonpb.Status{
1370
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1371 1372 1373 1374
			Reason:    err.Error(),
		}, nil
	}

1375 1376 1377
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1378
		zap.String("role", typeutil.ProxyRole),
1379 1380 1381 1382 1383 1384 1385
		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))

1386
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1387
		metrics.TotalLabel).Inc()
1388
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1389 1390
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1391
	return rpt.result, nil
1392 1393
}

1394
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1395
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1396 1397 1398 1399 1400
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1401 1402 1403 1404 1405

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

1406
	g := &getPartitionStatisticsTask{
1407 1408 1409
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1410
		dataCoord:                     node.dataCoord,
1411 1412
	}

1413
	method := "GetPartitionStatistics"
1414
	tr := timerecord.NewTimeRecorder(method)
1415 1416 1417 1418

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1419
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1420 1421 1422
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1423 1424 1425 1426 1427 1428

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1429
			zap.String("role", typeutil.ProxyRole),
1430 1431 1432 1433
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1434
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1435
			metrics.AbandonLabel).Inc()
1436

1437 1438 1439 1440 1441 1442 1443 1444
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1445 1446 1447
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1448
		zap.String("role", typeutil.ProxyRole),
1449 1450 1451
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1452 1453 1454
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1455 1456 1457 1458

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1459
			zap.Error(err),
1460
			zap.String("traceID", traceID),
1461
			zap.String("role", typeutil.ProxyRole),
1462 1463 1464
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1465 1466 1467 1468
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1469
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1470
			metrics.TotalLabel).Inc()
1471
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1472
			metrics.FailLabel).Inc()
1473

1474 1475 1476 1477 1478 1479 1480 1481
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1482 1483 1484
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1485
		zap.String("role", typeutil.ProxyRole),
1486 1487 1488 1489 1490 1491 1492
		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))

1493
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1494
		metrics.TotalLabel).Inc()
1495
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1496 1497
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1498
	return g.result, nil
1499 1500
}

1501
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1502
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1503 1504 1505 1506 1507
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1508 1509 1510 1511 1512

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

1513
	spt := &showPartitionsTask{
G
godchen 已提交
1514 1515 1516
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1517
		rootCoord:             node.rootCoord,
1518
		queryCoord:            node.queryCoord,
G
godchen 已提交
1519
		result:                nil,
1520 1521
	}

1522
	method := "ShowPartitions"
1523 1524 1525
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1526
		metrics.TotalLabel).Inc()
1527 1528 1529 1530

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1531
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1532
		zap.Any("request", request))
1533 1534 1535 1536 1537 1538

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

1542
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1543
			metrics.AbandonLabel).Inc()
1544

G
godchen 已提交
1545
		return &milvuspb.ShowPartitionsResponse{
1546
			Status: &commonpb.Status{
1547
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1548 1549 1550 1551 1552
				Reason:    err.Error(),
			},
		}, nil
	}

1553 1554 1555
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1556
		zap.String("role", typeutil.ProxyRole),
1557 1558 1559
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1560 1561
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1562 1563 1564 1565 1566
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1567
			zap.Error(err),
1568
			zap.String("traceID", traceID),
1569
			zap.String("role", typeutil.ProxyRole),
1570 1571 1572 1573 1574 1575
			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 已提交
1576

1577
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1578
			metrics.FailLabel).Inc()
1579

G
godchen 已提交
1580
		return &milvuspb.ShowPartitionsResponse{
1581
			Status: &commonpb.Status{
1582
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1583 1584 1585 1586
				Reason:    err.Error(),
			},
		}, nil
	}
1587 1588 1589 1590

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1591
		zap.String("role", typeutil.ProxyRole),
1592 1593 1594 1595 1596 1597 1598
		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))

1599
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1600 1601
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1602 1603 1604
	return spt.result, nil
}

1605
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1606
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1607 1608 1609
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1610 1611 1612 1613 1614

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

1615
	cit := &createIndexTask{
S
sunby 已提交
1616
		ctx:                ctx,
1617 1618
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1619
		rootCoord:          node.rootCoord,
1620 1621
	}

D
dragondriver 已提交
1622
	method := "CreateIndex"
1623
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1624 1625 1626 1627

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1628
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1629 1630 1631 1632
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1633 1634 1635 1636 1637 1638

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1639
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1640 1641 1642 1643 1644
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1645
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1646
			metrics.AbandonLabel).Inc()
1647

1648
		return &commonpb.Status{
1649
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1650 1651 1652 1653
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1654 1655 1656
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1657
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1658 1659 1660
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1661 1662 1663 1664
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1665 1666 1667 1668

	if err := cit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1669
			zap.Error(err),
D
dragondriver 已提交
1670
			zap.String("traceID", traceID),
1671
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1672 1673 1674
			zap.Int64("MsgID", cit.ID()),
			zap.Uint64("BeginTs", cit.BeginTs()),
			zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1675 1676 1677 1678 1679
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1680
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1681
			metrics.TotalLabel).Inc()
1682
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1683
			metrics.FailLabel).Inc()
1684

1685
		return &commonpb.Status{
1686
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1687 1688 1689 1690
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1691 1692 1693
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1694
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1695 1696 1697 1698 1699 1700 1701 1702
		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))

1703
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1704
		metrics.TotalLabel).Inc()
1705
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1706 1707
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1708 1709 1710
	return cit.result, nil
}

1711
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
1712
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
1713 1714 1715 1716 1717
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1718 1719 1720 1721 1722

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

1723
	dit := &describeIndexTask{
S
sunby 已提交
1724
		ctx:                  ctx,
1725 1726
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1727
		rootCoord:            node.rootCoord,
1728 1729
	}

1730 1731 1732
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
1733
	tr := timerecord.NewTimeRecorder(method)
1734 1735 1736 1737

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1738
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1739 1740 1741
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1742 1743 1744 1745 1746 1747 1748
		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),
1749
			zap.String("role", typeutil.ProxyRole),
1750 1751 1752 1753 1754
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

1755
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1756
			metrics.AbandonLabel).Inc()
1757

1758 1759
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1760
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1761 1762 1763 1764 1765
				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", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1773 1774 1775
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1776 1777 1778 1779 1780
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1781
			zap.Error(err),
1782
			zap.String("traceID", traceID),
1783
			zap.String("role", typeutil.ProxyRole),
1784 1785 1786
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1787 1788 1789
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
1790
			zap.String("index name", indexName))
D
dragondriver 已提交
1791

Z
zhenshan.cao 已提交
1792 1793 1794 1795
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
1796
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1797
			metrics.TotalLabel).Inc()
1798
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1799
			metrics.FailLabel).Inc()
1800

1801 1802
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
1803
				ErrorCode: errCode,
1804 1805 1806 1807 1808
				Reason:    err.Error(),
			},
		}, nil
	}

1809 1810 1811
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1812
		zap.String("role", typeutil.ProxyRole),
1813 1814 1815 1816 1817 1818 1819 1820
		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))

1821
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1822
		metrics.TotalLabel).Inc()
1823
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1824 1825
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1826 1827 1828
	return dit.result, nil
}

1829
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
1830
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
1831 1832 1833
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1834 1835 1836 1837 1838

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

1839
	dit := &dropIndexTask{
S
sunby 已提交
1840
		ctx:              ctx,
B
BossZou 已提交
1841 1842
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1843
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1844
	}
G
godchen 已提交
1845

D
dragondriver 已提交
1846
	method := "DropIndex"
1847
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1848 1849 1850 1851

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1852
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
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))

D
dragondriver 已提交
1858 1859 1860 1861 1862
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1863
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1864 1865 1866 1867
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1868
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1869
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
1870

B
BossZou 已提交
1871
		return &commonpb.Status{
1872
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1873 1874 1875
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1876

D
dragondriver 已提交
1877 1878 1879
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1880
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1881 1882 1883
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1884 1885 1886 1887
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
1888 1889 1890 1891

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1892
			zap.Error(err),
D
dragondriver 已提交
1893
			zap.String("traceID", traceID),
1894
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1895 1896 1897
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1898 1899 1900 1901 1902
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

1903
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1904
			metrics.TotalLabel).Inc()
1905
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1906
			metrics.FailLabel).Inc()
1907

B
BossZou 已提交
1908
		return &commonpb.Status{
1909
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1910 1911 1912
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1913 1914 1915 1916

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1917
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1918 1919 1920 1921 1922 1923 1924 1925
		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))

1926
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1927
		metrics.TotalLabel).Inc()
1928
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1929 1930
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
1931 1932 1933
	return dit.result, nil
}

1934 1935
// 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 已提交
1936
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1937 1938 1939 1940 1941
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1942 1943 1944 1945 1946

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

1947
	gibpt := &getIndexBuildProgressTask{
1948 1949 1950
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1951 1952
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1953
		dataCoord:                    node.dataCoord,
1954 1955
	}

1956
	method := "GetIndexBuildProgress"
1957
	tr := timerecord.NewTimeRecorder(method)
1958 1959 1960 1961

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1962
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1963 1964 1965 1966
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1967 1968 1969 1970 1971 1972

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1973
			zap.String("role", typeutil.ProxyRole),
1974 1975 1976 1977
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1978
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1979
			metrics.AbandonLabel).Inc()
1980

1981 1982 1983 1984 1985 1986 1987 1988
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1989 1990 1991
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1992
		zap.String("role", typeutil.ProxyRole),
1993 1994 1995
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
1996 1997 1998 1999
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2000 2001 2002 2003

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2004
			zap.Error(err),
2005
			zap.String("traceID", traceID),
2006
			zap.String("role", typeutil.ProxyRole),
2007 2008 2009
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2010 2011 2012 2013
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2014
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2015
			metrics.TotalLabel).Inc()
2016
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2017
			metrics.FailLabel).Inc()
2018 2019 2020 2021 2022 2023 2024 2025

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2026 2027 2028 2029

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2030
		zap.String("role", typeutil.ProxyRole),
2031 2032 2033 2034 2035 2036 2037 2038
		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))
2039

2040
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2041
		metrics.TotalLabel).Inc()
2042
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2043 2044
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2045
	return gibpt.result, nil
2046 2047
}

2048
// GetIndexState get the build-state of index.
C
Cai Yudong 已提交
2049
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2050 2051 2052 2053 2054
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2055 2056 2057 2058 2059

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

2060
	dipt := &getIndexStateTask{
G
godchen 已提交
2061 2062 2063
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2064 2065
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2066 2067
	}

2068
	method := "GetIndexState"
2069
	tr := timerecord.NewTimeRecorder(method)
2070 2071 2072 2073

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2074
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2075 2076 2077 2078
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2079 2080 2081 2082 2083 2084

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2085
			zap.String("role", typeutil.ProxyRole),
2086 2087 2088 2089 2090
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2091
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2092
			metrics.AbandonLabel).Inc()
2093

G
godchen 已提交
2094
		return &milvuspb.GetIndexStateResponse{
2095
			Status: &commonpb.Status{
2096
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2097 2098 2099 2100 2101
				Reason:    err.Error(),
			},
		}, nil
	}

2102 2103 2104
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2105
		zap.String("role", typeutil.ProxyRole),
2106 2107 2108
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2109 2110 2111 2112
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2113 2114 2115 2116

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2117
			zap.Error(err),
2118
			zap.String("traceID", traceID),
2119
			zap.String("role", typeutil.ProxyRole),
2120 2121 2122
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2123 2124 2125 2126 2127
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2128
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2129
			metrics.TotalLabel).Inc()
2130
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2131
			metrics.FailLabel).Inc()
2132

G
godchen 已提交
2133
		return &milvuspb.GetIndexStateResponse{
2134
			Status: &commonpb.Status{
2135
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2136 2137 2138 2139 2140
				Reason:    err.Error(),
			},
		}, nil
	}

2141 2142 2143
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2144
		zap.String("role", typeutil.ProxyRole),
2145 2146 2147 2148 2149 2150 2151 2152
		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))

2153
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2154
		metrics.TotalLabel).Inc()
2155
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2156 2157
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2158 2159 2160
	return dipt.result, nil
}

2161
// Insert insert records into collection.
C
Cai Yudong 已提交
2162
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2163 2164 2165 2166 2167 2168
	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))

2169 2170 2171 2172 2173
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2174 2175
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2176

2177
	it := &insertTask{
2178 2179
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2180
		// req:       request,
2181 2182 2183 2184
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2185
			InsertRequest: internalpb.InsertRequest{
2186
				Base: &commonpb.MsgBase{
X
xige-16 已提交
2187 2188 2189
					MsgType:  commonpb.MsgType_Insert,
					MsgID:    0,
					SourceID: Params.ProxyCfg.ProxyID,
2190 2191 2192
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2193 2194 2195
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2196
				// RowData: transfer column based request to this
2197 2198
			},
		},
2199
		rowIDAllocator: node.idAllocator,
2200
		segIDAssigner:  node.segAssigner,
2201
		chMgr:          node.chMgr,
2202
		chTicker:       node.chTicker,
2203
	}
2204 2205

	if len(it.PartitionName) <= 0 {
2206
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2207 2208
	}

X
Xiangyu Wang 已提交
2209
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2210
		numRows := request.NumRows
2211 2212 2213 2214
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2215

X
Xiangyu Wang 已提交
2216 2217 2218 2219 2220 2221 2222
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2223 2224
	}

X
Xiangyu Wang 已提交
2225
	log.Debug("Enqueue insert request in Proxy",
2226
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2227 2228 2229 2230 2231
		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)),
2232 2233
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2234

X
Xiangyu Wang 已提交
2235 2236
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2237
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2238
		return constructFailedResponse(err), nil
2239
	}
D
dragondriver 已提交
2240

X
Xiangyu Wang 已提交
2241
	log.Debug("Detail of insert request in Proxy",
2242
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2243
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2244 2245 2246 2247 2248
		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 已提交
2249 2250 2251 2252 2253
		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))
2254
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2255
			metrics.TotalLabel).Inc()
2256
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2257
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2258 2259 2260 2261 2262
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2263
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2264 2265 2266 2267 2268 2269 2270 2271 2272 2273 2274
			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
X
xige-16 已提交
2275
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2276

2277
	metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2278 2279 2280
		metrics.SuccessLabel).Inc()
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Add(float64(it.result.InsertCnt))
	metrics.ProxyInsertLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
2281 2282 2283
	return it.result, nil
}

2284
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2285
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2286 2287 2288
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2289 2290
	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))
2291

G
groot 已提交
2292 2293 2294 2295 2296 2297
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2298 2299 2300
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

C
Cai Yudong 已提交
2301 2302 2303 2304 2305 2306 2307
	deleteReq := &milvuspb.DeleteRequest{
		DbName:         request.DbName,
		CollectionName: request.CollectionName,
		PartitionName:  request.PartitionName,
		Expr:           request.Expr,
	}

2308
	dt := &deleteTask{
C
Cai Yudong 已提交
2309 2310 2311
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       deleteReq,
G
godchen 已提交
2312
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2313 2314 2315
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2316 2317 2318 2319 2320 2321 2322 2323
			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 已提交
2324 2325 2326 2327
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2328 2329
	}

2330
	log.Debug("Enqueue delete request in Proxy",
2331
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2332 2333 2334 2335
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2336 2337 2338 2339

	// 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))
2340
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2341
			metrics.FailLabel).Inc()
2342

G
groot 已提交
2343 2344 2345 2346 2347 2348 2349 2350
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2351
	log.Debug("Detail of delete request in Proxy",
2352
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2353 2354 2355 2356 2357
		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),
2358 2359
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2360

2361 2362
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
2363
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2364
			metrics.TotalLabel).Inc()
2365
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2366
			metrics.FailLabel).Inc()
G
groot 已提交
2367 2368 2369 2370 2371 2372 2373 2374
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2375
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2376
		metrics.TotalLabel).Inc()
2377
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2378 2379
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2380 2381 2382
	return dt.result, nil
}

2383
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2384
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2385 2386 2387 2388 2389
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2390 2391
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2392

C
cai.zhang 已提交
2393 2394
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2395 2396
	traceID, _, _ := trace.InfoFromSpan(sp)

2397
	qt := &searchTask{
S
sunby 已提交
2398
		ctx:       ctx,
2399
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2400
		SearchRequest: &internalpb.SearchRequest{
2401
			Base: &commonpb.MsgBase{
2402
				MsgType:  commonpb.MsgType_Search,
2403
				SourceID: Params.ProxyCfg.ProxyID,
2404
			},
2405
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2406
		},
2407
		resultBuf: make(chan []*internalpb.SearchResults, 1),
2408 2409
		query:     request,
		chMgr:     node.chMgr,
2410
		qc:        node.queryCoord,
2411
		tr:        timerecord.NewTimeRecorder("search"),
2412 2413
	}

2414 2415 2416 2417 2418
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

	log.Debug(
		rpcReceived(method),
D
dragondriver 已提交
2419
		zap.String("traceID", traceID),
2420
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2421 2422 2423 2424 2425
		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)),
2426 2427 2428 2429
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2430

2431 2432 2433
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2434 2435
			zap.Error(err),
			zap.String("traceID", traceID),
2436
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2437 2438 2439 2440 2441 2442
			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),
2443 2444 2445
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2446

2447
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2448
			metrics.SearchLabel, metrics.AbandonLabel).Inc()
2449

2450 2451
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2452
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2453 2454 2455 2456 2457
				Reason:    err.Error(),
			},
		}, nil
	}

2458 2459
	log.Debug(
		rpcEnqueued(method),
D
dragondriver 已提交
2460
		zap.String("traceID", traceID),
2461
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2462
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2463 2464 2465 2466 2467
		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),
2468
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2469 2470 2471 2472
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2473

2474 2475 2476
	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2477
			zap.Error(err),
D
dragondriver 已提交
2478
			zap.String("traceID", traceID),
2479
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2480
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2481 2482 2483 2484
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2485
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2486 2487 2488 2489
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2490
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2491
			metrics.SearchLabel, metrics.TotalLabel).Inc()
2492

2493 2494
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			metrics.SearchLabel, metrics.FailLabel).Inc()
2495

2496 2497
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2498
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2499 2500 2501 2502 2503
				Reason:    err.Error(),
			},
		}, nil
	}

2504 2505
	log.Debug(
		rpcDone(method),
D
dragondriver 已提交
2506
		zap.String("traceID", traceID),
2507
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2508 2509 2510 2511 2512 2513
		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)),
2514 2515 2516 2517
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2518

2519
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2520
		metrics.SearchLabel, metrics.TotalLabel).Inc()
2521
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2522
		metrics.SearchLabel, metrics.SuccessLabel).Inc()
2523
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2524
		metrics.SearchLabel).Set(float64(qt.result.Results.NumQueries))
C
cai.zhang 已提交
2525
	searchDur := tr.ElapseSpan().Milliseconds()
2526
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2527 2528
		metrics.SearchLabel).Observe(float64(searchDur))
	metrics.ProxySearchLatencyPerNQ.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Observe(float64(searchDur) / float64(qt.result.Results.NumQueries))
2529 2530 2531
	return qt.result, nil
}

2532
// Flush notify data nodes to persist the data of collection.
2533 2534 2535 2536 2537 2538 2539
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2540
	if !node.checkHealthy() {
2541 2542
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2543
	}
D
dragondriver 已提交
2544 2545 2546 2547 2548

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

2549
	ft := &flushTask{
T
ThreadDao 已提交
2550 2551 2552
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2553
		dataCoord:    node.dataCoord,
2554 2555
	}

D
dragondriver 已提交
2556
	method := "Flush"
2557 2558
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2559 2560 2561 2562

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2563
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2564 2565
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2566 2567 2568 2569 2570 2571

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

2576 2577
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

2578 2579
		resp.Status.Reason = err.Error()
		return resp, nil
2580 2581
	}

D
dragondriver 已提交
2582 2583 2584
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2585
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2586 2587 2588
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2589 2590
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2591 2592 2593 2594

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2595
			zap.Error(err),
D
dragondriver 已提交
2596
			zap.String("traceID", traceID),
2597
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2598 2599 2600
			zap.Int64("MsgID", ft.ID()),
			zap.Uint64("BeginTs", ft.BeginTs()),
			zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2601 2602 2603
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

2604 2605
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

D
dragondriver 已提交
2606
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2607 2608
		resp.Status.Reason = err.Error()
		return resp, nil
2609 2610
	}

D
dragondriver 已提交
2611 2612 2613
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2614
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2615 2616 2617 2618 2619 2620
		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))

2621 2622
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2623
	return ft.result, nil
2624 2625
}

2626
// Query get the records by primary keys.
C
Cai Yudong 已提交
2627
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2628 2629 2630 2631 2632
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2633

D
dragondriver 已提交
2634 2635 2636
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2637
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2638

2639
	qt := &queryTask{
2640 2641 2642 2643 2644
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
2645
				SourceID: Params.ProxyCfg.ProxyID,
2646
			},
2647
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2648 2649
		},
		resultBuf: make(chan []*internalpb.RetrieveResults),
2650
		query:     request,
2651 2652
		chMgr:     node.chMgr,
		qc:        node.queryCoord,
2653 2654
	}

D
dragondriver 已提交
2655 2656 2657 2658 2659
	method := "Query"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2660
		zap.String("role", typeutil.ProxyRole),
2661 2662 2663
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
G
godchen 已提交
2664

D
dragondriver 已提交
2665 2666 2667 2668 2669 2670
	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),
2671 2672 2673
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2674

2675
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2676
			metrics.QueryLabel, metrics.FailLabel).Inc()
2677 2678 2679 2680 2681 2682
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2683 2684
	}

D
dragondriver 已提交
2685 2686 2687
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2688
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2689 2690 2691
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
2692 2693 2694
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2695 2696 2697 2698 2699 2700

	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2701
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2702 2703 2704
			zap.Int64("MsgID", qt.ID()),
			zap.Uint64("BeginTs", qt.BeginTs()),
			zap.Uint64("EndTs", qt.EndTs()),
2705 2706 2707
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
2708

2709
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2710
			metrics.QueryLabel, metrics.TotalLabel).Inc()
2711
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2712
			metrics.QueryLabel, metrics.FailLabel).Inc()
2713

2714 2715 2716 2717 2718 2719 2720
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2721

D
dragondriver 已提交
2722 2723 2724 2725 2726 2727 2728
	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()),
2729 2730 2731
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2732

2733
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2734
		metrics.QueryLabel, metrics.TotalLabel).Inc()
2735
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2736
		metrics.QueryLabel, metrics.SuccessLabel).Inc()
2737
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2738
		metrics.QueryLabel).Set(float64(len(qt.result.FieldsData)))
2739
	metrics.ProxySendMessageLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2740
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2741 2742 2743 2744 2745
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2746

2747
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
2748 2749 2750 2751
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2752 2753 2754 2755 2756

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

Y
Yusup 已提交
2757 2758 2759 2760 2761 2762 2763
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
2764
	method := "CreateAlias"
2765 2766
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2767 2768 2769 2770 2771 2772 2773 2774 2775 2776 2777 2778 2779 2780 2781 2782 2783 2784 2785

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

2786 2787
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

Y
Yusup 已提交
2788 2789 2790 2791 2792 2793
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2794 2795 2796
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2797
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2798 2799 2800 2801
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2802 2803
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2804 2805 2806 2807

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2808
			zap.Error(err),
D
dragondriver 已提交
2809
			zap.String("traceID", traceID),
2810
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2811 2812 2813 2814
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2815 2816
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
2817
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
2818 2819 2820 2821 2822 2823 2824

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

D
dragondriver 已提交
2825 2826 2827 2828 2829 2830 2831 2832 2833 2834 2835
	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))

2836 2837
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2838 2839 2840
	return cat.result, nil
}

2841
// DropAlias alter the alias of collection.
Y
Yusup 已提交
2842 2843 2844 2845
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2846 2847 2848 2849 2850

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

Y
Yusup 已提交
2851 2852 2853 2854 2855 2856 2857
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
2858
	method := "DropAlias"
2859 2860
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2861 2862 2863 2864 2865 2866 2867 2868 2869 2870 2871 2872 2873 2874 2875 2876

	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))
2877
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2878

Y
Yusup 已提交
2879 2880 2881 2882 2883 2884
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2885 2886 2887
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2888
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2889 2890 2891 2892
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2893
		zap.String("alias", request.Alias))
D
dragondriver 已提交
2894 2895 2896 2897

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2898
			zap.Error(err),
D
dragondriver 已提交
2899
			zap.String("traceID", traceID),
2900
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2901 2902 2903 2904
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2905 2906
			zap.String("alias", request.Alias))

2907 2908
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

Y
Yusup 已提交
2909 2910 2911 2912 2913 2914
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2915 2916 2917 2918 2919 2920 2921 2922 2923 2924
	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))

2925 2926
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2927 2928 2929
	return dat.result, nil
}

2930
// AlterAlias alter alias of collection.
Y
Yusup 已提交
2931 2932 2933 2934
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2935 2936 2937 2938 2939

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

Y
Yusup 已提交
2940 2941 2942 2943 2944 2945 2946
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
2947
	method := "AlterAlias"
2948 2949
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2950 2951 2952 2953 2954 2955 2956 2957 2958 2959 2960 2961 2962 2963 2964 2965 2966 2967

	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))
2968
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2969

Y
Yusup 已提交
2970 2971 2972 2973 2974 2975
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2976 2977 2978
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2979
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2980 2981 2982 2983
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2984 2985
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2986 2987 2988 2989

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2990
			zap.Error(err),
D
dragondriver 已提交
2991
			zap.String("traceID", traceID),
2992
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2993 2994 2995 2996
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2997 2998 2999
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

3000 3001
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

Y
Yusup 已提交
3002 3003 3004 3005 3006 3007
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3008 3009 3010 3011 3012 3013 3014 3015 3016 3017 3018
	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))

3019 3020
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3021 3022 3023
	return aat.result, nil
}

3024
// CalcDistance calculates the distances between vectors.
3025
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3026 3027 3028 3029 3030
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3031
	param, _ := funcutil.GetAttrByKeyFromRepeatedKV("metric", request.GetParams())
3032 3033 3034 3035 3036 3037 3038 3039
	metric, err := distance.ValidateMetricType(param)
	if err != nil {
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
3040 3041
	}

3042 3043 3044 3045
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3046 3047
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3048

3049 3050 3051 3052 3053
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3054 3055
		}

3056
		qt := &queryTask{
3057 3058 3059 3060 3061
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
3062
					SourceID: Params.ProxyCfg.ProxyID,
3063
				},
3064
				ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
3065
			},
3066
			resultBuf: make(chan []*internalpb.RetrieveResults),
3067
			query:     queryRequest,
3068
			chMgr:     node.chMgr,
3069
			qc:        node.queryCoord,
Y
yukun 已提交
3070
			ids:       ids.IdArray,
3071 3072
		}

3073
		err := node.sched.dqQueue.Enqueue(qt)
3074
		if err != nil {
3075 3076 3077
			log.Debug("CalcDistance queryTask failed to enqueue",
				zap.Error(err),
				zap.String("traceID", traceID),
3078
				zap.String("role", typeutil.ProxyRole),
3079 3080 3081 3082
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames))

3083 3084 3085 3086 3087
			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3088
			}, err
3089
		}
3090 3091 3092

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
3093
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3094 3095 3096 3097 3098 3099
			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))
3100 3101 3102 3103

		err = qt.WaitToFinish()
		if err != nil {
			log.Debug("CalcDistance queryTask failed to WaitToFinish",
G
godchen 已提交
3104
				zap.Error(err),
3105
				zap.String("traceID", traceID),
3106
				zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3107 3108 3109 3110 3111 3112
				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))
3113 3114 3115 3116 3117 3118

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3119
			}, err
3120
		}
3121 3122 3123

		log.Debug("CalcDistance queryTask Done",
			zap.String("traceID", traceID),
3124
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3125 3126 3127 3128 3129 3130
			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))
3131 3132

		return &milvuspb.QueryResults{
3133 3134
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3135 3136 3137
		}, nil
	}

3138 3139 3140 3141 3142 3143 3144 3145 3146 3147 3148 3149 3150 3151
	// 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 {
3152
			return nil, errors.New("failed to fetch vectors")
3153 3154 3155 3156 3157 3158 3159 3160 3161 3162 3163 3164 3165 3166 3167 3168
		}

		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))
3169
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3170 3171 3172 3173 3174 3175 3176 3177 3178 3179 3180 3181 3182 3183 3184 3185 3186 3187 3188 3189 3190 3191 3192 3193 3194 3195
				}
				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))
3196
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3197 3198 3199 3200 3201 3202 3203 3204 3205 3206 3207 3208
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

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

3209
		return nil, errors.New("failed to fetch vectors")
3210 3211
	}

3212 3213
	log.Debug("CalcDistance received",
		zap.String("traceID", traceID),
3214
		zap.String("role", typeutil.ProxyRole),
3215
		zap.String("metric", metric))
G
godchen 已提交
3216

3217 3218 3219
	vectorsLeft := request.GetOpLeft().GetDataArray()
	opLeft := request.GetOpLeft().GetIdArray()
	if opLeft != nil {
3220 3221
		log.Debug("OpLeft IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3222
			zap.String("role", typeutil.ProxyRole))
3223

3224
		result, err := query(opLeft)
3225
		if err != nil {
3226 3227 3228
			log.Debug("Failed to get left vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3229
				zap.String("role", typeutil.ProxyRole))
3230

3231 3232 3233 3234 3235 3236 3237 3238
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3239 3240
		log.Debug("OpLeft IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3241
			zap.String("role", typeutil.ProxyRole))
3242

3243 3244
		vectorsLeft, err = arrangeFunc(opLeft, result.FieldsData)
		if err != nil {
3245 3246 3247
			log.Debug("Failed to re-arrange left vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3248
				zap.String("role", typeutil.ProxyRole))
3249

3250 3251 3252 3253 3254 3255
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3256
		}
3257 3258 3259

		log.Debug("Re-arrange left vectors done",
			zap.String("traceID", traceID),
3260
			zap.String("role", typeutil.ProxyRole))
3261 3262
	}

G
groot 已提交
3263
	if vectorsLeft == nil {
3264 3265 3266
		msg := "Left vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3267
			zap.String("role", typeutil.ProxyRole))
3268

G
groot 已提交
3269 3270 3271
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3272
				Reason:    msg,
G
groot 已提交
3273 3274 3275 3276
			},
		}, nil
	}

3277 3278 3279
	vectorsRight := request.GetOpRight().GetDataArray()
	opRight := request.GetOpRight().GetIdArray()
	if opRight != nil {
3280 3281
		log.Debug("OpRight IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3282
			zap.String("role", typeutil.ProxyRole))
3283

3284
		result, err := query(opRight)
3285
		if err != nil {
3286 3287 3288
			log.Debug("Failed to get right vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3289
				zap.String("role", typeutil.ProxyRole))
3290

3291 3292 3293 3294 3295 3296 3297 3298
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3299 3300
		log.Debug("OpRight IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3301
			zap.String("role", typeutil.ProxyRole))
3302

3303 3304
		vectorsRight, err = arrangeFunc(opRight, result.FieldsData)
		if err != nil {
3305 3306 3307
			log.Debug("Failed to re-arrange right vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3308
				zap.String("role", typeutil.ProxyRole))
3309

3310 3311 3312 3313 3314 3315
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3316
		}
3317 3318 3319

		log.Debug("Re-arrange right vectors done",
			zap.String("traceID", traceID),
3320
			zap.String("role", typeutil.ProxyRole))
3321 3322
	}

G
groot 已提交
3323
	if vectorsRight == nil {
3324 3325 3326
		msg := "Right vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3327
			zap.String("role", typeutil.ProxyRole))
3328

G
groot 已提交
3329 3330 3331
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3332
				Reason:    msg,
G
groot 已提交
3333 3334 3335 3336
			},
		}, nil
	}

3337
	if vectorsLeft.Dim != vectorsRight.Dim {
3338 3339 3340
		msg := "Vectors dimension is not equal"
		log.Debug(msg,
			zap.String("traceID", traceID),
3341
			zap.String("role", typeutil.ProxyRole))
3342

3343 3344 3345
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3346
				Reason:    msg,
3347 3348 3349 3350 3351
			},
		}, nil
	}

	if vectorsLeft.GetFloatVector() != nil && vectorsRight.GetFloatVector() != nil {
3352 3353
		distances, err := distance.CalcFloatDistance(vectorsLeft.Dim, vectorsLeft.GetFloatVector().Data, vectorsRight.GetFloatVector().Data, metric)
		if err != nil {
3354 3355 3356
			log.Debug("Failed to CalcFloatDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3357
				zap.String("role", typeutil.ProxyRole))
3358

3359 3360 3361 3362 3363 3364 3365 3366
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3367 3368 3369
		log.Debug("CalcFloatDistance done",
			zap.Error(err),
			zap.String("traceID", traceID),
3370
			zap.String("role", typeutil.ProxyRole))
3371

3372 3373 3374 3375 3376 3377 3378 3379 3380 3381
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
			Array: &milvuspb.CalcDistanceResults_FloatDist{
				FloatDist: &schemapb.FloatArray{
					Data: distances,
				},
			},
		}, nil
	}

3382
	if vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetBinaryVector() != nil {
G
groot 已提交
3383
		hamming, err := distance.CalcHammingDistance(vectorsLeft.Dim, vectorsLeft.GetBinaryVector(), vectorsRight.GetBinaryVector())
3384
		if err != nil {
3385 3386 3387
			log.Debug("Failed to CalcHammingDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3388
				zap.String("role", typeutil.ProxyRole))
3389

3390 3391 3392 3393 3394 3395 3396 3397 3398
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

		if metric == distance.HAMMING {
3399 3400
			log.Debug("CalcHammingDistance done",
				zap.String("traceID", traceID),
3401
				zap.String("role", typeutil.ProxyRole))
3402

3403 3404 3405 3406
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_IntDist{
					IntDist: &schemapb.IntArray{
G
groot 已提交
3407
						Data: hamming,
3408 3409 3410 3411 3412 3413
					},
				},
			}, nil
		}

		if metric == distance.TANIMOTO {
G
groot 已提交
3414
			tanimoto, err := distance.CalcTanimotoCoefficient(vectorsLeft.Dim, hamming)
3415
			if err != nil {
3416 3417 3418
				log.Debug("Failed to CalcTanimotoCoefficient",
					zap.Error(err),
					zap.String("traceID", traceID),
3419
					zap.String("role", typeutil.ProxyRole))
3420

3421 3422 3423 3424 3425 3426 3427 3428
				return &milvuspb.CalcDistanceResults{
					Status: &commonpb.Status{
						ErrorCode: commonpb.ErrorCode_UnexpectedError,
						Reason:    err.Error(),
					},
				}, nil
			}

3429 3430
			log.Debug("CalcTanimotoCoefficient done",
				zap.String("traceID", traceID),
3431
				zap.String("role", typeutil.ProxyRole))
3432

3433 3434 3435 3436 3437 3438 3439 3440 3441 3442 3443
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_FloatDist{
					FloatDist: &schemapb.FloatArray{
						Data: tanimoto,
					},
				},
			}, nil
		}
	}

3444
	err = errors.New("unexpected error")
3445
	if (vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetFloatVector() != nil) || (vectorsLeft.GetFloatVector() != nil && vectorsRight.GetBinaryVector() != nil) {
3446
		err = errors.New("cannot calculate distance between binary vectors and float vectors")
3447 3448
	}

3449 3450 3451
	log.Debug("Failed to CalcDistance",
		zap.Error(err),
		zap.String("traceID", traceID),
3452
		zap.String("role", typeutil.ProxyRole))
3453

3454 3455 3456 3457 3458 3459 3460 3461
	return &milvuspb.CalcDistanceResults{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		},
	}, nil
}

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

3467
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3468
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3469
	log.Debug("GetPersistentSegmentInfo",
3470
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3471 3472 3473
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3474
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3475
		Status: &commonpb.Status{
3476
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3477 3478
		},
	}
3479 3480 3481 3482
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3483 3484 3485
	method := "GetPersistentSegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
3486
		metrics.TotalLabel).Inc()
G
godchen 已提交
3487
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
X
XuanYang-cn 已提交
3488
	if err != nil {
3489
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3490 3491
		return resp, nil
	}
3492
	infoResp, err := node.dataCoord.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
X
XuanYang-cn 已提交
3493
		Base: &commonpb.MsgBase{
3494
			MsgType:   commonpb.MsgType_SegmentInfo,
X
XuanYang-cn 已提交
3495 3496
			MsgID:     0,
			Timestamp: 0,
3497
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3498 3499 3500 3501
		},
		SegmentIDs: segments,
	})
	if err != nil {
3502
		log.Debug("GetPersistentSegmentInfo fail", zap.Error(err))
3503
		resp.Status.Reason = fmt.Errorf("dataCoord:GetSegmentInfo, err:%w", err).Error()
X
XuanYang-cn 已提交
3504 3505
		return resp, nil
	}
3506
	log.Debug("GetPersistentSegmentInfo ", zap.Int("len(infos)", len(infoResp.Infos)), zap.Any("status", infoResp.Status))
3507
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3508 3509 3510 3511 3512 3513
		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 已提交
3514
			SegmentID:    info.ID,
X
XuanYang-cn 已提交
3515 3516
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
S
sunby 已提交
3517
			NumRows:      info.NumOfRows,
X
XuanYang-cn 已提交
3518 3519 3520
			State:        info.State,
		}
	}
3521
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
3522 3523
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
3524
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
X
XuanYang-cn 已提交
3525 3526 3527 3528
	resp.Infos = persistentInfos
	return resp, nil
}

J
jingkl 已提交
3529
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3530
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3531
	log.Debug("GetQuerySegmentInfo",
3532
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3533 3534 3535
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3536
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3537
		Status: &commonpb.Status{
3538
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3539 3540
		},
	}
3541 3542 3543 3544
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
G
godchen 已提交
3545
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
Z
zhenshan.cao 已提交
3546 3547 3548 3549
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3550 3551 3552 3553 3554
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3555
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
Z
zhenshan.cao 已提交
3556
		Base: &commonpb.MsgBase{
3557
			MsgType:   commonpb.MsgType_SegmentInfo,
Z
zhenshan.cao 已提交
3558 3559
			MsgID:     0,
			Timestamp: 0,
3560
			SourceID:  Params.ProxyCfg.ProxyID,
Z
zhenshan.cao 已提交
3561
		},
3562 3563
		CollectionID: collID,
		SegmentIDs:   segments,
Z
zhenshan.cao 已提交
3564 3565
	})
	if err != nil {
3566
		log.Error("Failed to get segment info from QueryCoord",
3567
			zap.Int64s("segmentIDs", segments), zap.Error(err))
Z
zhenshan.cao 已提交
3568 3569 3570
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3571
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3572
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3573
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3574 3575 3576 3577 3578 3579 3580 3581 3582 3583 3584 3585 3586
		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,
3587
			NodeID:       info.NodeID,
X
xige-16 已提交
3588
			State:        info.SegmentState,
Z
zhenshan.cao 已提交
3589 3590
		}
	}
3591
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3592 3593 3594 3595
	resp.Infos = queryInfos
	return resp, nil
}

C
Cai Yudong 已提交
3596
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3597
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3598
		Base: &commonpb.MsgBase{
3599
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3600 3601
			MsgID:     0,
			Timestamp: 0,
3602
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3603 3604 3605 3606 3607 3608 3609
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3610
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3611 3612 3613
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3614
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3615
		Base: &commonpb.MsgBase{
3616
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3617 3618
			MsgID:     0,
			Timestamp: 0,
3619
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3620 3621 3622 3623 3624 3625 3626 3627
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3628
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3629 3630 3631 3632 3633
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3634
		showSegmentResponse, err := node.rootCoord.ShowSegments(ctx, &milvuspb.ShowSegmentsRequest{
X
XuanYang-cn 已提交
3635
			Base: &commonpb.MsgBase{
3636
				MsgType:   commonpb.MsgType_ShowSegments,
X
XuanYang-cn 已提交
3637 3638
				MsgID:     0,
				Timestamp: 0,
3639
				SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3640 3641 3642 3643 3644 3645 3646
			},
			CollectionID: collectionID,
			PartitionID:  partitionID,
		})
		if err != nil {
			return nil, err
		}
3647
		if showSegmentResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3648 3649 3650 3651 3652 3653
			return nil, errors.New(showSegmentResponse.Status.Reason)
		}
		ret = append(ret, showSegmentResponse.SegmentIDs...)
	}
	return ret, nil
}
3654

J
jingkl 已提交
3655
// Dummy handles dummy request
C
Cai Yudong 已提交
3656
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3657 3658 3659 3660 3661 3662 3663 3664 3665 3666 3667
	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
	}

3668 3669
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3670
		if err != nil {
3671
			log.Debug("Failed to parse dummy query request")
3672 3673 3674
			return failedResponse, nil
		}

3675
		request := &milvuspb.QueryRequest{
3676 3677 3678
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3679
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3680 3681
		}

3682
		_, err = node.Query(ctx, request)
3683
		if err != nil {
3684
			log.Debug("Failed to execute dummy query")
3685 3686
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3687 3688 3689 3690 3691 3692

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

3693 3694
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3695 3696
}

J
jingkl 已提交
3697
// RegisterLink registers a link
C
Cai Yudong 已提交
3698
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
G
godchen 已提交
3699
	code := node.stateCode.Load().(internalpb.StateCode)
D
dragondriver 已提交
3700
	log.Debug("RegisterLink",
3701
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3702
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3703

G
godchen 已提交
3704
	if code != internalpb.StateCode_Healthy {
3705 3706 3707
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3708
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3709
				Reason:    "proxy not healthy",
3710 3711 3712
			},
		}, nil
	}
3713
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Inc()
3714 3715 3716
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3717
			ErrorCode: commonpb.ErrorCode_Success,
3718
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3719 3720 3721
		},
	}, nil
}
3722

3723 3724 3725
// 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",
3726
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3727 3728 3729 3730
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
3731
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3732
			zap.String("req", req.Request),
3733
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.ProxyID)))
3734 3735 3736 3737

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3738
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.ProxyID),
3739 3740 3741 3742 3743 3744 3745 3746
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
3747
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3748 3749 3750 3751 3752 3753 3754 3755 3756 3757 3758 3759 3760 3761 3762
			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 已提交
3763 3764 3765 3766 3767 3768 3769 3770 3771 3772
	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,
3773
		SourceID:  Params.ProxyCfg.ProxyID,
D
dragondriver 已提交
3774 3775
	}

3776
	if metricType == metricsinfo.SystemInfoMetrics {
3777 3778 3779 3780 3781 3782 3783
		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))

3784
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3785 3786

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

3793 3794
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3795
		return metrics, nil
3796 3797 3798
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
3799
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3800 3801 3802 3803 3804 3805 3806 3807 3808 3809 3810 3811
		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 已提交
3812 3813 3814
// 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",
3815
		zap.Int64("proxy_id", Params.ProxyCfg.ProxyID),
B
bigsheeper 已提交
3816 3817 3818 3819 3820 3821 3822 3823 3824 3825 3826 3827 3828 3829
		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,
3830
			SourceID:  Params.ProxyCfg.ProxyID,
B
bigsheeper 已提交
3831 3832 3833
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3834
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3835 3836 3837 3838 3839 3840 3841 3842 3843 3844 3845 3846 3847 3848 3849 3850 3851 3852
		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 已提交
3853
//GetCompactionState gets the compaction state of multiple segments
3854 3855 3856 3857 3858 3859 3860 3861 3862 3863 3864 3865 3866
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
}

3867
// ManualCompaction invokes compaction on specified collection
3868 3869 3870 3871 3872 3873 3874 3875 3876 3877 3878 3879 3880 3881 3882 3883 3884 3885 3886 3887 3888 3889 3890 3891 3892 3893
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 已提交
3894 3895 3896
// 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))
3897
	var err error
B
Bingyi Sun 已提交
3898 3899 3900 3901 3902 3903 3904
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3905
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3906 3907 3908 3909
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3910 3911 3912 3913
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3914 3915
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3916 3917 3918 3919
	code := node.stateCode.Load().(internalpb.StateCode)
	return code == internalpb.StateCode_Healthy
}

J
jingkl 已提交
3920
//unhealthyStatus returns the proxy not healthy status
3921 3922 3923
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3924
		Reason:    "proxy not healthy",
3925 3926
	}
}
G
groot 已提交
3927 3928 3929

// Import data files(json, numpy, etc.) on MinIO/S3 storage, read and parse them into sealed segments
func (node *Proxy) Import(ctx context.Context, req *milvuspb.ImportRequest) (*milvuspb.ImportResponse, error) {
G
groot 已提交
3930
	log.Info("received import request")
G
groot 已提交
3931 3932 3933 3934 3935 3936
	resp := &milvuspb.ImportResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

G
groot 已提交
3937 3938 3939
	resp, err := node.rootCoord.Import(ctx, req)
	log.Info("received import response", zap.String("collectionName", req.GetCollectionName()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
G
groot 已提交
3940 3941 3942 3943
}

// Check import task state from datanode
func (node *Proxy) GetImportState(ctx context.Context, req *milvuspb.GetImportStateRequest) (*milvuspb.GetImportStateResponse, error) {
G
groot 已提交
3944
	log.Info("received get import state request", zap.Int64("taskID", req.GetTask()))
G
groot 已提交
3945 3946 3947 3948 3949 3950
	resp := &milvuspb.GetImportStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

G
groot 已提交
3951 3952 3953
	resp, err := node.rootCoord.GetImportState(ctx, req)
	log.Info("received get import state response", zap.Int64("taskID", req.GetTask()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
G
groot 已提交
3954
}