impl.go 159.5 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 45
	"github.com/milvus-io/milvus/internal/util/logutil"
	"github.com/milvus-io/milvus/internal/util/metricsinfo"
46
	"github.com/milvus-io/milvus/internal/util/timerecord"
47
	"github.com/milvus-io/milvus/internal/util/trace"
X
Xiangyu Wang 已提交
48
	"github.com/milvus-io/milvus/internal/util/typeutil"
49 50
	"go.uber.org/zap"
	"go.uber.org/zap/zapcore"
51 52
)

53 54
const moduleName = "Proxy"

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

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

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

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

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

132
	return &commonpb.Status{
133
		ErrorCode: commonpb.ErrorCode_Success,
134 135
		Reason:    "",
	}, nil
136 137
}

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

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

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

153
	cct := &createCollectionTask{
S
sunby 已提交
154
		ctx:                     ctx,
155 156
		Condition:               NewTaskCondition(ctx),
		CreateCollectionRequest: request,
157
		rootCoord:               node.rootCoord,
158 159
	}

160 161 162
	// avoid data race
	lenOfSchema := len(request.Schema)

163 164
	log.Debug(
		rpcReceived(method),
165
		zap.String("traceID", traceID),
166
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
167 168
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
169
		zap.Int("len(schema)", lenOfSchema),
170 171
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
172

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

185
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
186
		return &commonpb.Status{
187
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
188 189 190 191
			Reason:    err.Error(),
		}, nil
	}

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

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

221
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
222
		return &commonpb.Status{
223
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
224 225 226 227
			Reason:    err.Error(),
		}, nil
	}

228 229
	log.Debug(
		rpcDone(method),
230
		zap.String("traceID", traceID),
231
		zap.String("role", typeutil.ProxyRole),
232 233 234 235 236 237
		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),
238 239
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
240

241 242
	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()))
243 244 245
	return cct.result, nil
}

246
// DropCollection drop a collection.
C
Cai Yudong 已提交
247
func (node *Proxy) DropCollection(ctx context.Context, request *milvuspb.DropCollectionRequest) (*commonpb.Status, error) {
248 249 250
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
251 252 253 254

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

259
	dct := &dropCollectionTask{
S
sunby 已提交
260
		ctx:                   ctx,
261 262
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
263
		rootCoord:             node.rootCoord,
264
		chMgr:                 node.chMgr,
S
sunby 已提交
265
		chTicker:              node.chTicker,
266 267
	}

268 269
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
270
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
271 272
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
273 274 275 276 277

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

282
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
283
		return &commonpb.Status{
284
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
285 286 287 288
			Reason:    err.Error(),
		}, nil
	}

289 290
	log.Debug("DropCollection enqueued",
		zap.String("traceID", traceID),
291
		zap.String("role", typeutil.ProxyRole),
292 293 294
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
295 296
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
297 298 299

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DropCollection failed to WaitToFinish",
D
dragondriver 已提交
300
			zap.Error(err),
301
			zap.String("traceID", traceID),
302
			zap.String("role", typeutil.ProxyRole),
303 304 305
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTs", dct.BeginTs()),
			zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
306 307 308
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

309
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
310
		return &commonpb.Status{
311
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
312 313 314 315
			Reason:    err.Error(),
		}, nil
	}

316 317
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
318
		zap.String("role", typeutil.ProxyRole),
319 320 321 322 323 324
		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))

325 326
	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()))
327 328 329
	return dct.result, nil
}

330
// HasCollection check if the specific collection exists in Milvus.
C
Cai Yudong 已提交
331
func (node *Proxy) HasCollection(ctx context.Context, request *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
332 333 334 335 336
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
337 338 339 340

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
341 342
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
343
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
344
		metrics.TotalLabel).Inc()
345 346 347

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
348
		zap.String("role", typeutil.ProxyRole),
349 350 351
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

352
	hct := &hasCollectionTask{
S
sunby 已提交
353
		ctx:                  ctx,
354 355
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
356
		rootCoord:            node.rootCoord,
357 358
	}

359 360 361 362
	if err := node.sched.ddQueue.Enqueue(hct); err != nil {
		log.Warn("HasCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
363
			zap.String("role", typeutil.ProxyRole),
364 365 366
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

367
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
368
			metrics.AbandonLabel).Inc()
369 370
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
371
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
372 373 374 375 376
				Reason:    err.Error(),
			},
		}, nil
	}

377 378
	log.Debug("HasCollection enqueued",
		zap.String("traceID", traceID),
379
		zap.String("role", typeutil.ProxyRole),
380 381 382
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
383 384
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
385 386 387

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

397
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
398
			metrics.FailLabel).Inc()
399 400
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
401
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
402 403 404 405 406
				Reason:    err.Error(),
			},
		}, nil
	}

407 408
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
409
		zap.String("role", typeutil.ProxyRole),
410 411 412 413 414 415
		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))

416
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
417
		metrics.SuccessLabel).Inc()
418
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
419 420 421
	return hct.result, nil
}

422
// LoadCollection load a collection into query nodes.
C
Cai Yudong 已提交
423
func (node *Proxy) LoadCollection(ctx context.Context, request *milvuspb.LoadCollectionRequest) (*commonpb.Status, error) {
424 425 426
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
427 428 429 430

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

443 444
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
445
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
446 447
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
448 449 450 451 452

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

457
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
458
			metrics.AbandonLabel).Inc()
459
		return &commonpb.Status{
460
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
461 462 463
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
464

465 466
	log.Debug("LoadCollection enqueued",
		zap.String("traceID", traceID),
467
		zap.String("role", typeutil.ProxyRole),
468 469 470
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
471 472
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
473 474 475

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

492 493
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
494
		zap.String("role", typeutil.ProxyRole),
495 496 497 498 499 500
		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))

501
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
502
		metrics.SuccessLabel).Inc()
503
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
504
	return lct.result, nil
505 506
}

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

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

528 529
	log.Debug(
		rpcReceived(method),
530
		zap.String("traceID", traceID),
531
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
532 533
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
534 535

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

544
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
545
			metrics.AbandonLabel).Inc()
546
		return &commonpb.Status{
547
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
548 549 550 551
			Reason:    err.Error(),
		}, nil
	}

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

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

574
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
575
			metrics.FailLabel).Inc()
576
		return &commonpb.Status{
577
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
578 579 580 581
			Reason:    err.Error(),
		}, nil
	}

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

592
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
593
		metrics.SuccessLabel).Inc()
594
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
595
	return rct.result, nil
596 597
}

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

606
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
607 608
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
609 610
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
611 612
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
613

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

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

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

635
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
636
			metrics.AbandonLabel).Inc()
637 638
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
639
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
640 641 642 643 644
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

665
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
666
			metrics.FailLabel).Inc()
667

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

676 677
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
678
		zap.String("role", typeutil.ProxyRole),
679 680 681 682 683 684
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

685
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
686
		metrics.SuccessLabel).Inc()
687
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
688 689 690
	return dct.result, nil
}

691 692 693 694 695 696 697 698 699 700 701 702 703 704
// GetStatistics get the statistics, such as `num_rows`.
// WARNING: It is an experimental API
func (node *Proxy) GetStatistics(ctx context.Context, request *milvuspb.GetStatisticsRequest) (*milvuspb.GetStatisticsResponse, error) {
	if !node.checkHealthy() {
		return &milvuspb.GetStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetCollectionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
	method := "GetStatistics"
	tr := timerecord.NewTimeRecorder(method)
705 706
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734
	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))

735
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769
			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))

770
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790
			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))

791
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
792
		metrics.SuccessLabel).Inc()
793
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
794 795 796
	return g.result, nil
}

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

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

819 820
	log.Debug(
		rpcReceived(method),
821
		zap.String("traceID", traceID),
822
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
823 824
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
825 826

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

835
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
836
			metrics.AbandonLabel).Inc()
837

G
godchen 已提交
838
		return &milvuspb.GetCollectionStatisticsResponse{
839
			Status: &commonpb.Status{
840
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
841 842 843 844 845
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

868
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
869
			metrics.FailLabel).Inc()
870

G
godchen 已提交
871
		return &milvuspb.GetCollectionStatisticsResponse{
872
			Status: &commonpb.Status{
873
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
874 875 876 877 878
				Reason:    err.Error(),
			},
		}, nil
	}

