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

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

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

26 27
	"github.com/milvus-io/milvus/internal/util"

28
	"go.uber.org/zap"
G
groot 已提交
29
	"go.uber.org/zap/zapcore"
S
sunby 已提交
30

31
	"github.com/milvus-io/milvus/internal/common"
X
Xiangyu Wang 已提交
32
	"github.com/milvus-io/milvus/internal/log"
33
	"github.com/milvus-io/milvus/internal/metrics"
J
jaime 已提交
34
	"github.com/milvus-io/milvus/internal/mq/msgstream"
35

36
	"github.com/golang/protobuf/proto"
X
Xiangyu Wang 已提交
37 38 39 40 41 42
	"github.com/milvus-io/milvus/internal/proto/commonpb"
	"github.com/milvus-io/milvus/internal/proto/datapb"
	"github.com/milvus-io/milvus/internal/proto/internalpb"
	"github.com/milvus-io/milvus/internal/proto/milvuspb"
	"github.com/milvus-io/milvus/internal/proto/proxypb"
	"github.com/milvus-io/milvus/internal/proto/querypb"
43
	"github.com/milvus-io/milvus/internal/util/crypto"
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
)

51 52
const moduleName = "Proxy"

53
// UpdateStateCode updates the state code of Proxy.
C
Cai Yudong 已提交
54
func (node *Proxy) UpdateStateCode(code internalpb.StateCode) {
55
	node.stateCode.Store(code)
Z
zhenshan.cao 已提交
56 57
}

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

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

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

108
	collectionName := request.CollectionName
109
	collectionID := request.CollectionID
N
neza2017 已提交
110
	if globalMetaCache != nil {
111 112 113 114 115 116
		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 已提交
117
	}
118
	logutil.Logger(ctx).Info("complete to invalidate collection meta cache",
119
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
120
		zap.String("db", request.DbName),
121 122
		zap.String("collection", collectionName),
		zap.Int64("collectionID", collectionID))
D
dragondriver 已提交
123

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

130
// ReleaseDQLMessageStream release the query message stream of specific collection.
C
Cai Yudong 已提交
131
func (node *Proxy) ReleaseDQLMessageStream(ctx context.Context, request *proxypb.ReleaseDQLMessageStreamRequest) (*commonpb.Status, error) {
132 133
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to release DQL message strem",
134
		zap.Any("role", typeutil.ProxyRole),
135 136 137
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

138 139 140 141
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

142
	logutil.Logger(ctx).Debug("complete to release DQL message stream",
143
		zap.Any("role", typeutil.ProxyRole),
144 145 146 147 148 149 150 151 152
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

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

153
// CreateCollection create a collection by the schema.
154
// TODO(dragondriver): add more detailed ut for ConsistencyLevel, should we support multiple consistency level in Proxy?
C
Cai Yudong 已提交
155
func (node *Proxy) CreateCollection(ctx context.Context, request *milvuspb.CreateCollectionRequest) (*commonpb.Status, error) {
156 157 158
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
159 160 161 162

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

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

168
	cct := &createCollectionTask{
S
sunby 已提交
169
		ctx:                     ctx,
170 171
		Condition:               NewTaskCondition(ctx),
		CreateCollectionRequest: request,
172
		rootCoord:               node.rootCoord,
173 174
	}

175 176 177
	// avoid data race
	lenOfSchema := len(request.Schema)

178 179
	log.Debug(
		rpcReceived(method),
180
		zap.String("traceID", traceID),
181
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
182 183
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
184
		zap.Int("len(schema)", lenOfSchema),
185 186
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
187

188 189 190
	if err := node.sched.ddQueue.Enqueue(cct); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
191 192
			zap.Error(err),
			zap.String("traceID", traceID),
193
			zap.String("role", typeutil.ProxyRole),
194 195 196
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Int("len(schema)", lenOfSchema),
197 198
			zap.Int32("shards_num", request.ShardsNum),
			zap.String("consistency_level", request.ConsistencyLevel.String()))
199

X
Xiaofan 已提交
200
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
201
		return &commonpb.Status{
202
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
203 204 205 206
			Reason:    err.Error(),
		}, nil
	}

207 208
	log.Debug(
		rpcEnqueued(method),
209
		zap.String("traceID", traceID),
210
		zap.String("role", typeutil.ProxyRole),
211 212 213
		zap.Int64("MsgID", cct.ID()),
		zap.Uint64("BeginTs", cct.BeginTs()),
		zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
214 215
		zap.Uint64("timestamp", request.Base.Timestamp),
		zap.String("db", request.DbName),
216 217
		zap.String("collection", request.CollectionName),
		zap.Int("len(schema)", lenOfSchema),
218 219
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
220

221 222 223
	if err := cct.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
224
			zap.Error(err),
225
			zap.String("traceID", traceID),
226
			zap.String("role", typeutil.ProxyRole),
227 228 229
			zap.Int64("MsgID", cct.ID()),
			zap.Uint64("BeginTs", cct.BeginTs()),
			zap.Uint64("EndTs", cct.EndTs()),
D
dragondriver 已提交
230 231
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
232
			zap.Int("len(schema)", lenOfSchema),
233 234
			zap.Int32("shards_num", request.ShardsNum),
			zap.String("consistency_level", request.ConsistencyLevel.String()))
D
dragondriver 已提交
235

X
Xiaofan 已提交
236
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
237
		return &commonpb.Status{
238
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
239 240 241 242
			Reason:    err.Error(),
		}, nil
	}

243 244
	log.Debug(
		rpcDone(method),
245
		zap.String("traceID", traceID),
246
		zap.String("role", typeutil.ProxyRole),
247 248 249 250 251 252
		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),
253 254
		zap.Int32("shards_num", request.ShardsNum),
		zap.String("consistency_level", request.ConsistencyLevel.String()))
255

X
Xiaofan 已提交
256 257
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
258 259 260
	return cct.result, nil
}

261
// DropCollection drop a collection.
C
Cai Yudong 已提交
262
func (node *Proxy) DropCollection(ctx context.Context, request *milvuspb.DropCollectionRequest) (*commonpb.Status, error) {
263 264 265
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
266 267 268 269

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

274
	dct := &dropCollectionTask{
S
sunby 已提交
275
		ctx:                   ctx,
276 277
		Condition:             NewTaskCondition(ctx),
		DropCollectionRequest: request,
278
		rootCoord:             node.rootCoord,
279
		chMgr:                 node.chMgr,
S
sunby 已提交
280
		chTicker:              node.chTicker,
281 282
	}

283 284
	log.Debug("DropCollection received",
		zap.String("traceID", traceID),
285
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
286 287
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
288 289 290 291 292

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

X
Xiaofan 已提交
297
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
298
		return &commonpb.Status{
299
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
300 301 302 303
			Reason:    err.Error(),
		}, nil
	}

304 305
	log.Debug("DropCollection enqueued",
		zap.String("traceID", traceID),
306
		zap.String("role", typeutil.ProxyRole),
307 308 309
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
310 311
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
312 313 314

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DropCollection failed to WaitToFinish",
D
dragondriver 已提交
315
			zap.Error(err),
316
			zap.String("traceID", traceID),
317
			zap.String("role", typeutil.ProxyRole),
318 319 320
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTs", dct.BeginTs()),
			zap.Uint64("EndTs", dct.EndTs()),
D
dragondriver 已提交
321 322 323
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
324
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.FailLabel).Inc()
325
		return &commonpb.Status{
326
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
327 328 329 330
			Reason:    err.Error(),
		}, nil
	}

331 332
	log.Debug("DropCollection done",
		zap.String("traceID", traceID),
333
		zap.String("role", typeutil.ProxyRole),
334 335 336 337 338 339
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTs", dct.BeginTs()),
		zap.Uint64("EndTs", dct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
340 341
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
342 343 344
	return dct.result, nil
}

345
// HasCollection check if the specific collection exists in Milvus.
C
Cai Yudong 已提交
346
func (node *Proxy) HasCollection(ctx context.Context, request *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
347 348 349 350 351
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
352 353 354 355

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
356 357
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
358
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
359
		metrics.TotalLabel).Inc()
360 361 362

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
363
		zap.String("role", typeutil.ProxyRole),
364 365 366
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

367
	hct := &hasCollectionTask{
S
sunby 已提交
368
		ctx:                  ctx,
369 370
		Condition:            NewTaskCondition(ctx),
		HasCollectionRequest: request,
371
		rootCoord:            node.rootCoord,
372 373
	}

374 375 376 377
	if err := node.sched.ddQueue.Enqueue(hct); err != nil {
		log.Warn("HasCollection failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
378
			zap.String("role", typeutil.ProxyRole),
379 380 381
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
382
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
383
			metrics.AbandonLabel).Inc()
384 385
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
386
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
387 388 389 390 391
				Reason:    err.Error(),
			},
		}, nil
	}

392 393
	log.Debug("HasCollection enqueued",
		zap.String("traceID", traceID),
394
		zap.String("role", typeutil.ProxyRole),
395 396 397
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
398 399
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
400 401 402

	if err := hct.WaitToFinish(); err != nil {
		log.Warn("HasCollection failed to WaitToFinish",
D
dragondriver 已提交
403
			zap.Error(err),
404
			zap.String("traceID", traceID),
405
			zap.String("role", typeutil.ProxyRole),
406 407 408
			zap.Int64("MsgID", hct.ID()),
			zap.Uint64("BeginTS", hct.BeginTs()),
			zap.Uint64("EndTS", hct.EndTs()),
D
dragondriver 已提交
409 410 411
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
412
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
413
			metrics.FailLabel).Inc()
414 415
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
416
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
417 418 419 420 421
				Reason:    err.Error(),
			},
		}, nil
	}

