impl.go 156.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 27 28
	"github.com/golang/protobuf/proto"
	"github.com/milvus-io/milvus/api/commonpb"
	"github.com/milvus-io/milvus/api/milvuspb"
29
	"github.com/milvus-io/milvus/internal/common"
X
Xiangyu Wang 已提交
30
	"github.com/milvus-io/milvus/internal/log"
31
	"github.com/milvus-io/milvus/internal/metrics"
J
jaime 已提交
32
	"github.com/milvus-io/milvus/internal/mq/msgstream"
X
Xiangyu Wang 已提交
33 34 35 36
	"github.com/milvus-io/milvus/internal/proto/datapb"
	"github.com/milvus-io/milvus/internal/proto/internalpb"
	"github.com/milvus-io/milvus/internal/proto/proxypb"
	"github.com/milvus-io/milvus/internal/proto/querypb"
37
	"github.com/milvus-io/milvus/internal/util"
38
	"github.com/milvus-io/milvus/internal/util/crypto"
39
	"github.com/milvus-io/milvus/internal/util/errorutil"
40 41
	"github.com/milvus-io/milvus/internal/util/logutil"
	"github.com/milvus-io/milvus/internal/util/metricsinfo"
42
	"github.com/milvus-io/milvus/internal/util/timerecord"
43
	"github.com/milvus-io/milvus/internal/util/trace"
X
Xiangyu Wang 已提交
44
	"github.com/milvus-io/milvus/internal/util/typeutil"
45 46
	"go.uber.org/zap"
	"go.uber.org/zap/zapcore"
47 48
)

49 50
const moduleName = "Proxy"

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

56
// GetComponentStates get state of Proxy.
57 58
func (node *Proxy) GetComponentStates(ctx context.Context) (*milvuspb.ComponentStates, error) {
	stats := &milvuspb.ComponentStates{
G
godchen 已提交
59 60 61 62
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
	}
63
	code, ok := node.stateCode.Load().(commonpb.StateCode)
G
godchen 已提交
64 65 66 67 68 69
	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
	}
76
	info := &milvuspb.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
	ctx = logutil.WithModule(ctx, moduleName)
100
	logutil.Logger(ctx).Info("received request to invalidate collection meta cache",
101
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
102
		zap.String("db", request.DbName),
103 104
		zap.String("collectionName", request.CollectionName),
		zap.Int64("collectionID", request.CollectionID))
D
dragondriver 已提交
105

106
	collectionName := request.CollectionName
107
	collectionID := request.CollectionID
N
neza2017 已提交
108
	if globalMetaCache != nil {
109 110 111 112 113 114
		if collectionName != "" {
			globalMetaCache.RemoveCollection(ctx, collectionName) // no need to return error, though collection may be not cached
		}
		if request.CollectionID != UniqueID(0) {
			globalMetaCache.RemoveCollectionsByID(ctx, collectionID)
		}
N
neza2017 已提交
115
	}
116 117 118 119
	if request.GetBase().GetMsgType() == commonpb.MsgType_DropCollection {
		// no need to handle error, since this Proxy may not create dml stream for the collection.
		_ = node.chMgr.removeDMLStream(request.GetCollectionID())
	}
120
	logutil.Logger(ctx).Info("complete to invalidate collection meta cache",
121
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
122
		zap.String("db", request.DbName),
123 124
		zap.String("collection", collectionName),
		zap.Int64("collectionID", collectionID))
D
dragondriver 已提交
125

126
	return &commonpb.Status{
127
		ErrorCode: commonpb.ErrorCode_Success,
128 129
		Reason:    "",
	}, nil
130 131
}

132
// CreateCollection create a collection by the schema.
133
// TODO(dragondriver): add more detailed ut for ConsistencyLevel, should we support multiple consistency level in Proxy?
C
Cai Yudong 已提交
134
func (node *Proxy) CreateCollection(ctx context.Context, request *milvuspb.CreateCollectionRequest) (*commonpb.Status, error) {
135 136 137
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
138 139 140 141

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
142 143 144
	method := "CreateCollection"
	tr := timerecord.NewTimeRecorder(method)

X
Xiaofan 已提交
145
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
146

147
	cct := &createCollectionTask{
S
sunby 已提交
148
		ctx:                     ctx,
149 150
		Condition:               NewTaskCondition(ctx),
		CreateCollectionRequest: request,
151
		rootCoord:               node.rootCoord,
152 153
	}

154 155 156
	// avoid data race
	lenOfSchema := len(request.Schema)

157 158
	log.Debug(
		rpcReceived(method),
159
		zap.String("traceID", traceID),
160
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
161 162
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
163
		zap.Int("len(schema)", lenOfSchema),
164 165
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
166

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

X
Xiaofan 已提交
179
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
180
		return &commonpb.Status{
181
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
182 183 184 185
			Reason:    err.Error(),
		}, nil
	}

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

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

X
Xiaofan 已提交
215
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
216
		return &commonpb.Status{
217
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
218 219 220 221
			Reason:    err.Error(),
		}, nil
	}

222 223
	log.Debug(
		rpcDone(method),
224
		zap.String("traceID", traceID),
225
		zap.String("role", typeutil.ProxyRole),
226 227 228 229 230 231
		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),
232 233
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
234

X
Xiaofan 已提交
235 236
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
237 238 239
	return cct.result, nil
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
249 250
	method := "DropCollection"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
251
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
252

253
	dct := &dropCollectionTask{
S
sunby 已提交
254
		ctx:                   ctx,
255 256
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
257
		rootCoord:             node.rootCoord,
258
		chMgr:                 node.chMgr,
S
sunby 已提交
259
		chTicker:              node.chTicker,
260 261
	}

262 263
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
264
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
265 266
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
267 268 269 270 271

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

X
Xiaofan 已提交
276
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
277
		return &commonpb.Status{
278
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
279 280 281 282
			Reason:    err.Error(),
		}, nil
	}

283 284
	log.Debug("DropCollection enqueued",
		zap.String("traceID", traceID),
285
		zap.String("role", typeutil.ProxyRole),
286 287 288
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
289 290
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
291 292 293

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DropCollection failed to WaitToFinish",
D
dragondriver 已提交
294
			zap.Error(err),
295
			zap.String("traceID", traceID),
296
			zap.String("role", typeutil.ProxyRole),
297 298 299
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTs", dct.BeginTs()),
			zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
300 301 302
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
303
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
304
		return &commonpb.Status{
305
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
306 307 308 309
			Reason:    err.Error(),
		}, nil
	}

310 311
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
312
		zap.String("role", typeutil.ProxyRole),
313 314 315 316 317 318
		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))

X
Xiaofan 已提交
319 320
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
321 322 323
	return dct.result, nil
}

324
// HasCollection check if the specific collection exists in Milvus.
C
Cai Yudong 已提交
325
func (node *Proxy) HasCollection(ctx context.Context, request *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
326 327 328 329 330
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
331 332 333 334

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
335 336
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
337
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
338
		metrics.TotalLabel).Inc()
339 340 341

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
342
		zap.String("role", typeutil.ProxyRole),
343 344 345
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

346
	hct := &hasCollectionTask{
S
sunby 已提交
347
		ctx:                  ctx,
348 349
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
350
		rootCoord:            node.rootCoord,
351 352
	}

353 354 355 356
	if err := node.sched.ddQueue.Enqueue(hct); err != nil {
		log.Warn("HasCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
357
			zap.String("role", typeutil.ProxyRole),
358 359 360
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
361
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
362
			metrics.AbandonLabel).Inc()
363 364
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
365
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
366 367 368 369 370
				Reason:    err.Error(),
			},
		}, nil
	}

371 372
	log.Debug("HasCollection enqueued",
		zap.String("traceID", traceID),
373
		zap.String("role", typeutil.ProxyRole),
374 375 376
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
377 378
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
379 380 381

	if err := hct.WaitToFinish(); err != nil {
		log.Warn("HasCollection failed to WaitToFinish",
D
dragondriver 已提交
382
			zap.Error(err),
383
			zap.String("traceID", traceID),
384
			zap.String("role", typeutil.ProxyRole),
385 386 387
			zap.Int64("MsgID", hct.ID()),
			zap.Uint64("BeginTS", hct.BeginTs()),
			zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
388 389 390
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
391
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
392
			metrics.FailLabel).Inc()
393 394
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
395
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
396 397 398 399 400
				Reason:    err.Error(),
			},
		}, nil
	}

401 402
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
403
		zap.String("role", typeutil.ProxyRole),
404 405 406 407 408 409
		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))

X
Xiaofan 已提交
410
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
411
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
412
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
413 414 415
	return hct.result, nil
}

416
// LoadCollection load a collection into query nodes.
C
Cai Yudong 已提交
417
func (node *Proxy) LoadCollection(ctx context.Context, request *milvuspb.LoadCollectionRequest) (*commonpb.Status, error) {
418 419 420
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
421 422 423 424

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
425 426
	method := "LoadCollection"
	tr := timerecord.NewTimeRecorder(method)
427

428
	lct := &loadCollectionTask{
S
sunby 已提交
429
		ctx:                   ctx,
430 431
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
432
		queryCoord:            node.queryCoord,
433 434
	}

435 436
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
437
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
438 439
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
440 441 442 443 444

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

X
Xiaofan 已提交
449
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
450
			metrics.AbandonLabel).Inc()
451
		return &commonpb.Status{
452
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
453 454 455
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
456

457 458
	log.Debug("LoadCollection enqueued",
		zap.String("traceID", traceID),
459
		zap.String("role", typeutil.ProxyRole),
460 461 462
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
463 464
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
465 466 467

	if err := lct.WaitToFinish(); err != nil {
		log.Warn("LoadCollection failed to WaitToFinish",
D
dragondriver 已提交
468
			zap.Error(err),
469
			zap.String("traceID", traceID),
470
			zap.String("role", typeutil.ProxyRole),
471 472 473
			zap.Int64("MsgID", lct.ID()),
			zap.Uint64("BeginTS", lct.BeginTs()),
			zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
474 475 476
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
477
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
478
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
479
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
480
			metrics.FailLabel).Inc()
481
		return &commonpb.Status{
482
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
483 484 485 486
			Reason:    err.Error(),
		}, nil
	}

487 488
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
489
		zap.String("role", typeutil.ProxyRole),
490 491 492 493 494 495
		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))

X
Xiaofan 已提交
496
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
497
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
498
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
499
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
500
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
501
	return lct.result, nil
502 503
}

504
// ReleaseCollection remove the loaded collection from query nodes.
C
Cai Yudong 已提交
505
func (node *Proxy) ReleaseCollection(ctx context.Context, request *milvuspb.ReleaseCollectionRequest) (*commonpb.Status, error) {
506 507 508
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
509

510
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
511 512
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
513 514
	method := "ReleaseCollection"
	tr := timerecord.NewTimeRecorder(method)
515

516
	rct := &releaseCollectionTask{
S
sunby 已提交
517
		ctx:                      ctx,
518 519
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
520
		queryCoord:               node.queryCoord,
521
		chMgr:                    node.chMgr,
522 523
	}

524 525
	log.Debug(
		rpcReceived(method),
526
		zap.String("traceID", traceID),
527
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
528 529
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
530 531

	if err := node.sched.ddQueue.Enqueue(rct); err != nil {
532 533
		log.Warn(
			rpcFailedToEnqueue(method),
534 535
			zap.Error(err),
			zap.String("traceID", traceID),
536
			zap.String("role", typeutil.ProxyRole),
537 538 539
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
540
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
541
			metrics.AbandonLabel).Inc()
542
		return &commonpb.Status{
543
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
544 545 546 547
			Reason:    err.Error(),
		}, nil
	}

548 549
	log.Debug(
		rpcEnqueued(method),
550
		zap.String("traceID", traceID),
551
		zap.String("role", typeutil.ProxyRole),
552 553 554
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
555 556
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
557 558

	if err := rct.WaitToFinish(); err != nil {
559 560
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
561
			zap.Error(err),
562
			zap.String("traceID", traceID),
563
			zap.String("role", typeutil.ProxyRole),
564 565 566
			zap.Int64("MsgID", rct.ID()),
			zap.Uint64("BeginTS", rct.BeginTs()),
			zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
567 568 569
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
570
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
571
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
572
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
573
			metrics.FailLabel).Inc()
574
		return &commonpb.Status{
575
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
576 577 578 579
			Reason:    err.Error(),
		}, nil
	}

580 581
	log.Debug(
		rpcDone(method),
582
		zap.String("traceID", traceID),
583
		zap.String("role", typeutil.ProxyRole),
584 585 586 587 588 589
		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))

X
Xiaofan 已提交
590
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
591
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
592
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
593
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
594
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
595
	return rct.result, nil
596 597
}