879 880
	log.Debug(
		rpcDone(method),
881
		zap.String("traceID", traceID),
882
		zap.String("role", typeutil.ProxyRole),
883
		zap.Int64("msgID", g.ID()),
884 885 886 887 888
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

889
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
890
		metrics.SuccessLabel).Inc()
891
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
892
	return g.result, nil
893 894
}

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

906
	sct := &showCollectionsTask{
G
godchen 已提交
907 908 909
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
910
		queryCoord:             node.queryCoord,
911
		rootCoord:              node.rootCoord,
912 913
	}

914
	log.Debug("ShowCollections received",
915
		zap.String("role", typeutil.ProxyRole),
916 917 918 919 920 921
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

922
	err := node.sched.ddQueue.Enqueue(sct)
923
	if err != nil {
924 925
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
926
			zap.String("role", typeutil.ProxyRole),
927 928 929 930 931 932
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

933
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
G
godchen 已提交
934
		return &milvuspb.ShowCollectionsResponse{
935
			Status: &commonpb.Status{
936
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
937 938 939 940 941
				Reason:    err.Error(),
			},
		}, nil
	}

942
	log.Debug("ShowCollections enqueued",
943
		zap.String("role", typeutil.ProxyRole),
944
		zap.Int64("MsgID", sct.ID()),
945
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
946
		zap.Uint64("TimeStamp", request.TimeStamp),
947 948 949
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
950

951 952
	err = sct.WaitToFinish()
	if err != nil {
953 954
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
955
			zap.String("role", typeutil.ProxyRole),
956 957 958 959 960 961 962
			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),
		)

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

G
godchen 已提交
965
		return &milvuspb.ShowCollectionsResponse{
966
			Status: &commonpb.Status{
967
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
968 969 970 971 972
				Reason:    err.Error(),
			},
		}, nil
	}

973
	log.Debug("ShowCollections Done",
974
		zap.String("role", typeutil.ProxyRole),
975 976 977 978
		zap.Int64("MsgID", sct.ID()),
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
979 980
		zap.Int("len(CollectionNames)", len(request.CollectionNames)),
		zap.Int("num_collections", len(sct.result.CollectionNames)))
981

982 983
	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()))
984 985 986
	return sct.result, nil
}

J
jaime 已提交
987 988 989 990 991 992 993 994 995 996 997
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)

998
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
J
jaime 已提交
999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022

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

1023
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
J
jaime 已提交
1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034 1035 1036 1037 1038 1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052
		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))

1053
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
J
jaime 已提交
1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069
		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))

1070 1071
	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 已提交
1072 1073 1074
	return act.result, nil
}

1075
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
1076
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
1077 1078 1079
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1080

1081
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
1082 1083
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1084 1085
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
1086
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1087

1088
	cpt := &createPartitionTask{
S
sunby 已提交
1089
		ctx:                    ctx,
1090 1091
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
1092
		rootCoord:              node.rootCoord,
1093 1094 1095
		result:                 nil,
	}

1096 1097 1098
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
1099
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1100 1101 1102
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1103 1104 1105 1106 1107 1108

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
1109
			zap.String("role", typeutil.ProxyRole),
1110 1111 1112 1113
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1116
		return &commonpb.Status{
1117
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1118 1119 1120
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1121

1122 1123 1124
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
1125
		zap.String("role", typeutil.ProxyRole),
1126 1127 1128
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
1129 1130 1131
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1132 1133 1134 1135

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

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

1148
		return &commonpb.Status{
1149
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1150 1151 1152
			Reason:    err.Error(),
		}, nil
	}
1153 1154 1155 1156

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
1157
		zap.String("role", typeutil.ProxyRole),
1158 1159 1160 1161 1162 1163 1164
		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))

1165 1166
	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()))
1167 1168 1169
	return cpt.result, nil
}

1170
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
1171
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
1172 1173 1174
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1175

1176
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
1177 1178
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1179 1180
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
1181
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1182

1183
	dpt := &dropPartitionTask{
S
sunby 已提交
1184
		ctx:                  ctx,
1185 1186
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1187
		rootCoord:            node.rootCoord,
1188 1189 1190
		result:               nil,
	}

1191 1192 1193
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1194
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1195 1196 1197
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1198 1199 1200 1201 1202 1203

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1204
			zap.String("role", typeutil.ProxyRole),
1205 1206 1207 1208
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1211
		return &commonpb.Status{
1212
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1213 1214 1215
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1216

1217 1218 1219
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1220
		zap.String("role", typeutil.ProxyRole),
1221 1222 1223
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1224 1225 1226
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1227 1228 1229 1230

	if err := dpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1231
			zap.Error(err),
1232
			zap.String("traceID", traceID),
1233
			zap.String("role", typeutil.ProxyRole),
1234 1235 1236
			zap.Int64("MsgID", dpt.ID()),
			zap.Uint64("BeginTS", dpt.BeginTs()),
			zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1237 1238 1239 1240
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1243
		return &commonpb.Status{
1244
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1245 1246 1247
			Reason:    err.Error(),
		}, nil
	}
1248 1249 1250 1251

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1252
		zap.String("role", typeutil.ProxyRole),
1253 1254 1255 1256 1257 1258 1259
		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))

1260 1261
	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()))
1262 1263 1264
	return dpt.result, nil
}

1265
// HasPartition check if partition exist.
C
Cai Yudong 已提交
1266
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
1267 1268 1269 1270 1271
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1272

D
dragondriver 已提交
1273
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
1274 1275
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1276 1277 1278
	method := "HasPartition"
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
1279
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1280
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1281

1282
	hpt := &hasPartitionTask{
S
sunby 已提交
1283
		ctx:                 ctx,
1284 1285
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1286
		rootCoord:           node.rootCoord,
1287 1288 1289
		result:              nil,
	}

D
dragondriver 已提交
1290 1291 1292
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1293
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1294 1295 1296
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1297 1298 1299 1300 1301 1302

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

1308
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1309
			metrics.AbandonLabel).Inc()
1310

1311 1312
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1313
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1314 1315 1316 1317 1318
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1319

D
dragondriver 已提交
1320 1321 1322
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1323
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1324 1325 1326
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1327 1328 1329
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1330 1331 1332 1333

	if err := hpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1334
			zap.Error(err),
D
dragondriver 已提交
1335
			zap.String("traceID", traceID),
1336
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1337 1338 1339
			zap.Int64("MsgID", hpt.ID()),
			zap.Uint64("BeginTS", hpt.BeginTs()),
			zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1340 1341 1342 1343
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1344
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1345
			metrics.FailLabel).Inc()
1346

1347 1348
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1349
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1350 1351 1352 1353 1354
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1355 1356 1357 1358

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1359
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1360 1361 1362 1363 1364 1365 1366
		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))

1367
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1368
		metrics.SuccessLabel).Inc()
1369
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1370 1371 1372
	return hpt.result, nil
}

1373
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1374
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1375 1376 1377
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1378

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

1394 1395 1396
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1397
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1398 1399 1400
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1401 1402 1403 1404 1405 1406

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1407
			zap.String("role", typeutil.ProxyRole),
1408 1409 1410 1411
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1412
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1413
			metrics.AbandonLabel).Inc()
1414

1415
		return &commonpb.Status{
1416
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1417 1418 1419 1420
			Reason:    err.Error(),
		}, nil
	}

1421 1422 1423
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1424
		zap.String("role", typeutil.ProxyRole),
1425 1426 1427
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1428 1429 1430
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1431 1432 1433 1434

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

1445
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1446
			metrics.FailLabel).Inc()
1447

1448
		return &commonpb.Status{
1449
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1450 1451 1452 1453
			Reason:    err.Error(),
		}, nil
	}

1454 1455 1456
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1457
		zap.String("role", typeutil.ProxyRole),
1458 1459 1460 1461 1462 1463 1464
		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))

1465
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1466
		metrics.SuccessLabel).Inc()
1467
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1468
	return lpt.result, nil
1469 1470
}

1471
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1472
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1473 1474 1475
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1476 1477 1478 1479 1480

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

