hub

package
v0.0.0-...-a1c7e93 Latest Latest
Warning

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

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

Documentation

Overview

Package hub 提供中继网络的核心类型与接口。

Package hub 提供中继网络的核心类型与 DHT 接口。 DHT 接口定义见 core.go;内置实现见 internal/memdht.go。

Package hub 提供星型中继网络的 Hub 端实现。

Hub 维护节点路由表(NodeID → mux.Mux), 为中继请求提供目标节点查找和转发能力。

Index

Constants

View Source
const (
	RegisterAckOK  = "REG_OK"
	RegisterAckErr = "REG_ERR:"
)

注册 ACK 帧常量(xfer 层一条消息)。导出供 sclient 注册后等待。 节点声明 per-node-secret 能力(见 CapabilityPerNodeSecret)时,REG_OK 携带 独立 secret,线上格式为 "REG_OK:<base64url secret>";未声明时为纯 "REG_OK"。

View Source
const (
	DialResultOK    = "ok"
	DialResultError = "error"
)

DialResultFrame 取值。

View Source
const CapabilityPerNodeSecret = "per-node-secret"

CapabilityPerNodeSecret 是节点声明"希望获得 per-node 独立 secret"的能力标志。 hub 注册成功后生成 32B 随机 secret 存入 NodeInfo.Secret 并随 REG_OK 下发, 供后续批次(B3 服务端校验 / B2 客户端携带)的信令身份校验使用(I1)。

View Source
const DiscPrefix = "disc"

DiscPrefix 是 mesh 自动对等发现临时节点 ID 的基前缀(不含 '-',AutoRegister 的 Prefix 参数用)。完整节点 ID 形如 disc-<base>-<unixnano>,单一来源避免各处硬编码 失配(hub 注册校验与 mesh 拨号共用)。

View Source
const PollTimeout = 25 * time.Second

PollTimeout 是 poll 长轮询的单次最长等待。

依赖约束:服务端写入超时(server_timeouts.write,默认 30s)必须大于 PollTimeout,否则长轮询会被服务端在到达 PollTimeout 前掐断(I32/S9)。 客户端单次 poll 的 HTTP 超时(signaling_client.go 的 httpClient.Timeout, 60s)必须大于 PollTimeout + 网络余量,避免客户端先于服务端超时。

View Source
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

View Source
var (
	// ErrSignalQueueFull 表示全局信令消息总量已达上限,新消息被拒绝(I9)。
	ErrSignalQueueFull = errors.New("信令队列已满")
	// ErrSignalPerSenderCap 表示同一发送方到目标收件箱的未消费消息已达上限(I9)。
	ErrSignalPerSenderCap = errors.New("信令发送方配额已满")
)
View Source
var DHTRegistry = plugin.New("dht", DHT(newMemoryDHT()))

DHTRegistry 是节点发现实现的插件注册表。 默认使用内置的内存 DHT 实现(newMemoryDHT)。

View Source
var ErrInvalidAccessKey = errors.New("invalid access key")

ErrInvalidAccessKey 是 AK 未命中或 accessKeys 未配置(fail-closed)时返回的哨兵错误。

View Source
var ErrInvalidAccessKeyProof = errors.New("invalid access key proof")

ErrInvalidAccessKeyProof 是 HMAC proof 校验失败时返回的哨兵错误。

View Source
var ErrRegisterRejected = errors.New("注册失败")

ErrRegisterRejected 表示 hub 通过注册 ACK 明确拒绝本次注册(鉴权/格式错误)。 客户端 isTerminalRelayError 用 errors.Is 判定,可穿透任意 %w 包装,避免文案 改写或包装后终态判定静默失效。

View Source
var ErrReplayRegisterNonce = errors.New("replayed register nonce")

ErrReplayRegisterNonce 是注册 nonce 已被使用(重放)时返回的哨兵错误。

View Source
var ErrStaleRegisterProof = errors.New("stale register proof")

ErrStaleRegisterProof 是注册证明 ts 超出新鲜度窗口时返回的哨兵错误(防重放)。

Functions

func ComputeRegisterProof

func ComputeRegisterProof(skHex, nodeID string, ts int64, nonce string) (string, error)

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

func NormalizeEndpoints(hubURL, serverURL string) (httpBase, wsURL string, err error)

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

func ParseDiscNodeID(nodeID string) (string, bool)

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

func ParseRegisterAck(ackStr string) (secret string, err error)

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

func RegisterDHT(name string, dht DHT, priority int)

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 AccessKey

type AccessKey struct {
	Key    string
	Secret string
}

AccessKey 是 hub 准入用的 SproxySig 凭据(hub 包自建,勿 import pkg/server)。

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 与新鲜度:

  1. 空 accessKeys → ErrInvalidAccessKey(fail-closed)
  2. 遍历按 Key constant-time 匹配;未命中 → ErrInvalidAccessKey
  3. |now−ts| > registerProofMaxAge → ErrStaleRegisterProof(防重放)
  4. nonce 已用过 → ErrReplayRegisterNonce(防重放)
  5. 命中 → 用 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。

func NewDHT

func NewDHT() DHT

NewDHT 创建一个新的内存 DHT 实现,返回 DHT 接口。

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

func (s *HubServer) HandleConn(ctx context.Context, conn xfer.Conn) error