598
// DescribeCollection get the meta information of specific collection, such as schema, created timestamp and etc.
C
Cai Yudong 已提交
599
func (node *Proxy) DescribeCollection(ctx context.Context, request *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
600 601 602 603 604
	if !node.checkHealthy() {
		return &milvuspb.DescribeCollectionResponse{
			Status: unhealthyStatus(),
		}, nil
	}
605

606
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
607 608
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
609 610
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
611

612
	dct := &describeCollectionTask{
S
sunby 已提交
613
		ctx:                       ctx,
614 615
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
616
		rootCoord:                 node.rootCoord,
617 618
	}

619 620
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
621
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
622 623
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
624 625 626 627 628

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

X
Xiaofan 已提交
633
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
634
			metrics.AbandonLabel).Inc()
635 636
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
637
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
638 639 640 641 642
				Reason:    err.Error(),
			},
		}, nil
	}

643 644
	log.Debug("DescribeCollection enqueued",
		zap.String("traceID", traceID),
645
		zap.String("role", typeutil.ProxyRole),
646 647 648
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
649 650
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
651 652 653

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DescribeCollection failed to WaitToFinish",
D
dragondriver 已提交
654
			zap.Error(err),
655
			zap.String("traceID", traceID),
656
			zap.String("role", typeutil.ProxyRole),
657 658 659
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTS", dct.BeginTs()),
			zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
660 661 662
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
663
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
664
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
665
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
666
			metrics.FailLabel).Inc()
667

668 669
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
670
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
671 672 673 674 675
				Reason:    err.Error(),
			},
		}, nil
	}

676 677
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
678
		zap.String("role", typeutil.ProxyRole),
679 680 681 682 683 684
		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))

X
Xiaofan 已提交
685
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
686
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
687
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
688
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
689
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
690 691 692
	return dct.result, nil
}

693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801
// GetStatistics get the statistics, such as `num_rows`.
// WARNING: It is an experimental API
func (node *Proxy) GetStatistics(ctx context.Context, request *milvuspb.GetStatisticsRequest) (*milvuspb.GetStatisticsResponse, error) {
	if !node.checkHealthy() {
		return &milvuspb.GetStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}

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

	g := &getStatisticsTask{
		request:   request,
		Condition: NewTaskCondition(ctx),
		ctx:       ctx,
		tr:        tr,
		dc:        node.dataCoord,
		qc:        node.queryCoord,
		shardMgr:  node.shardMgr,
	}

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

	if err := node.sched.ddQueue.Enqueue(g); 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("collection", request.CollectionName),
			zap.Strings("partitions", request.PartitionNames))

		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()

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

	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		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.Strings("partitions", request.PartitionNames))

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			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.Strings("partitions", request.PartitionNames))

		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.TotalLabel).Inc()
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()

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

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		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))

	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
	return g.result, nil
}

802
// GetCollectionStatistics get the collection statistics, such as `num_rows`.
C
Cai Yudong 已提交
803
func (node *Proxy) GetCollectionStatistics(ctx context.Context, request *milvuspb.GetCollectionStatisticsRequest) (*milvuspb.GetCollectionStatisticsResponse, error) {
804 805 806 807 808
	if !node.checkHealthy() {
		return &milvuspb.GetCollectionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
809 810 811 812

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
813 814
	method := "GetCollectionStatistics"
	tr := timerecord.NewTimeRecorder(method)
815

816
	g := &getCollectionStatisticsTask{
G
godchen 已提交
817 818 819
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
820
		dataCoord:                      node.dataCoord,
821 822
	}

823 824
	log.Debug(
		rpcReceived(method),
825
		zap.String("traceID", traceID),
826
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
827 828
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
829 830

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
831 832
		log.Warn(
			rpcFailedToEnqueue(method),
833 834
			zap.Error(err),
			zap.String("traceID", traceID),
835
			zap.String("role", typeutil.ProxyRole),
836 837 838
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
839
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
840
			metrics.AbandonLabel).Inc()
841

G
godchen 已提交
842
		return &milvuspb.GetCollectionStatisticsResponse{
843
			Status: &commonpb.Status{
844
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
845 846 847 848 849
				Reason:    err.Error(),
			},
		}, nil
	}

850 851
	log.Debug(
		rpcEnqueued(method),
852
		zap.String("traceID", traceID),
853
		zap.String("role", typeutil.ProxyRole),
854
		zap.Int64("msgID", g.ID()),
855 856
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
857 858
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
859 860

	if err := g.WaitToFinish(); err != nil {
861 862
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
863
			zap.Error(err),
864
			zap.String("traceID", traceID),
865
			zap.String("role", typeutil.ProxyRole),
866 867 868
			zap.Int64("MsgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
869 870 871
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
872
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
873
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
874
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
875
			metrics.FailLabel).Inc()
876

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

885 886
	log.Debug(
		rpcDone(method),
887
		zap.String("traceID", traceID),
888
		zap.String("role", typeutil.ProxyRole),
889
		zap.Int64("msgID", g.ID()),
890 891 892 893 894
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
895
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
896
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
897
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
898
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
899
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
900
	return g.result, nil
901 902
}

903
// ShowCollections list all collections in Milvus.
C
Cai Yudong 已提交
904
func (node *Proxy) ShowCollections(ctx context.Context, request *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
905 906 907 908 909
	if !node.checkHealthy() {
		return &milvuspb.ShowCollectionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
910 911
	method := "ShowCollections"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
912
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
913

914
	sct := &showCollectionsTask{
G
godchen 已提交
915 916 917
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
918
		queryCoord:             node.queryCoord,
919
		rootCoord:              node.rootCoord,
920 921
	}

922
	log.Debug("ShowCollections received",
923
		zap.String("role", typeutil.ProxyRole),
924 925 926 927 928 929
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

930
	err := node.sched.ddQueue.Enqueue(sct)
931
	if err != nil {
932 933
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
934
			zap.String("role", typeutil.ProxyRole),
935 936 937 938 939 940
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

X
Xiaofan 已提交
941
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
G
godchen 已提交
942
		return &milvuspb.ShowCollectionsResponse{
943
			Status: &commonpb.Status{
944
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
945 946 947 948 949
				Reason:    err.Error(),
			},
		}, nil
	}

950
	log.Debug("ShowCollections enqueued",
951
		zap.String("role", typeutil.ProxyRole),
952
		zap.Int64("MsgID", sct.ID()),
953
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
954
		zap.Uint64("TimeStamp", request.TimeStamp),
955 956 957
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
958

959 960
	err = sct.WaitToFinish()
	if err != nil {
961 962
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
963
			zap.String("role", typeutil.ProxyRole),
964 965 966 967 968 969 970
			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),
		)

X
Xiaofan 已提交
971
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
972

G
godchen 已提交
973
		return &milvuspb.ShowCollectionsResponse{
974
			Status: &commonpb.Status{
975
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
976 977 978 979 980
				Reason:    err.Error(),
			},
		}, nil
	}

981
	log.Debug("ShowCollections Done",
982
		zap.String("role", typeutil.ProxyRole),
983 984 985 986
		zap.Int64("MsgID", sct.ID()),
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
987 988
		zap.Int("len(CollectionNames)", len(request.CollectionNames)),
		zap.Int("num_collections", len(sct.result.CollectionNames)))
989

X
Xiaofan 已提交
990 991
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
992 993 994
	return sct.result, nil
}

J
jaime 已提交
995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082
func (node *Proxy) AlterCollection(ctx context.Context, request *milvuspb.AlterCollectionRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

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

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

	act := &alterCollectionTask{
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		AlterCollectionRequest: request,
		rootCoord:              node.rootCoord,
	}

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

	if err := node.sched.ddQueue.Enqueue(act); 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("collection", request.CollectionName))

		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", act.ID()),
		zap.Uint64("BeginTs", act.BeginTs()),
		zap.Uint64("EndTs", act.EndTs()),
		zap.Uint64("timestamp", request.Base.Timestamp),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

	if err := act.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.Int64("MsgID", act.ID()),
			zap.Uint64("BeginTs", act.BeginTs()),
			zap.Uint64("EndTs", act.EndTs()),
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", act.ID()),
		zap.Uint64("BeginTs", act.BeginTs()),
		zap.Uint64("EndTs", act.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
	return act.result, nil
}

1083
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
1084
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
1085 1086 1087
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1088

1089
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
1090 1091
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1092 1093
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
1094
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1095

1096
	cpt := &createPartitionTask{
S
sunby 已提交
1097
		ctx:                    ctx,
1098 1099
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
1100
		rootCoord:              node.rootCoord,
1101 1102 1103
		result:                 nil,
	}

1104 1105 1106
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
1107
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1108 1109 1110
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1111 1112 1113 1114 1115 1116

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
1117
			zap.String("role", typeutil.ProxyRole),
1118 1119 1120 1121
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1122
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
1123

1124
		return &commonpb.Status{
1125
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1126 1127 1128
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1129

1130 1131 1132
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
1133
		zap.String("role", typeutil.ProxyRole),
1134 1135 1136
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
1137 1138 1139
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1140 1141 1142 1143

	if err := cpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish("CreatePartition"),
D
dragondriver 已提交
1144
			zap.Error(err),
1145
			zap.String("traceID", traceID),
1146
			zap.String("role", typeutil.ProxyRole),
1147 1148 1149
			zap.Int64("MsgID", cpt.ID()),
			zap.Uint64("BeginTS", cpt.BeginTs()),
			zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
1150 1151 1152 1153
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1154
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
1155

1156
		return &commonpb.Status{
1157
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1158 1159 1160
			Reason:    err.Error(),
		}, nil
	}
1161 1162 1163 1164

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
1165
		zap.String("role", typeutil.ProxyRole),
1166 1167 1168 1169 1170 1171 1172
		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))

X
Xiaofan 已提交
1173 1174
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1175 1176 1177
	return cpt.result, nil
}

1178
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
1179
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
1180 1181 1182
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1183

1184
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
1185 1186
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1187 1188
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
1189
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1190

1191
	dpt := &dropPartitionTask{
S
sunby 已提交
1192
		ctx:                  ctx,
1193 1194
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1195
		rootCoord:            node.rootCoord,
1196 1197 1198
		result:               nil,
	}

1199 1200 1201
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1202
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1203 1204 1205
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1206 1207 1208 1209 1210 1211

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1212
			zap.String("role", typeutil.ProxyRole),
1213 1214 1215 1216
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1217
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
1218

1219
		return &commonpb.Status{
1220
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1221 1222 1223
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1224

1225 1226 1227
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1228
		zap.String("role", typeutil.ProxyRole),
1229 1230 1231
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1232 1233 1234
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1235 1236 1237 1238

	if err := dpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1239
			zap.Error(err),
1240
			zap.String("traceID", traceID),
1241
			zap.String("role", typeutil.ProxyRole),
1242 1243 1244
			zap.Int64("MsgID", dpt.ID()),
			zap.Uint64("BeginTS", dpt.BeginTs()),
			zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1245 1246 1247 1248
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1249
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
1250

1251
		return &commonpb.Status{
1252
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1253 1254 1255
			Reason:    err.Error(),
		}, nil
	}
1256 1257 1258 1259

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1260
		zap.String("role", typeutil.ProxyRole),
1261 1262 1263 1264 1265 1266 1267
		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))

