fairy

package module
v0.0.0-...-e07b629 Latest Latest
Warning

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

Go to latest
Published: Dec 3, 2018 License: GPL-3.0 Imports: 10 Imported by: 0

README

Fairy library(WIP)

目标:高效,灵活,易用,易扩展的异步网络框架,设计上参考netty,mina,grizzly,使用责任链设计模式

  • 支持tcp,websocket,kcp协议
  • 支持protobuf,json,xml,gob编码
  • 支持默认的消息处理线程模型
  • 支持高效定时器

一:用例

package chat

import (
	"github.com/jeckbjy/fairy"
	"github.com/jeckbjy/fairy/codecs"
	"github.com/jeckbjy/fairy/filters"
	"github.com/jeckbjy/fairy/frames"
	"github.com/jeckbjy/fairy/identities"
	"github.com/jeckbjy/fairy/log"
	"github.com/jeckbjy/fairy/tcp"
	"github.com/jeckbjy/fairy/timer"
	"github.com/jeckbjy/fairy/util"
)

type ChatMsg struct {
	Content   string
	Timestamp int64
}

func StartServer() {
	log.Debug("start server")
	// step1: register message
	fairy.RegisterMessage(&ChatMsg{}, nil)

	// step2: register handler
	fairy.RegisterHandler(&ChatMsg{}, func(ctx *fairy.HandlerCtx) {
		req := ctx.Message().(*ChatMsg)
		log.Debug("client msg:%+v", req)

		rsp := &ChatMsg{}
		rsp.Content = "welcome boy!"
		rsp.Timestamp = util.Now()
		ctx.Send(rsp)
	})

	// step3: create transport and add filters
	tran := tcp.NewTran()
	tran.AddFilters(
		filters.NewLogging(),
		filters.NewFrame(frames.NewLine()),
		filters.NewPacket(identities.NewString(), codecs.NewJson()),
		filters.NewExecutor())

	// step4: listen or connect
	tran.Listen(":8080", 0)
}

func StartClient() {
	log.Debug("start client")
	// step1: register message
	fairy.RegisterMessage(&ChatMsg{}, nil)

	// step2: register handler
	fairy.RegisterHandler(&ChatMsg{}, func(ctx *fairy.HandlerCtx) {
		req := ctx.Message().(*ChatMsg)
		log.Debug("server msg:%+v", req)
	})

	var gConn fairy.IConn
	// step3: create transport and add filters
	tran := tcp.NewTran()
	tran.AddFilters(
		filters.NewLogging(),
		filters.NewFrame(frames.NewLine()),
		filters.NewPacket(identities.NewString(), codecs.NewJson()),
		filters.NewExecutor())

	tran.AddFilters(filters.NewConnect(func(conn fairy.IConn) {
		// send msg to server
		req := &ChatMsg{}
		req.Content = "hello word!"
		conn.Send(req)
		gConn = conn
	}))

	// add timer for send message
	timer.Start(timer.ModeLoop, 1000, func() {
		log.Debug("Ontimeout")
		req := &ChatMsg{}
		req.Content = "hello word!"
		req.Timestamp = util.Now()
		gConn.Send(req)
	})

	// step4: listen or connect
	tran.Connect("localhost:8080", 0)
}

二:一些建议

  • 服务器集群,对于复杂的服务器架构,直接使用默认的消息编码并不能满足需求,通常需要自定义IPacket和IIdentity来扩展,比如增加uid,消息源,目标类型等
  • PacketFilter在某些情况下并不是高效的,因为里边进行了Codec的编解码,如果仅仅是转发协议,则并不需要解析body数据,可以自定义Filter,通过判断是否有消息处理回调判断是否需要进行body解析
  • rpc调用,本库并没有直接支持,如果需要,可以自定义Packet,增加一个唯一rpc id,Call时报错id到回调的映射,在消息处理处判断rpc id是否存在回调,如果存在则直接调用。额外需要一个定时器做延迟判断,防止消息永远没有返回,永远不能被执行

三:原理

  • Transport和Connection
    • Transport:主要提供Listen和Connect两个接口,用于创建Connection,Connection默认会自动断线重连,如果不需要断线重连,可以通过SetOption关闭
    • Connection:类似于net.Conn,主要提供异步Read,Write,Close等接口

Tran和Conn

  • Filter
    • Filter 提供InBound和OutBound两种流向
      • InBound: HandleRead,HandleOpen,HandleError
      • OutBound:HandleWrite,HandleClose
    • FilterCtx 用于Filter之间数据传递,最常用的函数:GetData和SetData用于消息编解码,透传消息
    • 内置的filters
      • FrameFilter,PacketFilter,ExecutorFilter,LoggingFilter,TelnetFilter,ConnectFilter,RC4Filter
      • 自定义filter
        • filter应该是一个无状态的类,调用Next才会继续执行下一个,不调用将会终止传递
        • 如果需要数据,可以有两种方式:临时Filter之间传递数据,可以存储在FilterCtx中,长期持有的,可以存储在Connection中

FilterChain

  • 消息的编解码

    • 在大部分应用中,消息的编解码是主要的通信工作,我这里划分了以下几个概念,Frame,Packet(Identity,Codec)
      • Frame:用于消息的粘包处理,例如类似http协议,以\r\n分隔,或者头部使用整数标识消息长度
      • Packet:消息包内容,通常分为两个部分,消息头和消息体,分别用Identity和Codec表示
        • Identity:用于消息头的编解码并创建具体的Packet
          • Fixed16Identity:小端编码,2个字节保存消息ID
          • StringIdentity:冒号分隔消息名和消息体
          • 自定义消息头:实现IIdentity接口并创建对应的IPacket
        • Codec: 用于消息体的编解码,例如json,protobuf
  • 线程模型

    • Connection线程,每个Connection都会创建一个读和写协程
      • InBound在Connection的读线程中处理,直到转发到ExectorFilter逻辑线程中处理
      • Outbound在调用线程中处理,直到最终调用Write方法转到写协程中发送数据
    • 消息处理线程,并没有强制约定,可以自己继承Filter实现定制消息处理,默认发送到一个单独的消息处理协程中
      • 单线程模式:只需末尾添加ExectorFilter即可实现消息统一转发的Exector中的消息队列中执行
      • Executor可以不止一个线程,比如:某些复杂但又独立的业务操作,可以在注册消息回调时制定一个queueIndex,则可以实现该模块在独立的线程中执行,但要使用者自己保证线程安全
    • 其他线程:Log线程,Timer线程,Executor线程
      • log线程需要注意的是属性的初始化是非线程安全的,需要在主线程中设置属性,启动后将不能再修改
      • timer线程,在一个独立的线程中执行定时器,如果需要放到消息线程中处理,需要手动Dispatch
  • 其他辅助类

    • buffer:底层的数据流存储,使用list存储[]byte,数据非连续的,可以像stream一样操作数据,使用时需要注意当前位置,以及哪些函数会影响当前位置
    • registry:非线程安全,用于消息的注册,可通过名字,或者id注册查询,也可以通过类型查询名字和id
    • dispatcher:非线程安全,handler的注册和查询
  • 扩展:本项目不依赖任何库,均以插件的形式扩展

  • 参考框架

Documentation

Index

Constants

View Source
const (
	// AttrKindConf 配置使用
	AttrKindConf = "conf"
	// AttrKindConn Conn中存储数据
	AttrKindConn = "conn"
	// AttrKindCtx filterContext中使用
	AttrKindCtx = "context"
)
View Source
const QueueMainID = 0

QueueMainID 主队列ID

Variables

This section is empty.

Functions

func InvokeHandler

func InvokeHandler(conn IConn, packet IPacket, handler *Handler)

InvokeHandler 调用Handler

func RegisterMessage

func RegisterMessage(msg interface{}, key interface{})

RegisterMessage 注册消息

Types

type AttrKey

type AttrKey struct {
	// contains filtered or unexported fields
}

AttrKey 将string延迟映射到唯一索引,不同进程间,索引并不一定一样

func NewAttrKey

func NewAttrKey(kind string, name string) *AttrKey

NewAttrKey 创建NewAttrKey

func (*AttrKey) Index

func (attr *AttrKey) Index() int

Index return attr index

func (*AttrKey) Name

func (attr *AttrKey) Name() string

Name return attr name

func (*AttrKey) String

func (attr *AttrKey) String() string

String return attr stringify

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 NewBuffer

func NewBuffer() *Buffer

func (*Buffer) Append

func (self *Buffer) Append(data []byte)

末尾出肉数据,游标移动到末尾

func (*Buffer) Bytes

func (self *Buffer) Bytes() []byte

func (*Buffer) Clear

func (self *Buffer) Clear()

Clear 清空数据

func (*Buffer) Concat

func (self *Buffer) Concat()

合并成一个[]byte

func (*Buffer) Discard

func (self *Buffer) Discard()

删除当前位置之前的数据

func (*Buffer) Empty

func (self *Buffer) Empty() bool

func (*Buffer) Eof

func (self *Buffer) Eof() bool

func (*Buffer) ExtendSpace

func (self *Buffer) ExtendSpace(count int) error

ExtendSpace 将最后一个扩展count个字节,配合GetSpace使用

func (*Buffer) Front

func (self *Buffer) Front() *list.Element

func (*Buffer) GetMark

func (self *Buffer) GetMark() int

func (*Buffer) GetSpace

func (self *Buffer) GetSpace() []byte

GetSpace = back of (capacity - length)

func (*Buffer) HasRemain

func (self *Buffer) HasRemain(count int) bool

func (*Buffer) IndexOf

func (self *Buffer) IndexOf(key string) int

IndexOf 从当前位置开始查询字符串位置

func (*Buffer) IndexOfLimit

func (self *Buffer) IndexOfLimit(key string, limit int) int

IndexOfLimit 从当前位置查找key,最长搜索limit个字节(-1无限制)

