eventbus

package module
v1.0.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 16, 2021 License: Apache-2.0 Imports: 21 Imported by: 0

README

概述

事件总线。使用事件驱动的方式进行业务解耦。
建议使用集群的方式运行,这样可以支持编程多语言环境。

获取

go get github.com/aacfactory/eventbus

使用

本地

// 构建
eb := eventbus.NewEventbus()
// 挂载处理器

type Arg struct {
    Id       string    `json:"id,omitempty"`
    Num      int       `json:"num,omitempty"`
    Datetime time.Time `json:"datetime,omitempty"`
}

type Result struct {
    Value string `json:"value,omitempty"`
}

func HandlerReply(event eventbus.Event) (result interface{}, err error) {
    arg := &Arg{}
    _ = json.Unmarshal(event.Body(), arg)
    fmt.Println("handle reply", event.Head(), arg)
    if arg.Num < 0 {
        err = errors.InvalidArgumentErrorWithDetails("bad number", "num", "less than 0")
        return
    }
    result = &Result{
        Value: "result",
    }
    return
}

func HandlerVoid(event eventbus.Event) (result interface{}, err error) {
    arg := &Arg{}
    _ = json.Unmarshal(event.Body(), arg)
    fmt.Println("handle void", event.Head(), arg)
    return
}

_ = eb.RegisterHandler("void", HandlerVoid)
_ = eb.RegisterHandler("reply", HandlerReply)
// 启动 (启动必须晚于挂载处理器)
eb.Start(context.TODO())
// 执行
options := eventbus.NewDeliveryOptions()
options.Add("h1", "1")
options.Add("h2", "2")

sendErr := eb.Send("reply", &Arg{
    Id:       "id",
    Num:      10,
    Datetime: time.Now(),
}, options)

if sendErr != nil {
fmt.Println("send failed", sendErr)
}

for i := 0; i < 2; i++ {
    rf := eb.Request("reply", &Arg{
        Id:       "id",
        Num:      i - 1,
        Datetime: time.Now(),
    }, options)
    result := &Result{}
    requestErr := rf.Get(result)
    if requestErr != nil {
    	fmt.Println("request failed", requestErr)
    } else {
        fmt.Println("request succeed", result)
    }
}
// 优雅的关闭
eb.Close(context.TODO())

集群

TCP + 地址发现 模式

当前没有实现服务发现,请使用 aacfactory/cluster ,或自行实现。

// options
options := eventbus.ClusterEventbusOption{
    Host:                       "0.0.0.0", // 实际监听地址
    Port:                       9090, // 实际监听端口
    PublicHost:                 "127.0.0.1", // 注册地址,如果为空,则默认使用监听地址
    PublicPort:                 0, // 注册端口,如果为空,则默认使用监听端口
    Meta:                       &eventbus.EndpointMeta{}, // 注册源数据
    Tags:                       nil, // 标签,一般用于版本化与运行隔离化
    TLS:                        &eventbus.EndpointTLS{}, // TLS 配置
    EventChanCap:               64, // 事件 chan 的长度
    EventHandlerInstanceNumber: 2,  // 事件处理器的实例个数
    EnableLocal:                true, // 是否开启本地处理器,如果关闭,则直接使用远程模式,建议开启
}
// discovery
discovery := Foo{}
// 创建
bus, err = eventbus.NewClusterEventbus(discovery, options)
if err != nil {
    return
}
// 操作与本地Eventbus一样


NATS 模型
// todo

Documentation

Index

Constants

View Source
const (
	EndpointStatusRunning = EndpointStatus("RUNNING")
	EndpointStatusClosing = EndpointStatus("CLOSING")
)

Variables

This section is empty.

Functions

func NewEventbusWithOption

func NewEventbusWithOption(option LocaledEventbusOption) (eb *localedEventbus)

Types

type ClusterEventbusOption

type ClusterEventbusOption struct {
	Host         string        `json:"host,omitempty"`
	Port         int           `json:"port,omitempty"`
	PublicHost   string        `json:"publicHost,omitempty"`
	PublicPort   int           `json:"publicPort,omitempty"`
	Meta         *EndpointMeta `json:"meta,omitempty"`
	Tags         []string      `json:"tags,omitempty"`
	TLS          *EndpointTLS  `json:"tls,omitempty"`
	EventChanCap int           `json:"eventChanCap,omitempty"`
	EventWorkers int           `json:"eventWorkers,omitempty"`
}

type DeliveryOptions

type DeliveryOptions interface {
	Add(key string, value string)
	Put(key string, value []string)
	Get(key string) (string, bool)
	Keys() []string
	Empty() bool
	Values(key string) ([]string, bool)
	Remove(key string)
	AddTag(tags ...string)
}

func NewDeliveryOptions

func NewDeliveryOptions() DeliveryOptions

type EndpointMeta

type EndpointMeta map[string]string

func NewEndpointMeta

func NewEndpointMeta() EndpointMeta

func (EndpointMeta) Empty

func (meta EndpointMeta) Empty() bool

func (EndpointMeta) Get

func (meta EndpointMeta) Get(key string) (string, bool)

func (EndpointMeta) Keys

func (meta EndpointMeta) Keys() []string

func (EndpointMeta) Merge

func (meta EndpointMeta) Merge(o ...Meta)

func (EndpointMeta) Put

func (meta EndpointMeta) Put(key string, value string)

func (EndpointMeta) Rem

func (meta EndpointMeta) Rem(key string)

type EndpointStatus

type EndpointStatus string

func (EndpointStatus) Closing

func (s EndpointStatus) Closing() bool

func (EndpointStatus) Ok

func (s EndpointStatus) Ok() bool

type EndpointTLS

