- go.mod require/replace git.zkcoi.com/zkcoi/meshray/core -> git.zkcoi.com/zkcoi/meshray
- .go imports meshray/core/{engine,connect} -> meshray/{engine,connect} (5 files)
- sync CHANGELOG + adapter.go comments
310 lines
10 KiB
Go
310 lines
10 KiB
Go
// Package corebind 是 Meshray-Manager 侧的 wireguard-go 集成层(形态 B:增强传输)。
|
||
//
|
||
// 它实现 wireguard-go 的 conn.Bind 接口(EnhancedBind),完全接管 WG 的 UDP 收发,
|
||
// 把密文流量导向 meshray 核心(git.zkcoi.com/zkcoi/meshray)的 9 层传输。
|
||
// 该层属于「管理器 / UI」范畴,而非核心;核心本身零 wireguard 依赖。
|
||
package corebind
|
||
|
||
import (
|
||
"net"
|
||
"sync"
|
||
|
||
"git.zkcoi.com/zkcoi/meshray/engine"
|
||
"golang.zx2c4.com/wireguard/conn"
|
||
"go.uber.org/zap"
|
||
)
|
||
|
||
// EnhancedBind 是 corebind 提供的 conn.Bind 实现,完全接管 WireGuard 的 UDP 收发,
|
||
// 使所有 WG 密文流量都经由 meshray 核心的 9 层传输(Direct/TURN/WebRTC/...)而非裸 UDP。
|
||
//
|
||
// 设计要点:
|
||
// - Direct-UDP 由一个绑定在 listen_port 的*共享未连接 UDP socket*(directSock)承载,
|
||
// 从任意来源收包,把「真实源地址」作为 Endpoint 交给 WG,WG 用密文内的 receiver index
|
||
// 匹配 peer 并据此更新对端出向地址(学址)。由此对称直连与 NAT 穿透学址原生可通。
|
||
// - TURN/WebRTC/WS 等中继层经核心传输连接(connMgr)承载,Direct 不可用时补充。
|
||
// - 集成通过核心暴露的钩子接入:SetPeerRegistrar / SetRelayInboundHandler /
|
||
// SetDirectByBind,核心永不反向依赖本包。
|
||
type EnhancedBind struct {
|
||
inner conn.Bind
|
||
networkID string
|
||
engine *core.Engine
|
||
logger *zap.Logger
|
||
|
||
// listenPort 本端 WG 监听端口;Direct-UDP 共享监听 socket 绑定于此,形成对称直连。
|
||
listenPort int
|
||
|
||
mu sync.RWMutex
|
||
peerByAddr map[string]string // endpoint addr(str) -> peerKey(Send 反查中继连接)
|
||
peerEndpoint map[string]conn.Endpoint // peerKey -> 对端真实 Endpoint(中继入站回灌用)
|
||
|
||
// directSock 承载 Direct-UDP 的共享未连接 UDP socket(绑定 listenPort)。
|
||
// 从任意来源收包(学习对端经 NAT 后的真实源地址),并按目标地址发包。
|
||
directSock *net.UDPConn
|
||
|
||
inboundCh chan inboundPacket
|
||
closed bool
|
||
}
|
||
|
||
// inboundPacket 入站密文;携带来源信息以便回灌 WG 时定位对端 Endpoint。
|
||
type inboundPacket struct {
|
||
peerKey string // 中继入站时按 peerKey 查 Endpoint
|
||
ep conn.Endpoint // Direct 共享 socket 学得的真实源地址;非空时优先使用
|
||
data []byte
|
||
}
|
||
|
||
// NewEnhancedBind 创建增强 Bind,并关联指定组网的 Engine 单例;
|
||
// 同时经核心钩子把入站回灌接到本 Bind,使 9 层传输连接读到的包能回到 WG。
|
||
func NewEnhancedBind(listenPort int, networkID string, logger *zap.Logger) (*EnhancedBind, error) {
|
||
eng := getOrCreateEngine(networkID, logger)
|
||
inner := conn.NewDefaultBind()
|
||
b := &EnhancedBind{
|
||
inner: inner,
|
||
networkID: networkID,
|
||
engine: eng,
|
||
logger: logger,
|
||
listenPort: listenPort,
|
||
peerByAddr: make(map[string]string),
|
||
peerEndpoint: make(map[string]conn.Endpoint),
|
||
inboundCh: make(chan inboundPacket, 1024),
|
||
}
|
||
// 通过核心钩子接入,避免核心反向依赖本包:
|
||
eng.SetPeerRegistrar(b) // 收到对端候选时注册 endpoint→peerKey
|
||
eng.SetRelayInboundHandler(b.deliverInbound) // 传输连接入站回灌 WG
|
||
eng.SetDirectByBind(true) // Direct 由本 Bind 共享 socket 承载
|
||
eng.SetListenPort(listenPort)
|
||
return b, nil
|
||
}
|
||
|
||
// Engine 返回关联的引擎实例(standalone 二进制用于注入候选触发拨号)
|
||
func (b *EnhancedBind) Engine() *core.Engine { return b.engine }
|
||
|
||
// RegisterPeer 注册 endpoint 地址到 peerKey 的映射,并缓存对端真实 Endpoint 供回灌使用。
|
||
// 由 Engine.NotifyPeerInfo 在收到候选时自动调用;standalone 也可直接调用。
|
||
func (b *EnhancedBind) RegisterPeer(addr string, peerKey string) {
|
||
ep, err := b.inner.ParseEndpoint(addr)
|
||
if err != nil {
|
||
b.logger.Warn("解析对端 Endpoint 失败,回灌将使用 0.0.0.0:0",
|
||
zap.String("addr", addr), zap.Error(err))
|
||
ep = nil
|
||
}
|
||
b.mu.Lock()
|
||
b.peerByAddr[addr] = peerKey
|
||
if ep != nil {
|
||
b.peerEndpoint[peerKey] = ep
|
||
}
|
||
b.mu.Unlock()
|
||
}
|
||
|
||
// Open 创建 Direct-UDP 共享监听 socket(绑定 listen_port),并启动接收协程。
|
||
// 返回从 inboundCh 读取入站密文的接收函数:Direct 入站与中继入站统一经此回灌 WG。
|
||
func (b *EnhancedBind) Open(port uint16) ([]conn.ReceiveFunc, uint16, error) {
|
||
b.mu.Lock()
|
||
b.closed = false
|
||
// 端口优先用 WG 传入的 port,其次用构造时的 listenPort;0 表示由 OS 分配。
|
||
bindPort := int(port)
|
||
if bindPort == 0 {
|
||
bindPort = b.listenPort
|
||
}
|
||
// 关闭可能残留的旧 socket(重开场景)
|
||
if b.directSock != nil {
|
||
b.directSock.Close()
|
||
b.directSock = nil
|
||
}
|
||
b.mu.Unlock()
|
||
|
||
sock, err := net.ListenUDP("udp", &net.UDPAddr{Port: bindPort})
|
||
if err != nil {
|
||
return nil, 0, err
|
||
}
|
||
|
||
actual := sock.LocalAddr().(*net.UDPAddr).Port
|
||
|
||
b.mu.Lock()
|
||
b.directSock = sock
|
||
b.listenPort = actual
|
||
b.mu.Unlock()
|
||
|
||
// 回填引擎实际监听端口(listen_port=0 自动分配时同步)
|
||
if b.engine != nil {
|
||
b.engine.SetListenPort(actual)
|
||
}
|
||
|
||
go b.readDirectLoop(sock)
|
||
|
||
b.logger.Info("Direct-UDP 共享监听 socket 已就绪",
|
||
zap.Int("listen_port", actual))
|
||
|
||
// v4/v6 各一个接收函数,均从同一个 inboundCh 读取
|
||
return []conn.ReceiveFunc{b.receiveFunc, b.receiveFunc}, uint16(actual), nil
|
||
}
|
||
|
||
// readDirectLoop 从共享监听 socket 收包(任意来源),把「真实源地址」作为 Endpoint 投入 inboundCh。
|
||
// 这是 NAT 穿透学址的核心:公网侧据此学到对端 NAT 映射后的真实地址,回包得以正确返回。
|
||
func (b *EnhancedBind) readDirectLoop(sock *net.UDPConn) {
|
||
buf := make([]byte, 65535)
|
||
for {
|
||
n, src, err := sock.ReadFromUDP(buf)
|
||
if err != nil {
|
||
b.mu.RLock()
|
||
closed := b.closed
|
||
b.mu.RUnlock()
|
||
if closed {
|
||
return
|
||
}
|
||
b.logger.Debug("Direct 共享 socket 读取失败,停止接收", zap.Error(err))
|
||
return
|
||
}
|
||
data := make([]byte, n)
|
||
copy(data, buf[:n])
|
||
|
||
ep, perr := b.inner.ParseEndpoint(src.String())
|
||
if perr != nil {
|
||
b.logger.Debug("解析入站源地址失败,丢弃", zap.String("src", src.String()), zap.Error(perr))
|
||
continue
|
||
}
|
||
select {
|
||
case b.inboundCh <- inboundPacket{ep: ep, data: data}:
|
||
default:
|
||
b.logger.Warn("入站队列满,丢弃包", zap.String("src", src.String()))
|
||
}
|
||
}
|
||
}
|
||
|
||
// receiveFunc 从 inboundCh 取出入站密文,拷贝进 WG 提供的缓冲(保持 WG 缓冲所有权),
|
||
// 并填入来源 Endpoint 供 WG 握手更新出向地址:
|
||
// - Direct 入站:ep 为共享 socket 学得的真实源地址(优先);
|
||
// - 中继入站:ep 为空,按 peerKey 查缓存的对端 Endpoint。
|
||
func (b *EnhancedBind) receiveFunc(packets [][]byte, sizes []int, eps []conn.Endpoint) (int, error) {
|
||
first, ok := <-b.inboundCh
|
||
if !ok {
|
||
return 0, net.ErrClosed
|
||
}
|
||
n := 0
|
||
copy(packets[n], first.data)
|
||
sizes[n] = len(first.data)
|
||
eps[n] = b.endpointForPacket(first)
|
||
n++
|
||
for n < len(packets) {
|
||
select {
|
||
case pkt := <-b.inboundCh:
|
||
copy(packets[n], pkt.data)
|
||
sizes[n] = len(pkt.data)
|
||
eps[n] = b.endpointForPacket(pkt)
|
||
n++
|
||
default:
|
||
return n, nil
|
||
}
|
||
}
|
||
return n, nil
|
||
}
|
||
|
||
// endpointForPacket 返回入站包的来源 Endpoint:Direct 入站直接用其学得的 ep;
|
||
// 中继入站按 peerKey 查缓存,未知时退化为 0.0.0.0:0。
|
||
func (b *EnhancedBind) endpointForPacket(pkt inboundPacket) conn.Endpoint {
|
||
if pkt.ep != nil {
|
||
return pkt.ep
|
||
}
|
||
b.mu.RLock()
|
||
defer b.mu.RUnlock()
|
||
if ep, ok := b.peerEndpoint[pkt.peerKey]; ok {
|
||
return ep
|
||
}
|
||
ep, _ := b.inner.ParseEndpoint("0.0.0.0:0")
|
||
return ep
|
||
}
|
||
|
||
// Send 出站路径:
|
||
// - 若该对端已建立中继/信令层连接(TURN/WebRTC/WS,位于 connMgr),经其发送;
|
||
// - 否则(Direct-UDP)经共享监听 socket 直接发往目标地址(WG 维护的对端 Endpoint:
|
||
// 初始为配置 endpoint,学址后为对端真实地址)。
|
||
func (b *EnhancedBind) Send(bufs [][]byte, ep conn.Endpoint) error {
|
||
addr := ep.DstToString()
|
||
|
||
b.mu.RLock()
|
||
peerKey, mapped := b.peerByAddr[addr]
|
||
b.mu.RUnlock()
|
||
|
||
if mapped {
|
||
if c, ok := b.engine.GetConnMgr().Get(peerKey); ok && c != nil {
|
||
for _, buf := range bufs {
|
||
if len(buf) == 0 {
|
||
continue
|
||
}
|
||
if _, err := c.Write(buf); err != nil {
|
||
b.logger.Debug("经中继连接发送失败,转 Direct 共享 socket",
|
||
zap.String("peer", peerKey), zap.Error(err))
|
||
return b.directSendTo(addr, bufs)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
}
|
||
// Direct-UDP:共享监听 socket 直发(NAT 穿透原生支持)
|
||
return b.directSendTo(addr, bufs)
|
||
}
|
||
|
||
// directSendTo 经共享监听 socket 把数据发往目标地址。源端口固定为 listen_port,
|
||
// 满足对称直连;NAT 后主动侧发包后,NAT 记录会话,公网侧回包得以返回。
|
||
func (b *EnhancedBind) directSendTo(addr string, bufs [][]byte) error {
|
||
b.mu.RLock()
|
||
sock := b.directSock
|
||
b.mu.RUnlock()
|
||
if sock == nil {
|
||
return net.ErrClosed
|
||
}
|
||
|
||
udpAddr, err := net.ResolveUDPAddr("udp", addr)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
for _, buf := range bufs {
|
||
if len(buf) == 0 {
|
||
continue
|
||
}
|
||
if _, err := sock.WriteToUDP(buf, udpAddr); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// deliverInbound 由引擎中继传输连接读取到入站时调用,把密文投入 inboundCh 回灌 WG。
|
||
// 中继入站不携带来源地址,按 peerKey 在回灌时查对端 Endpoint。
|
||
func (b *EnhancedBind) deliverInbound(peerKey string, packet []byte) {
|
||
cp := make([]byte, len(packet))
|
||
copy(cp, packet)
|
||
select {
|
||
case b.inboundCh <- inboundPacket{peerKey: peerKey, data: cp}:
|
||
default:
|
||
b.logger.Warn("入站队列满,丢弃包", zap.String("peer", peerKey))
|
||
}
|
||
}
|
||
|
||
// Close 关闭共享监听 socket 与入站通道(使接收函数返回 net.ErrClosed)。
|
||
func (b *EnhancedBind) Close() error {
|
||
b.mu.Lock()
|
||
if !b.closed {
|
||
b.closed = true
|
||
close(b.inboundCh)
|
||
}
|
||
if b.directSock != nil {
|
||
b.directSock.Close()
|
||
b.directSock = nil
|
||
}
|
||
b.mu.Unlock()
|
||
return b.inner.Close()
|
||
}
|
||
|
||
// SetMark 本实现不涉及内核 socket,置空标记为 no-op
|
||
func (b *EnhancedBind) SetMark(mark uint32) error { return nil }
|
||
|
||
// ParseEndpoint 委托默认 Bind 解析
|
||
func (b *EnhancedBind) ParseEndpoint(s string) (conn.Endpoint, error) {
|
||
return b.inner.ParseEndpoint(s)
|
||
}
|
||
|
||
// BatchSize 返回默认 Bind 的批量大小
|
||
func (b *EnhancedBind) BatchSize() int {
|
||
return b.inner.BatchSize()
|
||
}
|