impl.go 142.6 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"
S
sunby 已提交
29

30
	"github.com/milvus-io/milvus/internal/common"
X
Xiangyu Wang 已提交
31
	"github.com/milvus-io/milvus/internal/log"
32
	"github.com/milvus-io/milvus/internal/metrics"
J
jaime 已提交
33
	"github.com/milvus-io/milvus/internal/mq/msgstream"
X
Xiangyu Wang 已提交
34 35 36 37 38 39
	"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"
40
	"github.com/milvus-io/milvus/internal/proto/schemapb"
41
	"github.com/milvus-io/milvus/internal/util/crypto"
42
	"github.com/milvus-io/milvus/internal/util/distance"
43 44 45
	"github.com/milvus-io/milvus/internal/util/funcutil"
	"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 102
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to invalidate collection meta cache",
103
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
104 105 106
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

107
	collectionName := request.CollectionName
N
neza2017 已提交
108 109 110
	if globalMetaCache != nil {
		globalMetaCache.RemoveCollection(ctx, collectionName) // no need to return error, though collection may be not cached
	}
111
	logutil.Logger(ctx).Debug("complete to invalidate collection meta cache",
112
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
113 114 115
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

116
	return &commonpb.Status{
117
		ErrorCode: commonpb.ErrorCode_Success,
118 119
		Reason:    "",
	}, nil
120 121
}

122
// ReleaseDQLMessageStream release the query message stream of specific collection.
C
Cai Yudong 已提交
123
func (node *Proxy) ReleaseDQLMessageStream(ctx context.Context, request *proxypb.ReleaseDQLMessageStreamRequest) (*commonpb.Status, error) {
124 125
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to release DQL message strem",
126
		zap.Any("role", typeutil.ProxyRole),
127 128 129
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

130 131 132 133
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}

134 135
	_ = node.chMgr.removeDQLStream(request.CollectionID)

136
	logutil.Logger(ctx).Debug("complete to release DQL message stream",
137
		zap.Any("role", typeutil.ProxyRole),
138 139 140 141 142 143 144 145 146
		zap.Any("db", request.DbID),
		zap.Any("collection", request.CollectionID))

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

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreateCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
157 158 159 160
	method := "CreateCollection"
	tr := timerecord.NewTimeRecorder(method)

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

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

169 170 171
	// avoid data race
	lenOfSchema := len(request.Schema)

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

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

194
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
195
		return &commonpb.Status{
196
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
197 198 199 200
			Reason:    err.Error(),
		}, nil
	}

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

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

230
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
231
		return &commonpb.Status{
232
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
233 234 235 236
			Reason:    err.Error(),
		}, nil
	}

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

250 251
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
252 253 254
	return cct.result, nil
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
264 265 266
	method := "DropCollection"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
267

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

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

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

291
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
292
		return &commonpb.Status{
293
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
294 295 296 297
			Reason:    err.Error(),
		}, nil
	}

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

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

318
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
319
		return &commonpb.Status{
320
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
321 322 323 324
			Reason:    err.Error(),
		}, nil
	}

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

334 335
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
336 337 338
	return dct.result, nil
}

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

	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasCollection")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
350 351 352
	method := "HasCollection"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
353
		metrics.TotalLabel).Inc()
354 355 356

	log.Debug("HasCollection received",
		zap.String("traceID", traceID),
357
		zap.String("role", typeutil.ProxyRole),
358 359 360
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))

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

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

376
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
377
			metrics.AbandonLabel).Inc()
378 379
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
380
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
381 382 383 384 385
				Reason:    err.Error(),
			},
		}, nil
	}

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

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

406
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
407
			metrics.FailLabel).Inc()
408 409
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
410
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
411 412 413 414 415
				Reason:    err.Error(),
			},
		}, nil
	}

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

425
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
426 427
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
428 429 430
	return hct.result, nil
}

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

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

443
	lct := &loadCollectionTask{
S
sunby 已提交
444
		ctx:                   ctx,
445 446
		Condition:             NewTaskCondition(ctx),
		LoadCollectionRequest: request,
447
		queryCoord:            node.queryCoord,
448 449
	}

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

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

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

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

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

492
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
493
			metrics.TotalLabel).Inc()
494
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
495
			metrics.FailLabel).Inc()
496
		return &commonpb.Status{
497
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
498 499 500 501
			Reason:    err.Error(),
		}, nil
	}

502 503
	log.Debug("LoadCollection done",
		zap.String("traceID", traceID),
504
		zap.String("role", typeutil.ProxyRole),
505 506 507 508 509 510
		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))

511
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
512
		metrics.TotalLabel).Inc()
513
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
514 515
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
516
	return lct.result, nil
517 518
}

519
// ReleaseCollection remove the loaded collection from query nodes.
C
Cai Yudong 已提交
520
func (node *Proxy) ReleaseCollection(ctx context.Context, request *milvuspb.ReleaseCollectionRequest) (*commonpb.Status, error) {
521 522 523
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
524

525
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-ReleaseCollection")
526 527
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
528 529
	method := "ReleaseCollection"
	tr := timerecord.NewTimeRecorder(method)
530

531
	rct := &releaseCollectionTask{
S
sunby 已提交
532
		ctx:                      ctx,
533 534
		Condition:                NewTaskCondition(ctx),
		ReleaseCollectionRequest: request,
535
		queryCoord:               node.queryCoord,
536
		chMgr:                    node.chMgr,
537 538
	}

539 540
	log.Debug(
		rpcReceived(method),
541
		zap.String("traceID", traceID),
542
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
543 544
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
545 546

	if err := node.sched.ddQueue.Enqueue(rct); err != nil {
547 548
		log.Warn(
			rpcFailedToEnqueue(method),
549 550
			zap.Error(err),
			zap.String("traceID", traceID),
551
			zap.String("role", typeutil.ProxyRole),
552 553 554
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName))

555
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
556
			metrics.AbandonLabel).Inc()
557
		return &commonpb.Status{
558
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
559 560 561 562
			Reason:    err.Error(),
		}, nil
	}

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

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

585
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
586
			metrics.TotalLabel).Inc()
587
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
588
			metrics.FailLabel).Inc()
589
		return &commonpb.Status{
590
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
591 592 593 594
			Reason:    err.Error(),
		}, nil
	}

595 596
	log.Debug(
		rpcDone(method),
597
		zap.String("traceID", traceID),
598
		zap.String("role", typeutil.ProxyRole),
599 600 601 602 603 604
		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))

605
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
606
		metrics.TotalLabel).Inc()
607
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
608 609
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
610
	return rct.result, nil
611 612
}

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

621
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DescribeCollection")
622 623
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
624 625
	method := "DescribeCollection"
	tr := timerecord.NewTimeRecorder(method)
626

627
	dct := &describeCollectionTask{
S
sunby 已提交
628
		ctx:                       ctx,
629 630
		Condition:                 NewTaskCondition(ctx),
		DescribeCollectionRequest: request,
631
		rootCoord:                 node.rootCoord,
632 633
	}

634 635
	log.Debug("DescribeCollection received",
		zap.String("traceID", traceID),
636
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
637 638
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
639 640 641 642 643

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

648
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
649
			metrics.AbandonLabel).Inc()
650 651
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
652
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
653 654 655 656 657
				Reason:    err.Error(),
			},
		}, nil
	}

658 659
	log.Debug("DescribeCollection enqueued",
		zap.String("traceID", traceID),
660
		zap.String("role", typeutil.ProxyRole),
661 662 663
		zap.Int64("MsgID", dct.ID()),
		zap.Uint64("BeginTS", dct.BeginTs()),
		zap.Uint64("EndTS", dct.EndTs()),
D
dragondriver 已提交
664 665
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
666 667 668

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

678
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
679
			metrics.TotalLabel).Inc()
680
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
681
			metrics.FailLabel).Inc()
682

683 684
		return &milvuspb.DescribeCollectionResponse{
			Status: &commonpb.Status{
685
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
686 687 688 689 690
				Reason:    err.Error(),
			},
		}, nil
	}

691 692
	log.Debug("DescribeCollection done",
		zap.String("traceID", traceID),
693
		zap.String("role", typeutil.ProxyRole),
694 695 696 697 698 699
		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))

700
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
701
		metrics.TotalLabel).Inc()
702
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
703 704
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
705 706 707
	return dct.result, nil
}

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

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

722
	g := &getCollectionStatisticsTask{
G
godchen 已提交
723 724 725
		ctx:                            ctx,
		Condition:                      NewTaskCondition(ctx),
		GetCollectionStatisticsRequest: request,
726
		dataCoord:                      node.dataCoord,
727 728
	}

729 730
	log.Debug("GetCollectionStatistics received",
		zap.String("traceID", traceID),
731
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
732 733
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
734 735 736 737 738

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

743
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
744
			metrics.AbandonLabel).Inc()
745

G
godchen 已提交
746
		return &milvuspb.GetCollectionStatisticsResponse{
747
			Status: &commonpb.Status{
748
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
749 750 751 752 753
				Reason:    err.Error(),
			},
		}, nil
	}

754 755
	log.Debug("GetCollectionStatistics enqueued",
		zap.String("traceID", traceID),
756
		zap.String("role", typeutil.ProxyRole),
757 758 759
		zap.Int64("MsgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
D
dragondriver 已提交
760 761
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName))
762 763 764

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

774
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
775
			metrics.TotalLabel).Inc()
776
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
777
			metrics.FailLabel).Inc()
778

G
godchen 已提交
779
		return &milvuspb.GetCollectionStatisticsResponse{
780
			Status: &commonpb.Status{
781
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
782 783 784 785 786
				Reason:    err.Error(),
			},
		}, nil
	}

787 788
	log.Debug("GetCollectionStatistics done",
		zap.String("traceID", traceID),
789
		zap.String("role", typeutil.ProxyRole),
790 791 792 793 794 795
		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))

796
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
797
		metrics.TotalLabel).Inc()
798
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
799 800
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
801
	return g.result, nil
802 803
}

804
// ShowCollections list all collections in Milvus.
C
Cai Yudong 已提交
805
func (node *Proxy) ShowCollections(ctx context.Context, request *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
806 807 808 809 810
	if !node.checkHealthy() {
		return &milvuspb.ShowCollectionsResponse{
			Status: unhealthyStatus(),
		}, nil
	}
811 812 813 814
	method := "ShowCollections"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()

815
	sct := &showCollectionsTask{
G
godchen 已提交
816 817 818
		ctx:                    ctx,
		Condition:              NewTaskCondition(ctx),
		ShowCollectionsRequest: request,
819
		queryCoord:             node.queryCoord,
820
		rootCoord:              node.rootCoord,
821 822
	}

823
	log.Debug("ShowCollections received",
824
		zap.String("role", typeutil.ProxyRole),
825 826 827 828 829 830
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
		zap.Any("CollectionNames", request.CollectionNames),
	)

831
	err := node.sched.ddQueue.Enqueue(sct)
832
	if err != nil {
833 834
		log.Warn("ShowCollections failed to enqueue",
			zap.Error(err),
835
			zap.String("role", typeutil.ProxyRole),
836 837 838 839 840 841
			zap.String("DbName", request.DbName),
			zap.Uint64("TimeStamp", request.TimeStamp),
			zap.String("ShowType", request.Type.String()),
			zap.Any("CollectionNames", request.CollectionNames),
		)

842
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
G
godchen 已提交
843
		return &milvuspb.ShowCollectionsResponse{
844
			Status: &commonpb.Status{
845
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
846 847 848 849 850
				Reason:    err.Error(),
			},
		}, nil
	}

851
	log.Debug("ShowCollections enqueued",
852
		zap.String("role", typeutil.ProxyRole),
853
		zap.Int64("MsgID", sct.ID()),
854
		zap.String("DbName", sct.ShowCollectionsRequest.DbName),
855
		zap.Uint64("TimeStamp", request.TimeStamp),
856 857 858
		zap.String("ShowType", sct.ShowCollectionsRequest.Type.String()),
		zap.Any("CollectionNames", sct.ShowCollectionsRequest.CollectionNames),
	)
D
dragondriver 已提交
859

860 861
	err = sct.WaitToFinish()
	if err != nil {
862 863
		log.Warn("ShowCollections failed to WaitToFinish",
			zap.Error(err),
864
			zap.String("role", typeutil.ProxyRole),
865 866 867 868 869 870 871
			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),
		)

872 873
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

G
godchen 已提交
874
		return &milvuspb.ShowCollectionsResponse{
875
			Status: &commonpb.Status{
876
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
877 878 879 880 881
				Reason:    err.Error(),
			},
		}, nil
	}