type EndpointTLS struct {
	Enable_     bool   `json:"enable,omitempty"`
	VerifySSL_  bool   `json:"verifySsl,omitempty"`
	CA_         string `json:"ca,omitempty"`
	ServerCert_ string `json:"serverCert,omitempty"`
	ServerKey_  string `json:"serverKey,omitempty"`
	ClientCert_ string `json:"clientCert,omitempty"`
	ClientKey_  string `json:"clientKey,omitempty"`
}

func (EndpointTLS) CA

func (s EndpointTLS) CA() string

func (EndpointTLS) ClientCert

func (s EndpointTLS) ClientCert() string

func (EndpointTLS) ClientKey

func (s EndpointTLS) ClientKey() string

func (EndpointTLS) Enable

func (s EndpointTLS) Enable() bool

func (EndpointTLS) ServerCert

func (s EndpointTLS) ServerCert() string

func (EndpointTLS) ServerKey

func (s EndpointTLS) ServerKey() string

func (EndpointTLS) ToClientTLSConfig

func (s EndpointTLS) ToClientTLSConfig() (config *tls.Config, err error)

func (EndpointTLS) ToServerTLSConfig

func (s EndpointTLS) ToServerTLSConfig() (config *tls.Config, err error)

func (EndpointTLS) VerifySSL

func (s EndpointTLS) VerifySSL() bool

type Event added in v1.0.1

type Event interface {
	Head() EventHead
	Body() []byte
}

type EventHandler

type EventHandler func(event Event) (result interface{}, err error)

type EventHead added in v1.0.1

type EventHead interface {
	Add(key string, value string)
	Put(key string, value []string)
	Get(key string) (string, bool)
	Keys() []string
	Empty() bool
	Values(key string) ([]string, bool)
	Remove(key string)
}

type Eventbus

type Eventbus interface {
	Send(address string, v interface{}, options ...DeliveryOptions) (err error)
	Request(address string, v interface{}, options ...DeliveryOptions) (reply ReplyFuture)
	RegisterHandler(address string, handler EventHandler, tags ...string) (err error)
	RegisterLocalHandler(address string, handler EventHandler, tags ...string) (err error)
	Start(context context.Context)
	Close(context context.Context)
}

func NewClusterEventbus

func NewClusterEventbus(discovery ServiceDiscovery, option ClusterEventbusOption) (bus Eventbus, err error)

func NewEventbus

func NewEventbus() Eventbus

type LocaledEventbusOption

type LocaledEventbusOption struct {
	EventChanCap int `json:"eventChanCap,omitempty"`
	EventWorkers int `json:"eventWorkers,omitempty"`
}

type Meta

type Meta interface {
	Put(key string, value string)
	Get(key string) (value string, has bool)
	Rem(key string)
	Keys() (keys []string)
	Empty() (ok bool)
	Merge(o ...Meta)
}

type MultiMap

type MultiMap map[string][]string

func (MultiMap) Add

func (h MultiMap) Add(key string, value string)

func (MultiMap) Empty

func (h MultiMap) Empty() bool

func (MultiMap) Get

func (h MultiMap) Get(key string) (string, bool)

func (MultiMap) Keys

func (h MultiMap) Keys() []string

func (MultiMap) Merge

func (h MultiMap) Merge(o ...MultiMap)

func (MultiMap) Put

func (h MultiMap) Put(key string, value []string)

func (MultiMap) Remove

func (h MultiMap) Remove(key string)

func (MultiMap) Values

func (h MultiMap) Values(key string) ([]string, bool)

type Registration

type Registration interface {
	NodeId() (nodeId string)
	NodeName() (nodeName string)
	Id() (id string)
	Group() (group string)
	Name() (name string)
	Status() (status Status)
	Protocol() (protocol string)
	Address() (address string)
	Tags() (tags []string)
	Meta() (meta Meta)
	TLS() (registrationTLS RegistrationTLS)
}

type RegistrationTLS

type RegistrationTLS interface {
	Enable() bool
	VerifySSL() bool
	CA() string
	ServerCert() string
	ServerKey() string
	ClientCert() string
	ClientKey() string
	ToServerTLSConfig() (config *tls.Config, err error)
	ToClientTLSConfig() (config *tls.Config, err error)
}

type ReplyError

type ReplyError struct {
	Id          string          `json:"id,omitempty"`
	FailureCode int             `json:"failureCode,omitempty"`
	Code        string          `json:"code,omitempty"`
	Message     string          `json:"message,omitempty"`
	Meta        errors.MultiMap `json:"meta,omitempty"`
}

func (*ReplyError) Error

func (e *ReplyError) Error() string

func (*ReplyError) GetMeta

func (e *ReplyError) GetMeta() errors.MultiMap

func (*ReplyError) GetStacktrace

func (e *ReplyError) GetStacktrace() (fn string, file string, line int)

func (*ReplyError) SetFailureCode

func (e *ReplyError) SetFailureCode(failureCode int) errors.CodeError

func (*ReplyError) SetId

func (e *ReplyError) SetId(id string) errors.CodeError

func (*ReplyError) String

func (e *ReplyError) String() string

func (*ReplyError) ToJson

func (e *ReplyError) ToJson() []byte

type ReplyFuture

type ReplyFuture interface {
	Get(v interface{}) (err error)
}

type ServiceDiscovery

type ServiceDiscovery interface {
	Publish(group string, name string, protocol string, address string, tags []string, meta Meta, registrationTLS RegistrationTLS) (registration Registration, err error)
	UnPublish(registration Registration) (err error)
	Get(group string, name string, tags ...string) (registration Registration, has bool, err error)
	GetALL(group string, name string, tags ...string) (registrations []Registration, has bool, err error)
}

type Status

type Status interface {
	Ok() bool
	Closing() bool
}

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL