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

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

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

S
sunby 已提交
16 17
	"github.com/golang/protobuf/proto"
	"github.com/zilliztech/milvus-distributed/internal/proto/masterpb"
S
sunby 已提交
18 19 20

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

S
sunby 已提交
21 22 23 24 25 26 27 28 29 30 31 32 33
	"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 已提交
34 35
const role = "dataservice"

S
sunby 已提交
36 37
type DataService interface {
	typeutil.Service
N
neza2017 已提交
38
	typeutil.Component
S
sunby 已提交
39 40 41 42 43 44 45
	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 已提交
46 47
	GetSegmentInfoChannel() (*milvuspb.StringResponse, error)
	GetInsertChannels(req *datapb.InsertChannelRequest) (*internalpb2.StringList, error)
S
sunby 已提交
48 49 50
	GetCollectionStatistics(req *datapb.CollectionStatsRequest) (*datapb.CollectionStatsResponse, error)
	GetPartitionStatistics(req *datapb.PartitionStatsRequest) (*datapb.PartitionStatsResponse, error)
	GetComponentStates() (*internalpb2.ComponentStates, error)
51
	GetCount(req *datapb.CollectionCountRequest) (*datapb.CollectionCountResponse, error)
X
XuanYang-cn 已提交
52
	GetSegmentInfo(req *datapb.SegmentInfoRequest) (*datapb.SegmentInfoResponse, error)
S
sunby 已提交
53 54
}

S
sunby 已提交
55 56 57 58 59 60 61 62 63 64
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 已提交
65 66 67 68 69 70 71
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 已提交
72
type (
S
sunby 已提交
73 74 75
	UniqueID  = typeutil.UniqueID
	Timestamp = typeutil.Timestamp
	Server    struct {
N
neza2017 已提交
76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94
		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
95
		insertChannels    []string
G
groot 已提交
96
		msFactory         msgstream.Factory
C
cai.zhang 已提交
97
		ttBarrier         timesync.TimeTickBarrier
S
sunby 已提交
98 99 100
	}
)

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

115 116 117 118 119 120 121 122 123
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 已提交
124 125
func (s *Server) SetMasterClient(masterClient MasterClient) {
	s.masterClient = masterClient
S
sunby 已提交
126 127 128
}

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

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

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

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

S
sunby 已提交
169 170 171 172 173 174 175
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 已提交
176
	s.meta, err = newMeta(etcdKV)
S
sunby 已提交
177 178 179 180 181 182
	if err != nil {
		return err
	}
	return nil
}

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

S
sunby 已提交
212 213
func (s *Server) loadMetaFromMaster() error {
	log.Println("loading collection meta from master")
S
sunby 已提交
214 215 216
	if err := s.checkMasterIsHealthy(); err != nil {
		return err
	}
N
neza2017 已提交
217 218 219 220 221 222 223
	if s.ddChannelName == "" {
		channel, err := s.masterClient.GetDdChannel()
		if err != nil {
			return err
		}
		s.ddChannelName = channel
	}
S
sunby 已提交
224 225 226 227 228
	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 已提交
229
			SourceID:  Params.NodeID,
S
sunby 已提交
230 231 232 233 234 235 236 237 238 239 240 241
		},
		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 已提交
242
				SourceID:  Params.NodeID,
S
sunby 已提交
243 244 245 246 247 248 249 250 251 252 253 254 255
			},
			DbName:         "",
			CollectionName: collectionName,
		})
		if err != nil {
			log.Println(err.Error())
			continue
		}
		partitions, err := s.masterClient.ShowPartitions(&milvuspb.ShowPartitionRequest{
			Base: &commonpb.MsgBase{
				MsgType:   commonpb.MsgType_kShowPartitions,
				MsgID:     -1, // todo
				Timestamp: 0,  // todo
S
sunby 已提交
256
				SourceID:  Params.NodeID,
S
sunby 已提交
257 258 259 260 261 262 263 264 265 266 267 268
			},
			DbName:         "",
			CollectionName: collectionName,
			CollectionID:   collection.CollectionID,
		})
		if err != nil {
			log.Println(err.Error())
			continue
		}
		err = s.meta.AddCollection(&collectionInfo{
			ID:         collection.CollectionID,
			Schema:     collection.Schema,
S
sunby 已提交
269
			Partitions: partitions.PartitionIDs,
S
sunby 已提交
270 271 272 273 274 275 276 277
		})
		if err != nil {
			log.Println(err.Error())
			continue
		}
	}
	log.Println("load collection meta from master complete")
	return nil