1481
	rpt := &releasePartitionsTask{
G
godchen 已提交
1482 1483 1484
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1485
		queryCoord:               node.queryCoord,
1486 1487
	}

1488
	method := "ReleasePartitions"
1489
	tr := timerecord.NewTimeRecorder(method)
1490 1491
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
1492 1493 1494
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1495
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1496 1497 1498
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1499 1500 1501 1502 1503 1504

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1505
			zap.String("role", typeutil.ProxyRole),
1506 1507 1508 1509
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1510
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1511
			metrics.AbandonLabel).Inc()
1512

1513
		return &commonpb.Status{
1514
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1515 1516 1517 1518
			Reason:    err.Error(),
		}, nil
	}

1519 1520 1521
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1522
		zap.String("role", typeutil.ProxyRole),
1523 1524 1525
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1526 1527 1528
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1529 1530 1531 1532

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1533
			zap.Error(err),
1534
			zap.String("traceID", traceID),
1535
			zap.String("role", typeutil.ProxyRole),
1536 1537 1538
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1539 1540 1541 1542
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1543
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1544
			metrics.FailLabel).Inc()
1545

1546
		return &commonpb.Status{
1547
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1548 1549 1550 1551
			Reason:    err.Error(),
		}, nil
	}

1552 1553 1554
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1555
		zap.String("role", typeutil.ProxyRole),
1556 1557 1558 1559 1560 1561 1562
		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))

1563
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1564
		metrics.SuccessLabel).Inc()
1565
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1566
	return rpt.result, nil
1567 1568
}

1569
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1570
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1571 1572 1573 1574 1575
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1576 1577 1578 1579

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-GetPartitionStatistics")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1580 1581
	method := "GetPartitionStatistics"
	tr := timerecord.NewTimeRecorder(method)
1582 1583
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
1584

1585
	g := &getPartitionStatisticsTask{
1586 1587 1588
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1589
		dataCoord:                     node.dataCoord,
1590 1591
	}

1592 1593 1594
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1595
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1596 1597 1598
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1599 1600 1601 1602 1603 1604

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1605
			zap.String("role", typeutil.ProxyRole),
1606 1607 1608 1609
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1610
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1611
			metrics.AbandonLabel).Inc()
1612

1613 1614 1615 1616 1617 1618 1619 1620
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1621 1622 1623
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1624
		zap.String("role", typeutil.ProxyRole),
1625 1626 1627
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1628 1629 1630
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1631 1632 1633 1634

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1635
			zap.Error(err),
1636
			zap.String("traceID", traceID),
1637
			zap.String("role", typeutil.ProxyRole),
1638 1639 1640
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1641 1642 1643 1644
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1645
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1646
			metrics.FailLabel).Inc()
1647

1648 1649 1650 1651 1652 1653 1654 1655
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1656 1657 1658
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1659
		zap.String("role", typeutil.ProxyRole),
1660 1661 1662 1663 1664 1665 1666
		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))

1667
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1668
		metrics.SuccessLabel).Inc()
1669
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1670
	return g.result, nil
1671 1672
}

1673
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1674
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1675 1676 1677 1678 1679
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1680 1681 1682 1683 1684

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

1685
	spt := &showPartitionsTask{
G
godchen 已提交
1686 1687 1688
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1689
		rootCoord:             node.rootCoord,
1690
		queryCoord:            node.queryCoord,
G
godchen 已提交
1691
		result:                nil,
1692 1693
	}

1694
	method := "ShowPartitions"
1695 1696
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
1697
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1698
		metrics.TotalLabel).Inc()
1699 1700 1701 1702

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1703
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1704
		zap.Any("request", request))
1705 1706 1707 1708 1709 1710

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

1714
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1715
			metrics.AbandonLabel).Inc()
1716

G
godchen 已提交
1717
		return &milvuspb.ShowPartitionsResponse{
1718
			Status: &commonpb.Status{
1719
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1720 1721 1722 1723 1724
				Reason:    err.Error(),
			},
		}, nil
	}

1725 1726 1727
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1728
		zap.String("role", typeutil.ProxyRole),
1729 1730 1731
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1732 1733
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1734 1735 1736 1737 1738
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1739
			zap.Error(err),
1740
			zap.String("traceID", traceID),
1741
			zap.String("role", typeutil.ProxyRole),
1742 1743 1744 1745 1746 1747
			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 已提交
1748

1749
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1750
			metrics.FailLabel).Inc()
1751

G
godchen 已提交
1752
		return &milvuspb.ShowPartitionsResponse{
1753
			Status: &commonpb.Status{
1754
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1755 1756 1757 1758
				Reason:    err.Error(),
			},
		}, nil
	}
1759 1760 1761 1762

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1763
		zap.String("role", typeutil.ProxyRole),
1764 1765 1766 1767 1768 1769 1770
		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))

1771
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1772
		metrics.SuccessLabel).Inc()
1773
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1774 1775 1776
	return spt.result, nil
}

S
SimFG 已提交
1777 1778
func (node *Proxy) getCollectionProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest, collectionID int64) (int64, error) {
	resp, err := node.queryCoord.ShowCollections(ctx, &querypb.ShowCollectionsRequest{
1779 1780 1781 1782
		Base: commonpbutil.NewMsgBaseCopy(
			request.Base,
			commonpbutil.WithMsgType(commonpb.MsgType_DescribeCollection),
		),
S
SimFG 已提交
1783 1784 1785 1786 1787 1788 1789 1790 1791 1792 1793 1794 1795 1796 1797 1798 1799 1800 1801 1802 1803 1804 1805
		CollectionIDs: []int64{collectionID},
	})
	if err != nil {
		return 0, err
	}
	if len(resp.InMemoryPercentages) == 0 {
		return 0, errors.New("fail to show collections from the querycoord, no data")
	}
	return resp.InMemoryPercentages[0], nil
}

func (node *Proxy) getPartitionProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest, collectionID int64) (int64, error) {
	IDs2Names := make(map[int64]string)
	partitionIDs := make([]int64, 0)
	for _, partitionName := range request.PartitionNames {
		partitionID, err := globalMetaCache.GetPartitionID(ctx, request.CollectionName, partitionName)
		if err != nil {
			return 0, err
		}
		IDs2Names[partitionID] = partitionName
		partitionIDs = append(partitionIDs, partitionID)
	}
	resp, err := node.queryCoord.ShowPartitions(ctx, &querypb.ShowPartitionsRequest{
1806 1807 1808 1809
		Base: commonpbutil.NewMsgBaseCopy(
			request.Base,
			commonpbutil.WithMsgType(commonpb.MsgType_ShowPartitions),
		),
S
SimFG 已提交
1810 1811 1812 1813 1814 1815 1816 1817 1818 1819 1820 1821 1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832 1833 1834 1835
		CollectionID: collectionID,
		PartitionIDs: partitionIDs,
	})
	if err != nil {
		return 0, err
	}
	if len(resp.InMemoryPercentages) != len(partitionIDs) {
		return 0, errors.New("fail to show partitions from the querycoord, invalid data num")
	}
	var progress int64
	for _, p := range resp.InMemoryPercentages {
		progress += p
	}
	progress /= int64(len(partitionIDs))
	return progress, nil
}

func (node *Proxy) GetLoadingProgress(ctx context.Context, request *milvuspb.GetLoadingProgressRequest) (*milvuspb.GetLoadingProgressResponse, error) {
	if !node.checkHealthy() {
		return &milvuspb.GetLoadingProgressResponse{Status: unhealthyStatus()}, nil
	}
	method := "GetLoadingProgress"
	tr := timerecord.NewTimeRecorder(method)
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ShowPartitions")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1836
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
S
SimFG 已提交
1837 1838 1839 1840 1841 1842 1843 1844
	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))
1845
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
S
SimFG 已提交
1846 1847 1848 1849 1850 1851 1852 1853 1854 1855 1856 1857 1858 1859
		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
	}
1860 1861 1862 1863 1864 1865
	msgBase := commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
S
SimFG 已提交
1866 1867 1868 1869 1870 1871 1872 1873 1874 1875 1876 1877 1878 1879 1880 1881 1882 1883 1884 1885 1886 1887 1888
	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))
1889 1890
	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 已提交
