impl.go 160.1 KB
Newer Older
1 2 3 4 5 6
// Licensed to the LF AI & Data foundation under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
7 8
// with the License. You may obtain a copy of the License at
//
9
//     http://www.apache.org/licenses/LICENSE-2.0
10
//
11 12 13 14 15
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
16

C
Cai Yudong 已提交
17
package proxy
18 19 20

import (
	"context"
21
	"errors"
22
	"fmt"
C
cai.zhang 已提交
23
	"os"
24
	"strconv"
25 26 27
	"sync"

	"golang.org/x/sync/errgroup"
28

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

54 55
const moduleName = "Proxy"

56
// UpdateStateCode updates the state code of Proxy.
57
func (node *Proxy) UpdateStateCode(code commonpb.StateCode) {
58
	node.stateCode.Store(code)
Z
zhenshan.cao 已提交
59 60
}

61
// GetComponentStates get state of Proxy.
62 63
func (node *Proxy) GetComponentStates(ctx context.Context) (*milvuspb.ComponentStates, error) {
	stats := &milvuspb.ComponentStates{
G
godchen 已提交
64 65 66 67
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
	}
68
	code, ok := node.stateCode.Load().(commonpb.StateCode)
G
godchen 已提交
69 70 71 72 73 74
	if !ok {
		errMsg := "unexpected error in type assertion"
		stats.Status = &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    errMsg,
		}
G
godchen 已提交
75
		return stats, nil
G
godchen 已提交
76
	}
77 78 79 80
	nodeID := common.NotRegisteredID
	if node.session != nil && node.session.Registered() {
		nodeID = node.session.ServerID
	}
81
	info := &milvuspb.ComponentInfo{
82 83
		// NodeID:    Params.ProxyID, // will race with Proxy.Register()
		NodeID:    nodeID,
C
Cai Yudong 已提交
84
		Role:      typeutil.ProxyRole,
G
godchen 已提交
85 86 87 88 89 90
		StateCode: code,
	}
	stats.State = info
	return stats, nil
}

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

102
// InvalidateCollectionMetaCache invalidate the meta cache of specific collection.
C
Cai Yudong 已提交
103
func (node *Proxy) InvalidateCollectionMetaCache(ctx context.Context, request *proxypb.InvalidateCollMetaCacheRequest) (*commonpb.Status, error) {
104
	ctx = logutil.WithModule(ctx, moduleName)
105
	logutil.Logger(ctx).Info("received request to invalidate collection meta cache",
106
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
107
		zap.String("db", request.DbName),
108 109
		zap.String("collectionName", request.CollectionName),
		zap.Int64("collectionID", request.CollectionID))
D
dragondriver 已提交
110

111
	collectionName := request.CollectionName
112
	collectionID := request.CollectionID
X
Xiaofan 已提交
113 114

	var aliasName []string
N
neza2017 已提交
115
	if globalMetaCache != nil {
116 117 118 119
		if collectionName != "" {
			globalMetaCache.RemoveCollection(ctx, collectionName) // no need to return error, though collection may be not cached
		}
		if request.CollectionID != UniqueID(0) {
X
Xiaofan 已提交
120
			aliasName = globalMetaCache.RemoveCollectionsByID(ctx, collectionID)
121
		}
N
neza2017 已提交
122
	}
123 124
	if request.GetBase().GetMsgType() == commonpb.MsgType_DropCollection {
		// no need to handle error, since this Proxy may not create dml stream for the collection.
125 126 127
		node.chMgr.removeDMLStream(request.GetCollectionID())
		// clean up collection level metrics
		metrics.CleanupCollectionMetrics(Params.ProxyCfg.GetNodeID(), collectionName)
X
Xiaofan 已提交
128 129 130
		for _, alias := range aliasName {
			metrics.CleanupCollectionMetrics(Params.ProxyCfg.GetNodeID(), alias)
		}
131
	}
132
	logutil.Logger(ctx).Info("complete to invalidate collection meta cache",
133
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
134
		zap.String("db", request.DbName),
135 136
		zap.String("collection", collectionName),
		zap.Int64("collectionID", collectionID))
D
dragondriver 已提交
137

138
	return &commonpb.Status{
139
		ErrorCode: commonpb.ErrorCode_Success,
140 141
		Reason:    "",
	}, nil
142 143
}

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

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

157
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
158

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

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

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

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

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

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

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

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

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

247 248
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
249 250 251
	return cct.result, nil
}

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

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

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

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

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

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

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

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

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

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

331 332
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
333 334 335
	return dct.result, nil
}

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

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

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

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

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

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

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

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

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

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

422
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
423
		metrics.SuccessLabel).Inc()
424
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
425 426 427
	return hct.result, nil
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
437 438
	method := "LoadCollection"
	tr := timerecord.NewTimeRecorder(method)
439 440
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
441
	lct := &loadCollectionTask{
S
sunby 已提交
442
		ctx:                   ctx,
443 444
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
445
		queryCoord:            node.queryCoord,
C
cai.zhang 已提交
446
		indexCoord:            node.indexCoord,
447 448
	}

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

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

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

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

	if err := lct.WaitToFinish(); err != nil {
		log.Warn("LoadCollection failed to WaitToFinish",
D
dragondriver 已提交
482
			zap.Error(err),
483
			zap.String("traceID", traceID),
484
			zap.String("role", typeutil.ProxyRole),
485 486 487
			zap.Int64("MsgID", lct.ID()),
			zap.Uint64("BeginTS", lct.BeginTs()),
			zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
488 489
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))
490
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
491
			metrics.FailLabel).Inc()
492
		return &commonpb.Status{
493
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
494 495 496 497
			Reason:    err.Error(),
		}, nil
	}

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

507
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
508
		metrics.SuccessLabel).Inc()
509
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
510
	return lct.result, nil
511 512
}

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

519
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
520 521
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
522 523
	method := "ReleaseCollection"
	tr := timerecord.NewTimeRecorder(method)
524 525
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
526
	rct := &releaseCollectionTask{
S
sunby 已提交
527
		ctx:                      ctx,
528 529
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
530
		queryCoord:               node.queryCoord,
531
		chMgr:                    node.chMgr,
532 533
	}

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

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

550
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
551
			metrics.AbandonLabel).Inc()
552
		return &commonpb.Status{
553
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
554 555 556 557
			Reason:    err.Error(),
		}, nil
	}

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

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

580
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
581
			metrics.FailLabel).Inc()
582
		return &commonpb.Status{
583
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
584 585 586 587
			Reason:    err.Error(),
		}, nil
	}

588 589
	log.Debug(
		rpcDone(method),
590
		zap.String("traceID", traceID),
591
		zap.String("role", typeutil.ProxyRole),
592 593 594 595 596 597
		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))

598
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
599
		metrics.SuccessLabel).Inc()
600
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
601
	return rct.result, nil
602 603
}

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

612
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
613 614
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
615 616
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
617 618
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
619

