impl.go 136.0 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 351
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName, 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 375
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, 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 405
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, 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 424 425 426
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName, metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
427 428 429
	return hct.result, nil
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

679 680 681 682 683
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dct.CollectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dct.CollectionID, 10), metrics.FailLabel).Inc()

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

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

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

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
721 722
	method := "GetCollectionStatistics"
	tr := timerecord.NewTimeRecorder(method)
723

724
	g := &getCollectionStatisticsTask{
G
godchen 已提交
725 726 727
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
728
		dataCoord:                      node.dataCoord,
729 730
	}

731 732
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
733
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
734 735
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
736 737 738 739 740

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

745 746 747
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

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

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

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

776 777 778 779 780
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(g.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(g.collectionID, 10), metrics.FailLabel).Inc()

G
godchen 已提交
781
		return &milvuspb.GetCollectionStatisticsResponse{
782
			Status: &commonpb.Status{
783
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
784 785 786 787 788
				Reason:    err.Error(),
			},
		}, nil
	}

789 790
	log.Debug("GetCollectionStatistics done",
		zap.String("traceID", traceID),
791
		zap.String("role", typeutil.ProxyRole),
792 793 794 795 796 797
		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))

798 799 800 801 802 803
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
804
	return g.result, nil
805 806
}

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

818
	sct := &showCollectionsTask{
G
godchen 已提交
819 820 821
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
822
		queryCoord:             node.queryCoord,
823
		rootCoord:              node.rootCoord,
824 825
	}

826
	log.Debug("ShowCollections received",
827
		zap.String("role", typeutil.ProxyRole),
828 829 830 831 832 833
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

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

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

854
	log.Debug("ShowCollections enqueued",
855
		zap.String("role", typeutil.ProxyRole),
856
		zap.Int64("MsgID", sct.ID()),
857
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
858
		zap.Uint64("TimeStamp", request.TimeStamp),
859 860 861
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
862

863 864
	err = sct.WaitToFinish()
	if err != nil {
865 866
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
867
			zap.String("role", typeutil.ProxyRole),
868 869 870 871 872 873 874
			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),
		)

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

G
godchen 已提交
877
		return &milvuspb.ShowCollectionsResponse{
878
			Status: &commonpb.Status{
879
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
880 881 882 883 884
				Reason:    err.Error(),
			},
		}, nil
	}

885
	log.Debug("ShowCollections Done",
886
		zap.String("role", typeutil.ProxyRole),
887 888 889 890 891 892 893 894
		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),
	)

895 896
	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()))
897 898 899
	return sct.result, nil
}

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

906
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
907 908
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
909 910 911
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
912

913
	cpt := &createPartitionTask{
S
sunby 已提交
914
		ctx:                    ctx,
915 916
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
917
		rootCoord:              node.rootCoord,
918 919 920
		result:                 nil,
	}

921 922 923
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
924
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
925 926 927
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
928 929 930 931 932 933

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
934
			zap.String("role", typeutil.ProxyRole),
935 936 937 938
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

941
		return &commonpb.Status{
942
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
943 944 945
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
946

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

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

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

973
		return &commonpb.Status{
974
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
975 976 977
			Reason:    err.Error(),
		}, nil
	}
978 979 980 981

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
982
		zap.String("role", typeutil.ProxyRole),
983 984 985 986 987 988 989
		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))

990 991
	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()))
992 993 994
	return cpt.result, nil
}

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

1001
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
1002 1003
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1004 1005 1006
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
1007

1008
	dpt := &dropPartitionTask{
S
sunby 已提交
1009
		ctx:                  ctx,
1010 1011
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1012
		rootCoord:            node.rootCoord,
1013 1014 1015
		result:               nil,
	}

1016 1017 1018
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1019
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1020 1021 1022
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1023 1024 1025 1026 1027 1028

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1029
			zap.String("role", typeutil.ProxyRole),
1030 1031 1032 1033
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1036
		return &commonpb.Status{
1037
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1038 1039 1040
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1041

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

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

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

1068
		return &commonpb.Status{
1069
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1070 1071 1072
			Reason:    err.Error(),
		}, nil
	}
1073 1074 1075 1076

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1077
		zap.String("role", typeutil.ProxyRole),
1078 1079 1080 1081 1082 1083 1084
		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))

1085 1086
	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()))
1087 1088 1089
	return dpt.result, nil
}

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

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

1107
	hpt := &hasPartitionTask{
S
sunby 已提交
1108
		ctx:                 ctx,
1109 1110
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1111
		rootCoord:           node.rootCoord,
1112 1113 1114
		result:              nil,
	}

D
dragondriver 已提交
1115 1116 1117
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1118
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1119 1120 1121
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1122 1123 1124 1125 1126 1127

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

1133 1134 1135
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1136 1137
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1138
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1139 1140 1141 1142 1143
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1144

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

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

1169 1170 1171
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.FailLabel).Inc()

1172 1173
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1174
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1175 1176 1177 1178 1179
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1180 1181 1182 1183

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1184
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1185 1186 1187 1188 1189 1190 1191
		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))

1192 1193 1194 1195
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName, metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
1196 1197 1198
	return hpt.result, nil
}

1199
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1200
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1201 1202 1203
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1204

D
dragondriver 已提交
1205
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1206 1207
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1208 1209
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1210

1211
	lpt := &loadPartitionsTask{
G
godchen 已提交
1212 1213 1214
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1215
		queryCoord:            node.queryCoord,
1216 1217
	}

1218 1219 1220
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1221
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1222 1223 1224
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1225 1226 1227 1228 1229 1230

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1231
			zap.String("role", typeutil.ProxyRole),
1232 1233 1234 1235
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1236 1237 1238
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1239
		return &commonpb.Status{
1240
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1241 1242 1243 1244
			Reason:    err.Error(),
		}, nil
	}

1245 1246 1247
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1248
		zap.String("role", typeutil.ProxyRole),
1249 1250 1251
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1252 1253 1254
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1255 1256 1257 1258

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

1269 1270 1271 1272 1273
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(lpt.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(lpt.collectionID, 10), metrics.FailLabel).Inc()