882
	log.Debug("ShowCollections Done",
883
		zap.String("role", typeutil.ProxyRole),
884 885 886 887
		zap.Int64("MsgID", sct.ID()),
		zap.String("DbName", request.DbName),
		zap.Uint64("TimeStamp", request.TimeStamp),
		zap.String("ShowType", request.Type.String()),
888 889
		zap.Int("len(CollectionNames)", len(request.CollectionNames)),
		zap.Int("num_collections", len(sct.result.CollectionNames)))
890

891 892
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
893 894 895
	return sct.result, nil
}

896
// CreatePartition create a partition in specific collection.
C
Cai Yudong 已提交
897
func (node *Proxy) CreatePartition(ctx context.Context, request *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
898 899 900
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
901

902
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CreatePartition")
903 904
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
905 906 907
	method := "CreatePartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
908

909
	cpt := &createPartitionTask{
S
sunby 已提交
910
		ctx:                    ctx,
911 912
		Condition:              NewTaskCondition(ctx),
		CreatePartitionRequest: request,
913
		rootCoord:              node.rootCoord,
914 915 916
		result:                 nil,
	}

917 918 919
	log.Debug(
		rpcReceived("CreatePartition"),
		zap.String("traceID", traceID),
920
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
921 922 923
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
924 925 926 927 928 929

	if err := node.sched.ddQueue.Enqueue(cpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue("CreatePartition"),
			zap.Error(err),
			zap.String("traceID", traceID),
930
			zap.String("role", typeutil.ProxyRole),
931 932 933 934
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

935 936
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

937
		return &commonpb.Status{
938
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
939 940 941
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
942

943 944 945
	log.Debug(
		rpcEnqueued("CreatePartition"),
		zap.String("traceID", traceID),
946
		zap.String("role", typeutil.ProxyRole),
947 948 949
		zap.Int64("MsgID", cpt.ID()),
		zap.Uint64("BeginTS", cpt.BeginTs()),
		zap.Uint64("EndTS", cpt.EndTs()),
D
dragondriver 已提交
950 951 952
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
953 954 955 956

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

967 968
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

969
		return &commonpb.Status{
970
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
971 972 973
			Reason:    err.Error(),
		}, nil
	}
974 975 976 977

	log.Debug(
		rpcDone("CreatePartition"),
		zap.String("traceID", traceID),
978
		zap.String("role", typeutil.ProxyRole),
979 980 981 982 983 984 985
		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))

986 987
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
988 989 990
	return cpt.result, nil
}

991
// DropPartition drop a partition in specific collection.
C
Cai Yudong 已提交
992
func (node *Proxy) DropPartition(ctx context.Context, request *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
993 994 995
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
996

997
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-DropPartition")
998 999
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1000 1001 1002
	method := "DropPartition"
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
1003

1004
	dpt := &dropPartitionTask{
S
sunby 已提交
1005
		ctx:                  ctx,
1006 1007
		Condition:            NewTaskCondition(ctx),
		DropPartitionRequest: request,
1008
		rootCoord:            node.rootCoord,
1009 1010 1011
		result:               nil,
	}

1012 1013 1014
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1015
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1016 1017 1018
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1019 1020 1021 1022 1023 1024

	if err := node.sched.ddQueue.Enqueue(dpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1025
			zap.String("role", typeutil.ProxyRole),
1026 1027 1028 1029
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1030 1031
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

1032
		return &commonpb.Status{
1033
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1034 1035 1036
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1037

1038 1039 1040
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1041
		zap.String("role", typeutil.ProxyRole),
1042 1043 1044
		zap.Int64("MsgID", dpt.ID()),
		zap.Uint64("BeginTS", dpt.BeginTs()),
		zap.Uint64("EndTS", dpt.EndTs()),
D
dragondriver 已提交
1045 1046 1047
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1048 1049 1050 1051

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

1062 1063
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

1064
		return &commonpb.Status{
1065
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1066 1067 1068
			Reason:    err.Error(),
		}, nil
	}
1069 1070 1071 1072

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1073
		zap.String("role", typeutil.ProxyRole),
1074 1075 1076 1077 1078 1079 1080
		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))

1081 1082
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1083 1084 1085
	return dpt.result, nil
}

1086
// HasPartition check if partition exist.
C
Cai Yudong 已提交
1087
func (node *Proxy) HasPartition(ctx context.Context, request *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
1088 1089 1090 1091 1092
	if !node.checkHealthy() {
		return &milvuspb.BoolResponse{
			Status: unhealthyStatus(),
		}, nil
	}
D
dragondriver 已提交
1093

D
dragondriver 已提交
1094
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-HasPartition")
D
dragondriver 已提交
1095 1096
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1097 1098 1099 1100
	method := "HasPartition"
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1101
		metrics.TotalLabel).Inc()
D
dragondriver 已提交
1102

1103
	hpt := &hasPartitionTask{
S
sunby 已提交
1104
		ctx:                 ctx,
1105 1106
		Condition:           NewTaskCondition(ctx),
		HasPartitionRequest: request,
1107
		rootCoord:           node.rootCoord,
1108 1109 1110
		result:              nil,
	}

D
dragondriver 已提交
1111 1112 1113
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1114
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1115 1116 1117
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1118 1119 1120 1121 1122 1123

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

1129
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1130
			metrics.AbandonLabel).Inc()
1131

1132 1133
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1134
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1135 1136 1137 1138 1139
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1140

D
dragondriver 已提交
1141 1142 1143
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1144
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1145 1146 1147
		zap.Int64("MsgID", hpt.ID()),
		zap.Uint64("BeginTS", hpt.BeginTs()),
		zap.Uint64("EndTS", hpt.EndTs()),
D
dragondriver 已提交
1148 1149 1150
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
D
dragondriver 已提交
1151 1152 1153 1154

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

1165
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1166
			metrics.FailLabel).Inc()
1167

1168 1169
		return &milvuspb.BoolResponse{
			Status: &commonpb.Status{
1170
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1171 1172 1173 1174 1175
				Reason:    err.Error(),
			},
			Value: false,
		}, nil
	}
D
dragondriver 已提交
1176 1177 1178 1179

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1180
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1181 1182 1183 1184 1185 1186 1187
		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))

1188
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1189 1190
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1191 1192 1193
	return hpt.result, nil
}

1194
// LoadPartitions load specific partitions into query nodes.
C
Cai Yudong 已提交
1195
func (node *Proxy) LoadPartitions(ctx context.Context, request *milvuspb.LoadPartitionsRequest) (*commonpb.Status, error) {
1196 1197 1198
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1199

D
dragondriver 已提交
1200
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-LoadPartitions")
1201 1202
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
1203 1204
	method := "LoadPartitions"
	tr := timerecord.NewTimeRecorder(method)
1205

1206
	lpt := &loadPartitionsTask{
G
godchen 已提交
1207 1208 1209
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		LoadPartitionsRequest: request,
1210
		queryCoord:            node.queryCoord,
1211 1212
	}

1213 1214 1215
	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1216
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1217 1218 1219
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1220 1221 1222 1223 1224 1225

	if err := node.sched.ddQueue.Enqueue(lpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1226
			zap.String("role", typeutil.ProxyRole),
1227 1228 1229 1230
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1231
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1232
			metrics.AbandonLabel).Inc()
1233

1234
		return &commonpb.Status{
1235
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1236 1237 1238 1239
			Reason:    err.Error(),
		}, nil
	}

1240 1241 1242
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1243
		zap.String("role", typeutil.ProxyRole),
1244 1245 1246
		zap.Int64("MsgID", lpt.ID()),
		zap.Uint64("BeginTS", lpt.BeginTs()),
		zap.Uint64("EndTS", lpt.EndTs()),
D
dragondriver 已提交
1247 1248 1249
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1250 1251 1252 1253

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

1264
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1265
			metrics.TotalLabel).Inc()
1266
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1267
			metrics.FailLabel).Inc()
1268

1269
		return &commonpb.Status{
1270
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1271 1272 1273 1274
			Reason:    err.Error(),
		}, nil
	}

1275 1276 1277
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1278
		zap.String("role", typeutil.ProxyRole),
1279 1280 1281 1282 1283 1284 1285
		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))

1286
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1287
		metrics.TotalLabel).Inc()
1288
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1289 1290
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1291
	return lpt.result, nil
1292 1293
}

1294
// ReleasePartitions release specific partitions from query nodes.
C
Cai Yudong 已提交
1295
func (node *Proxy) ReleasePartitions(ctx context.Context, request *milvuspb.ReleasePartitionsRequest) (*commonpb.Status, error) {
1296 1297 1298
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
1299 1300 1301 1302 1303

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

1304
	rpt := &releasePartitionsTask{
G
godchen 已提交
1305 1306 1307
		ctx:                      ctx,
		Condition:                NewTaskCondition(ctx),
		ReleasePartitionsRequest: request,
1308
		queryCoord:               node.queryCoord,
1309 1310
	}

1311
	method := "ReleasePartitions"
1312
	tr := timerecord.NewTimeRecorder(method)
1313 1314 1315 1316

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1317
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1318 1319 1320
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1321 1322 1323 1324 1325 1326

	if err := node.sched.ddQueue.Enqueue(rpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1327
			zap.String("role", typeutil.ProxyRole),
1328 1329 1330 1331
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))

1332
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1333
			metrics.AbandonLabel).Inc()
1334

1335
		return &commonpb.Status{
1336
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1337 1338 1339 1340
			Reason:    err.Error(),
		}, nil
	}

1341 1342 1343
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1344
		zap.String("role", typeutil.ProxyRole),
1345 1346 1347
		zap.Int64("msgID", rpt.Base.MsgID),
		zap.Uint64("BeginTS", rpt.BeginTs()),
		zap.Uint64("EndTS", rpt.EndTs()),
D
dragondriver 已提交
1348 1349 1350
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
1351 1352 1353 1354

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

1365
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1366
			metrics.TotalLabel).Inc()
1367
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1368
			metrics.FailLabel).Inc()
1369

1370
		return &commonpb.Status{
1371
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1372 1373 1374 1375
			Reason:    err.Error(),
		}, nil
	}

1376 1377 1378
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1379
		zap.String("role", typeutil.ProxyRole),
1380 1381 1382 1383 1384 1385 1386
		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))

1387
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1388
		metrics.TotalLabel).Inc()
1389
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1390 1391
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1392
	return rpt.result, nil
1393 1394
}

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

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

1407
	g := &getPartitionStatisticsTask{
1408 1409 1410
		ctx:                           ctx,
		Condition:                     NewTaskCondition(ctx),
		GetPartitionStatisticsRequest: request,
1411
		dataCoord:                     node.dataCoord,
1412 1413
	}

1414
	method := "GetPartitionStatistics"
1415
	tr := timerecord.NewTimeRecorder(method)
1416 1417 1418 1419

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1420
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1421 1422 1423
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1424 1425 1426 1427 1428 1429

	if err := node.sched.ddQueue.Enqueue(g); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1430
			zap.String("role", typeutil.ProxyRole),
1431 1432 1433 1434
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1435
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1436
			metrics.AbandonLabel).Inc()
1437

1438 1439 1440 1441 1442 1443 1444 1445
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1446 1447 1448
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1449
		zap.String("role", typeutil.ProxyRole),
1450 1451 1452
		zap.Int64("msgID", g.ID()),
		zap.Uint64("BeginTS", g.BeginTs()),
		zap.Uint64("EndTS", g.EndTs()),
1453 1454 1455
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName))
1456 1457 1458 1459

	if err := g.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
1460
			zap.Error(err),
1461
			zap.String("traceID", traceID),
1462
			zap.String("role", typeutil.ProxyRole),
1463 1464 1465
			zap.Int64("msgID", g.ID()),
			zap.Uint64("BeginTS", g.BeginTs()),
			zap.Uint64("EndTS", g.EndTs()),
1466 1467 1468 1469
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("partition", request.PartitionName))