620
	dct := &describeCollectionTask{
S
sunby 已提交
621
		ctx:                       ctx,
622 623
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
624
		rootCoord:                 node.rootCoord,
625 626
	}

627 628
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
629
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
630 631
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
632 633 634 635 636

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

641
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
642
			metrics.AbandonLabel).Inc()
643 644
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
645
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
646 647 648 649 650
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

671
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
672
			metrics.FailLabel).Inc()
673

674 675
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
676
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
677 678 679 680 681
				Reason:    err.Error(),
			},
		}, nil
	}

682 683
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
684
		zap.String("role", typeutil.ProxyRole),
685 686 687 688 689 690
		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))

691
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
692
		metrics.SuccessLabel).Inc()
693
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
694 695 696
	return dct.result, nil
}

697 698 699 700 701 702 703 704 705
// 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
	}

706
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetStatistics")
707 708 709 710
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
	method := "GetStatistics"
	tr := timerecord.NewTimeRecorder(method)
711 712
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
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
	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))

741
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
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
			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))

776
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796
			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))

797
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
798
		metrics.SuccessLabel).Inc()
799
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
800 801 802
	return g.result, nil
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
814 815
	method := "GetCollectionStatistics"
	tr := timerecord.NewTimeRecorder(method)
816 817
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
818
	g := &getCollectionStatisticsTask{
G
godchen 已提交
819 820 821
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
822
		dataCoord:                      node.dataCoord,
823 824
	}

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

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

841
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
842
			metrics.AbandonLabel).Inc()
843

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

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

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

874
		metrics.ProxyFunctionCall.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))

895
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
896
		metrics.SuccessLabel).Inc()
897
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
898
	return g.result, nil
899 900
}

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

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

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

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

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

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

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

969
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
970

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

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

988 989
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
990 991 992
	return sct.result, nil
}

J
jaime 已提交
993 994 995 996 997 998 999 1000 1001 1002 1003
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)

1004
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
J
jaime 已提交
1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028

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

1029
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
J
jaime 已提交
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
		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))

1059
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
J
jaime 已提交
1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075
		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))

1076 1077
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
J
jaime 已提交
1078 1079 1080
	return act.result, nil
}

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

1087
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
1088 1089
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1090 1091
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
1092
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1093

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

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

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

1120
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
1121

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

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

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

1152
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
1153

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

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

1171 1172
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1173 1174 1175
	return cpt.result, nil
}

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

1182
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
1183 1184
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1185 1186
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
1187
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1188

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

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

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

1216
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
1217

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

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

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

1248
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
1249

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

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

1267 1268
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1269 1270 1271
	return dpt.result, nil
}

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

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

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

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

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

1315
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1316
			metrics.AbandonLabel).Inc()
1317

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

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

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

1351
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1352
			metrics.FailLabel).Inc()
1353

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

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

1374
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1375
		metrics.SuccessLabel).Inc()
1376
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1377 1378 1379
	return hpt.result, nil
}

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

D
dragondriver 已提交
1386
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1387 1388
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1389 1390
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1391 1392
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
1393
	lpt := &loadPartitionsTask{
G
godchen 已提交
1394 1395 1396
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1397
		queryCoord:            node.queryCoord,
C
cai.zhang 已提交
1398
		indexCoord:            node.indexCoord,
1399 1400
	}

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

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

1419
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1420
			metrics.AbandonLabel).Inc()
1421

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

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

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

1452
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1453
			metrics.FailLabel).Inc()
1454

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

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

1472
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1473
		metrics.SuccessLabel).Inc()
1474
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1475
	return lpt.result, nil
1476 1477
}

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

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

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

1495
	method := "ReleasePartitions"
1496
	tr := timerecord.NewTimeRecorder(method)
1497 1498
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
1499 1500 1501
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1502
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1503 1504 1505
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1506 1507 1508 1509 1510 1511

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

1517
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1518
			metrics.AbandonLabel).Inc()
1519

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

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

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

1550
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1551
			metrics.FailLabel).Inc()
1552

1553
		return &commonpb.Status{
1554
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1555 1556 1557 1558
			Reason:    err.Error(),
		}, nil
	}

1559 1560 1561
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1562
		zap.String("role", typeutil.ProxyRole),
1563 1564 1565 1566 1567 1568 1569
		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))

1570
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1571
		metrics.SuccessLabel).Inc()
1572
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1573
	return rpt.result, nil
1574 1575
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetPartitionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1587 1588
	method := "GetPartitionStatistics"
	tr := timerecord.NewTimeRecorder(method)
1589 1590
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
1591

1592
	g := &getPartitionStatisticsTask{
1593 1594 1595
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1596
		dataCoord:                     node.dataCoord,
1597 1598
	}

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

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

1617
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1618
			metrics.AbandonLabel).Inc()
1619

1620 1621 1622 1623 1624 1625 1626 1627
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

1652
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1653
			metrics.FailLabel).Inc()
1654

1655 1656 1657 1658 1659 1660 1661 1662
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1663 1664 1665
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1666
		zap.String("role", typeutil.ProxyRole),
1667 1668 1669 1670 1671 1672 1673
		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))

1674
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1675
		metrics.SuccessLabel).Inc()
1676
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1677
	return g.result, nil
1678 1679
}

1680
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1681
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1682 1683 1684 1685 1686
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1687 1688 1689 1690 1691

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

1692
	spt := &showPartitionsTask{
G
godchen 已提交
1693 1694 1695
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1696
		rootCoord:             node.rootCoord,
1697
		queryCoord:            node.queryCoord,
G
godchen 已提交
1698
		result:                nil,
1699 1700
	}

1701
	method := "ShowPartitions"
1702 1703
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
1704
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1705
		metrics.TotalLabel).Inc()
1706 1707 1708 1709

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1710
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1711
		zap.Any("request", request))
1712 1713 1714 1715 1716 1717

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

1721
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1722
			metrics.AbandonLabel).Inc()
1723

G
godchen 已提交
1724
		return &milvuspb.ShowPartitionsResponse{
1725
			Status: &commonpb.Status{
1726
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1727 1728 1729 1730 1731
				Reason:    err.Error(),
			},
		}, nil
	}

1732 1733 1734
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1735
		zap.String("role", typeutil.ProxyRole),
1736 1737 1738
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1739 1740
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1741 1742 1743 1744 1745
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1746
			zap.Error(err),
1747
			zap.String("traceID", traceID),
1748
			zap.String("role", typeutil.ProxyRole),
1749 1750 1751 1752 1753 1754
			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 已提交
1755

1756
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1757
			metrics.FailLabel).Inc()
1758

G
godchen 已提交
1759
		return &milvuspb.ShowPartitionsResponse{
1760
			Status: &commonpb.Status{
1761
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1762 1763 1764 1765
				Reason:    err.Error(),
			},
		}, nil
	}
1766 1767 1768 1769

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1770
		zap.String("role", typeutil.ProxyRole),
1771 1772 1773 1774 1775 1776 1777
		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))

