server.go 23.2 KB
Newer Older
S
sunby 已提交
1 2
package dataservice

S
sunby 已提交
3 4
import (
	"context"
S
sunby 已提交
5
	"errors"
S
sunby 已提交
6
	"fmt"
S
sunby 已提交
7 8
	"path"
	"strconv"
S
sunby 已提交
9
	"sync"
S
sunby 已提交
10 11 12
	"sync/atomic"
	"time"

13 14 15 16
	"go.uber.org/zap"

	"github.com/zilliztech/milvus-distributed/internal/log"

S
sunby 已提交
17 18
	"github.com/zilliztech/milvus-distributed/internal/distributed/datanode"

S
sunby 已提交
19 20
	"github.com/golang/protobuf/proto"
	"github.com/zilliztech/milvus-distributed/internal/proto/masterpb"
S
sunby 已提交
21 22 23

	"github.com/zilliztech/milvus-distributed/internal/msgstream"

S
sunby 已提交
24 25 26 27 28 29 30 31 32 33 34 35 36
	"github.com/zilliztech/milvus-distributed/internal/proto/milvuspb"

	"github.com/zilliztech/milvus-distributed/internal/timesync"

	etcdkv "github.com/zilliztech/milvus-distributed/internal/kv/etcd"
	"go.etcd.io/etcd/clientv3"

	"github.com/zilliztech/milvus-distributed/internal/proto/commonpb"
	"github.com/zilliztech/milvus-distributed/internal/proto/datapb"
	"github.com/zilliztech/milvus-distributed/internal/proto/internalpb2"
	"github.com/zilliztech/milvus-distributed/internal/util/typeutil"
)

S
sunby 已提交
37 38
const role = "dataservice"

S
sunby 已提交
39 40
type DataService interface {
	typeutil.Service
N
neza2017 已提交
41
	typeutil.Component
S
sunby 已提交
42 43 44 45 46 47 48
	RegisterNode(req *datapb.RegisterNodeRequest) (*datapb.RegisterNodeResponse, error)
	Flush(req *datapb.FlushRequest) (*commonpb.Status, error)

	AssignSegmentID(req *datapb.AssignSegIDRequest) (*datapb.AssignSegIDResponse, error)
	ShowSegments(req *datapb.ShowSegmentRequest) (*datapb.ShowSegmentResponse, error)
	GetSegmentStates(req *datapb.SegmentStatesRequest) (*datapb.SegmentStatesResponse, error)
	GetInsertBinlogPaths(req *datapb.InsertBinlogPathRequest) (*datapb.InsertBinlogPathsResponse, error)
G
godchen 已提交
49 50
	GetSegmentInfoChannel() (*milvuspb.StringResponse, error)
	GetInsertChannels(req *datapb.InsertChannelRequest) (*internalpb2.StringList, error)
S
sunby 已提交
51 52 53
	GetCollectionStatistics(req *datapb.CollectionStatsRequest) (*datapb.CollectionStatsResponse, error)
	GetPartitionStatistics(req *datapb.PartitionStatsRequest) (*datapb.PartitionStatsResponse, error)
	GetComponentStates() (*internalpb2.ComponentStates, error)
54
	GetCount(req *datapb.CollectionCountRequest) (*datapb.CollectionCountResponse, error)
X
XuanYang-cn 已提交
55
	GetSegmentInfo(req *datapb.SegmentInfoRequest) (*datapb.SegmentInfoResponse, error)
S
sunby 已提交
56 57
}

S
sunby 已提交
58 59 60 61 62 63 64 65 66 67
type MasterClient interface {
	ShowCollections(in *milvuspb.ShowCollectionRequest) (*milvuspb.ShowCollectionResponse, error)
	DescribeCollection(in *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error)
	ShowPartitions(in *milvuspb.ShowPartitionRequest) (*milvuspb.ShowPartitionResponse, error)
	GetDdChannel() (string, error)
	AllocTimestamp(in *masterpb.TsoRequest) (*masterpb.TsoResponse, error)
	AllocID(in *masterpb.IDRequest) (*masterpb.IDResponse, error)
	GetComponentStates() (*internalpb2.ComponentStates, error)
}

S
sunby 已提交
68 69 70 71 72 73 74
type DataNodeClient interface {
	WatchDmChannels(in *datapb.WatchDmChannelRequest) (*commonpb.Status, error)
	GetComponentStates(empty *commonpb.Empty) (*internalpb2.ComponentStates, error)
	FlushSegments(in *datapb.FlushSegRequest) (*commonpb.Status, error)
	Stop() error
}