1470
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1471
			metrics.TotalLabel).Inc()
1472
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1473
			metrics.FailLabel).Inc()
1474

1475 1476 1477 1478 1479 1480 1481 1482
		return &milvuspb.GetPartitionStatisticsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1483 1484 1485
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1486
		zap.String("role", typeutil.ProxyRole),
1487 1488 1489 1490 1491 1492 1493
		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))

1494
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1495
		metrics.TotalLabel).Inc()
1496
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1497 1498
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1499
	return g.result, nil
1500 1501
}

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

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

1514
	spt := &showPartitionsTask{
G
godchen 已提交
1515 1516 1517
		ctx:                   ctx,
		Condition:             NewTaskCondition(ctx),
		ShowPartitionsRequest: request,
1518
		rootCoord:             node.rootCoord,
1519
		queryCoord:            node.queryCoord,
G
godchen 已提交
1520
		result:                nil,
1521 1522
	}

1523
	method := "ShowPartitions"
1524 1525 1526
	tr := timerecord.NewTimeRecorder(method)
	//TODO: use collectionID instead of collectionName
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1527
		metrics.TotalLabel).Inc()
1528 1529 1530 1531

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1532
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1533
		zap.Any("request", request))
1534 1535 1536 1537 1538 1539

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

1543
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1544
			metrics.AbandonLabel).Inc()
1545

G
godchen 已提交
1546
		return &milvuspb.ShowPartitionsResponse{
1547
			Status: &commonpb.Status{
1548
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1549 1550 1551 1552 1553
				Reason:    err.Error(),
			},
		}, nil
	}

1554 1555 1556
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1557
		zap.String("role", typeutil.ProxyRole),
1558 1559 1560
		zap.Int64("msgID", spt.ID()),
		zap.Uint64("BeginTS", spt.BeginTs()),
		zap.Uint64("EndTS", spt.EndTs()),
1561 1562
		zap.String("db", spt.ShowPartitionsRequest.DbName),
		zap.String("collection", spt.ShowPartitionsRequest.CollectionName),
1563 1564 1565 1566 1567
		zap.Any("partitions", spt.ShowPartitionsRequest.PartitionNames))

	if err := spt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
1568
			zap.Error(err),
1569
			zap.String("traceID", traceID),
1570
			zap.String("role", typeutil.ProxyRole),
1571 1572 1573 1574 1575 1576
			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 已提交
1577

1578
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1579
			metrics.FailLabel).Inc()
1580

G
godchen 已提交
1581
		return &milvuspb.ShowPartitionsResponse{
1582
			Status: &commonpb.Status{
1583
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1584 1585 1586 1587
				Reason:    err.Error(),
			},
		}, nil
	}
1588 1589 1590 1591

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1592
		zap.String("role", typeutil.ProxyRole),
1593 1594 1595 1596 1597 1598 1599
		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))

1600
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1601 1602
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1603 1604 1605
	return spt.result, nil
}

1606
// CreateIndex create index for collection.
C
Cai Yudong 已提交
1607
func (node *Proxy) CreateIndex(ctx context.Context, request *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
1608 1609 1610
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1611 1612 1613 1614 1615

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

1616
	cit := &createIndexTask{
S
sunby 已提交
1617
		ctx:                ctx,
1618 1619
		Condition:          NewTaskCondition(ctx),
		CreateIndexRequest: request,
1620
		rootCoord:          node.rootCoord,
1621 1622
	}

D
dragondriver 已提交
1623
	method := "CreateIndex"
1624
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1625 1626 1627 1628

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1629
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1630 1631 1632 1633
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1634 1635 1636 1637 1638 1639

	if err := node.sched.ddQueue.Enqueue(cit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1640
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1641 1642 1643 1644 1645
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.Any("extra_params", request.ExtraParams))

1646
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1647
			metrics.AbandonLabel).Inc()
1648

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

D
dragondriver 已提交
1655 1656 1657
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1658
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1659 1660 1661
		zap.Int64("MsgID", cit.ID()),
		zap.Uint64("BeginTs", cit.BeginTs()),
		zap.Uint64("EndTs", cit.EndTs()),
D
dragondriver 已提交
1662 1663 1664 1665
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.Any("extra_params", request.ExtraParams))
D
dragondriver 已提交
1666 1667 1668 1669

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

1681
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1682
			metrics.TotalLabel).Inc()
1683
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1684
			metrics.FailLabel).Inc()
1685

1686
		return &commonpb.Status{
1687
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
1688 1689 1690 1691
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
1692 1693 1694
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1695
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1696 1697 1698 1699 1700 1701 1702 1703
		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))

1704
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1705
		metrics.TotalLabel).Inc()
1706
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1707 1708
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1709 1710 1711
	return cit.result, nil
}

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

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

1724
	dit := &describeIndexTask{
S
sunby 已提交
1725
		ctx:                  ctx,
1726 1727
		Condition:            NewTaskCondition(ctx),
		DescribeIndexRequest: request,
1728
		rootCoord:            node.rootCoord,
1729 1730
	}

1731 1732 1733
	method := "DescribeIndex"
	// avoid data race
	indexName := request.IndexName
1734
	tr := timerecord.NewTimeRecorder(method)
1735 1736 1737 1738

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

1756
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1757
			metrics.AbandonLabel).Inc()
1758

1759 1760
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
1761
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
1762 1763 1764 1765 1766
				Reason:    err.Error(),
			},
		}, nil
	}

1767 1768 1769
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1770
		zap.String("role", typeutil.ProxyRole),
1771 1772 1773
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1774 1775 1776
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
1777 1778 1779 1780 1781
		zap.String("index name", indexName))

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

Z
zhenshan.cao 已提交
1793 1794 1795 1796
		errCode := commonpb.ErrorCode_UnexpectedError
		if dit.result != nil {
			errCode = dit.result.Status.GetErrorCode()
		}
1797
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1798
			metrics.TotalLabel).Inc()
1799
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1800
			metrics.FailLabel).Inc()
1801

1802 1803
		return &milvuspb.DescribeIndexResponse{
			Status: &commonpb.Status{
Z
zhenshan.cao 已提交
1804
				ErrorCode: errCode,
1805 1806 1807 1808 1809
				Reason:    err.Error(),
			},
		}, nil
	}

1810 1811 1812
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1813
		zap.String("role", typeutil.ProxyRole),
1814 1815 1816 1817 1818 1819 1820 1821
		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))

1822
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1823
		metrics.TotalLabel).Inc()
1824
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1825 1826
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
1827 1828 1829
	return dit.result, nil
}

1830
// DropIndex drop the index of collection.
C
Cai Yudong 已提交
1831
func (node *Proxy) DropIndex(ctx context.Context, request *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
1832 1833 1834
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
1835 1836 1837 1838 1839

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

1840
	dit := &dropIndexTask{
S
sunby 已提交
1841
		ctx:              ctx,
B
BossZou 已提交
1842 1843
		Condition:        NewTaskCondition(ctx),
		DropIndexRequest: request,
1844
		rootCoord:        node.rootCoord,
B
BossZou 已提交
1845
	}
G
godchen 已提交
1846

D
dragondriver 已提交
1847
	method := "DropIndex"
1848
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
1849 1850 1851 1852

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1853
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1854 1855 1856 1857 1858
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))

D
dragondriver 已提交
1859 1860 1861 1862 1863
	if err := node.sched.ddQueue.Enqueue(dit); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1864
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1865 1866 1867 1868
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1869
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1870
			metrics.AbandonLabel).Inc()
D
dragondriver 已提交
1871

B
BossZou 已提交
1872
		return &commonpb.Status{
1873
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1874 1875 1876
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1877

D
dragondriver 已提交
1878 1879 1880
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1881
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1882 1883 1884
		zap.Int64("MsgID", dit.ID()),
		zap.Uint64("BeginTs", dit.BeginTs()),
		zap.Uint64("EndTs", dit.EndTs()),
D
dragondriver 已提交
1885 1886 1887 1888
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
D
dragondriver 已提交
1889 1890 1891 1892

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

1904
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1905
			metrics.TotalLabel).Inc()
1906
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1907
			metrics.FailLabel).Inc()
1908

B
BossZou 已提交
1909
		return &commonpb.Status{
1910
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
B
BossZou 已提交
1911 1912 1913
			Reason:    err.Error(),
		}, nil
	}
D
dragondriver 已提交
1914 1915 1916 1917

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
1918
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
1919 1920 1921 1922 1923 1924 1925 1926
		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))

1927
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1928
		metrics.TotalLabel).Inc()
1929
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1930 1931
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
B
BossZou 已提交
1932 1933 1934
	return dit.result, nil
}

1935 1936
// 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 已提交
1937
func (node *Proxy) GetIndexBuildProgress(ctx context.Context, request *milvuspb.GetIndexBuildProgressRequest) (*milvuspb.GetIndexBuildProgressResponse, error) {
1938 1939 1940 1941 1942
	if !node.checkHealthy() {
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: unhealthyStatus(),
		}, nil
	}
1943 1944 1945 1946 1947

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

1948
	gibpt := &getIndexBuildProgressTask{
1949 1950 1951
		ctx:                          ctx,
		Condition:                    NewTaskCondition(ctx),
		GetIndexBuildProgressRequest: request,
1952 1953
		indexCoord:                   node.indexCoord,
		rootCoord:                    node.rootCoord,
1954
		dataCoord:                    node.dataCoord,
1955 1956
	}

1957
	method := "GetIndexBuildProgress"
1958
	tr := timerecord.NewTimeRecorder(method)
1959 1960 1961 1962

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
1963
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
1964 1965 1966 1967
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
1968 1969 1970 1971 1972 1973

	if err := node.sched.ddQueue.Enqueue(gibpt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
1974
			zap.String("role", typeutil.ProxyRole),
1975 1976 1977 1978
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
1979
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
1980
			metrics.AbandonLabel).Inc()
1981

1982 1983 1984 1985 1986 1987 1988 1989
		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

1990 1991 1992
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
1993
		zap.String("role", typeutil.ProxyRole),
1994 1995 1996
		zap.Int64("MsgID", gibpt.ID()),
		zap.Uint64("BeginTs", gibpt.BeginTs()),
		zap.Uint64("EndTs", gibpt.EndTs()),
1997 1998 1999 2000
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2001 2002 2003 2004

	if err := gibpt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
2005
			zap.Error(err),
2006
			zap.String("traceID", traceID),
2007
			zap.String("role", typeutil.ProxyRole),
2008 2009 2010
			zap.Int64("MsgID", gibpt.ID()),
			zap.Uint64("BeginTs", gibpt.BeginTs()),
			zap.Uint64("EndTs", gibpt.EndTs()),
2011 2012 2013 2014
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))
2015
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2016
			metrics.TotalLabel).Inc()
2017
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2018
			metrics.FailLabel).Inc()
2019 2020 2021 2022 2023 2024 2025 2026

		return &milvuspb.GetIndexBuildProgressResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2027 2028 2029 2030

	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2031
		zap.String("role", typeutil.ProxyRole),
2032 2033 2034 2035 2036 2037 2038 2039
		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))
2040

2041
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2042
		metrics.TotalLabel).Inc()
2043
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2044 2045
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2046
	return gibpt.result, nil
2047 2048
}

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

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

2061
	dipt := &getIndexStateTask{
G
godchen 已提交
2062 2063 2064
		ctx:                  ctx,
		Condition:            NewTaskCondition(ctx),
		GetIndexStateRequest: request,
2065 2066
		indexCoord:           node.indexCoord,
		rootCoord:            node.rootCoord,
2067 2068
	}

2069
	method := "GetIndexState"
2070
	tr := timerecord.NewTimeRecorder(method)
2071 2072 2073 2074

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2075
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2076 2077 2078 2079
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2080 2081 2082 2083 2084 2085

	if err := node.sched.ddQueue.Enqueue(dipt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2086
			zap.String("role", typeutil.ProxyRole),
2087 2088 2089 2090 2091
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.String("field", request.FieldName),
			zap.String("index name", request.IndexName))

2092
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2093
			metrics.AbandonLabel).Inc()
2094

G
godchen 已提交
2095
		return &milvuspb.GetIndexStateResponse{
2096
			Status: &commonpb.Status{
2097
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2098 2099 2100 2101 2102
				Reason:    err.Error(),
			},
		}, nil
	}