X
Xiaofan 已提交
1268 1269
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1270 1271 1272
	return dpt.result, nil
}

1273
// HasPartition check if partition exist.
C
Cai Yudong 已提交
1274
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
1275 1276 1277 1278 1279
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1280

D
dragondriver 已提交
1281
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
1282 1283
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1284 1285 1286
	method := "HasPartition"
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
X
Xiaofan 已提交
1287
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1288
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1289

1290
	hpt := &hasPartitionTask{
S
sunby 已提交
1291
		ctx:                 ctx,
1292 1293
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1294
		rootCoord:           node.rootCoord,
1295 1296 1297
		result:              nil,
	}

D
dragondriver 已提交
1298 1299 1300
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1301
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1302 1303 1304
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1305 1306 1307 1308 1309 1310

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

X
Xiaofan 已提交
1316
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1317
			metrics.AbandonLabel).Inc()
1318

1319 1320
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1321
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1322 1323 1324 1325 1326
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1327

D
dragondriver 已提交
1328 1329 1330
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1331
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1332 1333 1334
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1335 1336 1337
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1338 1339 1340 1341

	if err := hpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1342
			zap.Error(err),
D
dragondriver 已提交
1343
			zap.String("traceID", traceID),
1344
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1345 1346 1347
			zap.Int64("MsgID", hpt.ID()),
			zap.Uint64("BeginTS", hpt.BeginTs()),
			zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1348 1349 1350 1351
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1352
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1353
			metrics.FailLabel).Inc()
1354

1355 1356
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1357
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1358 1359 1360 1361 1362
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1363 1364 1365 1366

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1367
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1368 1369 1370 1371 1372 1373 1374
		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))

X
Xiaofan 已提交
1375
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1376
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1377
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1378 1379 1380
	return hpt.result, nil
}

1381
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1382
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1383 1384 1385
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1386

D
dragondriver 已提交
1387
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1388 1389
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1390 1391
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1392

1393
	lpt := &loadPartitionsTask{
G
godchen 已提交
1394 1395 1396
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1397
		queryCoord:            node.queryCoord,
1398 1399
	}

1400 1401 1402
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1403
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1404 1405 1406
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1407 1408 1409 1410 1411 1412

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1413
			zap.String("role", typeutil.ProxyRole),
1414 1415 1416 1417
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1418
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1419
			metrics.AbandonLabel).Inc()
1420

1421
		return &commonpb.Status{
1422
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1423 1424 1425 1426
			Reason:    err.Error(),
		}, nil
	}

1427 1428 1429
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1430
		zap.String("role", typeutil.ProxyRole),
1431 1432 1433
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1434 1435 1436
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1437 1438 1439 1440

	if err := lpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1441
			zap.Error(err),
1442
			zap.String("traceID", traceID),
1443
			zap.String("role", typeutil.ProxyRole),
1444 1445 1446
			zap.Int64("MsgID", lpt.ID()),
			zap.Uint64("BeginTS", lpt.BeginTs()),
			zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1447 1448 1449 1450
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1451
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1452
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1453
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1454
			metrics.FailLabel).Inc()
1455

1456
		return &commonpb.Status{
1457
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1458 1459 1460 1461
			Reason:    err.Error(),
		}, nil
	}

1462 1463 1464
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1465
		zap.String("role", typeutil.ProxyRole),
1466 1467 1468 1469 1470 1471 1472
		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))

X
Xiaofan 已提交
1473
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1474
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1475
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1476
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1477
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1478
	return lpt.result, nil
1479 1480
}

1481
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1482
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1483 1484 1485
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1486 1487 1488 1489 1490

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

1491
	rpt := &releasePartitionsTask{
G
godchen 已提交
1492 1493 1494
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1495
		queryCoord:               node.queryCoord,
1496 1497
	}

1498
	method := "ReleasePartitions"
1499
	tr := timerecord.NewTimeRecorder(method)
1500 1501 1502 1503

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1504
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1505 1506 1507
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1508 1509 1510 1511 1512 1513

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1514
			zap.String("role", typeutil.ProxyRole),
1515 1516 1517 1518
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1519
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1520
			metrics.AbandonLabel).Inc()
1521

1522
		return &commonpb.Status{
1523
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1524 1525 1526 1527
			Reason:    err.Error(),
		}, nil
	}

1528 1529 1530
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1531
		zap.String("role", typeutil.ProxyRole),
1532 1533 1534
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1535 1536 1537
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1538 1539 1540 1541

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1542
			zap.Error(err),
1543
			zap.String("traceID", traceID),
1544
			zap.String("role", typeutil.ProxyRole),
1545 1546 1547
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1548 1549 1550 1551
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1552
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1553
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1554
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1555
			metrics.FailLabel).Inc()
1556

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

1563 1564 1565
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1566
		zap.String("role", typeutil.ProxyRole),
1567 1568 1569 1570 1571 1572 1573
		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))

X
Xiaofan 已提交
1574
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1575
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1576
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1577
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1578
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1579
	return rpt.result, nil
1580 1581
}

1582
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1583
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1584 1585 1586 1587 1588
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1589 1590 1591 1592

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetPartitionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1593 1594
	method := "GetPartitionStatistics"
	tr := timerecord.NewTimeRecorder(method)
1595

1596
	g := &getPartitionStatisticsTask{
1597 1598 1599
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1600
		dataCoord:                     node.dataCoord,
1601 1602
	}

1603 1604 1605
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1606
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1607 1608 1609
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1610 1611 1612 1613 1614 1615

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1616
			zap.String("role", typeutil.ProxyRole),
1617 1618 1619 1620
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1621
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1622
			metrics.AbandonLabel).Inc()
1623

1624 1625 1626 1627 1628 1629 1630 1631
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1632 1633 1634
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1635
		zap.String("role", typeutil.ProxyRole),
1636 1637 1638
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1639 1640 1641
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1642 1643 1644 1645

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1646
			zap.Error(err),
1647
			zap.String("traceID", traceID),
1648
			zap.String("role", typeutil.ProxyRole),
1649 1650 1651
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1652 1653 1654 1655
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1656
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1657
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1658
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1659
			metrics.FailLabel).Inc()
1660

1661 1662 1663 1664 1665 1666 1667 1668
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1669 1670 1671
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1672
		zap.String("role", typeutil.ProxyRole),
1673 1674 1675 1676 1677 1678 1679
		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))

X
Xiaofan 已提交
1680
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1681
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1682
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1683
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1684
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1685
	return g.result, nil
1686 1687
}

1688
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1689
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1690 1691 1692 1693 1694
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1695 1696 1697 1698 1699

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

1700
	spt := &showPartitionsTask{
G
godchen 已提交
1701 1702 1703
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1704
		rootCoord:             node.rootCoord,
1705
		queryCoord:            node.queryCoord,
G
godchen 已提交
1706
		result:                nil,
1707 1708
	}

1709
	method := "ShowPartitions"
1710 1711
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
X
Xiaofan 已提交
1712
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1713
		metrics.TotalLabel).Inc()
1714 1715 1716 1717

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1718
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1719
		zap.Any("request", request))
1720 1721 1722 1723 1724 1725

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

X
Xiaofan 已提交
1729
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1730
			metrics.AbandonLabel).Inc()
1731

G
godchen 已提交
1732
		return &milvuspb.ShowPartitionsResponse{
1733
			Status: &commonpb.Status{
1734
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1735 1736 1737 1738 1739
				Reason:    err.Error(),
			},
		}, nil
	}

1740 1741 1742
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1743
		zap.String("role", typeutil.ProxyRole),
1744 1745 1746
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1747 1748
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1749 1750 1751 1752 1753
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1754
			zap.Error(err),
1755
			zap.String("traceID", traceID),
1756
			zap.String("role", typeutil.ProxyRole),
1757 1758 1759 1760 1761 1762
			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 已提交
1763

X
Xiaofan 已提交
1764
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1765
			metrics.FailLabel).Inc()
1766

G
godchen 已提交
1767
		return &milvuspb.ShowPartitionsResponse{
1768
			Status: &commonpb.Status{
1769
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1770 1771 1772 1773
				Reason:    err.Error(),
			},
		}, nil
	}
1774 1775 1776 1777

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1778
		zap.String("role", typeutil.ProxyRole),
1779 1780 1781 1782 1783 1784 1785
		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))

X
Xiaofan 已提交
1786
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1787
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1788
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1789 1790 1791
	return spt.result, nil
}

S
SimFG 已提交
1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844 1845 1846 1847 1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877
func (node *Proxy) getCollectionProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest, collectionID int64) (int64, error) {
	resp, err := node.queryCoord.ShowCollections(ctx, &querypb.ShowCollectionsRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_ShowCollections,
			MsgID:     request.Base.MsgID,
			Timestamp: request.Base.Timestamp,
			SourceID:  request.Base.SourceID,
		},
		CollectionIDs: []int64{collectionID},
	})
	if err != nil {
		return 0, err
	}
	if len(resp.InMemoryPercentages) == 0 {
		return 0, errors.New("fail to show collections from the querycoord, no data")
	}
	return resp.InMemoryPercentages[0], nil
}

func (node *Proxy) getPartitionProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest, collectionID int64) (int64, error) {
	IDs2Names := make(map[int64]string)
	partitionIDs := make([]int64, 0)
	for _, partitionName := range request.PartitionNames {
		partitionID, err := globalMetaCache.GetPartitionID(ctx, request.CollectionName, partitionName)
		if err != nil {
			return 0, err
		}
		IDs2Names[partitionID] = partitionName
		partitionIDs = append(partitionIDs, partitionID)
	}
	resp, err := node.queryCoord.ShowPartitions(ctx, &querypb.ShowPartitionsRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_ShowPartitions,
			MsgID:     request.Base.MsgID,
			Timestamp: request.Base.Timestamp,
			SourceID:  request.Base.SourceID,
		},
		CollectionID: collectionID,
		PartitionIDs: partitionIDs,
	})
	if err != nil {
		return 0, err
	}
	if len(resp.InMemoryPercentages) != len(partitionIDs) {
		return 0, errors.New("fail to show partitions from the querycoord, invalid data num")
	}
	var progress int64
	for _, p := range resp.InMemoryPercentages {
		progress += p
	}
	progress /= int64(len(partitionIDs))
	return progress, nil
}

func (node *Proxy) GetLoadingProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest) (*milvuspb.GetLoadingProgressResponse, error) {
	if !node.checkHealthy() {
		return &milvuspb.GetLoadingProgressResponse{Status: unhealthyStatus()}, nil
	}
	method := "GetLoadingProgress"
	tr := timerecord.NewTimeRecorder(method)
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ShowPartitions")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

	logger.Info(
		rpcReceived(method),
		zap.String("traceID", traceID),
		zap.Any("request", request))

	getErrResponse := func(err error) *milvuspb.GetLoadingProgressResponse {
		logger.Error("fail to get loading progress", zap.String("collection_name", request.CollectionName),
			zap.Strings("partition_name", request.PartitionNames), zap.Error(err))
		return &milvuspb.GetLoadingProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}
	}
	if err := validateCollectionName(request.CollectionName); err != nil {
		return getErrResponse(err), nil
	}
	collectionID, err := globalMetaCache.GetCollectionID(ctx, request.CollectionName)
	if err != nil {
		return getErrResponse(err), nil
	}
