Files
zkcoi e9ca2f7d70 fix: rename manager module to git.zkcoi.com/zkcoi/Meshray-Manager
- go.mod module path matches repo zkcoi/Meshray-Manager (case-distinct from zkcoi/Meshray/core)
- rewrite internal imports meshray/{internal,web,pkg} -> Meshray-Manager/... (core refs kept)
- sync README.md / install.sh repo URLs; add CHANGELOG entry
2026-07-15 16:25:46 +08:00

226 lines
5.9 KiB
Go

package scheduler
import (
"context"
"fmt"
"sync"
"time"
"git.zkcoi.com/zkcoi/Meshray-Manager/internal/model"
"git.zkcoi.com/zkcoi/Meshray-Manager/internal/service"
"go.uber.org/zap"
"gorm.io/gorm"
)
// DDNSUpdaterService DDNS 自动更新调度服务
type DDNSUpdaterService struct {
db *gorm.DB
logger *zap.Logger
ipDetection *service.IPDetectionService
ddnsOperation *service.DDNSOperationService
ctx context.Context
cancel context.CancelFunc
wg sync.WaitGroup
checkInterval time.Duration // 检测间隔
updateThreshold int // IP 变化阈值(连续多少次不同才更新)
}
// NewDDNSUpdaterService 创建 DDNS 自动更新服务
func NewDDNSUpdaterService(
db *gorm.DB,
logger *zap.Logger,
checkInterval time.Duration,
) *DDNSUpdaterService {
ctx, cancel := context.WithCancel(context.Background())
return &DDNSUpdaterService{
db: db,
logger: logger,
ipDetection: service.NewIPDetectionService(),
ddnsOperation: service.NewDDNSOperationService(logger, db),
ctx: ctx,
cancel: cancel,
checkInterval: checkInterval,
updateThreshold: 2, // 默认连续 2 次检测到不同 IP 才更新
}
}
// Start 启动后台自动更新任务
func (s *DDNSUpdaterService) Start() error {
s.logger.Info("启动 DDNS 自动更新服务",
zap.Duration("检测间隔", s.checkInterval),
zap.Int("变化阈值", s.updateThreshold))
s.wg.Add(1)
go s.runUpdater()
return nil
}
// Stop 停止后台自动更新任务
func (s *DDNSUpdaterService) Stop() {
s.logger.Info("停止 DDNS 自动更新服务")
s.cancel()
s.wg.Wait()
}
// runUpdater 运行自动更新协程
func (s *DDNSUpdaterService) runUpdater() {
defer s.wg.Done()
ticker := time.NewTicker(s.checkInterval)
defer ticker.Stop()
// 用于跟踪每个 DDNS 服务的连续不同 IP 次数
consecutiveChanges := make(map[string]int)
lastIPs := make(map[string]string)
for {
select {
case <-ticker.C:
s.checkAndUpdate(consecutiveChanges, lastIPs)
case <-s.ctx.Done():
s.logger.Info("DDNS 自动更新服务已停止")
return
}
}
}
// checkAndUpdate 检测并更新 DDNS 记录
func (s *DDNSUpdaterService) checkAndUpdate(
consecutiveChanges map[string]int,
lastIPs map[string]string,
) {
s.logger.Debug("开始检测 DDNS 服务 IP 变化")
// 查询所有启用的 DDNS 全功能模式服务
var ddnsServices []model.Service
if err := s.db.Where("type = ? AND config_mode = ? AND enabled = ?",
"DDNS", "fullservice", true).Find(&ddnsServices).Error; err != nil {
s.logger.Error("查询 DDNS 服务失败", zap.Error(err))
return
}
s.logger.Info("找到启用的 DDNS 全功能服务", zap.Int("数量", len(ddnsServices)))
for _, svc := range ddnsServices {
s.wg.Add(1)
go func(service model.Service) {
defer s.wg.Done()
s.processSingleService(&service, consecutiveChanges, lastIPs)
}(svc)
}
}
// processSingleService 处理单个 DDNS 服务
func (s *DDNSUpdaterService) processSingleService(
service *model.Service,
consecutiveChanges map[string]int,
lastIPs map[string]string,
) {
serviceID := service.ID
recordType := service.RecordType
// 只处理 A 和 AAAA 记录(需要 IP 检测)
if recordType != "A" && recordType != "AAAA" {
return
}
// 确定要比较的字段
var currentConfigIP string
switch recordType {
case "A":
currentConfigIP = service.TargetIP
case "AAAA":
currentConfigIP = service.TargetIP
}
if currentConfigIP == "" {
s.logger.Warn("DDNS 服务目标 IP 为空,跳过检测",
zap.String("service_id", serviceID))
return
}
// 检测当前公网 IP
detectedIP, err := s.ipDetection.DetectIP(recordType)
if err != nil {
s.logger.Error("检测 IP 失败",
zap.String("service_id", serviceID),
zap.String("record_type", recordType),
zap.Error(err))
return
}
s.logger.Debug("IP 检测结果",
zap.String("service_id", serviceID),
zap.String("record_type", recordType),
zap.String("配置 IP", currentConfigIP),
zap.String("检测 IP", detectedIP))
// 检查 IP 是否变化
mapKey := fmt.Sprintf("%s_%s", serviceID, recordType)
if detectedIP != currentConfigIP {
// IP 不同,增加计数
consecutiveChanges[mapKey]++
s.logger.Debug("IP 不一致",
zap.String("service_id", serviceID),
zap.Int("连续次数", consecutiveChanges[mapKey]),
zap.Int("阈值", s.updateThreshold))
// 达到阈值才更新
if consecutiveChanges[mapKey] >= s.updateThreshold {
s.logger.Info("IP 变化达到阈值,开始更新 DNS 记录",
zap.String("service_id", serviceID),
zap.String("旧 IP", currentConfigIP),
zap.String("新 IP", detectedIP))
// 更新 DNS 记录
err := s.updateDNSRecord(service, detectedIP)
if err != nil {
s.logger.Error("更新 DNS 记录失败",
zap.String("service_id", serviceID),
zap.Error(err))
} else {
s.logger.Info("DNS 记录更新成功",
zap.String("service_id", serviceID),
zap.String("新 IP", detectedIP))
// 重置计数
consecutiveChanges[mapKey] = 0
lastIPs[mapKey] = detectedIP
}
}
} else {
// IP 相同,重置计数
if consecutiveChanges[mapKey] > 0 {
s.logger.Debug("IP 恢复一致,重置计数器",
zap.String("service_id", serviceID))
consecutiveChanges[mapKey] = 0
}
}
// 更新最后检测的 IP
lastIPs[mapKey] = detectedIP
}
// updateDNSRecord 更新 DNS 记录
func (s *DDNSUpdaterService) updateDNSRecord(svc *model.Service, newIP string) error {
// 使用 DDNSOperationService 的 UpdateDNSRecord 方法
err := s.ddnsOperation.UpdateDNSRecord(svc, svc.RecordType, svc.Subdomain, newIP, 300)
if err != nil {
return fmt.Errorf("更新 DNS 记录失败:%w", err)
}
// 更新数据库中的 IP
if err := s.db.Model(&model.Service{}).Where("id = ?", svc.ID).Update("target_ip", newIP).Error; err != nil {
s.logger.Warn("更新数据库中的 IP 失败",
zap.String("service_id", svc.ID),
zap.Error(err))
}
return nil
}