2103 2104 2105
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2106
		zap.String("role", typeutil.ProxyRole),
2107 2108 2109
		zap.Int64("MsgID", dipt.ID()),
		zap.Uint64("BeginTs", dipt.BeginTs()),
		zap.Uint64("EndTs", dipt.EndTs()),
D
dragondriver 已提交
2110 2111 2112 2113
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("field", request.FieldName),
		zap.String("index name", request.IndexName))
2114 2115 2116 2117

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

2129
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2130
			metrics.TotalLabel).Inc()
2131
		metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2132
			metrics.FailLabel).Inc()
2133

G
godchen 已提交
2134
		return &milvuspb.GetIndexStateResponse{
2135
			Status: &commonpb.Status{
2136
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2137 2138 2139 2140 2141
				Reason:    err.Error(),
			},
		}, nil
	}

2142 2143 2144
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2145
		zap.String("role", typeutil.ProxyRole),
2146 2147 2148 2149 2150 2151 2152 2153
		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))

2154
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2155
		metrics.TotalLabel).Inc()
2156
	metrics.ProxyDQLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2157 2158
		metrics.SuccessLabel).Inc()
	metrics.ProxyDQLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2159 2160 2161
	return dipt.result, nil
}

2162
// Insert insert records into collection.
C
Cai Yudong 已提交
2163
func (node *Proxy) Insert(ctx context.Context, request *milvuspb.InsertRequest) (*milvuspb.MutationResult, error) {
X
Xiangyu Wang 已提交
2164 2165 2166 2167 2168 2169
	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))

2170 2171 2172 2173 2174
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}
2175 2176
	method := "Insert"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2177

2178
	it := &insertTask{
2179 2180
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
X
xige-16 已提交
2181
		// req:       request,
2182 2183 2184 2185
		BaseInsertTask: BaseInsertTask{
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2186
			InsertRequest: internalpb.InsertRequest{
2187
				Base: &commonpb.MsgBase{
X
xige-16 已提交
2188 2189 2190
					MsgType:  commonpb.MsgType_Insert,
					MsgID:    0,
					SourceID: Params.ProxyCfg.ProxyID,
2191 2192 2193
				},
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
X
xige-16 已提交
2194 2195 2196
				FieldsData:     request.FieldsData,
				NumRows:        uint64(request.NumRows),
				Version:        internalpb.InsertDataVersion_ColumnBased,
2197
				// RowData: transfer column based request to this
2198 2199
			},
		},
2200
		rowIDAllocator: node.idAllocator,
2201
		segIDAssigner:  node.segAssigner,
2202
		chMgr:          node.chMgr,
2203
		chTicker:       node.chTicker,
2204
	}
2205 2206

	if len(it.PartitionName) <= 0 {
2207
		it.PartitionName = Params.CommonCfg.DefaultPartitionName
2208 2209
	}

X
Xiangyu Wang 已提交
2210
	constructFailedResponse := func(err error) *milvuspb.MutationResult {
X
xige-16 已提交
2211
		numRows := request.NumRows
2212 2213 2214 2215
		errIndex := make([]uint32, numRows)
		for i := uint32(0); i < numRows; i++ {
			errIndex[i] = i
		}
2216

X
Xiangyu Wang 已提交
2217 2218 2219 2220 2221 2222 2223
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
			ErrIndex: errIndex,
		}
2224 2225
	}

X
Xiangyu Wang 已提交
2226
	log.Debug("Enqueue insert request in Proxy",
2227
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2228 2229 2230 2231 2232
		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)),
2233 2234
		zap.Uint32("NumRows", request.NumRows),
		zap.String("traceID", traceID))
D
dragondriver 已提交
2235

X
Xiangyu Wang 已提交
2236 2237
	if err := node.sched.dmQueue.Enqueue(it); err != nil {
		log.Debug("Failed to enqueue insert task: " + err.Error())
2238
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), metrics.AbandonLabel).Inc()
X
Xiangyu Wang 已提交
2239
		return constructFailedResponse(err), nil
2240
	}
D
dragondriver 已提交
2241

X
Xiangyu Wang 已提交
2242
	log.Debug("Detail of insert request in Proxy",
2243
		zap.String("role", typeutil.ProxyRole),
X
Xiangyu Wang 已提交
2244
		zap.Int64("msgID", it.Base.MsgID),
D
dragondriver 已提交
2245 2246 2247 2248 2249
		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 已提交
2250 2251 2252 2253 2254
		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))
2255
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2256
			metrics.TotalLabel).Inc()
2257
		metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2258
			metrics.FailLabel).Inc()
X
Xiangyu Wang 已提交
2259 2260 2261 2262 2263
		return constructFailedResponse(err), nil
	}

	if it.result.Status.ErrorCode != commonpb.ErrorCode_Success {
		setErrorIndex := func() {
X
xige-16 已提交
2264
			numRows := request.NumRows
X
Xiangyu Wang 已提交
2265 2266 2267 2268 2269 2270 2271 2272 2273 2274 2275
			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 已提交
2276
	it.result.InsertCnt = int64(request.NumRows)
D
dragondriver 已提交
2277

2278
	metrics.ProxyInsertCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2279 2280 2281
		metrics.SuccessLabel).Inc()
	metrics.ProxyInsertVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Add(float64(it.result.InsertCnt))
	metrics.ProxyInsertLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Observe(float64(tr.ElapseSpan().Milliseconds()))
2282 2283 2284
	return it.result, nil
}

2285
// Delete delete records from collection, then these records cannot be searched.
G
groot 已提交
2286
func (node *Proxy) Delete(ctx context.Context, request *milvuspb.DeleteRequest) (*milvuspb.MutationResult, error) {
2287 2288 2289
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Delete")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2290 2291
	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))
2292

G
groot 已提交
2293 2294 2295 2296 2297 2298
	if !node.checkHealthy() {
		return &milvuspb.MutationResult{
			Status: unhealthyStatus(),
		}, nil
	}

2299 2300 2301
	method := "Delete"
	tr := timerecord.NewTimeRecorder(method)

2302
	dt := &deleteTask{
X
xige-16 已提交
2303 2304 2305
		ctx:        ctx,
		Condition:  NewTaskCondition(ctx),
		deleteExpr: request.Expr,
G
godchen 已提交
2306
		BaseDeleteTask: BaseDeleteTask{
G
godchen 已提交
2307 2308 2309
			BaseMsg: msgstream.BaseMsg{
				HashValues: request.HashKeys,
			},
G
godchen 已提交
2310 2311 2312 2313 2314
			DeleteRequest: internalpb.DeleteRequest{
				Base: &commonpb.MsgBase{
					MsgType: commonpb.MsgType_Delete,
					MsgID:   0,
				},
X
xige-16 已提交
2315
				DbName:         request.DbName,
G
godchen 已提交
2316 2317 2318
				CollectionName: request.CollectionName,
				PartitionName:  request.PartitionName,
				// RowData: transfer column based request to this
C
Cai Yudong 已提交
2319 2320 2321 2322
			},
		},
		chMgr:    node.chMgr,
		chTicker: node.chTicker,
G
groot 已提交
2323 2324
	}

2325
	log.Debug("Enqueue delete request in Proxy",
2326
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2327 2328 2329 2330
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.String("partition", request.PartitionName),
		zap.String("expr", request.Expr))
2331 2332 2333 2334

	// 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))
2335
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2336
			metrics.FailLabel).Inc()
2337

G
groot 已提交
2338 2339 2340 2341 2342 2343 2344 2345
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2346
	log.Debug("Detail of delete request in Proxy",
2347
		zap.String("role", typeutil.ProxyRole),
G
groot 已提交
2348 2349 2350 2351 2352
		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),
2353 2354
		zap.String("expr", request.Expr),
		zap.String("traceID", traceID))
G
groot 已提交
2355

2356 2357
	if err := dt.WaitToFinish(); err != nil {
		log.Error("Failed to execute delete task in task scheduler: "+err.Error(), zap.String("traceID", traceID))
2358
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2359
			metrics.TotalLabel).Inc()
2360
		metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2361
			metrics.FailLabel).Inc()
G
groot 已提交
2362 2363 2364 2365 2366 2367 2368 2369
		return &milvuspb.MutationResult{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}

2370
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2371
		metrics.TotalLabel).Inc()
2372
	metrics.ProxyDMLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method,
2373 2374
		metrics.SuccessLabel).Inc()
	metrics.ProxyDMLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
G
groot 已提交
2375 2376 2377
	return dt.result, nil
}

2378
// Search search the most similar records of requests.
C
Cai Yudong 已提交
2379
func (node *Proxy) Search(ctx context.Context, request *milvuspb.SearchRequest) (*milvuspb.SearchResults, error) {
2380 2381 2382 2383 2384
	if !node.checkHealthy() {
		return &milvuspb.SearchResults{
			Status: unhealthyStatus(),
		}, nil
	}
2385 2386
	method := "Search"
	tr := timerecord.NewTimeRecorder(method)
D
dragondriver 已提交
2387

C
cai.zhang 已提交
2388 2389
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Search")
	defer sp.Finish()
D
dragondriver 已提交
2390 2391
	traceID, _, _ := trace.InfoFromSpan(sp)

2392
	qt := &searchTask{
S
sunby 已提交
2393
		ctx:       ctx,
2394
		Condition: NewTaskCondition(ctx),
G
godchen 已提交
2395
		SearchRequest: &internalpb.SearchRequest{
2396
			Base: &commonpb.MsgBase{
2397
				MsgType:  commonpb.MsgType_Search,
2398
				SourceID: Params.ProxyCfg.ProxyID,
2399
			},
2400
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2401
		},
2402
		resultBuf: make(chan []*internalpb.SearchResults, 1),
2403 2404
		query:     request,
		chMgr:     node.chMgr,
2405
		qc:        node.queryCoord,
2406
		tr:        timerecord.NewTimeRecorder("search"),
2407 2408
	}

2409 2410 2411 2412 2413
	travelTs := request.TravelTimestamp
	guaranteeTs := request.GuaranteeTimestamp

	log.Debug(
		rpcReceived(method),
D
dragondriver 已提交
2414
		zap.String("traceID", traceID),
2415
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2416 2417 2418 2419 2420
		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)),
2421 2422 2423 2424
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2425

2426 2427 2428
	if err := node.sched.dqQueue.Enqueue(qt); err != nil {
		log.Warn(
			rpcFailedToEnqueue(method),
D
dragondriver 已提交
2429 2430
			zap.Error(err),
			zap.String("traceID", traceID),
2431
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2432 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)),
			zap.Any("OutputFields", request.OutputFields),
2438 2439 2440
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2441

2442
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2443
			metrics.SearchLabel, metrics.AbandonLabel).Inc()
2444

2445 2446
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2447
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2448 2449 2450 2451 2452
				Reason:    err.Error(),
			},
		}, nil
	}

2453 2454
	log.Debug(
		rpcEnqueued(method),
D
dragondriver 已提交
2455
		zap.String("traceID", traceID),
2456
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2457
		zap.Int64("msgID", qt.ID()),
D
dragondriver 已提交
2458 2459 2460 2461 2462
		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),
2463
		zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2464 2465 2466 2467
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2468

2469 2470 2471
	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2472
			zap.Error(err),
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
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames),
			zap.Any("dsl", request.Dsl),
2480
			zap.Any("len(PlaceholderGroup)", len(request.PlaceholderGroup)),
2481 2482 2483 2484
			zap.Any("OutputFields", request.OutputFields),
			zap.Any("search_params", request.SearchParams),
			zap.Uint64("travel_timestamp", travelTs),
			zap.Uint64("guarantee_timestamp", guaranteeTs))
2485
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2486
			metrics.SearchLabel, metrics.TotalLabel).Inc()
2487

2488 2489
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
			metrics.SearchLabel, metrics.FailLabel).Inc()
2490

2491 2492
		return &milvuspb.SearchResults{
			Status: &commonpb.Status{
2493
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
2494 2495 2496 2497 2498
				Reason:    err.Error(),
			},
		}, nil
	}

2499 2500
	log.Debug(
		rpcDone(method),
D
dragondriver 已提交
2501
		zap.String("traceID", traceID),
2502
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2503 2504 2505 2506 2507 2508
		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)),