1778
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1779
		metrics.SuccessLabel).Inc()
1780
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1781 1782 1783
	return spt.result, nil
}

S
SimFG 已提交
1784 1785
func (node *Proxy) getCollectionProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest, collectionID int64) (int64, error) {
	resp, err := node.queryCoord.ShowCollections(ctx, &querypb.ShowCollectionsRequest{
S
smellthemoon 已提交
1786
		Base: commonpbutil.UpdateMsgBase(
1787 1788 1789
			request.Base,
			commonpbutil.WithMsgType(commonpb.MsgType_DescribeCollection),
		),
S
SimFG 已提交
1790 1791 1792 1793 1794
		CollectionIDs: []int64{collectionID},
	})
	if err != nil {
		return 0, err
	}
X
xige-16 已提交
1795 1796 1797 1798 1799

	if resp.Status.ErrorCode != commonpb.ErrorCode_Success {
		return 0, errors.New(resp.Status.Reason)
	}

S
SimFG 已提交
1800 1801 1802 1803 1804 1805 1806 1807 1808 1809 1810 1811 1812 1813 1814 1815 1816 1817
	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{
S
smellthemoon 已提交
1818
		Base: commonpbutil.UpdateMsgBase(
1819 1820 1821
			request.Base,
			commonpbutil.WithMsgType(commonpb.MsgType_ShowPartitions),
		),
S
SimFG 已提交
1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841 1842 1843 1844
		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)
1845
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetLoadingProgress")
S
SimFG 已提交
1846 1847
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1848
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
S
SimFG 已提交
1849 1850 1851 1852 1853 1854 1855 1856
	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))
1857
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
S
SimFG 已提交
1858 1859 1860 1861 1862 1863 1864 1865 1866 1867 1868 1869 1870 1871
		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
	}
1872 1873 1874 1875 1876 1877
	msgBase := commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
S
SimFG 已提交
1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888 1889 1890 1891 1892 1893 1894 1895 1896 1897 1898 1899 1900
	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))
1901 1902
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
S
SimFG 已提交
1903 1904 1905 1906 1907 1908 1909 1910
	return &milvuspb.GetLoadingProgressResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
		Progress: progress,
	}, nil
}

1911
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1912
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1913 1914 1915
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1916

1917
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateIndex")
D
dragondriver 已提交
1918 1919 1920
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

1921
	cit := &createIndexTask{
Z
zhenshan.cao 已提交
1922 1923 1924 1925 1926
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		req:        request,
		rootCoord:  node.rootCoord,
		indexCoord: node.indexCoord,
1927
		queryCoord: node.queryCoord,
1928 1929
	}

D
dragondriver 已提交
1930
	method := "CreateIndex"
1931
	tr := timerecord.NewTimeRecorder(method)
1932 1933
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1934 1935 1936
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1937
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1938 1939 1940 1941
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1942 1943 1944 1945 1946 1947

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

1954
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1955
			metrics.AbandonLabel).Inc()
1956

1957
		return &commonpb.Status{
1958
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1959 1960 1961 1962
			Reason:    err.Error(),
		}, nil
	}

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

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

1989
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1990
			metrics.FailLabel).Inc()
1991

1992
		return &commonpb.Status{
1993
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1994 1995 1996 1997
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1998 1999 2000
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2001
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2002 2003 2004 2005 2006 2007 2008 2009
		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))

2010
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2011
		metrics.SuccessLabel).Inc()
2012
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2013 2014 2015
	return cit.result, nil
}

2016
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
2017
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
2018 2019 2020 2021 2022
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2023 2024 2025 2026 2027

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

2028
	dit := &describeIndexTask{
S
sunby 已提交
2029
		ctx:                  ctx,
2030 2031
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
2032
		indexCoord:           node.indexCoord,
2033 2034
	}

2035 2036 2037
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
2038
	tr := timerecord.NewTimeRecorder(method)
2039 2040
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2041 2042 2043
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2044
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2045 2046 2047
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2048 2049 2050 2051 2052 2053 2054
		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),
2055
			zap.String("role", typeutil.ProxyRole),
2056 2057 2058 2059 2060
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

2061
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2062
			metrics.AbandonLabel).Inc()
2063

2064 2065
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
2066
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2067 2068 2069 2070 2071
				Reason:    err.Error(),
			},
		}, nil
	}

2072 2073 2074
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2075
		zap.String("role", typeutil.ProxyRole),
2076 2077 2078
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2079 2080 2081
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2082 2083 2084 2085 2086
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2087
			zap.Error(err),
2088
			zap.String("traceID", traceID),
2089
			zap.String("role", typeutil.ProxyRole),
2090 2091 2092
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2093 2094 2095
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
2096
			zap.String("index name", indexName))
D
dragondriver 已提交
2097

Z
zhenshan.cao 已提交
2098 2099 2100 2101
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
2102
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2103
			metrics.FailLabel).Inc()
2104

2105 2106
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
2107
				ErrorCode: errCode,
2108 2109 2110 2111 2112
				Reason:    err.Error(),
			},
		}, nil
	}

2113 2114 2115
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2116
		zap.String("role", typeutil.ProxyRole),
2117 2118 2119 2120 2121 2122 2123 2124
		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))

2125
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2126
		metrics.SuccessLabel).Inc()
2127
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2128 2129 2130
	return dit.result, nil
}

2131
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
2132
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
2133 2134 2135
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2136 2137 2138 2139 2140

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

2141
	dit := &dropIndexTask{
S
sunby 已提交
2142
		ctx:              ctx,
B
BossZou 已提交
2143 2144
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
2145
		indexCoord:       node.indexCoord,
2146
		queryCoord:       node.queryCoord,
B
BossZou 已提交
2147
	}
G
godchen 已提交
2148

D
dragondriver 已提交
2149
	method := "DropIndex"
2150
	tr := timerecord.NewTimeRecorder(method)
2151 2152
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2153 2154 2155 2156

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2157
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2158 2159 2160 2161 2162
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
2163 2164 2165 2166 2167
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2168
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2169 2170 2171 2172
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2173
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2174
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2175

B
BossZou 已提交
2176
		return &commonpb.Status{
2177
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2178 2179 2180
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2181

D
dragondriver 已提交
2182 2183 2184
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2185
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2186 2187 2188
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2189 2190 2191 2192
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
2193 2194 2195 2196

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2197
			zap.Error(err),
D
dragondriver 已提交
2198
			zap.String("traceID", traceID),
2199
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2200 2201 2202
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2203 2204 2205 2206 2207
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2208
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2209
			metrics.FailLabel).Inc()
2210

B
BossZou 已提交
2211
		return &commonpb.Status{
2212
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2213 2214 2215
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2216 2217 2218 2219

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2220
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2221 2222 2223 2224 2225 2226 2227 2228
		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))

2229
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2230
		metrics.SuccessLabel).Inc()
2231
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
2232 2233 2234
	return dit.result, nil
}