S
sunby 已提交
278
}
S
sunby 已提交
279 280 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

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
}

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

func (s *Server) startStatsChannel(ctx context.Context) {
	defer s.serverLoopWg.Done()
G
groot 已提交
319
	statsStream, _ := s.msFactory.NewMsgStream(ctx)
Z
zhenshan.cao 已提交
320
	statsStream.AsConsumer([]string{Params.StatisticsChannelName}, Params.DataServiceSubscriptionName)
321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343
	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 {
					log.Println(err.Error())
					continue
				}
			}
		}
	}
}

func (s *Server) startSegmentFlushChannel(ctx context.Context) {
	defer s.serverLoopWg.Done()
G
groot 已提交
344
	flushStream, _ := s.msFactory.NewMsgStream(ctx)
Z
zhenshan.cao 已提交
345
	flushStream.AsConsumer([]string{Params.SegmentInfoChannelName}, Params.DataServiceSubscriptionName)
346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367
	flushStream.Start()
	defer flushStream.Close()
	for {
		select {
		case <-ctx.Done():
			log.Println("segment flush channel shut down")
			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 {
				log.Println(err.Error())
				continue
			}
			segmentInfo.FlushedTime = realMsg.BeginTimestamp
S
sunby 已提交
368
			segmentInfo.State = commonpb.SegmentState_SegmentFlushed
369 370 371 372 373 374 375 376
			if err = s.meta.UpdateSegment(segmentInfo); err != nil {
				log.Println(err.Error())
				continue
			}
		}
	}
}

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

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

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

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

S
sunby 已提交
422
func (s *Server) GetComponentStates() (*internalpb2.ComponentStates, error) {
S
sunby 已提交
423 424 425 426
	resp := &internalpb2.ComponentStates{
		State: &internalpb2.ComponentInfo{
			NodeID:    Params.NodeID,
			Role:      role,
S
sunby 已提交
427
			StateCode: s.state.Load().(internalpb2.StateCode),
S
sunby 已提交
428 429 430 431 432 433 434 435 436 437 438 439 440
		},
		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 已提交
441 442
}

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

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

func (s *Server) RegisterNode(req *datapb.RegisterNodeRequest) (*datapb.RegisterNodeResponse, error) {
S
sunby 已提交
462
	ret := &datapb.RegisterNodeResponse{
S
sunby 已提交
463
		Status: &commonpb.Status{
S
sunby 已提交
464
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
S
sunby 已提交
465
		},
S
sunby 已提交
466
	}
S
sunby 已提交
467 468 469 470 471
	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 已提交
472
	if s.ddChannelName == "" {
N
neza2017 已提交
473
		resp, err := s.masterClient.GetDdChannel()
S
sunby 已提交
474 475 476 477
		if err != nil {
			ret.Status.Reason = err.Error()
			return ret, err
		}
N
neza2017 已提交
478
		s.ddChannelName = resp
S
sunby 已提交
479 480 481 482 483 484 485 486
	}
	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 已提交
487
			{Key: "CompleteFlushChannelName", Value: Params.SegmentInfoChannelName},
S
sunby 已提交
488 489 490
		},
	}
	return ret, nil
S
sunby 已提交
491 492
}

S
sunby 已提交
493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511
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 已提交
512
func (s *Server) Flush(req *datapb.FlushRequest) (*commonpb.Status, error) {
S
sunby 已提交
513 514 515 516 517 518
	if !s.checkStateIsHealthy() {
		return &commonpb.Status{
			ErrorCode: commonpb.ErrorCode_UNEXPECTED_ERROR,
			Reason:    "server is initializing",
		}, nil
	}
N
neza2017 已提交
519
	s.segAllocator.SealAllSegments(req.CollectionID)
S
sunby 已提交
520 521 522
	return &commonpb.Status{
		ErrorCode: commonpb.ErrorCode_SUCCESS,
	}, nil
S
sunby 已提交
523 524 525 526 527 528 529 530 531
}

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 已提交
532 533 534 535 536
	if !s.checkStateIsHealthy() {
		resp.Status.ErrorCode = commonpb.ErrorCode_UNEXPECTED_ERROR
		resp.Status.Reason = "server is initializing"
		return resp, nil
	}
S
sunby 已提交
537 538 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
			}

			log.Printf("no enough space for allocation of Collection %d, Partition %d, Channel %s, Count %d",
				r.CollectionID, r.PartitionID, r.ChannelName, r.Count)
			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
}