422 423
	log.Debug("HasCollection done",
		zap.String("traceID", traceID),
424
		zap.String("role", typeutil.ProxyRole),
425 426 427 428 429 430
		zap.Int64("MsgID", hct.ID()),
		zap.Uint64("BeginTS", hct.BeginTs()),
		zap.Uint64("EndTS", hct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
431
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
432
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
433
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
434 435 436
	return hct.result, nil
}

437
// LoadCollection load a collection into query nodes.
C
Cai Yudong 已提交
438
func (node *Proxy) LoadCollection(ctx context.Context, request *milvuspb.LoadCollectionRequest) (*commonpb.Status, error) {
439 440 441
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
442 443 444 445

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

449
	lct := &loadCollectionTask{
S
sunby 已提交
450
		ctx:                   ctx,
451 452
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
453
		queryCoord:            node.queryCoord,
454 455
	}

456 457
	log.Debug("LoadCollection received",
		zap.String("traceID", traceID),
458
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
459 460
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
461 462 463 464 465

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

X
Xiaofan 已提交
470
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
471
			metrics.AbandonLabel).Inc()
472
		return &commonpb.Status{
473
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
474 475 476
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
477

478 479
	log.Debug("LoadCollection enqueued",
		zap.String("traceID", traceID),
480
		zap.String("role", typeutil.ProxyRole),
481 482 483
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
484 485
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
486 487 488

	if err := lct.WaitToFinish(); err != nil {
		log.Warn("LoadCollection failed to WaitToFinish",
D
dragondriver 已提交
489
			zap.Error(err),
490
			zap.String("traceID", traceID),
491
			zap.String("role", typeutil.ProxyRole),
492 493 494
			zap.Int64("MsgID", lct.ID()),
			zap.Uint64("BeginTS", lct.BeginTs()),
			zap.Uint64("EndTS", lct.EndTs()),
D
dragondriver 已提交
495 496 497
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
498
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
499
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
500
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
501
			metrics.FailLabel).Inc()
502
		return &commonpb.Status{
503
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
504 505 506 507
			Reason:    err.Error(),
		}, nil
	}

508 509
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
510
		zap.String("role", typeutil.ProxyRole),
511 512 513 514 515 516
		zap.Int64("MsgID", lct.ID()),
		zap.Uint64("BeginTS", lct.BeginTs()),
		zap.Uint64("EndTS", lct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
517
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
518
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
519
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
520
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
521
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
522
	return lct.result, nil
523 524
}

525
// ReleaseCollection remove the loaded collection from query nodes.
C
Cai Yudong 已提交
526
func (node *Proxy) ReleaseCollection(ctx context.Context, request *milvuspb.ReleaseCollectionRequest) (*commonpb.Status, error) {
527 528 529
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
530

531
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
532 533
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
534 535
	method := "ReleaseCollection"
	tr := timerecord.NewTimeRecorder(method)
536

537
	rct := &releaseCollectionTask{
S
sunby 已提交
538
		ctx:                      ctx,
539 540
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
541
		queryCoord:               node.queryCoord,
542
		chMgr:                    node.chMgr,
543 544
	}

545 546
	log.Debug(
		rpcReceived(method),
547
		zap.String("traceID", traceID),
548
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
549 550
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
551 552

	if err := node.sched.ddQueue.Enqueue(rct); err != nil {
553 554
		log.Warn(
			rpcFailedToEnqueue(method),
555 556
			zap.Error(err),
			zap.String("traceID", traceID),
557
			zap.String("role", typeutil.ProxyRole),
558 559 560
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
561
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
562
			metrics.AbandonLabel).Inc()
563
		return &commonpb.Status{
564
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
565 566 567 568
			Reason:    err.Error(),
		}, nil
	}

569 570
	log.Debug(
		rpcEnqueued(method),
571
		zap.String("traceID", traceID),
572
		zap.String("role", typeutil.ProxyRole),
573 574 575
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
576 577
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
578 579

	if err := rct.WaitToFinish(); err != nil {
580 581
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
582
			zap.Error(err),
583
			zap.String("traceID", traceID),
584
			zap.String("role", typeutil.ProxyRole),
585 586 587
			zap.Int64("MsgID", rct.ID()),
			zap.Uint64("BeginTS", rct.BeginTs()),
			zap.Uint64("EndTS", rct.EndTs()),
D
dragondriver 已提交
588 589 590
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
591
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
592
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
593
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
594
			metrics.FailLabel).Inc()
595
		return &commonpb.Status{
596
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
597 598 599 600
			Reason:    err.Error(),
		}, nil
	}

601 602
	log.Debug(
		rpcDone(method),
603
		zap.String("traceID", traceID),
604
		zap.String("role", typeutil.ProxyRole),
605 606 607 608 609 610
		zap.Int64("MsgID", rct.ID()),
		zap.Uint64("BeginTS", rct.BeginTs()),
		zap.Uint64("EndTS", rct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
611
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
612
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
613
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
614
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
615
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
616
	return rct.result, nil
617 618
}

619
// DescribeCollection get the meta information of specific collection, such as schema, created timestamp and etc.
C
Cai Yudong 已提交
620
func (node *Proxy) DescribeCollection(ctx context.Context, request *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
621 622 623 624 625
	if !node.checkHealthy() {
		return &milvuspb.DescribeCollectionResponse{
			Status: unhealthyStatus(),
		}, nil
	}
626

627
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
628 629
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
630 631
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
632

633
	dct := &describeCollectionTask{
S
sunby 已提交
634
		ctx:                       ctx,
635 636
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
637
		rootCoord:                 node.rootCoord,
638 639
	}

640 641
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
642
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
643 644
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
645 646 647 648 649

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

X
Xiaofan 已提交
654
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
655
			metrics.AbandonLabel).Inc()
656 657
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
658
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
659 660 661 662 663
				Reason:    err.Error(),
			},
		}, nil
	}

664 665
	log.Debug("DescribeCollection enqueued",
		zap.String("traceID", traceID),
666
		zap.String("role", typeutil.ProxyRole),
667 668 669
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
670 671
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
672 673 674

	if err := dct.WaitToFinish(); err != nil {
		log.Warn("DescribeCollection failed to WaitToFinish",
D
dragondriver 已提交
675
			zap.Error(err),
676
			zap.String("traceID", traceID),
677
			zap.String("role", typeutil.ProxyRole),
678 679 680
			zap.Int64("MsgID", dct.ID()),
			zap.Uint64("BeginTS", dct.BeginTs()),
			zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
681 682 683
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
684
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
685
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
686
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
687
			metrics.FailLabel).Inc()
688

689 690
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
691
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
692 693 694 695 696
				Reason:    err.Error(),
			},
		}, nil
	}

697 698
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
699
		zap.String("role", typeutil.ProxyRole),
700 701 702 703 704 705
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
706
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
707
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
708
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
709
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
710
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
711 712 713
	return dct.result, nil
}

714
// GetCollectionStatistics get the collection statistics, such as `num_rows`.
C
Cai Yudong 已提交
715
func (node *Proxy) GetCollectionStatistics(ctx context.Context, request *milvuspb.GetCollectionStatisticsRequest) (*milvuspb.GetCollectionStatisticsResponse, error) {
716 717 718 719 720
	if !node.checkHealthy() {
		return &milvuspb.GetCollectionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
721 722 723 724

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

728
	g := &getCollectionStatisticsTask{
G
godchen 已提交
729 730 731
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
732
		dataCoord:                      node.dataCoord,
733 734
	}

735 736
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
737
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
738 739
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
740 741 742 743 744

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn("GetCollectionStatistics failed to enqueue",
			zap.Error(err),
			zap.String("traceID", traceID),
745
			zap.String("role", typeutil.ProxyRole),
746 747 748
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
749
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
750
			metrics.AbandonLabel).Inc()
751

G
godchen 已提交
752
		return &milvuspb.GetCollectionStatisticsResponse{
753
			Status: &commonpb.Status{
754
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
755 756 757 758 759
				Reason:    err.Error(),
			},
		}, nil
	}

760 761
	log.Debug("GetCollectionStatistics enqueued",
		zap.String("traceID", traceID),
762
		zap.String("role", typeutil.ProxyRole),
763 764 765
		zap.Int64("MsgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
766 767
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
768 769 770

	if err := g.WaitToFinish(); err != nil {
		log.Warn("GetCollectionStatistics failed to WaitToFinish",
D
dragondriver 已提交
771
			zap.Error(err),
772
			zap.String("traceID", traceID),
773
			zap.String("role", typeutil.ProxyRole),
774 775 776
			zap.Int64("MsgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
777 778 779
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
780
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
781
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
782
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
783
			metrics.FailLabel).Inc()
784

G
godchen 已提交
785
		return &milvuspb.GetCollectionStatisticsResponse{
786
			Status: &commonpb.Status{
787
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
788 789 790 791 792
				Reason:    err.Error(),
			},
		}, nil
	}

793 794
	log.Debug("GetCollectionStatistics done",
		zap.String("traceID", traceID),
795
		zap.String("role", typeutil.ProxyRole),
796 797 798 799 800 801
		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))

X
Xiaofan 已提交
802
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
803
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
804
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
805
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
806
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
807
	return g.result, nil
808 809
}

810
// ShowCollections list all collections in Milvus.
C
Cai Yudong 已提交
811
func (node *Proxy) ShowCollections(ctx context.Context, request *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
812 813 814 815 816
	if !node.checkHealthy() {
		return &milvuspb.ShowCollectionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
817 818
	method := "ShowCollections"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
819
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
820

821
	sct := &showCollectionsTask{
G
godchen 已提交
822 823 824
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
825
		queryCoord:             node.queryCoord,
826
		rootCoord:              node.rootCoord,
827 828
	}

829
	log.Debug("ShowCollections received",
830
		zap.String("role", typeutil.ProxyRole),
831 832 833 834 835 836
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

837
	err := node.sched.ddQueue.Enqueue(sct)
838
	if err != nil {
839 840
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
841
			zap.String("role", typeutil.ProxyRole),
842 843 844 845 846 847
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

X
Xiaofan 已提交
848
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
G
godchen 已提交
849
		return &milvuspb.ShowCollectionsResponse{
850
			Status: &commonpb.Status{
851
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
852 853 854 855 856
				Reason:    err.Error(),
			},
		}, nil
	}

857
	log.Debug("ShowCollections enqueued",
858
		zap.String("role", typeutil.ProxyRole),
859
		zap.Int64("MsgID", sct.ID()),
860
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
861
		zap.Uint64("TimeStamp", request.TimeStamp),
862 863 864
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
865

866 867
	err = sct.WaitToFinish()
	if err != nil {
868 869
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
870
			zap.String("role", typeutil.ProxyRole),
871 872 873 874 875 876 877
			zap.Int64("MsgID", sct.ID()),
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

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

G
godchen 已提交
880
		return &milvuspb.ShowCollectionsResponse{
881
			Status: &commonpb.Status{
882
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
883 884 885 886 887
				Reason:    err.Error(),
			},
		}, nil
	}

888
	log.Debug("ShowCollections Done",
889
		zap.String("role", typeutil.ProxyRole),
890 891 892 893
		zap.Int64("MsgID", sct.ID()),
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
894 895
		zap.Int("len(CollectionNames)", len(request.CollectionNames)),
		zap.Int("num_collections", len(sct.result.CollectionNames)))
896

X
Xiaofan 已提交
897 898
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
899 900 901
	return sct.result, nil
}

902
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
903
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
904 905 906
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
907

908
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
909 910
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
911 912
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
913
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
914

915
	cpt := &createPartitionTask{
S
sunby 已提交
916
		ctx:                    ctx,
917 918
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
919
		rootCoord:              node.rootCoord,
920 921 922
		result:                 nil,
	}

923 924 925
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
926
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
927 928 929
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
930 931 932 933 934 935

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
936
			zap.String("role", typeutil.ProxyRole),
937 938 939 940
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

943
		return &commonpb.Status{
944
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
945 946 947
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
948

949 950 951
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
952
		zap.String("role", typeutil.ProxyRole),
953 954 955
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
956 957 958
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
959 960 961 962

	if err := cpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish("CreatePartition"),
D
dragondriver 已提交
963
			zap.Error(err),
964
			zap.String("traceID", traceID),
965
			zap.String("role", typeutil.ProxyRole),
966 967 968
			zap.Int64("MsgID", cpt.ID()),
			zap.Uint64("BeginTS", cpt.BeginTs()),
			zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
969 970 971 972
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

975
		return &commonpb.Status{
976
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
977 978 979
			Reason:    err.Error(),
		}, nil
	}
980 981 982 983

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
984
		zap.String("role", typeutil.ProxyRole),
985 986 987 988 989 990 991
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

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

997
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
998
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
999 1000 1001
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1002

1003
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
1004 1005
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1006 1007
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
1008
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
1009

1010
	dpt := &dropPartitionTask{
S
sunby 已提交
1011
		ctx:                  ctx,
1012 1013
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1014
		rootCoord:            node.rootCoord,
1015 1016 1017
		result:               nil,
	}

1018 1019 1020
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1021
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1022 1023 1024
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1025 1026 1027 1028 1029 1030

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1031
			zap.String("role", typeutil.ProxyRole),
1032 1033 1034 1035
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1038
		return &commonpb.Status{
1039
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1040 1041 1042
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1043

1044 1045 1046
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1047
		zap.String("role", typeutil.ProxyRole),
1048 1049 1050
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1051 1052 1053
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1054 1055 1056 1057

	if err := dpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1058
			zap.Error(err),
1059
			zap.String("traceID", traceID),
1060
			zap.String("role", typeutil.ProxyRole),
1061 1062 1063
			zap.Int64("MsgID", dpt.ID()),
			zap.Uint64("BeginTS", dpt.BeginTs()),
			zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1064 1065 1066 1067
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

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

1070
		return &commonpb.Status{
1071
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1072 1073 1074
			Reason:    err.Error(),
		}, nil
	}
1075 1076 1077 1078

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1079
		zap.String("role", typeutil.ProxyRole),
1080 1081 1082 1083 1084 1085 1086
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1087 1088
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1089 1090 1091
	return dpt.result, nil
}

1092
// HasPartition check if partition exist.
C
Cai Yudong 已提交
1093
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
1094 1095 1096 1097 1098
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1099

D
dragondriver 已提交
1100
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
1101 1102
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1103 1104 1105
	method := "HasPartition"
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
X
Xiaofan 已提交
1106
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1107
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1108

1109
	hpt := &hasPartitionTask{
S
sunby 已提交
1110
		ctx:                 ctx,
1111 1112
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1113
		rootCoord:           node.rootCoord,
1114 1115 1116
		result:              nil,
	}

D
dragondriver 已提交
1117 1118 1119
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1120
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1121 1122 1123
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1124 1125 1126 1127 1128 1129

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

X
Xiaofan 已提交
1135
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1136
			metrics.AbandonLabel).Inc()
1137

1138 1139
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1140
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1141 1142 1143 1144 1145
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1146

D
dragondriver 已提交
1147 1148 1149
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1150
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1151 1152 1153
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1154 1155 1156
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1157 1158 1159 1160

	if err := hpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1161
			zap.Error(err),
D
dragondriver 已提交
1162
			zap.String("traceID", traceID),
1163
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1164 1165 1166
			zap.Int64("MsgID", hpt.ID()),
			zap.Uint64("BeginTS", hpt.BeginTs()),
			zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1167 1168 1169 1170
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1171
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1172
			metrics.FailLabel).Inc()
1173

1174 1175
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1176
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1177 1178 1179 1180 1181
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1182 1183 1184 1185

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1186
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1187 1188 1189 1190 1191 1192 1193
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1194
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1195
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1196
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1197 1198 1199
	return hpt.result, nil
}

1200
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1201
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1202 1203 1204
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1205

D
dragondriver 已提交
1206
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1207 1208
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1209 1210
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1211

1212
	lpt := &loadPartitionsTask{
G
godchen 已提交
1213 1214 1215
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1216
		queryCoord:            node.queryCoord,
1217 1218
	}

1219 1220 1221
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1222
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1223 1224 1225
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1226 1227 1228 1229 1230 1231

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1232
			zap.String("role", typeutil.ProxyRole),
1233 1234 1235 1236
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1237
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1238
			metrics.AbandonLabel).Inc()
1239

1240
		return &commonpb.Status{
1241
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1242 1243 1244 1245
			Reason:    err.Error(),
		}, nil
	}

1246 1247 1248
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1249
		zap.String("role", typeutil.ProxyRole),
1250 1251 1252
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1253 1254 1255
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1256 1257 1258 1259

	if err := lpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1260
			zap.Error(err),
1261
			zap.String("traceID", traceID),
1262
			zap.String("role", typeutil.ProxyRole),
1263 1264 1265
			zap.Int64("MsgID", lpt.ID()),
			zap.Uint64("BeginTS", lpt.BeginTs()),
			zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1266 1267 1268 1269
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1270
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1271
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1272
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1273
			metrics.FailLabel).Inc()