S
sunby 已提交
75
type (
S
sunby 已提交
76 77 78
	UniqueID  = typeutil.UniqueID
	Timestamp = typeutil.Timestamp
	Server    struct {
N
neza2017 已提交
79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97
		ctx               context.Context
		serverLoopCtx     context.Context
		serverLoopCancel  context.CancelFunc
		serverLoopWg      sync.WaitGroup
		state             atomic.Value
		client            *etcdkv.EtcdKV
		meta              *meta
		segAllocator      segmentAllocator
		statsHandler      *statsHandler
		ddHandler         *ddHandler
		allocator         allocator
		cluster           *dataNodeCluster
		msgProducer       *timesync.MsgProducer
		registerFinishCh  chan struct{}
		masterClient      MasterClient
		ttMsgStream       msgstream.MsgStream
		k2sMsgStream      msgstream.MsgStream
		ddChannelName     string
		segmentInfoStream msgstream.MsgStream
98
		insertChannels    []string
G
groot 已提交
99
		msFactory         msgstream.Factory
C
cai.zhang 已提交
100
		ttBarrier         timesync.TimeTickBarrier
S
sunby 已提交
101 102 103
	}
)

G
groot 已提交
104
func CreateServer(ctx context.Context, factory msgstream.Factory) (*Server, error) {
S
sunby 已提交
105
	ch := make(chan struct{})
S
sunby 已提交
106
	s := &Server{
S
sunby 已提交
107
		ctx:              ctx,
S
sunby 已提交
108 109
		registerFinishCh: ch,
		cluster:          newDataNodeCluster(ch),
G
groot 已提交
110
		msFactory:        factory,
S
sunby 已提交
111
	}
112
	s.insertChannels = s.getInsertChannels()
N
neza2017 已提交
113
	s.state.Store(internalpb2.StateCode_INITIALIZING)
S
sunby 已提交
114 115 116
	return s, nil
}

117 118 119 120 121 122 123 124 125
func (s *Server) getInsertChannels() []string {
	channels := make([]string, Params.InsertChannelNum)
	var i int64 = 0
	for ; i < Params.InsertChannelNum; i++ {
		channels[i] = Params.InsertChannelPrefixName + strconv.FormatInt(i, 10)
	}
	return channels
}

S
sunby 已提交
126 127
func (s *Server) SetMasterClient(masterClient MasterClient) {
	s.masterClient = masterClient
S
sunby 已提交
128 129 130
}

func (s *Server) Init() error {
S
sunby 已提交
131 132 133 134
	return nil
}

func (s *Server) Start() error {
135
	var err error
G
groot 已提交
136 137 138 139 140 141 142 143 144
	m := map[string]interface{}{
		"PulsarAddress":  Params.PulsarAddress,
		"ReceiveBufSize": 1024,
		"PulsarBufSize":  1024}
	err = s.msFactory.SetParams(m)
	if err != nil {
		return err
	}

S
sunby 已提交
145
	s.allocator = newAllocatorImpl(s.masterClient)
146
	if err = s.initMeta(); err != nil {
S
sunby 已提交
147 148 149
		return err
	}
	s.statsHandler = newStatsHandler(s.meta)
Y
yukun 已提交
150
	s.segAllocator = newSegmentAllocator(s.meta, s.allocator)
N
neza2017 已提交
151
	s.ddHandler = newDDHandler(s.meta, s.segAllocator)
152 153
	s.initSegmentInfoChannel()
	if err = s.loadMetaFromMaster(); err != nil {
S
sunby 已提交
154 155
		return err
	}
156
	s.waitDataNodeRegister()
157
	s.cluster.WatchInsertChannels(s.insertChannels)
158 159 160
	if err = s.initMsgProducer(); err != nil {
		return err
	}
S
sunby 已提交
161
	s.startServerLoop()
S
sunby 已提交
162
	s.state.Store(internalpb2.StateCode_HEALTHY)
163
	log.Debug("start success")
S
sunby 已提交
164 165 166
	return nil
}

S
sunby 已提交
167 168 169 170
func (s *Server) checkStateIsHealthy() bool {
	return s.state.Load().(internalpb2.StateCode) == internalpb2.StateCode_HEALTHY
}

