Documentation
¶
Overview ¶
Package router 提供多语义消息路由:把一条消息发给一组人(presence IDs)、 一个流(stream)的全部成员、或所有人,并支持一帧内多条消息攒批下发。
它是 pkg/stream.Broadcaster 的增强版:Broadcaster 只做"扇出给订阅者", Router 额外支持"按在场 ID 定点投递"和"按流投递(借助 presence.Tracker)", 以及攒批(SendDeferred)减少系统调用。
零值不可用,用 New 构造。Router 并发安全。
Index ¶
- type Forwarder
- type Message
- type Option
- type Router
- func (r *Router) FlushDeferred() int
- func (r *Router) QueueDeferred(targets []string, m Message)
- func (r *Router) SendToAll(allSessionIDs func() []string, m Message) int
- func (r *Router) SendToPresenceIDs(ids []presence.ID, m Message) int
- func (r *Router) SendToSessionIDs(sessionIDs []string, m Message) int
- func (r *Router) SendToStream(stream presence.Stream, m Message, includeHidden bool) int
- func (r *Router) SetTracker(t *presence.Tracker)
- type Sink
- type SinkRegistry
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Forwarder ¶
Forwarder 把消息转发到其他节点。跨节点投递时调用:把目标 ids + 消息交给业务, 业务通过 RPC/消息总线送达目标节点,目标节点再用本地 Router 投递。
type Option ¶
type Option func(*config)
Option 配置 Router。
func WithForwarder ¶
WithForwarder 设置跨节点转发器。配合 WithLocalNode 使用。
func WithLocalNode ¶
WithLocalNode 设置本节点名。非空时启用跨节点路由:id.Node != localNode 的 presence 走 Forwarder 转发。空(默认)= 不启用跨节点,所有 id 视为本地。
type Router ¶
type Router struct {
// contains filtered or unexported fields
}
Router 按 presence IDs / stream / 全员 投递消息,并支持攒批。
func New ¶
func New(registry SinkRegistry, tracker *presence.Tracker, opts ...Option) *Router
New 创建 Router。tracker 可为 nil(此时 SendToStream 不可用,会返回 0)。
func (*Router) FlushDeferred ¶
FlushDeferred 投递所有攒批消息并清空队列。返回总投递数。 对同一 session 的多条消息会按顺序投递(保持 FIFO)。
func (*Router) QueueDeferred ¶
QueueDeferred 把一条消息加入攒批队列,稍后由 FlushDeferred 一次性投递。 适合一帧内产生多条消息时减少重复 Lookup/系统调用。 targets 为目标 session IDs;为空时消息不会投递(需广播请用 SendToAll)。
func (*Router) SendToAll ¶
SendToAll 广播给全员。需要业务提供全员 session ID 列表(由 AllSessionIDs 返回)。 这是保守设计:避免 Router 隐式持有全局 session 表导致耦合。
func (*Router) SendToPresenceIDs ¶
SendToPresenceIDs 把消息投递给指定 presence IDs 对应的会话。 返回成功投递的数量(含本地 Sink 成功 + 远端 Forwarder 报告的成功数)。 id.Node 为空或等于 localNode 视为本地;否则走 Forwarder(未配置则丢弃)。
func (*Router) SendToSessionIDs ¶
SendToSessionIDs 是 SendToPresenceIDs 的简化版:直接按 session ID 投递。
func (*Router) SendToStream ¶
SendToStream 把消息投递给某流的全部在场成员。 依赖 presence.Tracker;未配置时返回 0。includeHidden 控制是否含隐藏成员。 跨节点成员(id.Node != localNode)走 Forwarder 转发。
func (*Router) SetTracker ¶
SetTracker 在构造后注入/替换 tracker(用于解耦初始化顺序)。
type Sink ¶
Sink 是消息的最终投递目标,通常是一个 WebSocket/gRPC 会话的 Send 封装。 返回 false 表示接收者已不可达(下线),Router 会从 presence 清理(若配置了 Tracker)。
type SinkRegistry ¶
SinkRegistry 把 presence.ID 解析到 Sink。业务维护 session -> Sink 的映射。 只负责本节点 session:跨节点 session( id.Node != localNode)由 Forwarder 处理。