router

package
v0.3.6 Latest Latest
Warning

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

Go to latest
Published: Aug 6, 2026 License: Apache-2.0 Imports: 2 Imported by: 0

Documentation

Overview

Package router 提供多语义消息路由:把一条消息发给一组人(presence IDs)、 一个流(stream)的全部成员、或所有人,并支持一帧内多条消息攒批下发。

它是 pkg/stream.Broadcaster 的增强版:Broadcaster 只做"扇出给订阅者", Router 额外支持"按在场 ID 定点投递"和"按流投递(借助 presence.Tracker)", 以及攒批(SendDeferred)减少系统调用。

零值不可用,用 New 构造。Router 并发安全。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Forwarder

type Forwarder interface {
	Forward(node string, ids []presence.ID, m Message) int
}

Forwarder 把消息转发到其他节点。跨节点投递时调用:把目标 ids + 消息交给业务, 业务通过 RPC/消息总线送达目标节点,目标节点再用本地 Router 投递。

type Message

type Message struct {
	Data     []byte
	Reliable bool // 可靠投递:对慢接收者阻塞而非丢弃(若 Sink 支持)
}

Message 是一条待路由的消息:任意负载 + 可靠性标记。 业务自行约定 Data 的编码(JSON/protobuf),Router 不关心。

type Option

type Option func(*config)

Option 配置 Router。

func WithForwarder

func WithForwarder(f Forwarder) Option

WithForwarder 设置跨节点转发器。配合 WithLocalNode 使用。

func WithLocalNode

func WithLocalNode(name string) Option

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

func (r *Router) FlushDeferred() int

FlushDeferred 投递所有攒批消息并清空队列。返回总投递数。 对同一 session 的多条消息会按顺序投递(保持 FIFO)。

func (*Router) QueueDeferred

func (r *Router) QueueDeferred(targets []string, m Message)

QueueDeferred 把一条消息加入攒批队列,稍后由 FlushDeferred 一次性投递。 适合一帧内产生多条消息时减少重复 Lookup/系统调用。 targets 为目标 session IDs;为空时消息不会投递(需广播请用 SendToAll)。

func (*Router) SendToAll

func (r *Router) SendToAll(allSessionIDs func() []string, m Message) int

SendToAll 广播给全员。需要业务提供全员 session ID 列表(由 AllSessionIDs 返回)。 这是保守设计:避免 Router 隐式持有全局 session 表导致耦合。

func (*Router) SendToPresenceIDs

func (r *Router) SendToPresenceIDs(ids []presence.ID, m Message) int

SendToPresenceIDs 把消息投递给指定 presence IDs 对应的会话。 返回成功投递的数量(含本地 Sink 成功 + 远端 Forwarder 报告的成功数)。 id.Node 为空或等于 localNode 视为本地;否则走 Forwarder(未配置则丢弃)。

func (*Router) SendToSessionIDs

func (r *Router) SendToSessionIDs(sessionIDs []string, m Message) int

SendToSessionIDs 是 SendToPresenceIDs 的简化版:直接按 session ID 投递。

func (*Router) SendToStream

func (r *Router) SendToStream(stream presence.Stream, m Message, includeHidden bool) int

SendToStream 把消息投递给某流的全部在场成员。 依赖 presence.Tracker;未配置时返回 0。includeHidden 控制是否含隐藏成员。 跨节点成员(id.Node != localNode)走 Forwarder 转发。

func (*Router) SetTracker

func (r *Router) SetTracker(t *presence.Tracker)

SetTracker 在构造后注入/替换 tracker(用于解耦初始化顺序)。

type Sink

type Sink func(m Message) bool

Sink 是消息的最终投递目标,通常是一个 WebSocket/gRPC 会话的 Send 封装。 返回 false 表示接收者已不可达(下线),Router 会从 presence 清理(若配置了 Tracker)。

type SinkRegistry

type SinkRegistry interface {
	Lookup(sessionID string) Sink
}

SinkRegistry 把 presence.ID 解析到 Sink。业务维护 session -> Sink 的映射。 只负责本节点 session:跨节点 session( id.Node != localNode)由 Forwarder 处理。

Jump to

Keyboard shortcuts

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