From a7f022ee614c002d138c5e09c745965dba525ab8 Mon Sep 17 00:00:00 2001 From: hequan2017 Date: Tue, 10 Feb 2026 18:06:58 +0800 Subject: [PATCH] feat(server): add pcdn scheduling and dispatch task pipeline --- server/initialize/gorm_biz.go | 3 +- server/initialize/timer.go | 14 +++ server/model/pcdn/dispatch_task.go | 40 ++++++ server/service/enter.go | 2 + server/service/pcdn/dispatcher.go | 50 ++++++++ server/service/pcdn/enter.go | 5 + server/service/pcdn/fallback.go | 56 +++++++++ server/service/pcdn/policy_engine.go | 48 +++++++ server/service/pcdn/runtime.go | 147 ++++++++++++++++++++++ server/service/pcdn/scheduler.go | 180 +++++++++++++++++++++++++++ server/task/pcdn_dispatch.go | 24 ++++ 11 files changed, 568 insertions(+), 1 deletion(-) create mode 100644 server/model/pcdn/dispatch_task.go create mode 100644 server/service/pcdn/dispatcher.go create mode 100644 server/service/pcdn/enter.go create mode 100644 server/service/pcdn/fallback.go create mode 100644 server/service/pcdn/policy_engine.go create mode 100644 server/service/pcdn/runtime.go create mode 100644 server/service/pcdn/scheduler.go create mode 100644 server/task/pcdn_dispatch.go diff --git a/server/initialize/gorm_biz.go b/server/initialize/gorm_biz.go index e5f2f97..9774fab 100644 --- a/server/initialize/gorm_biz.go +++ b/server/initialize/gorm_biz.go @@ -5,12 +5,13 @@ import ( "github.com/flipped-aurora/gin-vue-admin/server/model/computenode" "github.com/flipped-aurora/gin-vue-admin/server/model/imageregistry" "github.com/flipped-aurora/gin-vue-admin/server/model/instance" + "github.com/flipped-aurora/gin-vue-admin/server/model/pcdn" "github.com/flipped-aurora/gin-vue-admin/server/model/product" ) func bizModel() error { db := global.GVA_DB - err := db.AutoMigrate(imageregistry.ImageRegistry{}, computenode.ComputeNode{}, product.ProductSpec{}, instance.Instance{}) + err := db.AutoMigrate(imageregistry.ImageRegistry{}, computenode.ComputeNode{}, product.ProductSpec{}, instance.Instance{}, pcdn.PcdnDispatchTask{}) if err != nil { return err } diff --git a/server/initialize/timer.go b/server/initialize/timer.go index 1464c4a..0fd3247 100644 --- a/server/initialize/timer.go +++ b/server/initialize/timer.go @@ -32,6 +32,20 @@ func Timer() { // 其他定时任务定在这里 参考上方使用方法 + _, err = global.GVA_Timer.AddTaskByFuncWithSecond("PCDNDispatch", "*/10 * * * * *", func() { + task.ProcessPcdnDispatchTasks() + }, "PCDN调度任务执行", option...) + if err != nil { + fmt.Println("add PCDN dispatch timer error:", err) + } + + _, err = global.GVA_Timer.AddTaskByFuncWithSecond("PCDNDispatch", "*/12 * * * * *", func() { + task.SyncPcdnDispatchStatus() + }, "PCDN调度状态同步", option...) + if err != nil { + fmt.Println("add PCDN sync timer error:", err) + } + //_, err := global.GVA_Timer.AddTaskByFunc("定时任务标识", "corn表达式", func() { // 具体执行内容... // ...... diff --git a/server/model/pcdn/dispatch_task.go b/server/model/pcdn/dispatch_task.go new file mode 100644 index 0000000..7a890fa --- /dev/null +++ b/server/model/pcdn/dispatch_task.go @@ -0,0 +1,40 @@ +package pcdn + +import ( + "github.com/flipped-aurora/gin-vue-admin/server/global" + "github.com/flipped-aurora/gin-vue-admin/server/model/common" +) + +const ( + DispatchStatusPending = "pending" + DispatchStatusRunning = "running" + DispatchStatusSuccess = "success" + DispatchStatusFailed = "failed" + DispatchStatusRetrying = "retrying" +) + +// PcdnDispatchTask 记录调度计划及下发执行结果,支持审计与重放。 +type PcdnDispatchTask struct { + global.GVA_MODEL + TaskID string `json:"taskId" gorm:"column:task_id;size:64;not null;uniqueIndex:uk_task_trace,priority:1;index:idx_task_trace,priority:1"` + TraceID string `json:"traceId" gorm:"column:trace_id;size:64;not null;uniqueIndex:uk_task_trace,priority:2;index:idx_task_trace,priority:2"` + ContentID string `json:"contentId" gorm:"column:content_id;size:128;not null;index"` + UserRegion string `json:"userRegion" gorm:"column:user_region;size:64;not null;index"` + UserISP string `json:"userIsp" gorm:"column:user_isp;size:64;default:''"` + TopN int `json:"topN" gorm:"column:top_n;default:3"` + Status string `json:"status" gorm:"column:status;size:32;not null;index"` + RetryCount int `json:"retryCount" gorm:"column:retry_count;default:0"` + MaxRetry int `json:"maxRetry" gorm:"column:max_retry;default:3"` + TimeoutSeconds int `json:"timeoutSeconds" gorm:"column:timeout_seconds;default:8"` + NextRetryUnix int64 `json:"nextRetryUnix" gorm:"column:next_retry_unix;default:0;index"` + LastError string `json:"lastError" gorm:"column:last_error;size:1024;default:''"` + Candidates common.JSONMap `json:"candidates" gorm:"column:candidates;type:json"` + MetricsSnapshot common.JSONMap `json:"metricsSnapshot" gorm:"column:metrics_snapshot;type:json"` + DispatchProtocol string `json:"dispatchProtocol" gorm:"column:dispatch_protocol;size:64;default:'mock'"` + PrimaryNodeID uint `json:"primaryNodeId" gorm:"column:primary_node_id;default:0"` + CurrentNodeID uint `json:"currentNodeId" gorm:"column:current_node_id;default:0"` +} + +func (PcdnDispatchTask) TableName() string { + return "pcdn_dispatch_task" +} diff --git a/server/service/enter.go b/server/service/enter.go index 79a8fd6..ad2616b 100644 --- a/server/service/enter.go +++ b/server/service/enter.go @@ -5,6 +5,7 @@ import ( "github.com/flipped-aurora/gin-vue-admin/server/service/example" "github.com/flipped-aurora/gin-vue-admin/server/service/imageregistry" "github.com/flipped-aurora/gin-vue-admin/server/service/instance" + "github.com/flipped-aurora/gin-vue-admin/server/service/pcdn" "github.com/flipped-aurora/gin-vue-admin/server/service/product" "github.com/flipped-aurora/gin-vue-admin/server/service/system" ) @@ -18,4 +19,5 @@ type ServiceGroup struct { ComputenodeServiceGroup computenode.ServiceGroup ProductServiceGroup product.ServiceGroup InstanceServiceGroup instance.ServiceGroup + PcdnServiceGroup pcdn.ServiceGroup } diff --git a/server/service/pcdn/dispatcher.go b/server/service/pcdn/dispatcher.go new file mode 100644 index 0000000..3eb561d --- /dev/null +++ b/server/service/pcdn/dispatcher.go @@ -0,0 +1,50 @@ +package pcdn + +import ( + "context" + "fmt" + "time" +) + +// DispatchRequest 下发请求。 +type DispatchRequest struct { + TaskID string + TraceID string + ContentID string + TargetNode uint + TimeoutSec int +} + +// DispatchResult 下发结果。 +type DispatchResult struct { + Success bool + Message string +} + +// Dispatcher 抽象下发协议接口,便于后续扩展 HTTP/gRPC/MQ。 +type Dispatcher interface { + Dispatch(ctx context.Context, request DispatchRequest) DispatchResult + Protocol() string +} + +// MockDispatcher 为默认实现。 +type MockDispatcher struct{} + +func (d *MockDispatcher) Protocol() string { return "mock" } + +func (d *MockDispatcher) Dispatch(ctx context.Context, request DispatchRequest) DispatchResult { + if request.TimeoutSec <= 0 { + request.TimeoutSec = 8 + } + timeoutCtx, cancel := context.WithTimeout(ctx, time.Duration(request.TimeoutSec)*time.Second) + defer cancel() + select { + case <-timeoutCtx.Done(): + return DispatchResult{Success: false, Message: timeoutCtx.Err().Error()} + case <-time.After(50 * time.Millisecond): + if request.TargetNode == 0 { + return DispatchResult{Success: false, Message: "invalid target node"} + } + return DispatchResult{Success: true, Message: fmt.Sprintf("dispatched to node %d", request.TargetNode)} + } +} diff --git a/server/service/pcdn/enter.go b/server/service/pcdn/enter.go new file mode 100644 index 0000000..0306c0a --- /dev/null +++ b/server/service/pcdn/enter.go @@ -0,0 +1,5 @@ +package pcdn + +type ServiceGroup struct { + SchedulerService SchedulerService +} diff --git a/server/service/pcdn/fallback.go b/server/service/pcdn/fallback.go new file mode 100644 index 0000000..81ca560 --- /dev/null +++ b/server/service/pcdn/fallback.go @@ -0,0 +1,56 @@ +package pcdn + +import pcdnModel "github.com/flipped-aurora/gin-vue-admin/server/model/pcdn" + +// NextFallbackNode 主节点失败时选择次优候选。 +func NextFallbackNode(task pcdnModel.PcdnDispatchTask) uint { + raw, ok := task.Candidates["list"] + if ok { + if nodeID := findNextNode(raw, task.CurrentNodeID); nodeID != 0 { + return nodeID + } + } + return findNextNode(task.Candidates, task.CurrentNodeID) +} + +func findNextNode(raw any, current uint) uint { + switch list := raw.(type) { + case []any: + for _, item := range list { + nodeID := parseNodeID(item) + if nodeID != 0 && nodeID != current { + return nodeID + } + } + case map[string]any: + nodeID := parseNodeID(list) + if nodeID != 0 && nodeID != current { + return nodeID + } + } + return 0 +} + +func parseNodeID(v any) uint { + m, ok := v.(map[string]any) + if !ok { + return 0 + } + raw, ok := m["node_id"] + if !ok { + return 0 + } + switch id := raw.(type) { + case uint: + return id + case int: + if id > 0 { + return uint(id) + } + case float64: + if id > 0 { + return uint(id) + } + } + return 0 +} diff --git a/server/service/pcdn/policy_engine.go b/server/service/pcdn/policy_engine.go new file mode 100644 index 0000000..f21f692 --- /dev/null +++ b/server/service/pcdn/policy_engine.go @@ -0,0 +1,48 @@ +package pcdn + +import "math" + +// RealtimeNodeMetric 节点实时指标。 +type RealtimeNodeMetric struct { + NodeID uint + LatencyMS float64 + UnitCost float64 + LoadPercent float64 + HealthScore float64 + Online bool + PolicyDisabled bool + ISP string +} + +// ScoreWeights 策略权重。 +type ScoreWeights struct { + Latency float64 + Cost float64 + Load float64 + Health float64 +} + +// PolicyEngine 执行加权评分。 +type PolicyEngine struct { + weights ScoreWeights +} + +func NewPolicyEngine(weights ScoreWeights) *PolicyEngine { + if weights == (ScoreWeights{}) { + weights = ScoreWeights{Latency: 0.35, Cost: 0.2, Load: 0.2, Health: 0.25} + } + return &PolicyEngine{weights: weights} +} + +// Score 将指标归一化后进行加权,分数越高越优。 +func (p *PolicyEngine) Score(metric RealtimeNodeMetric) float64 { + latencyScore := 1 / (1 + math.Max(metric.LatencyMS, 0)/50) + costScore := 1 / (1 + math.Max(metric.UnitCost, 0)) + loadScore := 1 - math.Min(math.Max(metric.LoadPercent, 0), 100)/100 + healthScore := math.Min(math.Max(metric.HealthScore, 0), 100) / 100 + + return p.weights.Latency*latencyScore + + p.weights.Cost*costScore + + p.weights.Load*loadScore + + p.weights.Health*healthScore +} diff --git a/server/service/pcdn/runtime.go b/server/service/pcdn/runtime.go new file mode 100644 index 0000000..98d852d --- /dev/null +++ b/server/service/pcdn/runtime.go @@ -0,0 +1,147 @@ +package pcdn + +import ( + "context" + "fmt" + "time" + + "github.com/flipped-aurora/gin-vue-admin/server/global" + pcdnModel "github.com/flipped-aurora/gin-vue-admin/server/model/pcdn" + "go.uber.org/zap" +) + +type RuntimeService struct { + dispatcher Dispatcher +} + +func NewRuntimeService(dispatcher Dispatcher) *RuntimeService { + if dispatcher == nil { + dispatcher = &MockDispatcher{} + } + return &RuntimeService{dispatcher: dispatcher} +} + +// ProcessPendingTasks 扫描待处理任务,支持超时、重试和幂等状态推进。 +func (r *RuntimeService) ProcessPendingTasks(ctx context.Context, limit int) error { + if limit <= 0 { + limit = 20 + } + now := time.Now().Unix() + var tasks []pcdnModel.PcdnDispatchTask + err := global.GVA_DB.WithContext(ctx). + Where("status IN ?", []string{pcdnModel.DispatchStatusPending, pcdnModel.DispatchStatusRetrying}). + Where("next_retry_unix = 0 OR next_retry_unix <= ?", now). + Order("created_at ASC"). + Limit(limit). + Find(&tasks).Error + if err != nil { + return err + } + for i := range tasks { + if err := r.handleTask(ctx, &tasks[i]); err != nil { + global.GVA_LOG.Warn("处理PCDN任务失败", zap.Uint("id", tasks[i].ID), zap.Error(err)) + } + } + return nil +} + +func (r *RuntimeService) handleTask(ctx context.Context, task *pcdnModel.PcdnDispatchTask) error { + locked, err := r.markRunning(ctx, task.ID) + if err != nil || !locked { + return err + } + + result := r.dispatcher.Dispatch(ctx, DispatchRequest{ + TaskID: task.TaskID, + TraceID: task.TraceID, + ContentID: task.ContentID, + TargetNode: task.CurrentNodeID, + TimeoutSec: task.TimeoutSeconds, + }) + + if result.Success { + return global.GVA_DB.WithContext(ctx). + Model(&pcdnModel.PcdnDispatchTask{}). + Where("id = ?", task.ID). + Updates(map[string]any{"status": pcdnModel.DispatchStatusSuccess, "last_error": ""}).Error + } + + nextNode := NextFallbackNode(*task) + retryCount := task.RetryCount + 1 + updates := map[string]any{ + "retry_count": retryCount, + "last_error": result.Message, + } + if nextNode != 0 { + updates["current_node_id"] = nextNode + } + if retryCount <= task.MaxRetry { + updates["status"] = pcdnModel.DispatchStatusRetrying + updates["next_retry_unix"] = time.Now().Unix() + int64(retryCount*2) + } else { + updates["status"] = pcdnModel.DispatchStatusFailed + if nextNode == 0 { + updates["last_error"] = fmt.Sprintf("%s; no fallback node available", result.Message) + } + } + return global.GVA_DB.WithContext(ctx). + Model(&pcdnModel.PcdnDispatchTask{}). + Where("id = ?", task.ID). + Updates(updates).Error +} + +func (r *RuntimeService) markRunning(ctx context.Context, id uint) (bool, error) { + res := global.GVA_DB.WithContext(ctx). + Model(&pcdnModel.PcdnDispatchTask{}). + Where("id = ? AND status IN ?", id, []string{pcdnModel.DispatchStatusPending, pcdnModel.DispatchStatusRetrying}). + Updates(map[string]any{"status": pcdnModel.DispatchStatusRunning, "next_retry_unix": 0}) + if res.Error != nil { + return false, res.Error + } + return res.RowsAffected == 1, nil +} + +// SyncTimeoutTasks 将长时间 Running 的任务回收为重试态。 +func (r *RuntimeService) SyncTimeoutTasks(ctx context.Context) error { + var running []pcdnModel.PcdnDispatchTask + if err := global.GVA_DB.WithContext(ctx).Where("status = ?", pcdnModel.DispatchStatusRunning).Find(&running).Error; err != nil { + return err + } + now := time.Now().Unix() + for _, t := range running { + deadline := t.UpdatedAt.Unix() + int64(maxInt(t.TimeoutSeconds, 8)) + if now <= deadline { + continue + } + status := pcdnModel.DispatchStatusRetrying + if t.RetryCount >= t.MaxRetry { + status = pcdnModel.DispatchStatusFailed + } + nextRetry := int64(0) + if status == pcdnModel.DispatchStatusRetrying { + nextRetry = now + 2 + } + err := global.GVA_DB.WithContext(ctx). + Model(&pcdnModel.PcdnDispatchTask{}). + Where("id = ?", t.ID). + Updates(map[string]any{ + "status": status, + "retry_count": t.RetryCount + 1, + "next_retry_unix": nextRetry, + "last_error": "dispatch timeout", + }).Error + if err != nil { + global.GVA_LOG.Warn("回收PCDN超时任务失败", zap.Uint("task_id", t.ID), zap.Error(err)) + } + } + return nil +} + +func maxInt(a, b int) int { + if a > b { + return a + } + return b +} + +var PcdnRuntimeService = NewRuntimeService(nil) diff --git a/server/service/pcdn/scheduler.go b/server/service/pcdn/scheduler.go new file mode 100644 index 0000000..1652af7 --- /dev/null +++ b/server/service/pcdn/scheduler.go @@ -0,0 +1,180 @@ +package pcdn + +import ( + "context" + "errors" + "sort" + "strconv" + + "github.com/flipped-aurora/gin-vue-admin/server/global" + computenodeModel "github.com/flipped-aurora/gin-vue-admin/server/model/computenode" + pcdnModel "github.com/flipped-aurora/gin-vue-admin/server/model/pcdn" + "gorm.io/gorm" +) + +// ScheduleInput 调度输入。 +type ScheduleInput struct { + TaskID string + TraceID string + ContentID string + UserRegion string + UserISP string + TopN int + MaxRetry int + TimeoutSec int + Protocol string + NodeMetrics map[uint]RealtimeNodeMetric +} + +// CandidateNode 调度候选结果。 +type CandidateNode struct { + NodeID uint `json:"nodeId"` + Name string `json:"name"` + Region string `json:"region"` + ISP string `json:"isp"` + Score float64 `json:"score"` +} + +type SchedulerService struct { + policy *PolicyEngine +} + +var PcdnSchedulerService = NewSchedulerService() + +func NewSchedulerService() *SchedulerService { + return &SchedulerService{policy: NewPolicyEngine(ScoreWeights{})} +} + +// Schedule 执行“约束过滤 + 加权评分”,并记录到 PcdnDispatchTask。 +func (s *SchedulerService) Schedule(ctx context.Context, in ScheduleInput) ([]CandidateNode, error) { + if in.TaskID == "" || in.TraceID == "" { + return nil, errors.New("task_id and trace_id are required") + } + if in.ContentID == "" || in.UserRegion == "" { + return nil, errors.New("content_id and user_region are required") + } + if in.TopN <= 0 { + in.TopN = 3 + } + if in.MaxRetry <= 0 { + in.MaxRetry = 3 + } + if in.TimeoutSec <= 0 { + in.TimeoutSec = 8 + } + if in.Protocol == "" { + in.Protocol = "mock" + } + + var nodes []computenodeModel.ComputeNode + if err := global.GVA_DB.WithContext(ctx).Where("deleted_at IS NULL").Find(&nodes).Error; err != nil { + return nil, err + } + + candidates := make([]CandidateNode, 0, len(nodes)) + for _, node := range nodes { + if node.ID == 0 || node.IsOnShelf == nil || !*node.IsOnShelf { + continue + } + m, ok := in.NodeMetrics[node.ID] + if !ok { + m = RealtimeNodeMetric{NodeID: node.ID, Online: true, HealthScore: 80} + } + if !m.Online || m.PolicyDisabled || m.LoadPercent >= 90 { + continue + } + if node.DockerStatus != nil && *node.DockerStatus == "failed" { + continue + } + + regionAffinity := 0.92 + if node.Region != nil && *node.Region == in.UserRegion { + regionAffinity = 1.08 + } + ispAffinity := 0.95 + if in.UserISP == "" || m.ISP == "" || m.ISP == in.UserISP { + ispAffinity = 1.03 + } + + score := s.policy.Score(m) * regionAffinity * ispAffinity + candidates = append(candidates, CandidateNode{ + NodeID: node.ID, + Name: safeString(node.Name), + Region: safeString(node.Region), + ISP: m.ISP, + Score: score, + }) + } + + sort.Slice(candidates, func(i, j int) bool { + return candidates[i].Score > candidates[j].Score + }) + if len(candidates) > in.TopN { + candidates = candidates[:in.TopN] + } + + if err := s.upsertTask(ctx, in, candidates); err != nil { + return nil, err + } + return candidates, nil +} + +func (s *SchedulerService) upsertTask(ctx context.Context, in ScheduleInput, candidates []CandidateNode) error { + var task pcdnModel.PcdnDispatchTask + err := global.GVA_DB.WithContext(ctx).Where("task_id = ? AND trace_id = ?", in.TaskID, in.TraceID).First(&task).Error + if err == nil { + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + + candidateList := make([]any, 0, len(candidates)) + for _, c := range candidates { + candidateList = append(candidateList, map[string]any{ + "node_id": c.NodeID, + "score": c.Score, + "region": c.Region, + "isp": c.ISP, + }) + } + metricsSnapshot := map[string]any{} + for nodeID, m := range in.NodeMetrics { + metricsSnapshot[strconv.FormatUint(uint64(nodeID), 10)] = map[string]any{ + "latency_ms": m.LatencyMS, + "unit_cost": m.UnitCost, + "load_percent": m.LoadPercent, + "health_score": m.HealthScore, + "online": m.Online, + "policy_disabled": m.PolicyDisabled, + "isp": m.ISP, + } + } + + newTask := pcdnModel.PcdnDispatchTask{ + TaskID: in.TaskID, + TraceID: in.TraceID, + ContentID: in.ContentID, + UserRegion: in.UserRegion, + UserISP: in.UserISP, + TopN: in.TopN, + Status: pcdnModel.DispatchStatusPending, + MaxRetry: in.MaxRetry, + TimeoutSeconds: in.TimeoutSec, + Candidates: map[string]any{"list": candidateList}, + MetricsSnapshot: metricsSnapshot, + DispatchProtocol: in.Protocol, + } + if len(candidates) > 0 { + newTask.PrimaryNodeID = candidates[0].NodeID + newTask.CurrentNodeID = candidates[0].NodeID + } + return global.GVA_DB.WithContext(ctx).Create(&newTask).Error +} + +func safeString(v *string) string { + if v == nil { + return "" + } + return *v +} diff --git a/server/task/pcdn_dispatch.go b/server/task/pcdn_dispatch.go new file mode 100644 index 0000000..3f0fbb3 --- /dev/null +++ b/server/task/pcdn_dispatch.go @@ -0,0 +1,24 @@ +package task + +import ( + "context" + + pcdnService "github.com/flipped-aurora/gin-vue-admin/server/service/pcdn" + "go.uber.org/zap" + + "github.com/flipped-aurora/gin-vue-admin/server/global" +) + +// ProcessPcdnDispatchTasks 执行PCDN调度任务下发。 +func ProcessPcdnDispatchTasks() { + if err := pcdnService.PcdnRuntimeService.ProcessPendingTasks(context.Background(), 50); err != nil { + global.GVA_LOG.Error("处理PCDN调度任务失败", zap.Error(err)) + } +} + +// SyncPcdnDispatchStatus 同步处理超时和重试状态。 +func SyncPcdnDispatchStatus() { + if err := pcdnService.PcdnRuntimeService.SyncTimeoutTasks(context.Background()); err != nil { + global.GVA_LOG.Error("同步PCDN调度状态失败", zap.Error(err)) + } +}