1274

1275
		return &commonpb.Status{
1276
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1277 1278 1279 1280
			Reason:    err.Error(),
		}, nil
	}

1281 1282 1283
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1284
		zap.String("role", typeutil.ProxyRole),
1285 1286 1287 1288 1289 1290 1291
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1292
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1293
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1294
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1295
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1296
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1297
	return lpt.result, nil
1298 1299
}

1300
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1301
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1302 1303 1304
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1305 1306 1307 1308 1309

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

1310
	rpt := &releasePartitionsTask{
G
godchen 已提交
1311 1312 1313
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1314
		queryCoord:               node.queryCoord,
1315 1316
	}

1317
	method := "ReleasePartitions"
1318
	tr := timerecord.NewTimeRecorder(method)
1319 1320 1321 1322

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1323
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1324 1325 1326
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1327 1328 1329 1330 1331 1332

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1333
			zap.String("role", typeutil.ProxyRole),
1334 1335 1336 1337
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1338
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1339
			metrics.AbandonLabel).Inc()
1340

1341
		return &commonpb.Status{
1342
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1343 1344 1345 1346
			Reason:    err.Error(),
		}, nil
	}

1347 1348 1349
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1350
		zap.String("role", typeutil.ProxyRole),
1351 1352 1353
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1354 1355 1356
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1357 1358 1359 1360

	if err := rpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1361
			zap.Error(err),
1362
			zap.String("traceID", traceID),
1363
			zap.String("role", typeutil.ProxyRole),
1364 1365 1366
			zap.Int64("msgID", rpt.Base.MsgID),
			zap.Uint64("BeginTS", rpt.BeginTs()),
			zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1367 1368 1369 1370
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1371
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1372
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1373
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1374
			metrics.FailLabel).Inc()
1375

1376
		return &commonpb.Status{
1377
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1378 1379 1380 1381
			Reason:    err.Error(),
		}, nil
	}

1382 1383 1384
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1385
		zap.String("role", typeutil.ProxyRole),
1386 1387 1388 1389 1390 1391 1392
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))

X
Xiaofan 已提交
1393
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1394
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1395
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1396
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1397
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1398
	return rpt.result, nil
1399 1400
}