2235 2236
// 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.
2237
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2238
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
2239 2240 2241 2242 2243
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2244 2245 2246 2247 2248

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

2249
	gibpt := &getIndexBuildProgressTask{
2250 2251 2252
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
2253 2254
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
2255
		dataCoord:                    node.dataCoord,
2256 2257
	}

2258
	method := "GetIndexBuildProgress"
2259
	tr := timerecord.NewTimeRecorder(method)
2260 2261
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2262 2263 2264
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2265
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2266 2267 2268 2269
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2270 2271 2272 2273 2274 2275

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2276
			zap.String("role", typeutil.ProxyRole),
2277 2278 2279 2280
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2281
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2282
			metrics.AbandonLabel).Inc()
2283

2284 2285 2286 2287 2288 2289 2290 2291
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2292 2293 2294
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2295
		zap.String("role", typeutil.ProxyRole),
2296 2297 2298
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
2299 2300 2301 2302
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2303 2304 2305 2306

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2307
			zap.Error(err),
2308
			zap.String("traceID", traceID),
2309
			zap.String("role", typeutil.ProxyRole),
2310 2311 2312
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2313 2314 2315 2316
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2317
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2318
			metrics.FailLabel).Inc()
2319 2320 2321 2322 2323 2324 2325 2326

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2327 2328 2329 2330

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2331
		zap.String("role", typeutil.ProxyRole),
2332 2333 2334 2335 2336 2337 2338 2339
		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))
2340

2341
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2342
		metrics.SuccessLabel).Inc()
2343
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2344
	return gibpt.result, nil
2345 2346
}

2347
// GetIndexState get the build-state of index.
2348
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2349
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2350 2351 2352 2353 2354
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2355 2356 2357 2358 2359

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

2360
	dipt := &getIndexStateTask{
G
godchen 已提交
2361 2362 2363
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2364 2365
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2366 2367
	}

2368
	method := "GetIndexState"
2369
	tr := timerecord.NewTimeRecorder(method)
2370 2371
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2372 2373 2374
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2375
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2376 2377 2378 2379
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2380 2381 2382 2383 2384 2385

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2386
			zap.String("role", typeutil.ProxyRole),
2387 2388 2389 2390 2391
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2392
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2393
			metrics.AbandonLabel).Inc()
2394

G
godchen 已提交
2395
		return &milvuspb.GetIndexStateResponse{
2396
			Status: &commonpb.Status{
2397
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2398 2399 2400 2401 2402
				Reason:    err.Error(),
			},
		}, nil
	}

2403 2404 2405
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2406
		zap.String("role", typeutil.ProxyRole),
2407 2408 2409
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2410 2411 2412 2413
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2414 2415 2416 2417

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2418
			zap.Error(err),
2419
			zap.String("traceID", traceID),
2420
			zap.String("role", typeutil.ProxyRole),
2421 2422 2423
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2424 2425 2426 2427
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2428
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2429
			metrics.FailLabel).Inc()
2430

G
godchen 已提交
2431
		return &milvuspb.GetIndexStateResponse{
2432
			Status: &commonpb.Status{
2433
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2434 2435 2436 2437 2438
				Reason:    err.Error(),
			},
		}, nil
	}

2439 2440 2441
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2442
		zap.String("role", typeutil.ProxyRole),
2443 2444 2445 2446 2447 2448 2449 2450
		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))

2451
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2452
		metrics.SuccessLabel).Inc()
2453
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2454 2455 2456
	return dipt.result, nil
}

2457
// Insert insert records into collection.
C
Cai Yudong 已提交
2458
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2459 2460 2461 2462 2463 2464
	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))

2465 2466 2467 2468 2469
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2470 2471
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
2472
	receiveSize := proto.Size(request)
2473 2474
	rateCol.Add(internalpb.RateType_DMLInsert.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Add(float64(receiveSize))
D
dragondriver 已提交
2475

2476
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
2477
	it := &insertTask{
2478 2479
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2480
		// req:       request,
2481 2482 2483 2484
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2485
			InsertRequest: internalpb.InsertRequest{
2486 2487 2488 2489 2490
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Insert),
					commonpbutil.WithMsgID(0),
					commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
				),
2491 2492
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2493 2494 2495
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2496
				// RowData: transfer column based request to this
2497 2498
			},
		},
2499
		idAllocator:   node.rowIDAllocator,
2500 2501 2502
		segIDAssigner: node.segAssigner,
		chMgr:         node.chMgr,
		chTicker:      node.chTicker,
2503
	}
2504 2505

	if len(it.PartitionName) <= 0 {
2506
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2507 2508
	}

X
Xiangyu Wang 已提交
2509
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2510
		numRows := request.NumRows
2511 2512 2513 2514
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2515

X
Xiangyu Wang 已提交
2516 2517 2518 2519 2520 2521 2522
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2523 2524
	}

X
Xiangyu Wang 已提交
2525
	log.Debug("Enqueue insert request in Proxy",
2526
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2527 2528 2529 2530 2531
		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)),
2532 2533
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2534

X
Xiangyu Wang 已提交
2535 2536
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2537
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2538
			metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2539
		return constructFailedResponse(err), nil
2540
	}
D
dragondriver 已提交
2541

X
Xiangyu Wang 已提交
2542
	log.Debug("Detail of insert request in Proxy",
2543
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2544
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2545 2546 2547 2548 2549
		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 已提交
2550 2551 2552 2553 2554
		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))
2555
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2556
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2557 2558 2559 2560 2561
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2562
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2563 2564 2565 2566 2567 2568 2569 2570 2571 2572 2573
			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 已提交
2574
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2575

2576
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2577
		metrics.SuccessLabel).Inc()
2578 2579
	successCnt := it.result.InsertCnt - int64(len(it.result.ErrIndex))
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(successCnt))
2580
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2581
	metrics.ProxyCollectionMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
2582 2583 2584
	return it.result, nil
}

2585
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2586
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2587 2588 2589
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2590 2591
	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))
2592

2593
	receiveSize := proto.Size(request)
2594 2595
	rateCol.Add(internalpb.RateType_DMLDelete.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Add(float64(receiveSize))
2596

G
groot 已提交
2597 2598 2599 2600 2601 2602
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2603 2604 2605
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

2606
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2607
		metrics.TotalLabel).Inc()
2608
	dt := &deleteTask{
X
xige-16 已提交
2609 2610 2611
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		deleteExpr: request.Expr,
G
godchen 已提交
2612
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2613 2614 2615
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2616
			DeleteRequest: internalpb.DeleteRequest{
2617 2618 2619 2620
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Delete),
					commonpbutil.WithMsgID(0),
				),
X
xige-16 已提交
2621
				DbName:         request.DbName,
G
godchen 已提交
2622 2623 2624
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2625 2626 2627 2628
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2629 2630
	}