S
sunby 已提交
171 172 173 174 175 176 177
func (s *Server) initMeta() error {
	etcdClient, err := clientv3.New(clientv3.Config{Endpoints: []string{Params.EtcdAddress}})
	if err != nil {
		return err
	}
	etcdKV := etcdkv.NewEtcdKV(etcdClient, Params.MetaRootPath)
	s.client = etcdKV
N
neza2017 已提交
178
	s.meta, err = newMeta(etcdKV)
S
sunby 已提交
179 180 181 182 183 184
	if err != nil {
		return err
	}
	return nil
}

185
func (s *Server) initSegmentInfoChannel() {
G
groot 已提交
186
	segmentInfoStream, _ := s.msFactory.NewMsgStream(s.ctx)
Z
zhenshan.cao 已提交
187
	segmentInfoStream.AsProducer([]string{Params.SegmentInfoChannelName})
188 189
	s.segmentInfoStream = segmentInfoStream
	s.segmentInfoStream.Start()
S
sunby 已提交
190
}
S
sunby 已提交
191
func (s *Server) initMsgProducer() error {
C
cai.zhang 已提交
192
	var err error
193
	if s.ttMsgStream, err = s.msFactory.NewMsgStream(s.ctx); err != nil {
C
cai.zhang 已提交
194 195 196
		return err
	}
	s.ttMsgStream.AsConsumer([]string{Params.TimeTickChannelName}, Params.DataServiceSubscriptionName)
S
sunby 已提交
197
	s.ttMsgStream.Start()
S
sunby 已提交
198 199
	s.ttBarrier = timesync.NewHardTimeTickBarrier(s.ctx, s.ttMsgStream, s.cluster.GetNodeIDs())
	s.ttBarrier.Start()
G
groot 已提交
200
	if s.k2sMsgStream, err = s.msFactory.NewMsgStream(s.ctx); err != nil {
C
cai.zhang 已提交
201 202 203
		return err
	}
	s.k2sMsgStream.AsProducer(Params.K2SChannelNames)
204
	s.k2sMsgStream.Start()
C
cai.zhang 已提交
205
	dataNodeTTWatcher := newDataNodeTimeTickWatcher(s.meta, s.segAllocator, s.cluster)
206
	k2sMsgWatcher := timesync.NewMsgTimeTickWatcher(s.k2sMsgStream)
C
cai.zhang 已提交
207
	if s.msgProducer, err = timesync.NewTimeSyncMsgProducer(s.ttBarrier, dataNodeTTWatcher, k2sMsgWatcher); err != nil {
S
sunby 已提交
208 209 210 211 212
		return err
	}
	s.msgProducer.Start(s.ctx)
	return nil
}
S
sunby 已提交
213

S
sunby 已提交
214
func (s *Server) loadMetaFromMaster() error {
215
	log.Debug("loading collection meta from master")
S
sunby 已提交
216 217 218
	if err := s.checkMasterIsHealthy(); err != nil {
		return err
	}
N
neza2017 已提交
219 220 221 222 223 224 225
	if s.ddChannelName == "" {
		channel, err := s.masterClient.GetDdChannel()
		if err != nil {
			return err
		}
		s.ddChannelName = channel
	}
S
sunby 已提交
226 227 228 229 230
	collections, err := s.masterClient.ShowCollections(&milvuspb.ShowCollectionRequest{
		Base: &commonpb.MsgBase{
			MsgType:   commonpb.MsgType_kShowCollections,
			MsgID:     -1, // todo add msg id
			Timestamp: 0,  // todo
S
sunby 已提交
231
			SourceID:  Params.NodeID,
S
sunby 已提交
232 233 234 235 236 237 238 239 240 241 242 243
		},
		DbName: "",
	})
	if err != nil {
		return err
	}
	for _, collectionName := range collections.CollectionNames {
		collection, err := s.masterClient.DescribeCollection(&milvuspb.DescribeCollectionRequest{
			Base: &commonpb.MsgBase{
				MsgType:   commonpb.MsgType_kDescribeCollection,
				MsgID:     -1, // todo
				Timestamp: 0,  // todo
S
sunby 已提交
244
				SourceID:  Params.NodeID,
S
sunby 已提交
245 246 247 248 249
			},
			DbName:         "",
			CollectionName: collectionName,
		})
		if err != nil {
250
			log.Error("describe collection error", zap.Error(err))
S
sunby 已提交
251 252 253 254 255 256 257
			continue
		}
		partitions, err := s.masterClient.ShowPartitions(&milvuspb.ShowPartitionRequest{
			Base: &commonpb.MsgBase{
				MsgType:   commonpb.MsgType_kShowPartitions,
				MsgID:     -1, // todo
				Timestamp: 0,  // todo
S
sunby 已提交
258
				SourceID:  Params.NodeID,
S
sunby 已提交
259 260 261 262 263 264
			},
			DbName:         "",
			CollectionName: collectionName,
			CollectionID:   collection.CollectionID,
		})
		if err != nil {
265
			log.Error("show partitions error", zap.Error(err))
S
sunby 已提交
266 267 268 269 270
			continue
		}
		err = s.meta.AddCollection(&collectionInfo{
			ID:         collection.CollectionID,
			Schema:     collection.Schema,
S
sunby 已提交
271
			Partitions: partitions.PartitionIDs,
S
sunby 已提交
272 273
		})
		if err != nil {
274
			log.Error("add collection error", zap.Error(err))
S
sunby 已提交
275 276 277
			continue
		}
	}
278
	log.Debug("load collection meta from master complete")
S
sunby 已提交
279
	return nil
S
sunby 已提交
280
}
S
sunby 已提交
281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310