1401
// GetPartitionStatistics get the statistics of partition, such as num_rows.
C
Cai Yudong 已提交
1402
func (node *Proxy) GetPartitionStatistics(ctx context.Context, request *milvuspb.GetPartitionStatisticsRequest) (*milvuspb.GetPartitionStatisticsResponse, error) {
1403 1404 1405 1406 1407
	if !node.checkHealthy() {
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1408 1409 1410 1411 1412

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

1413
	g := &getPartitionStatisticsTask{
1414 1415 1416
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1417
		dataCoord:                     node.dataCoord,
1418 1419
	}

1420
	method := "GetPartitionStatistics"
1421
	tr := timerecord.NewTimeRecorder(method)
1422 1423 1424 1425

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1426
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1427 1428 1429
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1430 1431 1432 1433 1434 1435

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1436
			zap.String("role", typeutil.ProxyRole),
1437 1438 1439 1440
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1441
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1442
			metrics.AbandonLabel).Inc()
1443

1444 1445 1446 1447 1448 1449 1450 1451
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1452 1453 1454
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1455
		zap.String("role", typeutil.ProxyRole),
1456 1457 1458
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1459 1460 1461
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1462 1463 1464 1465

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1466
			zap.Error(err),
1467
			zap.String("traceID", traceID),
1468
			zap.String("role", typeutil.ProxyRole),
1469 1470 1471
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1472 1473 1474 1475
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1476
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1477
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1478
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1479
			metrics.FailLabel).Inc()
1480

1481 1482 1483 1484 1485 1486 1487 1488
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1489 1490 1491
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1492
		zap.String("role", typeutil.ProxyRole),
1493 1494 1495 1496 1497 1498 1499
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))

X
Xiaofan 已提交
1500
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1501
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1502
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1503
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1504
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1505
	return g.result, nil
1506 1507
}

1508
// ShowPartitions list all partitions in the specific collection.
C
Cai Yudong 已提交
1509
func (node *Proxy) ShowPartitions(ctx context.Context, request *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
1510 1511 1512 1513 1514
	if !node.checkHealthy() {
		return &milvuspb.ShowPartitionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1515 1516 1517 1518 1519

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

1520
	spt := &showPartitionsTask{
G
godchen 已提交
1521 1522 1523
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1524
		rootCoord:             node.rootCoord,
1525
		queryCoord:            node.queryCoord,
G
godchen 已提交
1526
		result:                nil,
1527 1528
	}

1529
	method := "ShowPartitions"
1530 1531
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
X
Xiaofan 已提交
1532
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1533
		metrics.TotalLabel).Inc()
1534 1535 1536 1537

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1538
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1539
		zap.Any("request", request))
1540 1541 1542 1543 1544 1545

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

X
Xiaofan 已提交
1549
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1550
			metrics.AbandonLabel).Inc()
1551

G
godchen 已提交
1552
		return &milvuspb.ShowPartitionsResponse{
1553
			Status: &commonpb.Status{
1554
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1555 1556 1557 1558 1559
				Reason:    err.Error(),
			},
		}, nil
	}

1560 1561 1562
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1563
		zap.String("role", typeutil.ProxyRole),
1564 1565 1566
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1567 1568
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1569 1570 1571 1572 1573
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1574
			zap.Error(err),
1575
			zap.String("traceID", traceID),
1576
			zap.String("role", typeutil.ProxyRole),
1577 1578 1579 1580 1581 1582
			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 已提交
1583

X
Xiaofan 已提交
1584
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1585
			metrics.FailLabel).Inc()
1586

G
godchen 已提交
1587
		return &milvuspb.ShowPartitionsResponse{
1588
			Status: &commonpb.Status{
1589
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1590 1591 1592 1593
				Reason:    err.Error(),
			},
		}, nil
	}
1594 1595 1596 1597

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1598
		zap.String("role", typeutil.ProxyRole),
1599 1600 1601 1602 1603 1604 1605
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

X
Xiaofan 已提交
1606
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1607
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1608
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1609 1610 1611
	return spt.result, nil
}

1612
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1613
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1614 1615 1616
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1617 1618 1619 1620 1621

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

1622
	cit := &createIndexTask{
S
sunby 已提交
1623
		ctx:                ctx,
1624 1625
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1626
		rootCoord:          node.rootCoord,
1627 1628
	}

D
dragondriver 已提交
1629
	method := "CreateIndex"
1630
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1631 1632 1633 1634

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1635
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1636 1637 1638 1639
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1640 1641 1642 1643 1644 1645

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1646
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1647 1648 1649 1650 1651
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

X
Xiaofan 已提交
1652
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1653
			metrics.AbandonLabel).Inc()
1654

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

D
dragondriver 已提交
1661 1662 1663
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1664
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1665 1666 1667
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1668 1669 1670 1671
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1672 1673 1674 1675

	if err := cit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1676
			zap.Error(err),
D
dragondriver 已提交
1677
			zap.String("traceID", traceID),
1678
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1679 1680 1681
			zap.Int64("MsgID", cit.ID()),
			zap.Uint64("BeginTs", cit.BeginTs()),
			zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1682 1683 1684 1685 1686
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

X
Xiaofan 已提交
1687
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1688
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1689
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1690
			metrics.FailLabel).Inc()
1691

1692
		return &commonpb.Status{
1693
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1694 1695 1696 1697
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1698 1699 1700
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1701
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1702 1703 1704 1705 1706 1707 1708 1709
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))

X
Xiaofan 已提交
1710
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1711
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1712
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1713
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1714
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1715 1716 1717
	return cit.result, nil
}

1718
// DescribeIndex get the meta information of index, such as index state, index id and etc.
C
Cai Yudong 已提交
1719
func (node *Proxy) DescribeIndex(ctx context.Context, request *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
1720 1721 1722 1723 1724
	if !node.checkHealthy() {
		return &milvuspb.DescribeIndexResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1725 1726 1727 1728 1729

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

1730
	dit := &describeIndexTask{
S
sunby 已提交
1731
		ctx:                  ctx,
1732 1733
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1734
		rootCoord:            node.rootCoord,
1735 1736
	}

1737 1738 1739
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
1740
	tr := timerecord.NewTimeRecorder(method)
1741 1742 1743 1744

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1745
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1746 1747 1748
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1749 1750 1751 1752 1753 1754 1755
		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),
1756
			zap.String("role", typeutil.ProxyRole),
1757 1758 1759 1760 1761
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", indexName))

X
Xiaofan 已提交
1762
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1763
			metrics.AbandonLabel).Inc()
1764

1765 1766
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1767
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1768 1769 1770 1771 1772
				Reason:    err.Error(),
			},
		}, nil
	}

1773 1774 1775
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1776
		zap.String("role", typeutil.ProxyRole),
1777 1778 1779
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1780 1781 1782
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1783 1784 1785 1786 1787
		zap.String("index name", indexName))

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1788
			zap.Error(err),
1789
			zap.String("traceID", traceID),
1790
			zap.String("role", typeutil.ProxyRole),
1791 1792 1793
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1794 1795 1796
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
1797
			zap.String("index name", indexName))
D
dragondriver 已提交
1798

Z
zhenshan.cao 已提交
1799 1800 1801 1802
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
X
Xiaofan 已提交
1803
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1804
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1805
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1806
			metrics.FailLabel).Inc()
1807

1808 1809
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
1810
				ErrorCode: errCode,
1811 1812 1813 1814 1815
				Reason:    err.Error(),
			},
		}, nil
	}

1816 1817 1818
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1819
		zap.String("role", typeutil.ProxyRole),
1820 1821 1822 1823 1824 1825 1826 1827
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", indexName))

X
Xiaofan 已提交
1828
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1829
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1830
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1831
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1832
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1833 1834 1835
	return dit.result, nil
}

1836
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
1837
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
1838 1839 1840
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1841 1842 1843 1844 1845

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

1846
	dit := &dropIndexTask{
S
sunby 已提交
1847
		ctx:              ctx,
B
BossZou 已提交
1848 1849
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1850
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1851
	}
G
godchen 已提交
1852

D
dragondriver 已提交
1853
	method := "DropIndex"
1854
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1855 1856 1857 1858

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1859
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1860 1861 1862 1863 1864
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
1865 1866 1867 1868 1869
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1870
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1871 1872 1873 1874
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
1875
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1876
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
1877

B
BossZou 已提交
1878
		return &commonpb.Status{
1879
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1880 1881 1882
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1883

D
dragondriver 已提交
1884 1885 1886
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1887
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1888 1889 1890
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1891 1892 1893 1894
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
1895 1896 1897 1898

	if err := dit.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1899
			zap.Error(err),
D
dragondriver 已提交
1900
			zap.String("traceID", traceID),
1901
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1902 1903 1904
			zap.Int64("MsgID", dit.ID()),
			zap.Uint64("BeginTs", dit.BeginTs()),
			zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1905 1906 1907 1908 1909
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
1910
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1911
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1912
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1913
			metrics.FailLabel).Inc()
1914

B
BossZou 已提交
1915
		return &commonpb.Status{
1916
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1917 1918 1919
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1920 1921 1922 1923

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1924
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1925 1926 1927 1928 1929 1930 1931 1932
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

X
Xiaofan 已提交
1933
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1934
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
1935
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1936
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
1937
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
1938 1939 1940
	return dit.result, nil
}

1941 1942
// GetIndexBuildProgress gets index build progress with filed_name and index_name.
// IndexRows is the num of indexed rows. And TotalRows is the total number of segment rows.
C
Cai Yudong 已提交
1943
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1944 1945 1946 1947 1948
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1949 1950 1951 1952 1953

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

1954
	gibpt := &getIndexBuildProgressTask{
1955 1956 1957
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1958 1959
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1960
		dataCoord:                    node.dataCoord,
1961 1962
	}

1963
	method := "GetIndexBuildProgress"
1964
	tr := timerecord.NewTimeRecorder(method)
