Files
zkcoi 59e3059246 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
2026-07-15 16:08:13 +08:00

411 lines
12 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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"`
}