rts-sim-testing-service/message_server/ms_manage/manage.go
walker eb86063724 重构消息服务代码结构
实现部分信号布置图的状态消息采集(其他都待重构)
2023-10-26 16:41:18 +08:00

73 lines
1.5 KiB
Go

package ms_manage
import (
"context"
"log/slog"
"runtime/debug"
"time"
apiproto "joylink.club/bj-rtsts-server/grpcproto"
"joylink.club/bj-rtsts-server/message_server/ms_api"
)
type MsgServer struct {
ms_api.IMsgServer
ctx context.Context
cancelFn context.CancelFunc
}
// 消息服务管理map
var servers map[string]*MsgServer = make(map[string]*MsgServer)
// 注册消息服务
func Register(server ms_api.IMsgServer) *MsgServer {
ms := &MsgServer{
IMsgServer: server,
}
ctx, cancelFn := context.WithCancel(context.Background())
ms.ctx = ctx
ms.cancelFn = cancelFn
go run(ms)
servers[server.GetChannel()] = ms
return ms
}
// 注销消息服务
func Unregister(server ms_api.IMsgServer) {
if server == nil {
return
}
s := servers[server.GetChannel()]
s.cancelFn()
delete(servers, server.GetChannel())
}
// 消息服务运行
func run(server *MsgServer) {
defer func() {
if err := recover(); err != nil {
slog.Error("消息服务运行异常", "channel", server.GetChannel(), "error", err, "stack", string(debug.Stack()))
debug.PrintStack()
}
}()
for {
select {
case <-server.ctx.Done():
slog.Info("消息服务退出", "channel", server.GetChannel())
return
default:
}
topicMsgs, err := server.OnTick()
if err != nil {
slog.Error("消息服务构建定时发送消息错误", "channel", server.GetChannel(), "error", err)
continue
}
if len(topicMsgs) > 0 {
for _, msg := range topicMsgs {
apiproto.PublishMsg(msg.Channel, msg.Data)
}
}
time.Sleep(server.GetInterval())
}
}