func (s *Server) checkMasterIsHealthy() error {
	ticker := time.NewTicker(300 * time.Millisecond)
	ctx, cancel := context.WithTimeout(s.ctx, 30*time.Second)
	defer func() {
		ticker.Stop()
		cancel()
	}()
	for {
		var resp *internalpb2.ComponentStates
		var err error
		select {
		case <-ctx.Done():
			return fmt.Errorf("master is not healthy")
		case <-ticker.C:
			resp, err = s.masterClient.GetComponentStates()
			if err != nil {
				return err
			}
			if resp.Status.ErrorCode != commonpb.ErrorCode_SUCCESS {
				return errors.New(resp.Status.Reason)
			}
		}
		if resp.State.StateCode == internalpb2.StateCode_HEALTHY {
			break
		}
	}
	return nil
}

311 312
func (s *Server) startServerLoop() {
	s.serverLoopCtx, s.serverLoopCancel = context.WithCancel(s.ctx)
S
sunby 已提交
313
	s.serverLoopWg.Add(3)
314 315
	go s.startStatsChannel(s.serverLoopCtx)
	go s.startSegmentFlushChannel(s.serverLoopCtx)
N
neza2017 已提交
316
	go s.startDDChannel(s.serverLoopCtx)
317 318 319 320
}

func (s *Server) startStatsChannel(ctx context.Context) {
	defer s.serverLoopWg.Done()
G
groot 已提交
321
	statsStream, _ := s.msFactory.NewMsgStream(ctx)
Z
zhenshan.cao 已提交
322
	statsStream.AsConsumer([]string{Params.StatisticsChannelName}, Params.DataServiceSubscriptionName)
323 324 325 326 327 328 329 330 331 332 333 334 335
	statsStream.Start()
	defer statsStream.Close()
	for {
		select {
		case <-ctx.Done():
			return
		default:
		}
		msgPack := statsStream.Consume()
		for _, msg := range msgPack.Msgs {
			statistics := msg.(*msgstream.SegmentStatisticsMsg)
			for _, stat := range statistics.SegStats {
				if err := s.statsHandler.HandleSegmentStat(stat); err != nil {
336
					log.Error("handle segment stat error", zap.Error(err))
337 338 339 340 341 342 343 344 345
					continue
				}
			}
		}
	}
}

func (s *Server) startSegmentFlushChannel(ctx context.Context) {
	defer s.serverLoopWg.Done()
G
groot 已提交
346
	flushStream, _ := s.msFactory.NewMsgStream(ctx)
Z
zhenshan.cao 已提交
347
	flushStream.AsConsumer([]string{Params.SegmentInfoChannelName}, Params.DataServiceSubscriptionName)
348 349 350 351 352
	flushStream.Start()
	defer flushStream.Close()
	for {
		select {
		case <-ctx.Done():
353
			log.Debug("segment flush channel shut down")
354 355 356 357 358 359 360 361 362 363 364 365
			return
		default:
		}
		msgPack := flushStream.Consume()
		for _, msg := range msgPack.Msgs {
			if msg.Type() != commonpb.MsgType_kSegmentFlushDone {
				continue
			}
			realMsg := msg.(*msgstream.FlushCompletedMsg)

			segmentInfo, err := s.meta.GetSegment(realMsg.SegmentID)
			if err != nil {
366
				log.Error("get segment error", zap.Error(err))
367 368 369
				continue
			}
			segmentInfo.FlushedTime = realMsg.BeginTimestamp
S
sunby 已提交
370
			segmentInfo.State = commonpb.SegmentState_SegmentFlushed
371
			if err = s.meta.UpdateSegment(segmentInfo); err != nil {
372
				log.Error("update segment error", zap.Error(err))
373 374 375 376 377 378
				continue
			}
		}
	}
}