1965 1966 1967 1968

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1969
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1970 1971 1972 1973
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1974 1975 1976 1977 1978 1979

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1980
			zap.String("role", typeutil.ProxyRole),
1981 1982 1983 1984
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
1985
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
1986
			metrics.AbandonLabel).Inc()
1987

1988 1989 1990 1991 1992 1993 1994 1995
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1996 1997 1998
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1999
		zap.String("role", typeutil.ProxyRole),
2000 2001 2002
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
2003 2004 2005 2006
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2007 2008 2009 2010

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2011
			zap.Error(err),
2012
			zap.String("traceID", traceID),
2013
			zap.String("role", typeutil.ProxyRole),
2014 2015 2016
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2017 2018 2019 2020
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
X
Xiaofan 已提交
2021
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2022
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2023
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2024
			metrics.FailLabel).Inc()
2025 2026 2027 2028 2029 2030 2031 2032

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2033 2034 2035 2036

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2037
		zap.String("role", typeutil.ProxyRole),
2038 2039 2040 2041 2042 2043 2044 2045
		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))
2046

X
Xiaofan 已提交
2047
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2048
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2049
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2050
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2051
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2052
	return gibpt.result, nil
2053 2054
}

2055
// GetIndexState get the build-state of index.
C
Cai Yudong 已提交
2056
func (node *Proxy) GetIndexState(ctx context.Context, request *milvuspb.GetIndexStateRequest) (*milvuspb.GetIndexStateResponse, error) {
2057 2058 2059 2060 2061
	if !node.checkHealthy() {
		return &milvuspb.GetIndexStateResponse{
			Status: unhealthyStatus(),
		}, nil
	}
2062 2063 2064 2065 2066

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

2067
	dipt := &getIndexStateTask{
G
godchen 已提交
2068 2069 2070
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2071 2072
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2073 2074
	}

2075
	method := "GetIndexState"
2076
	tr := timerecord.NewTimeRecorder(method)
2077 2078 2079 2080

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2081
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2082 2083 2084 2085
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2086 2087 2088 2089 2090 2091

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2092
			zap.String("role", typeutil.ProxyRole),
2093 2094 2095 2096 2097
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2098
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2099
			metrics.AbandonLabel).Inc()
2100

G
godchen 已提交
2101
		return &milvuspb.GetIndexStateResponse{
2102
			Status: &commonpb.Status{
2103
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2104 2105 2106 2107 2108
				Reason:    err.Error(),
			},
		}, nil
	}

2109 2110 2111
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2112
		zap.String("role", typeutil.ProxyRole),
2113 2114 2115
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2116 2117 2118 2119
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2120 2121 2122 2123

	if err := dipt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2124
			zap.Error(err),
2125
			zap.String("traceID", traceID),
2126
			zap.String("role", typeutil.ProxyRole),
2127 2128 2129
			zap.Int64("MsgID", dipt.ID()),
			zap.Uint64("BeginTs", dipt.BeginTs()),
			zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2130 2131 2132 2133 2134
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2135
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2136
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2137
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2138
			metrics.FailLabel).Inc()
2139

G
godchen 已提交
2140
		return &milvuspb.GetIndexStateResponse{
2141
			Status: &commonpb.Status{
2142
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2143 2144 2145 2146 2147
				Reason:    err.Error(),
			},
		}, nil
	}

2148 2149 2150
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2151
		zap.String("role", typeutil.ProxyRole),
2152 2153 2154 2155 2156 2157 2158 2159
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

X
Xiaofan 已提交
2160
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2161
		metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2162
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2163
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
2164
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2165 2166 2167
	return dipt.result, nil
}

2168
// Insert insert records into collection.
C
Cai Yudong 已提交
2169
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2170 2171 2172 2173 2174 2175
	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))

2176 2177 2178 2179 2180
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2181 2182
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
2183 2184
	receiveSize := proto.Size(request)
	metrics.ProxyMutationReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(receiveSize))
D
dragondriver 已提交
2185

2186 2187 2188 2189 2190
	defer func() {
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.TotalLabel).Inc()
	}()

2191
	it := &insertTask{
2192 2193
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2194
		// req:       request,
2195 2196 2197 2198
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2199
			InsertRequest: internalpb.InsertRequest{
2200
				Base: &commonpb.MsgBase{
X
xige-16 已提交
2201 2202
					MsgType:  commonpb.MsgType_Insert,
					MsgID:    0,
X
Xiaofan 已提交
2203
					SourceID: Params.ProxyCfg.GetNodeID(),
2204 2205 2206
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2207 2208 2209
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2210
				// RowData: transfer column based request to this
2211 2212
			},
		},
2213 2214 2215 2216
		idAllocator:   node.idAllocator,
		segIDAssigner: node.segAssigner,
		chMgr:         node.chMgr,
		chTicker:      node.chTicker,
2217
	}
2218 2219

	if len(it.PartitionName) <= 0 {
2220
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2221 2222
	}

X
Xiangyu Wang 已提交
2223
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2224
		numRows := request.NumRows
2225 2226 2227 2228
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2229

X
Xiangyu Wang 已提交
2230 2231 2232 2233 2234 2235 2236
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2237 2238
	}

X
Xiangyu Wang 已提交
2239
	log.Debug("Enqueue insert request in Proxy",
2240
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2241 2242 2243 2244 2245
		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)),
2246 2247
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2248

X
Xiangyu Wang 已提交
2249 2250
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2251 2252
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
			metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2253
		return constructFailedResponse(err), nil
2254
	}
D
dragondriver 已提交
2255

X
Xiangyu Wang 已提交
2256
	log.Debug("Detail of insert request in Proxy",
2257
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2258
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2259 2260 2261 2262 2263
		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 已提交
2264 2265 2266 2267 2268
		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))
2269
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2270
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2271 2272 2273 2274 2275
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2276
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287
			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 已提交
2288
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2289

2290
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2291
		metrics.SuccessLabel).Inc()
2292 2293
	successCnt := it.result.InsertCnt - int64(len(it.result.ErrIndex))
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(successCnt))
2294
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.InsertLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2295 2296 2297
	return it.result, nil
}

2298
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2299
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2300 2301 2302
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2303 2304
	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))
2305

2306 2307 2308
	receiveSize := proto.Size(request)
	metrics.ProxyMutationReceiveBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(receiveSize))

G
groot 已提交
2309 2310 2311 2312 2313 2314
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2315 2316 2317
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

2318 2319
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
2320
	dt := &deleteTask{
X
xige-16 已提交
2321 2322 2323
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		deleteExpr: request.Expr,
G
godchen 已提交
2324
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2325 2326 2327
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2328 2329 2330 2331 2332
			DeleteRequest: internalpb.DeleteRequest{
				Base: &commonpb.MsgBase{
					MsgType: commonpb.MsgType_Delete,
					MsgID:   0,
				},
X
xige-16 已提交
2333
				DbName:         request.DbName,
G
godchen 已提交
2334 2335 2336
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2337 2338 2339 2340
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2341 2342
	}

2343
	log.Debug("Enqueue delete request in Proxy",
2344
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2345 2346 2347 2348
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2349 2350 2351 2352

	// MsgID will be set by Enqueue()
	if err := node.sched.dmQueue.Enqueue(dt); err != nil {
		log.Error("Failed to enqueue delete task: "+err.Error(), zap.String("traceID", traceID))
X
Xiaofan 已提交
2353
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2354
			metrics.FailLabel).Inc()
2355

G
groot 已提交
2356 2357 2358 2359 2360 2361 2362 2363
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2364
	log.Debug("Detail of delete request in Proxy",
2365
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2366 2367 2368 2369 2370
		zap.Int64("msgID", dt.Base.MsgID),
		zap.Uint64("timestamp", dt.Base.Timestamp),
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
2371 2372
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2373

2374 2375
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
X
Xiaofan 已提交
2376
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2377
			metrics.TotalLabel).Inc()
X
Xiaofan 已提交
2378
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2379
			metrics.FailLabel).Inc()
G
groot 已提交
2380 2381 2382 2383 2384 2385 2386 2387
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

X
Xiaofan 已提交
2388
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
2389
		metrics.SuccessLabel).Inc()
2390
	metrics.ProxyMutationLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), metrics.DeleteLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2391 2392 2393
	return dt.result, nil
}

2394
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2395
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2396 2397 2398 2399 2400
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2401 2402
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
2403 2404
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
2405

C
cai.zhang 已提交
2406 2407
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2408 2409
	traceID, _, _ := trace.InfoFromSpan(sp)

2410
	qt := &searchTask{
S
sunby 已提交
2411
		ctx:       ctx,
2412
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2413
		SearchRequest: &internalpb.SearchRequest{
2414
			Base: &commonpb.MsgBase{
2415
				MsgType:  commonpb.MsgType_Search,
X
Xiaofan 已提交
2416
				SourceID: Params.ProxyCfg.GetNodeID(),
2417
			},
2418
			ReqID: Params.ProxyCfg.GetNodeID(),
2419
		},
2420 2421 2422 2423
		request:  request,
		qc:       node.queryCoord,
		tr:       timerecord.NewTimeRecorder("search"),
		shardMgr: node.shardMgr,
2424 2425
	}

2426 2427 2428 2429 2430
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

	log.Debug(
		rpcReceived(method),
D
dragondriver 已提交
2431
		zap.String("traceID", traceID),
2432
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2433 2434 2435 2436 2437
		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)),
2438 2439 2440 2441
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2442

2443 2444 2445
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2446 2447
			zap.Error(err),
			zap.String("traceID", traceID),
2448
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2449 2450 2451 2452 2453 2454
			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),
2455 2456 2457
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2458

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

2462 2463
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2464
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2465 2466 2467 2468
				Reason:    err.Error(),
			},
		}, nil
	}
2469
	tr.Record("search request enqueue")
2470

