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

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

import (
	"context"
21
	"errors"
22
	"fmt"
C
cai.zhang 已提交
23
	"os"
24
	"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 1928
	}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

2924 2925
	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()))
2926
	return ft.result, nil
2927 2928
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Z
Zach 已提交
3031
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3032 3033
		rpcDone(method),
		zap.String("role", typeutil.ProxyRole),
3034
		zap.Int64("msgID", qt.ID()),
3035 3036 3037
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
3038

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

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

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

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

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

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

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

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

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

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

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

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

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

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

3146 3147
	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 已提交
3148 3149 3150
	return cat.result, nil
}

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

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

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

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

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

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

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

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

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

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

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

3235 3236
	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 已提交
3237 3238 3239
	return dat.result, nil
}

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

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

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

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

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

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

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

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

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

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

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

3329 3330
	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 已提交
3331 3332 3333
	return aat.result, nil
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

	// list segments
	collectionID, err := globalMetaCache.GetCollectionID(ctx, req.GetCollectionName())
X
XuanYang-cn 已提交
3451
	if err != nil {
3452 3453 3454 3455 3456 3457 3458 3459 3460 3461 3462 3463 3464
		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()
3465
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3466 3467
		return resp, nil
	}
3468 3469

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

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

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

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

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

	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()))
3578
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3579 3580 3581 3582
	resp.Infos = queryInfos
	return resp, nil
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

3715 3716
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

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

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

3734 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
// 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
	}

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

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

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

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

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

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

3859 3860 3861 3862 3863 3864 3865 3866 3867
// 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 已提交
3868 3869 3870 3871
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_GetReplicas),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3872 3873 3874 3875 3876 3877 3878 3879 3880 3881 3882 3883

	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
}

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

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

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

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

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

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

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

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

3981 3982 3983 3984 3985 3986 3987 3988
	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
	}

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

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

	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()))
4006
	return respFromRC, nil
G
groot 已提交
4007 4008
}

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

	resp, err := node.rootCoord.GetImportState(ctx, req)
4023 4024 4025 4026 4027 4028 4029 4030 4031 4032 4033 4034
	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 已提交
4035 4036 4037 4038 4039 4040 4041 4042 4043 4044
}

// 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
	}
4045 4046 4047 4048
	method := "ListImportTasks"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
G
groot 已提交
4049
	resp, err := node.rootCoord.ListImportTasks(ctx, req)
4050 4051 4052 4053 4054
	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 已提交
4055 4056 4057
		return resp, nil
	}

4058 4059 4060
	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 已提交
4061 4062 4063
	return resp, err
}

4064 4065 4066 4067 4068 4069
// 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))
4070
	if !node.checkHealthy() {
4071
		return unhealthyStatus(), nil
4072
	}
4073 4074 4075 4076 4077 4078 4079 4080 4081 4082 4083 4084 4085 4086 4087 4088 4089 4090 4091 4092 4093

	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))
4094
	if !node.checkHealthy() {
4095
		return unhealthyStatus(), nil
4096
	}
4097 4098

	credInfo := &internalpb.CredentialInfo{
4099 4100
		Username:       request.Username,
		Sha256Password: request.Password,
4101 4102 4103 4104 4105 4106 4107 4108 4109 4110 4111 4112 4113 4114 4115
	}
	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) {
4116 4117
	log.Debug("CreateCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4118
		return unhealthyStatus(), nil
4119
	}
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
	// 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
	}
4151

4152 4153 4154
	credInfo := &internalpb.CredentialInfo{
		Username:          req.Username,
		EncryptedPassword: encryptedPassword,
4155
		Sha256Password:    crypto.SHA256(rawPassword, req.Username),
4156 4157 4158 4159 4160 4161 4162 4163 4164 4165 4166 4167
	}
	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 已提交
4168
func (node *Proxy) UpdateCredential(ctx context.Context, req *milvuspb.UpdateCredentialRequest) (*commonpb.Status, error) {
4169 4170
	log.Debug("UpdateCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4171
		return unhealthyStatus(), nil
4172
	}
C
codeman 已提交
4173 4174 4175 4176 4177 4178 4179 4180 4181
	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)
4182 4183 4184 4185 4186 4187 4188
	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 已提交
4189 4190
	// valid new password
	if err = ValidatePassword(rawNewPassword); err != nil {
4191 4192 4193 4194 4195 4196
		log.Error("illegal password", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
		}, nil
	}
4197 4198

	if !passwordVerify(ctx, req.Username, rawOldPassword, globalMetaCache) {
C
codeman 已提交
4199 4200 4201 4202 4203 4204 4205
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "old password is not correct:" + req.Username,
		}, nil
	}
	// update meta data
	encryptedPassword, err := crypto.PasswordEncrypt(rawNewPassword)
4206 4207 4208 4209 4210 4211 4212
	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 已提交
4213
	updateCredReq := &internalpb.CredentialInfo{
4214
		Username:          req.Username,
4215
		Sha256Password:    crypto.SHA256(rawNewPassword, req.Username),
4216 4217
		EncryptedPassword: encryptedPassword,
	}