2631
	log.Debug("Enqueue delete request in Proxy",
2632
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2633 2634 2635 2636
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2637 2638 2639 2640

	// 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))
2641 2642
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
2643

G
groot 已提交
2644 2645 2646 2647 2648 2649 2650 2651
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2652
	log.Debug("Detail of delete request in Proxy",
2653
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2654 2655 2656 2657 2658
		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),
2659 2660
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2661

2662 2663
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
2664
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2665
			metrics.FailLabel).Inc()
G
groot 已提交
2666 2667 2668 2669 2670 2671 2672 2673
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2674
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2675
		metrics.SuccessLabel).Inc()
2676
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2677
	metrics.ProxyCollectionMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2678 2679 2680
	return dt.result, nil
}

2681
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2682
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2683 2684 2685 2686 2687
	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()))

2688 2689 2690 2691 2692
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2693 2694
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
2695
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2696
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2697

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

2701
	qt := &searchTask{
S
sunby 已提交
2702
		ctx:       ctx,
2703
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2704
		SearchRequest: &internalpb.SearchRequest{
2705 2706 2707 2708
			Base: commonpbutil.NewMsgBase(
				commonpbutil.WithMsgType(commonpb.MsgType_Search),
				commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
			),
2709
			ReqID: Params.ProxyCfg.GetNodeID(),
2710
		},
2711 2712 2713 2714
		request:  request,
		qc:       node.queryCoord,
		tr:       timerecord.NewTimeRecorder("search"),
		shardMgr: node.shardMgr,
2715 2716
	}

2717 2718 2719
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

Z
Zach 已提交
2720
	log.Ctx(ctx).Info(
2721
		rpcReceived(method),
2722
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2723 2724 2725 2726 2727
		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)),
2728 2729 2730 2731
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2732

2733
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
2734
		log.Ctx(ctx).Warn(
2735
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2736
			zap.Error(err),
2737
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2738 2739 2740 2741 2742 2743
			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),
2744 2745 2746
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2747

2748
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2749
			metrics.AbandonLabel).Inc()
2750

2751 2752
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2753
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2754 2755 2756 2757
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
2758
	tr.CtxRecord(ctx, "search request enqueue")
2759

Z
Zach 已提交
2760
	log.Ctx(ctx).Debug(
2761
		rpcEnqueued(method),
2762
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2763
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2764 2765 2766 2767 2768
		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),
2769
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2770 2771 2772 2773
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2774

2775
	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
2776
		log.Ctx(ctx).Warn(
2777
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2778
			zap.Error(err),
2779
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2780
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2781 2782 2783 2784
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2785
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2786 2787 2788 2789
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2790

2791
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2792
			metrics.FailLabel).Inc()
2793

2794 2795
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2796
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2797 2798 2799 2800 2801
				Reason:    err.Error(),
			},
		}, nil
	}

Z
Zach 已提交
2802
	span := tr.CtxRecord(ctx, "wait search result")
2803 2804
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel).Observe(float64(span.Milliseconds()))
2805
	tr.CtxRecord(ctx, "wait search result")
Z
Zach 已提交
2806
	log.Ctx(ctx).Debug(
2807
		rpcDone(method),
2808
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2809 2810 2811 2812 2813 2814
		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)),
2815 2816 2817 2818
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2819

2820
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2821 2822
		metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(qt.result.GetResults().GetNumQueries()))
C
cai.zhang 已提交
2823
	searchDur := tr.ElapseSpan().Milliseconds()
2824
	metrics.ProxySQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
2825
		metrics.SearchLabel).Observe(float64(searchDur))
2826 2827
	metrics.ProxyCollectionSQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel, request.CollectionName).Observe(float64(searchDur))
2828 2829 2830
	if qt.result != nil {
		sentSize := proto.Size(qt.result)
		metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
2831
		rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
2832
	}
2833 2834 2835
	return qt.result, nil
}

2836
// Flush notify data nodes to persist the data of collection.
2837 2838 2839 2840 2841 2842 2843
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2844
	if !node.checkHealthy() {
2845 2846
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2847
	}
D
dragondriver 已提交
2848 2849 2850 2851 2852

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

2853
	ft := &flushTask{
T
ThreadDao 已提交
2854 2855 2856
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2857
		dataCoord:    node.dataCoord,
2858 2859
	}

D
dragondriver 已提交
2860
	method := "Flush"
2861
	tr := timerecord.NewTimeRecorder(method)
2862
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2863 2864 2865 2866

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2867
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2868 2869
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2870 2871 2872 2873 2874 2875

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

2880
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
2881

2882 2883
		resp.Status.Reason = err.Error()
		return resp, nil
2884 2885
	}

D
dragondriver 已提交
2886 2887 2888
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2889
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2890 2891 2892
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2893 2894
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2895 2896 2897 2898

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

2908
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
2909

D
dragondriver 已提交
2910
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2911 2912
		resp.Status.Reason = err.Error()
		return resp, nil
2913 2914
	}

D
dragondriver 已提交
2915 2916 2917
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2918
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2919 2920 2921 2922 2923 2924
		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))

2925 2926
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2927
	return ft.result, nil
2928 2929
}

2930
// Query get the records by primary keys.
C
Cai Yudong 已提交
2931
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2932 2933 2934 2935 2936
	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)

2937 2938 2939 2940 2941
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2942

D
dragondriver 已提交
2943 2944
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
2945
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2946

2947
	qt := &queryTask{
2948 2949 2950
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
2951 2952 2953 2954
			Base: commonpbutil.NewMsgBase(
				commonpbutil.WithMsgType(commonpb.MsgType_Retrieve),
				commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
			),
2955
			ReqID: Params.ProxyCfg.GetNodeID(),
2956
		},
2957 2958
		request:          request,
		qc:               node.queryCoord,
2959
		queryShardPolicy: mergeRoundRobinPolicy,
2960
		shardMgr:         node.shardMgr,
2961 2962
	}

D
dragondriver 已提交
2963 2964
	method := "Query"

2965
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2966 2967
		metrics.TotalLabel).Inc()

Z
Zach 已提交
2968
	log.Ctx(ctx).Info(
D
dragondriver 已提交
2969
		rpcReceived(method),
2970
		zap.String("role", typeutil.ProxyRole),
2971 2972
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2973 2974 2975 2976 2977
		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 已提交
2978

D
dragondriver 已提交
2979
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
2980
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
2981 2982 2983
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("role", typeutil.ProxyRole),
2984 2985 2986
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2987

2988 2989
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
2990

2991 2992 2993 2994 2995 2996
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2997
	}
Z
Zach 已提交
2998
	tr.CtxRecord(ctx, "query request enqueue")
2999

Z
Zach 已提交
3000
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3001
		rpcEnqueued(method),
3002
		zap.String("role", typeutil.ProxyRole),
3003
		zap.Int64("msgID", qt.ID()),
3004 3005
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
3006
		zap.Strings("partitions", request.PartitionNames))
D
dragondriver 已提交
3007 3008

	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
3009
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
3010 3011
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
3012
			zap.String("role", typeutil.ProxyRole),
3013
			zap.Int64("msgID", qt.ID()),