1891 1892 1893 1894 1895 1896 1897 1898
	return &milvuspb.GetLoadingProgressResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
		},
		Progress: progress,
	}, nil
}

1899
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1900
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1901 1902 1903
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1904 1905 1906 1907 1908

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

1909
	cit := &createIndexTask{
Z
zhenshan.cao 已提交
1910 1911 1912 1913 1914
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		req:        request,
		rootCoord:  node.rootCoord,
		indexCoord: node.indexCoord,
1915 1916
	}

D
dragondriver 已提交
1917
	method := "CreateIndex"
1918
	tr := timerecord.NewTimeRecorder(method)
1919 1920
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1921 1922 1923
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1924
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1925 1926 1927 1928
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1929 1930 1931 1932 1933 1934

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1935
			zap.String("role", typeutil.ProxyRole),
Z
zhenshan.cao 已提交
1936 1937 1938 1939
			zap.String("db", request.GetDbName()),
			zap.String("collection", request.GetCollectionName()),
			zap.String("field", request.GetFieldName()),
			zap.Any("extra_params", request.GetExtraParams()))
D
dragondriver 已提交
1940

1941
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1942
			metrics.AbandonLabel).Inc()
1943

1944
		return &commonpb.Status{
1945
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1946 1947 1948 1949
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1950 1951 1952
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1953
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1954 1955 1956
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1957 1958 1959 1960
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1961 1962 1963 1964

	if err := cit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1965
			zap.Error(err),
D
dragondriver 已提交
1966
			zap.String("traceID", traceID),
1967
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1968 1969 1970
			zap.Int64("MsgID", cit.ID()),
			zap.Uint64("BeginTs", cit.BeginTs()),
			zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1971 1972 1973 1974 1975
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1976
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1977
			metrics.FailLabel).Inc()
1978

1979
		return &commonpb.Status{
1980
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1981 1982 1983 1984
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1985 1986 1987
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1988
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1989 1990 1991 1992 1993 1994 1995 1996
		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))

1997
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1998
		metrics.SuccessLabel).Inc()
1999
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2000 2001 2002
	return cit.result, nil
}

2003
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
2004
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
2005 2006 2007 2008 2009
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2010 2011 2012 2013 2014

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

2015
	dit := &describeIndexTask{
S
sunby 已提交
2016
		ctx:                  ctx,
2017 2018
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
2019
		indexCoord:           node.indexCoord,
2020 2021
	}

2022 2023 2024
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
2025
	tr := timerecord.NewTimeRecorder(method)
2026 2027
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2028 2029 2030
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2031
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2032 2033 2034
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2035 2036 2037 2038 2039 2040 2041
		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),
2042
			zap.String("role", typeutil.ProxyRole),
2043 2044 2045 2046 2047
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

2048
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2049
			metrics.AbandonLabel).Inc()
2050

2051 2052
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
2053
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2054 2055 2056 2057 2058
				Reason:    err.Error(),
			},
		}, nil
	}

2059 2060 2061
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2062
		zap.String("role", typeutil.ProxyRole),
2063 2064 2065
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2066 2067 2068
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
2069 2070 2071 2072 2073
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2074
			zap.Error(err),
2075
			zap.String("traceID", traceID),
2076
			zap.String("role", typeutil.ProxyRole),
2077 2078 2079
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2080 2081 2082
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
2083
			zap.String("index name", indexName))
D
dragondriver 已提交
2084

Z
zhenshan.cao 已提交
2085 2086 2087 2088
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
2089
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2090
			metrics.FailLabel).Inc()
2091

2092 2093
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
2094
				ErrorCode: errCode,
2095 2096 2097 2098 2099
				Reason:    err.Error(),
			},
		}, nil
	}

2100 2101 2102
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2103
		zap.String("role", typeutil.ProxyRole),
2104 2105 2106 2107 2108 2109 2110 2111
		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))

2112
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2113
		metrics.SuccessLabel).Inc()
2114
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2115 2116 2117
	return dit.result, nil
}

2118
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
2119
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
2120 2121 2122
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2123 2124 2125 2126 2127

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

2128
	dit := &dropIndexTask{
S
sunby 已提交
2129
		ctx:              ctx,
B
BossZou 已提交
2130 2131
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
2132
		indexCoord:       node.indexCoord,
2133
		queryCoord:       node.queryCoord,
B
BossZou 已提交
2134
	}
G
godchen 已提交
2135

D
dragondriver 已提交
2136
	method := "DropIndex"
2137
	tr := timerecord.NewTimeRecorder(method)
2138 2139
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2140 2141 2142 2143

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2144
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2145 2146 2147 2148 2149
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
2150 2151 2152 2153 2154
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2155
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2156 2157 2158 2159
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2160
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2161
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2162

B
BossZou 已提交
2163
		return &commonpb.Status{
2164
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2165 2166 2167
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2168

D
dragondriver 已提交
2169 2170 2171
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2172
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2173 2174 2175
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2176 2177 2178 2179
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
2180 2181 2182 2183

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2184
			zap.Error(err),
D
dragondriver 已提交
2185
			zap.String("traceID", traceID),
2186
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2187 2188 2189
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
2190 2191 2192 2193 2194
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2195
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2196
			metrics.FailLabel).Inc()
2197

B
BossZou 已提交
2198
		return &commonpb.Status{
2199
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
2200 2201 2202
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
2203 2204 2205 2206

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2207
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2208 2209 2210 2211 2212 2213 2214 2215
		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))

2216
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2217
		metrics.SuccessLabel).Inc()
2218
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
2219 2220 2221
	return dit.result, nil
}

2222 2223
// 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.
2224
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2225
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
2226 2227 2228 2229 2230
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2231 2232 2233 2234 2235

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

2236
	gibpt := &getIndexBuildProgressTask{
2237 2238 2239
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
2240 2241
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
2242
		dataCoord:                    node.dataCoord,
2243 2244
	}

2245
	method := "GetIndexBuildProgress"
2246
	tr := timerecord.NewTimeRecorder(method)
2247 2248
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2249 2250 2251
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2252
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2253 2254 2255 2256
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2257 2258 2259 2260 2261 2262

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2263
			zap.String("role", typeutil.ProxyRole),
2264 2265 2266 2267
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2268
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2269
			metrics.AbandonLabel).Inc()
2270

2271 2272 2273 2274 2275 2276 2277 2278
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2279 2280 2281
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2282
		zap.String("role", typeutil.ProxyRole),
2283 2284 2285
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
2286 2287 2288 2289
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2290 2291 2292 2293

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2294
			zap.Error(err),
2295
			zap.String("traceID", traceID),
2296
			zap.String("role", typeutil.ProxyRole),
2297 2298 2299
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2300 2301 2302 2303
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2304
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2305
			metrics.FailLabel).Inc()
2306 2307 2308 2309 2310 2311 2312 2313

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2314 2315 2316 2317

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2318
		zap.String("role", typeutil.ProxyRole),
2319 2320 2321 2322 2323 2324 2325 2326
		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))
2327

2328
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2329
		metrics.SuccessLabel).Inc()
2330
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2331
	return gibpt.result, nil
2332 2333
}

2334
// GetIndexState get the build-state of index.
2335
// Deprecated: use DescribeIndex instead
C
Cai Yudong 已提交
2336
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2337 2338 2339 2340 2341
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2342 2343 2344 2345 2346

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

2347
	dipt := &getIndexStateTask{
G
godchen 已提交
2348 2349 2350
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2351 2352
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2353 2354
	}

2355
	method := "GetIndexState"
2356
	tr := timerecord.NewTimeRecorder(method)
2357 2358
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2359 2360 2361
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2362
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2363 2364 2365 2366
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2367 2368 2369 2370 2371 2372

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2373
			zap.String("role", typeutil.ProxyRole),
2374 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
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2380
			metrics.AbandonLabel).Inc()
2381

G
godchen 已提交
2382
		return &milvuspb.GetIndexStateResponse{
2383
			Status: &commonpb.Status{
2384
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2385 2386 2387 2388 2389
				Reason:    err.Error(),
			},
		}, nil
	}

2390 2391 2392
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2393
		zap.String("role", typeutil.ProxyRole),