1274
		return &commonpb.Status{
1275
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1276 1277 1278 1279
			Reason:    err.Error(),
		}, nil
	}

1280 1281 1282
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1283
		zap.String("role", typeutil.ProxyRole),
1284 1285 1286 1287 1288 1289 1290
		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))

1291 1292 1293 1294 1295 1296
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(lpt.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(lpt.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(lpt.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
1297
	return lpt.result, nil
1298 1299
}

1300
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1301
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1302 1303 1304
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1305 1306 1307 1308 1309

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

1310
	rpt := &releasePartitionsTask{
G
godchen 已提交
1311 1312 1313
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1314
		queryCoord:               node.queryCoord,
1315 1316
	}

1317
	method := "ReleasePartitions"
1318
	tr := timerecord.NewTimeRecorder(method)
1319 1320 1321 1322

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1323
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1324 1325 1326
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1327 1328 1329 1330 1331 1332

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1333
			zap.String("role", typeutil.ProxyRole),
1334 1335 1336 1337
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1338 1339 1340
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1341
		return &commonpb.Status{
1342
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1343 1344 1345 1346
			Reason:    err.Error(),
		}, nil
	}

1347 1348 1349
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1350
		zap.String("role", typeutil.ProxyRole),
1351 1352 1353
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1354 1355 1356
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1357 1358 1359 1360

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1361
			zap.Error(err),
1362
			zap.String("traceID", traceID),
1363
			zap.String("role", typeutil.ProxyRole),
1364 1365 1366
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1367 1368 1369 1370
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1371 1372 1373 1374 1375
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(rpt.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(rpt.collectionID, 10), metrics.FailLabel).Inc()

1376
		return &commonpb.Status{
1377
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1378 1379 1380 1381
			Reason:    err.Error(),
		}, nil
	}

1382 1383 1384
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1385
		zap.String("role", typeutil.ProxyRole),
1386 1387 1388 1389 1390 1391 1392
		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))

1393 1394 1395 1396 1397 1398
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(rpt.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(rpt.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(rpt.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
1399
	return rpt.result, nil
1400 1401
}

1402
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1403
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1404 1405 1406 1407 1408
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1409 1410 1411 1412 1413

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

1414
	g := &getPartitionStatisticsTask{
1415 1416 1417
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1418
		dataCoord:                     node.dataCoord,
1419 1420
	}

1421
	method := "GetPartitionStatistics"
1422
	tr := timerecord.NewTimeRecorder(method)
1423 1424 1425 1426

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1427
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1428 1429 1430
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1431 1432 1433 1434 1435 1436

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1437
			zap.String("role", typeutil.ProxyRole),
1438 1439 1440 1441
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1442 1443 1444
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1445 1446 1447 1448 1449 1450 1451 1452
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1453 1454 1455
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1456
		zap.String("role", typeutil.ProxyRole),
1457 1458 1459
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1460 1461 1462
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1463 1464 1465 1466

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1467
			zap.Error(err),
1468
			zap.String("traceID", traceID),
1469
			zap.String("role", typeutil.ProxyRole),
1470 1471 1472
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1473 1474 1475 1476
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1477 1478 1479 1480 1481
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(g.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(g.collectionID, 10), metrics.FailLabel).Inc()

1482 1483 1484 1485 1486 1487 1488 1489
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1490 1491 1492
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1493
		zap.String("role", typeutil.ProxyRole),
1494 1495 1496 1497 1498 1499 1500
		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))

1501 1502 1503 1504 1505 1506
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(g.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
1507
	return g.result, nil
1508 1509
}

1510
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1511
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1512 1513 1514 1515 1516
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1517 1518 1519 1520 1521

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

1522
	spt := &showPartitionsTask{
G
godchen 已提交
1523 1524 1525
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1526
		rootCoord:             node.rootCoord,
1527
		queryCoord:            node.queryCoord,
G
godchen 已提交
1528
		result:                nil,
1529 1530
	}

1531
	method := "ShowPartitions"
1532 1533 1534 1535
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName, metrics.TotalLabel).Inc()
1536 1537 1538 1539

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1540
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1541
		zap.Any("request", request))
1542 1543 1544 1545 1546 1547

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

1551 1552 1553
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

G
godchen 已提交
1554
		return &milvuspb.ShowPartitionsResponse{
1555
			Status: &commonpb.Status{
1556
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1557 1558 1559 1560 1561
				Reason:    err.Error(),
			},
		}, nil
	}

1562 1563 1564
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1565
		zap.String("role", typeutil.ProxyRole),
1566 1567 1568
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1569 1570
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1571 1572 1573 1574 1575
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1576
			zap.Error(err),
1577
			zap.String("traceID", traceID),
1578
			zap.String("role", typeutil.ProxyRole),
1579 1580 1581 1582 1583 1584
			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 已提交
1585

1586 1587 1588
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.FailLabel).Inc()

G
godchen 已提交
1589
		return &milvuspb.ShowPartitionsResponse{
1590
			Status: &commonpb.Status{
1591
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1592 1593 1594 1595
				Reason:    err.Error(),
			},
		}, nil
	}
1596 1597 1598 1599

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1600
		zap.String("role", typeutil.ProxyRole),
1601 1602 1603 1604 1605 1606 1607
		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))

1608 1609 1610 1611
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName, metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
1612 1613 1614
	return spt.result, nil
}

1615
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1616
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1617 1618 1619
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1620 1621 1622 1623 1624

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

1625
	cit := &createIndexTask{
S
sunby 已提交
1626
		ctx:                ctx,
1627 1628
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1629
		rootCoord:          node.rootCoord,
1630 1631
	}

D
dragondriver 已提交
1632
	method := "CreateIndex"
1633
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1634 1635 1636 1637

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1638
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1639 1640 1641 1642
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1643 1644 1645 1646 1647 1648

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1649
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1650 1651 1652 1653 1654
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1655 1656 1657
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1658
		return &commonpb.Status{
1659
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1660 1661 1662 1663
			Reason:    err.Error(),
		}, nil
	}

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

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

1690 1691 1692 1693 1694
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(cit.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(cit.collectionID, 10), metrics.FailLabel).Inc()