func (*Buffer) Length

func (self *Buffer) Length() int

func (*Buffer) Merge

func (self *Buffer) Merge(other *Buffer)

Merge 合并两个buffer为一个

func (*Buffer) Peek

func (self *Buffer) Peek(buffer []byte) (int, error)

Peek 获取数据但不修改当前位置

func (*Buffer) Position

func (self *Buffer) Position() int

func (*Buffer) Prepend

func (self *Buffer) Prepend(data []byte)

前边插入数据,游标改为初始位置

func (*Buffer) Read

func (self *Buffer) Read(buffer []byte) (int, error)

Read 读取数据并移动当前位置

func (*Buffer) ReadAll

func (buffer *Buffer) ReadAll(reader io.Reader) error

func (*Buffer) ReadByte

func (self *Buffer) ReadByte() (byte, error)

ReadByte 实现接口io.ByteReader

func (*Buffer) ReadLine

func (self *Buffer) ReadLine() (*Buffer, error)

ReadLine 读取到\n或\r\n为止

func (*Buffer) ReadToEnd

func (self *Buffer) ReadToEnd() []byte

func (*Buffer) ReadUntil

func (self *Buffer) ReadUntil(key byte) (string, error)

ReadUntil 读取到key位置

func (*Buffer) Rewind

func (self *Buffer) Rewind()

Rewind 回到头部

func (*Buffer) Seek

func (self *Buffer) Seek(offset int, whence int) error

Seek 移动当前位置

func (*Buffer) SetMark

func (self *Buffer) SetMark(value int)

func (*Buffer) Split

func (self *Buffer) Split(result *Buffer)

从当前位置分隔成两个

func (*Buffer) String

func (self *Buffer) String() string

func (*Buffer) Swap

func (self *Buffer) Swap(other *Buffer)

Swap 交换两个buffer

func (*Buffer) Visit

func (self *Buffer) Visit(cb func([]byte) bool)

遍历整个数组

func (*Buffer) Write

func (self *Buffer) Write(bufffer []byte) (int, error)

Write 实现io.Writer接口

func (*Buffer) WriteAll

func (buffer *Buffer) WriteAll(writer io.Writer) error

type Callback

type Callback func()

Callback 回调函数

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队列作为主循环

func GetExecutor

func GetExecutor() *Executor

GetExecutor 获取全局调度器

func NewExecutor

func NewExecutor() *Executor

NewExecutor 创建调度器

func (*Executor) Dispatch

func (exec *Executor) Dispatch(queueId uint, task Callback) error

Dispatch 分发到相应队列执行任务

func (*Executor) EnsureQueue

func (exec *Executor) EnsureQueue(queueId uint)

EnsureQueue 确保队列存在

func (*Executor) Stop

func (exec *Executor) Stop()

Stop 关闭队列,等待所有任务结束

type Handler

type Handler struct {
	Func    HandlerCB
	QueueId uint
	Data    interface{}
}

Handler 消息处理

func RegisterHandler

func RegisterHandler(key interface{}, cb HandlerCB) *Handler

RegisterHandler 注册消息回调

type HandlerCB

type HandlerCB func(ctx *HandlerCtx)

HandlerCB 消息回调

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 IFrame

type IFrame interface {
	Encode(buffer *Buffer) error
	Decode(buffer *Buffer) (*Buffer, error)
}

IFrame 帧序列,用于粘包处理

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
}

////////////////////////////////////////////////////// 迭代器实现 //////////////////////////////////////////////////////

func (*Iterator) Create

func (self *Iterator) Create(elem *list.Element, offset int, length int)

func (*Iterator) MoveEnd

func (self *Iterator) MoveEnd()

func (*Iterator) Next

func (self *Iterator) Next() bool

type Option

type Option interface {
}

Option 配置

type ReconnectOption

type ReconnectOption struct {
	Count    int // 重连次数,0标识无需断线重连
	Interval int // 重连间隔
}

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

func GetRegistry

func GetRegistry() *Registry

GetRegistry 获取全局Registry

func NewRegistry

func NewRegistry() *Registry

NewRegistry 创建Registry

func (*Registry) Create

func (r *Registry) Create(msgid uint, name string) interface{}

Create 通过Id或者名字创建

func (*Registry) GetInfo

func (r *Registry) GetInfo(msg interface{}) (uint, string, bool)

GetInfo 通过消息类型,查询信息

func (*Registry) Register

func (r *Registry) Register(msg interface{}, key interface{}) error

Register 注册消息元信息,id为字符串或者整数,nil则反射msg名字注册

func (*Registry) Remove

func (r *Registry) Remove(msg interface{}) bool

Remove 删除消息注册

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

type TagOption

type TagOption struct {
	Tag string
}

TagOption 传递Tag使用

func WithTagOption

func WithTagOption(tag string) *TagOption

Directories

Path Synopsis
container
inlist
Package inlist implements an intrusive doubly linked list.
Package inlist implements an intrusive doubly linked list.

Jump to

Keyboard shortcuts

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