2394 2395 2396
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2397 2398 2399 2400
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2401 2402 2403 2404

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2405
			zap.Error(err),
2406
			zap.String("traceID", traceID),
2407
			zap.String("role", typeutil.ProxyRole),
2408 2409 2410
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2411 2412 2413 2414
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2415
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2416
			metrics.FailLabel).Inc()
2417

G
godchen 已提交
2418
		return &milvuspb.GetIndexStateResponse{
2419
			Status: &commonpb.Status{
2420
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2421 2422 2423 2424 2425
				Reason:    err.Error(),
			},
		}, nil
	}

2426 2427 2428
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2429
		zap.String("role", typeutil.ProxyRole),
2430 2431 2432 2433 2434 2435 2436 2437
		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))

2438
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2439
		metrics.SuccessLabel).Inc()
2440
	metrics.ProxyReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2441 2442 2443
	return dipt.result, nil
}

2444
// Insert insert records into collection.
C
Cai Yudong 已提交
2445
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2446 2447 2448 2449 2450 2451
	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))

2452 2453 2454 2455 2456
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2457 2458
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
2459
	receiveSize := proto.Size(request)
2460 2461
	rateCol.Add(internalpb.RateType_DMLInsert.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Add(float64(receiveSize))
D
dragondriver 已提交
2462

2463
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
2464
	it := &insertTask{
2465 2466
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2467
		// req:       request,
2468 2469 2470 2471
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2472
			InsertRequest: internalpb.InsertRequest{
2473 2474 2475 2476 2477
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Insert),
					commonpbutil.WithMsgID(0),
					commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
				),
2478 2479
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2480 2481 2482
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2483
				// RowData: transfer column based request to this
2484 2485
			},
		},
2486
		idAllocator:   node.rowIDAllocator,
2487 2488 2489
		segIDAssigner: node.segAssigner,
		chMgr:         node.chMgr,
		chTicker:      node.chTicker,
2490
	}
2491 2492

	if len(it.PartitionName) <= 0 {
2493
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2494 2495
	}

X
Xiangyu Wang 已提交
2496
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2497
		numRows := request.NumRows
2498 2499 2500 2501
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2502

X
Xiangyu Wang 已提交
2503 2504 2505 2506 2507 2508 2509
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2510 2511
	}

X
Xiangyu Wang 已提交
2512
	log.Debug("Enqueue insert request in Proxy",
2513
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2514 2515 2516 2517 2518
		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)),
2519 2520
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2521

X
Xiangyu Wang 已提交
2522 2523
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2524
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2525
			metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2526
		return constructFailedResponse(err), nil
2527
	}
D
dragondriver 已提交
2528

X
Xiangyu Wang 已提交
2529
	log.Debug("Detail of insert request in Proxy",
2530
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2531
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2532 2533 2534 2535 2536
		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 已提交
2537 2538 2539 2540 2541
		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))
2542
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2543
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2544 2545 2546 2547 2548
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2549
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2550 2551 2552 2553 2554 2555 2556 2557 2558 2559 2560
			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 已提交
2561
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2562

2563
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2564
		metrics.SuccessLabel).Inc()
2565 2566
	successCnt := it.result.InsertCnt - int64(len(it.result.ErrIndex))
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(successCnt))
2567
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2568
	metrics.ProxyCollectionMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
2569 2570 2571
	return it.result, nil
}

2572
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2573
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2574 2575 2576
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2577 2578
	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))
2579

2580
	receiveSize := proto.Size(request)
2581 2582
	rateCol.Add(internalpb.RateType_DMLDelete.String(), float64(receiveSize))
	metrics.ProxyReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Add(float64(receiveSize))
2583

G
groot 已提交
2584 2585 2586 2587 2588 2589
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2590 2591 2592
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

2593
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2594
		metrics.TotalLabel).Inc()
2595
	dt := &deleteTask{
X
xige-16 已提交
2596 2597 2598
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		deleteExpr: request.Expr,
G
godchen 已提交
2599
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2600 2601 2602
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2603
			DeleteRequest: internalpb.DeleteRequest{
2604 2605 2606 2607
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Delete),
					commonpbutil.WithMsgID(0),
				),
X
xige-16 已提交
2608
				DbName:         request.DbName,
G
godchen 已提交
2609 2610 2611
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2612 2613 2614 2615
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2616 2617
	}

2618
	log.Debug("Enqueue delete request in Proxy",
2619
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2620 2621 2622 2623
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2624 2625 2626 2627

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

G
groot 已提交
2631 2632 2633 2634 2635 2636 2637 2638
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2639
	log.Debug("Detail of delete request in Proxy",
2640
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2641 2642 2643 2644 2645
		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),
2646 2647
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2648

2649 2650
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
2651
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2652
			metrics.FailLabel).Inc()
G
groot 已提交
2653 2654 2655 2656 2657 2658 2659 2660
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2661
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2662
		metrics.SuccessLabel).Inc()
2663
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2664
	metrics.ProxyCollectionMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2665 2666 2667
	return dt.result, nil
}

2668
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2669
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2670 2671 2672 2673 2674
	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()))

2675 2676 2677 2678 2679
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2680 2681
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
2682
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2683
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2684

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

2688
	qt := &searchTask{
S
sunby 已提交
2689
		ctx:       ctx,
2690
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2691
		SearchRequest: &internalpb.SearchRequest{
2692 2693 2694 2695
			Base: commonpbutil.NewMsgBase(
				commonpbutil.WithMsgType(commonpb.MsgType_Search),
				commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
			),
2696
			ReqID: Params.ProxyCfg.GetNodeID(),
2697
		},
2698 2699 2700 2701
		request:  request,
		qc:       node.queryCoord,
		tr:       timerecord.NewTimeRecorder("search"),
		shardMgr: node.shardMgr,
2702 2703
	}

2704 2705 2706
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

Z
Zach 已提交
2707
	log.Ctx(ctx).Info(
2708
		rpcReceived(method),
2709
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2710 2711 2712 2713 2714
		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)),
2715 2716 2717 2718
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2719

2720
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
2721
		log.Ctx(ctx).Warn(
2722
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2723
			zap.Error(err),
2724
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2725 2726 2727 2728 2729 2730
			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),
2731 2732 2733
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2734

2735
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2736
			metrics.AbandonLabel).Inc()
2737

2738 2739
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2740
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2741 2742 2743 2744
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
2745
	tr.CtxRecord(ctx, "search request enqueue")
2746

Z
Zach 已提交
2747
	log.Ctx(ctx).Debug(
2748
		rpcEnqueued(method),
2749
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2750
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2751 2752 2753 2754 2755
		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),
2756
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2757 2758 2759 2760
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2761

2762
	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
2763
		log.Ctx(ctx).Warn(
2764
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2765
			zap.Error(err),
2766
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2767
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2768 2769 2770 2771
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2772
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2773 2774 2775 2776
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2777

2778
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2779
			metrics.FailLabel).Inc()
2780

2781 2782
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2783
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2784 2785 2786 2787 2788
				Reason:    err.Error(),
			},
		}, nil
	}

Z
Zach 已提交
2789
	span := tr.CtxRecord(ctx, "wait search result")
2790 2791
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel).Observe(float64(span.Milliseconds()))
2792
	tr.CtxRecord(ctx, "wait search result")
Z
Zach 已提交
2793
	log.Ctx(ctx).Debug(
2794
		rpcDone(method),
2795
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2796 2797 2798 2799 2800 2801
		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)),
2802 2803 2804 2805
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2806

2807
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2808 2809
		metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(qt.result.GetResults().GetNumQueries()))
C
cai.zhang 已提交
2810
	searchDur := tr.ElapseSpan().Milliseconds()
2811
	metrics.ProxySQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
2812
		metrics.SearchLabel).Observe(float64(searchDur))
2813 2814
	metrics.ProxyCollectionSQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel, request.CollectionName).Observe(float64(searchDur))
2815 2816 2817
	if qt.result != nil {
		sentSize := proto.Size(qt.result)
		metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
2818
		rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
2819
	}
2820 2821 2822
	return qt.result, nil
}

2823
// Flush notify data nodes to persist the data of collection.
2824 2825 2826 2827 2828 2829 2830
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2831
	if !node.checkHealthy() {
2832 2833
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2834
	}
