Documentation
¶
Index ¶
- Constants
- func InvokeHandler(conn IConn, packet IPacket, handler *Handler)
- func RegisterMessage(msg interface{}, key interface{})
- type AttrKey
- type AttrMap
- type Buffer
- func (self *Buffer) Append(data []byte)
- func (self *Buffer) Bytes() []byte
- func (self *Buffer) Clear()
- func (self *Buffer) Concat()
- func (self *Buffer) Discard()
- func (self *Buffer) Empty() bool
- func (self *Buffer) Eof() bool
- func (self *Buffer) ExtendSpace(count int) error
- func (self *Buffer) Front() *list.Element
- func (self *Buffer) GetMark() int
- func (self *Buffer) GetSpace() []byte
- func (self *Buffer) HasRemain(count int) bool
- func (self *Buffer) IndexOf(key string) int
- func (self *Buffer) IndexOfLimit(key string, limit int) int
- func (self *Buffer) Length() int
- func (self *Buffer) Merge(other *Buffer)
- func (self *Buffer) Peek(buffer []byte) (int, error)
- func (self *Buffer) Position() int
- func (self *Buffer) Prepend(data []byte)
- func (self *Buffer) Read(buffer []byte) (int, error)
- func (buffer *Buffer) ReadAll(reader io.Reader) error
- func (self *Buffer) ReadByte() (byte, error)
- func (self *Buffer) ReadLine() (*Buffer, error)
- func (self *Buffer) ReadToEnd() []byte
- func (self *Buffer) ReadUntil(key byte) (string, error)
- func (self *Buffer) Rewind()
- func (self *Buffer) Seek(offset int, whence int) error
- func (self *Buffer) SetMark(value int)
- func (self *Buffer) Split(result *Buffer)
- func (self *Buffer) String() string
- func (self *Buffer) Swap(other *Buffer)
- func (self *Buffer) Visit(cb func([]byte) bool)
- func (self *Buffer) Write(bufffer []byte) (int, error)
- func (buffer *Buffer) WriteAll(writer io.Writer) error
- type Callback
- type Dispatcher
- func (d *Dispatcher) GetHandler(id uint, name string) *Handler
- func (d *Dispatcher) GetHandlerById(id uint) *Handler
- func (d *Dispatcher) GetHandlerByName(name string) *Handler
- func (d *Dispatcher) Middlewares() []HandlerCB
- func (d *Dispatcher) Register(key interface{}, cb HandlerCB, queueId uint) *Handler
- func (d *Dispatcher) Use(cb HandlerCB)
- type Executor
- type Handler
- type HandlerCB
- type HandlerCtx
- func (ctx *HandlerCtx) Add(v interface{})
- func (ctx *HandlerCtx) Get(index int) interface{}
- func (ctx *HandlerCtx) Init(conn IConn, packet IPacket, handler *Handler, chain []HandlerCB)
- func (ctx *HandlerCtx) Len() int
- func (ctx *HandlerCtx) Message() interface{}
- func (ctx *HandlerCtx) MsgID() uint
- func (ctx *HandlerCtx) MsgName() string
- func (ctx *HandlerCtx) Next()
- func (ctx *HandlerCtx) Packet() IPacket
- type ICodec
- type IConn
- type IFilter
- type IFilterChain
- type IFilterCtx
- type IFrame
- type IIdentity
- type IPacket
- type ITran
- type Iterator
- type Option
- type ReconnectOption
- type Registry
- type ReverseIterator
- type SyncOption
- type TagOption
Constants ¶
const ( // AttrKindConf 配置使用 AttrKindConf = "conf" // AttrKindConn Conn中存储数据 AttrKindConn = "conn" // AttrKindCtx filterContext中使用 AttrKindCtx = "context" )
const QueueMainID = 0
QueueMainID 主队列ID
Variables ¶
This section is empty.
Functions ¶
func InvokeHandler ¶
InvokeHandler 调用Handler
Types ¶
type AttrKey ¶
type AttrKey struct {
// contains filtered or unexported fields
}
AttrKey 将string延迟映射到唯一索引,不同进程间,索引并不一定一样
type AttrMap ¶
type AttrMap interface {
HasAttr(key *AttrKey) bool
SetAttr(key *AttrKey, val interface{})
GetAttr(key *AttrKey) interface{}
GetAttrEx(key *AttrKey, defVal interface{}) interface{}
}
AttrMap 额外数据
type Buffer ¶
type Buffer struct {
// contains filtered or unexported fields
}
func (*Buffer) ExtendSpace ¶
ExtendSpace 将最后一个扩展count个字节,配合GetSpace使用
func (*Buffer) IndexOfLimit ¶
IndexOfLimit 从当前位置查找key,最长搜索limit个字节(-1无限制)
type Dispatcher ¶
type Dispatcher struct {
// contains filtered or unexported fields
}
Dispatcher 消息回调管理,id不能超过uint16 支持添加中间件
func GetDispatcher ¶
func GetDispatcher() *Dispatcher
func NewDispatcher ¶
func NewDispatcher() *Dispatcher
func (*Dispatcher) GetHandler ¶
func (d *Dispatcher) GetHandler(id uint, name string) *Handler
func (*Dispatcher) GetHandlerById ¶
func (d *Dispatcher) GetHandlerById(id uint) *Handler
func (*Dispatcher) GetHandlerByName ¶
func (d *Dispatcher) GetHandlerByName(name string) *Handler
func (*Dispatcher) Middlewares ¶
func (d *Dispatcher) Middlewares() []HandlerCB
func (*Dispatcher) Register ¶
func (d *Dispatcher) Register(key interface{}, cb HandlerCB, queueId uint) *Handler
Register 注册消息回调,可以指定队列
func (*Dispatcher) Use ¶
func (d *Dispatcher) Use(cb HandlerCB)
type Executor ¶
type Executor struct {
// contains filtered or unexported fields
}
Executor 任务调度器 一个Main,多个Work队列,main和work队列不会同时执行消息 如果不需要Main与Work互斥,则可以将其中一个work队列作为主循环
type HandlerCtx ¶
type HandlerCtx struct {
IConn
// contains filtered or unexported fields
}
HandlerCtx 支持中间件,在Dispatcher中注册
func (*HandlerCtx) Add ¶
func (ctx *HandlerCtx) Add(v interface{})
func (*HandlerCtx) Get ¶
func (ctx *HandlerCtx) Get(index int) interface{}
func (*HandlerCtx) Init ¶
func (ctx *HandlerCtx) Init(conn IConn, packet IPacket, handler *Handler, chain []HandlerCB)
func (*HandlerCtx) Len ¶
func (ctx *HandlerCtx) Len() int
func (*HandlerCtx) Message ¶
func (ctx *HandlerCtx) Message() interface{}
func (*HandlerCtx) MsgID ¶
func (ctx *HandlerCtx) MsgID() uint
func (*HandlerCtx) MsgName ¶
func (ctx *HandlerCtx) MsgName() string
func (*HandlerCtx) Next ¶
func (ctx *HandlerCtx) Next()
func (*HandlerCtx) Packet ¶
func (ctx *HandlerCtx) Packet() IPacket
type ICodec ¶
type ICodec interface {
Encode(buffer *Buffer, msg interface{}) error
Decode(buffer *Buffer, msg interface{}) error
}
ICodec 用于消息的序列化
type IConn ¶
type IConn interface {
AttrMap
GetId() uint
GetTag() string
SetTag(tag string)
GetData() interface{}
SetData(data interface{})
IsActive() bool // 连接是否正常
IsConnector() bool // 是否通过调用Connect产生,否则Listen产生
LocalAddr() net.Addr // 本地地址
RemoteAddr() net.Addr // 远程地址
Close() // 异步关闭,会等待数据发送完
Read() *Buffer // 异步读缓存
Write(buffer *Buffer) error // 异步写数据
Send(msg interface{}) error // 异步发消息
}
IConn asynchronous connection,
type IFilter ¶
type IFilter interface {
Name() string // unique name for filter
HandleRead(ctx IFilterCtx) // read data
HandleWrite(ctx IFilterCtx) // send data
HandleOpen(ctx IFilterCtx) // open by connect or listen
HandleClose(ctx IFilterCtx)
HandleError(ctx IFilterCtx)
}
IFilter must be stateless,if need data, can get from Conn or FilteCtx
type IFilterChain ¶
type IFilterChain interface {
Len() int
IndexOf(name string) int // 通过名字查询索引
AddFirst(filters ...IFilter) // 前边插入,Prepend
AddLast(filters ...IFilter) // 后边插入,Append
HandleOpen(conn IConn)
HandleClose(conn IConn)
HandleRead(conn IConn)
HandleWrite(conn IConn, msg interface{})
HandleError(conn IConn, err error)
}
IFilterChain 递归执行每一个Filter
type IFilterCtx ¶
type IFilterCtx interface {
AttrMap
GetConn() IConn
SetData(data interface{}) // 用于Filter间透传数据
GetData() interface{} // 获取数据
Error(err error) // 抛出错误
Next() // 执行下一个
Jump(index int) error // 跳转到某个索引,可以负索引
JumpBy(name string) error // 通过名字跳转
}
IFilterCtx filter上下文
type IIdentity ¶
type IIdentity interface {
Encode(buffer *Buffer, pkt IPacket) error
Decode(buffer *Buffer) (IPacket, error)
}
IIdentity 用于消息头的序列化,并创建相应Packet,注:无需创建Message
type IPacket ¶
type IPacket interface {
GetId() uint
SetId(id uint)
GetName() string
SetName(name string)
GetMessage() interface{}
SetMessage(msg interface{})
}
IPacket 用于定义消息包,分为消息头和消息体 消息头:Id,Name,用于唯一标识消息,Id为0时使用Name 消息体:可以是string,json,protobu编码,具体编解码由Codec实现
type ITran ¶
type ITran interface {
GetChain() IFilterChain
SetOptions(option ...Option)
AddFilters(filters ...IFilter)
Listen(host string, options ...Option) error // 支持TagOption
Connect(host string, options ...Option) error // 支持TagOption,SyncOption,ReconnectOption
Start()
Stop()
}
ITran Transport,用于创建Connection
type Iterator ¶
type Iterator struct {
// contains filtered or unexported fields
}
////////////////////////////////////////////////////// 迭代器实现 //////////////////////////////////////////////////////
type ReconnectOption ¶
ReconnectOption 重连配置
func WithCloseReconnectOption ¶
func WithCloseReconnectOption() *ReconnectOption
WithCloseReconnectOption 关闭断线重现
func WithReconnectOption ¶
func WithReconnectOption(count, interval int) *ReconnectOption
type Registry ¶
type Registry struct {
// contains filtered or unexported fields
}
Registry 用于注册消息的id和name
type ReverseIterator ¶
type ReverseIterator struct {
// contains filtered or unexported fields
}
////////////////////////////////////////////////////// 反向迭代器 //////////////////////////////////////////////////////
func (*ReverseIterator) Create ¶
func (self *ReverseIterator) Create(elem *list.Element, offset int, length int)
func (*ReverseIterator) MoveEnd ¶
func (self *ReverseIterator) MoveEnd()
func (*ReverseIterator) Next ¶
func (self *ReverseIterator) Next() bool
type SyncOption ¶
type SyncOption struct {
Flag bool
}
SyncOption 阻塞调用
func WithSyncOption ¶
func WithSyncOption(flag bool) *SyncOption