2471 2472
	log.Debug(
		rpcEnqueued(method),
D
dragondriver 已提交
2473
		zap.String("traceID", traceID),
2474
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2475
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2476 2477 2478 2479 2480
		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),
2481
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2482 2483 2484 2485
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2486

2487 2488 2489
	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2490
			zap.Error(err),
D
dragondriver 已提交
2491
			zap.String("traceID", traceID),
2492
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2493
			zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2494 2495 2496 2497
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2498
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2499 2500 2501 2502
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2503

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

2507 2508
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2509
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2510 2511 2512 2513 2514
				Reason:    err.Error(),
			},
		}, nil
	}

2515 2516 2517
	span := tr.Record("wait search result")
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.SearchLabel).Observe(float64(span.Milliseconds()))
2518 2519
	log.Debug(
		rpcDone(method),
D
dragondriver 已提交
2520
		zap.String("traceID", traceID),
2521
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2522 2523 2524 2525 2526 2527
		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)),
2528 2529 2530 2531
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2532

2533 2534 2535
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.SuccessLabel).Inc()
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(qt.result.GetResults().GetNumQueries()))
C
cai.zhang 已提交
2536
	searchDur := tr.ElapseSpan().Milliseconds()
X
Xiaofan 已提交
2537
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
2538
		metrics.SearchLabel).Observe(float64(searchDur))
2539 2540 2541 2542 2543

	if qt.result != nil {
		sentSize := proto.Size(qt.result)
		metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
	}
2544 2545 2546
	return qt.result, nil
}

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

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

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

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

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

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

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

2593 2594
		resp.Status.Reason = err.Error()
		return resp, nil
2595 2596
	}

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

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

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

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

D
dragondriver 已提交
2626 2627 2628
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2629
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2630 2631 2632 2633 2634 2635
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))

X
Xiaofan 已提交
2636 2637
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2638
	return ft.result, nil
2639 2640
}

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

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

2654
	qt := &queryTask{
2655 2656 2657 2658 2659
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
X
Xiaofan 已提交
2660
				SourceID: Params.ProxyCfg.GetNodeID(),
2661
			},
2662
			ReqID: Params.ProxyCfg.GetNodeID(),
2663
		},
2664 2665
		request:          request,
		qc:               node.queryCoord,
2666
		queryShardPolicy: mergeRoundRobinPolicy,
2667
		shardMgr:         node.shardMgr,
2668 2669
	}

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

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

D
dragondriver 已提交
2675 2676 2677
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2678
		zap.String("role", typeutil.ProxyRole),
2679 2680
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2681 2682 2683 2684 2685
		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 已提交
2686

D
dragondriver 已提交
2687 2688 2689 2690 2691 2692
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
2693 2694 2695
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2696

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

2700 2701 2702 2703 2704 2705
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2706
	}
2707
	tr.Record("query request enqueue")
2708

D
dragondriver 已提交
2709 2710 2711
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2712
		zap.String("role", typeutil.ProxyRole),
2713
		zap.Int64("msgID", qt.ID()),
2714 2715
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
2716
		zap.Strings("partitions", request.PartitionNames))
D
dragondriver 已提交
2717 2718 2719 2720 2721 2722

	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2723
			zap.String("role", typeutil.ProxyRole),
2724
			zap.Int64("msgID", qt.ID()),
2725 2726 2727
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
2728

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

2732 2733 2734 2735 2736 2737 2738
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2739 2740 2741
	span := tr.Record("wait query result")
	metrics.ProxyWaitForSearchResultLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
		metrics.QueryLabel).Observe(float64(span.Milliseconds()))
D
dragondriver 已提交
2742 2743 2744 2745
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
2746
		zap.Int64("msgID", qt.ID()),
2747 2748 2749
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2750

2751 2752 2753 2754
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
		metrics.SuccessLabel).Inc()

	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10),
2755
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2756 2757

	ret := &milvuspb.QueryResults{
2758 2759
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
2760 2761 2762 2763
	}
	sentSize := proto.Size(qt.result)
	metrics.ProxyReadReqSendBytes.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Add(float64(sentSize))
	return ret, nil
2764
}
2765

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

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

Y
Yusup 已提交
2776 2777 2778 2779 2780 2781 2782
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
2783
	method := "CreateAlias"
2784
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
2785
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2786 2787 2788 2789 2790 2791 2792 2793 2794 2795 2796 2797 2798 2799 2800 2801 2802 2803 2804

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

	if err := node.sched.ddQueue.Enqueue(cat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

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

Y
Yusup 已提交
2807 2808 2809 2810 2811 2812
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2813 2814 2815
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2816
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2817 2818 2819 2820
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2821 2822
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2823 2824 2825 2826

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

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

D
dragondriver 已提交
2844 2845 2846 2847 2848 2849 2850 2851 2852 2853 2854
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
2855 2856
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2857 2858 2859
	return cat.result, nil
}

2860
// DropAlias alter the alias of collection.
Y
Yusup 已提交
2861 2862 2863 2864
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2865 2866 2867 2868 2869

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

Y
Yusup 已提交
2870 2871 2872 2873 2874 2875 2876
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
2877
	method := "DropAlias"
2878
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
2879
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2880 2881 2882 2883 2884 2885 2886 2887 2888 2889 2890 2891 2892 2893 2894 2895

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

	if err := node.sched.ddQueue.Enqueue(dat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias))
X
Xiaofan 已提交
2896
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2897

Y
Yusup 已提交
2898 2899 2900 2901 2902 2903
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2904 2905 2906
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2907
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2908 2909 2910 2911
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2912
		zap.String("alias", request.Alias))
D
dragondriver 已提交
2913 2914 2915 2916

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2917
			zap.Error(err),
D
dragondriver 已提交
2918
			zap.String("traceID", traceID),
2919
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2920 2921 2922 2923
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2924 2925
			zap.String("alias", request.Alias))

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

Y
Yusup 已提交
2928 2929 2930 2931 2932 2933
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2934 2935 2936 2937 2938 2939 2940 2941 2942 2943
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias))

X
Xiaofan 已提交
2944 2945
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2946 2947 2948
	return dat.result, nil
}

2949
// AlterAlias alter alias of collection.
Y
Yusup 已提交
2950 2951 2952 2953
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2954 2955 2956 2957 2958

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

Y
Yusup 已提交
2959 2960 2961 2962 2963 2964 2965
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
2966
	method := "AlterAlias"
2967
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
2968
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2969 2970 2971 2972 2973 2974 2975 2976 2977 2978 2979 2980 2981 2982 2983 2984 2985 2986

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

	if err := node.sched.ddQueue.Enqueue(aat); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
			zap.String("role", typeutil.ProxyRole),
			zap.String("db", request.DbName),
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
X
Xiaofan 已提交
2987
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2988

Y
Yusup 已提交
2989 2990 2991 2992 2993 2994
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2995 2996 2997
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2998
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2999 3000 3001 3002
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
3003 3004
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
3005 3006 3007 3008

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

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

Y
Yusup 已提交
3021 3022 3023 3024 3025 3026
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3027 3028 3029 3030 3031 3032 3033 3034 3035 3036 3037
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))

X
Xiaofan 已提交
3038 3039
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3040 3041 3042
	return aat.result, nil
}

3043
// CalcDistance calculates the distances between vectors.
3044
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3045 3046 3047 3048 3049
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3050

3051 3052 3053 3054
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3055 3056
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3057

3058 3059 3060 3061 3062
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3063 3064
		}

3065
		qt := &queryTask{
3066 3067 3068 3069 3070
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
X
Xiaofan 已提交
3071
					SourceID: Params.ProxyCfg.GetNodeID(),
3072
				},
3073
				ReqID: Params.ProxyCfg.GetNodeID(),
3074
			},
3075 3076 3077 3078
			request: queryRequest,
			qc:      node.queryCoord,
			ids:     ids.IdArray,

3079
			queryShardPolicy: mergeRoundRobinPolicy,
3080
			shardMgr:         node.shardMgr,
3081 3082
		}

G
groot 已提交
3083 3084 3085 3086 3087 3088
		items := []zapcore.Field{
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields),
		}

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

3093 3094 3095 3096 3097
			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3098
			}, err
3099
		}
3100

G
groot 已提交
3101
		log.Debug("CalcDistance queryTask enqueued", items...)
3102 3103 3104

		err = qt.WaitToFinish()
		if err != nil {
G
groot 已提交
3105
			log.Error("CalcDistance queryTask failed to WaitToFinish", append(items, zap.Error(err))...)
3106 3107 3108 3109 3110 3111

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3112
			}, err
3113
		}
3114

G
groot 已提交
3115
		log.Debug("CalcDistance queryTask Done", items...)
3116 3117

		return &milvuspb.QueryResults{
3118 3119
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3120 3121 3122
		}, nil
	}

G
groot 已提交
3123 3124 3125 3126
	// calcDistanceTask is not a standard task, no need to enqueue
	task := &calcDistanceTask{
		traceID:   traceID,
		queryFunc: query,
3127 3128
	}

G
groot 已提交
3129
	return task.Execute(ctx, request)
3130 3131
}

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

3137
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3138
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3139
	log.Debug("GetPersistentSegmentInfo",
3140
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3141 3142 3143
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3144
	resp := &milvuspb.GetPersistentSegmentInfoResponse{
X
XuanYang-cn 已提交
3145
		Status: &commonpb.Status{
3146
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
XuanYang-cn 已提交
3147 3148
		},
	}
3149 3150 3151 3152
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3153 3154
	method := "GetPersistentSegmentInfo"
	tr := timerecord.NewTimeRecorder(method)
X
Xiaofan 已提交
3155
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3156
		metrics.TotalLabel).Inc()
