- 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
411 lines
12 KiB
Go
411 lines
12 KiB
Go
package ctr
|
||
|
||
import (
|
||
"fmt"
|
||
"strconv"
|
||
"sync"
|
||
|
||
contract "git.zkcoi.com/zkcoi/meshray-contract"
|
||
"go.uber.org/zap"
|
||
)
|
||
|
||
// Ctr meshray-ctr 调度中心
|
||
type Ctr struct {
|
||
name string // 组网名称(用于日志)
|
||
networkID uint64 // 组网 ID(雪花算法 ID)
|
||
config *CtrConfig
|
||
logger *zap.Logger
|
||
|
||
// WireGuard 管理器
|
||
wgManager *WGManager
|
||
|
||
// 增强引擎实例(由 meshray-core 提供,按 networkID 索引);原生模式为空
|
||
engines map[string]contract.EngineController
|
||
|
||
mu sync.RWMutex
|
||
}
|
||
|
||
// CtrConfig 配置
|
||
type CtrConfig struct {
|
||
// 空配置,保留结构体以备未来扩展
|
||
}
|
||
|
||
// NewCtr 创建调度中心
|
||
func NewCtr(name string, networkID uint64, config *CtrConfig, logger *zap.Logger) (*Ctr, error) {
|
||
ctr := &Ctr{
|
||
name: name,
|
||
networkID: networkID,
|
||
config: config,
|
||
logger: logger,
|
||
engines: make(map[string]contract.EngineController),
|
||
}
|
||
|
||
// 初始化 WireGuard 管理器(两种模式都需要)
|
||
ctr.wgManager = NewWGManager(logger)
|
||
|
||
return ctr, nil
|
||
}
|
||
|
||
// Start 启动调度中心
|
||
func (c *Ctr) Start() error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
|
||
c.logger.Info("启动 meshray-ctr",
|
||
zap.String("name", c.name),
|
||
zap.Uint64("network_id", c.networkID))
|
||
|
||
return nil
|
||
}
|
||
|
||
// Stop 停止调度中心
|
||
func (c *Ctr) Stop() error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
|
||
c.logger.Info("停止 meshray-ctr",
|
||
zap.Uint64("network_id", c.networkID))
|
||
|
||
// 停止所有增强引擎
|
||
for networkID, engineCtl := range c.engines {
|
||
if err := engineCtl.Stop(); err != nil {
|
||
c.logger.Warn("停止增强引擎失败",
|
||
zap.String("network_id", networkID), zap.Error(err))
|
||
}
|
||
delete(c.engines, networkID)
|
||
}
|
||
|
||
// 停止 WireGuard 管理
|
||
if err := c.wgManager.Stop(); err != nil {
|
||
c.logger.Error("停止 WireGuard 管理失败", zap.Error(err))
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// CreateNetwork 创建网络(根据组网模式决定行为)
|
||
// - 原生模式:仅创建 WG 设备(wireguard-go 默认 Bind 直连)
|
||
// - 增强模式:WG 设备使用 meshray-core 的自定义 conn.Bind 接管收发;并启动增强引擎控制面
|
||
func (c *Ctr) CreateNetwork(networkID uint64, subnet string, listenPort int, meshMode string) error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
|
||
c.logger.Info("开始创建网络",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("subnet", subnet),
|
||
zap.Int("listen_port", listenPort),
|
||
zap.String("mesh_mode", meshMode))
|
||
|
||
// 1. 创建 WireGuard 设备(两种模式都需要;增强模式使用 core 注册的自定义 Bind)
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
if err := c.wgManager.CreateDevice(networkIDStr, subnet, listenPort, meshMode); err != nil {
|
||
return fmt.Errorf("创建 WireGuard 设备失败:%w", err)
|
||
}
|
||
c.logger.Debug("WireGuard 设备创建成功",
|
||
zap.Uint64("network_id", networkID))
|
||
|
||
// 2. 仅增强模式需要启动增强引擎(控制面)
|
||
if meshMode == "enhanced" {
|
||
engineCtl, err := contract.NewEngineController(networkIDStr, c.logger)
|
||
if err != nil {
|
||
c.wgManager.DeleteDevice(networkIDStr) // 回滚 WG 设备
|
||
return fmt.Errorf("增强模式需 meshray-core:%w", err)
|
||
}
|
||
if err := engineCtl.Start(); err != nil {
|
||
c.wgManager.DeleteDevice(networkIDStr) // 回滚 WG 设备
|
||
return fmt.Errorf("启动增强引擎失败:%w", err)
|
||
}
|
||
c.engines[networkIDStr] = engineCtl
|
||
c.logger.Info("增强模式:引擎已启动(收发由 meshray-core conn.Bind 接管)",
|
||
zap.Uint64("network_id", networkID))
|
||
} else {
|
||
c.logger.Info("原生模式:仅创建 WG 设备,标准直连",
|
||
zap.Uint64("network_id", networkID))
|
||
}
|
||
|
||
c.logger.Info("网络创建成功",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("mesh_mode", meshMode))
|
||
|
||
return nil
|
||
}
|
||
|
||
// DeleteNetwork 删除网络
|
||
func (c *Ctr) DeleteNetwork(networkID uint64) error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
|
||
c.logger.Info("开始删除网络",
|
||
zap.Uint64("network_id", networkID))
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
|
||
// 停止并清理增强引擎(若存在)
|
||
if engineCtl, ok := c.engines[networkIDStr]; ok {
|
||
if err := engineCtl.Stop(); err != nil {
|
||
c.logger.Warn("停止增强引擎失败", zap.String("network_id", networkIDStr), zap.Error(err))
|
||
}
|
||
delete(c.engines, networkIDStr)
|
||
}
|
||
|
||
// 删除 WireGuard 设备
|
||
if err := c.wgManager.DeleteDevice(networkIDStr); err != nil {
|
||
return fmt.Errorf("删除 WireGuard 设备失败:%w", err)
|
||
}
|
||
|
||
c.logger.Info("网络删除成功",
|
||
zap.Uint64("network_id", networkID))
|
||
|
||
return nil
|
||
}
|
||
|
||
// AddPeer 添加 Peer
|
||
func (c *Ctr) AddPeer(networkID uint64, publicKey, allowedIP string) error {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
c.logger.Info("添加 Peer",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("public_key", publicKey[:8]+"..."))
|
||
|
||
// 1. 添加到 WireGuard
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
if err := c.wgManager.AddPeer(networkIDStr, publicKey, allowedIP); err != nil {
|
||
return fmt.Errorf("添加 Peer 失败:%w", err)
|
||
}
|
||
|
||
// 2. 增强模式:meshray-core 的 conn.Bind 已接管该组网收发,peer Endpoint 保持真实地址,
|
||
// 无需改写回环端口;候选地址后续由信令层经 NotifyPeerCandidates 注入触发建连。
|
||
if _, ok := c.engines[networkIDStr]; ok {
|
||
c.logger.Info("增强模式:动态 Peer 已交由 meshray-core 接管",
|
||
zap.String("public_key", publicKey[:8]+"..."))
|
||
} else {
|
||
c.logger.Debug("原生模式无需代理 Peer", zap.Uint64("network_id", networkID))
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// RemovePeer 移除 Peer
|
||
func (c *Ctr) RemovePeer(networkID uint64, publicKey string) error {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
c.logger.Info("移除 Peer",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("public_key", publicKey[:8]+"..."))
|
||
|
||
// 从 WireGuard 移除
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
if err := c.wgManager.RemovePeer(networkIDStr, publicKey); err != nil {
|
||
return fmt.Errorf("移除 Peer 失败:%w", err)
|
||
}
|
||
|
||
// 增强模式:meshray-core 的 conn.Bind 按 WG peer 状态自动管理收发,
|
||
// 移除 WG peer 后不再有该 peer 的包,无需显式解绑
|
||
if _, ok := c.engines[networkIDStr]; ok {
|
||
c.logger.Info("增强模式:Core 已随 WG Peer 移除停止对该 Peer 的转发",
|
||
zap.String("public_key", publicKey[:8]+"..."))
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// SetSTUNTURNConfig 为指定网络设置 STUN/TURN 配置(仅增强模式生效)
|
||
func (c *Ctr) SetSTUNTURNConfig(networkID uint64, stunServers []string, turnServers []string) error {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
|
||
// 获取增强引擎
|
||
engineCtl, ok := c.engines[networkIDStr]
|
||
if !ok {
|
||
c.logger.Debug("网络未启动增强模式,跳过 STUN/TURN 配置",
|
||
zap.Uint64("network_id", networkID))
|
||
return nil // 无需错误,因为原生模式不需要
|
||
}
|
||
|
||
// 更新 WebRTC 工厂的 ICE 配置
|
||
if err := engineCtl.SetICEConfig(contract.ICEConfig{
|
||
STUNServers: stunServers,
|
||
TURNServers: turnServers,
|
||
}); err != nil {
|
||
return fmt.Errorf("设置 ICE 配置失败:%w", err)
|
||
}
|
||
|
||
c.logger.Info("STUN/TURN 配置已设置",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.Int("stun_count", len(stunServers)),
|
||
zap.Int("turn_count", len(turnServers)))
|
||
|
||
return nil
|
||
}
|
||
|
||
func (c *Ctr) GetStatus(networkID uint64) (*NetworkStatus, error) {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
status := &NetworkStatus{
|
||
NetworkID: networkIDStr,
|
||
}
|
||
|
||
// 获取 WireGuard 状态
|
||
wgStatus, err := c.wgManager.GetStatus(networkIDStr)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
status.WGStatus = wgStatus
|
||
|
||
// 获取增强引擎状态
|
||
if engineCtl, ok := c.engines[networkIDStr]; ok {
|
||
if coreStatus, err := engineCtl.GetStatus(); err == nil {
|
||
status.CoreStatus = coreStatus
|
||
}
|
||
}
|
||
|
||
return status, nil
|
||
}
|
||
|
||
// SwitchMode 切换传输模式(原生→增强)
|
||
// 仅支持 "native" → "enhanced" 单向切换。增强模式由 meshray-core 的 conn.Bind 接管收发,
|
||
// 不再改写 peer Endpoint 为回环端口(根治入站接不通的缺陷)。
|
||
func (c *Ctr) SwitchMode(networkID uint64, mode string) error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
|
||
c.logger.Info("开始切换传输模式",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("target_mode", mode))
|
||
|
||
// 1. 验证目标模式
|
||
if mode != "enhanced" {
|
||
return fmt.Errorf("仅支持切换到 enhanced 模式,当前请求:%s", mode)
|
||
}
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
|
||
// 2. 获取当前 WG 状态(获取 Peer 数量)
|
||
wgStatus, err := c.wgManager.GetStatus(networkIDStr)
|
||
if err != nil {
|
||
return fmt.Errorf("获取 WG 状态失败:%w", err)
|
||
}
|
||
|
||
c.logger.Info("获取到 WG 状态",
|
||
zap.Int("peer_count", wgStatus.PeerCount))
|
||
|
||
// 3. 创建并启动增强引擎(控制面)
|
||
engineCtl, err := contract.NewEngineController(networkIDStr, c.logger)
|
||
if err != nil {
|
||
return fmt.Errorf("创建增强引擎失败(需 meshray-core):%w", err)
|
||
}
|
||
if err := engineCtl.Start(); err != nil {
|
||
return fmt.Errorf("启动增强引擎失败:%w", err)
|
||
}
|
||
c.engines[networkIDStr] = engineCtl
|
||
|
||
c.logger.Info("增强引擎已启动(收发由 meshray-core conn.Bind 接管)",
|
||
zap.Uint64("network_id", networkID))
|
||
|
||
// 4. 获取 Peer 列表
|
||
peers, err := c.wgManager.ListPeers(networkIDStr)
|
||
if err != nil {
|
||
return fmt.Errorf("获取 Peer 列表失败:%w", err)
|
||
}
|
||
|
||
// 5. 把已知真实 Endpoint 作为候选注入,触发 P2P 建连
|
||
// (不再改写 peer Endpoint 为回环端口;peer Endpoint 保持真实地址,WG 入站匹配成功)
|
||
for _, peer := range peers {
|
||
if peer.Endpoint != "" {
|
||
candidates := []contract.Candidate{
|
||
{Addr: peer.Endpoint, Type: "host", Priority: 1, Protocol: "udp"},
|
||
}
|
||
if err := engineCtl.NotifyPeerCandidates(peer.PublicKey, candidates); err != nil {
|
||
c.logger.Warn("注入对端候选地址失败",
|
||
zap.String("public_key", peer.PublicKey[:8]+"..."),
|
||
zap.Error(err))
|
||
}
|
||
}
|
||
}
|
||
|
||
c.logger.Info("传输模式切换成功(增强模式:meshray-core 接管收发)",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.String("mode", mode))
|
||
|
||
return nil
|
||
}
|
||
|
||
// CoreConfigUpdate 核心配置更新请求(由 server 层在策略变更时构造并传入)
|
||
type CoreConfigUpdate struct {
|
||
STUNServers []string // STUN 服务器列表
|
||
TURNServers []string // TURN 服务器列表(URL)
|
||
EnabledLayers []contract.Layer // 启用的传输层顺序(为空保持默认 9 层)
|
||
}
|
||
|
||
// NotifyPeerCandidates 由信令层调用:注入对端候选地址并触发 P2P 建连
|
||
func (c *Ctr) NotifyPeerCandidates(networkID uint64, publicKey string, candidates []contract.Candidate) error {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
engineCtl, ok := c.engines[networkIDStr]
|
||
if !ok {
|
||
return fmt.Errorf("网络未启用增强模式")
|
||
}
|
||
return engineCtl.NotifyPeerCandidates(publicKey, candidates)
|
||
}
|
||
|
||
// UpdateCoreConfig 更新 Core 配置(策略修改后同步)
|
||
// 真实生效:更新 ICE/STUN/TURN 配置到 WebRTC 工厂,并按需调整传输层优先级顺序
|
||
func (c *Ctr) UpdateCoreConfig(networkID uint64, config interface{}) error {
|
||
c.mu.RLock()
|
||
defer c.mu.RUnlock()
|
||
|
||
update, ok := config.(*CoreConfigUpdate)
|
||
if !ok {
|
||
return fmt.Errorf("不支持的配置类型:%T", config)
|
||
}
|
||
|
||
networkIDStr := strconv.FormatUint(networkID, 10)
|
||
engineCtl, ok := c.engines[networkIDStr]
|
||
if !ok {
|
||
return fmt.Errorf("网络未启用增强模式")
|
||
}
|
||
|
||
// 1. 更新 ICE/STUN/TURN 配置(下发到 WebRTC 工厂)
|
||
if err := engineCtl.SetICEConfig(contract.ICEConfig{
|
||
STUNServers: update.STUNServers,
|
||
TURNServers: update.TURNServers,
|
||
}); err != nil {
|
||
return fmt.Errorf("设置 ICE 配置失败:%w", err)
|
||
}
|
||
|
||
// 2. 更新传输层优先级顺序(策略核心配置)
|
||
if len(update.EnabledLayers) > 0 {
|
||
if err := engineCtl.SetLayerOrder(update.EnabledLayers); err != nil {
|
||
return fmt.Errorf("设置传输层顺序失败:%w", err)
|
||
}
|
||
}
|
||
|
||
c.logger.Info("Core 配置已更新",
|
||
zap.Uint64("network_id", networkID),
|
||
zap.Int("stun_count", len(update.STUNServers)),
|
||
zap.Int("turn_count", len(update.TURNServers)),
|
||
zap.Int("enabled_layers", len(update.EnabledLayers)))
|
||
|
||
return nil
|
||
}
|
||
|
||
// GetWGMode 获取系统 WG 模式
|
||
func (c *Ctr) GetWGMode() string {
|
||
return c.wgManager.GetWGMode()
|
||
}
|
||
|
||
// NetworkStatus 网络状态
|
||
type NetworkStatus struct {
|
||
NetworkID string `json:"network_id"`
|
||
WGStatus *WGStatus `json:"wg_status"`
|
||
CoreStatus *contract.EngineStatus `json:"core_status,omitempty"`
|
||
}
|