N
neza2017 已提交
379 380
func (s *Server) startDDChannel(ctx context.Context) {
	defer s.serverLoopWg.Done()
G
groot 已提交
381
	ddStream, _ := s.msFactory.NewMsgStream(ctx)
Z
zhenshan.cao 已提交
382
	ddStream.AsConsumer([]string{s.ddChannelName}, Params.DataServiceSubscriptionName)
N
neza2017 已提交
383 384 385 386 387
	ddStream.Start()
	defer ddStream.Close()
	for {
		select {
		case <-ctx.Done():
388
			log.Debug("dd channel shut down")
N
neza2017 已提交
389 390 391 392 393 394
			return
		default:
		}
		msgPack := ddStream.Consume()
		for _, msg := range msgPack.Msgs {
			if err := s.ddHandler.HandleDDMsg(msg); err != nil {
395
				log.Error("handle dd msg error", zap.Error(err))
N
neza2017 已提交
396 397 398 399 400 401
				continue
			}
		}
	}
}

402
func (s *Server) waitDataNodeRegister() {
403
	log.Debug("waiting data node to register")
404
	<-s.registerFinishCh
405
	log.Debug("all data nodes register")
406
}
S
sunby 已提交
407 408

func (s *Server) Stop() error {
S
sunby 已提交
409
	s.cluster.ShutDownClients()
S
sunby 已提交
410
	s.ttBarrier.Close()
S
sunby 已提交
411
	s.ttMsgStream.Close()
412
	s.k2sMsgStream.Close()
S
sunby 已提交
413
	s.msgProducer.Close()
N
neza2017 已提交
414
	s.segmentInfoStream.Close()
S
sunby 已提交
415
	s.stopServerLoop()
S
sunby 已提交
416 417 418
	return nil
}

S
sunby 已提交
419 420 421 422 423
func (s *Server) stopServerLoop() {
	s.serverLoopCancel()
	s.serverLoopWg.Wait()
}

S
sunby 已提交
424
func (s *Server) GetComponentStates() (*internalpb2.ComponentStates, error) {
S
sunby 已提交
425 426 427 428
	resp := &internalpb2.ComponentStates{
		State: &internalpb2.ComponentInfo{
			NodeID:    Params.NodeID,
			Role:      role,
S
sunby 已提交
429
			StateCode: s.state.Load().(internalpb2.StateCode),
S
sunby 已提交
430 431 432 433 434 435 436 437 438 439 440 441 442
		},
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
	dataNodeStates, err := s.cluster.GetDataNodeStates()
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	resp.SubcomponentStates = dataNodeStates
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	return resp, nil
S
sunby 已提交
443 444
}

G
godchen 已提交
445 446 447 448 449 450 451
func (s *Server) GetTimeTickChannel() (*milvuspb.StringResponse, error) {
	return &milvuspb.StringResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_SUCCESS,
		},
		Value: Params.TimeTickChannelName,
	}, nil
S
sunby 已提交
452 453
}

G
godchen 已提交
454 455 456 457 458 459 460
func (s *Server) GetStatisticsChannel() (*milvuspb.StringResponse, error) {
	return &milvuspb.StringResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_SUCCESS,
		},
		Value: Params.StatisticsChannelName,
	}, nil
S
sunby 已提交
461 462 463
}