1695
		return &commonpb.Status{
1696
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1697 1698 1699 1700
			Reason:    err.Error(),
		}, nil
	}

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

1713 1714 1715 1716 1717 1718
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(cit.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(cit.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(cit.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
1719 1720 1721
	return cit.result, nil
}

1722
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
1723
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
1724 1725 1726 1727 1728
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1729 1730 1731 1732 1733

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

1734
	dit := &describeIndexTask{
S
sunby 已提交
1735
		ctx:                  ctx,
1736 1737
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1738
		rootCoord:            node.rootCoord,
1739 1740
	}

1741 1742 1743
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
1744
	tr := timerecord.NewTimeRecorder(method)
1745 1746 1747 1748

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1749
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1750 1751 1752
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1753 1754 1755 1756 1757 1758 1759
		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),
1760
			zap.String("role", typeutil.ProxyRole),
1761 1762 1763 1764 1765
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

1766 1767 1768
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

1769 1770
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1771
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1772 1773 1774 1775 1776
				Reason:    err.Error(),
			},
		}, nil
	}

1777 1778 1779
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1780
		zap.String("role", typeutil.ProxyRole),
1781 1782 1783
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1784 1785 1786
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1787 1788 1789 1790 1791
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1792
			zap.Error(err),
1793
			zap.String("traceID", traceID),
1794
			zap.String("role", typeutil.ProxyRole),
1795 1796 1797
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1798 1799 1800
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
1801
			zap.String("index name", indexName))
D
dragondriver 已提交
1802

Z
zhenshan.cao 已提交
1803 1804 1805 1806
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
1807 1808 1809 1810 1811
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dit.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dit.collectionID, 10), metrics.FailLabel).Inc()

1812 1813
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
1814
				ErrorCode: errCode,
1815 1816 1817 1818 1819
				Reason:    err.Error(),
			},
		}, nil
	}

1820 1821 1822
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1823
		zap.String("role", typeutil.ProxyRole),
1824 1825 1826 1827 1828 1829 1830 1831
		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))

1832 1833 1834 1835 1836 1837
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
1838 1839 1840
	return dit.result, nil
}

1841
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
1842
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
1843 1844 1845
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1846 1847 1848 1849 1850

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

1851
	dit := &dropIndexTask{
S
sunby 已提交
1852
		ctx:              ctx,
B
BossZou 已提交
1853 1854
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1855
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1856
	}
G
godchen 已提交
1857

D
dragondriver 已提交
1858
	method := "DropIndex"
1859
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1860 1861 1862 1863

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1864
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1865 1866 1867 1868 1869
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
1870 1871 1872 1873 1874
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1875
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1876 1877 1878 1879
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1880 1881
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
1882

B
BossZou 已提交
1883
		return &commonpb.Status{
1884
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1885 1886 1887
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1888

D
dragondriver 已提交
1889 1890 1891
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1892
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1893 1894 1895
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1896 1897 1898 1899
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
1900 1901 1902 1903

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1904
			zap.Error(err),
D
dragondriver 已提交
1905
			zap.String("traceID", traceID),
1906
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1907 1908 1909
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1910 1911 1912 1913 1914
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

1915 1916 1917 1918 1919
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dit.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dit.collectionID, 10), metrics.FailLabel).Inc()

B
BossZou 已提交
1920
		return &commonpb.Status{
1921
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1922 1923 1924
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1925 1926 1927 1928

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1929
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1930 1931 1932 1933 1934 1935 1936 1937
		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))

1938 1939 1940 1941 1942 1943
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dit.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
1944 1945 1946
	return dit.result, nil
}

1947 1948
// 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 已提交
1949
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1950 1951 1952 1953 1954
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1955 1956 1957 1958 1959

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

1960
	gibpt := &getIndexBuildProgressTask{
1961 1962 1963
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1964 1965
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1966
		dataCoord:                    node.dataCoord,
1967 1968
	}

1969
	method := "GetIndexBuildProgress"
1970
	tr := timerecord.NewTimeRecorder(method)
1971 1972 1973 1974

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1975
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1976 1977 1978 1979
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1980 1981 1982 1983 1984 1985

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1986
			zap.String("role", typeutil.ProxyRole),
1987 1988 1989 1990
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1991 1992
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()
1993

1994 1995 1996 1997 1998 1999 2000 2001
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2002 2003 2004
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2005
		zap.String("role", typeutil.ProxyRole),
2006 2007 2008
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
2009 2010 2011 2012
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2013 2014 2015 2016

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2017
			zap.Error(err),
2018
			zap.String("traceID", traceID),
2019
			zap.String("role", typeutil.ProxyRole),
2020 2021 2022
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2023 2024 2025 2026
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2027 2028 2029 2030
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(gibpt.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(gibpt.collectionID, 10), metrics.FailLabel).Inc()
2031 2032 2033 2034 2035 2036 2037 2038

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2039 2040 2041 2042

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2043
		zap.String("role", typeutil.ProxyRole),
2044 2045 2046 2047 2048 2049 2050 2051
		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))
2052

2053 2054 2055 2056 2057 2058
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(gibpt.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(gibpt.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(gibpt.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
2059
	return gibpt.result, nil
2060 2061
}

2062
// GetIndexState get the build-state of index.
C
Cai Yudong 已提交
2063
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2064 2065 2066 2067 2068
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2069 2070 2071 2072 2073

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

2074
	dipt := &getIndexStateTask{
G
godchen 已提交
2075 2076 2077
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2078 2079
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2080 2081
	}

2082
	method := "GetIndexState"
2083
	tr := timerecord.NewTimeRecorder(method)
2084 2085 2086 2087

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2088
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2089 2090 2091 2092
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2093 2094 2095 2096 2097 2098

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2099
			zap.String("role", typeutil.ProxyRole),
2100 2101 2102 2103 2104
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2105 2106 2107
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.AbandonLabel).Inc()

G
godchen 已提交
2108
		return &milvuspb.GetIndexStateResponse{
2109
			Status: &commonpb.Status{
2110
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2111 2112 2113 2114 2115
				Reason:    err.Error(),
			},
		}, nil
	}