2509 2510 2511 2512
		zap.Any("OutputFields", request.OutputFields),
		zap.Any("search_params", request.SearchParams),
		zap.Uint64("travel_timestamp", travelTs),
		zap.Uint64("guarantee_timestamp", guaranteeTs))
D
dragondriver 已提交
2513

2514
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2515
		metrics.SearchLabel, metrics.TotalLabel).Inc()
2516
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2517
		metrics.SearchLabel, metrics.SuccessLabel).Inc()
2518
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2519
		metrics.SearchLabel).Set(float64(qt.result.Results.NumQueries))
C
cai.zhang 已提交
2520
	searchDur := tr.ElapseSpan().Milliseconds()
2521
	metrics.ProxySearchLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2522 2523
		metrics.SearchLabel).Observe(float64(searchDur))
	metrics.ProxySearchLatencyPerNQ.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Observe(float64(searchDur) / float64(qt.result.Results.NumQueries))
2524 2525 2526
	return qt.result, nil
}

2527
// Flush notify data nodes to persist the data of collection.
2528 2529 2530 2531 2532 2533 2534
func (node *Proxy) Flush(ctx context.Context, request *milvuspb.FlushRequest) (*milvuspb.FlushResponse, error) {
	resp := &milvuspb.FlushResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    "",
		},
	}
2535
	if !node.checkHealthy() {
2536 2537
		resp.Status.Reason = "proxy is not healthy"
		return resp, nil
2538
	}
D
dragondriver 已提交
2539 2540 2541 2542 2543

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

2544
	ft := &flushTask{
T
ThreadDao 已提交
2545 2546 2547
		ctx:          ctx,
		Condition:    NewTaskCondition(ctx),
		FlushRequest: request,
2548
		dataCoord:    node.dataCoord,
2549 2550
	}

D
dragondriver 已提交
2551
	method := "Flush"
2552 2553
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2554 2555 2556 2557

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2558
		zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
2559 2560
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2561 2562 2563 2564 2565 2566

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

2571 2572
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

2573 2574
		resp.Status.Reason = err.Error()
		return resp, nil
2575 2576
	}

D
dragondriver 已提交
2577 2578 2579
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2580
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2581 2582 2583
		zap.Int64("MsgID", ft.ID()),
		zap.Uint64("BeginTs", ft.BeginTs()),
		zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2584 2585
		zap.String("db", request.DbName),
		zap.Any("collections", request.CollectionNames))
D
dragondriver 已提交
2586 2587 2588 2589

	if err := ft.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
D
dragondriver 已提交
2590
			zap.Error(err),
D
dragondriver 已提交
2591
			zap.String("traceID", traceID),
2592
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2593 2594 2595
			zap.Int64("MsgID", ft.ID()),
			zap.Uint64("BeginTs", ft.BeginTs()),
			zap.Uint64("EndTs", ft.EndTs()),
D
dragondriver 已提交
2596 2597 2598
			zap.String("db", request.DbName),
			zap.Any("collections", request.CollectionNames))

2599 2600
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

D
dragondriver 已提交
2601
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
2602 2603
		resp.Status.Reason = err.Error()
		return resp, nil
2604 2605
	}

D
dragondriver 已提交
2606 2607 2608
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
2609
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2610 2611 2612 2613 2614 2615
		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))

2616 2617
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
2618
	return ft.result, nil
2619 2620
}

2621
// Query get the records by primary keys.
C
Cai Yudong 已提交
2622
func (node *Proxy) Query(ctx context.Context, request *milvuspb.QueryRequest) (*milvuspb.QueryResults, error) {
2623 2624 2625 2626 2627
	if !node.checkHealthy() {
		return &milvuspb.QueryResults{
			Status: unhealthyStatus(),
		}, nil
	}
2628

D
dragondriver 已提交
2629 2630 2631
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-Query")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)
2632
	tr := timerecord.NewTimeRecorder("Query")
D
dragondriver 已提交
2633

2634
	qt := &queryTask{
2635 2636 2637 2638 2639
		ctx:       ctx,
		Condition: NewTaskCondition(ctx),
		RetrieveRequest: &internalpb.RetrieveRequest{
			Base: &commonpb.MsgBase{
				MsgType:  commonpb.MsgType_Retrieve,
2640
				SourceID: Params.ProxyCfg.ProxyID,
2641
			},
2642
			ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2643 2644
		},
		resultBuf: make(chan []*internalpb.RetrieveResults),
2645
		query:     request,
2646 2647
		chMgr:     node.chMgr,
		qc:        node.queryCoord,
2648 2649
	}

D
dragondriver 已提交
2650 2651 2652 2653 2654
	method := "Query"

	log.Debug(
		rpcReceived(method),
		zap.String("traceID", traceID),
2655
		zap.String("role", typeutil.ProxyRole),
2656 2657 2658
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
G
godchen 已提交
2659

D
dragondriver 已提交
2660 2661 2662 2663 2664 2665
	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),
2666 2667 2668
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2669

2670
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2671
			metrics.QueryLabel, metrics.FailLabel).Inc()
2672 2673 2674 2675 2676 2677
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
2678 2679
	}

D
dragondriver 已提交
2680 2681 2682
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2683
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2684 2685 2686
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
2687 2688 2689
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2690 2691 2692 2693 2694 2695

	if err := qt.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
			zap.Error(err),
			zap.String("traceID", traceID),
2696
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2697 2698 2699
			zap.Int64("MsgID", qt.ID()),
			zap.Uint64("BeginTs", qt.BeginTs()),
			zap.Uint64("EndTs", qt.EndTs()),
2700 2701 2702
			zap.String("db", request.DbName),
			zap.String("collection", request.CollectionName),
			zap.Any("partitions", request.PartitionNames))
2703

2704
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2705
			metrics.QueryLabel, metrics.TotalLabel).Inc()
2706
		metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2707
			metrics.QueryLabel, metrics.FailLabel).Inc()
2708

2709 2710 2711 2712 2713 2714 2715
		return &milvuspb.QueryResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
	}
2716

D
dragondriver 已提交
2717 2718 2719 2720 2721 2722 2723
	log.Debug(
		rpcDone(method),
		zap.String("traceID", traceID),
		zap.String("role", typeutil.ProxyRole),
		zap.Int64("MsgID", qt.ID()),
		zap.Uint64("BeginTs", qt.BeginTs()),
		zap.Uint64("EndTs", qt.EndTs()),
2724 2725 2726
		zap.String("db", request.DbName),
		zap.String("collection", request.CollectionName),
		zap.Any("partitions", request.PartitionNames))
D
dragondriver 已提交
2727

2728
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2729
		metrics.QueryLabel, metrics.TotalLabel).Inc()
2730
	metrics.ProxySearchCount.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2731
		metrics.QueryLabel, metrics.SuccessLabel).Inc()
2732
	metrics.ProxySearchVectors.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2733
		metrics.QueryLabel).Set(float64(len(qt.result.FieldsData)))
2734
	metrics.ProxySendMessageLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
2735
		metrics.QueryLabel).Observe(float64(tr.ElapseSpan().Milliseconds()))
2736 2737 2738 2739 2740
	return &milvuspb.QueryResults{
		Status:     qt.result.Status,
		FieldsData: qt.result.FieldsData,
	}, nil
}
2741

2742
// CreateAlias create alias for collection, then you can search the collection with alias.
Y
Yusup 已提交
2743 2744 2745 2746
func (node *Proxy) CreateAlias(ctx context.Context, request *milvuspb.CreateAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2747 2748 2749 2750 2751

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

Y
Yusup 已提交
2752 2753 2754 2755 2756 2757 2758
	cat := &CreateAliasTask{
		ctx:                ctx,
		Condition:          NewTaskCondition(ctx),
		CreateAliasRequest: request,
		rootCoord:          node.rootCoord,
	}

D
dragondriver 已提交
2759
	method := "CreateAlias"
2760 2761
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2762 2763 2764 2765 2766 2767 2768 2769 2770 2771 2772 2773 2774 2775 2776 2777 2778 2779 2780

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

2781 2782
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()

Y
Yusup 已提交
2783 2784 2785 2786 2787 2788
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2789 2790 2791
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2792
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2793 2794 2795 2796
		zap.Int64("MsgID", cat.ID()),
		zap.Uint64("BeginTs", cat.BeginTs()),
		zap.Uint64("EndTs", cat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2797 2798
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2799 2800 2801 2802

	if err := cat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2803
			zap.Error(err),
D
dragondriver 已提交
2804
			zap.String("traceID", traceID),
2805
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2806 2807 2808 2809
			zap.Int64("MsgID", cat.ID()),
			zap.Uint64("BeginTs", cat.BeginTs()),
			zap.Uint64("EndTs", cat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2810 2811
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))
2812
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()
Y
Yusup 已提交
2813 2814 2815 2816 2817 2818 2819

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

D
dragondriver 已提交
2820 2821 2822 2823 2824 2825 2826 2827 2828 2829 2830
	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))

2831 2832
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2833 2834 2835
	return cat.result, nil
}

2836
// DropAlias alter the alias of collection.
Y
Yusup 已提交
2837 2838 2839 2840
func (node *Proxy) DropAlias(ctx context.Context, request *milvuspb.DropAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2841 2842 2843 2844 2845

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

Y
Yusup 已提交
2846 2847 2848 2849 2850 2851 2852
	dat := &DropAliasTask{
		ctx:              ctx,
		Condition:        NewTaskCondition(ctx),
		DropAliasRequest: request,
		rootCoord:        node.rootCoord,
	}

D
dragondriver 已提交
2853
	method := "DropAlias"
2854 2855
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2856 2857 2858 2859 2860 2861 2862 2863 2864 2865 2866 2867 2868 2869 2870 2871

	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))
2872
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2873

Y
Yusup 已提交
2874 2875 2876 2877 2878 2879
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2880 2881 2882
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2883
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2884 2885 2886 2887
		zap.Int64("MsgID", dat.ID()),
		zap.Uint64("BeginTs", dat.BeginTs()),
		zap.Uint64("EndTs", dat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2888
		zap.String("alias", request.Alias))
D
dragondriver 已提交
2889 2890 2891 2892

	if err := dat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2893
			zap.Error(err),
D
dragondriver 已提交
2894
			zap.String("traceID", traceID),
2895
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2896 2897 2898 2899
			zap.Int64("MsgID", dat.ID()),
			zap.Uint64("BeginTs", dat.BeginTs()),
			zap.Uint64("EndTs", dat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2900 2901
			zap.String("alias", request.Alias))

2902 2903
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

Y
Yusup 已提交
2904 2905 2906 2907 2908 2909
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2910 2911 2912 2913 2914 2915 2916 2917 2918 2919
	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))

2920 2921
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
2922 2923 2924
	return dat.result, nil
}

2925
// AlterAlias alter alias of collection.
Y
Yusup 已提交
2926 2927 2928 2929
func (node *Proxy) AlterAlias(ctx context.Context, request *milvuspb.AlterAliasRequest) (*commonpb.Status, error) {
	if !node.checkHealthy() {
		return unhealthyStatus(), nil
	}
D
dragondriver 已提交
2930 2931 2932 2933 2934

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

Y
Yusup 已提交
2935 2936 2937 2938 2939 2940 2941
	aat := &AlterAliasTask{
		ctx:               ctx,
		Condition:         NewTaskCondition(ctx),
		AlterAliasRequest: request,
		rootCoord:         node.rootCoord,
	}

D
dragondriver 已提交
2942
	method := "AlterAlias"
2943 2944
	tr := timerecord.NewTimeRecorder(method)
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.TotalLabel).Inc()
D
dragondriver 已提交
2945 2946 2947 2948 2949 2950 2951 2952 2953 2954 2955 2956 2957 2958 2959 2960 2961 2962

	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))
2963
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.AbandonLabel).Inc()
D
dragondriver 已提交
2964

Y
Yusup 已提交
2965 2966 2967 2968 2969 2970
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
2971 2972 2973
	log.Debug(
		rpcEnqueued(method),
		zap.String("traceID", traceID),
2974
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2975 2976 2977 2978
		zap.Int64("MsgID", aat.ID()),
		zap.Uint64("BeginTs", aat.BeginTs()),
		zap.Uint64("EndTs", aat.EndTs()),
		zap.String("db", request.DbName),
Y
Yusup 已提交
2979 2980
		zap.String("alias", request.Alias),
		zap.String("collection", request.CollectionName))