func (s *Server) RegisterNode(req *datapb.RegisterNodeRequest) (*datapb.RegisterNodeResponse, error) {
S
sunby 已提交
464
	ret := &datapb.RegisterNodeResponse{
S
sunby 已提交
465
		Status: &commonpb.Status{
S
sunby 已提交
466
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
S
sunby 已提交
467
		},
S
sunby 已提交
468
	}
S
sunby 已提交
469 470 471 472 473
	node, err := s.newDataNode(req.Address.Ip, req.Address.Port, req.Base.SourceID)
	if err != nil {
		return nil, err
	}
	s.cluster.Register(node)
N
neza2017 已提交
474
	if s.ddChannelName == "" {
N
neza2017 已提交
475
		resp, err := s.masterClient.GetDdChannel()
S
sunby 已提交
476 477 478 479
		if err != nil {
			ret.Status.Reason = err.Error()
			return ret, err
		}
N
neza2017 已提交
480
		s.ddChannelName = resp
S
sunby 已提交
481 482 483 484 485 486 487 488
	}
	ret.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	ret.InitParams = &internalpb2.InitParams{
		NodeID: Params.NodeID,
		StartParams: []*commonpb.KeyValuePair{
			{Key: "DDChannelName", Value: s.ddChannelName},
			{Key: "SegmentStatisticsChannelName", Value: Params.StatisticsChannelName},
			{Key: "TimeTickChannelName", Value: Params.TimeTickChannelName},
N
neza2017 已提交
489
			{Key: "CompleteFlushChannelName", Value: Params.SegmentInfoChannelName},
S
sunby 已提交
490 491 492
		},
	}
	return ret, nil
S
sunby 已提交
493 494
}

S
sunby 已提交
495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513
func (s *Server) newDataNode(ip string, port int64, id UniqueID) (*dataNode, error) {
	client := datanode.NewClient(fmt.Sprintf("%s:%d", ip, port))
	if err := client.Init(); err != nil {
		return nil, err
	}
	if err := client.Start(); err != nil {
		return nil, err
	}
	return &dataNode{
		id: id,
		address: struct {
			ip   string
			port int64
		}{ip: ip, port: port},
		client:     client,
		channelNum: 0,
	}, nil
}

S
sunby 已提交
514
func (s *Server) Flush(req *datapb.FlushRequest) (*commonpb.Status, error) {
S
sunby 已提交
515 516 517 518 519 520
	if !s.checkStateIsHealthy() {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
			Reason:    "server is initializing",
		}, nil
	}
N
neza2017 已提交
521
	s.segAllocator.SealAllSegments(req.CollectionID)
S
sunby 已提交
522 523 524
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_SUCCESS,
	}, nil
S
sunby 已提交
525 526 527 528 529 530 531 532 533
}

func (s *Server) AssignSegmentID(req *datapb.AssignSegIDRequest) (*datapb.AssignSegIDResponse, error) {
	resp := &datapb.AssignSegIDResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_SUCCESS,
		},
		SegIDAssignments: make([]*datapb.SegIDAssignment, 0),
	}
S
sunby 已提交
534 535 536 537 538
	if !s.checkStateIsHealthy() {
		resp.Status.ErrorCode = commonpb.ErrorCode_UNEXPECTED_ERROR
		resp.Status.Reason = "server is initializing"
		return resp, nil
	}
S
sunby 已提交
539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582
	for _, r := range req.SegIDRequests {
		result := &datapb.SegIDAssignment{
			Status: &commonpb.Status{
				ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
			},
		}
		segmentID, retCount, expireTs, err := s.segAllocator.AllocSegment(r.CollectionID, r.PartitionID, r.ChannelName, int(r.Count))
		if err != nil {
			if _, ok := err.(errRemainInSufficient); !ok {
				result.Status.Reason = fmt.Sprintf("allocation of Collection %d, Partition %d, Channel %s, Count %d error:  %s",
					r.CollectionID, r.PartitionID, r.ChannelName, r.Count, err.Error())
				resp.SegIDAssignments = append(resp.SegIDAssignments, result)
				continue
			}

			if err = s.openNewSegment(r.CollectionID, r.PartitionID, r.ChannelName); err != nil {
				result.Status.Reason = fmt.Sprintf("open new segment of Collection %d, Partition %d, Channel %s, Count %d error:  %s",
					r.CollectionID, r.PartitionID, r.ChannelName, r.Count, err.Error())
				resp.SegIDAssignments = append(resp.SegIDAssignments, result)
				continue
			}

			segmentID, retCount, expireTs, err = s.segAllocator.AllocSegment(r.CollectionID, r.PartitionID, r.ChannelName, int(r.Count))
			if err != nil {
				result.Status.Reason = fmt.Sprintf("retry allocation of Collection %d, Partition %d, Channel %s, Count %d error:  %s",
					r.CollectionID, r.PartitionID, r.ChannelName, r.Count, err.Error())
				resp.SegIDAssignments = append(resp.SegIDAssignments, result)
				continue
			}
		}

		result.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
		result.CollectionID = r.CollectionID
		result.SegID = segmentID
		result.PartitionID = r.PartitionID
		result.Count = uint32(retCount)
		result.ExpireTime = expireTs
		result.ChannelName = r.ChannelName
		resp.SegIDAssignments = append(resp.SegIDAssignments, result)
	}
	return resp, nil
}