1878 1879 1880 1881 1882
	msgBase := &commonpb.MsgBase{
		MsgType:   commonpb.MsgType_SystemInfo,
		MsgID:     0,
		Timestamp: 0,
		SourceID:  Params.ProxyCfg.GetNodeID(),
S
SimFG 已提交
1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900 1901 1902 1903 1904 1905 1906 1907 1908 1909 1910 1911 1912 1913 1914 1915 1916 1917
	}
	if request.Base == nil {
		request.Base = msgBase
	} else {
		request.Base.MsgID = msgBase.MsgID
		request.Base.Timestamp = msgBase.Timestamp
		request.Base.SourceID = msgBase.SourceID
	}

	var progress int64
	if len(request.GetPartitionNames()) == 0 {
		if progress, err = node.getCollectionProgress(ctx, request, collectionID); err != nil {
			return getErrResponse(err), nil
		}
	} else {
		if progress, err = node.getPartitionProgress(ctx, request, collectionID); err != nil {
			return getErrResponse(err), nil
		}
	}

	logger.Info(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.Any("request", request))
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
	return &milvuspb.GetLoadingProgressResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
		Progress: progress,
	}, nil
}

1918
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1919
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1920 1921 1922
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1923 1924 1925 1926 1927

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

1928
	cit := &createIndexTask{
Z
zhenshan.cao 已提交
1929 1930 1931 1932 1933
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		req:        request,
		rootCoord:  node.rootCoord,
		indexCoord: node.indexCoord,
1934 1935
	}

D
dragondriver 已提交
1936
	method := "CreateIndex"
1937
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1938 1939 1940 1941

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1942
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1943 1944 1945 1946
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1947 1948 1949 1950 1951 1952

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1953
			zap.String("role", typeutil.ProxyRole),
Z
zhenshan.cao 已提交
1954 1955 1956 1957
			zap.String("db", request.GetDbName()),
			zap.String("collection", request.GetCollectionName()),
			zap.String("field", request.GetFieldName()),
			zap.Any("extra_params", request.GetExtraParams()))
D
dragondriver 已提交
1958

X
Xiaofan 已提交
1959
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1960
			metrics.AbandonLabel).Inc()
1961

1962
		return &commonpb.Status{
1963
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1964 1965 1966 1967
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1968 1969 1970
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1971
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1972 1973 1974
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1975 1976 1977 1978
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1979 1980 1981 1982

	if err := cit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1983
			zap.Error(err),
D
dragondriver 已提交
1984
			zap.String("traceID", traceID),
1985
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1986 1987 1988
			zap.Int64("MsgID", cit.ID()),
			zap.Uint64("BeginTs", cit.BeginTs()),
			zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1989 1990 1991 1992 1993
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

X
Xiaofan 已提交
1994
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1995
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1996
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1997
			metrics.FailLabel).Inc()
1998

1999
		return &commonpb.Status{
2000
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
2001 2002 2003 2004
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2005 2006 2007
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2008
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2009 2010 2011 2012 2013 2014 2015 2016
		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))

X
Xiaofan 已提交
2017
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2018
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2019
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2020
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2021
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2022 2023 2024
	return cit.result, nil
}

2025
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
2026
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
2027 2028 2029 2030 2031
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2032 2033 2034 2035 2036

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

2037
	dit := &describeIndexTask{
S
sunby 已提交
2038
		ctx:                  ctx,
2039 2040
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
2041
		indexCoord:           node.indexCoord,
2042 2043
	}

2044 2045 2046
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
2047
	tr := timerecord.NewTimeRecorder(method)
2048 2049 2050 2051

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2052
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2053 2054 2055
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2056 2057 2058 2059 2060 2061 2062
		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),
2063
			zap.String("role", typeutil.ProxyRole),
2064 2065 2066 2067 2068
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

X
Xiaofan 已提交
2069
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2070
			metrics.AbandonLabel).Inc()
2071

2072 2073
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
2074
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2075 2076 2077 2078 2079
				Reason:    err.Error(),
			},
		}, nil
	}

2080 2081 2082
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2083
		zap.String("role", typeutil.ProxyRole),
2084 2085 2086
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2087 2088 2089
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2090 2091 2092 2093 2094
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2095
			zap.Error(err),
2096
			zap.String("traceID", traceID),
2097
			zap.String("role", typeutil.ProxyRole),
2098 2099 2100
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2101 2102 2103
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
2104
			zap.String("index name", indexName))
D
dragondriver 已提交
2105

Z
zhenshan.cao 已提交
2106 2107 2108 2109
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
X
Xiaofan 已提交
2110
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2111
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2112
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2113
			metrics.FailLabel).Inc()
2114

2115 2116
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
2117
				ErrorCode: errCode,
2118 2119 2120 2121 2122
				Reason:    err.Error(),
			},
		}, nil
	}

2123 2124 2125
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2126
		zap.String("role", typeutil.ProxyRole),
2127 2128 2129 2130 2131 2132 2133 2134
		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))

X
Xiaofan 已提交
2135
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2136
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2137
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2138
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2139
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2140 2141 2142
	return dit.result, nil
}

2143
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
2144
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
2145 2146 2147
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2148 2149 2150 2151 2152

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

2153
	dit := &dropIndexTask{
S
sunby 已提交
2154
		ctx:              ctx,
B
BossZou 已提交
2155 2156
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
2157
		indexCoord:       node.indexCoord,
B
BossZou 已提交
2158
	}
G
godchen 已提交
2159

D
dragondriver 已提交
2160
	method := "DropIndex"
2161
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2162 2163 2164 2165

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2166
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2167 2168 2169 2170 2171
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
2172 2173 2174 2175 2176
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2177
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2178 2179 2180 2181
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
2182
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2183
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2184

B
BossZou 已提交
2185
		return &commonpb.Status{
2186
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2187 2188 2189
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2190

D
dragondriver 已提交
2191 2192 2193
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2194
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2195 2196 2197
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2198 2199 2200 2201
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
2202 2203 2204 2205

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2206
			zap.Error(err),
D
dragondriver 已提交
2207
			zap.String("traceID", traceID),
2208
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2209 2210 2211
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2212 2213 2214 2215 2216
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2217
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2218
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2219
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2220
			metrics.FailLabel).Inc()
2221

B
BossZou 已提交
2222
		return &commonpb.Status{
2223
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2224 2225 2226
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2227 2228 2229 2230

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2231
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2232 2233 2234 2235 2236 2237 2238 2239
		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))

X
Xiaofan 已提交
2240
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2241
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2242
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2243
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2244
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
2245 2246 2247
	return dit.result, nil
}

2248 2249
// 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.
2250
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2251
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
2252 2253 2254 2255 2256
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2257 2258 2259 2260 2261

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

2262
	gibpt := &getIndexBuildProgressTask{
2263 2264 2265
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
2266 2267
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
2268
		dataCoord:                    node.dataCoord,
2269 2270
	}

2271
	method := "GetIndexBuildProgress"
2272
	tr := timerecord.NewTimeRecorder(method)
2273 2274 2275 2276

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2277
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2278 2279 2280 2281
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2282 2283 2284 2285 2286 2287

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2288
			zap.String("role", typeutil.ProxyRole),
2289 2290 2291 2292
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
2293
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2294
			metrics.AbandonLabel).Inc()
2295

2296 2297 2298 2299 2300 2301 2302 2303
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2304 2305 2306
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2307
		zap.String("role", typeutil.ProxyRole),
2308 2309 2310
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
2311 2312 2313 2314
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2315 2316 2317 2318

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2319
			zap.Error(err),
2320
			zap.String("traceID", traceID),
2321
			zap.String("role", typeutil.ProxyRole),
2322 2323 2324
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2325 2326 2327 2328
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
2329
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2330
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2331
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2332
			metrics.FailLabel).Inc()
2333 2334 2335 2336 2337 2338 2339 2340

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2341 2342 2343 2344

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2345
		zap.String("role", typeutil.ProxyRole),
2346 2347 2348 2349 2350 2351 2352 2353
		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))
2354

X
Xiaofan 已提交
2355
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2356
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2357
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2358
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2359
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2360
	return gibpt.result, nil
2361 2362
}

2363
// GetIndexState get the build-state of index.
2364
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2365
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2366 2367 2368 2369 2370
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2371 2372 2373 2374 2375

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

2376
	dipt := &getIndexStateTask{
G
godchen 已提交
2377 2378 2379
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2380 2381
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2382 2383
	}

2384
	method := "GetIndexState"
2385
	tr := timerecord.NewTimeRecorder(method)
2386 2387 2388 2389

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2390
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2391 2392 2393 2394
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2395 2396 2397 2398 2399 2400

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2401
			zap.String("role", typeutil.ProxyRole),
2402 2403 2404 2405 2406
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2407
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2408
			metrics.AbandonLabel).Inc()
2409

G
godchen 已提交
2410
		return &milvuspb.GetIndexStateResponse{
2411
			Status: &commonpb.Status{
2412
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2413 2414 2415 2416 2417
				Reason:    err.Error(),
			},
		}, nil
	}

2418 2419 2420
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2421
		zap.String("role", typeutil.ProxyRole),
2422 2423 2424
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2425 2426 2427 2428
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2429 2430 2431 2432

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2433
			zap.Error(err),
2434
			zap.String("traceID", traceID),
2435
			zap.String("role", typeutil.ProxyRole),
2436 2437 2438
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2439 2440 2441 2442 2443
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2444
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2445
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2446
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2447
			metrics.FailLabel).Inc()
2448

G
godchen 已提交
2449
		return &milvuspb.GetIndexStateResponse{
2450
			Status: &commonpb.Status{
2451
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2452 2453 2454 2455 2456
				Reason:    err.Error(),
			},
		}, nil
	}

2457 2458 2459
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2460
		zap.String("role", typeutil.ProxyRole),
2461 2462 2463 2464 2465 2466 2467 2468
		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))

X
Xiaofan 已提交
2469
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2470
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2471
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2472
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2473
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2474 2475 2476
	return dipt.result, nil
}

2477
// Insert insert records into collection.
C
Cai Yudong 已提交
2478
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2479 2480 2481 2482 2483 2484
	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))

2485 2486 2487 2488 2489
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2490 2491
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
2492
	receiveSize := proto.Size(request)
2493 2494
	rateCol.Add(internalpb.RateType_DMLInsert.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Add(float64(receiveSize))
D
dragondriver 已提交
2495

2496 2497 2498 2499 2500
	defer func() {
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.TotalLabel).Inc()
	}()

2501
	it := &insertTask{
2502 2503
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2504
		// req:       request,
2505 2506 2507 2508
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2509
			InsertRequest: internalpb.InsertRequest{
2510
				Base: &commonpb.MsgBase{
X
xige-16 已提交
2511 2512
					MsgType:  commonpb.MsgType_Insert,
					MsgID:    0,
X
Xiaofan 已提交
2513
					SourceID: Params.ProxyCfg.GetNodeID(),
2514 2515 2516
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2517 2518 2519
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2520
				// RowData: transfer column based request to this
2521 2522
			},
		},
2523
		idAllocator:   node.rowIDAllocator,
2524 2525 2526
		segIDAssigner: node.segAssigner,
		chMgr:         node.chMgr,
		chTicker:      node.chTicker,
2527
	}
2528 2529

	if len(it.PartitionName) <= 0 {
2530
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2531 2532
	}

X
Xiangyu Wang 已提交
2533
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2534
		numRows := request.NumRows
2535 2536 2537 2538
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2539

X
Xiangyu Wang 已提交
2540 2541 2542 2543 2544 2545 2546
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2547 2548
	}

X
Xiangyu Wang 已提交
2549
	log.Debug("Enqueue insert request in Proxy",
2550
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2551 2552 2553 2554 2555
		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)),
2556 2557
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2558

X
Xiangyu Wang 已提交
2559 2560
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2561 2562
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2563
		return constructFailedResponse(err), nil
2564
	}
D
dragondriver 已提交
2565

X
Xiangyu Wang 已提交
2566
	log.Debug("Detail of insert request in Proxy",
2567
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2568
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2569 2570 2571 2572 2573
		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 已提交
2574 2575 2576 2577 2578
		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))
2579
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2580
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2581 2582 2583 2584 2585
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2586
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2587 2588 2589 2590 2591 2592 2593 2594 2595 2596 2597
			errIndex := make([]uint32, numRows)
			for i := uint32(0); i < numRows; i++ {
				errIndex[i] = i
			}
			it.result.ErrIndex = errIndex
		}

		setErrorIndex()
	}

	// InsertCnt always equals to the number of entities in the request