D
dragondriver 已提交
2981 2982 2983 2984

	if err := aat.WaitToFinish(); err != nil {
		log.Warn(
			rpcFailedToWaitToFinish(method),
Y
Yusup 已提交
2985
			zap.Error(err),
D
dragondriver 已提交
2986
			zap.String("traceID", traceID),
2987
			zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
2988 2989 2990 2991
			zap.Int64("MsgID", aat.ID()),
			zap.Uint64("BeginTs", aat.BeginTs()),
			zap.Uint64("EndTs", aat.EndTs()),
			zap.String("db", request.DbName),
Y
Yusup 已提交
2992 2993 2994
			zap.String("alias", request.Alias),
			zap.String("collection", request.CollectionName))

2995 2996
		metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.FailLabel).Inc()

Y
Yusup 已提交
2997 2998 2999 3000 3001 3002
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		}, nil
	}

D
dragondriver 已提交
3003 3004 3005 3006 3007 3008 3009 3010 3011 3012 3013
	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))

3014 3015
	metrics.ProxyDDLFunctionCall.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method, metrics.SuccessLabel).Inc()
	metrics.ProxyDDLReqLatency.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10), method).Observe(float64(tr.ElapseSpan().Milliseconds()))
Y
Yusup 已提交
3016 3017 3018
	return aat.result, nil
}

3019
// CalcDistance calculates the distances between vectors.
3020
func (node *Proxy) CalcDistance(ctx context.Context, request *milvuspb.CalcDistanceRequest) (*milvuspb.CalcDistanceResults, error) {
3021 3022 3023 3024 3025
	if !node.checkHealthy() {
		return &milvuspb.CalcDistanceResults{
			Status: unhealthyStatus(),
		}, nil
	}
3026
	param, _ := funcutil.GetAttrByKeyFromRepeatedKV("metric", request.GetParams())
3027 3028 3029 3030 3031 3032 3033 3034
	metric, err := distance.ValidateMetricType(param)
	if err != nil {
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
				Reason:    err.Error(),
			},
		}, nil
3035 3036
	}

3037 3038 3039 3040
	sp, ctx := trace.StartSpanFromContextWithOperationName(ctx, "Proxy-CalcDistance")
	defer sp.Finish()
	traceID, _, _ := trace.InfoFromSpan(sp)

3041 3042
	query := func(ids *milvuspb.VectorIDs) (*milvuspb.QueryResults, error) {
		outputFields := []string{ids.FieldName}
3043

3044 3045 3046 3047 3048
		queryRequest := &milvuspb.QueryRequest{
			DbName:         "",
			CollectionName: ids.CollectionName,
			PartitionNames: ids.PartitionNames,
			OutputFields:   outputFields,
3049 3050
		}

3051
		qt := &queryTask{
3052 3053 3054 3055 3056
			ctx:       ctx,
			Condition: NewTaskCondition(ctx),
			RetrieveRequest: &internalpb.RetrieveRequest{
				Base: &commonpb.MsgBase{
					MsgType:  commonpb.MsgType_Retrieve,
3057
					SourceID: Params.ProxyCfg.ProxyID,
3058
				},
3059
				ResultChannelID: strconv.FormatInt(Params.ProxyCfg.ProxyID, 10),
3060
			},
3061
			resultBuf: make(chan []*internalpb.RetrieveResults),
3062
			query:     queryRequest,
3063
			chMgr:     node.chMgr,
3064
			qc:        node.queryCoord,
Y
yukun 已提交
3065
			ids:       ids.IdArray,
3066 3067
		}

3068
		err := node.sched.dqQueue.Enqueue(qt)
3069
		if err != nil {
3070 3071 3072
			log.Debug("CalcDistance queryTask failed to enqueue",
				zap.Error(err),
				zap.String("traceID", traceID),
3073
				zap.String("role", typeutil.ProxyRole),
3074 3075 3076 3077
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames))

3078 3079 3080 3081 3082
			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3083
			}, err
3084
		}
3085 3086 3087

		log.Debug("CalcDistance queryTask enqueued",
			zap.String("traceID", traceID),
3088
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3089 3090 3091 3092 3093 3094
			zap.Int64("msgID", qt.Base.MsgID),
			zap.Uint64("timestamp", qt.Base.Timestamp),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields))
3095 3096 3097 3098

		err = qt.WaitToFinish()
		if err != nil {
			log.Debug("CalcDistance queryTask failed to WaitToFinish",
G
godchen 已提交
3099
				zap.Error(err),
3100
				zap.String("traceID", traceID),
3101
				zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3102 3103 3104 3105 3106 3107
				zap.Int64("msgID", qt.Base.MsgID),
				zap.Uint64("timestamp", qt.Base.Timestamp),
				zap.String("db", queryRequest.DbName),
				zap.String("collection", queryRequest.CollectionName),
				zap.Any("partitions", queryRequest.PartitionNames),
				zap.Any("OutputFields", queryRequest.OutputFields))
3108 3109 3110 3111 3112 3113

			return &milvuspb.QueryResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
3114
			}, err
3115
		}
3116 3117 3118

		log.Debug("CalcDistance queryTask Done",
			zap.String("traceID", traceID),
3119
			zap.String("role", typeutil.ProxyRole),
G
godchen 已提交
3120 3121 3122 3123 3124 3125
			zap.Int64("msgID", qt.Base.MsgID),
			zap.Uint64("timestamp", qt.Base.Timestamp),
			zap.String("db", queryRequest.DbName),
			zap.String("collection", queryRequest.CollectionName),
			zap.Any("partitions", queryRequest.PartitionNames),
			zap.Any("OutputFields", queryRequest.OutputFields))
3126 3127

		return &milvuspb.QueryResults{
3128 3129
			Status:     qt.result.Status,
			FieldsData: qt.result.FieldsData,
3130 3131 3132
		}, nil
	}

3133 3134 3135 3136 3137 3138 3139 3140 3141 3142 3143 3144 3145 3146
	// the vectors retrieved are random order, we need re-arrange the vectors by the order of input ids
	arrangeFunc := func(ids *milvuspb.VectorIDs, retrievedFields []*schemapb.FieldData) (*schemapb.VectorField, error) {
		var retrievedIds *schemapb.ScalarField
		var retrievedVectors *schemapb.VectorField
		for _, fieldData := range retrievedFields {
			if fieldData.FieldName == ids.FieldName {
				retrievedVectors = fieldData.GetVectors()
			}
			if fieldData.Type == schemapb.DataType_Int64 {
				retrievedIds = fieldData.GetScalars()
			}
		}

		if retrievedIds == nil || retrievedVectors == nil {
3147
			return nil, errors.New("failed to fetch vectors")
3148 3149 3150 3151 3152 3153 3154 3155 3156 3157 3158 3159 3160 3161 3162 3163
		}

		dict := make(map[int64]int)
		for index, id := range retrievedIds.GetLongData().Data {
			dict[id] = index
		}

		inputIds := ids.IdArray.GetIntId().Data
		if retrievedVectors.GetFloatVector() != nil {
			floatArr := retrievedVectors.GetFloatVector().Data
			element := retrievedVectors.GetDim()
			result := make([]float32, 0, int64(len(inputIds))*element)
			for _, id := range inputIds {
				index, ok := dict[id]
				if !ok {
					log.Error("id not found in CalcDistance", zap.Int64("id", id))
3164
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3165 3166 3167 3168 3169 3170 3171 3172 3173 3174 3175 3176 3177 3178 3179 3180 3181 3182 3183 3184 3185 3186 3187 3188 3189 3190
				}
				result = append(result, floatArr[int64(index)*element:int64(index+1)*element]...)
			}

			return &schemapb.VectorField{
				Dim: element,
				Data: &schemapb.VectorField_FloatVector{
					FloatVector: &schemapb.FloatArray{
						Data: result,
					},
				},
			}, nil
		}

		if retrievedVectors.GetBinaryVector() != nil {
			binaryArr := retrievedVectors.GetBinaryVector()
			element := retrievedVectors.GetDim()
			if element%8 != 0 {
				element = element + 8 - element%8
			}

			result := make([]byte, 0, int64(len(inputIds))*element)
			for _, id := range inputIds {
				index, ok := dict[id]
				if !ok {
					log.Error("id not found in CalcDistance", zap.Int64("id", id))
3191
					return nil, errors.New("failed to fetch vectors by id: " + fmt.Sprintln(id))
3192 3193 3194 3195 3196 3197 3198 3199 3200 3201 3202 3203
				}
				result = append(result, binaryArr[int64(index)*element:int64(index+1)*element]...)
			}

			return &schemapb.VectorField{
				Dim: element * 8,
				Data: &schemapb.VectorField_BinaryVector{
					BinaryVector: result,
				},
			}, nil
		}

3204
		return nil, errors.New("failed to fetch vectors")
3205 3206
	}

3207 3208
	log.Debug("CalcDistance received",
		zap.String("traceID", traceID),
3209
		zap.String("role", typeutil.ProxyRole),
3210
		zap.String("metric", metric))
G
godchen 已提交
3211

3212 3213 3214
	vectorsLeft := request.GetOpLeft().GetDataArray()
	opLeft := request.GetOpLeft().GetIdArray()
	if opLeft != nil {
3215 3216
		log.Debug("OpLeft IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3217
			zap.String("role", typeutil.ProxyRole))
3218

3219
		result, err := query(opLeft)
3220
		if err != nil {
3221 3222 3223
			log.Debug("Failed to get left vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3224
				zap.String("role", typeutil.ProxyRole))
3225

3226 3227 3228 3229 3230 3231 3232 3233
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3234 3235
		log.Debug("OpLeft IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3236
			zap.String("role", typeutil.ProxyRole))
3237

3238 3239
		vectorsLeft, err = arrangeFunc(opLeft, result.FieldsData)
		if err != nil {
3240 3241 3242
			log.Debug("Failed to re-arrange left vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3243
				zap.String("role", typeutil.ProxyRole))
3244

3245 3246 3247 3248 3249 3250
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3251
		}
3252 3253 3254

		log.Debug("Re-arrange left vectors done",
			zap.String("traceID", traceID),
3255
			zap.String("role", typeutil.ProxyRole))
3256 3257
	}

G
groot 已提交
3258
	if vectorsLeft == nil {
3259 3260 3261
		msg := "Left vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3262
			zap.String("role", typeutil.ProxyRole))
3263

G
groot 已提交
3264 3265 3266
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3267
				Reason:    msg,
G
groot 已提交
3268 3269 3270 3271
			},
		}, nil
	}

3272 3273 3274
	vectorsRight := request.GetOpRight().GetDataArray()
	opRight := request.GetOpRight().GetIdArray()
	if opRight != nil {
3275 3276
		log.Debug("OpRight IdArray not empty, Get vectors by id",
			zap.String("traceID", traceID),
3277
			zap.String("role", typeutil.ProxyRole))
3278

3279
		result, err := query(opRight)
3280
		if err != nil {
3281 3282 3283
			log.Debug("Failed to get right vectors by id",
				zap.Error(err),
				zap.String("traceID", traceID),
3284
				zap.String("role", typeutil.ProxyRole))
3285

3286 3287 3288 3289 3290 3291 3292 3293
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3294 3295
		log.Debug("OpRight IdArray not empty, Get vectors by id done",
			zap.String("traceID", traceID),
3296
			zap.String("role", typeutil.ProxyRole))
3297

3298 3299
		vectorsRight, err = arrangeFunc(opRight, result.FieldsData)
		if err != nil {
3300 3301 3302
			log.Debug("Failed to re-arrange right vectors",
				zap.Error(err),
				zap.String("traceID", traceID),
3303
				zap.String("role", typeutil.ProxyRole))
3304

3305 3306 3307 3308 3309 3310
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
3311
		}
3312 3313 3314

		log.Debug("Re-arrange right vectors done",
			zap.String("traceID", traceID),
3315
			zap.String("role", typeutil.ProxyRole))
3316 3317
	}

G
groot 已提交
3318
	if vectorsRight == nil {
3319 3320 3321
		msg := "Right vectors array is empty"
		log.Debug(msg,
			zap.String("traceID", traceID),
3322
			zap.String("role", typeutil.ProxyRole))
3323

G
groot 已提交
3324 3325 3326
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3327
				Reason:    msg,
G
groot 已提交
3328 3329 3330 3331
			},
		}, nil
	}