D
dragondriver 已提交
2835 2836 2837 2838 2839

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

2840
	ft := &flushTask{
T
ThreadDao 已提交
2841 2842 2843
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2844
		dataCoord:    node.dataCoord,
2845 2846
	}

D
dragondriver 已提交
2847
	method := "Flush"
2848
	tr := timerecord.NewTimeRecorder(method)
2849
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2850 2851 2852 2853

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2854
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2855 2856
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2857 2858 2859 2860 2861 2862

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

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

2869 2870
		resp.Status.Reason = err.Error()
		return resp, nil
2871 2872
	}

D
dragondriver 已提交
2873 2874 2875
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2876
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2877 2878 2879
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2880 2881
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2882 2883 2884 2885

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2886
			zap.Error(err),
D
dragondriver 已提交
2887
			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 2894
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

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

D
dragondriver 已提交
2897
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2898 2899
		resp.Status.Reason = err.Error()
		return resp, nil
2900 2901
	}

D
dragondriver 已提交
2902 2903 2904
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2905
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2906 2907 2908 2909 2910 2911
		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))

2912 2913
	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()))
2914
	return ft.result, nil
2915 2916
}

2917
// Query get the records by primary keys.
C
Cai Yudong 已提交
2918
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2919 2920 2921 2922 2923
	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)

2924 2925 2926 2927 2928
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2929

D
dragondriver 已提交
2930 2931
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
2932
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2933

2934
	qt := &queryTask{
2935 2936 2937
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
2938 2939 2940 2941
			Base: commonpbutil.NewMsgBase(
				commonpbutil.WithMsgType(commonpb.MsgType_Retrieve),
				commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
			),
2942
			ReqID: Params.ProxyCfg.GetNodeID(),
2943
		},
2944 2945
		request:          request,
		qc:               node.queryCoord,
2946
		queryShardPolicy: mergeRoundRobinPolicy,
2947
		shardMgr:         node.shardMgr,
2948 2949
	}

D
dragondriver 已提交
2950 2951
	method := "Query"

2952
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2953 2954
		metrics.TotalLabel).Inc()

Z
Zach 已提交
2955
	log.Ctx(ctx).Info(
D
dragondriver 已提交
2956
		rpcReceived(method),
2957
		zap.String("role", typeutil.ProxyRole),
2958 2959
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2960 2961 2962 2963 2964
		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 已提交
2965

D
dragondriver 已提交
2966
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
Z
Zach 已提交
2967
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
2968 2969 2970
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("role", typeutil.ProxyRole),
2971 2972 2973
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2974

2975 2976
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
2977

2978 2979 2980 2981 2982 2983
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2984
	}
Z
Zach 已提交
2985
	tr.CtxRecord(ctx, "query request enqueue")
2986

Z
Zach 已提交
2987
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
2988
		rpcEnqueued(method),
2989
		zap.String("role", typeutil.ProxyRole),
2990
		zap.Int64("msgID", qt.ID()),
2991 2992
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2993
		zap.Strings("partitions", request.PartitionNames))
D
dragondriver 已提交
2994 2995

	if err := qt.WaitToFinish(); err != nil {
Z
Zach 已提交
2996
		log.Ctx(ctx).Warn(
D
dragondriver 已提交
2997 2998
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
2999
			zap.String("role", typeutil.ProxyRole),
3000
			zap.Int64("msgID", qt.ID()),
3001 3002 3003
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
3004

3005
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3006
			metrics.FailLabel).Inc()
3007

3008 3009 3010 3011 3012 3013 3014
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
Z
Zach 已提交
3015
	span := tr.CtxRecord(ctx, "wait query result")
3016 3017
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel).Observe(float64(span.Milliseconds()))
3018

Z
Zach 已提交
3019
	log.Ctx(ctx).Debug(
D
dragondriver 已提交
3020 3021
		rpcDone(method),
		zap.String("role", typeutil.ProxyRole),
3022
		zap.Int64("msgID", qt.ID()),
3023 3024 3025
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
3026

3027
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3028 3029
		metrics.SuccessLabel).Inc()

3030
	metrics.ProxySQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
3031
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
3032 3033
	metrics.ProxyCollectionSQLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel, request.CollectionName).Observe(float64(tr.ElapseSpan().Milliseconds()))
3034 3035

	ret := &milvuspb.QueryResults{
3036 3037
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
3038 3039
	}
	sentSize := proto.Size(qt.result)
3040
	rateCol.Add(metricsinfo.ReadResultThroughput, float64(sentSize))
3041 3042
	metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
	return ret, nil
3043
}
3044

3045
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
3046 3047 3048 3049
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3050 3051 3052 3053 3054

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

Y
Yusup 已提交
3055 3056 3057 3058 3059 3060 3061
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
3062
	method := "CreateAlias"
3063
	tr := timerecord.NewTimeRecorder(method)
3064
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3065 3066 3067 3068 3069 3070 3071 3072 3073 3074 3075 3076 3077 3078 3079 3080 3081 3082 3083

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

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

Y
Yusup 已提交
3086 3087 3088 3089 3090 3091
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3092 3093 3094
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3095
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3096 3097 3098 3099
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3100 3101
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3102 3103 3104 3105

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3106
			zap.Error(err),
D
dragondriver 已提交
3107
			zap.String("traceID", traceID),
3108
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3109 3110 3111 3112
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3113 3114
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
3115
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
3116 3117 3118 3119 3120 3121 3122

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

D
dragondriver 已提交
3123 3124 3125 3126 3127 3128 3129 3130 3131 3132 3133
	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))

3134 3135
	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 已提交
3136 3137 3138
	return cat.result, nil
}

3139
// DropAlias alter the alias of collection.
Y
Yusup 已提交
3140 3141 3142 3143
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3144 3145 3146 3147 3148

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

Y
Yusup 已提交
3149 3150 3151 3152 3153 3154 3155
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
3156
	method := "DropAlias"
3157
	tr := timerecord.NewTimeRecorder(method)
3158
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3159 3160 3161 3162 3163 3164 3165 3166 3167 3168 3169 3170 3171 3172 3173 3174

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

Y
Yusup 已提交
3177 3178 3179 3180 3181 3182
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3183 3184 3185
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3186
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3187 3188 3189 3190
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3191
		zap.String("alias", request.Alias))
D
dragondriver 已提交
3192 3193 3194 3195

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3196
			zap.Error(err),
D
dragondriver 已提交
3197
			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 3204
			zap.String("alias", request.Alias))

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

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

D
dragondriver 已提交
3213 3214 3215 3216 3217 3218 3219 3220 3221 3222
	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))

3223 3224
	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 已提交
3225 3226 3227
	return dat.result, nil
}

3228
// AlterAlias alter alias of collection.
Y
Yusup 已提交
3229 3230 3231 3232
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
3233 3234 3235 3236 3237

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

Y
Yusup 已提交
3238 3239 3240 3241 3242 3243 3244
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
3245
	method := "AlterAlias"
3246
	tr := timerecord.NewTimeRecorder(method)
3247
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
3248 3249 3250 3251 3252 3253 3254 3255 3256 3257 3258 3259 3260 3261 3262 3263 3264 3265

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

Y
Yusup 已提交
3268 3269 3270 3271 3272 3273
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3274 3275 3276
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
3277
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3278 3279 3280 3281
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3282 3283
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3284 3285 3286 3287

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
3288
			zap.Error(err),
D
dragondriver 已提交
3289
			zap.String("traceID", traceID),
3290
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3291 3292 3293 3294
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
3295 3296 3297
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

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

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

D
dragondriver 已提交
3306 3307 3308 3309 3310 3311 3312 3313 3314 3315 3316
	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))

3317 3318
	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 已提交
3319 3320 3321
	return aat.result, nil
}

3322
// CalcDistance calculates the distances between vectors.
3323
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3324 3325 3326 3327 3328
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3329

3330 3331 3332 3333
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3334 3335
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3336

3337 3338 3339 3340 3341
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3342 3343
		}

3344
		qt := &queryTask{
3345 3346 3347
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
3348 3349 3350 3351
				Base: commonpbutil.NewMsgBase(
					commonpbutil.WithMsgType(commonpb.MsgType_Retrieve),
					commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
				),
3352
				ReqID: Params.ProxyCfg.GetNodeID(),
3353
			},