G
godchen 已提交
3157
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
X
XuanYang-cn 已提交
3158
	if err != nil {
3159
		resp.Status.Reason = fmt.Errorf("getSegmentsOfCollection, err:%w", err).Error()
X
XuanYang-cn 已提交
3160 3161
		return resp, nil
	}
3162
	infoResp, err := node.dataCoord.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
X
XuanYang-cn 已提交
3163
		Base: &commonpb.MsgBase{
3164
			MsgType:   commonpb.MsgType_SegmentInfo,
X
XuanYang-cn 已提交
3165 3166
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3167
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3168 3169 3170 3171
		},
		SegmentIDs: segments,
	})
	if err != nil {
3172
		log.Debug("GetPersistentSegmentInfo fail", zap.Error(err))
3173
		resp.Status.Reason = fmt.Errorf("dataCoord:GetSegmentInfo, err:%w", err).Error()
X
XuanYang-cn 已提交
3174 3175
		return resp, nil
	}
3176
	log.Debug("GetPersistentSegmentInfo ", zap.Int("len(infos)", len(infoResp.Infos)), zap.Any("status", infoResp.Status))
3177
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3178 3179 3180 3181 3182 3183
		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 已提交
3184
			SegmentID:    info.ID,
X
XuanYang-cn 已提交
3185 3186
			CollectionID: info.CollectionID,
			PartitionID:  info.PartitionID,
S
sunby 已提交
3187
			NumRows:      info.NumOfRows,
X
XuanYang-cn 已提交
3188 3189 3190
			State:        info.State,
		}
	}
X
Xiaofan 已提交
3191
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method,
3192
		metrics.SuccessLabel).Inc()
X
Xiaofan 已提交
3193
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
3194
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
X
XuanYang-cn 已提交
3195 3196 3197 3198
	resp.Infos = persistentInfos
	return resp, nil
}

J
jingkl 已提交
3199
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3200
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3201
	log.Debug("GetQuerySegmentInfo",
3202
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3203 3204 3205
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3206
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3207
		Status: &commonpb.Status{
3208
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3209 3210
		},
	}
3211 3212 3213 3214
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3215

3216 3217 3218 3219 3220
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3221
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
Z
zhenshan.cao 已提交
3222
		Base: &commonpb.MsgBase{
3223
			MsgType:   commonpb.MsgType_SegmentInfo,
Z
zhenshan.cao 已提交
3224 3225
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3226
			SourceID:  Params.ProxyCfg.GetNodeID(),
Z
zhenshan.cao 已提交
3227
		},
3228
		CollectionID: collID,
Z
zhenshan.cao 已提交
3229 3230
	})
	if err != nil {
3231
		log.Error("Failed to get segment info from QueryCoord",
3232
			zap.Error(err))
Z
zhenshan.cao 已提交
3233 3234 3235
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3236
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3237
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3238
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3239 3240 3241 3242 3243 3244 3245 3246 3247 3248 3249 3250 3251
		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 已提交
3252
			State:        info.SegmentState,
3253
			NodeIds:      info.NodeIds,
Z
zhenshan.cao 已提交
3254 3255
		}
	}
3256
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3257 3258 3259 3260
	resp.Infos = queryInfos
	return resp, nil
}

C
Cai Yudong 已提交
3261
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3262
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3263
		Base: &commonpb.MsgBase{
3264
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3265 3266
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3267
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3268 3269 3270 3271 3272 3273 3274
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3275
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3276 3277 3278
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3279
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3280
		Base: &commonpb.MsgBase{
3281
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3282 3283
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3284
			SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3285 3286 3287 3288 3289 3290 3291 3292
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3293
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3294 3295 3296 3297 3298
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3299
		showSegmentResponse, err := node.rootCoord.ShowSegments(ctx, &milvuspb.ShowSegmentsRequest{
X
XuanYang-cn 已提交
3300
			Base: &commonpb.MsgBase{
3301
				MsgType:   commonpb.MsgType_ShowSegments,
X
XuanYang-cn 已提交
3302 3303
				MsgID:     0,
				Timestamp: 0,
X
Xiaofan 已提交
3304
				SourceID:  Params.ProxyCfg.GetNodeID(),
X
XuanYang-cn 已提交
3305 3306 3307 3308 3309 3310 3311
			},
			CollectionID: collectionID,
			PartitionID:  partitionID,
		})
		if err != nil {
			return nil, err
		}
3312
		if showSegmentResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3313 3314 3315 3316 3317 3318
			return nil, errors.New(showSegmentResponse.Status.Reason)
		}
		ret = append(ret, showSegmentResponse.SegmentIDs...)
	}
	return ret, nil
}
3319

J
jingkl 已提交
3320
// Dummy handles dummy request
C
Cai Yudong 已提交
3321
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3322 3323 3324 3325 3326 3327 3328 3329 3330 3331 3332
	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
	}

3333 3334
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3335
		if err != nil {
3336
			log.Debug("Failed to parse dummy query request")
3337 3338 3339
			return failedResponse, nil
		}

3340
		request := &milvuspb.QueryRequest{
3341 3342 3343
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3344
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3345 3346
		}

3347
		_, err = node.Query(ctx, request)
3348
		if err != nil {
3349
			log.Debug("Failed to execute dummy query")
3350 3351
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3352 3353 3354 3355 3356 3357

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

3358 3359
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3360 3361
}

J
jingkl 已提交
3362
// RegisterLink registers a link
C
Cai Yudong 已提交
3363
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
G
godchen 已提交
3364
	code := node.stateCode.Load().(internalpb.StateCode)
D
dragondriver 已提交
3365
	log.Debug("RegisterLink",
3366
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3367
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3368

G
godchen 已提交
3369
	if code != internalpb.StateCode_Healthy {
3370 3371 3372
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3373
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3374
				Reason:    "proxy not healthy",
3375 3376 3377
			},
		}, nil
	}
X
Xiaofan 已提交
3378
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.GetNodeID(), 10)).Inc()
3379 3380 3381
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3382
			ErrorCode: commonpb.ErrorCode_Success,
3383
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3384 3385 3386
		},
	}, nil
}
3387

3388
// GetMetrics gets the metrics of proxy
3389 3390 3391
// 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 已提交
3392
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3393 3394 3395 3396
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
X
Xiaofan 已提交
3397
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3398
			zap.String("req", req.Request),
X
Xiaofan 已提交
3399
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.GetNodeID())))
3400 3401 3402 3403

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
X
Xiaofan 已提交
3404
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.GetNodeID()),
3405 3406 3407 3408 3409 3410 3411 3412
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
X
Xiaofan 已提交
3413
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3414 3415 3416 3417 3418 3419 3420 3421 3422 3423 3424 3425 3426 3427 3428
			zap.String("req", req.Request),
			zap.Error(err))

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

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

D
dragondriver 已提交
3429 3430 3431 3432 3433 3434 3435 3436 3437 3438
	msgID := UniqueID(0)
	msgID, err = node.idAllocator.AllocOne()
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to allocate id",
			zap.Error(err))
	}
	req.Base = &commonpb.MsgBase{
		MsgType:   commonpb.MsgType_SystemInfo,
		MsgID:     msgID,
		Timestamp: 0,
X
Xiaofan 已提交
3439
		SourceID:  Params.ProxyCfg.GetNodeID(),
D
dragondriver 已提交
3440 3441
	}

3442
	if metricType == metricsinfo.SystemInfoMetrics {
3443 3444 3445 3446 3447 3448 3449
		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))

3450
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3451 3452

		log.Debug("Proxy.GetMetrics",
X
Xiaofan 已提交
3453
			zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3454 3455 3456 3457 3458
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3459 3460
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3461
		return metrics, nil
3462 3463 3464
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
X
Xiaofan 已提交
3465
		zap.Int64("node_id", Params.ProxyCfg.GetNodeID()),
3466 3467 3468 3469 3470 3471 3472 3473 3474 3475 3476 3477
		zap.String("req", req.Request),
		zap.String("metric_type", metricType))

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

B
bigsheeper 已提交
3478 3479 3480
// 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 已提交
3481
		zap.Int64("proxy_id", Params.ProxyCfg.GetNodeID()),
B
bigsheeper 已提交
3482 3483 3484 3485 3486 3487 3488 3489 3490
		zap.Any("req", req))

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

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
3491 3492 3493 3494 3495 3496 3497

	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 已提交
3498 3499 3500 3501 3502
	infoResp, err := node.queryCoord.LoadBalance(ctx, &querypb.LoadBalanceRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_LoadBalanceSegments,
			MsgID:     0,
			Timestamp: 0,
X
Xiaofan 已提交
3503
			SourceID:  Params.ProxyCfg.GetNodeID(),
B
bigsheeper 已提交
3504 3505 3506
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3507
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3508
		SealedSegmentIDs: req.SealedSegmentIDs,
3509
		CollectionID:     collectionID,
B
bigsheeper 已提交
3510 3511 3512 3513 3514 3515 3516 3517 3518 3519 3520 3521 3522 3523 3524 3525 3526
	})
	if err != nil {
		log.Error("Failed to LoadBalance from Query Coordinator",
			zap.Any("req", req), zap.Error(err))
		status.Reason = err.Error()
		return status, nil
	}
	if infoResp.ErrorCode != commonpb.ErrorCode_Success {
		log.Error("Failed to LoadBalance from Query Coordinator", zap.String("errMsg", infoResp.Reason))
		status.Reason = infoResp.Reason
		return status, nil
	}
	log.Debug("LoadBalance Done", zap.Any("req", req), zap.Any("status", infoResp))
	status.ErrorCode = commonpb.ErrorCode_Success
	return status, nil
}

J
jingkl 已提交
3527
//GetCompactionState gets the compaction state of multiple segments
3528 3529 3530 3531 3532 3533 3534 3535 3536 3537 3538 3539 3540
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
}