3014 3015 3016
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
3017

3018
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3019
			metrics.FailLabel).Inc()
3020

3021 3022 3023 3024 3025 3026 3027
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
3028
	span := tr.CtxRecord(ctx, "wait query result")
3029 3030
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel).Observe(float64(span.Milliseconds()))
3031

Z
Zach 已提交
3032
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3033 3034
		rpcDone(method),
		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))
D
dragondriver 已提交
3039

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

3043
	metrics.ProxySQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
3044
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
3045 3046
	metrics.ProxyCollectionSQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
3047 3048

	ret := &milvuspb.QueryResults{
3049 3050
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
3051 3052
	}
	sentSize := proto.Size(qt.result)
3053
	rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
3054 3055
	metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
	return ret, nil
3056
}
3057

3058
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
3059 3060 3061 3062
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3063 3064 3065 3066 3067

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

Y
Yusup 已提交
3068 3069 3070 3071 3072 3073 3074
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
3075
	method := "CreateAlias"
3076
	tr := timerecord.NewTimeRecorder(method)
3077
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3078 3079 3080 3081 3082 3083 3084 3085 3086 3087 3088 3089 3090 3091 3092 3093 3094 3095 3096

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

3097
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
3098

Y
Yusup 已提交
3099 3100 3101 3102 3103 3104
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3105 3106 3107
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3108
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3109 3110 3111 3112
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3113 3114
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3115 3116 3117 3118

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3119
			zap.Error(err),
D
dragondriver 已提交
3120
			zap.String("traceID", traceID),
3121
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3122 3123 3124 3125
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3126 3127
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
3128
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
3129 3130 3131 3132 3133 3134 3135

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

D
dragondriver 已提交
3136 3137 3138 3139 3140 3141 3142 3143 3144 3145 3146
	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))

3147 3148
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3149 3150 3151
	return cat.result, nil
}

3152
// DropAlias alter the alias of collection.
Y
Yusup 已提交
3153 3154 3155 3156
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3157 3158 3159 3160 3161

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

Y
Yusup 已提交
3162 3163 3164 3165 3166 3167 3168
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
3169
	method := "DropAlias"
3170
	tr := timerecord.NewTimeRecorder(method)
3171
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3172 3173 3174 3175 3176 3177 3178 3179 3180 3181 3182 3183 3184 3185 3186 3187

	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))
3188
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
3189

Y
Yusup 已提交
3190 3191 3192 3193 3194 3195
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3196 3197 3198
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3199
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3200 3201 3202 3203
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3204
		zap.String("alias", request.Alias))
D
dragondriver 已提交
3205 3206 3207 3208

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3209
			zap.Error(err),
D
dragondriver 已提交
3210
			zap.String("traceID", traceID),
3211
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3212 3213 3214 3215
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3216 3217
			zap.String("alias", request.Alias))

3218
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3219

Y
Yusup 已提交
3220 3221 3222 3223 3224 3225
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3226 3227 3228 3229 3230 3231 3232 3233 3234 3235
	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))

3236 3237
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3238 3239 3240
	return dat.result, nil
}

3241
// AlterAlias alter alias of collection.
Y
Yusup 已提交
3242 3243 3244 3245
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3246 3247 3248 3249 3250

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

Y
Yusup 已提交
3251 3252 3253 3254 3255 3256 3257
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
3258
	method := "AlterAlias"
3259
	tr := timerecord.NewTimeRecorder(method)
3260
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3261 3262 3263 3264 3265 3266 3267 3268 3269 3270 3271 3272 3273 3274 3275 3276 3277 3278

	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))
3279
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
3280

Y
Yusup 已提交
3281 3282 3283 3284 3285 3286
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3287 3288 3289
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3290
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3291 3292 3293 3294
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3295 3296
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3297 3298 3299 3300

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3301
			zap.Error(err),
D
dragondriver 已提交
3302
			zap.String("traceID", traceID),
3303
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3304 3305 3306 3307
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3308 3309 3310
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

3311
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3312

Y
Yusup 已提交
3313 3314 3315 3316 3317 3318
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3319 3320 3321 3322 3323 3324 3325 3326 3327 3328 3329
	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))

3330 3331
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3332 3333 3334
	return aat.result, nil
}

3335
// CalcDistance calculates the distances between vectors.
3336
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3337 3338 3339 3340 3341
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3342

3343 3344 3345 3346
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3347 3348
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3349

3350 3351 3352 3353 3354
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3355 3356
		}

3357
		qt := &queryTask{
3358 3359 3360
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
3361 3362 3363 3364
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Retrieve),
					commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
				),
3365
				ReqID: Params.ProxyCfg.GetNodeID(),
3366
			},
3367 3368 3369 3370
			request: queryRequest,
			qc:      node.queryCoord,
			ids:     ids.IdArray,

3371
			queryShardPolicy: mergeRoundRobinPolicy,
3372
			shardMgr:         node.shardMgr,
3373 3374
		}

G
groot 已提交
3375 3376 3377 3378 3379 3380
		items := []zapcore.Field{
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields),
		}

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

3385 3386 3387 3388 3389
			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3390
			}, err
3391
		}
3392

G
groot 已提交
3393
		log.Debug("CalcDistance queryTask enqueued", items...)
3394 3395 3396

		err = qt.WaitToFinish()
		if err != nil {
G
groot 已提交
3397
			log.Error("CalcDistance queryTask failed to WaitToFinish", append(items, zap.Error(err))...)
3398 3399 3400 3401 3402 3403

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

G
groot 已提交
3407
		log.Debug("CalcDistance queryTask Done", items...)
3408 3409

		return &milvuspb.QueryResults{
3410 3411
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3412 3413 3414
		}, nil
	}

G
groot 已提交
3415 3416 3417 3418
	// calcDistanceTask is not a standard task, no need to enqueue
	task := &calcDistanceTask{
		traceID:   traceID,
		queryFunc: query,
3419 3420
	}

G
groot 已提交
3421
	return task.Execute(ctx, request)
3422 3423
}

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

3429
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3430
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3431
	log.Debug("GetPersistentSegmentInfo",
3432
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3433 3434 3435
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3436
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3437
		Status: &commonpb.Status{
3438
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3439 3440
		},
	}
3441 3442 3443 3444
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3445 3446
	method := "GetPersistentSegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
3447
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3448
		metrics.TotalLabel).Inc()
3449 3450 3451

	// list segments
	collectionID, err := globalMetaCache.GetCollectionID(ctx, req.GetCollectionName())
X
XuanYang-cn 已提交
3452
	if err != nil {
3453 3454 3455 3456 3457 3458 3459 3460 3461 3462 3463 3464 3465
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
		resp.Status.Reason = fmt.Errorf("getCollectionID failed, err:%w", err).Error()
		return resp, nil
	}

	getSegmentsByStatesResponse, err := node.dataCoord.GetSegmentsByStates(ctx, &datapb.GetSegmentsByStatesRequest{
		CollectionID: collectionID,
		// -1 means list all partition segemnts
		PartitionID: -1,
		States:      []commonpb.SegmentState{commonpb.SegmentState_Flushing, commonpb.SegmentState_Flushed, commonpb.SegmentState_Sealed},
	})
	if err != nil {
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3466
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3467 3468
		return resp, nil
	}