3354 3355 3356 3357
			request: queryRequest,
			qc:      node.queryCoord,
			ids:     ids.IdArray,

3358
			queryShardPolicy: mergeRoundRobinPolicy,
3359
			shardMgr:         node.shardMgr,
3360 3361
		}

G
groot 已提交
3362 3363 3364 3365 3366 3367
		items := []zapcore.Field{
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields),
		}

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

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

G
groot 已提交
3380
		log.Debug("CalcDistance queryTask enqueued", items...)
3381 3382 3383

		err = qt.WaitToFinish()
		if err != nil {
G
groot 已提交
3384
			log.Error("CalcDistance queryTask failed to WaitToFinish", append(items, zap.Error(err))...)
3385 3386 3387 3388 3389 3390

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3391
			}, err
3392
		}
3393

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

		return &milvuspb.QueryResults{
3397 3398
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3399 3400 3401
		}, nil
	}

G
groot 已提交
3402 3403 3404 3405
	// calcDistanceTask is not a standard task, no need to enqueue
	task := &calcDistanceTask{
		traceID:   traceID,
		queryFunc: query,
3406 3407
	}

G
groot 已提交
3408
	return task.Execute(ctx, request)
3409 3410
}

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

3416
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3417
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3418
	log.Debug("GetPersistentSegmentInfo",
3419
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3420 3421 3422
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3423
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3424
		Status: &commonpb.Status{
3425
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3426 3427
		},
	}
3428 3429 3430 3431
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3432 3433
	method := "GetPersistentSegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
3434
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3435
		metrics.TotalLabel).Inc()
3436 3437 3438

	// list segments
	collectionID, err := globalMetaCache.GetCollectionID(ctx, req.GetCollectionName())
X
XuanYang-cn 已提交
3439
	if err != nil {
3440 3441 3442 3443 3444 3445 3446 3447 3448 3449 3450 3451 3452
		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()
3453
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3454 3455
		return resp, nil
	}
3456 3457

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

J
jingkl 已提交
3499
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3500
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3501
	log.Debug("GetQuerySegmentInfo",
3502
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3503 3504 3505
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3506
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3507
		Status: &commonpb.Status{
3508
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3509 3510
		},
	}
3511 3512 3513 3514
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3515

3516 3517 3518 3519 3520
	method := "GetQuerySegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()

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

	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()))
3566
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3567 3568 3569 3570
	resp.Infos = queryInfos
	return resp, nil
}

J
jingkl 已提交
3571
// Dummy handles dummy request
C
Cai Yudong 已提交
3572
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3573 3574 3575 3576 3577 3578 3579 3580 3581 3582 3583
	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
	}

3584 3585
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3586
		if err != nil {
3587
			log.Debug("Failed to parse dummy query request")
3588 3589 3590
			return failedResponse, nil
		}

3591
		request := &milvuspb.QueryRequest{
3592 3593 3594
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3595
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3596 3597
		}

3598
		_, err = node.Query(ctx, request)
3599
		if err != nil {
3600
			log.Debug("Failed to execute dummy query")
3601 3602
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3603 3604 3605 3606 3607 3608

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

3609 3610
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3611 3612
}

J
jingkl 已提交
3613
// RegisterLink registers a link
C
Cai Yudong 已提交
3614
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
3615
	code := node.stateCode.Load().(commonpb.StateCode)
D
dragondriver 已提交
3616
	log.Debug("RegisterLink",
3617
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3618
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3619

3620
	if code != commonpb.StateCode_Healthy {
3621 3622 3623
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3624
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3625
				Reason:    "proxy not healthy",
3626 3627 3628
			},
		}, nil
	}
X
Xiaofan 已提交
3629
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Inc()
3630 3631 3632
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3633
			ErrorCode: commonpb.ErrorCode_Success,
3634
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3635 3636 3637
		},
	}, nil
}
3638

3639
// GetMetrics gets the metrics of proxy
3640 3641 3642
// 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 已提交
3643
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3644 3645 3646 3647
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
X
Xiaofan 已提交
3648
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3649
			zap.String("req", req.Request),
X
Xiaofan 已提交
3650
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.GetNodeID())))
3651 3652 3653 3654

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
Xiaofan 已提交
3655
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.GetNodeID()),
3656 3657 3658 3659 3660 3661 3662 3663
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
X
Xiaofan 已提交
3664
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3665 3666 3667 3668 3669 3670 3671 3672 3673 3674 3675 3676 3677 3678 3679
			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))

3680 3681 3682 3683 3684 3685
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3686
	if metricType == metricsinfo.SystemInfoMetrics {
3687 3688 3689 3690 3691 3692 3693
		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))

3694
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3695 3696

		log.Debug("Proxy.GetMetrics",
X
Xiaofan 已提交
3697
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3698 3699 3700 3701 3702
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3703 3704
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3705
		return metrics, nil
3706 3707 3708
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
X
Xiaofan 已提交
3709
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3710 3711 3712 3713 3714 3715 3716 3717 3718 3719 3720 3721
		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
}

3722 3723 3724 3725 3726 3727 3728 3729 3730 3731 3732 3733 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
// GetProxyMetrics gets the metrics of proxy, it's an internal interface which is different from GetMetrics interface,
// because it only obtains the metrics of Proxy, not including the topological metrics of Query cluster and Data cluster.
func (node *Proxy) GetProxyMetrics(ctx context.Context, req *milvuspb.GetMetricsRequest) (*milvuspb.GetMetricsResponse, error) {
	log.Debug("Proxy.GetProxyMetrics",
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
		zap.String("req", req.Request))

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

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

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

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

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

3761 3762 3763 3764 3765 3766
	req.Base = commonpbutil.NewMsgBase(
		commonpbutil.WithMsgType(commonpb.MsgType_SystemInfo),
		commonpbutil.WithMsgID(0),
		commonpbutil.WithTimeStamp(0),
		commonpbutil.WithSourceID(Params.ProxyCfg.GetNodeID()),
	)
3767 3768 3769 3770 3771 3772 3773 3774 3775 3776 3777 3778 3779 3780 3781 3782 3783 3784 3785 3786 3787 3788 3789 3790 3791 3792 3793 3794 3795 3796 3797 3798 3799 3800 3801 3802 3803 3804 3805

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

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

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

		return proxyMetrics, nil
	}

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

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

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

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

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
3819 3820 3821 3822 3823 3824 3825

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

3855 3856 3857 3858 3859 3860 3861 3862 3863 3864 3865 3866 3867 3868 3869 3870 3871 3872 3873 3874 3875 3876 3877 3878 3879
// 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
	}

	req.Base = &commonpb.MsgBase{
		MsgType:  commonpb.MsgType_GetReplicas,
		SourceID: Params.ProxyCfg.GetNodeID(),
	}

	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
}

3880
// GetCompactionState gets the compaction state of multiple segments
3881 3882 3883 3884 3885 3886 3887 3888 3889 3890 3891 3892 3893
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
}

3894
// ManualCompaction invokes compaction on specified collection
3895 3896 3897 3898 3899 3900 3901 3902 3903 3904 3905 3906 3907
func (node *Proxy) ManualCompaction(ctx context.Context, req *milvuspb.ManualCompactionRequest) (*milvuspb.ManualCompactionResponse, error) {
	log.Info("received ManualCompaction request", zap.Int64("collectionID", req.GetCollectionID()))
	resp := &milvuspb.ManualCompactionResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

	resp, err := node.dataCoord.ManualCompaction(ctx, req)
	log.Info("received ManualCompaction response", zap.Int64("collectionID", req.GetCollectionID()), zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

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

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

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

3948 3949 3950
func (node *Proxy) checkHealthyAndReturnCode() (commonpb.StateCode, bool) {
	code := node.stateCode.Load().(commonpb.StateCode)
	return code, code == commonpb.StateCode_Healthy
3951 3952
}

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

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

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

3982
	// Call rootCoord to finish import.
3983 3984
	respFromRC, err := node.rootCoord.Import(ctx, req)
	if err != nil {
3985
		metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
3986 3987 3988 3989 3990
		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
	}
3991 3992 3993

	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()))
