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 2538 2539 2540 2541
	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))
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10), metrics.SearchLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
	metrics.ProxySearchLatencyPerNQ.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
		strconv.FormatInt(qt.CollectionID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()) / float64(qt.result.Results.NumQueries))
2542 2543 2544
	return qt.result, nil
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

2722 2723 2724 2725 2726
		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()

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

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

2746 2747 2748 2749 2750 2751 2752 2753
	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()))
2754 2755 2756 2757 2758
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2759

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
3106
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3107 3108 3109 3110 3111 3112
			zap.Int64("msgID", qt.Base.MsgID),
			zap.Uint64("timestamp", qt.Base.Timestamp),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields))
3113 3114 3115 3116

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

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

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

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

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

		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))
3182
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3183 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
				}
				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))
3209
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3210 3211 3212 3213 3214 3215 3216 3217 3218 3219 3220 3221
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

3807 3808
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

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

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

3881
// ManualCompaction invokes compaction on specified collection
3882 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
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 已提交
3908 3909 3910
// 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))
3911
	var err error
B
Bingyi Sun 已提交
3912 3913 3914 3915 3916 3917 3918
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

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

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

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