3469 3470

	// get Segment info
3471
	infoResp, err := node.dataCoord.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
3472 3473 3474 3475 3476 3477
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_SegmentInfo),
			commonpbutil.WithMsgID(0),
			commonpbutil.WithTimeStamp(0),
			commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
		),
3478
		SegmentIDs: getSegmentsByStatesResponse.Segments,
X
XuanYang-cn 已提交
3479 3480
	})
	if err != nil {
3481 3482 3483
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()
		log.Warn("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 {
3489 3490
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.FailLabel).Inc()
X
XuanYang-cn 已提交
3491 3492 3493 3494 3495 3496
		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 已提交
3497
			SegmentID:    info.ID,
X
XuanYang-cn 已提交
3498 3499
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
S
sunby 已提交
3500
			NumRows:      info.NumOfRows,
X
XuanYang-cn 已提交
3501 3502 3503
			State:        info.State,
		}
	}
3504
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3505
		metrics.SuccessLabel).Inc()
3506
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
3507
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
X
XuanYang-cn 已提交
3508 3509 3510 3511
	resp.Infos = persistentInfos
	return resp, nil
}

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

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

3529 3530 3531 3532 3533
	method := "GetQuerySegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()

3534 3535
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
3536
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3537 3538 3539
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3540
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
3541 3542 3543 3544 3545 3546
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_SegmentInfo),
			commonpbutil.WithMsgID(0),
			commonpbutil.WithTimeStamp(0),
			commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
		),
3547
		CollectionID: collID,
Z
zhenshan.cao 已提交
3548 3549
	})
	if err != nil {
3550 3551
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
		log.Error("Failed to get segment info from QueryCoord", zap.Error(err))
Z
zhenshan.cao 已提交
3552 3553 3554
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3555
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3556
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3557
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3558
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3559 3560 3561 3562 3563 3564 3565 3566 3567 3568 3569 3570 3571
		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 已提交
3572
			State:        info.SegmentState,
3573
			NodeIds:      info.NodeIds,
Z
zhenshan.cao 已提交
3574 3575
		}
	}
3576 3577 3578

	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
3579
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3580 3581 3582 3583
	resp.Infos = queryInfos
	return resp, nil
}

J
jingkl 已提交
3584
// Dummy handles dummy request
C
Cai Yudong 已提交
3585
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3586 3587 3588 3589 3590 3591 3592 3593 3594 3595 3596
	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
	}

3597 3598
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3599
		if err != nil {
3600
			log.Debug("Failed to parse dummy query request")
3601 3602 3603
			return failedResponse, nil
		}

3604
		request := &milvuspb.QueryRequest{
3605 3606 3607
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3608
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3609 3610
		}

3611
		_, err = node.Query(ctx, request)
3612
		if err != nil {
3613
			log.Debug("Failed to execute dummy query")
3614 3615
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3616 3617 3618 3619 3620 3621

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

3622 3623
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3624 3625
}

J
jingkl 已提交
3626
// RegisterLink registers a link
C
Cai Yudong 已提交
3627
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
3628
	code := node.stateCode.Load().(commonpb.StateCode)
D
dragondriver 已提交
3629
	log.Debug("RegisterLink",
3630
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3631
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3632

3633
	if code != commonpb.StateCode_Healthy {
3634 3635 3636
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3637
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3638
				Reason:    "proxy not healthy",
3639 3640 3641
			},
		}, nil
	}
X
Xiaofan 已提交
3642
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Inc()
3643 3644 3645
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3646
			ErrorCode: commonpb.ErrorCode_Success,
3647
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3648 3649 3650
		},
	}, nil
}
3651

3652
// GetMetrics gets the metrics of proxy
3653 3654 3655
// 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 已提交
3656
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3657 3658 3659 3660
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
X
Xiaofan 已提交
3661
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3662
			zap.String("req", req.Request),
X
Xiaofan 已提交
3663
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.GetNodeID())))
3664 3665 3666 3667

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
Xiaofan 已提交
3668
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.GetNodeID()),
3669 3670 3671 3672 3673 3674 3675 3676
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
X
Xiaofan 已提交
3677
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3678 3679 3680 3681 3682 3683 3684 3685 3686 3687 3688 3689 3690 3691 3692
			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))

3693 3694 3695 3696 3697 3698
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3699
	if metricType == metricsinfo.SystemInfoMetrics {
3700 3701 3702 3703 3704 3705 3706
		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))

3707
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3708 3709

		log.Debug("Proxy.GetMetrics",
X
Xiaofan 已提交
3710
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3711 3712 3713 3714 3715
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3716 3717
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3718
		return metrics, nil
3719 3720 3721
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
X
Xiaofan 已提交
3722
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3723 3724 3725 3726 3727 3728 3729 3730 3731 3732 3733 3734
		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
}

3735 3736 3737 3738 3739 3740 3741 3742 3743 3744 3745 3746 3747 3748 3749 3750 3751 3752 3753 3754 3755 3756 3757 3758 3759 3760 3761 3762 3763 3764 3765 3766
// 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) {
	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
	}

3767 3768 3769 3770 3771 3772
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3773 3774 3775 3776 3777 3778 3779 3780 3781 3782 3783 3784 3785 3786 3787 3788 3789 3790 3791 3792

	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),
3793
			zap.String("metric_type", metricType))
3794 3795 3796 3797 3798 3799 3800 3801 3802 3803 3804 3805 3806 3807 3808 3809 3810

		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 已提交
3811 3812 3813
// 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 已提交
3814
		zap.Int64("proxy_id", Params.ProxyCfg.GetNodeID()),
B
bigsheeper 已提交
3815 3816 3817 3818 3819 3820 3821 3822 3823
		zap.Any("req", req))

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

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
3824 3825 3826 3827 3828 3829 3830

	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 已提交
3831
	infoResp, err := node.queryCoord.LoadBalance(ctx, &querypb.LoadBalanceRequest{
3832 3833 3834 3835 3836 3837
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_LoadBalanceSegments),
			commonpbutil.WithMsgID(0),
			commonpbutil.WithTimeStamp(0),
			commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
		),
B
bigsheeper 已提交
3838 3839
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3840
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3841
		SealedSegmentIDs: req.SealedSegmentIDs,
3842
		CollectionID:     collectionID,
B
bigsheeper 已提交
3843 3844 3845 3846 3847 3848 3849 3850 3851 3852 3853 3854 3855 3856 3857 3858 3859
	})
	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
}

3860 3861 3862 3863 3864 3865 3866 3867 3868
// 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
	}