3541
// ManualCompaction invokes compaction on specified collection
3542 3543 3544 3545 3546 3547 3548 3549 3550 3551 3552 3553 3554
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
}

3555
// GetCompactionStateWithPlans returns the compactions states with the given plan ID
3556 3557 3558 3559 3560 3561 3562 3563 3564 3565 3566 3567 3568
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 已提交
3569 3570 3571
// 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))
3572
	var err error
B
Bingyi Sun 已提交
3573 3574 3575 3576 3577 3578 3579
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3580
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3581 3582 3583 3584
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3585 3586 3587 3588
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3589 3590
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3591 3592 3593 3594
	code := node.stateCode.Load().(internalpb.StateCode)
	return code == internalpb.StateCode_Healthy
}

J
jingkl 已提交
3595
//unhealthyStatus returns the proxy not healthy status
3596 3597 3598
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3599
		Reason:    "proxy not healthy",
3600 3601
	}
}
G
groot 已提交
3602 3603 3604

// 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) {
3605 3606 3607
	log.Info("received import request",
		zap.String("collection name", req.GetCollectionName()),
		zap.Bool("row-based", req.GetRowBased()))
3608 3609 3610 3611 3612 3613
	resp := &milvuspb.ImportResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
	}
G
groot 已提交
3614 3615 3616 3617
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
3618 3619 3620 3621 3622 3623 3624 3625
	// Get collection ID and then channel names.
	collID, err := globalMetaCache.GetCollectionID(ctx, req.GetCollectionName())
	if err != nil {
		log.Error("collection ID not found",
			zap.String("collection name", req.GetCollectionName()),
			zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
3626
		return resp, nil
3627 3628 3629
	}
	chNames, err := node.chMgr.getVChannels(collID)
	if err != nil {
3630 3631 3632 3633 3634 3635 3636
		log.Error("failed to get virtual channels",
			zap.Error(err),
			zap.String("collection", req.GetCollectionName()),
			zap.Int64("collection_id", collID))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, nil
3637 3638
	}
	req.ChannelNames = chNames
3639 3640 3641
	if req.GetPartitionName() == "" {
		req.PartitionName = Params.CommonCfg.DefaultPartitionName
	}
3642
	// Call rootCoord to finish import.
3643 3644 3645 3646 3647 3648 3649 3650
	respFromRC, err := node.rootCoord.Import(ctx, req)
	if err != nil {
		log.Error("failed to execute bulk load request", zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	return respFromRC, nil
G
groot 已提交
3651 3652
}

G
groot 已提交
3653 3654 3655 3656 3657 3658 3659 3660 3661 3662 3663 3664 3665 3666 3667 3668 3669 3670 3671 3672 3673 3674 3675 3676 3677 3678 3679 3680
// GetImportState checks import task state from datanode
func (node *Proxy) GetImportState(ctx context.Context, req *milvuspb.GetImportStateRequest) (*milvuspb.GetImportStateResponse, error) {
	log.Info("received get import state request", zap.Int64("taskID", req.GetTask()))
	resp := &milvuspb.GetImportStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

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

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

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

X
XuanYang-cn 已提交
3681 3682 3683 3684 3685 3686 3687 3688 3689
// 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
	}

3690 3691
	req.Base = &commonpb.MsgBase{
		MsgType:  commonpb.MsgType_GetReplicas,
X
Xiaofan 已提交
3692
		SourceID: Params.ProxyCfg.GetNodeID(),
3693 3694
	}

X
XuanYang-cn 已提交
3695 3696 3697 3698 3699
	resp, err := node.queryCoord.GetReplicas(ctx, req)
	log.Info("received get replicas response", zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

3700 3701 3702 3703 3704 3705 3706 3707 3708 3709 3710 3711 3712 3713 3714 3715 3716 3717 3718 3719 3720 3721 3722 3723 3724 3725 3726 3727 3728
// 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))

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

	credInfo := &internalpb.CredentialInfo{
3729 3730
		Username:       request.Username,
		Sha256Password: request.Password,
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 3761 3762 3763 3764 3765 3766 3767 3768 3769 3770 3771 3772 3773 3774 3775 3776 3777 3778 3779 3780
	}
	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) {
	log.Debug("CreateCredential", zap.String("role", typeutil.RootCoordRole), zap.String("username", req.Username))
	// 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
	}
	credInfo := &internalpb.CredentialInfo{
		Username:          req.Username,
		EncryptedPassword: encryptedPassword,
3781
		Sha256Password:    crypto.SHA256(rawPassword, req.Username),
3782 3783 3784 3785 3786 3787 3788 3789 3790 3791 3792 3793
	}
	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 已提交
3794
func (node *Proxy) UpdateCredential(ctx context.Context, req *milvuspb.UpdateCredentialRequest) (*commonpb.Status, error) {
3795
	log.Debug("UpdateCredential", zap.String("role", typeutil.RootCoordRole), zap.String("username", req.Username))
C
codeman 已提交
3796 3797 3798 3799 3800 3801 3802 3803 3804
	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)
3805 3806 3807 3808 3809 3810 3811
	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 已提交
3812 3813
	// valid new password
	if err = ValidatePassword(rawNewPassword); err != nil {
3814 3815 3816 3817 3818 3819
		log.Error("illegal password", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_IllegalArgument,
			Reason:    err.Error(),
		}, nil
	}
C
codeman 已提交
3820 3821 3822 3823 3824 3825 3826 3827 3828
	// check old password is correct
	oldCredInfo, err := globalMetaCache.GetCredentialInfo(ctx, req.Username)
	if err != nil {
		log.Error("found no credential", zap.String("username", req.Username), zap.Error(err))
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "found no credential:" + req.Username,
		}, nil
	}
3829
	if !crypto.PasswordVerify(rawOldPassword, oldCredInfo) {
C
codeman 已提交
3830 3831 3832 3833 3834 3835 3836
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "old password is not correct:" + req.Username,
		}, nil
	}
	// update meta data
	encryptedPassword, err := crypto.PasswordEncrypt(rawNewPassword)
3837 3838 3839 3840 3841 3842 3843
	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 已提交
3844
	updateCredReq := &internalpb.CredentialInfo{
3845
		Username:          req.Username,
3846
		Sha256Password:    crypto.SHA256(rawNewPassword, req.Username),
3847 3848
		EncryptedPassword: encryptedPassword,
	}
C
codeman 已提交
3849
	result, err := node.rootCoord.UpdateCredential(ctx, updateCredReq)
3850 3851 3852 3853 3854 3855 3856 3857 3858 3859 3860 3861
	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) {
	log.Debug("DeleteCredential", zap.String("role", typeutil.RootCoordRole), zap.String("username", req.Username))
3862 3863 3864 3865 3866 3867
	if req.Username == util.UserRoot {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_DeleteCredentialFailure,
			Reason:    "user root cannot be deleted",
		}, nil
	}
3868 3869 3870 3871 3872 3873 3874 3875 3876 3877 3878 3879 3880
	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) {
	log.Debug("ListCredUsers", zap.String("role", typeutil.RootCoordRole))
3881 3882 3883 3884 3885 3886
	rootCoordReq := &milvuspb.ListCredUsersRequest{
		Base: &commonpb.MsgBase{
			MsgType: commonpb.MsgType_ListCredUsernames,
		},
	}
	resp, err := node.rootCoord.ListCredUsers(ctx, rootCoordReq)
3887 3888 3889 3890 3891 3892 3893 3894 3895 3896 3897 3898
	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,
		},
3899
		Usernames: resp.Usernames,
3900 3901
	}, nil
}
3902 3903 3904 3905 3906 3907 3908 3909 3910 3911 3912 3913 3914 3915 3916 3917

// SendSearchResult needs to be removed TODO
func (node *Proxy) SendSearchResult(ctx context.Context, req *internalpb.SearchResults) (*commonpb.Status, error) {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
		Reason:    "Not implemented",
	}, nil
}

// SendRetrieveResult needs to be removed TODO
func (node *Proxy) SendRetrieveResult(ctx context.Context, req *internalpb.RetrieveResults) (*commonpb.Status, error) {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
		Reason:    "Not implemented",
	}, nil
}
3918 3919 3920 3921 3922 3923 3924 3925 3926 3927 3928 3929 3930 3931 3932 3933 3934 3935 3936 3937 3938 3939 3940 3941 3942 3943 3944 3945 3946 3947 3948 3949 3950 3951 3952 3953 3954 3955 3956 3957 3958 3959 3960 3961 3962

func (node *Proxy) CreateRole(ctx context.Context, request *milvuspb.CreateRoleRequest) (*commonpb.Status, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) DropRole(ctx context.Context, request *milvuspb.DropRoleRequest) (*commonpb.Status, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) OperateUserRole(ctx context.Context, request *milvuspb.OperateUserRoleRequest) (*commonpb.Status, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) SelectRole(ctx context.Context, request *milvuspb.SelectRoleRequest) (*milvuspb.SelectRoleResponse, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) SelectUser(ctx context.Context, request *milvuspb.SelectUserRequest) (*milvuspb.SelectUserResponse, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) SelectResource(ctx context.Context, request *milvuspb.SelectResourceRequest) (*milvuspb.SelectResourceResponse, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) OperatePrivilege(ctx context.Context, request *milvuspb.OperatePrivilegeRequest) (*commonpb.Status, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) SelectGrant(ctx context.Context, request *milvuspb.SelectGrantRequest) (*milvuspb.SelectGrantResponse, error) {
	//TODO implement me
	panic("implement me")
}

func (node *Proxy) RefreshPolicyInfoCache(ctx context.Context, request *proxypb.RefreshPolicyInfoCacheRequest) (*commonpb.Status, error) {
	//TODO implement me
	panic("implement me")
}