X
xige-16 已提交
2598
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2599

2600
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2601
		metrics.SuccessLabel).Inc()
2602 2603
	successCnt := it.result.InsertCnt - int64(len(it.result.ErrIndex))
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(successCnt))
2604
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2605 2606 2607
	return it.result, nil
}

2608
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2609
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2610 2611 2612
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2613 2614
	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))
2615

2616
	receiveSize := proto.Size(request)
2617 2618
	rateCol.Add(internalpb.RateType_DMLDelete.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Add(float64(receiveSize))
2619

G
groot 已提交
2620 2621 2622 2623 2624 2625
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2626 2627 2628
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

2629 2630
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2631
	dt := &deleteTask{
X
xige-16 已提交
2632 2633 2634
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		deleteExpr: request.Expr,
G
godchen 已提交
2635
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2636 2637 2638
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2639 2640 2641 2642 2643
			DeleteRequest: internalpb.DeleteRequest{
				Base: &commonpb.MsgBase{
					MsgType: commonpb.MsgType_Delete,
					MsgID:   0,
				},
X
xige-16 已提交
2644
				DbName:         request.DbName,
G
godchen 已提交
2645 2646 2647
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2648 2649 2650 2651
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2652 2653
	}

2654
	log.Debug("Enqueue delete request in Proxy",
2655
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2656 2657 2658 2659
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2660 2661 2662 2663

	// 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))
X
Xiaofan 已提交
2664
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2665
			metrics.FailLabel).Inc()
2666

G
groot 已提交
2667 2668 2669 2670 2671 2672 2673 2674
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2675
	log.Debug("Detail of delete request in Proxy",
2676
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2677 2678 2679 2680 2681
		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),
2682 2683
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2684

2685 2686
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
X
Xiaofan 已提交
2687
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2688
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2689
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2690
			metrics.FailLabel).Inc()
G
groot 已提交
2691 2692 2693 2694 2695 2696 2697 2698
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

X
Xiaofan 已提交
2699
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2700
		metrics.SuccessLabel).Inc()
2701
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2702 2703 2704
	return dt.result, nil
}

2705
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2706
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2707 2708 2709 2710 2711
	receiveSize := proto.Size(request)
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.SearchLabel).Add(float64(receiveSize))

	rateCol.Add(internalpb.RateType_DQLSearch.String(), float64(request.GetNq()))

2712 2713 2714 2715 2716
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2717 2718
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
2719 2720
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2721

C
cai.zhang 已提交
2722 2723
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2724

2725
	qt := &searchTask{
S
sunby 已提交
2726
		ctx:       ctx,
2727
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2728
		SearchRequest: &internalpb.SearchRequest{
2729
			Base: &commonpb.MsgBase{
2730
				MsgType:  commonpb.MsgType_Search,
X
Xiaofan 已提交
2731
				SourceID: Params.ProxyCfg.GetNodeID(),
2732
			},
2733
			ReqID: Params.ProxyCfg.GetNodeID(),
2734
		},
2735 2736 2737 2738
		request:  request,
		qc:       node.queryCoord,
		tr:       timerecord.NewTimeRecorder("search"),
		shardMgr: node.shardMgr,
2739 2740
	}

2741 2742 2743
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

Z
Zach 已提交
2744
	log.Ctx(ctx).Info(
2745
		rpcReceived(method),
2746
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2747 2748 2749 2750 2751
		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)),
2752 2753 2754 2755
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2756

2757
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
2758
		log.Ctx(ctx).Warn(
2759
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2760
			zap.Error(err),
2761
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2762 2763 2764 2765 2766 2767
			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),
2768 2769 2770
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2771

2772 2773
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
2774

2775 2776
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2777
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2778 2779 2780 2781
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
2782
	tr.CtxRecord(ctx, "search request enqueue")
2783

Z
Zach 已提交
2784
	log.Ctx(ctx).Debug(
2785
		rpcEnqueued(method),
2786
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2787
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2788 2789 2790 2791 2792
		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),
2793
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2794 2795 2796 2797
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2798

2799
	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
2800
		log.Ctx(ctx).Warn(
2801
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2802
			zap.Error(err),
2803
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2804
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2805 2806 2807 2808
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2809
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2810 2811 2812 2813
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2814

2815 2816
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()
2817

2818 2819
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2820
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2821 2822 2823 2824 2825
				Reason:    err.Error(),
			},
		}, nil
	}

Z
Zach 已提交
2826
	span := tr.CtxRecord(ctx, "wait search result")
2827 2828
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel).Observe(float64(span.Milliseconds()))
Z
Zach 已提交
2829
	log.Ctx(ctx).Debug(
2830
		rpcDone(method),
2831
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2832 2833 2834 2835 2836 2837
		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)),
2838 2839 2840 2841
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2842

2843 2844 2845
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(qt.result.GetResults().GetNumQueries()))
C
cai.zhang 已提交
2846
	searchDur := tr.ElapseSpan().Milliseconds()
X
Xiaofan 已提交
2847
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
2848
		metrics.SearchLabel).Observe(float64(searchDur))
2849 2850 2851 2852

	if qt.result != nil {
		sentSize := proto.Size(qt.result)
		metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
2853
		rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
2854
	}
2855 2856 2857
	return qt.result, nil
}

2858
// Flush notify data nodes to persist the data of collection.
2859 2860 2861 2862 2863 2864 2865
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2866
	if !node.checkHealthy() {
2867 2868
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2869
	}
D
dragondriver 已提交
2870 2871 2872 2873 2874

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

2875
	ft := &flushTask{
T
ThreadDao 已提交
2876 2877 2878
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2879
		dataCoord:    node.dataCoord,
2880 2881
	}

D
dragondriver 已提交
2882
	method := "Flush"
2883
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
2884
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2885 2886 2887 2888

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2889
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2890 2891
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2892 2893 2894 2895 2896 2897

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

X
Xiaofan 已提交
2902
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
2903

2904 2905
		resp.Status.Reason = err.Error()
		return resp, nil
2906 2907
	}

D
dragondriver 已提交
2908 2909 2910
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2911
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2912 2913 2914
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2915 2916
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2917 2918 2919 2920

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2921
			zap.Error(err),
D
dragondriver 已提交
2922
			zap.String("traceID", traceID),
2923
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2924 2925 2926
			zap.Int64("MsgID", ft.ID()),
			zap.Uint64("BeginTs", ft.BeginTs()),
			zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2927 2928 2929
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

X
Xiaofan 已提交
2930
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
2931

D
dragondriver 已提交
2932
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2933 2934
		resp.Status.Reason = err.Error()
		return resp, nil
2935 2936
	}

D
dragondriver 已提交
2937 2938 2939
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2940
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2941 2942 2943 2944 2945 2946
		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))

X
Xiaofan 已提交
2947 2948
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2949
	return ft.result, nil
2950 2951
}

2952
// Query get the records by primary keys.
C
Cai Yudong 已提交
2953
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2954 2955 2956 2957 2958
	receiveSize := proto.Size(request)
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.QueryLabel).Add(float64(receiveSize))

	rateCol.Add(internalpb.RateType_DQLQuery.String(), 1)

2959 2960 2961 2962 2963
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2964

D
dragondriver 已提交
2965 2966
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
2967
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2968

2969
	qt := &queryTask{
2970 2971 2972 2973 2974
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
X
Xiaofan 已提交
2975
				SourceID: Params.ProxyCfg.GetNodeID(),
2976
			},
2977
			ReqID: Params.ProxyCfg.GetNodeID(),
2978
		},
2979 2980
		request:          request,
		qc:               node.queryCoord,
2981
		queryShardPolicy: mergeRoundRobinPolicy,
2982
		shardMgr:         node.shardMgr,
2983 2984
	}

D
dragondriver 已提交
2985 2986
	method := "Query"

2987 2988 2989
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()

Z
Zach 已提交
2990
	log.Ctx(ctx).Info(
D
dragondriver 已提交
2991
		rpcReceived(method),
2992
		zap.String("role", typeutil.ProxyRole),
2993 2994
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2995 2996 2997 2998 2999
		zap.Strings("partitions", request.PartitionNames),
		zap.String("expr", request.Expr),
		zap.Strings("OutputFields", request.OutputFields),
		zap.Uint64("travel_timestamp", request.TravelTimestamp),
		zap.Uint64("guarantee_timestamp", request.GuaranteeTimestamp))
G
godchen 已提交
3000

D
dragondriver 已提交
3001
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
3002
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
3003 3004 3005
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("role", typeutil.ProxyRole),
3006 3007 3008
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
3009

3010 3011 3012
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()

3013 3014 3015 3016 3017 3018
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
3019
	}
Z
Zach 已提交
3020
	tr.CtxRecord(ctx, "query request enqueue")
3021

Z
Zach 已提交
3022
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3023
		rpcEnqueued(method),
3024
		zap.String("role", typeutil.ProxyRole),
3025
		zap.Int64("msgID", qt.ID()),
3026 3027
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
3028
		zap.Strings("partitions", request.PartitionNames))
D
dragondriver 已提交
3029 3030

	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
3031
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
3032 3033
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
3034
			zap.String("role", typeutil.ProxyRole),
3035
			zap.Int64("msgID", qt.ID()),
3036 3037 3038
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
3039

3040 3041
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()
3042

3043 3044 3045 3046 3047 3048 3049
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
3050
	span := tr.CtxRecord(ctx, "wait query result")
3051 3052
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel).Observe(float64(span.Milliseconds()))
Z
Zach 已提交
3053
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3054 3055
		rpcDone(method),
		zap.String("role", typeutil.ProxyRole),
3056
		zap.Int64("msgID", qt.ID()),
3057 3058 3059
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
3060

3061 3062 3063 3064
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.SuccessLabel).Inc()

	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
3065
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
3066 3067

	ret := &milvuspb.QueryResults{
3068 3069
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
3070 3071
	}
	sentSize := proto.Size(qt.result)
3072
	rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
3073 3074
	metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
	return ret, nil
3075
}
3076

3077
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
3078 3079 3080 3081
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3082 3083 3084 3085 3086

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

Y
Yusup 已提交
3087 3088 3089 3090 3091 3092 3093
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
3094
	method := "CreateAlias"
3095
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
3096
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3097 3098 3099 3100 3101 3102 3103 3104 3105 3106 3107 3108 3109 3110 3111 3112 3113 3114 3115

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

X
Xiaofan 已提交
3116
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
3117

Y
Yusup 已提交
3118 3119 3120 3121 3122 3123
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3124 3125 3126
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3127
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3128 3129 3130 3131
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3132 3133
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3134 3135 3136 3137

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3138
			zap.Error(err),
D
dragondriver 已提交
3139
			zap.String("traceID", traceID),
3140
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3141 3142 3143 3144
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3145 3146
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
X
Xiaofan 已提交
3147
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
3148 3149 3150 3151 3152 3153 3154

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

D
dragondriver 已提交
3155 3156 3157 3158 3159 3160 3161 3162 3163 3164 3165
	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))

X
Xiaofan 已提交
3166 3167
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3168 3169 3170
	return cat.result, nil
}

3171
// DropAlias alter the alias of collection.
Y
Yusup 已提交
3172 3173 3174 3175
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3176 3177 3178 3179 3180

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

Y
Yusup 已提交
3181 3182 3183 3184 3185 3186 3187
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
3188
	method := "DropAlias"
3189
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
3190
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3191 3192 3193 3194 3195 3196 3197 3198 3199 3200 3201 3202 3203 3204 3205 3206

	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))
X
Xiaofan 已提交
3207
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
3208

Y
Yusup 已提交
3209 3210 3211 3212 3213 3214
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3215 3216 3217
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3218
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3219 3220 3221 3222
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3223
		zap.String("alias", request.Alias))
D
dragondriver 已提交
3224 3225 3226 3227

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3228
			zap.Error(err),
D
dragondriver 已提交
3229
			zap.String("traceID", traceID),