C
codeman 已提交
4218
	result, err := node.rootCoord.UpdateCredential(ctx, updateCredReq)
4219 4220 4221 4222 4223 4224 4225 4226 4227 4228 4229
	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) {
4230 4231
	log.Debug("DeleteCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4232
		return unhealthyStatus(), nil
4233 4234
	}

4235 4236 4237 4238 4239 4240
	if req.Username == util.UserRoot {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_DeleteCredentialFailure,
			Reason:    "user root cannot be deleted",
		}, nil
	}
4241 4242 4243 4244 4245 4246 4247 4248 4249 4250 4251 4252
	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) {
4253 4254
	log.Debug("ListCredUsers", zap.String("role", typeutil.ProxyRole))
	if !node.checkHealthy() {
4255
		return &milvuspb.ListCredUsersResponse{Status: unhealthyStatus()}, nil
4256
	}
4257
	rootCoordReq := &milvuspb.ListCredUsersRequest{
4258 4259 4260
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_ListCredUsernames),
		),
4261 4262
	}
	resp, err := node.rootCoord.ListCredUsers(ctx, rootCoordReq)
4263 4264 4265 4266 4267 4268 4269 4270 4271 4272 4273 4274
	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,
		},
4275
		Usernames: resp.Usernames,
4276 4277
	}, nil
}
4278

4279 4280 4281
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 {
4282
		return errorutil.UnhealthyStatus(code), nil
4283 4284 4285 4286 4287 4288 4289 4290 4291 4292
	}

	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(),
4293
		}, nil
4294 4295 4296 4297 4298 4299 4300 4301
	}

	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(),
4302
		}, nil
4303 4304
	}
	return result, nil
4305 4306
}

4307 4308 4309
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 {
4310
		return errorutil.UnhealthyStatus(code), nil
4311 4312 4313 4314 4315
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4316
		}, nil
4317
	}
4318 4319 4320 4321 4322
	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,
4323
		}, nil
4324
	}
4325 4326 4327 4328 4329 4330
	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(),
4331
		}, nil
4332 4333
	}
	return result, nil
4334 4335
}

4336 4337 4338
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 {
4339
		return errorutil.UnhealthyStatus(code), nil
4340 4341 4342 4343 4344
	}
	if err := ValidateUsername(req.Username); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4345
		}, nil
4346 4347 4348 4349 4350
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4351
		}, nil
4352 4353 4354 4355 4356 4357 4358 4359
	}

	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(),
4360
		}, nil
4361 4362
	}
	return result, nil
4363 4364
}

4365 4366 4367
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 {
4368
		return &milvuspb.SelectRoleResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4369 4370 4371 4372 4373 4374 4375 4376 4377
	}

	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(),
				},
4378
			}, nil
4379 4380 4381 4382 4383 4384 4385 4386 4387 4388 4389
		}
	}

	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(),
			},
4390
		}, nil
4391 4392
	}
	return result, nil
4393 4394
}

4395 4396 4397
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 {
4398
		return &milvuspb.SelectUserResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4399 4400 4401 4402 4403 4404 4405 4406 4407
	}

	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(),
				},
4408
			}, nil
4409 4410 4411 4412 4413 4414 4415 4416 4417 4418 4419
		}
	}

	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(),
			},
4420
		}, nil
4421 4422
	}
	return result, nil
4423 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
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
4455 4456
}

4457 4458 4459
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 {
4460
		return errorutil.UnhealthyStatus(code), nil
4461 4462 4463 4464 4465
	}
	if err := node.validPrivilegeParams(req); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4466
		}, nil
4467 4468 4469 4470 4471 4472
	}
	curUser, err := GetCurUserFromContext(ctx)
	if err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4473
		}, nil
4474 4475 4476 4477 4478 4479 4480 4481
	}
	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(),
4482
		}, nil
4483 4484
	}
	return result, nil
4485 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
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 {
4516
		return &milvuspb.SelectGrantResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4517 4518 4519 4520 4521 4522 4523 4524
	}

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

	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(),
			},
4536
		}, nil
4537 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
	}
	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
4565
}
4566 4567 4568 4569 4570 4571 4572 4573 4574 4575 4576 4577 4578 4579 4580 4581 4582 4583 4584 4585

// 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
}
4586 4587 4588 4589

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")
4590 4591 4592 4593
		return &milvuspb.CheckHealthResponse{
			Status:    unhealthyStatus(),
			IsHealthy: false,
			Reasons:   []string{reason}}, nil
4594 4595 4596 4597 4598 4599 4600 4601 4602 4603 4604
	}

	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 {
4605
			log.Warn("check health fail", zap.String("role", role), zap.Error(err))
4606 4607 4608 4609 4610
			errReasons = append(errReasons, fmt.Sprintf("check health fail for %s", role))
			return err
		}

		if !resp.IsHealthy {
4611
			log.Warn("check health fail", zap.String("role", role))
4612 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
			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
	}

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