mirror of
https://github.com/duanhf2012/origin.git
synced 2026-02-12 22:54:43 +08:00
Compare commits
16 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8e0ed62fca | ||
|
|
7116b509e9 | ||
|
|
73d384361d | ||
|
|
ce56b19fe8 | ||
|
|
1367d776e6 | ||
|
|
987d35ff15 | ||
|
|
d225bb4bd2 | ||
|
|
ea37fb5081 | ||
|
|
0a92f48d0b | ||
|
|
f5e86fee02 | ||
|
|
166facc959 | ||
|
|
5bb747201b | ||
|
|
1014bc54e4 | ||
|
|
9c26c742fe | ||
|
|
d1935b1bbc | ||
|
|
90d54bf3e2 |
123
README.md
123
README.md
@@ -661,51 +661,7 @@ Module1 Release.
|
|||||||
第四章:事件使用
|
第四章:事件使用
|
||||||
----------------
|
----------------
|
||||||
|
|
||||||
事件是origin中一个重要的组成部分,可以在同一个node中的service与service或者与module之间进行事件通知。系统内置的几个服务,如:TcpService/HttpService等都是通过事件功能实现。他也是一个典型的观察者设计模型。在event中有两个类型的interface,一个是event.IEventProcessor它提供注册与卸载功能,另一个是event.IEventHandler提供消息广播等功能。
|
事件是origin中一个重要的组成部分,可以在服务与各module之间进行事件通知。它也是一个典型的观察者设计模型。在event中有两个类型的interface,一个是event.IEventProcessor它提供注册与卸载功能,另一个是event.IEventHandler提供消息广播等功能。
|
||||||
|
|
||||||
在目录simple_event/TestService4.go中
|
|
||||||
|
|
||||||
```
|
|
||||||
package simple_event
|
|
||||||
|
|
||||||
import (
|
|
||||||
"github.com/duanhf2012/origin/v2/event"
|
|
||||||
"github.com/duanhf2012/origin/v2/node"
|
|
||||||
"github.com/duanhf2012/origin/v2/service"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
const (
|
|
||||||
//自定义事件类型,必需从event.Sys_Event_User_Define开始
|
|
||||||
//event.Sys_Event_User_Define以内给系统预留
|
|
||||||
EVENT1 event.EventType =event.Sys_Event_User_Define+1
|
|
||||||
)
|
|
||||||
|
|
||||||
func init(){
|
|
||||||
node.Setup(&TestService4{})
|
|
||||||
}
|
|
||||||
|
|
||||||
type TestService4 struct {
|
|
||||||
service.Service
|
|
||||||
}
|
|
||||||
|
|
||||||
func (slf *TestService4) OnInit() error {
|
|
||||||
//10秒后触发广播事件
|
|
||||||
slf.AfterFunc(time.Second*10,slf.TriggerEvent)
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (slf *TestService4) TriggerEvent(){
|
|
||||||
//广播事件,传入event.Event对象,类型为EVENT1,Data可以自定义任何数据
|
|
||||||
//这样,所有监听者都可以收到该事件
|
|
||||||
slf.GetEventHandler().NotifyEvent(&event.Event{
|
|
||||||
Type: EVENT1,
|
|
||||||
Data: "event data.",
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
```
|
|
||||||
|
|
||||||
在目录simple_event/TestService5.go中
|
在目录simple_event/TestService5.go中
|
||||||
|
|
||||||
@@ -713,53 +669,68 @@ func (slf *TestService4) TriggerEvent(){
|
|||||||
package simple_event
|
package simple_event
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"github.com/duanhf2012/origin/v2/event"
|
"github.com/duanhf2012/origin/v2/event"
|
||||||
"github.com/duanhf2012/origin/v2/node"
|
"github.com/duanhf2012/origin/v2/node"
|
||||||
"github.com/duanhf2012/origin/v2/service"
|
"github.com/duanhf2012/origin/v2/service"
|
||||||
|
"github.com/duanhf2012/origin/v2/util/timer"
|
||||||
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
func init(){
|
func init() {
|
||||||
node.Setup(&TestService5{})
|
node.Setup(&TestService5{})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
//自定义事件类型,必需从event.Sys_Event_User_Define开始
|
||||||
|
//event.Sys_Event_User_Define以内给系统预留
|
||||||
|
EVENT1 event.EventType = event.Sys_Event_User_Define + 1
|
||||||
|
)
|
||||||
|
|
||||||
type TestService5 struct {
|
type TestService5 struct {
|
||||||
service.Service
|
service.Service
|
||||||
}
|
}
|
||||||
|
|
||||||
type TestModule struct {
|
type TestModule struct {
|
||||||
service.Module
|
service.Module
|
||||||
}
|
}
|
||||||
|
|
||||||
func (slf *TestModule) OnInit() error{
|
func (slf *TestModule) OnInit() error {
|
||||||
//在当前node中查找TestService4
|
//在TestModule中注册监听EVENT1事件
|
||||||
pService := node.GetService("TestService4")
|
slf.GetEventProcessor().RegEventReceiverFunc(EVENT1, slf.GetEventHandler(), slf.OnModuleEvent)
|
||||||
|
|
||||||
//在TestModule中,往TestService4中注册EVENT1类型事件监听
|
return nil
|
||||||
pService.(*TestService4).GetEventProcessor().RegEventReciverFunc(EVENT1,slf.GetEventHandler(),slf.OnModuleEvent)
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (slf *TestModule) OnModuleEvent(ev event.IEvent){
|
// OnModuleEvent 模块监听事件回调
|
||||||
event := ev.(*event.Event)
|
func (slf *TestModule) OnModuleEvent(ev event.IEvent) {
|
||||||
fmt.Printf("OnModuleEvent type :%d data:%+v\n",event.GetEventType(),event.Data)
|
event := ev.(*event.Event)
|
||||||
|
fmt.Printf("OnModuleEvent type :%d data:%+v\n", event.GetEventType(), event.Data)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// OnInit 服务初始化函数,在安装服务时,服务将自动调用OnInit函数
|
||||||
//服务初始化函数,在安装服务时,服务将自动调用OnInit函数
|
|
||||||
func (slf *TestService5) OnInit() error {
|
func (slf *TestService5) OnInit() error {
|
||||||
//通过服务名获取服务对象
|
//在服务中注册监听EVENT1类型事件
|
||||||
pService := node.GetService("TestService4")
|
slf.RegEventReceiverFunc(EVENT1, slf.GetEventHandler(), slf.OnServiceEvent)
|
||||||
|
slf.AddModule(&TestModule{})
|
||||||
|
|
||||||
////在TestModule中,往TestService4中注册EVENT1类型事件监听
|
slf.AfterFunc(time.Second*10, slf.TriggerEvent)
|
||||||
pService.(*TestService4).GetEventProcessor().RegEventReciverFunc(EVENT1,slf.GetEventHandler(),slf.OnServiceEvent)
|
return nil
|
||||||
slf.AddModule(&TestModule{})
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (slf *TestService5) OnServiceEvent(ev event.IEvent){
|
// OnServiceEvent 服务监听事件回调
|
||||||
event := ev.(*event.Event)
|
func (slf *TestService5) OnServiceEvent(ev event.IEvent) {
|
||||||
fmt.Printf("OnServiceEvent type :%d data:%+v\n",event.Type,event.Data)
|
event := ev.(*event.Event)
|
||||||
|
fmt.Printf("OnServiceEvent type :%d data:%+v\n", event.Type, event.Data)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (slf *TestService5) TriggerEvent(t *timer.Timer) {
|
||||||
|
//广播事件,传入event.Event对象,类型为EVENT1,Data可以自定义任何数据
|
||||||
|
//这样,所有监听者都可以收到该事件
|
||||||
|
slf.GetEventHandler().NotifyEvent(&event.Event{
|
||||||
|
Type: EVENT1,
|
||||||
|
Data: "event data.",
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -768,8 +739,8 @@ func (slf *TestService5) OnServiceEvent(ev event.IEvent){
|
|||||||
程序运行10秒后,调用slf.TriggerEvent函数广播事件,于是在TestService5中会收到
|
程序运行10秒后,调用slf.TriggerEvent函数广播事件,于是在TestService5中会收到
|
||||||
|
|
||||||
```
|
```
|
||||||
OnServiceEvent type :1001 data:event data.
|
OnServiceEvent type :2 data:event data.
|
||||||
OnModuleEvent type :1001 data:event data.
|
OnModuleEvent type :2 data:event data.
|
||||||
```
|
```
|
||||||
|
|
||||||
在上面的TestModule中监听的事情,当这个Module被Release时监听会自动卸载。
|
在上面的TestModule中监听的事情,当这个Module被Release时监听会自动卸载。
|
||||||
@@ -1181,6 +1152,8 @@ func (slf *TestTcpService) OnRequest (clientid string,msg proto.Message){
|
|||||||
* log/log.go:日志的封装,可以使用它构建对象记录业务文件日志
|
* log/log.go:日志的封装,可以使用它构建对象记录业务文件日志
|
||||||
* util:在该目录下,有常用的uuid,hash,md5,协程封装等工具库
|
* util:在该目录下,有常用的uuid,hash,md5,协程封装等工具库
|
||||||
* https://github.com/duanhf2012/originservice: 其他扩展支持的服务可以在该工程上看到,目前支持firebase推送的封装。
|
* https://github.com/duanhf2012/originservice: 其他扩展支持的服务可以在该工程上看到,目前支持firebase推送的封装。
|
||||||
|
* https://github.com/duanhf2012/origingame: 基础游戏服务器的框架
|
||||||
|
* etcd与nats开发环境搭建可以从https://github.com/duanhf2012/originserver_v2下的docker-compose获取
|
||||||
|
|
||||||
备注:
|
备注:
|
||||||
-----
|
-----
|
||||||
|
|||||||
@@ -8,6 +8,8 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"github.com/duanhf2012/origin/v2/event"
|
"github.com/duanhf2012/origin/v2/event"
|
||||||
|
"errors"
|
||||||
|
"reflect"
|
||||||
)
|
)
|
||||||
|
|
||||||
var configDir = "./config/"
|
var configDir = "./config/"
|
||||||
@@ -433,6 +435,25 @@ func (cls *Cluster) GetGlobalCfg() interface{} {
|
|||||||
return cls.globalCfg
|
return cls.globalCfg
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
func (cls *Cluster) ParseGlobalCfg(cfg interface{}) error{
|
||||||
|
if cls.globalCfg == nil {
|
||||||
|
return errors.New("no service configuration found")
|
||||||
|
}
|
||||||
|
|
||||||
|
rv := reflect.ValueOf(cls.globalCfg)
|
||||||
|
if rv.Kind() == reflect.Ptr && rv.IsNil() {
|
||||||
|
return errors.New("no service configuration found")
|
||||||
|
}
|
||||||
|
|
||||||
|
bytes,err := json.Marshal(cls.globalCfg)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return json.Unmarshal(bytes,cfg)
|
||||||
|
}
|
||||||
|
|
||||||
func (cls *Cluster) GetNodeInfo(nodeId string) (NodeInfo,bool) {
|
func (cls *Cluster) GetNodeInfo(nodeId string) (NodeInfo,bool) {
|
||||||
cls.locker.RLock()
|
cls.locker.RLock()
|
||||||
defer cls.locker.RUnlock()
|
defer cls.locker.RUnlock()
|
||||||
@@ -448,6 +469,11 @@ func (cls *Cluster) GetNodeInfo(nodeId string) (NodeInfo,bool) {
|
|||||||
func (dc *Cluster) CanDiscoveryService(fromMasterNodeId string,serviceName string) bool{
|
func (dc *Cluster) CanDiscoveryService(fromMasterNodeId string,serviceName string) bool{
|
||||||
canDiscovery := true
|
canDiscovery := true
|
||||||
|
|
||||||
|
splitServiceName := strings.Split(serviceName,":")
|
||||||
|
if len(splitServiceName) == 2 {
|
||||||
|
serviceName = splitServiceName[0]
|
||||||
|
}
|
||||||
|
|
||||||
for i:=0;i<len(dc.GetLocalNodeInfo().DiscoveryService);i++{
|
for i:=0;i<len(dc.GetLocalNodeInfo().DiscoveryService);i++{
|
||||||
masterNodeId := dc.GetLocalNodeInfo().DiscoveryService[i].MasterNodeId
|
masterNodeId := dc.GetLocalNodeInfo().DiscoveryService[i].MasterNodeId
|
||||||
//无效的配置,则跳过
|
//无效的配置,则跳过
|
||||||
|
|||||||
@@ -160,6 +160,12 @@ func (dc *OriginDiscoveryMaster) OnNatsDisconnect(){
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (ds *OriginDiscoveryMaster) OnNodeConnected(nodeId string) {
|
func (ds *OriginDiscoveryMaster) OnNodeConnected(nodeId string) {
|
||||||
|
var notifyDiscover rpc.SubscribeDiscoverNotify
|
||||||
|
notifyDiscover.IsFull = true
|
||||||
|
notifyDiscover.NodeInfo = ds.nodeInfo
|
||||||
|
notifyDiscover.MasterNodeId = cluster.GetLocalNodeInfo().NodeId
|
||||||
|
|
||||||
|
ds.GoNode(nodeId, SubServiceDiscover, ¬ifyDiscover)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ds *OriginDiscoveryMaster) OnNodeDisconnect(nodeId string) {
|
func (ds *OriginDiscoveryMaster) OnNodeDisconnect(nodeId string) {
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
|
|
||||||
"github.com/duanhf2012/origin/v2/log"
|
"github.com/duanhf2012/origin/v2/log"
|
||||||
"github.com/duanhf2012/origin/v2/util/queue"
|
"github.com/duanhf2012/origin/v2/util/queue"
|
||||||
|
"context"
|
||||||
)
|
)
|
||||||
|
|
||||||
var idleTimeout = int64(2 * time.Second)
|
var idleTimeout = int64(2 * time.Second)
|
||||||
@@ -30,6 +31,9 @@ type dispatch struct {
|
|||||||
|
|
||||||
waitWorker sync.WaitGroup
|
waitWorker sync.WaitGroup
|
||||||
waitDispatch sync.WaitGroup
|
waitDispatch sync.WaitGroup
|
||||||
|
|
||||||
|
cancelContext context.Context
|
||||||
|
cancel context.CancelFunc
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *dispatch) open(minGoroutineNum int32, maxGoroutineNum int32, tasks chan task, cbChannel chan func(error)) {
|
func (d *dispatch) open(minGoroutineNum int32, maxGoroutineNum int32, tasks chan task, cbChannel chan func(error)) {
|
||||||
@@ -40,7 +44,7 @@ func (d *dispatch) open(minGoroutineNum int32, maxGoroutineNum int32, tasks chan
|
|||||||
d.workerQueue = make(chan task)
|
d.workerQueue = make(chan task)
|
||||||
d.cbChannel = cbChannel
|
d.cbChannel = cbChannel
|
||||||
d.queueIdChannel = make(chan int64, cap(tasks))
|
d.queueIdChannel = make(chan int64, cap(tasks))
|
||||||
|
d.cancelContext,d.cancel = context.WithCancel(context.Background())
|
||||||
d.waitDispatch.Add(1)
|
d.waitDispatch.Add(1)
|
||||||
go d.run()
|
go d.run()
|
||||||
}
|
}
|
||||||
@@ -64,10 +68,12 @@ func (d *dispatch) run() {
|
|||||||
d.processqueueEvent(queueId)
|
d.processqueueEvent(queueId)
|
||||||
case <-timeout.C:
|
case <-timeout.C:
|
||||||
d.processTimer()
|
d.processTimer()
|
||||||
if atomic.LoadInt32(&d.minConcurrentNum) == -1 && len(d.tasks) == 0 {
|
case <- d.cancelContext.Done():
|
||||||
atomic.StoreInt64(&idleTimeout,int64(time.Millisecond * 10))
|
atomic.StoreInt64(&idleTimeout,int64(time.Millisecond * 5))
|
||||||
}
|
|
||||||
timeout.Reset(time.Duration(atomic.LoadInt64(&idleTimeout)))
|
timeout.Reset(time.Duration(atomic.LoadInt64(&idleTimeout)))
|
||||||
|
for i:=int32(0);i<d.workerNum;i++{
|
||||||
|
d.processIdle()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -166,6 +172,8 @@ func (c *dispatch) pushAsyncDoCallbackEvent(cb func(err error)) {
|
|||||||
|
|
||||||
func (d *dispatch) close() {
|
func (d *dispatch) close() {
|
||||||
atomic.StoreInt32(&d.minConcurrentNum, -1)
|
atomic.StoreInt32(&d.minConcurrentNum, -1)
|
||||||
|
d.cancel()
|
||||||
|
|
||||||
|
|
||||||
breakFor:
|
breakFor:
|
||||||
for {
|
for {
|
||||||
|
|||||||
@@ -10,16 +10,20 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const defaultSkip = 7
|
||||||
type IOriginHandler interface {
|
type IOriginHandler interface {
|
||||||
slog.Handler
|
slog.Handler
|
||||||
Lock()
|
Lock()
|
||||||
UnLock()
|
UnLock()
|
||||||
|
SetSkip(skip int)
|
||||||
|
GetSkip() int
|
||||||
}
|
}
|
||||||
|
|
||||||
type BaseHandler struct {
|
type BaseHandler struct {
|
||||||
addSource bool
|
addSource bool
|
||||||
w io.Writer
|
w io.Writer
|
||||||
locker sync.Mutex
|
locker sync.Mutex
|
||||||
|
skip int
|
||||||
}
|
}
|
||||||
|
|
||||||
type OriginTextHandler struct {
|
type OriginTextHandler struct {
|
||||||
@@ -32,6 +36,14 @@ type OriginJsonHandler struct {
|
|||||||
*slog.JSONHandler
|
*slog.JSONHandler
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (bh *BaseHandler) SetSkip(skip int){
|
||||||
|
bh.skip = skip
|
||||||
|
}
|
||||||
|
|
||||||
|
func (bh *BaseHandler) GetSkip() int{
|
||||||
|
return bh.skip
|
||||||
|
}
|
||||||
|
|
||||||
func getStrLevel(level slog.Level) string{
|
func getStrLevel(level slog.Level) string{
|
||||||
switch level {
|
switch level {
|
||||||
case LevelTrace:
|
case LevelTrace:
|
||||||
@@ -78,6 +90,7 @@ func NewOriginTextHandler(level slog.Level,w io.Writer,addSource bool,replaceAtt
|
|||||||
ReplaceAttr: replaceAttr,
|
ReplaceAttr: replaceAttr,
|
||||||
})
|
})
|
||||||
|
|
||||||
|
textHandler.skip = defaultSkip
|
||||||
return &textHandler
|
return &textHandler
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -124,6 +137,7 @@ func NewOriginJsonHandler(level slog.Level,w io.Writer,addSource bool,replaceAtt
|
|||||||
ReplaceAttr: replaceAttr,
|
ReplaceAttr: replaceAttr,
|
||||||
})
|
})
|
||||||
|
|
||||||
|
jsonHandler.skip = defaultSkip
|
||||||
return &jsonHandler
|
return &jsonHandler
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -141,7 +155,7 @@ func (oh *OriginJsonHandler) Handle(context context.Context, record slog.Record)
|
|||||||
func (b *BaseHandler) Fill(context context.Context, record *slog.Record) {
|
func (b *BaseHandler) Fill(context context.Context, record *slog.Record) {
|
||||||
if b.addSource {
|
if b.addSource {
|
||||||
var pcs [1]uintptr
|
var pcs [1]uintptr
|
||||||
runtime.Callers(7, pcs[:])
|
runtime.Callers(b.skip, pcs[:])
|
||||||
record.PC = pcs[0]
|
record.PC = pcs[0]
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -239,6 +239,10 @@ func (iw *IoWriter) swichFile() error{
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func GetDefaultHandler() IOriginHandler{
|
||||||
|
return gLogger.(*Logger).Slogger.Handler().(IOriginHandler)
|
||||||
|
}
|
||||||
|
|
||||||
func NewTextLogger(level slog.Level,pathName string,filePrefix string,addSource bool,logChannelCap int) (ILogger,error){
|
func NewTextLogger(level slog.Level,pathName string,filePrefix string,addSource bool,logChannelCap int) (ILogger,error){
|
||||||
var logger Logger
|
var logger Logger
|
||||||
logger.ioWriter.filePath = pathName
|
logger.ioWriter.filePath = pathName
|
||||||
|
|||||||
@@ -45,8 +45,10 @@ func (jsonProcessor *JsonProcessor) SetByteOrder(littleEndian bool) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
func (jsonProcessor *JsonProcessor ) MsgRoute(clientId string,msg interface{}) error{
|
func (jsonProcessor *JsonProcessor ) MsgRoute(clientId string,msg interface{},recyclerReaderBytes func(data []byte)) error{
|
||||||
pPackInfo := msg.(*JsonPackInfo)
|
pPackInfo := msg.(*JsonPackInfo)
|
||||||
|
defer recyclerReaderBytes(pPackInfo.rawMsg)
|
||||||
|
|
||||||
v,ok := jsonProcessor.mapMsg[pPackInfo.typ]
|
v,ok := jsonProcessor.mapMsg[pPackInfo.typ]
|
||||||
if ok == false {
|
if ok == false {
|
||||||
return fmt.Errorf("cannot find msgtype %d is register!",pPackInfo.typ)
|
return fmt.Errorf("cannot find msgtype %d is register!",pPackInfo.typ)
|
||||||
@@ -58,7 +60,6 @@ func (jsonProcessor *JsonProcessor ) MsgRoute(clientId string,msg interface{}) e
|
|||||||
|
|
||||||
func (jsonProcessor *JsonProcessor) Unmarshal(clientId string,data []byte) (interface{}, error) {
|
func (jsonProcessor *JsonProcessor) Unmarshal(clientId string,data []byte) (interface{}, error) {
|
||||||
typeStruct := struct {Type int `json:"typ"`}{}
|
typeStruct := struct {Type int `json:"typ"`}{}
|
||||||
defer jsonProcessor.ReleaseBytes(data)
|
|
||||||
err := json.Unmarshal(data, &typeStruct)
|
err := json.Unmarshal(data, &typeStruct)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -76,7 +77,7 @@ func (jsonProcessor *JsonProcessor) Unmarshal(clientId string,data []byte) (inte
|
|||||||
return nil,err
|
return nil,err
|
||||||
}
|
}
|
||||||
|
|
||||||
return &JsonPackInfo{typ:msgType,msg:msgData},nil
|
return &JsonPackInfo{typ:msgType,msg:msgData,rawMsg: data},nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (jsonProcessor *JsonProcessor) Marshal(clientId string,msg interface{}) ([]byte, error) {
|
func (jsonProcessor *JsonProcessor) Marshal(clientId string,msg interface{}) ([]byte, error) {
|
||||||
@@ -104,7 +105,8 @@ func (jsonProcessor *JsonProcessor) MakeRawMsg(msgType uint16,msg []byte) *JsonP
|
|||||||
return &JsonPackInfo{typ:msgType,rawMsg:msg}
|
return &JsonPackInfo{typ:msgType,rawMsg:msg}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (jsonProcessor *JsonProcessor) UnknownMsgRoute(clientId string,msg interface{}){
|
func (jsonProcessor *JsonProcessor) UnknownMsgRoute(clientId string,msg interface{},recyclerReaderBytes func(data []byte)){
|
||||||
|
defer recyclerReaderBytes(msg.([]byte))
|
||||||
if jsonProcessor.unknownMessageHandler==nil {
|
if jsonProcessor.unknownMessageHandler==nil {
|
||||||
log.Debug("Unknown message",log.String("clientId",clientId))
|
log.Debug("Unknown message",log.String("clientId",clientId))
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -54,8 +54,10 @@ func (slf *PBPackInfo) GetMsg() proto.Message {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
func (pbProcessor *PBProcessor) MsgRoute(clientId string, msg interface{}) error {
|
func (pbProcessor *PBProcessor) MsgRoute(clientId string, msg interface{},recyclerReaderBytes func(data []byte)) error {
|
||||||
pPackInfo := msg.(*PBPackInfo)
|
pPackInfo := msg.(*PBPackInfo)
|
||||||
|
defer recyclerReaderBytes(pPackInfo.rawMsg)
|
||||||
|
|
||||||
v, ok := pbProcessor.mapMsg[pPackInfo.typ]
|
v, ok := pbProcessor.mapMsg[pPackInfo.typ]
|
||||||
if ok == false {
|
if ok == false {
|
||||||
return fmt.Errorf("Cannot find msgtype %d is register!", pPackInfo.typ)
|
return fmt.Errorf("Cannot find msgtype %d is register!", pPackInfo.typ)
|
||||||
@@ -67,7 +69,6 @@ func (pbProcessor *PBProcessor) MsgRoute(clientId string, msg interface{}) error
|
|||||||
|
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
func (pbProcessor *PBProcessor) Unmarshal(clientId string, data []byte) (interface{}, error) {
|
func (pbProcessor *PBProcessor) Unmarshal(clientId string, data []byte) (interface{}, error) {
|
||||||
defer pbProcessor.ReleaseBytes(data)
|
|
||||||
return pbProcessor.UnmarshalWithOutRelease(clientId, data)
|
return pbProcessor.UnmarshalWithOutRelease(clientId, data)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -91,7 +92,7 @@ func (pbProcessor *PBProcessor) UnmarshalWithOutRelease(clientId string, data []
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return &PBPackInfo{typ: msgType, msg: protoMsg}, nil
|
return &PBPackInfo{typ: msgType, msg: protoMsg,rawMsg:data}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
@@ -133,8 +134,9 @@ func (pbProcessor *PBProcessor) MakeRawMsg(msgType uint16, msg []byte) *PBPackIn
|
|||||||
return &PBPackInfo{typ: msgType, rawMsg: msg}
|
return &PBPackInfo{typ: msgType, rawMsg: msg}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pbProcessor *PBProcessor) UnknownMsgRoute(clientId string, msg interface{}) {
|
func (pbProcessor *PBProcessor) UnknownMsgRoute(clientId string, msg interface{},recyclerReaderBytes func(data []byte)) {
|
||||||
pbProcessor.unknownMessageHandler(clientId, msg.([]byte))
|
pbProcessor.unknownMessageHandler(clientId, msg.([]byte))
|
||||||
|
recyclerReaderBytes(msg.([]byte))
|
||||||
}
|
}
|
||||||
|
|
||||||
// connect event
|
// connect event
|
||||||
|
|||||||
@@ -38,9 +38,11 @@ func (pbRawProcessor *PBRawProcessor) SetByteOrder(littleEndian bool) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
func (pbRawProcessor *PBRawProcessor ) MsgRoute(clientId string, msg interface{}) error{
|
func (pbRawProcessor *PBRawProcessor ) MsgRoute(clientId string, msg interface{},recyclerReaderBytes func(data []byte)) error{
|
||||||
pPackInfo := msg.(*PBRawPackInfo)
|
pPackInfo := msg.(*PBRawPackInfo)
|
||||||
pbRawProcessor.msgHandler(clientId,pPackInfo.typ,pPackInfo.rawMsg)
|
pbRawProcessor.msgHandler(clientId,pPackInfo.typ,pPackInfo.rawMsg)
|
||||||
|
recyclerReaderBytes(pPackInfo.rawMsg)
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -80,7 +82,8 @@ func (pbRawProcessor *PBRawProcessor) MakeRawMsg(msgType uint16,msg []byte,pbRaw
|
|||||||
pbRawPackInfo.rawMsg = msg
|
pbRawPackInfo.rawMsg = msg
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pbRawProcessor *PBRawProcessor) UnknownMsgRoute(clientId string,msg interface{}){
|
func (pbRawProcessor *PBRawProcessor) UnknownMsgRoute(clientId string,msg interface{},recyclerReaderBytes func(data []byte)){
|
||||||
|
defer recyclerReaderBytes(msg.([]byte))
|
||||||
if pbRawProcessor.unknownMessageHandler == nil {
|
if pbRawProcessor.unknownMessageHandler == nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,9 +3,9 @@ package processor
|
|||||||
|
|
||||||
type IProcessor interface {
|
type IProcessor interface {
|
||||||
// must goroutine safe
|
// must goroutine safe
|
||||||
MsgRoute(clientId string,msg interface{}) error
|
MsgRoute(clientId string,msg interface{},recyclerReaderBytes func(data []byte)) error
|
||||||
//must goroutine safe
|
//must goroutine safe
|
||||||
UnknownMsgRoute(clientId string,msg interface{})
|
UnknownMsgRoute(clientId string,msg interface{},recyclerReaderBytes func(data []byte))
|
||||||
// connect event
|
// connect event
|
||||||
ConnectedRoute(clientId string)
|
ConnectedRoute(clientId string)
|
||||||
DisConnectedRoute(clientId string)
|
DisConnectedRoute(clientId string)
|
||||||
|
|||||||
@@ -129,6 +129,13 @@ func (tcpConn *TCPConn) ReadMsg() ([]byte, error) {
|
|||||||
return tcpConn.msgParser.Read(tcpConn)
|
return tcpConn.msgParser.Read(tcpConn)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tcpConn *TCPConn) GetRecyclerReaderBytes() func (data []byte) {
|
||||||
|
bytePool := tcpConn.msgParser.IBytesMempool
|
||||||
|
return func(data []byte) {
|
||||||
|
bytePool.ReleaseBytes(data)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (tcpConn *TCPConn) ReleaseReadMsg(byteBuff []byte){
|
func (tcpConn *TCPConn) ReleaseReadMsg(byteBuff []byte){
|
||||||
tcpConn.msgParser.ReleaseBytes(byteBuff)
|
tcpConn.msgParser.ReleaseBytes(byteBuff)
|
||||||
}
|
}
|
||||||
|
|||||||
14
node/node.go
14
node/node.go
@@ -29,6 +29,7 @@ var preSetupTemplateService []func()service.IService
|
|||||||
var profilerInterval time.Duration
|
var profilerInterval time.Duration
|
||||||
var bValid bool
|
var bValid bool
|
||||||
var configDir = "./config/"
|
var configDir = "./config/"
|
||||||
|
var NodeIsRun = false
|
||||||
|
|
||||||
const(
|
const(
|
||||||
SingleStop syscall.Signal = 10
|
SingleStop syscall.Signal = 10
|
||||||
@@ -57,7 +58,7 @@ func init() {
|
|||||||
console.RegisterCommandString("loglevel", "debug", "<-loglevel debug|release|warning|error|fatal> Set loglevel.", setLevel)
|
console.RegisterCommandString("loglevel", "debug", "<-loglevel debug|release|warning|error|fatal> Set loglevel.", setLevel)
|
||||||
console.RegisterCommandString("logpath", "", "<-logpath path> Set log file path.", setLogPath)
|
console.RegisterCommandString("logpath", "", "<-logpath path> Set log file path.", setLogPath)
|
||||||
console.RegisterCommandInt("logsize", 0, "<-logsize size> Set log size(MB).", setLogSize)
|
console.RegisterCommandInt("logsize", 0, "<-logsize size> Set log size(MB).", setLogSize)
|
||||||
console.RegisterCommandInt("logchannelcap", 0, "<-logchannelcap num> Set log channel cap.", setLogChannelCapNum)
|
console.RegisterCommandInt("logchannelcap", -1, "<-logchannelcap num> Set log channel cap.", setLogChannelCapNum)
|
||||||
console.RegisterCommandString("pprof", "", "<-pprof ip:port> Open performance analysis.", setPprof)
|
console.RegisterCommandString("pprof", "", "<-pprof ip:port> Open performance analysis.", setPprof)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -354,13 +355,14 @@ func startNode(args interface{}) error {
|
|||||||
cluster.GetCluster().Start()
|
cluster.GetCluster().Start()
|
||||||
|
|
||||||
//6.监听程序退出信号&性能报告
|
//6.监听程序退出信号&性能报告
|
||||||
bRun := true
|
|
||||||
var pProfilerTicker *time.Ticker = &time.Ticker{}
|
var pProfilerTicker *time.Ticker = &time.Ticker{}
|
||||||
if profilerInterval > 0 {
|
if profilerInterval > 0 {
|
||||||
pProfilerTicker = time.NewTicker(profilerInterval)
|
pProfilerTicker = time.NewTicker(profilerInterval)
|
||||||
}
|
}
|
||||||
|
|
||||||
for bRun {
|
NodeIsRun = true
|
||||||
|
for NodeIsRun {
|
||||||
select {
|
select {
|
||||||
case s := <-sig:
|
case s := <-sig:
|
||||||
signal := s.(syscall.Signal)
|
signal := s.(syscall.Signal)
|
||||||
@@ -368,7 +370,7 @@ func startNode(args interface{}) error {
|
|||||||
log.Info("receipt retire signal.")
|
log.Info("receipt retire signal.")
|
||||||
notifyAllServiceRetire()
|
notifyAllServiceRetire()
|
||||||
}else {
|
}else {
|
||||||
bRun = false
|
NodeIsRun = false
|
||||||
log.Info("receipt stop signal.")
|
log.Info("receipt stop signal.")
|
||||||
}
|
}
|
||||||
case <-pProfilerTicker.C:
|
case <-pProfilerTicker.C:
|
||||||
@@ -504,6 +506,10 @@ func setLogChannelCapNum(args interface{}) error {
|
|||||||
return errors.New("param logsize is error")
|
return errors.New("param logsize is error")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if logChannelCap == -1 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
log.LogChannelCap = logChannelCap
|
log.LogChannelCap = logChannelCap
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -110,7 +110,6 @@ type IRpcHandler interface {
|
|||||||
GoNode(nodeId string, serviceMethod string, args interface{}) error
|
GoNode(nodeId string, serviceMethod string, args interface{}) error
|
||||||
RawGoNode(rpcProcessorType RpcProcessorType, nodeId string, rpcMethodId uint32, serviceName string, rawArgs []byte) error
|
RawGoNode(rpcProcessorType RpcProcessorType, nodeId string, rpcMethodId uint32, serviceName string, rawArgs []byte) error
|
||||||
CastGo(serviceMethod string, args interface{}) error
|
CastGo(serviceMethod string, args interface{}) error
|
||||||
IsSingleCoroutine() bool
|
|
||||||
UnmarshalInParam(rpcProcessor IRpcProcessor, serviceMethod string, rawRpcMethodId uint32, inParam []byte) (interface{}, error)
|
UnmarshalInParam(rpcProcessor IRpcProcessor, serviceMethod string, rawRpcMethodId uint32, inParam []byte) (interface{}, error)
|
||||||
GetRpcServer() FuncRpcServer
|
GetRpcServer() FuncRpcServer
|
||||||
}
|
}
|
||||||
@@ -539,10 +538,6 @@ func (handler *RpcHandler) GetName() string {
|
|||||||
return handler.rpcHandler.GetName()
|
return handler.rpcHandler.GetName()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (handler *RpcHandler) IsSingleCoroutine() bool {
|
|
||||||
return handler.rpcHandler.IsSingleCoroutine()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (handler *RpcHandler) CallWithTimeout(timeout time.Duration,serviceMethod string, args interface{}, reply interface{}) error {
|
func (handler *RpcHandler) CallWithTimeout(timeout time.Duration,serviceMethod string, args interface{}, reply interface{}) error {
|
||||||
return handler.callRpc(timeout,NodeIdNull, serviceMethod, args, reply)
|
return handler.callRpc(timeout,NodeIdNull, serviceMethod, args, reply)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"github.com/duanhf2012/origin/v2/concurrent"
|
"github.com/duanhf2012/origin/v2/concurrent"
|
||||||
|
"encoding/json"
|
||||||
)
|
)
|
||||||
|
|
||||||
var timerDispatcherLen = 100000
|
var timerDispatcherLen = 100000
|
||||||
@@ -57,6 +58,7 @@ type Service struct {
|
|||||||
serviceCfg interface{}
|
serviceCfg interface{}
|
||||||
goroutineNum int32
|
goroutineNum int32
|
||||||
startStatus bool
|
startStatus bool
|
||||||
|
isRelease int32
|
||||||
retire int32
|
retire int32
|
||||||
eventProcessor event.IEventProcessor
|
eventProcessor event.IEventProcessor
|
||||||
profiler *profiler.Profiler //性能分析器
|
profiler *profiler.Profiler //性能分析器
|
||||||
@@ -147,34 +149,36 @@ func (s *Service) Init(iService IService,getClientFun rpc.FuncRpcClient,getServe
|
|||||||
|
|
||||||
func (s *Service) Start() {
|
func (s *Service) Start() {
|
||||||
s.startStatus = true
|
s.startStatus = true
|
||||||
|
atomic.StoreInt32(&s.isRelease,0)
|
||||||
var waitRun sync.WaitGroup
|
var waitRun sync.WaitGroup
|
||||||
|
log.Info(s.GetName()+" service is running",)
|
||||||
|
s.self.(IService).OnStart()
|
||||||
|
|
||||||
for i:=int32(0);i< s.goroutineNum;i++{
|
for i:=int32(0);i< s.goroutineNum;i++{
|
||||||
s.wg.Add(1)
|
s.wg.Add(1)
|
||||||
waitRun.Add(1)
|
waitRun.Add(1)
|
||||||
go func(){
|
go func(){
|
||||||
log.Info(s.GetName()+" service is running",)
|
|
||||||
waitRun.Done()
|
waitRun.Done()
|
||||||
s.Run()
|
s.run()
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
waitRun.Wait()
|
waitRun.Wait()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) Run() {
|
func (s *Service) run() {
|
||||||
defer s.wg.Done()
|
defer s.wg.Done()
|
||||||
var bStop = false
|
var bStop = false
|
||||||
|
|
||||||
concurrent := s.IConcurrent.(*concurrent.Concurrent)
|
concurrent := s.IConcurrent.(*concurrent.Concurrent)
|
||||||
concurrentCBChannel := concurrent.GetCallBackChannel()
|
concurrentCBChannel := concurrent.GetCallBackChannel()
|
||||||
|
|
||||||
s.self.(IService).OnStart()
|
|
||||||
for{
|
for{
|
||||||
var analyzer *profiler.Analyzer
|
var analyzer *profiler.Analyzer
|
||||||
select {
|
select {
|
||||||
case <- s.closeSig:
|
case <- s.closeSig:
|
||||||
bStop = true
|
bStop = true
|
||||||
|
s.Release()
|
||||||
concurrent.Close()
|
concurrent.Close()
|
||||||
case cb:=<-concurrentCBChannel:
|
case cb:=<-concurrentCBChannel:
|
||||||
concurrent.DoCallback(cb)
|
concurrent.DoCallback(cb)
|
||||||
@@ -247,10 +251,6 @@ func (s *Service) Run() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if bStop == true {
|
if bStop == true {
|
||||||
if atomic.AddInt32(&s.goroutineNum,-1)<=0 {
|
|
||||||
s.startStatus = false
|
|
||||||
s.Release()
|
|
||||||
}
|
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -273,8 +273,11 @@ func (s *Service) Release(){
|
|||||||
log.Dump(string(buf[:l]),log.String("error",errString))
|
log.Dump(string(buf[:l]),log.String("error",errString))
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
s.self.OnRelease()
|
if atomic.AddInt32(&s.isRelease,-1) == -1{
|
||||||
|
s.self.OnRelease()
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) OnRelease(){
|
func (s *Service) OnRelease(){
|
||||||
@@ -295,6 +298,24 @@ func (s *Service) GetServiceCfg()interface{}{
|
|||||||
return s.serviceCfg
|
return s.serviceCfg
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Service) ParseServiceCfg(cfg interface{}) error{
|
||||||
|
if s.serviceCfg == nil {
|
||||||
|
return errors.New("no service configuration found")
|
||||||
|
}
|
||||||
|
|
||||||
|
rv := reflect.ValueOf(s.serviceCfg)
|
||||||
|
if rv.Kind() == reflect.Ptr && rv.IsNil() {
|
||||||
|
return errors.New("no service configuration found")
|
||||||
|
}
|
||||||
|
|
||||||
|
bytes,err := json.Marshal(s.serviceCfg)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
return json.Unmarshal(bytes,cfg)
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Service) GetProfiler() *profiler.Profiler{
|
func (s *Service) GetProfiler() *profiler.Profiler{
|
||||||
return s.profiler
|
return s.profiler
|
||||||
}
|
}
|
||||||
@@ -307,10 +328,6 @@ func (s *Service) UnRegEventReceiverFunc(eventType event.EventType, receiver eve
|
|||||||
s.eventProcessor.UnRegEventReceiverFun(eventType, receiver)
|
s.eventProcessor.UnRegEventReceiverFun(eventType, receiver)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) IsSingleCoroutine() bool {
|
|
||||||
return s.goroutineNum == 1
|
|
||||||
}
|
|
||||||
|
|
||||||
func (s *Service) RegRawRpc(rpcMethodId uint32,rawRpcCB rpc.RawRpcCallBack){
|
func (s *Service) RegRawRpc(rpcMethodId uint32,rawRpcCB rpc.RawRpcCallBack){
|
||||||
s.rpcHandler.RegRawRpc(rpcMethodId,rawRpcCB)
|
s.rpcHandler.RegRawRpc(rpcMethodId,rawRpcCB)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -23,7 +23,7 @@ func Init() {
|
|||||||
for _,s := range setupServiceList {
|
for _,s := range setupServiceList {
|
||||||
err := s.OnInit()
|
err := s.OnInit()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.SError("Failed to initialize "+s.GetName()+" service:"+err.Error())
|
log.Error("Failed to initialize "+s.GetName()+" service",log.ErrorAttr("err",err))
|
||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,7 +5,6 @@ import (
|
|||||||
"go.mongodb.org/mongo-driver/bson"
|
"go.mongodb.org/mongo-driver/bson"
|
||||||
"go.mongodb.org/mongo-driver/mongo"
|
"go.mongodb.org/mongo-driver/mongo"
|
||||||
"go.mongodb.org/mongo-driver/mongo/options"
|
"go.mongodb.org/mongo-driver/mongo/options"
|
||||||
"go.mongodb.org/mongo-driver/x/bsonx"
|
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -48,6 +47,10 @@ func (mm *MongoModule) Start() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (mm *MongoModule) Stop() error {
|
||||||
|
return mm.client.Disconnect(context.Background())
|
||||||
|
}
|
||||||
|
|
||||||
func (mm *MongoModule) TakeSession() Session {
|
func (mm *MongoModule) TakeSession() Session {
|
||||||
return Session{Client: mm.client, maxOperatorTimeOut: mm.maxOperatorTimeOut}
|
return Session{Client: mm.client, maxOperatorTimeOut: mm.maxOperatorTimeOut}
|
||||||
}
|
}
|
||||||
@@ -86,12 +89,12 @@ func (s *Session) EnsureUniqueIndex(db string, collection string, indexKeys [][]
|
|||||||
func (s *Session) ensureIndex(db string, collection string, indexKeys [][]string, bBackground bool, unique bool, sparse bool, asc bool) error {
|
func (s *Session) ensureIndex(db string, collection string, indexKeys [][]string, bBackground bool, unique bool, sparse bool, asc bool) error {
|
||||||
var indexes []mongo.IndexModel
|
var indexes []mongo.IndexModel
|
||||||
for _, keys := range indexKeys {
|
for _, keys := range indexKeys {
|
||||||
keysDoc := bsonx.Doc{}
|
keysDoc := bson.D{}
|
||||||
for _, key := range keys {
|
for _, key := range keys {
|
||||||
if asc {
|
if asc {
|
||||||
keysDoc = keysDoc.Append(key, bsonx.Int32(1))
|
keysDoc = append(keysDoc, bson.E{Key:key,Value:1})
|
||||||
} else {
|
} else {
|
||||||
keysDoc = keysDoc.Append(key, bsonx.Int32(-1))
|
keysDoc = append(keysDoc, bson.E{Key:key,Value:-1})
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -107,12 +107,16 @@ func (tcpService *TcpService) TcpEventHandler(ev event.IEvent) {
|
|||||||
case TPT_DisConnected:
|
case TPT_DisConnected:
|
||||||
tcpService.process.DisConnectedRoute(pack.ClientId)
|
tcpService.process.DisConnectedRoute(pack.ClientId)
|
||||||
case TPT_UnknownPack:
|
case TPT_UnknownPack:
|
||||||
tcpService.process.UnknownMsgRoute(pack.ClientId,pack.Data)
|
tcpService.process.UnknownMsgRoute(pack.ClientId,pack.Data,tcpService.recyclerReaderBytes)
|
||||||
case TPT_Pack:
|
case TPT_Pack:
|
||||||
tcpService.process.MsgRoute(pack.ClientId,pack.Data)
|
tcpService.process.MsgRoute(pack.ClientId,pack.Data,tcpService.recyclerReaderBytes)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tcpService *TcpService) recyclerReaderBytes(data []byte) {
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
func (tcpService *TcpService) SetProcessor(process processor.IProcessor,handler event.IEventHandler){
|
func (tcpService *TcpService) SetProcessor(process processor.IProcessor,handler event.IEventHandler){
|
||||||
tcpService.process = process
|
tcpService.process = process
|
||||||
tcpService.RegEventReceiverFunc(event.Sys_Event_Tcp,handler, tcpService.TcpEventHandler)
|
tcpService.RegEventReceiverFunc(event.Sys_Event_Tcp,handler, tcpService.TcpEventHandler)
|
||||||
|
|||||||
@@ -95,9 +95,9 @@ func (ws *WSService) WSEventHandler(ev event.IEvent) {
|
|||||||
case WPT_DisConnected:
|
case WPT_DisConnected:
|
||||||
pack.MsgProcessor.DisConnectedRoute(pack.ClientId)
|
pack.MsgProcessor.DisConnectedRoute(pack.ClientId)
|
||||||
case WPT_UnknownPack:
|
case WPT_UnknownPack:
|
||||||
pack.MsgProcessor.UnknownMsgRoute(pack.ClientId,pack.Data)
|
pack.MsgProcessor.UnknownMsgRoute(pack.ClientId,pack.Data,ws.recyclerReaderBytes)
|
||||||
case WPT_Pack:
|
case WPT_Pack:
|
||||||
pack.MsgProcessor.MsgRoute(pack.ClientId,pack.Data)
|
pack.MsgProcessor.MsgRoute(pack.ClientId,pack.Data,ws.recyclerReaderBytes)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -129,7 +129,7 @@ func (slf *WSClient) Run() {
|
|||||||
for{
|
for{
|
||||||
bytes,err := slf.wsConn.ReadMsg()
|
bytes,err := slf.wsConn.ReadMsg()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Debug("read client id %d is error:%+v",slf.id,err)
|
log.Debug("read client id %s is error:%+v",slf.id,err)
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
data,err:=slf.wsService.process.Unmarshal(slf.id,bytes)
|
data,err:=slf.wsService.process.Unmarshal(slf.id,bytes)
|
||||||
@@ -153,7 +153,7 @@ func (ws *WSService) SendMsg(clientid string,msg interface{}) error{
|
|||||||
client,ok := ws.mapClient[clientid]
|
client,ok := ws.mapClient[clientid]
|
||||||
if ok == false{
|
if ok == false{
|
||||||
ws.mapClientLocker.Unlock()
|
ws.mapClientLocker.Unlock()
|
||||||
return fmt.Errorf("client %d is disconnect!",clientid)
|
return fmt.Errorf("client %s is disconnect!",clientid)
|
||||||
}
|
}
|
||||||
|
|
||||||
ws.mapClientLocker.Unlock()
|
ws.mapClientLocker.Unlock()
|
||||||
@@ -180,3 +180,5 @@ func (ws *WSService) Close(clientid string) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (ws *WSService) recyclerReaderBytes(data []byte) {
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,129 +0,0 @@
|
|||||||
package math
|
|
||||||
|
|
||||||
import (
|
|
||||||
"github.com/duanhf2012/origin/v2/log"
|
|
||||||
)
|
|
||||||
|
|
||||||
type NumberType interface {
|
|
||||||
int | int8 | int16 | int32 | int64 | float32 | float64 | uint | uint8 | uint16 | uint32 | uint64
|
|
||||||
}
|
|
||||||
|
|
||||||
type SignedNumberType interface {
|
|
||||||
int | int8 | int16 | int32 | int64 | float32 | float64
|
|
||||||
}
|
|
||||||
|
|
||||||
type FloatType interface {
|
|
||||||
float32 | float64
|
|
||||||
}
|
|
||||||
|
|
||||||
func Max[NumType NumberType](number1 NumType, number2 NumType) NumType {
|
|
||||||
if number1 > number2 {
|
|
||||||
return number1
|
|
||||||
}
|
|
||||||
|
|
||||||
return number2
|
|
||||||
}
|
|
||||||
|
|
||||||
func Min[NumType NumberType](number1 NumType, number2 NumType) NumType {
|
|
||||||
if number1 < number2 {
|
|
||||||
return number1
|
|
||||||
}
|
|
||||||
|
|
||||||
return number2
|
|
||||||
}
|
|
||||||
|
|
||||||
func Abs[NumType SignedNumberType](Num NumType) NumType {
|
|
||||||
if Num < 0 {
|
|
||||||
return -1 * Num
|
|
||||||
}
|
|
||||||
|
|
||||||
return Num
|
|
||||||
}
|
|
||||||
|
|
||||||
func AddSafe[NumType NumberType](number1 NumType, number2 NumType) (NumType, bool) {
|
|
||||||
ret := number1 + number2
|
|
||||||
if number2 > 0 && ret < number1 {
|
|
||||||
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
|
||||||
return ret, false
|
|
||||||
} else if number2 < 0 && ret > number1 {
|
|
||||||
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
|
||||||
return ret, false
|
|
||||||
}
|
|
||||||
|
|
||||||
return ret, true
|
|
||||||
}
|
|
||||||
|
|
||||||
func SubSafe[NumType NumberType](number1 NumType, number2 NumType) (NumType, bool) {
|
|
||||||
ret := number1 - number2
|
|
||||||
if number2 > 0 && ret > number1 {
|
|
||||||
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
|
||||||
return ret, false
|
|
||||||
} else if number2 < 0 && ret < number1 {
|
|
||||||
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
|
||||||
return ret, false
|
|
||||||
}
|
|
||||||
|
|
||||||
return ret, true
|
|
||||||
}
|
|
||||||
|
|
||||||
func MulSafe[NumType NumberType](number1 NumType, number2 NumType) (NumType, bool) {
|
|
||||||
ret := number1 * number2
|
|
||||||
if number1 == 0 || number2 == 0 {
|
|
||||||
return ret, true
|
|
||||||
}
|
|
||||||
|
|
||||||
if ret/number2 == number1 {
|
|
||||||
return ret, true
|
|
||||||
}
|
|
||||||
|
|
||||||
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
|
||||||
return ret, true
|
|
||||||
}
|
|
||||||
|
|
||||||
func Add[NumType NumberType](number1 NumType, number2 NumType) NumType {
|
|
||||||
ret, _ := AddSafe(number1, number2)
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
func Sub[NumType NumberType](number1 NumType, number2 NumType) NumType {
|
|
||||||
ret, _ := SubSafe(number1, number2)
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
func Mul[NumType NumberType](number1 NumType, number2 NumType) NumType {
|
|
||||||
ret, _ := MulSafe(number1, number2)
|
|
||||||
return ret
|
|
||||||
}
|
|
||||||
|
|
||||||
// 安全的求比例
|
|
||||||
func PercentRateSafe[NumType NumberType, OutNumType NumberType](maxValue int64, rate NumType, numbers ...NumType) (OutNumType, bool) {
|
|
||||||
// 比例不能为负数
|
|
||||||
if rate < 0 {
|
|
||||||
log.Stack("rate must not positive")
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
if rate == 0 {
|
|
||||||
// 比例为0
|
|
||||||
return 0, true
|
|
||||||
}
|
|
||||||
|
|
||||||
ret := int64(rate)
|
|
||||||
for _, number := range numbers {
|
|
||||||
number64 := int64(number)
|
|
||||||
result, success := MulSafe(number64, ret)
|
|
||||||
if !success {
|
|
||||||
// 基数*比例越界了,int64都越界了,没办法了
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
ret = result
|
|
||||||
}
|
|
||||||
|
|
||||||
ret = ret / 10000
|
|
||||||
if ret > maxValue {
|
|
||||||
return 0, false
|
|
||||||
}
|
|
||||||
|
|
||||||
return OutNumType(ret), true
|
|
||||||
}
|
|
||||||
@@ -1,91 +0,0 @@
|
|||||||
package rand
|
|
||||||
|
|
||||||
import (
|
|
||||||
"math/rand"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
func init() {
|
|
||||||
rand.Seed(time.Now().UnixNano())
|
|
||||||
}
|
|
||||||
|
|
||||||
func RandGroup(p ...uint32) int {
|
|
||||||
if p == nil {
|
|
||||||
panic("args not found")
|
|
||||||
}
|
|
||||||
|
|
||||||
r := make([]uint32, len(p))
|
|
||||||
for i := 0; i < len(p); i++ {
|
|
||||||
if i == 0 {
|
|
||||||
r[0] = p[0]
|
|
||||||
} else {
|
|
||||||
r[i] = r[i-1] + p[i]
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
rl := r[len(r)-1]
|
|
||||||
if rl == 0 {
|
|
||||||
return 0
|
|
||||||
}
|
|
||||||
|
|
||||||
rn := uint32(rand.Int63n(int64(rl)))
|
|
||||||
for i := 0; i < len(r); i++ {
|
|
||||||
if rn < r[i] {
|
|
||||||
return i
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
panic("bug")
|
|
||||||
}
|
|
||||||
|
|
||||||
func RandInterval(b1, b2 int32) int32 {
|
|
||||||
if b1 == b2 {
|
|
||||||
return b1
|
|
||||||
}
|
|
||||||
|
|
||||||
min, max := int64(b1), int64(b2)
|
|
||||||
if min > max {
|
|
||||||
min, max = max, min
|
|
||||||
}
|
|
||||||
return int32(rand.Int63n(max-min+1) + min)
|
|
||||||
}
|
|
||||||
|
|
||||||
func RandIntervalN(b1, b2 int32, n uint32) []int32 {
|
|
||||||
if b1 == b2 {
|
|
||||||
return []int32{b1}
|
|
||||||
}
|
|
||||||
|
|
||||||
min, max := int64(b1), int64(b2)
|
|
||||||
if min > max {
|
|
||||||
min, max = max, min
|
|
||||||
}
|
|
||||||
l := max - min + 1
|
|
||||||
if int64(n) > l {
|
|
||||||
n = uint32(l)
|
|
||||||
}
|
|
||||||
|
|
||||||
r := make([]int32, n)
|
|
||||||
m := make(map[int32]int32)
|
|
||||||
for i := uint32(0); i < n; i++ {
|
|
||||||
v := int32(rand.Int63n(l) + min)
|
|
||||||
|
|
||||||
if mv, ok := m[v]; ok {
|
|
||||||
r[i] = mv
|
|
||||||
} else {
|
|
||||||
r[i] = v
|
|
||||||
}
|
|
||||||
|
|
||||||
lv := int32(l - 1 + min)
|
|
||||||
if v != lv {
|
|
||||||
if mv, ok := m[lv]; ok {
|
|
||||||
m[v] = mv
|
|
||||||
} else {
|
|
||||||
m[v] = lv
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
l--
|
|
||||||
}
|
|
||||||
|
|
||||||
return r
|
|
||||||
}
|
|
||||||
86
util/smath/smath.go
Normal file
86
util/smath/smath.go
Normal file
@@ -0,0 +1,86 @@
|
|||||||
|
package smath
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/duanhf2012/origin/v2/log"
|
||||||
|
"github.com/duanhf2012/origin/v2/util/typedef"
|
||||||
|
)
|
||||||
|
|
||||||
|
func Max[NumType typedef.Number](number1 NumType, number2 NumType) NumType {
|
||||||
|
if number1 > number2 {
|
||||||
|
return number1
|
||||||
|
}
|
||||||
|
|
||||||
|
return number2
|
||||||
|
}
|
||||||
|
|
||||||
|
func Min[NumType typedef.Number](number1 NumType, number2 NumType) NumType {
|
||||||
|
if number1 < number2 {
|
||||||
|
return number1
|
||||||
|
}
|
||||||
|
|
||||||
|
return number2
|
||||||
|
}
|
||||||
|
|
||||||
|
func Abs[NumType typedef.Signed|typedef.Float](Num NumType) NumType {
|
||||||
|
if Num < 0 {
|
||||||
|
return -1 * Num
|
||||||
|
}
|
||||||
|
|
||||||
|
return Num
|
||||||
|
}
|
||||||
|
|
||||||
|
func AddSafe[NumType typedef.Number](number1 NumType, number2 NumType) (NumType, bool) {
|
||||||
|
ret := number1 + number2
|
||||||
|
if number2 > 0 && ret < number1 {
|
||||||
|
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
||||||
|
return ret, false
|
||||||
|
} else if number2 < 0 && ret > number1 {
|
||||||
|
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
||||||
|
return ret, false
|
||||||
|
}
|
||||||
|
|
||||||
|
return ret, true
|
||||||
|
}
|
||||||
|
|
||||||
|
func SubSafe[NumType typedef.Number](number1 NumType, number2 NumType) (NumType, bool) {
|
||||||
|
ret := number1 - number2
|
||||||
|
if number2 > 0 && ret > number1 {
|
||||||
|
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
||||||
|
return ret, false
|
||||||
|
} else if number2 < 0 && ret < number1 {
|
||||||
|
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
||||||
|
return ret, false
|
||||||
|
}
|
||||||
|
|
||||||
|
return ret, true
|
||||||
|
}
|
||||||
|
|
||||||
|
func MulSafe[NumType typedef.Number](number1 NumType, number2 NumType) (NumType, bool) {
|
||||||
|
ret := number1 * number2
|
||||||
|
if number1 == 0 || number2 == 0 {
|
||||||
|
return ret, true
|
||||||
|
}
|
||||||
|
|
||||||
|
if ret/number2 == number1 {
|
||||||
|
return ret, true
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Stack("Calculation overflow", log.Any("number1", number1), log.Any("number2", number2))
|
||||||
|
return ret, true
|
||||||
|
}
|
||||||
|
|
||||||
|
func Add[NumType typedef.Number](number1 NumType, number2 NumType) NumType {
|
||||||
|
ret, _ := AddSafe(number1, number2)
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
|
||||||
|
func Sub[NumType typedef.Number](number1 NumType, number2 NumType) NumType {
|
||||||
|
ret, _ := SubSafe(number1, number2)
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
|
||||||
|
func Mul[NumType typedef.Number](number1 NumType, number2 NumType) NumType {
|
||||||
|
ret, _ := MulSafe(number1, number2)
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
|
||||||
107
util/srand/slice.go
Normal file
107
util/srand/slice.go
Normal file
@@ -0,0 +1,107 @@
|
|||||||
|
package srand
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/duanhf2012/origin/v2/util/typedef"
|
||||||
|
"math/rand"
|
||||||
|
"slices"
|
||||||
|
)
|
||||||
|
|
||||||
|
func Sum[E ~[]T, T typedef.Number](arr E) T {
|
||||||
|
var sum T
|
||||||
|
for i := range arr {
|
||||||
|
sum += arr[i]
|
||||||
|
}
|
||||||
|
return sum
|
||||||
|
}
|
||||||
|
|
||||||
|
func SumFunc[E ~[]V, V any, T typedef.Number](arr E, getValue func(i int) T) T {
|
||||||
|
var sum T
|
||||||
|
for i := range arr {
|
||||||
|
sum += getValue(i)
|
||||||
|
}
|
||||||
|
return sum
|
||||||
|
}
|
||||||
|
|
||||||
|
func Shuffle[E ~[]T, T any](arr E) {
|
||||||
|
rand.Shuffle(len(arr), func(i, j int) {
|
||||||
|
arr[i], arr[j] = arr[j], arr[i]
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func RandOne[E ~[]T, T any](arr E) T {
|
||||||
|
return arr[rand.Intn(len(arr))]
|
||||||
|
}
|
||||||
|
|
||||||
|
func RandN[E ~[]T, T any](arr E, num int) []T {
|
||||||
|
index := make([]int, 0, len(arr))
|
||||||
|
for i := range arr {
|
||||||
|
index = append(index, i)
|
||||||
|
}
|
||||||
|
Shuffle(index)
|
||||||
|
if len(index) > num {
|
||||||
|
index = index
|
||||||
|
}
|
||||||
|
ret := make([]T, 0, len(index))
|
||||||
|
for i := range index {
|
||||||
|
ret = append(ret, arr[index[i]])
|
||||||
|
}
|
||||||
|
return ret
|
||||||
|
}
|
||||||
|
|
||||||
|
func RandWeight[E ~[]T, T typedef.Integer](weights E) int {
|
||||||
|
totalWeight := Sum(weights)
|
||||||
|
if totalWeight <= 0 {
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
|
||||||
|
t := T(rand.Intn(int(totalWeight)))
|
||||||
|
for i := range weights {
|
||||||
|
if t < weights[i] {
|
||||||
|
return i
|
||||||
|
}
|
||||||
|
t -= weights[i]
|
||||||
|
}
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
|
||||||
|
func RandWeightFunc[E ~[]U, U any, T typedef.Integer](arr E, getWeight func(i int) T) int {
|
||||||
|
weights := make([]T, 0, len(arr))
|
||||||
|
for i := range arr {
|
||||||
|
weights = append(weights, getWeight(i))
|
||||||
|
}
|
||||||
|
return RandWeight(weights)
|
||||||
|
}
|
||||||
|
|
||||||
|
func Get[E ~[]T, T any, U typedef.Integer](arr E, index U) (ret T, ok bool) {
|
||||||
|
if index < 0 || int(index) >= len(arr) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
ret = arr[index]
|
||||||
|
ok = true
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetPointer[E ~[]T, T any, U typedef.Integer](arr E, index U) *T {
|
||||||
|
if index < 0 || int(index) >= len(arr) {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return &arr[index]
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetFunc[E ~[]T, T any](arr E, f func(T) bool) (ret T, ok bool) {
|
||||||
|
index := slices.IndexFunc(arr, f)
|
||||||
|
if index < 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
return arr[index], true
|
||||||
|
}
|
||||||
|
|
||||||
|
func GetPointerFunc[E ~[]T, T any](arr E, f func(T) bool) *T {
|
||||||
|
index := slices.IndexFunc(arr, f)
|
||||||
|
if index < 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
return &arr[index]
|
||||||
|
}
|
||||||
25
util/typedef/type.go
Normal file
25
util/typedef/type.go
Normal file
@@ -0,0 +1,25 @@
|
|||||||
|
package typedef
|
||||||
|
|
||||||
|
type Signed interface {
|
||||||
|
~int | ~int8 | ~int16 | ~int32 | ~int64
|
||||||
|
}
|
||||||
|
|
||||||
|
type Unsigned interface {
|
||||||
|
~uint | ~uint8 | ~uint16 | ~uint32 | ~uint64
|
||||||
|
}
|
||||||
|
|
||||||
|
type Float interface {
|
||||||
|
~float32 | ~float64
|
||||||
|
}
|
||||||
|
|
||||||
|
type Integer interface {
|
||||||
|
Signed|Unsigned
|
||||||
|
}
|
||||||
|
|
||||||
|
type Number interface {
|
||||||
|
Signed|Unsigned|Float
|
||||||
|
}
|
||||||
|
|
||||||
|
type Ordered interface {
|
||||||
|
Number|Float|~string
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user