2116 2117 2118
	log.Debug(
		rpcEnqueued(method),
		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
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2127 2128 2129 2130

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2131
			zap.Error(err),
2132
			zap.String("traceID", traceID),
2133
			zap.String("role", typeutil.ProxyRole),
2134 2135 2136
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2137 2138 2139 2140 2141
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2142 2143 2144 2145 2146
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dipt.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dipt.collectionID, 10), metrics.FailLabel).Inc()

G
godchen 已提交
2147
		return &milvuspb.GetIndexStateResponse{
2148
			Status: &commonpb.Status{
2149
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2150 2151 2152 2153 2154
				Reason:    err.Error(),
			},
		}, nil
	}

2155 2156 2157
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2158
		zap.String("role", typeutil.ProxyRole),
2159 2160 2161 2162 2163 2164 2165 2166
		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))

2167 2168 2169 2170 2171 2172
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dipt.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dipt.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dipt.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
2173 2174 2175
	return dipt.result, nil
}

2176
// Insert insert records into collection.
C
Cai Yudong 已提交
2177
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2178 2179 2180 2181 2182 2183
	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))

2184 2185 2186 2187 2188
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2189 2190
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2191

2192
	it := &insertTask{
2193 2194 2195
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       request,
2196 2197 2198 2199
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2200
			InsertRequest: internalpb.InsertRequest{
2201
				Base: &commonpb.MsgBase{
2202
					MsgType: commonpb.MsgType_Insert,
2203 2204 2205 2206
					MsgID:   0,
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
2207
				// RowData: transfer column based request to this
2208 2209
			},
		},
2210
		rowIDAllocator: node.idAllocator,
2211
		segIDAssigner:  node.segAssigner,
2212
		chMgr:          node.chMgr,
2213
		chTicker:       node.chTicker,
2214
	}
2215 2216

	if len(it.PartitionName) <= 0 {
2217
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2218 2219
	}

X
Xiangyu Wang 已提交
2220
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
2221 2222 2223 2224 2225
		numRows := it.req.NumRows
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2226

X
Xiangyu Wang 已提交
2227 2228 2229 2230 2231 2232 2233
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2234 2235
	}

X
Xiangyu Wang 已提交
2236
	log.Debug("Enqueue insert request in Proxy",
2237
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2238 2239 2240 2241 2242
		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)),
2243 2244
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2245

X
Xiangyu Wang 已提交
2246 2247
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2248
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), request.CollectionName, metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2249
		return constructFailedResponse(err), nil
2250
	}
D
dragondriver 已提交
2251

X
Xiangyu Wang 已提交
2252
	log.Debug("Detail of insert request in Proxy",
2253
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2254
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2255 2256 2257 2258 2259
		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 已提交
2260 2261 2262 2263 2264
		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))
2265 2266 2267 2268
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			strconv.FormatInt(it.CollectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			strconv.FormatInt(it.CollectionID, 10), metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2269 2270 2271 2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286
		return constructFailedResponse(err), nil
	}

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

		setErrorIndex()
	}

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

2288 2289 2290 2291 2292 2293
	metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(it.CollectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(it.CollectionID, 10)).Add(float64(it.result.InsertCnt))
	metrics.ProxyInsertLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(it.CollectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
2294 2295 2296
	return it.result, nil
}

2297
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2298
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2299 2300 2301
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2302 2303
	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))
2304

G
groot 已提交
2305 2306 2307 2308 2309 2310
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2311 2312 2313
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

C
Cai Yudong 已提交
2314 2315 2316 2317 2318 2319 2320
	deleteReq := &milvuspb.DeleteRequest{
		DbName:         request.DbName,
		CollectionName: request.CollectionName,
		PartitionName:  request.PartitionName,
		Expr:           request.Expr,
	}

2321
	dt := &deleteTask{
C
Cai Yudong 已提交
2322 2323 2324
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		req:       deleteReq,
G
godchen 已提交
2325
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2326 2327 2328
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2329 2330 2331 2332 2333 2334 2335 2336
			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 已提交
2337 2338 2339 2340
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2341 2342
	}

2343
	log.Debug("Enqueue delete request in Proxy",
2344
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2345 2346 2347 2348
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2349 2350 2351 2352

	// 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))
2353 2354 2355
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			request.CollectionName, metrics.FailLabel).Inc()

G
groot 已提交
2356 2357 2358 2359 2360 2361 2362 2363
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2364
	log.Debug("Detail of delete request in Proxy",
2365
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2366 2367 2368 2369 2370
		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),
2371 2372
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2373

2374 2375
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
2376 2377 2378 2379
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dt.collectionID, 10), metrics.TotalLabel).Inc()
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
			strconv.FormatInt(dt.collectionID, 10), metrics.FailLabel).Inc()
G
groot 已提交
2380 2381 2382 2383 2384 2385 2386 2387
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2388 2389 2390 2391 2392 2393
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dt.collectionID, 10), metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dt.collectionID, 10), metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
		strconv.FormatInt(dt.collectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2394 2395 2396
	return dt.result, nil
}

2397
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2398
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2399 2400 2401 2402 2403
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2404 2405
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2406

C
cai.zhang 已提交
2407 2408
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2409 2410
	traceID, _, _ := trace.InfoFromSpan(sp)

2411
	qt := &searchTask{
S
sunby 已提交
2412
		ctx:       ctx,
2413
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2414
		SearchRequest: &internalpb.SearchRequest{
2415
			Base: &commonpb.MsgBase{
2416
				MsgType:  commonpb.MsgType_Search,
2417
				SourceID: Params.ProxyCfg.ProxyID,
2418
			},
2419
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2420
		},
2421
		resultBuf: make(chan []*internalpb.SearchResults, 1),
2422 2423
		query:     request,
		chMgr:     node.chMgr,
2424
		qc:        node.queryCoord,
2425
		tr:        timerecord.NewTimeRecorder("search"),
2426 2427
	}

2428 2429 2430 2431 2432
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

	log.Debug(
		rpcReceived(method),
D
dragondriver 已提交
2433
		zap.String("traceID", traceID),
2434
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2435 2436 2437 2438 2439
		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)),