3332
	if vectorsLeft.Dim != vectorsRight.Dim {
3333 3334 3335
		msg := "Vectors dimension is not equal"
		log.Debug(msg,
			zap.String("traceID", traceID),
3336
			zap.String("role", typeutil.ProxyRole))
3337

3338 3339 3340
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3341
				Reason:    msg,
3342 3343 3344 3345 3346
			},
		}, nil
	}

	if vectorsLeft.GetFloatVector() != nil && vectorsRight.GetFloatVector() != nil {
3347 3348
		distances, err := distance.CalcFloatDistance(vectorsLeft.Dim, vectorsLeft.GetFloatVector().Data, vectorsRight.GetFloatVector().Data, metric)
		if err != nil {
3349 3350 3351
			log.Debug("Failed to CalcFloatDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3352
				zap.String("role", typeutil.ProxyRole))
3353

3354 3355 3356 3357 3358 3359 3360 3361
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{
					ErrorCode: commonpb.ErrorCode_UnexpectedError,
					Reason:    err.Error(),
				},
			}, nil
		}

3362 3363 3364
		log.Debug("CalcFloatDistance done",
			zap.Error(err),
			zap.String("traceID", traceID),
3365
			zap.String("role", typeutil.ProxyRole))
3366

3367 3368 3369 3370 3371 3372 3373 3374 3375 3376
		return &milvuspb.CalcDistanceResults{
			Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
			Array: &milvuspb.CalcDistanceResults_FloatDist{
				FloatDist: &schemapb.FloatArray{
					Data: distances,
				},
			},
		}, nil
	}

3377
	if vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetBinaryVector() != nil {
G
groot 已提交
3378
		hamming, err := distance.CalcHammingDistance(vectorsLeft.Dim, vectorsLeft.GetBinaryVector(), vectorsRight.GetBinaryVector())
3379
		if err != nil {
3380 3381 3382
			log.Debug("Failed to CalcHammingDistance",
				zap.Error(err),
				zap.String("traceID", traceID),
3383
				zap.String("role", typeutil.ProxyRole))
3384

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

		if metric == distance.HAMMING {
3394 3395
			log.Debug("CalcHammingDistance done",
				zap.String("traceID", traceID),
3396
				zap.String("role", typeutil.ProxyRole))
3397

3398 3399 3400 3401
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_IntDist{
					IntDist: &schemapb.IntArray{
G
groot 已提交
3402
						Data: hamming,
3403 3404 3405 3406 3407 3408
					},
				},
			}, nil
		}

		if metric == distance.TANIMOTO {
G
groot 已提交
3409
			tanimoto, err := distance.CalcTanimotoCoefficient(vectorsLeft.Dim, hamming)
3410
			if err != nil {
3411 3412 3413
				log.Debug("Failed to CalcTanimotoCoefficient",
					zap.Error(err),
					zap.String("traceID", traceID),
3414
					zap.String("role", typeutil.ProxyRole))
3415

3416 3417 3418 3419 3420 3421 3422 3423
				return &milvuspb.CalcDistanceResults{
					Status: &commonpb.Status{
						ErrorCode: commonpb.ErrorCode_UnexpectedError,
						Reason:    err.Error(),
					},
				}, nil
			}

3424 3425
			log.Debug("CalcTanimotoCoefficient done",
				zap.String("traceID", traceID),
3426
				zap.String("role", typeutil.ProxyRole))
3427

3428 3429 3430 3431 3432 3433 3434 3435 3436 3437 3438
			return &milvuspb.CalcDistanceResults{
				Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success, Reason: ""},
				Array: &milvuspb.CalcDistanceResults_FloatDist{
					FloatDist: &schemapb.FloatArray{
						Data: tanimoto,
					},
				},
			}, nil
		}
	}

3439
	err = errors.New("unexpected error")
3440
	if (vectorsLeft.GetBinaryVector() != nil && vectorsRight.GetFloatVector() != nil) || (vectorsLeft.GetFloatVector() != nil && vectorsRight.GetBinaryVector() != nil) {
3441
		err = errors.New("cannot calculate distance between binary vectors and float vectors")
3442 3443
	}

3444 3445 3446
	log.Debug("Failed to CalcDistance",
		zap.Error(err),
		zap.String("traceID", traceID),
3447
		zap.String("role", typeutil.ProxyRole))
3448

3449 3450 3451 3452 3453 3454 3455 3456
	return &milvuspb.CalcDistanceResults{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
			Reason:    err.Error(),
		},
	}, nil
}

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

3462
// GetPersistentSegmentInfo get the information of sealed segment.
C
Cai Yudong 已提交
3463
func (node *Proxy) GetPersistentSegmentInfo(ctx context.Context, req *milvuspb.GetPersistentSegmentInfoRequest) (*milvuspb.GetPersistentSegmentInfoResponse, error) {
D
dragondriver 已提交
3464
	log.Debug("GetPersistentSegmentInfo",
3465
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3466 3467 3468
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

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

J
jingkl 已提交
3524
// GetQuerySegmentInfo gets segment information from QueryCoord.
C
Cai Yudong 已提交
3525
func (node *Proxy) GetQuerySegmentInfo(ctx context.Context, req *milvuspb.GetQuerySegmentInfoRequest) (*milvuspb.GetQuerySegmentInfoResponse, error) {
D
dragondriver 已提交
3526
	log.Debug("GetQuerySegmentInfo",
3527
		zap.String("role", typeutil.ProxyRole),
D
dragondriver 已提交
3528 3529 3530
		zap.String("db", req.DbName),
		zap.Any("collection", req.CollectionName))

G
godchen 已提交
3531
	resp := &milvuspb.GetQuerySegmentInfoResponse{
Z
zhenshan.cao 已提交
3532
		Status: &commonpb.Status{
3533
			ErrorCode: commonpb.ErrorCode_UnexpectedError,
Z
zhenshan.cao 已提交
3534 3535
		},
	}
3536 3537 3538 3539
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
G
godchen 已提交
3540
	segments, err := node.getSegmentsOfCollection(ctx, req.DbName, req.CollectionName)
Z
zhenshan.cao 已提交
3541 3542 3543 3544
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3545 3546 3547 3548 3549
	collID, err := globalMetaCache.GetCollectionID(ctx, req.CollectionName)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3550
	infoResp, err := node.queryCoord.GetSegmentInfo(ctx, &querypb.GetSegmentInfoRequest{
Z
zhenshan.cao 已提交
3551
		Base: &commonpb.MsgBase{
3552
			MsgType:   commonpb.MsgType_SegmentInfo,
Z
zhenshan.cao 已提交
3553 3554
			MsgID:     0,
			Timestamp: 0,
3555
			SourceID:  Params.ProxyCfg.ProxyID,
Z
zhenshan.cao 已提交
3556
		},
3557 3558
		CollectionID: collID,
		SegmentIDs:   segments,
Z
zhenshan.cao 已提交
3559 3560
	})
	if err != nil {
3561
		log.Error("Failed to get segment info from QueryCoord",
3562
			zap.Int64s("segmentIDs", segments), zap.Error(err))
Z
zhenshan.cao 已提交
3563 3564 3565
		resp.Status.Reason = err.Error()
		return resp, nil
	}
3566
	log.Debug("GetQuerySegmentInfo ", zap.Any("infos", infoResp.Infos), zap.Any("status", infoResp.Status))
3567
	if infoResp.Status.ErrorCode != commonpb.ErrorCode_Success {
3568
		log.Error("Failed to get segment info from QueryCoord", zap.String("errMsg", infoResp.Status.Reason))
Z
zhenshan.cao 已提交
3569 3570 3571 3572 3573 3574 3575 3576 3577 3578 3579 3580 3581
		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,
3582
			NodeID:       info.NodeID,
X
xige-16 已提交
3583
			State:        info.SegmentState,
Z
zhenshan.cao 已提交
3584 3585
		}
	}
3586
	resp.Status.ErrorCode = commonpb.ErrorCode_Success
Z
zhenshan.cao 已提交
3587 3588 3589 3590
	resp.Infos = queryInfos
	return resp, nil
}

C
Cai Yudong 已提交
3591
func (node *Proxy) getSegmentsOfCollection(ctx context.Context, dbName string, collectionName string) ([]UniqueID, error) {
3592
	describeCollectionResponse, err := node.rootCoord.DescribeCollection(ctx, &milvuspb.DescribeCollectionRequest{
X
XuanYang-cn 已提交
3593
		Base: &commonpb.MsgBase{
3594
			MsgType:   commonpb.MsgType_DescribeCollection,
X
XuanYang-cn 已提交
3595 3596
			MsgID:     0,
			Timestamp: 0,
3597
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3598 3599 3600 3601 3602 3603 3604
		},
		DbName:         dbName,
		CollectionName: collectionName,
	})
	if err != nil {
		return nil, err
	}
3605
	if describeCollectionResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3606 3607 3608
		return nil, errors.New(describeCollectionResponse.Status.Reason)
	}
	collectionID := describeCollectionResponse.CollectionID
3609
	showPartitionsResp, err := node.rootCoord.ShowPartitions(ctx, &milvuspb.ShowPartitionsRequest{
X
XuanYang-cn 已提交
3610
		Base: &commonpb.MsgBase{
3611
			MsgType:   commonpb.MsgType_ShowPartitions,
X
XuanYang-cn 已提交
3612 3613
			MsgID:     0,
			Timestamp: 0,
3614
			SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3615 3616 3617 3618 3619 3620 3621 3622
		},
		DbName:         dbName,
		CollectionName: collectionName,
		CollectionID:   collectionID,
	})
	if err != nil {
		return nil, err
	}
3623
	if showPartitionsResp.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3624 3625 3626 3627 3628
		return nil, errors.New(showPartitionsResp.Status.Reason)
	}

	ret := make([]UniqueID, 0)
	for _, partitionID := range showPartitionsResp.PartitionIDs {
3629
		showSegmentResponse, err := node.rootCoord.ShowSegments(ctx, &milvuspb.ShowSegmentsRequest{
X
XuanYang-cn 已提交
3630
			Base: &commonpb.MsgBase{
3631
				MsgType:   commonpb.MsgType_ShowSegments,
X
XuanYang-cn 已提交
3632 3633
				MsgID:     0,
				Timestamp: 0,
3634
				SourceID:  Params.ProxyCfg.ProxyID,
X
XuanYang-cn 已提交
3635 3636 3637 3638 3639 3640 3641
			},
			CollectionID: collectionID,
			PartitionID:  partitionID,
		})
		if err != nil {
			return nil, err
		}
3642
		if showSegmentResponse.Status.ErrorCode != commonpb.ErrorCode_Success {
X
XuanYang-cn 已提交
3643 3644 3645 3646 3647 3648
			return nil, errors.New(showSegmentResponse.Status.Reason)
		}
		ret = append(ret, showSegmentResponse.SegmentIDs...)
	}
	return ret, nil
}
3649

J
jingkl 已提交
3650
// Dummy handles dummy request
C
Cai Yudong 已提交
3651
func (node *Proxy) Dummy(ctx context.Context, req *milvuspb.DummyRequest) (*milvuspb.DummyResponse, error) {
3652 3653 3654 3655 3656 3657 3658 3659 3660 3661 3662
	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
	}

3663 3664
	if drt.RequestType == "query" {
		drr, err := parseDummyQueryRequest(req.RequestType)
3665
		if err != nil {
3666
			log.Debug("Failed to parse dummy query request")
3667 3668 3669
			return failedResponse, nil
		}

3670
		request := &milvuspb.QueryRequest{
3671 3672 3673
			DbName:         drr.DbName,
			CollectionName: drr.CollectionName,
			PartitionNames: drr.PartitionNames,
3674
			OutputFields:   drr.OutputFields,
X
Xiangyu Wang 已提交
3675 3676
		}

3677
		_, err = node.Query(ctx, request)
3678
		if err != nil {
3679
			log.Debug("Failed to execute dummy query")
3680 3681
			return failedResponse, err
		}
X
Xiangyu Wang 已提交
3682 3683 3684 3685 3686 3687

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

3688 3689
	log.Debug("cannot find specify dummy request type")
	return failedResponse, nil
X
Xiangyu Wang 已提交
3690 3691
}