3994
	return respFromRC, nil
G
groot 已提交
3995 3996
}

3997
// GetImportState checks import task state from RootCoord.
G
groot 已提交
3998 3999 4000 4001 4002 4003 4004
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
	}
4005 4006 4007 4008
	method := "GetImportState"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
G
groot 已提交
4009 4010

	resp, err := node.rootCoord.GetImportState(ctx, req)
4011 4012 4013 4014 4015 4016 4017 4018 4019 4020 4021 4022
	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 已提交
4023 4024 4025 4026 4027 4028 4029 4030 4031 4032
}

// 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
	}
4033 4034 4035 4036
	method := "ListImportTasks"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
G
groot 已提交
4037
	resp, err := node.rootCoord.ListImportTasks(ctx, req)
4038 4039 4040 4041 4042
	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 已提交
4043 4044 4045
		return resp, nil
	}

4046 4047 4048
	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 已提交
4049 4050 4051
	return resp, err
}

4052 4053 4054 4055 4056 4057
// 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))
4058
	if !node.checkHealthy() {
4059
		return unhealthyStatus(), nil
4060
	}
4061 4062 4063 4064 4065 4066 4067 4068 4069 4070 4071 4072 4073 4074 4075 4076 4077 4078 4079 4080 4081

	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))
4082
	if !node.checkHealthy() {
4083
		return unhealthyStatus(), nil
4084
	}
4085 4086

	credInfo := &internalpb.CredentialInfo{
4087 4088
		Username:       request.Username,
		Sha256Password: request.Password,
4089 4090 4091 4092 4093 4094 4095 4096 4097 4098 4099 4100 4101 4102 4103
	}
	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) {
4104 4105
	log.Debug("CreateCredential", zap.String("role", typeutil.ProxyRole), zap.String("username", req.Username))
	if !node.checkHealthy() {
4106
		return unhealthyStatus(), nil
4107
	}
4108 4109 4110 4111 4112 4113 4114 4115 4116 4117 4118 4119 4120 4121 4122 4123 4124 4125 4126 4127 4128 4129 4130 4131 4132 4133 4134 4135 4136 4137 4138
	// 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
	}
4139

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

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

4223 4224 4225 4226 4227 4228
	if req.Username == util.UserRoot {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_DeleteCredentialFailure,
			Reason:    "user root cannot be deleted",
		}, nil
	}
4229 4230 4231 4232 4233 4234 4235 4236 4237 4238 4239 4240
	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) {
4241 4242
	log.Debug("ListCredUsers", zap.String("role", typeutil.ProxyRole))
	if !node.checkHealthy() {
4243
		return &milvuspb.ListCredUsersResponse{Status: unhealthyStatus()}, nil
4244
	}
4245
	rootCoordReq := &milvuspb.ListCredUsersRequest{
4246 4247 4248
		Base: commonpbutil.NewMsgBase(
			commonpbutil.WithMsgType(commonpb.MsgType_ListCredUsernames),
		),
4249 4250
	}
	resp, err := node.rootCoord.ListCredUsers(ctx, rootCoordReq)
4251 4252 4253 4254 4255 4256 4257 4258 4259 4260 4261 4262
	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,
		},
4263
		Usernames: resp.Usernames,
4264 4265
	}, nil
}
4266

4267 4268 4269
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 {
4270
		return errorutil.UnhealthyStatus(code), nil
4271 4272 4273 4274 4275 4276 4277 4278 4279 4280
	}

	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(),
4281
		}, nil
4282 4283 4284 4285 4286 4287 4288 4289
	}

	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(),
4290
		}, nil
4291 4292
	}
	return result, nil
4293 4294
}

4295 4296 4297
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 {
4298
		return errorutil.UnhealthyStatus(code), nil
4299 4300 4301 4302 4303
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4304
		}, nil
4305
	}
4306 4307 4308 4309 4310
	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,
4311
		}, nil
4312
	}
4313 4314 4315 4316 4317 4318
	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(),
4319
		}, nil
4320 4321
	}
	return result, nil
4322 4323
}

4324 4325 4326
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 {
4327
		return errorutil.UnhealthyStatus(code), nil
4328 4329 4330 4331 4332
	}
	if err := ValidateUsername(req.Username); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4333
		}, nil
4334 4335 4336 4337 4338
	}
	if err := ValidateRoleName(req.RoleName); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4339
		}, nil
4340 4341 4342 4343 4344 4345 4346 4347
	}

	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(),
4348
		}, nil
4349 4350
	}
	return result, nil
4351 4352
}

4353 4354 4355
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 {
4356
		return &milvuspb.SelectRoleResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4357 4358 4359 4360 4361 4362 4363 4364 4365
	}

	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(),
				},
4366
			}, nil
4367 4368 4369 4370 4371 4372 4373 4374 4375 4376 4377
		}
	}

	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(),
			},
4378
		}, nil
4379 4380
	}
	return result, nil
4381 4382
}

4383 4384 4385
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 {
4386
		return &milvuspb.SelectUserResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4387 4388 4389 4390 4391 4392 4393 4394 4395
	}

	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(),
				},
4396
			}, nil
4397 4398 4399 4400 4401 4402 4403 4404 4405 4406 4407
		}
	}

	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(),
			},
4408
		}, nil
4409 4410
	}
	return result, nil
4411 4412
}

4413 4414 4415 4416 4417 4418 4419 4420 4421 4422 4423 4424 4425 4426 4427 4428 4429 4430 4431 4432 4433 4434 4435 4436 4437 4438 4439 4440 4441 4442
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
4443 4444
}

4445 4446 4447
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 {
4448
		return errorutil.UnhealthyStatus(code), nil
4449 4450 4451 4452 4453
	}
	if err := node.validPrivilegeParams(req); err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
4454
		}, nil
4455 4456 4457 4458 4459 4460
	}
	curUser, err := GetCurUserFromContext(ctx)
	if err != nil {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
4461
		}, nil
4462 4463 4464 4465 4466 4467 4468 4469
	}
	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(),
4470
		}, nil
4471 4472
	}
	return result, nil
4473 4474
}

4475 4476 4477 4478 4479 4480 4481 4482 4483 4484 4485 4486 4487 4488 4489 4490 4491 4492 4493 4494 4495 4496 4497 4498 4499 4500 4501 4502 4503
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 {
4504
		return &milvuspb.SelectGrantResponse{Status: errorutil.UnhealthyStatus(code)}, nil
4505 4506 4507 4508 4509 4510 4511 4512
	}

	if err := node.validGrantParams(req); err != nil {
		return &milvuspb.SelectGrantResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_IllegalArgument,
				Reason:    err.Error(),
			},
4513
		}, nil
4514 4515 4516 4517 4518 4519 4520 4521 4522 4523
	}

	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(),
			},
4524
		}, nil
4525 4526 4527 4528 4529 4530 4531 4532 4533 4534 4535 4536 4537 4538 4539 4540 4541 4542 4543 4544 4545 4546 4547 4548 4549 4550 4551 4552
	}
	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
4553
}
4554 4555 4556 4557 4558 4559 4560 4561 4562 4563 4564 4565 4566 4567 4568 4569 4570 4571 4572 4573 4574

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

	err := node.multiRateLimiter.globalRateLimiter.setRates(request.GetRates())
	// TODO: set multiple rate limiter rates
	if err != nil {
		resp.Reason = err.Error()
		return resp, nil
	}
	resp.ErrorCode = commonpb.ErrorCode_Success
	return resp, nil
}
4575 4576 4577 4578 4579 4580 4581 4582 4583 4584 4585 4586 4587 4588 4589 4590 4591 4592 4593 4594 4595 4596 4597 4598 4599 4600 4601 4602 4603 4604 4605 4606 4607 4608 4609 4610 4611 4612 4613 4614 4615 4616 4617 4618 4619 4620 4621 4622 4623 4624 4625 4626 4627 4628 4629 4630 4631 4632

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")
		return &milvuspb.CheckHealthResponse{IsHealthy: false, Reasons: []string{reason}}, nil
	}

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

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

	return &milvuspb.CheckHealthResponse{IsHealthy: true}, nil
}