2440 2441 2442 2443
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2444

2445 2446 2447
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2448 2449
			zap.Error(err),
			zap.String("traceID", traceID),
2450
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2451 2452 2453 2454 2455 2456
			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),
2457 2458 2459
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2460

2461 2462 2463
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			request.CollectionName, metrics.SearchLabel, metrics.AbandonLabel).Inc()

2464 2465
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2466
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2467 2468 2469 2470 2471
				Reason:    err.Error(),
			},
		}, nil
	}

2472 2473
	log.Debug(
		rpcEnqueued(method),
D
dragondriver 已提交
2474
		zap.String("traceID", traceID),
2475
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2476
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2477 2478 2479 2480 2481
		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),
2482
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2483 2484 2485 2486
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2487

2488 2489 2490
	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2491
			zap.Error(err),
D
dragondriver 已提交
2492
			zap.String("traceID", traceID),
2493
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2494
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2495 2496 2497 2498
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2499
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2500 2501 2502 2503
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2504 2505 2506 2507
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel, metrics.TotalLabel).Inc()

		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), request.CollectionName, metrics.SearchLabel, metrics.FailLabel).Inc()
2508

2509 2510
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2511
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2512 2513 2514 2515 2516
				Reason:    err.Error(),
			},
		}, nil
	}

2517 2518
	log.Debug(
		rpcDone(method),
D
dragondriver 已提交
2519
		zap.String("traceID", traceID),
2520
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2521 2522 2523 2524 2525 2526
		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)),
2527 2528 2529 2530
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2531

2532 2533 2534 2535 2536 2537
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel, metrics.TotalLabel).Inc()
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel, metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel).Set(float64(qt.result.Results.NumQueries))
C
cai.zhang 已提交
2538
	searchDur := tr.ElapseSpan().Milliseconds()
2539
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
C
cai.zhang 已提交
2540
		strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel).Observe(float64(searchDur))
2541
	metrics.ProxySearchLatencyPerNQ.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
C
cai.zhang 已提交
2542
		strconv.FormatInt(qt.CollectionID, 10)).Observe(float64(searchDur) / float64(qt.result.Results.NumQueries))
2543 2544 2545
	return qt.result, nil
}

2546
// Flush notify data nodes to persist the data of collection.
2547 2548 2549 2550 2551 2552 2553
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2554
	if !node.checkHealthy() {
2555 2556
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2557
	}
D
dragondriver 已提交
2558 2559 2560 2561 2562

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

2563
	ft := &flushTask{
T
ThreadDao 已提交
2564 2565 2566
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2567
		dataCoord:    node.dataCoord,
2568 2569
	}

D
dragondriver 已提交
2570
	method := "Flush"
2571 2572
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2573 2574 2575 2576

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2577
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2578 2579
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2580 2581 2582 2583 2584 2585

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

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

2592 2593
		resp.Status.Reason = err.Error()
		return resp, nil
2594 2595
	}

D
dragondriver 已提交
2596 2597 2598
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2599
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2600 2601 2602
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2603 2604
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2605 2606 2607 2608

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2609
			zap.Error(err),
D
dragondriver 已提交
2610
			zap.String("traceID", traceID),
2611
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2612 2613 2614
			zap.Int64("MsgID", ft.ID()),
			zap.Uint64("BeginTs", ft.BeginTs()),
			zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2615 2616 2617
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

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

D
dragondriver 已提交
2620
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2621 2622
		resp.Status.Reason = err.Error()
		return resp, nil
2623 2624
	}

D
dragondriver 已提交
2625 2626 2627
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2628
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2629 2630 2631 2632 2633 2634
		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))

2635 2636
	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()))
2637
	return ft.result, nil
2638 2639
}

2640
// Query get the records by primary keys.
C
Cai Yudong 已提交
2641
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2642 2643 2644 2645 2646
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2647

D
dragondriver 已提交
2648 2649 2650
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2651
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2652

2653
	qt := &queryTask{
2654 2655 2656 2657 2658
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
2659
				SourceID: Params.ProxyCfg.ProxyID,
2660
			},
2661
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2662 2663
		},
		resultBuf: make(chan []*internalpb.RetrieveResults),
2664
		query:     request,
2665 2666
		chMgr:     node.chMgr,
		qc:        node.queryCoord,
2667 2668
	}

D
dragondriver 已提交
2669 2670 2671 2672 2673
	method := "Query"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2674
		zap.String("role", typeutil.ProxyRole),
2675 2676 2677
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
G
godchen 已提交
2678

D
dragondriver 已提交
2679 2680 2681 2682 2683 2684
	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),
2685 2686 2687
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2688

2689 2690
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			request.CollectionName, metrics.QueryLabel, metrics.FailLabel).Inc()
2691 2692 2693 2694 2695 2696
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2697 2698
	}

D
dragondriver 已提交
2699 2700 2701
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2702
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2703 2704 2705
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
2706 2707 2708
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2709 2710 2711 2712 2713 2714

	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2715
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2716 2717 2718
			zap.Int64("MsgID", qt.ID()),
			zap.Uint64("BeginTs", qt.BeginTs()),
			zap.Uint64("EndTs", qt.EndTs()),
2719 2720 2721
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
2722

2723 2724 2725 2726 2727
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel, metrics.TotalLabel).Inc()
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel, metrics.FailLabel).Inc()

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

D
dragondriver 已提交
2736 2737 2738 2739 2740 2741 2742
	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()),
2743 2744 2745
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2746

2747 2748 2749 2750 2751 2752 2753 2754
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel, metrics.TotalLabel).Inc()
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel, metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel).Set(float64(len(qt.result.FieldsData)))
	metrics.ProxySendMessageLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2755 2756 2757 2758 2759
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2760

2761
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
2762 2763 2764 2765
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2766 2767 2768 2769 2770

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

Y
Yusup 已提交
2771 2772 2773 2774 2775 2776 2777
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
2778
	method := "CreateAlias"
2779 2780
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2781 2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792 2793 2794 2795 2796 2797 2798 2799

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

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