func (s *Server) openNewSegment(collectionID UniqueID, partitionID UniqueID, channelName string) error {
N
neza2017 已提交
583 584 585 586
	id, err := s.allocator.allocID()
	if err != nil {
		return err
	}
S
sunby 已提交
587
	segmentInfo, err := BuildSegment(collectionID, partitionID, id, channelName)
S
sunby 已提交
588 589 590 591 592 593
	if err != nil {
		return err
	}
	if err = s.meta.AddSegment(segmentInfo); err != nil {
		return err
	}
S
sunby 已提交
594
	if err = s.segAllocator.OpenSegment(segmentInfo); err != nil {
S
sunby 已提交
595 596
		return err
	}
597
	infoMsg := &msgstream.SegmentInfoMsg{
S
sunby 已提交
598 599 600
		BaseMsg: msgstream.BaseMsg{
			HashValues: []uint32{0},
		},
601 602 603 604
		SegmentMsg: datapb.SegmentMsg{
			Base: &commonpb.MsgBase{
				MsgType:   commonpb.MsgType_kSegmentInfo,
				MsgID:     0,
N
neza2017 已提交
605 606
				Timestamp: 0,
				SourceID:  Params.NodeID,
607 608 609 610
			},
			Segment: segmentInfo,
		},
	}
G
groot 已提交
611
	msgPack := &msgstream.MsgPack{
612 613 614 615 616
		Msgs: []msgstream.TsMsg{infoMsg},
	}
	if err = s.segmentInfoStream.Produce(msgPack); err != nil {
		return err
	}
S
sunby 已提交
617 618 619 620
	return nil
}

func (s *Server) ShowSegments(req *datapb.ShowSegmentRequest) (*datapb.ShowSegmentResponse, error) {
S
sunby 已提交
621 622 623 624 625
	resp := &datapb.ShowSegmentResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
S
sunby 已提交
626
	if !s.checkStateIsHealthy() {
S
sunby 已提交
627 628
		resp.Status.Reason = "server is initializing"
		return resp, nil
S
sunby 已提交
629
	}
630
	ids := s.meta.GetSegmentsOfPartition(req.CollectionID, req.PartitionID)
S
sunby 已提交
631 632 633
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	resp.SegmentIDs = ids
	return resp, nil
S
sunby 已提交
634 635 636 637 638 639 640 641
}

func (s *Server) GetSegmentStates(req *datapb.SegmentStatesRequest) (*datapb.SegmentStatesResponse, error) {
	resp := &datapb.SegmentStatesResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
S
sunby 已提交
642 643 644 645
	if !s.checkStateIsHealthy() {
		resp.Status.Reason = "server is initializing"
		return resp, nil
	}
S
sunby 已提交
646

Z
zhenshan.cao 已提交
647 648 649 650 651 652 653 654 655 656 657 658 659 660 661
	for _, segmentID := range req.SegmentIDs {
		state := &datapb.SegmentStateInfo{
			Status:    &commonpb.Status{},
			SegmentID: segmentID,
		}
		segmentInfo, err := s.meta.GetSegment(segmentID)
		if err != nil {
			state.Status.ErrorCode = commonpb.ErrorCode_UNEXPECTED_ERROR
			state.Status.Reason = "get segment states error: " + err.Error()
		} else {
			state.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
			state.State = segmentInfo.State
			state.CreateTime = segmentInfo.OpenTime
			state.SealedTime = segmentInfo.SealedTime
			state.FlushedTime = segmentInfo.FlushedTime
X
XuanYang-cn 已提交
662 663
			state.StartPosition = segmentInfo.StartPosition
			state.EndPosition = segmentInfo.EndPosition
Z
zhenshan.cao 已提交
664 665
		}
		resp.States = append(resp.States, state)
S
sunby 已提交
666
	}
S
sunby 已提交
667
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
Z
zhenshan.cao 已提交
668

S
sunby 已提交
669 670 671 672
	return resp, nil
}