3230
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3231 3232 3233 3234
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3235 3236
			zap.String("alias", request.Alias))

X
Xiaofan 已提交
3237
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3238

Y
Yusup 已提交
3239 3240 3241 3242 3243 3244
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3245 3246 3247 3248 3249 3250 3251 3252 3253 3254
	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))

X
Xiaofan 已提交
3255 3256
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3257 3258 3259
	return dat.result, nil
}

3260
// AlterAlias alter alias of collection.
Y
Yusup 已提交
3261 3262 3263 3264
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3265 3266 3267 3268 3269

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

Y
Yusup 已提交
3270 3271 3272 3273 3274 3275 3276
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
3277
	method := "AlterAlias"
3278
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
3279
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3280 3281 3282 3283 3284 3285 3286 3287 3288 3289 3290 3291 3292 3293 3294 3295 3296 3297

	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))
X
Xiaofan 已提交
3298
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
3299

Y
Yusup 已提交
3300 3301 3302 3303 3304 3305
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3306 3307 3308
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3309
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3310 3311 3312 3313
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3314 3315
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3316 3317 3318 3319

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3320
			zap.Error(err),
D
dragondriver 已提交
3321
			zap.String("traceID", traceID),
3322
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3323 3324 3325 3326
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3327 3328 3329
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
3330
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3331

Y
Yusup 已提交
3332 3333 3334 3335 3336 3337
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3338 3339 3340 3341 3342 3343 3344 3345 3346 3347 3348
	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))

X
Xiaofan 已提交
3349 3350
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3351 3352 3353
	return aat.result, nil
}

3354
// CalcDistance calculates the distances between vectors.
3355
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3356 3357 3358 3359 3360
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3361

3362 3363 3364 3365
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3366 3367
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3368

3369 3370 3371 3372 3373
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3374 3375
		}

3376
		qt := &queryTask{
3377 3378 3379 3380 3381
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
X
Xiaofan 已提交
3382
					SourceID: Params.ProxyCfg.GetNodeID(),
3383
				},
3384
				ReqID: Params.ProxyCfg.GetNodeID(),
3385
			},
3386 3387 3388 3389
			request: queryRequest,
			qc:      node.queryCoord,
			ids:     ids.IdArray,

3390
			queryShardPolicy: mergeRoundRobinPolicy,
3391
			shardMgr:         node.shardMgr,
3392 3393
		}

G
groot 已提交
3394 3395 3396 3397 3398 3399
		items := []zapcore.Field{
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields),
		}

3400
		err := node.sched.dqQueue.Enqueue(qt)
3401
		if err != nil {
G
groot 已提交
3402
			log.Error("CalcDistance queryTask failed to enqueue", append(items, zap.Error(err))...)
3403

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

G
groot 已提交
3412
		log.Debug("CalcDistance queryTask enqueued", items...)
3413 3414 3415

		err = qt.WaitToFinish()
		if err != nil {
G
groot 已提交
3416
			log.Error("CalcDistance queryTask failed to WaitToFinish", append(items, zap.Error(err))...)
3417 3418 3419 3420 3421 3422

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3423
			}, err
3424
		}
3425

G
groot 已提交
3426
		log.Debug("CalcDistance queryTask Done", items...)
3427 3428

		return &milvuspb.QueryResults{
3429 3430
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3431 3432 3433
		}, nil
	}

G
groot 已提交
3434 3435 3436 3437
	// calcDistanceTask is not a standard task, no need to enqueue
	task := &calcDistanceTask{
		traceID:   traceID,
		queryFunc: query,
3438 3439
	}

G
groot 已提交
3440
	return task.Execute(ctx, request)
3441 3442
}

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

3448
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3449
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3450
	log.Debug("GetPersistentSegmentInfo",
3451
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3452 3453 3454
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3455
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3456
		Status: &commonpb.Status{
3457
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3458 3459
		},
	}
3460 3461 3462 3463
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3464 3465
	method := "GetPersistentSegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
3466
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3467
		metrics.TotalLabel).Inc()
G
godchen 已提交
3468
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
X
XuanYang-cn 已提交
3469
	if err != nil {
3470
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3471 3472
		return resp, nil
	}
3473
	infoResp, err := node.dataCoord.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
X
XuanYang-cn 已提交
3474
		Base: &commonpb.MsgBase{
3475
			MsgType:   commonpb.MsgType_SegmentInfo,
X
XuanYang-cn 已提交
3476 3477
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3478
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3479 3480 3481 3482
		},
		SegmentIDs: segments,
	})
	if err != nil {
3483
		log.Debug("GetPersistentSegmentInfo fail", zap.Error(err))
3484
		resp.Status.Reason = fmt.Errorf("dataCoord:GetSegmentInfo, err:%w", err).Error()
X
XuanYang-cn 已提交
3485 3486
		return resp, nil
	}
3487
	log.Debug("GetPersistentSegmentInfo ", zap.Int("len(infos)", len(infoResp.Infos)), zap.Any("status", infoResp.Status))
3488
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3489 3490 3491 3492 3493 3494
		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 已提交
3495
			SegmentID:    info.ID,
X
XuanYang-cn 已提交
3496 3497
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
S
sunby 已提交
3498
			NumRows:      info.NumOfRows,
X
XuanYang-cn 已提交
3499 3500 3501
			State:        info.State,
		}
	}
X
Xiaofan 已提交
3502
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3503
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
3504
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
3505
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
X
XuanYang-cn 已提交
3506 3507 3508 3509
	resp.Infos = persistentInfos
	return resp, nil
}

J
jingkl 已提交
3510
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3511
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3512
	log.Debug("GetQuerySegmentInfo",
3513
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3514 3515 3516
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3517
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3518
		Status: &commonpb.Status{
3519
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3520 3521
		},
	}
3522 3523 3524 3525
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3526

3527 3528 3529 3530 3531
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3532
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
Z
zhenshan.cao 已提交
3533
		Base: &commonpb.MsgBase{
3534
			MsgType:   commonpb.MsgType_SegmentInfo,
Z
zhenshan.cao 已提交
3535 3536
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3537
			SourceID:  Params.ProxyCfg.GetNodeID(),
Z
zhenshan.cao 已提交
3538
		},
3539
		CollectionID: collID,
Z
zhenshan.cao 已提交
3540 3541
	})
	if err != nil {
3542
		log.Error("Failed to get segment info from QueryCoord",
3543
			zap.Error(err))
Z
zhenshan.cao 已提交
3544 3545 3546
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3547
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3548
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3549
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3550 3551 3552 3553 3554 3555 3556 3557 3558 3559 3560 3561 3562
		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,
X
xige-16 已提交
3563
			State:        info.SegmentState,
3564
			NodeIds:      info.NodeIds,
Z
zhenshan.cao 已提交
3565 3566
		}
	}
3567
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3568 3569 3570 3571
	resp.Infos = queryInfos
	return resp, nil
}

C
Cai Yudong 已提交
3572
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3573
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3574
		Base: &commonpb.MsgBase{
3575
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3576 3577
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3578
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3579 3580 3581 3582 3583 3584 3585
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3586
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3587 3588 3589
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3590
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3591
		Base: &commonpb.MsgBase{
3592
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3593 3594
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3595
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3596 3597 3598 3599 3600 3601 3602 3603
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3604
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3605 3606 3607 3608 3609
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3610
		getSegmentsByStatesResponse, err := node.dataCoord.GetSegmentsByStates(ctx, &datapb.GetSegmentsByStatesRequest{
X
XuanYang-cn 已提交
3611 3612
			CollectionID: collectionID,
			PartitionID:  partitionID,
3613
			States:       []commonpb.SegmentState{commonpb.SegmentState_Flushing, commonpb.SegmentState_Flushed, commonpb.SegmentState_Sealed},
X
XuanYang-cn 已提交
3614 3615 3616 3617
		})
		if err != nil {
			return nil, err
		}
3618 3619
		if getSegmentsByStatesResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
			return nil, errors.New(getSegmentsByStatesResponse.Status.Reason)
X
XuanYang-cn 已提交
3620
		}
3621
		ret = append(ret, getSegmentsByStatesResponse.GetSegments()...)
X
XuanYang-cn 已提交
3622 3623 3624
	}
	return ret, nil
}
3625

J
jingkl 已提交
3626
// Dummy handles dummy request
C
Cai Yudong 已提交
3627
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3628 3629 3630 3631 3632 3633 3634 3635 3636 3637 3638
	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
	}

3639 3640
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3641
		if err != nil {
3642
			log.Debug("Failed to parse dummy query request")
3643 3644 3645
			return failedResponse, nil
		}

3646
		request := &milvuspb.QueryRequest{
3647 3648 3649
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3650
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3651 3652
		}

3653
		_, err = node.Query(ctx, request)
3654
		if err != nil {
3655
			log.Debug("Failed to execute dummy query")
3656 3657
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3658 3659 3660 3661 3662 3663

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

3664 3665
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3666 3667
}

J
jingkl 已提交
3668
// RegisterLink registers a link
C
Cai Yudong 已提交
3669
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
3670
	code := node.stateCode.Load().(commonpb.StateCode)
D
dragondriver 已提交
3671
	log.Debug("RegisterLink",
3672
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3673
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3674

3675
	if code != commonpb.StateCode_Healthy {
3676 3677 3678
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3679
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3680
				Reason:    "proxy not healthy",
3681 3682 3683
			},
		}, nil
	}
X
Xiaofan 已提交
3684
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Inc()
3685 3686 3687
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3688
			ErrorCode: commonpb.ErrorCode_Success,
3689
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3690 3691 3692
		},
	}, nil
}
3693

3694
// GetMetrics gets the metrics of proxy
3695 3696 3697
// 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",
X
Xiaofan 已提交
3698
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3699 3700 3701 3702
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
X
Xiaofan 已提交
3703
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3704
			zap.String("req", req.Request),
X
Xiaofan 已提交
3705
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.GetNodeID())))
3706 3707 3708 3709

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
Xiaofan 已提交
3710
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.GetNodeID()),
3711 3712 3713 3714 3715 3716 3717 3718
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
X
Xiaofan 已提交
3719
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3720 3721 3722 3723 3724 3725 3726 3727 3728 3729 3730 3731 3732 3733 3734
			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 已提交
3735 3736
	req.Base = &commonpb.MsgBase{
		MsgType:   commonpb.MsgType_SystemInfo,
3737
		MsgID:     0,
D
dragondriver 已提交
3738
		Timestamp: 0,
X
Xiaofan 已提交
3739
		SourceID:  Params.ProxyCfg.GetNodeID(),
D
dragondriver 已提交
3740 3741
	}

3742
	if metricType == metricsinfo.SystemInfoMetrics {
3743 3744 3745 3746 3747 3748 3749
		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))

3750
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3751 3752

		log.Debug("Proxy.GetMetrics",
X
Xiaofan 已提交
3753
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3754 3755 3756 3757 3758
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3759 3760
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3761
		return metrics, nil
3762 3763 3764
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
X
Xiaofan 已提交
3765
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3766 3767 3768 3769 3770 3771 3772 3773 3774 3775 3776 3777
		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
}