Y
Yusup 已提交
2802 2803 2804 2805 2806 2807
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2808 2809 2810
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2811
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2812 2813 2814 2815
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2816 2817
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2818 2819 2820 2821

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2822
			zap.Error(err),
D
dragondriver 已提交
2823
			zap.String("traceID", traceID),
2824
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2825 2826 2827 2828
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2829 2830
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
2831
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
2832 2833 2834 2835 2836 2837 2838

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

D
dragondriver 已提交
2839 2840 2841 2842 2843 2844 2845 2846 2847 2848 2849
	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))

2850 2851
	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 已提交
2852 2853 2854
	return cat.result, nil
}

2855
// DropAlias alter the alias of collection.
Y
Yusup 已提交
2856 2857 2858 2859
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2860 2861 2862 2863 2864

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

Y
Yusup 已提交
2865 2866 2867 2868 2869 2870 2871
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
2872
	method := "DropAlias"
2873 2874
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2875 2876 2877 2878 2879 2880 2881 2882 2883 2884 2885 2886 2887 2888 2889 2890

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

Y
Yusup 已提交
2893 2894 2895 2896 2897 2898
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2899 2900 2901
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2902
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2903 2904 2905 2906
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2907
		zap.String("alias", request.Alias))
D
dragondriver 已提交
2908 2909 2910 2911

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2912
			zap.Error(err),
D
dragondriver 已提交
2913
			zap.String("traceID", traceID),
2914
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2915 2916 2917 2918
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2919 2920
			zap.String("alias", request.Alias))

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

Y
Yusup 已提交
2923 2924 2925 2926 2927 2928
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2929 2930 2931 2932 2933 2934 2935 2936 2937 2938
	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))

2939 2940
	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 已提交
2941 2942 2943
	return dat.result, nil
}

2944
// AlterAlias alter alias of collection.
Y
Yusup 已提交
2945 2946 2947 2948
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2949 2950 2951 2952 2953

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

Y
Yusup 已提交
2954 2955 2956 2957 2958 2959 2960
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
2961
	method := "AlterAlias"
2962 2963
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2964 2965 2966 2967 2968 2969 2970 2971 2972 2973 2974 2975 2976 2977 2978 2979 2980 2981

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

Y
Yusup 已提交
2984 2985 2986 2987 2988 2989
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2990 2991 2992
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2993
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2994 2995 2996 2997
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2998 2999
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3000 3001 3002 3003

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3004
			zap.Error(err),
D
dragondriver 已提交
3005
			zap.String("traceID", traceID),
3006
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3007 3008 3009 3010
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3011 3012 3013
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

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

Y
Yusup 已提交
3016 3017 3018 3019 3020 3021
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3022 3023 3024 3025 3026 3027 3028 3029 3030 3031 3032
	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))

3033 3034
	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 已提交
3035 3036 3037
	return aat.result, nil
}

3038
// CalcDistance calculates the distances between vectors.
3039
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3040 3041 3042 3043 3044
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3045
	param, _ := funcutil.GetAttrByKeyFromRepeatedKV("metric", request.GetParams())
3046 3047 3048 3049 3050 3051 3052 3053
	metric, err := distance.ValidateMetricType(param)
	if err != nil {
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
3054 3055
	}

3056 3057 3058 3059
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3060 3061
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3062

3063 3064 3065 3066 3067
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3068 3069
		}

3070
		qt := &queryTask{
3071 3072 3073 3074 3075
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
3076
					SourceID: Params.ProxyCfg.ProxyID,
3077
				},
3078
				ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
3079
			},
3080
			resultBuf: make(chan []*internalpb.RetrieveResults),
3081
			query:     queryRequest,
3082
			chMgr:     node.chMgr,
3083
			qc:        node.queryCoord,
Y
yukun 已提交
3084
			ids:       ids.IdArray,
3085 3086
		}

3087
		err := node.sched.dqQueue.Enqueue(qt)
3088
		if err != nil {
3089 3090 3091
			log.Debug("CalcDistance queryTask failed to enqueue",
				zap.Error(err),
				zap.String("traceID", traceID),
3092
				zap.String("role", typeutil.ProxyRole),
3093 3094 3095 3096
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames))

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

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
3107
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3108 3109 3110 3111 3112 3113
			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))
3114 3115 3116 3117

		err = qt.WaitToFinish()
		if err != nil {
			log.Debug("CalcDistance queryTask failed to WaitToFinish",
G
godchen 已提交
3118
				zap.Error(err),
3119
				zap.String("traceID", traceID),
3120
				zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3121 3122 3123 3124 3125 3126
				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))
3127 3128 3129 3130 3131 3132

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3133
			}, err
3134
		}
3135 3136 3137

		log.Debug("CalcDistance queryTask Done",
			zap.String("traceID", traceID),
3138
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3139 3140 3141 3142 3143 3144
			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))
3145 3146

		return &milvuspb.QueryResults{
3147 3148
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3149 3150 3151
		}, nil
	}

3152 3153 3154 3155 3156 3157 3158 3159 3160 3161 3162 3163 3164 3165
	// 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 {
3166
			return nil, errors.New("failed to fetch vectors")
3167 3168 3169 3170 3171 3172 3173 3174 3175 3176 3177 3178 3179 3180 3181 3182
		}

		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))
3183
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3184 3185 3186 3187 3188 3189 3190 3191 3192 3193 3194 3195 3196 3197 3198 3199 3200 3201 3202 3203 3204 3205 3206 3207 3208 3209
				}
				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))
3210
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3211 3212 3213 3214 3215 3216 3217 3218 3219 3220 3221 3222
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

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

3223
		return nil, errors.New("failed to fetch vectors")
3224 3225
	}

3226 3227
	log.Debug("CalcDistance received",
		zap.String("traceID", traceID),
3228
		zap.String("role", typeutil.ProxyRole),
3229
		zap.String("metric", metric))
G
godchen 已提交
3230

