rebuild: split core into separate repo; move wg-go integration to internal/ctr/corebind
- go.mod: require git.zkcoi.com/zkcoi/meshray/core, replace => ../Meshray - internal/ctr/corebind: EnhancedBind + registry + adapter (wg-go integration isolated) - cmd/mr-wg migrated from old core/cmd/meshray-core - core/ removed; CHANGELOG updated; fix checkdb vet warning
This commit is contained in:
@@ -0,0 +1,67 @@
|
||||
// Package corebind 的契约类型转换:在 meshray-contract 与 meshray/core 内部类型间双向转换。
|
||||
package corebind
|
||||
|
||||
import (
|
||||
"git.zkcoi.com/zkcoi/meshray-contract"
|
||||
"git.zkcoi.com/zkcoi/meshray/core/connect"
|
||||
"git.zkcoi.com/zkcoi/meshray/core/engine"
|
||||
)
|
||||
|
||||
// 以下函数完成开源契约类型(contract 包)与 meshray/core 内部类型之间的双向转换。
|
||||
// 数值型 Layer 在 contract 与 connect 中定义一致,可直接转型。
|
||||
|
||||
func toCoreCandidates(in []contract.Candidate) []core.Candidate {
|
||||
out := make([]core.Candidate, 0, len(in))
|
||||
for _, c := range in {
|
||||
out = append(out, core.Candidate{
|
||||
Addr: c.Addr,
|
||||
Type: c.Type,
|
||||
Priority: c.Priority,
|
||||
Protocol: c.Protocol,
|
||||
})
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func toCoreICEConfig(in contract.ICEConfig) connect.ICEConfig {
|
||||
turn := make([]connect.TURNServerConfig, 0, len(in.TURNServers))
|
||||
for _, s := range in.TURNServers {
|
||||
turn = append(turn, connect.TURNServerConfig{URLs: []string{s}})
|
||||
}
|
||||
return connect.ICEConfig{
|
||||
STUNServers: in.STUNServers,
|
||||
TURNServers: turn,
|
||||
}
|
||||
}
|
||||
|
||||
func toCoreLayerOrder(in []contract.Layer) []connect.Layer {
|
||||
out := make([]connect.Layer, 0, len(in))
|
||||
for _, l := range in {
|
||||
out = append(out, connect.Layer(l))
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func toContractStatus(in *core.EngineStatus) *contract.EngineStatus {
|
||||
if in == nil {
|
||||
return nil
|
||||
}
|
||||
peers := make(map[string]*contract.PeerStatus, len(in.Peers))
|
||||
for k, v := range in.Peers {
|
||||
peers[k] = &contract.PeerStatus{
|
||||
PeerKey: v.PeerKey,
|
||||
Connected: v.Connected,
|
||||
Layer: v.Layer,
|
||||
}
|
||||
}
|
||||
return &contract.EngineStatus{
|
||||
PeerCount: in.PeerCount,
|
||||
Peers: peers,
|
||||
ActiveConnections: in.ActiveConnections,
|
||||
TotalConnections: in.TotalConnections,
|
||||
BytesSent: in.BytesSent,
|
||||
BytesReceived: in.BytesReceived,
|
||||
StrategyFallbacks: in.StrategyFallbacks,
|
||||
LastSwitchTime: in.LastSwitchTime,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,309 @@
|
||||
// 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/core/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()
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
// Package corebind 的契约注册:向开源契约(meshray-contract)注册增强 Bind 工厂与引擎提供方。
|
||||
package corebind
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"git.zkcoi.com/zkcoi/meshray-contract"
|
||||
"git.zkcoi.com/zkcoi/meshray/core/engine"
|
||||
"go.uber.org/zap"
|
||||
"golang.zx2c4.com/wireguard/conn"
|
||||
)
|
||||
|
||||
// engines 维护 networkID -> Engine 单例映射。
|
||||
// EnhancedBind(收发接管)与 engineController(控制面)均通过 networkID 关联同一个 Engine,
|
||||
// 使信令下发的候选地址、ICE 配置、层顺序都能落到接管该组网收发的 Engine 上。
|
||||
var (
|
||||
enginesMu sync.Mutex
|
||||
engines = make(map[string]*core.Engine)
|
||||
)
|
||||
|
||||
// getOrCreateEngine 获取或创建指定组网的 Engine 单例(线程安全)
|
||||
func getOrCreateEngine(networkID string, logger *zap.Logger) *core.Engine {
|
||||
enginesMu.Lock()
|
||||
defer enginesMu.Unlock()
|
||||
if e, ok := engines[networkID]; ok {
|
||||
return e
|
||||
}
|
||||
e := core.NewEngine(logger, core.NewMetrics())
|
||||
engines[networkID] = e
|
||||
return e
|
||||
}
|
||||
|
||||
// engineController 适配 contract.EngineController,以 networkID 关联 Engine 单例,
|
||||
// 供管理器在 enhanced 模式调用,无需直接依赖本包具体类型。
|
||||
type engineController struct {
|
||||
networkID string
|
||||
engine *core.Engine
|
||||
}
|
||||
|
||||
func (c *engineController) Start() error { return c.engine.Start() }
|
||||
|
||||
func (c *engineController) Stop() error { return c.engine.Stop() }
|
||||
|
||||
// NotifyPeerCandidates 由信令层调用:注入对端候选地址并触发 P2P 建连
|
||||
func (c *engineController) NotifyPeerCandidates(peerKey string, candidates []contract.Candidate) error {
|
||||
return c.engine.NotifyPeerInfo(peerKey, toCoreCandidates(candidates))
|
||||
}
|
||||
|
||||
// SetICEConfig 更新 ICE/STUN/TURN 配置
|
||||
func (c *engineController) SetICEConfig(cfg contract.ICEConfig) error {
|
||||
return c.engine.SetICEConfig(toCoreICEConfig(cfg))
|
||||
}
|
||||
|
||||
// SetLayerOrder 设置传输层优先级顺序(策略核心配置)
|
||||
func (c *engineController) SetLayerOrder(order []contract.Layer) error {
|
||||
c.engine.SetLayerOrder(toCoreLayerOrder(order))
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetStatus 查询引擎状态
|
||||
func (c *engineController) GetStatus() (*contract.EngineStatus, error) {
|
||||
st, err := c.engine.GetStatus()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return toContractStatus(st), nil
|
||||
}
|
||||
|
||||
// init 向开源契约注册表注册增强 Bind 工厂与引擎提供方。
|
||||
// 由管理器侧 core_enable.go 在 -tags meshray_core 构建下空白导入本包触发,
|
||||
// 此时管理器侧 contract.GetBindFactory(true) / NewEngineController 才能取到实现;
|
||||
// 默认构建不引入本包,原生直连(默认 Bind)始终可用。
|
||||
func init() {
|
||||
contract.RegisterEnhancedBindFactory(func(listenPort int, networkID string, logger *zap.Logger) (conn.Bind, error) {
|
||||
return NewEnhancedBind(listenPort, networkID, logger)
|
||||
})
|
||||
contract.RegisterEngineProvider(func(networkID string, logger *zap.Logger) (contract.EngineController, error) {
|
||||
return &engineController{networkID: networkID, engine: getOrCreateEngine(networkID, logger)}, nil
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user