S
smellthemoon 已提交
3869 3870 3871 3872
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_GetReplicas),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3873 3874 3875 3876 3877 3878 3879 3880 3881 3882 3883 3884

	resp, err := node.queryCoord.GetReplicas(ctx, req)
	if err != nil {
		log.Error("Failed to get replicas from Query Coordinator", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	log.Info("received get replicas response", zap.Any("resp", resp), zap.Error(err))
	return resp, nil
}

3885
// GetCompactionState gets the compaction state of multiple segments
3886 3887 3888 3889 3890 3891 3892 3893 3894 3895 3896 3897 3898
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
}

3899
// ManualCompaction invokes compaction on specified collection
3900 3901 3902 3903 3904 3905 3906 3907 3908 3909 3910 3911 3912
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
}

3913
// GetCompactionStateWithPlans returns the compactions states with the given plan ID
3914 3915 3916 3917 3918 3919 3920 3921 3922 3923 3924 3925 3926
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 已提交
3927 3928 3929
// 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))
3930
	var err error
B
Bingyi Sun 已提交
3931 3932 3933 3934 3935 3936 3937
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3938
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3939 3940 3941 3942
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3943 3944 3945 3946
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3947 3948
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3949 3950
	code := node.stateCode.Load().(commonpb.StateCode)
	return code == commonpb.StateCode_Healthy
3951 3952
}

3953 3954 3955
func (node *Proxy) checkHealthyAndReturnCode() (commonpb.StateCode, bool) {
	code := node.stateCode.Load().(commonpb.StateCode)
	return code, code == commonpb.StateCode_Healthy
3956 3957
}

3958
// unhealthyStatus returns the proxy not healthy status
3959 3960 3961
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3962
		Reason:    "proxy not healthy",
3963 3964
	}
}
G
groot 已提交
3965 3966 3967

// 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) {
3968 3969 3970
	log.Info("received import request",
		zap.String("collection name", req.GetCollectionName()),
		zap.Bool("row-based", req.GetRowBased()))
3971 3972 3973 3974 3975 3976
	resp := &milvuspb.ImportResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
	}
G
groot 已提交
3977 3978 3979 3980
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3981

3982 3983 3984 3985 3986 3987 3988 3989
	err := importutil.ValidateOptions(req.GetOptions())
	if err != nil {
		log.Error("failed to execute import request", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = "request options is not illegal    \n" + err.Error() + "    \nIllegal option format    \n" + importutil.OptionFormat
		return resp, nil
	}

3990 3991 3992 3993 3994
	method := "Import"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()

3995
	// Call rootCoord to finish import.
3996 3997
	respFromRC, err := node.rootCoord.Import(ctx, req)
	if err != nil {
3998
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3999 4000 4001 4002 4003
		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
	}
4004 4005 4006

	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
4007
	return respFromRC, nil
G
groot 已提交
4008 4009
}

4010
// GetImportState checks import task state from RootCoord.
G
groot 已提交
4011 4012 4013 4014 4015 4016 4017
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
	}
4018 4019 4020 4021
	method := "GetImportState"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
G
groot 已提交
4022 4023

	resp, err := node.rootCoord.GetImportState(ctx, req)
4024 4025 4026 4027 4028 4029 4030 4031 4032 4033 4034 4035
	if err != nil {
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
		log.Error("failed to execute get import state", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, nil
	}

	log.Info("successfully received get import state response", zap.Int64("taskID", req.GetTask()), zap.Any("resp", resp), zap.Error(err))
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
	return resp, nil
G
groot 已提交
4036 4037 4038 4039 4040 4041 4042 4043 4044 4045
}

// 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
	}
4046 4047 4048 4049
	method := "ListImportTasks"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
G
groot 已提交
4050
	resp, err := node.rootCoord.ListImportTasks(ctx, req)
4051 4052 4053 4054 4055
	if err != nil {
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
		log.Error("failed to execute list import tasks", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
X
XuanYang-cn 已提交
4056 4057 4058
		return resp, nil
	}

4059 4060 4061
	log.Info("successfully received list import tasks response", zap.String("collection", req.CollectionName), zap.Any("tasks", resp.Tasks))
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
X
XuanYang-cn 已提交
4062 4063 4064
	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
	rootCoordReq := &milvuspb.ListCredUsersRequest{
4259 4260 4261
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_ListCredUsernames),
		),
4262 4263
	}
	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

// SetRates limits the rates of requests.
func (node *Proxy) SetRates(ctx context.Context, request *proxypb.SetRatesRequest) (*commonpb.Status, error) {
	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
}
4587 4588 4589 4590

func (node *Proxy) CheckHealth(ctx context.Context, request *milvuspb.CheckHealthRequest) (*milvuspb.CheckHealthResponse, error) {
	if !node.checkHealthy() {
		reason := errorutil.UnHealthReason("proxy", node.session.ServerID, "proxy is unhealthy")
4591 4592 4593 4594
		return &milvuspb.CheckHealthResponse{
			Status:    unhealthyStatus(),
			IsHealthy: false,
			Reasons:   []string{reason}}, nil
4595 4596 4597 4598 4599 4600 4601 4602 4603 4604 4605
	}

	group, ctx := errgroup.WithContext(ctx)
	errReasons := make([]string, 0)

	mu := &sync.Mutex{}
	fn := func(role string, resp *milvuspb.CheckHealthResponse, err error) error {
		mu.Lock()
		defer mu.Unlock()

		if err != nil {
4606
			log.Warn("check health fail", zap.String("role", role), zap.Error(err))
4607 4608 4609 4610 4611
			errReasons = append(errReasons, fmt.Sprintf("check health fail for %s", role))
			return err
		}

		if !resp.IsHealthy {
4612
			log.Warn("check health fail", zap.String("role", role))
4613 4614 4615 4616 4617 4618 4619 4620 4621 4622 4623 4624 4625 4626 4627 4628 4629 4630 4631 4632 4633 4634 4635 4636 4637 4638 4639 4640 4641 4642 4643 4644 4645
			errReasons = append(errReasons, resp.Reasons...)
		}
		return nil
	}

	group.Go(func() error {
		resp, err := node.rootCoord.CheckHealth(ctx, request)
		return fn("rootcoord", resp, err)
	})

	group.Go(func() error {
		resp, err := node.queryCoord.CheckHealth(ctx, request)
		return fn("querycoord", resp, err)
	})

	group.Go(func() error {
		resp, err := node.dataCoord.CheckHealth(ctx, request)
		return fn("datacoord", resp, err)
	})

	group.Go(func() error {
		resp, err := node.indexCoord.CheckHealth(ctx, request)
		return fn("indexcoord", resp, err)
	})

	err := group.Wait()
	if err != nil || len(errReasons) != 0 {
		return &milvuspb.CheckHealthResponse{
			IsHealthy: false,
			Reasons:   errReasons,
		}, nil
	}

4646 4647 4648 4649 4650 4651 4652
	return &milvuspb.CheckHealthResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
		IsHealthy: true,
	}, nil
4653
}