3231 3232 3233
	vectorsLeft := request.GetOpLeft().GetDataArray()
	opLeft := request.GetOpLeft().GetIdArray()
	if opLeft != nil {
3234 3235
		log.Debug("OpLeft IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3236
			zap.String("role", typeutil.ProxyRole))
3237

3238
		result, err := query(opLeft)
3239
		if err != nil {
3240 3241 3242
			log.Debug("Failed to get left vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3243
				zap.String("role", typeutil.ProxyRole))
3244

3245 3246 3247 3248 3249 3250 3251 3252
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3253 3254
		log.Debug("OpLeft IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3255
			zap.String("role", typeutil.ProxyRole))
3256

3257 3258
		vectorsLeft, err = arrangeFunc(opLeft, result.FieldsData)
		if err != nil {
3259 3260 3261
			log.Debug("Failed to re-arrange left vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3262
				zap.String("role", typeutil.ProxyRole))
3263

3264 3265 3266 3267 3268 3269
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3270
		}
3271 3272 3273

		log.Debug("Re-arrange left vectors done",
			zap.String("traceID", traceID),
3274
			zap.String("role", typeutil.ProxyRole))
3275 3276
	}

G
groot 已提交
3277
	if vectorsLeft == nil {
3278 3279 3280
		msg := "Left vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3281
			zap.String("role", typeutil.ProxyRole))
3282

G
groot 已提交
3283 3284 3285
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3286
				Reason:    msg,
G
groot 已提交
3287 3288 3289 3290
			},
		}, nil
	}

3291 3292 3293
	vectorsRight := request.GetOpRight().GetDataArray()
	opRight := request.GetOpRight().GetIdArray()
	if opRight != nil {
3294 3295
		log.Debug("OpRight IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3296
			zap.String("role", typeutil.ProxyRole))
3297

3298
		result, err := query(opRight)
3299
		if err != nil {
3300 3301 3302
			log.Debug("Failed to get right vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3303
				zap.String("role", typeutil.ProxyRole))
3304

3305 3306 3307 3308 3309 3310 3311 3312
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3313 3314
		log.Debug("OpRight IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3315
			zap.String("role", typeutil.ProxyRole))
3316

3317 3318
		vectorsRight, err = arrangeFunc(opRight, result.FieldsData)
		if err != nil {
3319 3320 3321
			log.Debug("Failed to re-arrange right vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3322
				zap.String("role", typeutil.ProxyRole))
3323

3324 3325 3326 3327 3328 3329
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3330
		}
3331 3332 3333

		log.Debug("Re-arrange right vectors done",
			zap.String("traceID", traceID),
3334
			zap.String("role", typeutil.ProxyRole))
3335 3336
	}

G
groot 已提交
3337
	if vectorsRight == nil {
3338 3339 3340
		msg := "Right vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3341
			zap.String("role", typeutil.ProxyRole))
3342

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

3351
	if vectorsLeft.Dim != vectorsRight.Dim {
3352 3353 3354
		msg := "Vectors dimension is not equal"
		log.Debug(msg,
			zap.String("traceID", traceID),
3355
			zap.String("role", typeutil.ProxyRole))
3356

3357 3358 3359
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3360
				Reason:    msg,
3361 3362 3363 3364 3365
			},
		}, nil
	}

	if vectorsLeft.GetFloatVector() != nil && vectorsRight.GetFloatVector() != nil {
3366 3367
		distances, err := distance.CalcFloatDistance(vectorsLeft.Dim, vectorsLeft.GetFloatVector().Data, vectorsRight.GetFloatVector().Data, metric)
		if err != nil {
3368 3369 3370
			log.Debug("Failed to CalcFloatDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3371
				zap.String("role", typeutil.ProxyRole))
3372

3373 3374 3375 3376 3377 3378 3379 3380
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3381 3382 3383
		log.Debug("CalcFloatDistance done",
			zap.Error(err),
			zap.String("traceID", traceID),
3384
			zap.String("role", typeutil.ProxyRole))
3385

3386 3387 3388 3389 3390 3391 3392 3393 3394 3395
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
			Array: &milvuspb.CalcDistanceResults_FloatDist{
				FloatDist: &schemapb.FloatArray{
					Data: distances,
				},
			},
		}, nil
	}

3396
	if vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetBinaryVector() != nil {
G
groot 已提交
3397
		hamming, err := distance.CalcHammingDistance(vectorsLeft.Dim, vectorsLeft.GetBinaryVector(), vectorsRight.GetBinaryVector())
3398
		if err != nil {
3399 3400 3401
			log.Debug("Failed to CalcHammingDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3402
				zap.String("role", typeutil.ProxyRole))
3403

3404 3405 3406 3407 3408 3409 3410 3411 3412
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

		if metric == distance.HAMMING {
3413 3414
			log.Debug("CalcHammingDistance done",
				zap.String("traceID", traceID),
3415
				zap.String("role", typeutil.ProxyRole))
3416

3417 3418 3419 3420
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_IntDist{
					IntDist: &schemapb.IntArray{
G
groot 已提交
3421
						Data: hamming,
3422 3423 3424 3425 3426 3427
					},
				},
			}, nil
		}

		if metric == distance.TANIMOTO {
G
groot 已提交
3428
			tanimoto, err := distance.CalcTanimotoCoefficient(vectorsLeft.Dim, hamming)
3429
			if err != nil {
3430 3431 3432
				log.Debug("Failed to CalcTanimotoCoefficient",
					zap.Error(err),
					zap.String("traceID", traceID),
3433
					zap.String("role", typeutil.ProxyRole))
3434

3435 3436 3437 3438 3439 3440 3441 3442
				return &milvuspb.CalcDistanceResults{
					Status: &commonpb.Status{
						ErrorCode: commonpb.ErrorCode_UnexpectedError,
						Reason:    err.Error(),
					},
				}, nil
			}

3443 3444
			log.Debug("CalcTanimotoCoefficient done",
				zap.String("traceID", traceID),
3445
				zap.String("role", typeutil.ProxyRole))
3446

3447 3448 3449 3450 3451 3452 3453 3454 3455 3456 3457
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_FloatDist{
					FloatDist: &schemapb.FloatArray{
						Data: tanimoto,
					},
				},
			}, nil
		}
	}

3458
	err = errors.New("unexpected error")
3459
	if (vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetFloatVector() != nil) || (vectorsLeft.GetFloatVector() != nil && vectorsRight.GetBinaryVector() != nil) {
3460
		err = errors.New("cannot calculate distance between binary vectors and float vectors")
3461 3462
	}