3778 3779 3780 3781 3782 3783 3784 3785 3786 3787 3788 3789 3790 3791 3792 3793 3794 3795 3796 3797 3798 3799 3800 3801 3802 3803 3804 3805 3806 3807 3808 3809 3810 3811 3812 3813 3814 3815 3816 3817
// GetProxyMetrics gets the metrics of proxy, it's an internal interface which is different from GetMetrics interface,
// because it only obtains the metrics of Proxy, not including the topological metrics of Query cluster and Data cluster.
func (node *Proxy) GetProxyMetrics(ctx context.Context, req *milvuspb.GetMetricsRequest) (*milvuspb.GetMetricsResponse, error) {
	log.Debug("Proxy.GetProxyMetrics",
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
		zap.String("req", req.Request))

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

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.GetNodeID()),
			},
		}, nil
	}

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

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

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

	req.Base = &commonpb.MsgBase{
3818 3819 3820 3821
		MsgType:   commonpb.MsgType_SystemInfo,
		MsgID:     0,
		Timestamp: 0,
		SourceID:  Params.ProxyCfg.GetNodeID(),
3822 3823 3824 3825 3826 3827 3828 3829 3830 3831 3832 3833 3834 3835 3836 3837 3838 3839 3840 3841 3842 3843 3844 3845 3846 3847 3848 3849 3850 3851 3852 3853 3854 3855 3856 3857 3858 3859 3860 3861
	}

	if metricType == metricsinfo.SystemInfoMetrics {
		proxyMetrics, err := getProxyMetrics(ctx, req, node)
		if err != nil {
			log.Warn("Proxy.GetProxyMetrics failed to getProxyMetrics",
				zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
				zap.String("req", req.Request),
				zap.Error(err))

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

		log.Debug("Proxy.GetProxyMetrics",
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Error(err))

		return proxyMetrics, nil
	}

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

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

B
bigsheeper 已提交
3862 3863 3864
// 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",
X
Xiaofan 已提交
3865
		zap.Int64("proxy_id", Params.ProxyCfg.GetNodeID()),
B
bigsheeper 已提交
3866 3867 3868 3869 3870 3871 3872 3873 3874
		zap.Any("req", req))

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

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
3875 3876 3877 3878 3879 3880 3881

	collectionID, err := globalMetaCache.GetCollectionID(ctx, req.GetCollectionName())
	if err != nil {
		log.Error("failed to get collection id", zap.String("collection name", req.GetCollectionName()), zap.Error(err))
		status.Reason = err.Error()
		return status, nil
	}
B
bigsheeper 已提交
3882 3883 3884 3885 3886
	infoResp, err := node.queryCoord.LoadBalance(ctx, &querypb.LoadBalanceRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_LoadBalanceSegments,
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3887
			SourceID:  Params.ProxyCfg.GetNodeID(),
B
bigsheeper 已提交
3888 3889 3890
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3891
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3892
		SealedSegmentIDs: req.SealedSegmentIDs,
3893
		CollectionID:     collectionID,
B
bigsheeper 已提交
3894 3895 3896 3897 3898 3899 3900 3901 3902 3903 3904 3905 3906 3907 3908 3909 3910
	})
	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 已提交
3911
//GetCompactionState gets the compaction state of multiple segments
3912 3913 3914 3915 3916 3917 3918 3919 3920 3921 3922 3923 3924
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
}

3925
// ManualCompaction invokes compaction on specified collection
3926 3927 3928 3929 3930 3931 3932 3933 3934 3935 3936 3937 3938
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
}

3939
// GetCompactionStateWithPlans returns the compactions states with the given plan ID
3940 3941 3942 3943 3944 3945 3946 3947 3948 3949 3950 3951 3952
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 已提交
3953 3954 3955
// 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))
3956
	var err error
B
Bingyi Sun 已提交
3957 3958 3959 3960 3961 3962 3963
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3964
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3965 3966 3967 3968
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3969 3970 3971 3972
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3973 3974
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3975 3976
	code := node.stateCode.Load().(commonpb.StateCode)
	return code == commonpb.StateCode_Healthy
3977 3978
}

3979 3980 3981
func (node *Proxy) checkHealthyAndReturnCode() (commonpb.StateCode, bool) {
	code := node.stateCode.Load().(commonpb.StateCode)
	return code, code == commonpb.StateCode_Healthy
3982 3983
}

J
jingkl 已提交
3984
//unhealthyStatus returns the proxy not healthy status
3985 3986 3987
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3988
		Reason:    "proxy not healthy",
3989 3990
	}
}
G
groot 已提交
3991 3992 3993

// Import data files(json, numpy, etc.) on MinIO/S3 storage, read and parse them into sealed segments
func (node *Proxy) Import(ctx context.Context, req *milvuspb.ImportRequest) (*milvuspb.ImportResponse, error) {
3994 3995 3996
	log.Info("received import request",
		zap.String("collection name", req.GetCollectionName()),
		zap.Bool("row-based", req.GetRowBased()))
3997 3998 3999 4000 4001 4002
	resp := &milvuspb.ImportResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
	}
G
groot 已提交
4003 4004 4005 4006
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
4007
	// Call rootCoord to finish import.
4008 4009 4010 4011 4012 4013 4014 4015
	respFromRC, err := node.rootCoord.Import(ctx, req)
	if err != nil {
		log.Error("failed to execute bulk load request", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	return respFromRC, nil
G
groot 已提交
4016 4017
}

4018
// GetImportState checks import task state from RootCoord.
G
groot 已提交
4019 4020 4021 4022 4023 4024 4025 4026 4027 4028 4029 4030 4031 4032 4033 4034 4035 4036 4037 4038 4039 4040 4041 4042 4043 4044 4045
func (node *Proxy) GetImportState(ctx context.Context, req *milvuspb.GetImportStateRequest) (*milvuspb.GetImportStateResponse, error) {
	log.Info("received get import state request", zap.Int64("taskID", req.GetTask()))
	resp := &milvuspb.GetImportStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.rootCoord.GetImportState(ctx, req)
	log.Info("received get import state response", zap.Int64("taskID", req.GetTask()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

// ListImportTasks get id array of all import tasks from rootcoord
func (node *Proxy) ListImportTasks(ctx context.Context, req *milvuspb.ListImportTasksRequest) (*milvuspb.ListImportTasksResponse, error) {
	log.Info("received list import tasks request")
	resp := &milvuspb.ListImportTasksResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.rootCoord.ListImportTasks(ctx, req)
	log.Info("received list import tasks response")
	return resp, err
}

X
XuanYang-cn 已提交
4046 4047 4048 4049 4050 4051 4052 4053 4054
// GetReplicas gets replica info
func (node *Proxy) GetReplicas(ctx context.Context, req *milvuspb.GetReplicasRequest) (*milvuspb.GetReplicasResponse, error) {
	log.Info("received get replicas request")
	resp := &milvuspb.GetReplicasResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

4055 4056
	req.Base = &commonpb.MsgBase{
		MsgType:  commonpb.MsgType_GetReplicas,
X
Xiaofan 已提交
4057
		SourceID: Params.ProxyCfg.GetNodeID(),
4058 4059
	}

X
XuanYang-cn 已提交
4060 4061 4062 4063 4064
	resp, err := node.queryCoord.GetReplicas(ctx, req)
	log.Info("received get replicas response", zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

4065 4066 4067 4068 4069 4070
// InvalidateCredentialCache invalidate the credential cache of specified username.
func (node *Proxy) InvalidateCredentialCache(ctx context.Context, request *proxypb.InvalidateCredCacheRequest) (*commonpb.Status, error) {
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to invalidate credential cache",
		zap.String("role", typeutil.ProxyRole),
		zap.String("username", request.Username))
4071
	if !node.checkHealthy() {
4072
		return unhealthyStatus(), nil
4073
	}
4074 4075 4076 4077 4078 4079 4080 4081 4082 4083 4084 4085 4086 4087 4088 4089 4090 4091 4092 4093 4094

	username := request.Username
	if globalMetaCache != nil {
		globalMetaCache.RemoveCredential(username) // no need to return error, though credential may be not cached
	}
	logutil.Logger(ctx).Debug("complete to invalidate credential cache",
		zap.String("role", typeutil.ProxyRole),
		zap.String("username", request.Username))

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

// UpdateCredentialCache update the credential cache of specified username.
func (node *Proxy) UpdateCredentialCache(ctx context.Context, request *proxypb.UpdateCredCacheRequest) (*commonpb.Status, error) {
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to update credential cache",
		zap.String("role", typeutil.ProxyRole),
		zap.String("username", request.Username))
4095
	if !node.checkHealthy() {
4096
		return unhealthyStatus(), nil
4097
	}
4098 4099

	credInfo := &internalpb.CredentialInfo{
4100 4101
		Username:       request.Username,
		Sha256Password: request.Password,
4102 4103 4104 4105 4106 4107 4108 4109 4110 4111 4112 4113 4114 4115 4116
	}
	if globalMetaCache != nil {
		globalMetaCache.UpdateCredential(credInfo) // no need to return error, though credential may be not cached
	}
	logutil.Logger(ctx).Debug("complete to update credential cache",
		zap.String("role", typeutil.ProxyRole),
		zap.String("username", request.Username))

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

func (node *Proxy) CreateCredential(ctx context.Context, req *milvuspb.CreateCredentialRequest) (*commonpb.Status, error) {
4117 4118
	log.Debug("CreateCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4119
		return unhealthyStatus(), nil
4120
	}
4121 4122 4123 4124 4125 4126 4127 4128 4129 4130 4131 4132 4133 4134 4135 4136 4137 4138 4139 4140 4141 4142 4143 4144 4145 4146 4147 4148 4149 4150 4151
	// validate params
	username := req.Username
	if err := ValidateUsername(username); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
		}, nil
	}
	rawPassword, err := crypto.Base64Decode(req.Password)
	if err != nil {
		log.Error("decode password fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_CreateCredentialFailure,
			Reason:    "decode password fail key:" + req.Username,
		}, nil
	}
	if err = ValidatePassword(rawPassword); err != nil {
		log.Error("illegal password", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
		}, nil
	}
	encryptedPassword, err := crypto.PasswordEncrypt(rawPassword)
	if err != nil {
		log.Error("encrypt password fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_CreateCredentialFailure,
			Reason:    "encrypt password fail key:" + req.Username,
		}, nil
	}
4152

4153 4154 4155
	credInfo := &internalpb.CredentialInfo{
		Username:          req.Username,
		EncryptedPassword: encryptedPassword,
4156
		Sha256Password:    crypto.SHA256(rawPassword, req.Username),
4157 4158 4159 4160 4161 4162 4163 4164 4165 4166 4167 4168
	}
	result, err := node.rootCoord.CreateCredential(ctx, credInfo)
	if err != nil { // for error like conntext timeout etc.
		log.Error("create credential fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}
	return result, err
}

C
codeman 已提交
4169
func (node *Proxy) UpdateCredential(ctx context.Context, req *milvuspb.UpdateCredentialRequest) (*commonpb.Status, error) {
4170 4171
	log.Debug("UpdateCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4172
		return unhealthyStatus(), nil
4173
	}
C
codeman 已提交
4174 4175 4176 4177 4178 4179 4180 4181 4182
	rawOldPassword, err := crypto.Base64Decode(req.OldPassword)
	if err != nil {
		log.Error("decode old password fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "decode old password fail when updating:" + req.Username,
		}, nil
	}
	rawNewPassword, err := crypto.Base64Decode(req.NewPassword)
4183 4184 4185 4186 4187 4188 4189
	if err != nil {
		log.Error("decode password fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "decode password fail when updating:" + req.Username,
		}, nil
	}
C
codeman 已提交
4190 4191
	// valid new password
	if err = ValidatePassword(rawNewPassword); err != nil {
4192 4193 4194 4195 4196 4197
		log.Error("illegal password", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
		}, nil
	}
4198 4199

	if !passwordVerify(ctx, req.Username, rawOldPassword, globalMetaCache) {
C
codeman 已提交
4200 4201 4202 4203 4204 4205 4206
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "old password is not correct:" + req.Username,
		}, nil
	}
	// update meta data
	encryptedPassword, err := crypto.PasswordEncrypt(rawNewPassword)
4207 4208 4209 4210 4211 4212 4213
	if err != nil {
		log.Error("encrypt password fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "encrypt password fail when updating:" + req.Username,
		}, nil
	}
C
codeman 已提交
4214
	updateCredReq := &internalpb.CredentialInfo{
4215
		Username:          req.Username,
4216
		Sha256Password:    crypto.SHA256(rawNewPassword, req.Username),
4217 4218
		EncryptedPassword: encryptedPassword,
	}
C
codeman 已提交
4219
	result, err := node.rootCoord.UpdateCredential(ctx, updateCredReq)
4220 4221 4222 4223 4224 4225 4226 4227 4228 4229 4230
	if err != nil { // for error like conntext timeout etc.
		log.Error("update credential fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}
	return result, err
}

func (node *Proxy) DeleteCredential(ctx context.Context, req *milvuspb.DeleteCredentialRequest) (*commonpb.Status, error) {
4231 4232
	log.Debug("DeleteCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4233
		return unhealthyStatus(), nil
4234 4235
	}

4236 4237 4238 4239 4240 4241
	if req.Username == util.UserRoot {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_DeleteCredentialFailure,
			Reason:    "user root cannot be deleted",
		}, nil
	}
4242 4243 4244 4245 4246 4247 4248 4249 4250 4251 4252 4253
	result, err := node.rootCoord.DeleteCredential(ctx, req)
	if err != nil { // for error like conntext timeout etc.
		log.Error("delete credential fail", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}
	return result, err
}

func (node *Proxy) ListCredUsers(ctx context.Context, req *milvuspb.ListCredUsersRequest) (*milvuspb.ListCredUsersResponse, error) {
4254 4255
	log.Debug("ListCredUsers", zap.String("role", typeutil.ProxyRole))
	if !node.checkHealthy() {
4256
		return &milvuspb.ListCredUsersResponse{Status: unhealthyStatus()}, nil
4257
	}
4258 4259 4260 4261 4262 4263
	rootCoordReq := &milvuspb.ListCredUsersRequest{
		Base: &commonpb.MsgBase{
			MsgType: commonpb.MsgType_ListCredUsernames,
		},
	}
	resp, err := node.rootCoord.ListCredUsers(ctx, rootCoordReq)
4264 4265 4266 4267 4268 4269 4270 4271 4272 4273 4274 4275
	if err != nil {
		return &milvuspb.ListCredUsersResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
	return &milvuspb.ListCredUsersResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
4276
		Usernames: resp.Usernames,
4277 4278
	}, nil
}
4279

4280 4281 4282
func (node *Proxy) CreateRole(ctx context.Context, req *milvuspb.CreateRoleRequest) (*commonpb.Status, error) {
	logger.Debug("CreateRole", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4283
		return errorutil.UnhealthyStatus(code), nil
4284 4285 4286 4287 4288 4289 4290 4291 4292 4293
	}

	var roleName string
	if req.Entity != nil {
		roleName = req.Entity.Name
	}
	if err := ValidateRoleName(roleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4294
		}, nil
4295 4296 4297 4298 4299 4300 4301 4302
	}

	result, err := node.rootCoord.CreateRole(ctx, req)
	if err != nil {
		logger.Error("fail to create role", zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4303
		}, nil
4304 4305
	}
	return result, nil
4306 4307
}

4308 4309 4310
func (node *Proxy) DropRole(ctx context.Context, req *milvuspb.DropRoleRequest) (*commonpb.Status, error) {
	logger.Debug("DropRole", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4311
		return errorutil.UnhealthyStatus(code), nil
4312 4313 4314 4315 4316
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4317
		}, nil
4318
	}
4319 4320 4321 4322 4323
	if IsDefaultRole(req.RoleName) {
		errMsg := fmt.Sprintf("the role[%s] is a default role, which can't be droped", req.RoleName)
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    errMsg,
4324
		}, nil
4325
	}
4326 4327 4328 4329 4330 4331
	result, err := node.rootCoord.DropRole(ctx, req)
	if err != nil {
		logger.Error("fail to drop role", zap.String("role_name", req.RoleName), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4332
		}, nil
4333 4334
	}
	return result, nil
4335 4336
}

4337 4338 4339
func (node *Proxy) OperateUserRole(ctx context.Context, req *milvuspb.OperateUserRoleRequest) (*commonpb.Status, error) {
	logger.Debug("OperateUserRole", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4340
		return errorutil.UnhealthyStatus(code), nil
4341 4342 4343 4344 4345
	}
	if err := ValidateUsername(req.Username); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4346
		}, nil
4347 4348 4349 4350 4351
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4352
		}, nil
4353 4354 4355 4356 4357 4358 4359 4360
	}

	result, err := node.rootCoord.OperateUserRole(ctx, req)
	if err != nil {
		logger.Error("fail to operate user role", zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4361
		}, nil
4362 4363
	}
	return result, nil
4364 4365
}

4366 4367 4368
func (node *Proxy) SelectRole(ctx context.Context, req *milvuspb.SelectRoleRequest) (*milvuspb.SelectRoleResponse, error) {
	logger.Debug("SelectRole", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4369
		return &milvuspb.SelectRoleResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4370 4371 4372 4373 4374 4375 4376 4377 4378
	}

	if req.Role != nil {
		if err := ValidateRoleName(req.Role.Name); err != nil {
			return &milvuspb.SelectRoleResponse{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_IllegalArgument,
					Reason:    err.Error(),
				},
4379
			}, nil
4380 4381 4382 4383 4384 4385 4386 4387 4388 4389 4390
		}
	}

	result, err := node.rootCoord.SelectRole(ctx, req)
	if err != nil {
		logger.Error("fail to select role", zap.Error(err))
		return &milvuspb.SelectRoleResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
4391
		}, nil
4392 4393
	}
	return result, nil