func (s *Server) GetInsertBinlogPaths(req *datapb.InsertBinlogPathRequest) (*datapb.InsertBinlogPathsResponse, error) {
S
sunby 已提交
673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691
	resp := &datapb.InsertBinlogPathsResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
	p := path.Join(Params.SegmentFlushMetaPath, strconv.FormatInt(req.SegmentID, 10))
	value, err := s.client.Load(p)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	flushMeta := &datapb.SegmentFlushMeta{}
	err = proto.UnmarshalText(value, flushMeta)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	fields := make([]UniqueID, len(flushMeta.Fields))
	paths := make([]*internalpb2.StringList, len(flushMeta.Fields))
X
Xiangyu Wang 已提交
692 693 694
	for i, field := range flushMeta.Fields {
		fields[i] = field.FieldID
		paths[i] = &internalpb2.StringList{Values: field.BinlogPaths}
S
sunby 已提交
695
	}
S
sunby 已提交
696
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
S
sunby 已提交
697 698 699
	resp.FieldIDs = fields
	resp.Paths = paths
	return resp, nil
S
sunby 已提交
700 701
}

G
godchen 已提交
702 703 704 705 706 707 708
func (s *Server) GetInsertChannels(req *datapb.InsertChannelRequest) (*internalpb2.StringList, error) {
	return &internalpb2.StringList{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_SUCCESS,
		},
		Values: s.insertChannels,
	}, nil
S
sunby 已提交
709 710 711
}

func (s *Server) GetCollectionStatistics(req *datapb.CollectionStatsRequest) (*datapb.CollectionStatsResponse, error) {
712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730
	resp := &datapb.CollectionStatsResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
	nums, err := s.meta.GetNumRowsOfCollection(req.CollectionID)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	memsize, err := s.meta.GetMemSizeOfCollection(req.CollectionID)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	resp.Stats = append(resp.Stats, &commonpb.KeyValuePair{Key: "nums", Value: strconv.FormatInt(nums, 10)})
	resp.Stats = append(resp.Stats, &commonpb.KeyValuePair{Key: "memsize", Value: strconv.FormatInt(memsize, 10)})
	return resp, nil
S
sunby 已提交
731 732 733 734 735
}

func (s *Server) GetPartitionStatistics(req *datapb.PartitionStatsRequest) (*datapb.PartitionStatsResponse, error) {
	// todo implement
	return nil, nil
S
sunby 已提交
736
}
N
neza2017 已提交
737

G
godchen 已提交
738 739 740 741 742 743 744
func (s *Server) GetSegmentInfoChannel() (*milvuspb.StringResponse, error) {
	return &milvuspb.StringResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_SUCCESS,
		},
		Value: Params.SegmentInfoChannelName,
	}, nil
N
neza2017 已提交
745
}
746 747 748 749 750 751 752

func (s *Server) GetCount(req *datapb.CollectionCountRequest) (*datapb.CollectionCountResponse, error) {
	resp := &datapb.CollectionCountResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
X
XuanYang-cn 已提交
753 754 755 756
	if !s.checkStateIsHealthy() {
		resp.Status.Reason = "data service is not healthy"
		return resp, nil
	}
757 758 759 760 761 762 763 764 765
	nums, err := s.meta.GetNumRowsOfCollection(req.CollectionID)
	if err != nil {
		resp.Status.Reason = err.Error()
		return resp, nil
	}
	resp.Count = nums
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	return resp, nil
}
X
XuanYang-cn 已提交
766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789

func (s *Server) GetSegmentInfo(req *datapb.SegmentInfoRequest) (*datapb.SegmentInfoResponse, error) {
	resp := &datapb.SegmentInfoResponse{
		Status: &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
		},
	}
	if !s.checkStateIsHealthy() {
		resp.Status.Reason = "data service is not healthy"
		return resp, nil
	}
	infos := make([]*datapb.SegmentInfo, len(req.SegmentIDs))
	for i, id := range req.SegmentIDs {
		segmentInfo, err := s.meta.GetSegment(id)
		if err != nil {
			resp.Status.Reason = err.Error()
			return resp, nil
		}
		infos[i] = segmentInfo
	}
	resp.Status.ErrorCode = commonpb.ErrorCode_SUCCESS
	resp.Infos = infos
	return resp, nil
}