Documentation
¶
Overview ¶
Package hub 提供中继网络的核心类型与接口。
Package hub 提供中继网络的核心类型与 DHT 接口。 DHT 接口定义见 core.go;内置实现见 internal/memdht.go。
Package hub 提供星型中继网络的 Hub 端实现。
Hub 维护节点路由表(NodeID → mux.Mux), 为中继请求提供目标节点查找和转发能力。
Index ¶
- Constants
- Variables
- func ComputeRegisterProof(skHex, nodeID string, ts int64, nonce string) (string, error)
- func NewRegisterFrame(nodeID, ak, proof string, ts int64, nonce string, meta Meta, caps ...string) []byte
- func NewRegisterNonce() string
- func NormalizeEndpoints(hubURL, serverURL string) (httpBase, wsURL string, err error)
- func ParseDiscNodeID(nodeID string) (string, bool)
- func ParseRegisterAck(ackStr string) (secret string, err error)
- func RegisterDHT(name string, dht DHT, priority int)
- func RestoreFromSnapshot(mrt *MeshRouteTable, snap *Snapshot)
- func RestoreSignalQueue(q *SignalQueue, msgs []MessageSnap)
- type AccessKey
- type Authenticator
- type DHT
- type DialRequest
- type DialResultFrame
- type HubServer
- type HubSignaler
- func (s *HubSignaler) SendAnswer(to string, sdp string) error
- func (s *HubSignaler) SendOffer(to string, sdp string) error
- func (s *HubSignaler) SetAccessKeySecret(sk string)
- func (s *HubSignaler) SetContext(ctx context.Context)
- func (s *HubSignaler) SetHTTPClient(hc *http.Client)
- func (s *HubSignaler) WaitAnswer(ctx context.Context) (string, string, error)
- func (s *HubSignaler) WaitOffer(ctx context.Context) (string, string, error)
- type MeshRouteTable
- func (mrt *MeshRouteTable) Add(mesh string, info NodeInfo, svcs []Service)
- func (mrt *MeshRouteTable) AddNode(mesh string, id NodeID, m *mux.Mux)
- func (mrt *MeshRouteTable) AllMeshes() []string
- func (mrt *MeshRouteTable) Has(id NodeID) bool
- func (mrt *MeshRouteTable) List(mesh string) []NodeInfo
- func (mrt *MeshRouteTable) ListServices(mesh string) []NodeService
- func (mrt *MeshRouteTable) Lookup(id NodeID) *mux.Mux
- func (mrt *MeshRouteTable) LookupInfo(id NodeID) (NodeInfo, bool)
- func (mrt *MeshRouteTable) MeshOf(id NodeID) string
- func (mrt *MeshRouteTable) NodeCount(mesh string) int
- func (mrt *MeshRouteTable) Remove(id NodeID) bool
- func (mrt *MeshRouteTable) RemoveIfOwned(id NodeID, m *mux.Mux) bool
- func (mrt *MeshRouteTable) SetOnChange(fn func())
- func (mrt *MeshRouteTable) SetRemoveHook(fn func(NodeID))
- func (mrt *MeshRouteTable) Table(mesh string) *RouteTable
- type MessageSnap
- type Meta
- type NodeID
- type NodeInfo
- type NodeService
- type NodeSnap
- type PeerInfo
- type Persister
- type RegisterFrame
- type RouteTable
- func (rt *RouteTable) Add(id NodeID, m *mux.Mux)
- func (rt *RouteTable) AddWithInfo(info NodeInfo)
- func (rt *RouteTable) AddWithInfoAndServices(info NodeInfo, svcs []Service)
- func (rt *RouteTable) ClearServices(id NodeID)
- func (rt *RouteTable) Has(id NodeID) bool
- func (rt *RouteTable) HasService(id NodeID, name string) bool
- func (rt *RouteTable) List() []NodeInfo
- func (rt *RouteTable) ListServices() []NodeService
- func (rt *RouteTable) Lookup(id NodeID) *mux.Mux
- func (rt *RouteTable) LookupInfo(id NodeID) (NodeInfo, bool)
- func (rt *RouteTable) LookupService(name string) (NodeID, Service, bool)
- func (rt *RouteTable) NodeCount() int
- func (rt *RouteTable) Remove(id NodeID) bool
- func (rt *RouteTable) RemoveIfOwned(id NodeID, m *mux.Mux) bool
- func (rt *RouteTable) ServiceHosts(name string) []NodeID
- func (rt *RouteTable) ServicesOf(id NodeID) []Service
- func (rt *RouteTable) SetRemoveHook(fn func(NodeID))
- func (rt *RouteTable) SetServices(id NodeID, svcs []Service)
- type Service
- type SignalKind
- type SignalMsg
- type SignalQueue
- func (q *SignalQueue) Confirm(target, id string) bool
- func (q *SignalQueue) Peek(target string, kind SignalKind) *SignalMsg
- func (q *SignalQueue) Pop(target string) *SignalMsg
- func (q *SignalQueue) PopKind(target string, kind SignalKind) *SignalMsg
- func (q *SignalQueue) Purge(peerID string)
- func (q *SignalQueue) Push(m SignalMsg) error
- func (q *SignalQueue) Total() int
- func (q *SignalQueue) Wait(ctx context.Context, target string) error
- type Snapshot
Constants ¶
const ( RegisterAckOK = "REG_OK" RegisterAckErr = "REG_ERR:" )
注册 ACK 帧常量(xfer 层一条消息)。导出供 sclient 注册后等待。 节点声明 per-node-secret 能力(见 CapabilityPerNodeSecret)时,REG_OK 携带 独立 secret,线上格式为 "REG_OK:<base64url secret>";未声明时为纯 "REG_OK"。
const ( DialResultOK = "ok" DialResultError = "error" )
DialResultFrame 取值。
const CapabilityPerNodeSecret = "per-node-secret"
CapabilityPerNodeSecret 是节点声明"希望获得 per-node 独立 secret"的能力标志。 hub 注册成功后生成 32B 随机 secret 存入 NodeInfo.Secret 并随 REG_OK 下发, 供后续批次(B3 服务端校验 / B2 客户端携带)的信令身份校验使用(I1)。
const DiscPrefix = "disc"
DiscPrefix 是 mesh 自动对等发现临时节点 ID 的基前缀(不含 '-',AutoRegister 的 Prefix 参数用)。完整节点 ID 形如 disc-<base>-<unixnano>,单一来源避免各处硬编码 失配(hub 注册校验与 mesh 拨号共用)。
const PollTimeout = 25 * time.Second
PollTimeout 是 poll 长轮询的单次最长等待。
依赖约束:服务端写入超时(server_timeouts.write,默认 30s)必须大于 PollTimeout,否则长轮询会被服务端在到达 PollTimeout 前掐断(I32/S9)。 客户端单次 poll 的 HTTP 超时(signaling_client.go 的 httpClient.Timeout, 60s)必须大于 PollTimeout + 网络余量,避免客户端先于服务端超时。
const RegisterProofV2Context = "sproxy-hub-register/v2"
RegisterProofV2Context 是 hub 节点注册 HMAC 证明的上下文串(v2)。 证明 = HMAC-SHA256(SK, RegisterProofV2Context + "\n" + nodeID + "\n" + ts_ms + "\n" + nonce)。 v2 引入 ts+nonce 防重放(M-6):捕获的注册帧无法在新窗口内重放。
Variables ¶
var ( // ErrSignalQueueFull 表示全局信令消息总量已达上限,新消息被拒绝(I9)。 ErrSignalQueueFull = errors.New("信令队列已满") // ErrSignalPerSenderCap 表示同一发送方到目标收件箱的未消费消息已达上限(I9)。 ErrSignalPerSenderCap = errors.New("信令发送方配额已满") )
var DHTRegistry = plugin.New("dht", DHT(newMemoryDHT()))
DHTRegistry 是节点发现实现的插件注册表。 默认使用内置的内存 DHT 实现(newMemoryDHT)。
var ErrInvalidAccessKey = errors.New("invalid access key")
ErrInvalidAccessKey 是 AK 未命中或 accessKeys 未配置(fail-closed)时返回的哨兵错误。
var ErrInvalidAccessKeyProof = errors.New("invalid access key proof")
ErrInvalidAccessKeyProof 是 HMAC proof 校验失败时返回的哨兵错误。
var ErrRegisterRejected = errors.New("注册失败")
ErrRegisterRejected 表示 hub 通过注册 ACK 明确拒绝本次注册(鉴权/格式错误)。 客户端 isTerminalRelayError 用 errors.Is 判定,可穿透任意 %w 包装,避免文案 改写或包装后终态判定静默失效。
var ErrReplayRegisterNonce = errors.New("replayed register nonce")
ErrReplayRegisterNonce 是注册 nonce 已被使用(重放)时返回的哨兵错误。
var ErrStaleRegisterProof = errors.New("stale register proof")
ErrStaleRegisterProof 是注册证明 ts 超出新鲜度窗口时返回的哨兵错误(防重放)。
Functions ¶
func ComputeRegisterProof ¶
ComputeRegisterProof 计算节点注册的 HMAC 证明(v2): HMAC-SHA256(SK, "sproxy-hub-register/v2\n"+nodeID+"\n"+ts+"\n"+nonce)。 skHex 为 64 hex 字符的 SproxySig AccessKeySecret;nodeID 绑定节点防串用; ts 为 unix 毫秒、nonce 为一次性随机串,共同防重放。返回 64 hex 字符。
func NewRegisterFrame ¶
func NewRegisterFrame(nodeID, ak, proof string, ts int64, nonce string, meta Meta, caps ...string) []byte
NewRegisterFrame 构建注册帧。当无 meta/ak/proof/ts/nonce/caps 时退化为裸 nodeID, 保证与旧版 hub(仅接收裸 nodeID)兼容。 ts/nonce 为注册证明的防重放字段(M-6,与 ComputeRegisterProof 参数一致)。 caps 为可选变参:声明能力(如 CapabilityPerNodeSecret)后 hub 回 REG_OK 携带 per-node secret(I1);现有调用不传 caps 时行为不变。
func NewRegisterNonce ¶
func NewRegisterNonce() string
NewRegisterNonce 生成注册证明用的一次性 nonce(16 字节随机数,hex 编码)。 crypto/rand 失败概率极低;万一失败回退到时间戳+进程随机(仍满足一次性语义, 且 hub 端按 (ak, nonce) 去重 + 窗口校验,空/重复 nonce 均 fail-closed)。
func NormalizeEndpoints ¶
NormalizeEndpoints 将 hub 地址归一为信令 HTTP 基址与注册 WS 端点:
- httpBase(信令 post/poll 用,http(s)://host[:port],剥 path);
- wsURL(自动注册用,ws(s)://host[:port]/ws)。
hubURL 接受 http(s):// 或 ws(s)://(含 /ws 等 path);空串回退 serverURL。 畸形 URL / 未知 scheme 显式报错。
func ParseDiscNodeID ¶
ParseDiscNodeID 从 mesh discovery 临时节点 ID(disc-<base>-<suffix>)解析真实 node-id base(base 可含 '-';suffix 为 16 位随机 hex 尾段,保证并发拨号不碰撞)。 hub 注册校验与 mesh accept 侧解析共用同一实现,保证"hub 已验证的 base"与 "accept 侧使用的 base"一致。取最后一个 '-' 之前的全部为 base(base 可含 '-'); 尾段须为合法 hex(对齐 newTempSuffix 格式),避免歧义。
func ParseRegisterAck ¶
ParseRegisterAck 解析注册 ACK 帧(xfer 层一条消息)。返回节点 per-node secret:
- 纯 "REG_OK"(未声明能力)→ ("", nil);
- "REG_OK:<base64url secret>"(声明 per-node-secret 能力)→ (secret, nil);
- "REG_ERR:<reason>"(hub 明确拒绝)→ ("", 包装 ErrRegisterRejected 的错误);
- 未知响应 → ("", 包装 ErrRegisterRejected 的错误)。
用前缀匹配而非精确比较,避免声明能力后收 "REG_OK:<secret>" 被误判未知响应 导致 relay start 终止(B1 复检 bug 回归锁)。
func RegisterDHT ¶
RegisterDHT 注册一个命名的 DHT 实现到 DHTRegistry(cmd/sproxy 装配 Kademlia 用)。 priority > 0 时优先于内置内存 DHT(Active 返回最高优先级实现);同名重复注册 以后一次为准(注册顺序不变)。
func RestoreFromSnapshot ¶
func RestoreFromSnapshot(mrt *MeshRouteTable, snap *Snapshot)
RestoreFromSnapshot 将快照恢复到空的 MeshRouteTable(幂等调用会重复写入相似数据, 仅应在启动时对空表调用一次)。 恢复的节点没有在线 Mux(nil),等待客户端重连后由 HandleConn 重写 info(含新 Mux 与新 secret),保证重启后节点身份/服务不丢失而连接面在线。
func RestoreSignalQueue ¶
func RestoreSignalQueue(q *SignalQueue, msgs []MessageSnap)
RestoreSignalQueue 将快照的信令收件箱恢复到队列。 恢复时即过滤已过期消息(M3):过期死信不重投递(与惰性过期语义一致,重启后 不再投递已过期的消息),且不把过期消息计入 q.total——否则 q.total 在下次 惰性清理前被高估,白白占用全局配额。
Types ¶
type Authenticator ¶
type Authenticator struct {
// contains filtered or unexported fields
}
Authenticator 验证节点注册的 SproxySig AccessKey + HMAC proof(v2,含 ts/nonce 防重放)。 fail-closed:accessKeys 为空时拒绝所有注册。
func NewAuthenticator ¶
func NewAuthenticator(accessKeys []AccessKey) *Authenticator
NewAuthenticator 创建鉴权器。accessKeys 为空时创建 fail-closed 鉴权器(拒绝所有注册)。
func (*Authenticator) Authenticate ¶
func (a *Authenticator) Authenticate(ak, proof, nodeID string, ts int64, nonce string) error
Authenticate 验证 AK、HMAC proof 与新鲜度:
- 空 accessKeys → ErrInvalidAccessKey(fail-closed)
- 遍历按 Key constant-time 匹配;未命中 → ErrInvalidAccessKey
- |now−ts| > registerProofMaxAge → ErrStaleRegisterProof(防重放)
- nonce 已用过 → ErrReplayRegisterNonce(防重放)
- 命中 → 用 Secret 重算 ComputeRegisterProof(secret, nodeID, ts, nonce),与 proof constant-time 比对; 不匹配 → ErrInvalidAccessKeyProof;匹配 → nil
type DHT ¶
type DHT interface {
// Register 将本节点注册到 DHT 网络。
Register(ctx context.Context, node PeerInfo) error
// Lookup 按节点 ID 查找目标节点信息。
Lookup(ctx context.Context, nodeID string) (PeerInfo, error)
// GetClosestNodes 返回距离目标 ID 最近的 N 个节点。
// 距离算法由各实现定义(内置:词法排序;Kademlia:XOR 距离)。
GetClosestNodes(ctx context.Context, nodeID string, n int) ([]PeerInfo, error)
// Remove 从发现表中移除指定节点(节点断开/管理端踢出时调用,防幽灵节点残留)。
// 节点不存在时视为成功(幂等)。
Remove(ctx context.Context, nodeID string) error
// Bootstrap 连接到已知种子节点,加入 DHT 网络。
Bootstrap(ctx context.Context, seeds []string) error
// Close 退出 DHT 网络,释放资源。
Close() error
}
DHT 定义节点发现的最低接口。 内置实现是简单的线程安全内存 map;ext/kad 提供完整的 Kademlia。
type DialRequest ¶
type DialRequest struct {
Dial string `json:"dial,omitempty"` // 目标叶子出站连接的 TCP 地址
}
DialRequest 是发往叶子出口节点的导出结构(JSON), portal/relay 收到后向 addr 发起出站连接,随后进入字节中转。
type DialResultFrame ¶
type DialResultFrame struct {
DialResult string `json:"dial_result"`
Message string `json:"message,omitempty"`
}
DialResultFrame 是叶子出口拨号结果回帧(I27)。 hub 的 /api/relay/stream 在写 200 前读取 [4B len][DialResultFrame JSON], 据此决定返回 200(ok)或 502/504(error/超时)。 仅当叶子以 ServeOptions.DialResultFrames 模式(relay start 经 hub 中继)运行时 回写;webrtc 直连(p2p listen)不回写,避免结果帧污染数据流。
type HubServer ¶
type HubServer struct {
// contains filtered or unexported fields
}
HubServer 是 Hub 端的节点收口服务:接收 xfer.Conn 连接, 完成注册/鉴权并写入 RouteTable;连接断开后自动移除。
注册协议:节点连接建立后,在 xfer 层(mux 之上)直接发送一条 注册消息(JSON 或裸 nodeID),Hub 侧通过 conn.Receive 读取,不占用 mux 的流创建——这样后续 Tunnel.Serve 的 Open/Accept 与 TCP 流中继 可复用同一条 mux 而互不抢 acceptCh。
func NewHubServer ¶
func NewHubServer(rt *MeshRouteTable, auth *Authenticator, logger *slog.Logger, maxConns ...int) *HubServer
NewHubServer 创建节点收口服务。auth 为 nil 时视为 fail-closed(拒绝所有注册), 防止调用方"遗漏传 auth"走向开放注册(M-7)。 rt 为每 mesh 独立路由表的聚合(M-9),按注册 AK 解析的 mesh 分表隔离。 maxConns 为可选变参:传 >0 的值表示 Hub 同时处理的连接数上限(I30), 不传或 <=0 表示无上限。
func (*HubServer) HandleConn ¶
HandleConn 接收一个已建立的节点连接,注册并维护其生命周期。 阻塞直到连接断开后返回。
协议顺序(重要):先读取 xfer 层注册帧(节点连接后发送的一条消息, 不经过 mux 流),完成注册/鉴权后再创建 mux 走流协议。 若先建 mux,其 readLoop 会与注册帧的 Receive 竞争同一连接。
type HubSignaler ¶
type HubSignaler struct {
// contains filtered or unexported fields
}
HubSignaler 是通过 hub 的 /api/signal/* 桥实现 SDP 交换的网络信令器。 其方法签名与 xfer/ext/webrtc 的 Signaler 接口一致,可被直接传给 webrtc.DialWithSignaler / webrtc.ListenWithSignaler(无需 import webrtc)。
跨机器场景:两端的 SDP 都经 hub 队列存转;即使 Mac→hub 链路抖动, 信令体量 K 级、可重试,也能最终完成。
func NewHubSignaler ¶
func NewHubSignaler(baseURL, accessKey, nodeID string, secret ...string) *HubSignaler
NewHubSignaler 创建经 hub 信令桥的 Signaler。 accessKey 是本 mesh 的 SproxySig AccessKey(公开标识;配合 SetAccessKeySecret 签名)。 nodeID 是本节点在 hub 上注册的节点 ID(信令来源,服务端校验已注册)。 secret 为可选变参:传入 per-node secret 后 post/poll 携带 X-Node-Secret 头。
func (*HubSignaler) SendAnswer ¶
func (s *HubSignaler) SendAnswer(to string, sdp string) error
SendAnswer 向对端 to 发送 Answer SDP。
func (*HubSignaler) SendOffer ¶
func (s *HubSignaler) SendOffer(to string, sdp string) error
SendOffer 向对端 to 发送 Offer SDP。
func (*HubSignaler) SetAccessKeySecret ¶
func (s *HubSignaler) SetAccessKeySecret(sk string)
SetAccessKeySecret 设置 SproxySig AccessKeySecret(本地密钥,仅计算签名)。 空则信令请求不签名(无认证开发环境)。
func (*HubSignaler) SetContext ¶
func (s *HubSignaler) SetContext(ctx context.Context)
SetContext 注入 base context(I7)。SendOffer/SendAnswer/post 在注入后 使用该 ctx 而非 context.Background(),受调用方(mesh/p2p 命令 ctx)取消控制。 不调用则保持 context.Background()(向后兼容)。
func (*HubSignaler) SetHTTPClient ¶
func (s *HubSignaler) SetHTTPClient(hc *http.Client)
SetHTTPClient 注入自定义 http.Client(TLS 配置 / 超时)。nil 忽略(保留默认)。 对齐 SetContext 模式(I7):不调用则保持默认 &http.Client{Timeout:60s}(向后兼容)。 供 sclient --insecure 场景注入跳过证书校验的 client(自签 wss hub 信令链路)。
func (*HubSignaler) WaitAnswer ¶
WaitAnswer 阻塞等待发给本节点(s.nodeID)的 Answer SDP。 返回发送方节点 ID 与 SDP;调用方可校验 from 是否为目标对端。
type MeshRouteTable ¶
type MeshRouteTable struct {
// contains filtered or unexported fields
}
MeshRouteTable 是每 mesh 独立 RouteTable 的聚合:map[mesh]*RouteTable + nodeID→mesh 映射。 转发/列表/信令按 mesh 隔离;无 mesh("")为默认表,行为等价单 mesh。
func NewMeshRouteTable ¶
func NewMeshRouteTable() *MeshRouteTable
NewMeshRouteTable 创建每 mesh 独立路由表的聚合。
func (*MeshRouteTable) Add ¶
func (mrt *MeshRouteTable) Add(mesh string, info NodeInfo, svcs []Service)
Add 写入对应表(info.Mesh = mesh)。 若同一节点 ID 此前属于另一 mesh(跨 mesh 同名重注册),先从旧 mesh 表移除, 维持 nodeMesh 单一归属的隔离不变量(节点名不跨 mesh 共享)。 节点集合变更后在锁外触发 onChange 回调(持久化)。
func (*MeshRouteTable) AddNode ¶
func (mrt *MeshRouteTable) AddNode(mesh string, id NodeID, m *mux.Mux)
AddNode 低层节点注册(Add 的变体,不携带服务宣告)。
func (*MeshRouteTable) AllMeshes ¶
func (mrt *MeshRouteTable) AllMeshes() []string
AllMeshes 返回所有已存在的 mesh 列表(debug/管理用)。
func (*MeshRouteTable) Has ¶
func (mrt *MeshRouteTable) Has(id NodeID) bool
Has 检查节点是否存在(按 nodeMesh 自动定位所属 mesh)。
func (*MeshRouteTable) List ¶
func (mrt *MeshRouteTable) List(mesh string) []NodeInfo
List 返回某 mesh 的节点列表(/api/hub/nodes 用)。
func (*MeshRouteTable) ListServices ¶
func (mrt *MeshRouteTable) ListServices(mesh string) []NodeService
ListServices 返回某 mesh 的服务宣告(mesh 选路用)。
func (*MeshRouteTable) Lookup ¶
func (mrt *MeshRouteTable) Lookup(id NodeID) *mux.Mux
Lookup 按 ID 查找节点的 Mux 连接(转发用):查 nodeMesh[id] → 对应表 Lookup。 未找到时返回 nil。节点只属于一个 mesh,跨 mesh 的 ID 不可达。
func (*MeshRouteTable) LookupInfo ¶
func (mrt *MeshRouteTable) LookupInfo(id NodeID) (NodeInfo, bool)
LookupInfo 按 ID 查找节点的扩展信息(NodeInfo 含 Mesh,供 mesh 校验)。
func (*MeshRouteTable) MeshOf ¶
func (mrt *MeshRouteTable) MeshOf(id NodeID) string
MeshOf 返回节点所属 mesh(未注册返回空串)。
func (*MeshRouteTable) NodeCount ¶
func (mrt *MeshRouteTable) NodeCount(mesh string) int
NodeCount 返回某 mesh 的节点数(metrics 用)。
func (*MeshRouteTable) Remove ¶
func (mrt *MeshRouteTable) Remove(id NodeID) bool
Remove 按 ID 移除节点(自动按 nodeMesh 定位所属 mesh)。 移除成功后删除 nodeMesh 映射并清空该 mesh 的空表。 节点集合变更后在锁外触发 onChange 回调(持久化)。
func (*MeshRouteTable) RemoveIfOwned ¶
func (mrt *MeshRouteTable) RemoveIfOwned(id NodeID, m *mux.Mux) bool
RemoveIfOwned 仅当节点 ID 当前绑定到给定 mux(即本连接)时才从对应 mesh 移除。 防止旧连接断开时误删新注册的同名节点(stale identity 防护)。 节点集合变更后在锁外触发 onChange 回调(持久化)。
func (*MeshRouteTable) SetOnChange ¶
func (mrt *MeshRouteTable) SetOnChange(fn func())
SetOnChange 注册任意节点注册/移除后的回调(Add / Remove / RemoveIfOwned 成功路径)。 供 hub 状态持久化在节点集合变更后落盘。传 nil 清除回调。 回调在锁外同步调用,必须快速返回(不做阻塞 I/O),如需阻塞请自行 go。
func (*MeshRouteTable) SetRemoveHook ¶
func (mrt *MeshRouteTable) SetRemoveHook(fn func(NodeID))
SetRemoveHook 为每个内部 RouteTable 挂 onRemove 回调(SignalBroker 收件箱清理)。 已存在的表立即生效;之后惰性新建的表(Table)自动带上该回调。
func (*MeshRouteTable) Table ¶
func (mrt *MeshRouteTable) Table(mesh string) *RouteTable
Table 获取某 mesh 的路由表,不存在时惰性创建。 创建时若已注册 SetRemoveHook,新表同步挂上(收件箱清理全覆盖)。
type MessageSnap ¶
MessageSnap 是一个 peer 的信令收件箱快照。
func SnapshotSignalQueue ¶
func SnapshotSignalQueue(q *SignalQueue) []MessageSnap
SnapshotSignalQueue 捕获 SignalQueue 全部 per-peer 收件箱为 MessageSnap 序列。 M3:快照只保留未过期消息——队列的惰性过期(compactExpiredLocked)只在 Push/Peek/ Pop 时发生,一个无消息的空 poll 不会触发过期,若镜像原样复制会把已过期消息 残留在持久化文件里直到下次写盘。这里在生成快照时就按 signalMsgTTL 过滤, 保证**任何**持久化镜像(onChange / FlushSignal / 停服 Flush)都不含过期死信。
type Meta ¶
type Meta struct {
Addr string `json:"addr,omitempty"` // 节点自身可达地址(直连地址,可为空)
Services []Service `json:"services,omitempty"` // 本地服务宣告(本地监听或出口可达目标)
Tags []string `json:"tags,omitempty"` // 可选标签(如 "exit"、"trusted")
// RealNodeID 是 mesh 自动对等发现临时注册(disc-<base>-<nano>)代表的本节点真实
// node-id。hub 强制校验(见 registerNode):base 必须等于 RealNodeID 且持有该真实
// 节点 per-node secret 的 HMAC 证明(防冒充他人污染链路池)。
RealNodeID string `json:"real_node_id,omitempty"`
// RealNodeProof 是 HMAC-SHA256(真实节点 per-node secret, RealNodeID) 的 hex。
// 仅 mesh node 的 discovery 临时注册携带;hub 端用存储的真实节点 secret 复核。
RealNodeProof string `json:"real_node_proof,omitempty"`
}
Meta 是节点注册时宣告的附加信息,供 mesh 选路使用。
type NodeInfo ¶
type NodeInfo struct {
ID NodeID
Mux *mux.Mux
Connected time.Time // 连接时间
Addr string // 远端地址
Secret string // per-node 独立 secret(仅节点声明 per-node-secret 能力时下发;不落日志)
RealNodeID string // mesh discovery 临时注册(disc-)代表的本节点真实 node-id(hub 校验后记录)
Mesh string // 节点所属 mesh(由注册 AK 解析;"" 为默认 mesh)
}
NodeInfo 包含已注册节点的信息。
type NodeService ¶
NodeService 是 ListServices 返回的一个服务条目(节点 + 服务)。
type NodeSnap ¶
type NodeSnap struct {
ID NodeID `json:"id"`
Mesh string `json:"mesh,omitempty"`
Addr string `json:"addr,omitempty"`
Secret string `json:"secret,omitempty"`
RealNodeID string `json:"real_node_id,omitempty"`
Connected time.Time `json:"connected"`
Services []Service `json:"services,omitempty"`
}
NodeSnap 是持久化快照中单个节点的离线表示。 Mux 字段(在线连接)不持久化——重启后节点需重新建立连接。
type PeerInfo ¶
type PeerInfo struct {
ID string
Addrs []string // xfer transport 地址列表,如 ["tcp://192.168.1.1:9000"]
Meta map[string]string
}
PeerInfo 描述 DHT 中一个发现节点的网络位置信息。
type Persister ¶
type Persister struct {
// contains filtered or unexported fields
}
Persister 将 Hub 状态快照原子写入单个 JSON 文件,并以去抖方式合并高频变更。
并发语义:变更入口 Schedule 把最新的 snapshotFn 排入 pending 并启动/重置去抖周期, 到期后只执行**最后一次**排队的 snapshotFn(中间变更合并为一落盘)。Flush 同步 执行 pending(进程优雅停服前调用,确保最后一次变更不丢失)。线程安全。
注:本实现允许多个排队的 closure 逐个执行(非严格跳过),去抖合并的是 "多次 Schedule 中最后一次"——每次 timer 触发只执行当时的 pending。高扇出场景 (多个连接同时注册)由 200ms 窗口自然合并。
func (*Persister) Flush ¶
Flush 同步执行当前排队的快照(若存在)。用于进程优雅停服前确保状态不丢失; 无 pending 且 curr 为 nil 时是 no-op(返回 nil)。返回落盘错误(I2/M8:不静默吞)。 快照生成与落盘在同一临界区(I1)。
func (*Persister) FlushFn ¶
FlushFn 同步执行给定 snapshotFn 并落盘(curr nil 语义)。用于服务端在信令 变更后立即持久化当前收件箱状态。返回落盘错误(M8:不再以 bool 掩盖失败)。 快照生成与落盘在同一临界区(I1)。
func (*Persister) Load ¶
Load 读取快照文件并解码。
- 文件不存在(未持久化过)→ 返回空快照、无错误;
- 文件存在但损坏/非法 JSON,或超出 maxSnapshotBytes → 记录 warn、返回空快照、 无错误(hub 启动不因持久化文件损坏而失败,也不 panic);
- 其余 I/O 错误(如权限不足)→ 返回 error,由调用方决定是否中止。
type RegisterFrame ¶
type RegisterFrame struct {
NodeID string `json:"node_id"`
// Token 已废弃:不再用于准入(保留字段避免破坏旧客户端 JSON)。
Token string `json:"token,omitempty"`
// AccessKey 是 SproxySig 准入 AccessKey(与 access_keys 配置一致)。
AccessKey string `json:"access_key,omitempty"`
// AccessKeyProof 是 ComputeRegisterProof 输出(HMAC-SHA256 证明持有 SK)。
AccessKeyProof string `json:"access_key_proof,omitempty"`
// TS / Nonce 是注册证明的防重放字段(M-6):TS 为 unix 毫秒、Nonce 为一次性随机串,
// 均参与 ComputeRegisterProof 签名;hub 校验 TS 新鲜度 + nonce 去重。
TS int64 `json:"ts,omitempty"`
Nonce string `json:"nonce,omitempty"`
Meta Meta `json:"meta"`
Capabilities []string `json:"capabilities,omitempty"`
}
RegisterFrame 是节点连接后的注册帧(JSON)。 向后兼容:若首个流上收到的是非 JSON 裸字符串(旧版仅发 nodeID), 则等价于仅携带 NodeID 且无 token 的注册帧。
Capabilities 是节点声明的能力标志列表(可扩展:未来新增能力直接追加 字符串常量,hub 端用 hasCapability 判断,旧 hub/旧客户端忽略未知项)。
type RouteTable ¶
type RouteTable struct {
// contains filtered or unexported fields
}
RouteTable 是线程安全的节点路由表。
func (*RouteTable) Add ¶
func (rt *RouteTable) Add(id NodeID, m *mux.Mux)
Add 注册一个节点。如果节点 ID 已存在,先关闭旧连接再替换。
func (*RouteTable) AddWithInfo ¶
func (rt *RouteTable) AddWithInfo(info NodeInfo)
AddWithInfo 注册节点并保存扩展信息。
func (*RouteTable) AddWithInfoAndServices ¶
func (rt *RouteTable) AddWithInfoAndServices(info NodeInfo, svcs []Service)
AddWithInfoAndServices 原子地注册节点并写入服务宣告。 与分两次调用 AddWithInfo + SetServices 相比,消除了重连时短暂残留 旧服务宣告的非原子窗口(S4)。空/ nil svcs 等价于清除该节点的旧宣告。
func (*RouteTable) ClearServices ¶
func (rt *RouteTable) ClearServices(id NodeID)
ClearServices 清除节点宣告的服务(节点断开时调用)。
func (*RouteTable) HasService ¶
func (rt *RouteTable) HasService(id NodeID, name string) bool
HasService 返回节点是否宣告了指定名称的服务。
func (*RouteTable) ListServices ¶
func (rt *RouteTable) ListServices() []NodeService
ListServices 返回所有节点宣告的服务,按 (node, name) 稳定排序(I3)。 客户端据此确定性选路;多节点同名服务保持多候选(failover 语义不破坏)。
func (*RouteTable) Lookup ¶
func (rt *RouteTable) Lookup(id NodeID) *mux.Mux
Lookup 按 ID 查找节点的 Mux 连接。 未找到时返回 nil。
func (*RouteTable) LookupInfo ¶
func (rt *RouteTable) LookupInfo(id NodeID) (NodeInfo, bool)
LookupInfo 按 ID 查找节点的扩展信息。 同时确认 nodes 与 info 两表均存在该节点;节点不存在时返回 false。
func (*RouteTable) LookupService ¶
func (rt *RouteTable) LookupService(name string) (NodeID, Service, bool)
LookupService 返回宣告了指定服务名的节点与其服务地址。
func (*RouteTable) Remove ¶
func (rt *RouteTable) Remove(id NodeID) bool
Remove 移除一个节点。返回是否真正移除(节点存在)。 移除成功后在锁外触发 onRemove 回调(若注册)。
func (*RouteTable) RemoveIfOwned ¶
func (rt *RouteTable) RemoveIfOwned(id NodeID, m *mux.Mux) bool
RemoveIfOwned 仅当该节点 ID 当前绑定到给定 mux(即本连接)时才移除。 防止旧连接断开时误删新注册的同名节点(stale identity 防护)。 返回是否真正移除。移除成功后在锁外触发 onRemove 回调(若注册)。
func (*RouteTable) ServiceHosts ¶
func (rt *RouteTable) ServiceHosts(name string) []NodeID
ServiceHosts 返回宣告了指定服务名的节点 ID 列表。
func (*RouteTable) ServicesOf ¶
func (rt *RouteTable) ServicesOf(id NodeID) []Service
ServicesOf 返回指定节点宣告的服务列表。
func (*RouteTable) SetRemoveHook ¶
func (rt *RouteTable) SetRemoveHook(fn func(NodeID))
SetRemoveHook 注册节点移除回调:节点真正从路由表移除(Remove / RemoveIfOwned 成功)后、在锁外调用一次 fn(id)。传 nil 清除回调。回调应快速返回(不做阻塞 I/O),需异步自行 go。
func (*RouteTable) SetServices ¶
func (rt *RouteTable) SetServices(id NodeID, svcs []Service)
SetServices 记录节点宣告的服务。
type Service ¶
type Service struct {
Name string `json:"name"` // 服务名(mesh connect 的目标)
Addr string `json:"addr"` // 服务地址(127.0.0.1:22、或出口节点可达的 192.x.x.x:22)
Layer string `json:"layer,omitempty"` // 用途标注(debug/info)
}
Service 描述一个可被 mesh 寻址的服务。
type SignalKind ¶
type SignalKind string
SignalKind 标识信令消息类型。
const ( SignalOffer SignalKind = "offer" SignalAnswer SignalKind = "answer" SignalCandidate SignalKind = "candidate" )
type SignalMsg ¶
type SignalMsg struct {
ID string `json:"id,omitempty"` // 消息去重 ID(服务端 Push 时赋,见 I10)
Kind SignalKind `json:"kind"`
From string `json:"from"` // 发送方 peer
To string `json:"to"` // 目标 peer
SDP string `json:"sdp,omitempty"`
Cand string `json:"cand,omitempty"` // TODO(S8): candidate 端点在服务端无生产发送方,随 B3 删路由后一并移除
At int64 `json:"at"` // 毫秒时间戳(服务端设置;TTL 判定依据)
}
SignalMsg 是经 hub 存转的一条信令消息。
type SignalQueue ¶
type SignalQueue struct {
// contains filtered or unexported fields
}
SignalQueue 是 per-peer 的信令收件箱(有界队列)。 hub 收到 offer/answer/candidate 后 push 到目标 peer 的队列; 目标 peer 通过 poll 长轮询取走。
资源边界:per-peer 队列有上限(maxSignalInbox),并设全局消息总量上限 (maxSignalTotal)与 per-sender 未消费上限(maxSignalPerSender), 防止大量空转 peer 的收件箱无界累积内存、以及单 sender 灌流逐出目标消息(I9)。
func (*SignalQueue) Confirm ¶
func (q *SignalQueue) Confirm(target, id string) bool
Confirm 删除 target 收件箱中 ID 匹配的消息(Peek 后确认消费)。 返回是否成功删除;未找到(已被并发 poll 取走)返回 false。
func (*SignalQueue) Peek ¶
func (q *SignalQueue) Peek(target string, kind SignalKind) *SignalMsg
Peek 非阻塞查看 target 收件箱里第一条指定 kind 的消息(不删除,返回副本)。 配合 Confirm 实现「Encode 成功后才消费」的可靠 poll(I5)。
func (*SignalQueue) Pop ¶
func (q *SignalQueue) Pop(target string) *SignalMsg
Pop 非阻塞取走 target 收件箱里的一条消息(任意 kind);无则返回 nil。
func (*SignalQueue) PopKind ¶
func (q *SignalQueue) PopKind(target string, kind SignalKind) *SignalMsg
PopKind 非阻塞取走 target 收件箱里第一条指定 kind 的消息;无匹配返回 nil。 不匹配 kind 的消息保留在队列,交给对应 kind 的消费者(I9)—— 破坏性消费不再吞掉其他 kind 的信令。
func (*SignalQueue) Purge ¶
func (q *SignalQueue) Purge(peerID string)
Purge 清空 peer 的收件箱与 waiter(节点下线时调用,I6)。 幂等:peer 不存在时是 no-op。close-broadcast 唤醒阻塞中的 Wait(P1-8), 使其不再搁浅至 ctx 截止(最长 25s);唤醒后 Pop 为空,调用方自行收尾。 每通道只 close 一次:此处 close 后立即 delete,后续 Wait 新建通道。
func (*SignalQueue) Push ¶
func (q *SignalQueue) Push(m SignalMsg) error
Push 投递一条消息到目标 peer 的收件箱。 全局总量达上限返回 ErrSignalQueueFull;同一 sender 在目标收件箱的未消费 消息达上限返回 ErrSignalPerSenderCap(两者都不入队,供调用方感知并回 429)。 Push 会为消息赋去重 ID(若为空)并确保 At 时间戳存在(TTL 判定依据)。
func (*SignalQueue) Total ¶
func (q *SignalQueue) Total() int
Total 返回当前全局积压消息总数(含所有收件箱)。 供上层(server 联动测试 / 运维观察)确认节点下线 Purge 后配额已释放。
type Snapshot ¶
type Snapshot struct {
Nodes []NodeSnap `json:"nodes,omitempty"`
Messages []MessageSnap `json:"messages,omitempty"`
}
Snapshot 是 Hub 状态的完整快照:节点注册 + 信令收件箱。
func SnapshotRouteTable ¶
func SnapshotRouteTable(mrt *MeshRouteTable) *Snapshot
SnapshotRouteTable 捕获 MeshRouteTable 全部节点注册(含 mesh、服务、secret), 生成 Snapshot 的 Nodes 部分。Mux(在线连接)不捕获——重启后节点重连建立。