3463 3464 3465
	log.Debug("Failed to CalcDistance",
		zap.Error(err),
		zap.String("traceID", traceID),
3466
		zap.String("role", typeutil.ProxyRole))
3467

3468 3469 3470 3471 3472 3473 3474 3475
	return &milvuspb.CalcDistanceResults{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		},
	}, nil
}

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

3481
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3482
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3483
	log.Debug("GetPersistentSegmentInfo",
3484
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3485 3486 3487
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

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

J
jingkl 已提交
3544
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3545
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3546
	log.Debug("GetQuerySegmentInfo",
3547
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3548 3549 3550
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

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

C
Cai Yudong 已提交
3611
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3612
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3613
		Base: &commonpb.MsgBase{
3614
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3615 3616
			MsgID:     0,
			Timestamp: 0,
3617
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3618 3619 3620 3621 3622 3623 3624
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3625
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3626 3627 3628
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3629
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3630
		Base: &commonpb.MsgBase{
3631
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3632 3633
			MsgID:     0,
			Timestamp: 0,
3634
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3635 3636 3637 3638 3639 3640 3641 3642
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3643
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3644 3645 3646 3647 3648
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3649
		showSegmentResponse, err := node.rootCoord.ShowSegments(ctx, &milvuspb.ShowSegmentsRequest{
X
XuanYang-cn 已提交
3650
			Base: &commonpb.MsgBase{
3651
				MsgType:   commonpb.MsgType_ShowSegments,
X
XuanYang-cn 已提交
3652 3653
				MsgID:     0,
				Timestamp: 0,
3654
				SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3655 3656 3657 3658 3659 3660 3661
			},
			CollectionID: collectionID,
			PartitionID:  partitionID,
		})
		if err != nil {
			return nil, err
		}
3662
		if showSegmentResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3663 3664 3665 3666 3667 3668
			return nil, errors.New(showSegmentResponse.Status.Reason)
		}
		ret = append(ret, showSegmentResponse.SegmentIDs...)
	}
	return ret, nil
}
3669

J
jingkl 已提交
3670
// Dummy handles dummy request
C
Cai Yudong 已提交
3671
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3672 3673 3674 3675 3676 3677 3678 3679 3680 3681 3682
	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
	}

3683 3684
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3685
		if err != nil {
3686
			log.Debug("Failed to parse dummy query request")
3687 3688 3689
			return failedResponse, nil
		}

3690
		request := &milvuspb.QueryRequest{
3691 3692 3693
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3694
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3695 3696
		}

3697
		_, err = node.Query(ctx, request)
3698
		if err != nil {
3699
			log.Debug("Failed to execute dummy query")
3700 3701
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3702 3703 3704 3705 3706 3707

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

3708 3709
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3710 3711
}

J
jingkl 已提交
3712
// RegisterLink registers a link
C
Cai Yudong 已提交
3713
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
G
godchen 已提交
3714
	code := node.stateCode.Load().(internalpb.StateCode)
D
dragondriver 已提交
3715
	log.Debug("RegisterLink",
3716
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3717
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3718

G
godchen 已提交
3719
	if code != internalpb.StateCode_Healthy {
3720 3721 3722
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3723
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3724
				Reason:    "proxy not healthy",
3725 3726 3727
			},
		}, nil
	}
3728
	metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Inc()
3729 3730 3731
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3732
			ErrorCode: commonpb.ErrorCode_Success,
3733
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3734 3735 3736
		},
	}, nil
}
3737

3738 3739 3740
// 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",
3741
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3742 3743 3744 3745
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
3746
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3747
			zap.String("req", req.Request),
3748
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.ProxyID)))
3749 3750 3751 3752

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3753
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.ProxyID),
3754 3755 3756 3757 3758 3759 3760 3761
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
3762
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3763 3764 3765 3766 3767 3768 3769 3770 3771 3772 3773 3774 3775 3776 3777
			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 已提交
3778 3779 3780 3781 3782 3783 3784 3785 3786 3787
	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,
3788
		SourceID:  Params.ProxyCfg.ProxyID,
D
dragondriver 已提交
3789 3790
	}

3791
	if metricType == metricsinfo.SystemInfoMetrics {
3792 3793 3794 3795 3796 3797 3798
		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))

3799
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3800 3801

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

3808 3809
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3810
		return metrics, nil
3811 3812 3813
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
3814
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3815 3816 3817 3818 3819 3820 3821 3822 3823 3824 3825 3826
		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 已提交
3827 3828 3829
// 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",
3830
		zap.Int64("proxy_id", Params.ProxyCfg.ProxyID),
B
bigsheeper 已提交
3831 3832 3833 3834 3835 3836 3837 3838 3839 3840 3841 3842 3843 3844
		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,
3845
			SourceID:  Params.ProxyCfg.ProxyID,
B
bigsheeper 已提交
3846 3847 3848
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3849
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3850 3851 3852 3853 3854 3855 3856 3857 3858 3859 3860 3861 3862 3863 3864 3865 3866 3867
		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 已提交
3868
//GetCompactionState gets the compaction state of multiple segments
3869 3870 3871 3872 3873 3874 3875 3876 3877 3878 3879 3880 3881
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
}

3882
// ManualCompaction invokes compaction on specified collection
3883 3884 3885 3886 3887 3888 3889 3890 3891 3892 3893 3894 3895 3896 3897 3898 3899 3900 3901 3902 3903 3904 3905 3906 3907 3908
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 已提交
3909 3910 3911
// 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))
3912
	var err error
B
Bingyi Sun 已提交
3913 3914 3915 3916 3917 3918 3919
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3920
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3921 3922 3923 3924
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3925 3926 3927 3928
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3929 3930
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3931 3932 3933 3934
	code := node.stateCode.Load().(internalpb.StateCode)
	return code == internalpb.StateCode_Healthy
}

J
jingkl 已提交
3935
//unhealthyStatus returns the proxy not healthy status
3936 3937 3938
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3939
		Reason:    "proxy not healthy",
3940 3941
	}
}