aboutsummaryrefslogtreecommitdiff
path: root/weed/mq/client/publisher.go
blob: 826947721156a788a04ec8834988a0bd2deb8572 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
package client

import (
	"github.com/seaweedfs/seaweedfs/weed/mq/messages"
	"github.com/seaweedfs/seaweedfs/weed/pb"
	"time"
)

type PublishProcessor interface {
	AddMessage(m *messages.Message) error
	Shutdown() error
}

type PublisherOption struct {
	Masters string
	Topic   string
}

type Publisher struct {
	option    *PublisherOption
	masters   []pb.ServerAddress
	processor *PublishStreamProcessor
}

func NewPublisher(option *PublisherOption) *Publisher {
	p := &Publisher{
		masters:   pb.ServerAddresses(option.Masters).ToAddresses(),
		option:    option,
		processor: NewPublishStreamProcessor(3, 887*time.Millisecond),
	}
	return p
}

func (p Publisher) Publish(m *messages.Message) error {
	return p.processor.AddMessage(m)
}

func (p Publisher) Shutdown() error {
	return p.processor.Shutdown()
}