J
jingkl 已提交
3692
// RegisterLink registers a link
C
Cai Yudong 已提交
3693
func (node *Proxy) RegisterLink(ctx context.Context, req *milvuspb.RegisterLinkRequest) (*milvuspb.RegisterLinkResponse, error) {
G
godchen 已提交
3694
	code := node.stateCode.Load().(internalpb.StateCode)
D
dragondriver 已提交
3695
	log.Debug("RegisterLink",
3696
		zap.String("role", typeutil.ProxyRole),
C
Cai Yudong 已提交
3697
		zap.Any("state code of proxy", code))
D
dragondriver 已提交
3698

G
godchen 已提交
3699
	if code != internalpb.StateCode_Healthy {
3700 3701 3702
		return &milvuspb.RegisterLinkResponse{
			Address: nil,
			Status: &commonpb.Status{
3703
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3704
				Reason:    "proxy not healthy",
3705 3706 3707
			},
		}, nil
	}
3708
	//metrics.ProxyLinkedSDKs.WithLabelValues(strconv.FormatInt(Params.ProxyCfg.ProxyID, 10)).Inc()
3709 3710 3711
	return &milvuspb.RegisterLinkResponse{
		Address: nil,
		Status: &commonpb.Status{
3712
			ErrorCode: commonpb.ErrorCode_Success,
3713
			Reason:    os.Getenv(metricsinfo.DeployModeEnvKey),
3714 3715 3716
		},
	}, nil
}
3717

3718 3719 3720
// 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",
3721
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3722 3723 3724 3725
		zap.String("req", req.Request))

	if !node.checkHealthy() {
		log.Warn("Proxy.GetMetrics failed",
3726
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3727
			zap.String("req", req.Request),
3728
			zap.Error(errProxyIsUnhealthy(Params.ProxyCfg.ProxyID)))
3729 3730 3731 3732

		return &milvuspb.GetMetricsResponse{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UnexpectedError,
3733
				Reason:    msgProxyIsUnhealthy(Params.ProxyCfg.ProxyID),
3734 3735 3736 3737 3738 3739 3740 3741
			},
			Response: "",
		}, nil
	}

	metricType, err := metricsinfo.ParseMetricType(req.Request)
	if err != nil {
		log.Warn("Proxy.GetMetrics failed to parse metric type",
3742
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3743 3744 3745 3746 3747 3748 3749 3750 3751 3752 3753 3754 3755 3756 3757
			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 已提交
3758 3759 3760 3761 3762 3763 3764 3765 3766 3767
	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,
3768
		SourceID:  Params.ProxyCfg.ProxyID,
D
dragondriver 已提交
3769 3770
	}

3771
	if metricType == metricsinfo.SystemInfoMetrics {
3772 3773 3774 3775 3776 3777 3778
		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))

3779
		metrics, err := getSystemInfoMetrics(ctx, req, node)
3780 3781

		log.Debug("Proxy.GetMetrics",
3782
			zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3783 3784 3785 3786 3787
			zap.String("req", req.Request),
			zap.String("metric_type", metricType),
			zap.Any("metrics", metrics), // TODO(dragondriver): necessary? may be very large
			zap.Error(err))

3788 3789
		node.metricsCacheManager.UpdateSystemInfoMetrics(metrics)

G
godchen 已提交
3790
		return metrics, nil
3791 3792 3793
	}

	log.Debug("Proxy.GetMetrics failed, request metric type is not implemented yet",
3794
		zap.Int64("node_id", Params.ProxyCfg.ProxyID),
3795 3796 3797 3798 3799 3800 3801 3802 3803 3804 3805 3806
		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 已提交
3807 3808 3809
// 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",
3810
		zap.Int64("proxy_id", Params.ProxyCfg.ProxyID),
B
bigsheeper 已提交
3811 3812 3813 3814 3815 3816 3817 3818 3819 3820 3821 3822 3823 3824
		zap.Any("req", req))

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

	status := &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
	}
	infoResp, err := node.queryCoord.LoadBalance(ctx, &querypb.LoadBalanceRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_LoadBalanceSegments,
			MsgID:     0,
			Timestamp: 0,
3825
			SourceID:  Params.ProxyCfg.ProxyID,
B
bigsheeper 已提交
3826 3827 3828
		},
		SourceNodeIDs:    []int64{req.SrcNodeID},
		DstNodeIDs:       req.DstNodeIDs,
X
xige-16 已提交
3829
		BalanceReason:    querypb.TriggerCondition_GrpcRequest,
B
bigsheeper 已提交
3830 3831 3832 3833 3834 3835 3836 3837 3838 3839 3840 3841 3842 3843 3844 3845 3846 3847
		SealedSegmentIDs: req.SealedSegmentIDs,
	})
	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 已提交
3848
//GetCompactionState gets the compaction state of multiple segments
3849 3850 3851 3852 3853 3854 3855 3856 3857 3858 3859 3860 3861
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
}

3862
// ManualCompaction invokes compaction on specified collection
3863 3864 3865 3866 3867 3868 3869 3870 3871 3872 3873 3874 3875 3876 3877 3878 3879 3880 3881 3882 3883 3884 3885 3886 3887 3888
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
}

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 已提交
3889 3890 3891
// 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))
3892
	var err error
B
Bingyi Sun 已提交
3893 3894 3895 3896 3897 3898 3899
	resp := &milvuspb.GetFlushStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		log.Info("unable to get flush state because of closed server")
		return resp, nil
	}

3900
	resp, err = node.dataCoord.GetFlushState(ctx, req)
X
Xiaofan 已提交
3901 3902 3903 3904
	if err != nil {
		log.Info("failed to get flush state response", zap.Error(err))
		return nil, err
	}
B
Bingyi Sun 已提交
3905 3906 3907 3908
	log.Info("received get flush state response", zap.Any("response", resp))
	return resp, err
}

C
Cai Yudong 已提交
3909 3910
// checkHealthy checks proxy state is Healthy
func (node *Proxy) checkHealthy() bool {
3911 3912 3913 3914
	code := node.stateCode.Load().(internalpb.StateCode)
	return code == internalpb.StateCode_Healthy
}

J
jingkl 已提交
3915
//unhealthyStatus returns the proxy not healthy status
3916 3917 3918
func unhealthyStatus() *commonpb.Status {
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_UnexpectedError,
C
Cai Yudong 已提交
3919
		Reason:    "proxy not healthy",
3920 3921
	}
}
G
groot 已提交
3922 3923 3924

// 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) {
G
groot 已提交
3925
	log.Info("received import request")
3926 3927 3928 3929 3930 3931
	resp := &milvuspb.ImportResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_Success,
			Reason:    "",
		},
	}
G
groot 已提交
3932 3933 3934 3935
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}
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
	// 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()
		return resp, err
	}
	chNames, err := node.chMgr.getVChannels(collID)
	if err != nil {
		log.Error("get vChannels failed",
			zap.Int64("collection ID", collID),
			zap.Error(err))
		resp.Status.ErrorCode = commonpb.ErrorCode_UnexpectedError
		resp.Status.Reason = err.Error()
		return resp, err
	}
	req.ChannelNames = chNames
	// Call rootCoord to finish import.
	resp, err = node.rootCoord.Import(ctx, req)
	log.Info("received import response",
		zap.String("collection name", req.GetCollectionName()),
		zap.Any("resp", resp),
		zap.Error(err))
G
groot 已提交
3962
	return resp, err
G
groot 已提交
3963 3964
}

X
XuanYang-cn 已提交
3965 3966 3967 3968 3969 3970 3971 3972 3973 3974 3975 3976 3977 3978
// 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
	}

	resp, err := node.queryCoord.GetReplicas(ctx, req)
	log.Info("received get replicas response", zap.Any("resp", resp), zap.Error(err))
	return resp, err
}

G
groot 已提交
3979 3980
// Check import task state from datanode
func (node *Proxy) GetImportState(ctx context.Context, req *milvuspb.GetImportStateRequest) (*milvuspb.GetImportStateResponse, error) {
G
groot 已提交
3981
	log.Info("received get import state request", zap.Int64("taskID", req.GetTask()))
G
groot 已提交
3982 3983 3984 3985 3986 3987
	resp := &milvuspb.GetImportStateResponse{}
	if !node.checkHealthy() {
		resp.Status = unhealthyStatus()
		return resp, nil
	}

G
groot 已提交
3988 3989 3990
	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
G
groot 已提交
3991
}
3992 3993 3994 3995 3996 3997 3998 3999 4000 4001 4002 4003 4004 4005 4006 4007 4008 4009 4010 4011 4012 4013 4014 4015 4016 4017 4018 4019 4020 4021 4022 4023 4024 4025 4026 4027 4028 4029 4030 4031 4032 4033 4034 4035 4036 4037 4038 4039 4040 4041 4042 4043 4044 4045 4046 4047 4048 4049 4050 4051 4052 4053 4054 4055 4056 4057 4058 4059 4060 4061 4062 4063 4064 4065 4066 4067 4068 4069 4070 4071 4072 4073 4074 4075 4076 4077 4078 4079 4080 4081 4082 4083 4084 4085 4086 4087 4088 4089 4090 4091 4092 4093 4094 4095 4096 4097 4098 4099 4100 4101 4102

// 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{
		Username:          request.Username,
		EncryptedPassword: request.Password,
	}
	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) ClearCredUsersCache(ctx context.Context, request *internalpb.ClearCredUsersCacheRequest) (*commonpb.Status, error) {
	ctx = logutil.WithModule(ctx, moduleName)
	logutil.Logger(ctx).Debug("received request to clear credential usernames cache",
		zap.String("role", typeutil.ProxyRole))

	if globalMetaCache != nil {
		globalMetaCache.ClearCredUsers() // no need to return error, though credential may be not cached
	}
	logutil.Logger(ctx).Debug("complete to clear credential usernames cache",
		zap.String("role", typeutil.ProxyRole))

	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,
	}
	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 已提交
4103
func (node *Proxy) UpdateCredential(ctx context.Context, req *milvuspb.UpdateCredentialRequest) (*commonpb.Status, error) {
4104
	log.Debug("UpdateCredential", zap.String("role", typeutil.RootCoordRole), zap.String("username", req.Username))
C
codeman 已提交
4105 4106 4107 4108 4109 4110 4111 4112 4113
	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)
4114 4115 4116 4117 4118 4119 4120
	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 已提交
4121 4122
	// valid new password
	if err = ValidatePassword(rawNewPassword); err != nil {
4123 4124 4125 4126 4127 4128
		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 已提交
4129 4130 4131 4132 4133 4134 4135 4136 4137 4138 4139 4140 4141 4142 4143 4144 4145
	// 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
	}
	if !crypto.PasswordVerify(rawOldPassword, oldCredInfo.EncryptedPassword) {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UpdateCredentialFailure,
			Reason:    "old password is not correct:" + req.Username,
		}, nil
	}
	// update meta data
	encryptedPassword, err := crypto.PasswordEncrypt(rawNewPassword)
4146 4147 4148 4149 4150 4151 4152
	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 已提交
4153
	updateCredReq := &internalpb.CredentialInfo{
4154 4155 4156
		Username:          req.Username,
		EncryptedPassword: encryptedPassword,
	}
C
codeman 已提交
4157
	result, err := node.rootCoord.UpdateCredential(ctx, updateCredReq)
4158 4159 4160 4161 4162 4163 4164 4165 4166 4167 4168 4169
	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))
4170 4171 4172 4173 4174 4175
	if req.Username == util.UserRoot {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_DeleteCredentialFailure,
			Reason:    "user root cannot be deleted",
		}, nil
	}
4176 4177 4178 4179 4180 4181 4182 4183 4184 4185 4186 4187 4188 4189 4190 4191 4192 4193 4194 4195 4196 4197 4198 4199 4200 4201 4202 4203 4204 4205
	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))
	// get from cache
	usernames, err := globalMetaCache.GetCredUsernames(ctx)
	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,
		},
		Usernames: usernames,
	}, nil
}