SLG 游戏的聊天系统是玩家社交的核心——世界频道刷屏、联盟频道指挥、私聊交易,每秒可能有上百条消息。slg-go 的聊天服务用 WebSocket 做实时推送,gRPC 做跨服务调用,Hub 模式管理频道订阅。

双协议设计

协议 用途 特点
WebSocket 客户端实时收发消息 长连接、低延迟、支持推送
gRPC Logic 服务调用聊天 API 短连接、可靠、用于系统消息

客户端通过 WebSocket 发送和接收消息,Logic 服务通过 gRPC 发送系统公告、战斗结果通知等。

Hub:频道管理中心

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
// internal/chat/hub.go
type Hub struct {
channels map[string]*Channel // 频道 ID → Channel
history *RecentCache // 最近消息缓存
filter *ChatFilter // 敏感词过滤器
mu sync.RWMutex
}

type Channel struct {
ChannelType uint8
ChannelID string
members map[int64]chan<- *Message // playerID → 消息通道
}

type Message struct {
MsgID int64
Channel string
SenderID int64
SenderName string
TargetID int64
Content string
SentAt time.Time
}

Hub 是整个聊天系统的核心——管理频道、维护订阅关系、分发消息。每个 Channel 用 map[int64]chan<- *Message 存储订阅者,消息通过 channel 异步推送给每个 WebSocket 连接。

频道类型

1
2
3
4
5
6
const (
ChannelWorld uint8 = 1 // 世界频道
ChannelAlliance uint8 = 2 // 联盟频道
ChannelPrivate uint8 = 3 // 私聊
ChannelSystem uint8 = 4 // 系统消息
)

世界频道所有玩家可见,联盟频道只限联盟成员,私聊只限两个玩家。GetChannel() 按需创建频道:

1
2
3
4
5
6
7
8
9
10
11
12
13
func (h *Hub) GetChannel(channelID string) *Channel {
h.mu.Lock()
defer h.mu.Unlock()
if ch, ok := h.channels[channelID]; ok {
return ch
}
ch := &Channel{
ChannelID: channelID,
members: make(map[int64]chan<- *Message),
}
h.channels[channelID] = ch
return ch
}

WebSocket 处理:读写分离

每个 WebSocket 连接启动两个 goroutine——读循环和写循环:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// cmd/chat/ws.go — 简化示意
func newWebSocketHandler(hub *chat.Hub, tokenSecret string) http.Handler {
return websocket.Handler(func(ws *websocket.Conn) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

msgCh := make(chan *chat.Message, 128)
controlCh := make(chan wsEnvelope, 64)

// 写循环:从 msgCh/controlCh 读取消息,发送给客户端
go func() {
defer writeWG.Done()
for {
select {
case <-ctx.Done():
return
case msg := <-msgCh:
websocket.JSON.Send(ws, ...)
case env := <-controlCh:
websocket.JSON.Send(ws, ...)
}
}
}()

// 读循环:从客户端接收命令
for {
var cmd wsCommand
websocket.JSON.Receive(ws, &cmd)
// 处理命令...
}
})
}

读循环在主 goroutine 中执行(阻塞读),写循环在子 goroutine 中执行。msgCh 缓冲 128 条消息——如果客户端消费太慢,超出缓冲的消息会被丢弃并触发告警。

消息广播:异步分发

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
func (h *Hub) Broadcast(channelID string, msg *Message) {
h.mu.RLock()
ch, ok := h.channels[channelID]
h.mu.RUnlock()
if !ok { return }

// 缓存最近消息
h.history.Add(channelID, msg)

// 异步推送给每个订阅者
for _, msgCh := range ch.members {
select {
case msgCh <- msg:
// 发送成功
default:
// 消费者太慢,丢弃消息
}
}
}

select + default 非阻塞发送——如果某个消费者的 channel 已满,直接丢弃而不是阻塞整个广播。这保证了慢消费者不会拖垮整个频道。

敏感词过滤:发送前拦截

消息发送前经过 ChatFilter(Aho-Corasick 自动机)过滤:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
func (h *Hub) Send(senderID int64, channelID, content string) (*Message, error) {
// 敏感词过滤
filtered := h.filter.Filter(content)
if filtered != content {
// 有过滤,可以记录日志或告警
}
msg := &Message{
MsgID: generateID(),
Channel: channelID,
SenderID: senderID,
Content: filtered,
SentAt: time.Now(),
}
h.Broadcast(channelID, msg)
return msg, nil
}

过滤在 Hub 层而非 WebSocket 层,保证无论是 WebSocket 还是 gRPC 发送的消息都会被过滤。

消息持久化:RecentCache

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// internal/chat/cache/recent_cache.go
type RecentCache struct {
mu sync.RWMutex
messages map[string][]*Message // channelID → 最近 N 条消息
limit int // 每频道最大缓存条数(默认 1000)
}

func (c *RecentCache) Add(channelID string, msg *Message) {
c.mu.Lock()
defer c.mu.Unlock()
c.messages[channelID] = append(c.messages[channelID], msg)
if len(c.messages[channelID]) > c.limit {
c.messages[channelID] = c.messages[channelID][len(c.messages[channelID])-c.limit:]
}
}

内存缓存最近 1000 条消息,新玩家加入频道时可以获取历史消息。不需要持久化到数据库——重启后消息丢失对 SLG 聊天来说可以接受。

gRPC 接口:跨服务调用

Logic 服务通过 gRPC 发送系统消息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// internal/chat/handler/chat_grpc.go
type ChatGrpcHandler struct {
hub *chat.Hub
}

func (h *ChatGrpcHandler) SendSystemMessage(ctx context.Context, req *pb.SendSystemMessageRequest) (*pb.SendSystemMessageResponse, error) {
msg := &chat.Message{
Channel: req.ChannelId,
SenderID: 0, // 系统 ID
SenderName: "System",
Content: req.Content,
SentAt: time.Now(),
}
h.hub.Broadcast(req.ChannelId, msg)
return &pb.SendSystemMessageResponse{MsgId: msg.MsgID}, nil
}

gRPC 接口复用同一个 Hub,保证消息格式和过滤逻辑一致。

慢消费者保护

当某个 WebSocket 连接消费太慢时,msgCh(缓冲 128 条)会溢出。Hub 的 Broadcast 方法用 select + default 丢弃溢出消息,同时触发告警:

1
2
3
4
5
6
7
8
9
10
// hub.go
func (h *Hub) shouldLogDropWarn(now time.Time) bool {
h.mu.Lock()
defer h.mu.Unlock()
if now.Sub(h.lastDropWarn) < h.dropWarnCooldown {
return false
}
h.lastDropWarn = now
return true
}

告警有冷却时间(默认 30 秒),避免日志被刷爆。运维可以根据告警频率判断是否需要扩容或优化客户端消费逻辑。

总结

设计决策 原因
WebSocket + gRPC 双协议 客户端用 WebSocket 实时推送,Logic 用 gRPC 系统调用
Hub 模式管理频道 中心化订阅关系,避免消息路由复杂度
channel 异步推送 非阻塞分发,慢消费者不拖垮整体
内存缓存历史消息 1000 条/频道,新玩家加入可获取上下文
Aho-Corasick 过滤 Hub 层统一过滤,WebSocket 和 gRPC 都经过

聊天系统不是”加个 WebSocket 就行”的事。频道管理、慢消费者保护、敏感词过滤、历史消息,每个环节都要考虑。