HandleConn 接收一个已建立的节点连接,注册并维护其生命周期。 阻塞直到连接断开后返回。

协议顺序(重要):先读取 xfer 层注册帧(节点连接后发送的一条消息, 不经过 mux 流),完成注册/鉴权后再创建 mux 走流协议。 若先建 mux,其 readLoop 会与注册帧的 Receive 竞争同一连接。

func (*HubServer) SetDHT

func (s *HubServer) SetDHT(dht DHT)

SetDHT 注入节点发现表(DHT)。nil 清除(恢复不启用 DHT 候选)。 须在服务器开始处理连接前调用。

func (*HubServer) TryHandleConn

func (s *HubServer) TryHandleConn(ctx context.Context, conn xfer.Conn) bool

TryHandleConn 非阻塞获取一个连接名额;成功时启动 goroutine 调用 HandleConn 处理连接, 并在处理结束后释放名额。信号量已满时返回 false,由调用方负责关闭 conn。 maxConns 未配置(nil)时始终接受,退化为无条件并发处理。

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

func (s *HubSignaler) WaitAnswer(ctx context.Context) (string, string, error)

WaitAnswer 阻塞等待发给本节点(s.nodeID)的 Answer SDP。 返回发送方节点 ID 与 SDP;调用方可校验 from 是否为目标对端。

func (*HubSignaler) WaitOffer

func (s *HubSignaler) WaitOffer(ctx context.Context) (string, string, error)

WaitOffer 阻塞等待发给本节点(s.nodeID)的 Offer SDP。 返回发送方节点 ID 与 SDP。任何已注册节点发来的 offer 都接受 (listener 无法预知拨号方,身份由服务端已校验 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

type MessageSnap struct {
	Peer string      `json:"peer"`
	Msgs []SignalMsg `json:"msgs"`
}

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 NodeID

type NodeID string

NodeID 是节点唯一标识符。

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

type NodeService struct {
	Node    NodeID
	Service Service
}

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 NewPersister

func NewPersister(path string) *Persister

NewPersister 创建指向 path 的持久化器。

func (*Persister) Flush

func (p *Persister) Flush(curr *Snapshot) error

Flush 同步执行当前排队的快照(若存在)。用于进程优雅停服前确保状态不丢失; 无 pending 且 curr 为 nil 时是 no-op(返回 nil)。返回落盘错误(I2/M8:不静默吞)。 快照生成与落盘在同一临界区(I1)。

func (*Persister) FlushFn

func (p *Persister) FlushFn(fn snapshotFn) error

FlushFn 同步执行给定 snapshotFn 并落盘(curr nil 语义)。用于服务端在信令 变更后立即持久化当前收件箱状态。返回落盘错误(M8:不再以 bool 掩盖失败)。 快照生成与落盘在同一临界区(I1)。

func (*Persister) Load

func (p *Persister) Load() (*Snapshot, error)

Load 读取快照文件并解码。

  • 文件不存在(未持久化过)→ 返回空快照、无错误;
  • 文件存在但损坏/非法 JSON,或超出 maxSnapshotBytes → 记录 warn、返回空快照、 无错误(hub 启动不因持久化文件损坏而失败,也不 panic);
  • 其余 I/O 错误(如权限不足)→ 返回 error,由调用方决定是否中止。

func (*Persister) Path

func (p *Persister) Path() string

Path 返回持久化文件路径(供服务层日志/诊断展示)。

func (*Persister) Save

func (p *Persister) Save(snap *Snapshot) error

Save 原子写快照到 p.path:先写同目录临时文件再 rename(不出现半写文件)。 父目录不存在时返回 error。写失败记录 error 日志(不静默吞掉,I2)。

func (*Persister) Schedule

func (p *Persister) Schedule(fn snapshotFn)

Schedule 排队一次快照并在去抖窗口后异步落盘。fn 返回 nil 表示无可持久化变化。 多次调用只会让**最后一次**的 fn 在窗口到期后执行(合并),避免注册/信令风暴 反复落盘。线程安全。

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 NewRouteTable

func NewRouteTable() *RouteTable

NewRouteTable 创建路由表。

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) Has

func (rt *RouteTable) Has(id NodeID) bool

Has 检查节点是否存在。

func (*RouteTable) HasService

func (rt *RouteTable) HasService(id NodeID, name string) bool

HasService 返回节点是否宣告了指定名称的服务。

func (*RouteTable) List

func (rt *RouteTable) List() []NodeInfo

List 返回所有已注册节点的列表。

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) NodeCount

func (rt *RouteTable) NodeCount() int

NodeCount 返回当前注册的节点数量。

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 NewSignalQueue

func NewSignalQueue() *SignalQueue

NewSignalQueue 创建信令队列。

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 后配额已释放。

func (*SignalQueue) Wait

func (q *SignalQueue) Wait(ctx context.Context, target string) error

Wait 阻塞等待 target 收件箱出现新消息(长轮询语义)。 返回后调用方应再 Pop 一次;有超时保护。 同一 target 的并发 Wait 共享同一个开放通道,Push/Purge close 后**全部**被唤醒 (P1-8);成功唤醒与 ctx 取消两条路径都会清理 waiter 条目,防止 waiters 表 无界增长(I6/S10)。

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(在线连接)不捕获——重启后节点重连建立。

Directories

Path Synopsis
ext
kad module

Jump to

Keyboard shortcuts

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