4394 4395
}

4396 4397 4398
func (node *Proxy) SelectUser(ctx context.Context, req *milvuspb.SelectUserRequest) (*milvuspb.SelectUserResponse, error) {
	logger.Debug("SelectUser", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4399
		return &milvuspb.SelectUserResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4400 4401 4402 4403 4404 4405 4406 4407 4408
	}

	if req.User != nil {
		if err := ValidateUsername(req.User.Name); err != nil {
			return &milvuspb.SelectUserResponse{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_IllegalArgument,
					Reason:    err.Error(),
				},
4409
			}, nil
4410 4411 4412 4413 4414 4415 4416 4417 4418 4419 4420
		}
	}

	result, err := node.rootCoord.SelectUser(ctx, req)
	if err != nil {
		logger.Error("fail to select user", zap.Error(err))
		return &milvuspb.SelectUserResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
4421
		}, nil
4422 4423
	}
	return result, nil
4424 4425
}

4426 4427 4428 4429 4430 4431 4432 4433 4434 4435 4436 4437 4438 4439 4440 4441 4442 4443 4444 4445 4446 4447 4448 4449 4450 4451 4452 4453 4454 4455
func (node *Proxy) validPrivilegeParams(req *milvuspb.OperatePrivilegeRequest) error {
	if req.Entity == nil {
		return fmt.Errorf("the entity in the request is nil")
	}
	if req.Entity.Grantor == nil {
		return fmt.Errorf("the grantor entity in the grant entity is nil")
	}
	if req.Entity.Grantor.Privilege == nil {
		return fmt.Errorf("the privilege entity in the grantor entity is nil")
	}
	if err := ValidatePrivilege(req.Entity.Grantor.Privilege.Name); err != nil {
		return err
	}
	if req.Entity.Object == nil {
		return fmt.Errorf("the resource entity in the grant entity is nil")
	}
	if err := ValidateObjectType(req.Entity.Object.Name); err != nil {
		return err
	}
	if err := ValidateObjectName(req.Entity.ObjectName); err != nil {
		return err
	}
	if req.Entity.Role == nil {
		return fmt.Errorf("the object entity in the grant entity is nil")
	}
	if err := ValidateRoleName(req.Entity.Role.Name); err != nil {
		return err
	}

	return nil
4456 4457
}

4458 4459 4460
func (node *Proxy) OperatePrivilege(ctx context.Context, req *milvuspb.OperatePrivilegeRequest) (*commonpb.Status, error) {
	logger.Debug("OperatePrivilege", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4461
		return errorutil.UnhealthyStatus(code), nil
4462 4463 4464 4465 4466
	}
	if err := node.validPrivilegeParams(req); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4467
		}, nil
4468 4469 4470 4471 4472 4473
	}
	curUser, err := GetCurUserFromContext(ctx)
	if err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4474
		}, nil
4475 4476 4477 4478 4479 4480 4481 4482
	}
	req.Entity.Grantor.User = &milvuspb.UserEntity{Name: curUser}
	result, err := node.rootCoord.OperatePrivilege(ctx, req)
	if err != nil {
		logger.Error("fail to operate privilege", zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4483
		}, nil
4484 4485
	}
	return result, nil
4486 4487
}

4488 4489 4490 4491 4492 4493 4494 4495 4496 4497 4498 4499 4500 4501 4502 4503 4504 4505 4506 4507 4508 4509 4510 4511 4512 4513 4514 4515 4516
func (node *Proxy) validGrantParams(req *milvuspb.SelectGrantRequest) error {
	if req.Entity == nil {
		return fmt.Errorf("the grant entity in the request is nil")
	}

	if req.Entity.Object != nil {
		if err := ValidateObjectType(req.Entity.Object.Name); err != nil {
			return err
		}

		if err := ValidateObjectName(req.Entity.ObjectName); err != nil {
			return err
		}
	}

	if req.Entity.Role == nil {
		return fmt.Errorf("the role entity in the grant entity is nil")
	}

	if err := ValidateRoleName(req.Entity.Role.Name); err != nil {
		return err
	}

	return nil
}

func (node *Proxy) SelectGrant(ctx context.Context, req *milvuspb.SelectGrantRequest) (*milvuspb.SelectGrantResponse, error) {
	logger.Debug("SelectGrant", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
4517
		return &milvuspb.SelectGrantResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4518 4519 4520 4521 4522 4523 4524 4525
	}

	if err := node.validGrantParams(req); err != nil {
		return &milvuspb.SelectGrantResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_IllegalArgument,
				Reason:    err.Error(),
			},
4526
		}, nil
4527 4528 4529 4530 4531 4532 4533 4534 4535 4536
	}

	result, err := node.rootCoord.SelectGrant(ctx, req)
	if err != nil {
		logger.Error("fail to select grant", zap.Error(err))
		return &milvuspb.SelectGrantResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
4537
		}, nil
4538 4539 4540 4541 4542 4543 4544 4545 4546 4547 4548 4549 4550 4551 4552 4553 4554 4555 4556 4557 4558 4559 4560 4561 4562 4563 4564 4565
	}
	return result, nil
}

func (node *Proxy) RefreshPolicyInfoCache(ctx context.Context, req *proxypb.RefreshPolicyInfoCacheRequest) (*commonpb.Status, error) {
	logger.Debug("RefreshPrivilegeInfoCache", zap.Any("req", req))
	if code, ok := node.checkHealthyAndReturnCode(); !ok {
		return errorutil.UnhealthyStatus(code), errorutil.UnhealthyError()
	}

	if globalMetaCache != nil {
		err := globalMetaCache.RefreshPolicyInfo(typeutil.CacheOp{
			OpType: typeutil.CacheOpType(req.OpType),
			OpKey:  req.OpKey,
		})
		if err != nil {
			log.Error("fail to refresh policy info", zap.Error(err))
			return &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_RefreshPolicyInfoCacheFailure,
				Reason:    err.Error(),
			}, err
		}
	}
	logger.Debug("RefreshPrivilegeInfoCache success")

	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_Success,
	}, nil
4566
}
4567 4568 4569 4570 4571 4572 4573 4574 4575 4576 4577 4578 4579 4580 4581 4582 4583 4584 4585 4586 4587

// SetRates limits the rates of requests.
func (node *Proxy) SetRates(ctx context.Context, request *proxypb.SetRatesRequest) (*commonpb.Status, error) {
	log.Debug("SetRates", zap.String("role", typeutil.ProxyRole), zap.Any("rates", request.GetRates()))
	resp := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
	if !node.checkHealthy() {
		resp = unhealthyStatus()
		return resp, nil
	}

	err := node.multiRateLimiter.globalRateLimiter.setRates(request.GetRates())
	// TODO: set multiple rate limiter rates
	if err != nil {
		resp.Reason = err.Error()
		return resp, nil
	}
	resp.ErrorCode = commonpb.ErrorCode_Success
	return resp, nil
}