mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-03 23:06:36 +08:00
feat(cloudflare): add DNS pointing integration
This commit is contained in:
@@ -102,6 +102,8 @@ func DispatchTask(c *gin.Context) {
|
||||
// @Security SessionCookie
|
||||
// @Param status query string false "状态筛选 (pending/running/succeeded/failed)"
|
||||
// @Param task_type query string false "任务类型筛选"
|
||||
// @Param task_type_prefix query string false "任务类型前缀筛选(与 task_type / task_types 互斥,精确类型优先)"
|
||||
// @Param task_types query string false "逗号分隔的精确任务类型列表(IN 筛选,优先于前缀)"
|
||||
// @Param page query int false "页码" default(1)
|
||||
// @Param page_size query int false "每页条数" default(20)
|
||||
// @Success 200 {object} response.Any{data=object} "任务执行记录列表"
|
||||
|
||||
@@ -394,9 +394,22 @@ func ListAvailableDomains(ctx context.Context) ([]AvailableDomain, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
zones, err := repository.ListZones(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
zoneRoots := make(map[uint]string, len(zones))
|
||||
for i := range zones {
|
||||
zoneRoots[zones[i].ID] = zones[i].Domain
|
||||
}
|
||||
items := make([]AvailableDomain, 0, len(domains))
|
||||
for _, domain := range domains {
|
||||
items = append(items, AvailableDomain{ID: domain.ID, ZoneID: domain.ZoneID, Domain: domain.Domain})
|
||||
items = append(items, AvailableDomain{
|
||||
ID: domain.ID,
|
||||
ZoneID: domain.ZoneID,
|
||||
Domain: domain.Domain,
|
||||
ZoneDomain: zoneRoots[domain.ZoneID],
|
||||
})
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
@@ -32,14 +33,62 @@ const (
|
||||
TaskTypeSyncByNode = "of_cloudflare_sync_by_node"
|
||||
)
|
||||
|
||||
// SyncMemberMeta describes one-member reconciliation.
|
||||
var SyncMemberMeta = task.TaskMeta{Type: TaskTypeSyncMember, AsynqTask: SyncMemberTask, Name: "Cloudflare 域名同步", Description: "同步单个域名的 Cloudflare A 记录", MaxRetry: 3, Queue: task.QueueDefault, Retryable: true, InternalOnly: true}
|
||||
// SyncMemberMeta describes one-member reconciliation (admin-dispatchable).
|
||||
var SyncMemberMeta = task.TaskMeta{
|
||||
Type: TaskTypeSyncMember,
|
||||
AsynqTask: SyncMemberTask,
|
||||
Name: "Cloudflare 域名同步",
|
||||
Description: "同步单个域名的 Cloudflare A 记录",
|
||||
SupportsTime: false,
|
||||
MaxRetry: 3,
|
||||
Queue: task.QueueDefault,
|
||||
Retryable: true,
|
||||
Params: []task.TaskParam{
|
||||
{
|
||||
Name: "member_id",
|
||||
Label: "成员 ID",
|
||||
Type: "number",
|
||||
Required: true,
|
||||
Placeholder: "请输入 Cloudflare 指向成员 ID",
|
||||
Description: "of_cf_pointing_members 表中的成员主键 ID",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
// SyncGroupMeta describes group reconciliation.
|
||||
var SyncGroupMeta = task.TaskMeta{Type: TaskTypeSyncGroup, AsynqTask: SyncGroupTask, Name: "Cloudflare 分组同步", Description: "同步指向分组内全部域名", MaxRetry: 2, Queue: task.QueueDefault, Retryable: true, InternalOnly: true}
|
||||
// SyncGroupMeta describes group reconciliation (admin-dispatchable).
|
||||
var SyncGroupMeta = task.TaskMeta{
|
||||
Type: TaskTypeSyncGroup,
|
||||
AsynqTask: SyncGroupTask,
|
||||
Name: "Cloudflare 分组同步",
|
||||
Description: "同步指向分组内全部域名",
|
||||
SupportsTime: false,
|
||||
MaxRetry: 2,
|
||||
Queue: task.QueueDefault,
|
||||
Retryable: true,
|
||||
Params: []task.TaskParam{
|
||||
{
|
||||
Name: "group_id",
|
||||
Label: "分组 ID",
|
||||
Type: "number",
|
||||
Required: true,
|
||||
Placeholder: "请输入 Cloudflare 指向分组 ID",
|
||||
Description: "of_cf_pointing_groups 表中的分组主键 ID",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
// SyncByNodeMeta describes node-triggered reconciliation.
|
||||
var SyncByNodeMeta = task.TaskMeta{Type: TaskTypeSyncByNode, AsynqTask: SyncByNodeTask, Name: "Cloudflare 节点同步", Description: "同步当前指向指定节点的全部域名", MaxRetry: 2, Queue: task.QueueDefault, Retryable: true, InternalOnly: true}
|
||||
// SyncByNodeMeta describes node-triggered reconciliation (internal only).
|
||||
var SyncByNodeMeta = task.TaskMeta{
|
||||
Type: TaskTypeSyncByNode,
|
||||
AsynqTask: SyncByNodeTask,
|
||||
Name: "Cloudflare 节点同步",
|
||||
Description: "同步当前指向指定节点的全部域名",
|
||||
SupportsTime: false,
|
||||
MaxRetry: 2,
|
||||
Queue: task.QueueDefault,
|
||||
Retryable: true,
|
||||
InternalOnly: true,
|
||||
}
|
||||
|
||||
// SyncMemberPayload identifies one member.
|
||||
type SyncMemberPayload struct {
|
||||
@@ -93,9 +142,15 @@ type SyncMemberTaskHandler struct{}
|
||||
|
||||
// ValidatePayload validates a one-member task payload.
|
||||
func (handler *SyncMemberTaskHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
||||
if len(payload) == 0 {
|
||||
return nil, errors.New("任务参数不能为空")
|
||||
}
|
||||
var input SyncMemberPayload
|
||||
if err := decodePayload(payload, &input); err != nil || input.MemberID == 0 {
|
||||
return nil, errors.New("无效的 Cloudflare 成员同步参数")
|
||||
if err := decodePayload(payload, &input); err != nil {
|
||||
return nil, fmt.Errorf("无效的 Cloudflare 成员同步参数: %w", err)
|
||||
}
|
||||
if input.MemberID == 0 {
|
||||
return nil, errors.New("成员 ID 不能为空或零")
|
||||
}
|
||||
return json.Marshal(input)
|
||||
}
|
||||
@@ -108,11 +163,47 @@ func (handler *SyncMemberTaskHandler) Execute(ctx context.Context, payload []byt
|
||||
}
|
||||
var input SyncMemberPayload
|
||||
_ = json.Unmarshal(normalized, &input)
|
||||
task.AppendLog(ctx, "正在同步 Cloudflare 成员 ID=%d", input.MemberID)
|
||||
|
||||
state, loadErr := repository.GetCFPointingMemberContext(ctx, input.MemberID)
|
||||
if loadErr != nil {
|
||||
task.AppendLog(ctx, "加载成员上下文失败: member_id=%d error=%v", input.MemberID, loadErr)
|
||||
} else {
|
||||
task.AppendLog(ctx,
|
||||
"开始域名同步: domain=%s zone=%s group=%s(#%d) node=%s(%s) proxied=%v member_id=%d",
|
||||
state.Domain.Domain,
|
||||
state.Zone.Domain,
|
||||
state.Group.Name,
|
||||
state.Group.ID,
|
||||
state.Node.Name,
|
||||
strings.TrimSpace(state.Node.IP),
|
||||
state.Member.Proxied,
|
||||
input.MemberID,
|
||||
)
|
||||
}
|
||||
|
||||
if err = ReconcileMember(ctx, input.MemberID); err != nil {
|
||||
if state != nil {
|
||||
task.AppendLog(ctx, "域名同步失败: domain=%s member_id=%d error=%v",
|
||||
state.Domain.Domain, input.MemberID, err)
|
||||
} else {
|
||||
task.AppendLog(ctx, "域名同步失败: member_id=%d error=%v", input.MemberID, err)
|
||||
}
|
||||
return nil, fmt.Errorf("%s: %w", errSyncFailed, err)
|
||||
}
|
||||
return &task.TaskResult{Message: "Cloudflare 域名同步成功"}, nil
|
||||
|
||||
message := "Cloudflare 域名同步成功"
|
||||
if state != nil {
|
||||
ip := strings.TrimSpace(state.Node.IP)
|
||||
message = fmt.Sprintf("Cloudflare 域名同步成功: %s → %s (proxied=%v)",
|
||||
state.Domain.Domain, ip, state.Member.Proxied)
|
||||
task.AppendLog(ctx,
|
||||
"域名同步成功: domain=%s desired_ip=%s proxied=%v group=%s node=%s",
|
||||
state.Domain.Domain, ip, state.Member.Proxied, state.Group.Name, state.Node.Name,
|
||||
)
|
||||
} else {
|
||||
task.AppendLog(ctx, "域名同步成功: member_id=%d", input.MemberID)
|
||||
}
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
// SyncGroupTaskHandler reconciles every member in a group.
|
||||
@@ -120,9 +211,15 @@ type SyncGroupTaskHandler struct{}
|
||||
|
||||
// ValidatePayload validates a group task payload.
|
||||
func (handler *SyncGroupTaskHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
||||
if len(payload) == 0 {
|
||||
return nil, errors.New("任务参数不能为空")
|
||||
}
|
||||
var input SyncGroupPayload
|
||||
if err := decodePayload(payload, &input); err != nil || input.GroupID == 0 {
|
||||
return nil, errors.New("无效的 Cloudflare 分组同步参数")
|
||||
if err := decodePayload(payload, &input); err != nil {
|
||||
return nil, fmt.Errorf("无效的 Cloudflare 分组同步参数: %w", err)
|
||||
}
|
||||
if input.GroupID == 0 {
|
||||
return nil, errors.New("分组 ID 不能为空或零")
|
||||
}
|
||||
return json.Marshal(input)
|
||||
}
|
||||
@@ -137,8 +234,27 @@ func (handler *SyncGroupTaskHandler) Execute(ctx context.Context, payload []byte
|
||||
if err = json.Unmarshal(normalized, &input); err != nil {
|
||||
return nil, task.PermanentError(err.Error())
|
||||
}
|
||||
|
||||
scopeName := fmt.Sprintf("#%d", input.GroupID)
|
||||
activeNode := ""
|
||||
if group, groupErr := repository.GetCFPointingGroup(ctx, input.GroupID); groupErr != nil {
|
||||
task.AppendLog(ctx, "加载分组失败: group_id=%d error=%v", input.GroupID, groupErr)
|
||||
} else {
|
||||
scopeName = group.Name
|
||||
if node, nodeErr := repository.GetOpenFlareNodeByID(ctx, group.ActiveNodeID); nodeErr != nil {
|
||||
task.AppendLog(ctx, "加载生效节点失败: group=%s active_node_id=%d error=%v",
|
||||
group.Name, group.ActiveNodeID, nodeErr)
|
||||
} else {
|
||||
activeNode = fmt.Sprintf("%s(%s)", node.Name, strings.TrimSpace(node.IP))
|
||||
}
|
||||
task.AppendLog(ctx,
|
||||
"准备分组同步: group=%s id=%d enabled=%v active_node=%s default_proxied=%v",
|
||||
group.Name, group.ID, group.Enabled, activeNode, group.DefaultProxied,
|
||||
)
|
||||
}
|
||||
|
||||
members, err := repository.ListCFPointingMembersByGroupID(ctx, input.GroupID)
|
||||
return executeBatchSync(ctx, members, err, "分组")
|
||||
return executeBatchSync(ctx, members, err, "分组", scopeName, input.GroupID, activeNode)
|
||||
}
|
||||
|
||||
// SyncByNodeTaskHandler reconciles every member targeting a node.
|
||||
@@ -163,20 +279,66 @@ func (handler *SyncByNodeTaskHandler) Execute(ctx context.Context, payload []byt
|
||||
if err = json.Unmarshal(normalized, &input); err != nil {
|
||||
return nil, task.PermanentError(err.Error())
|
||||
}
|
||||
|
||||
scopeName := fmt.Sprintf("#%d", input.NodeID)
|
||||
activeNode := ""
|
||||
if node, nodeErr := repository.GetOpenFlareNodeByID(ctx, input.NodeID); nodeErr != nil {
|
||||
task.AppendLog(ctx, "加载节点失败: node_id=%d error=%v", input.NodeID, nodeErr)
|
||||
} else {
|
||||
scopeName = node.Name
|
||||
activeNode = fmt.Sprintf("%s(%s)", node.Name, strings.TrimSpace(node.IP))
|
||||
task.AppendLog(ctx, "准备节点同步: node=%s id=%d ip=%s",
|
||||
node.Name, node.ID, strings.TrimSpace(node.IP))
|
||||
}
|
||||
|
||||
members, err := repository.ListCFPointingMembersByActiveNodeID(ctx, input.NodeID)
|
||||
return executeBatchSync(ctx, members, err, "节点")
|
||||
return executeBatchSync(ctx, members, err, "节点", scopeName, input.NodeID, activeNode)
|
||||
}
|
||||
|
||||
func executeBatchSync(ctx context.Context, members []model.CFPointingMember, listErr error, scope string) (*task.TaskResult, error) {
|
||||
func executeBatchSync(
|
||||
ctx context.Context,
|
||||
members []model.CFPointingMember,
|
||||
listErr error,
|
||||
scope, scopeName string,
|
||||
scopeID uint,
|
||||
activeNode string,
|
||||
) (*task.TaskResult, error) {
|
||||
if listErr != nil {
|
||||
task.AppendLog(ctx, "列出%s成员失败: name=%s id=%d error=%v",
|
||||
scope, scopeName, scopeID, listErr)
|
||||
return nil, listErr
|
||||
}
|
||||
for _, member := range members {
|
||||
|
||||
task.AppendLog(ctx, "开始%s同步: name=%s id=%d active_node=%s 域名数=%d",
|
||||
scope, scopeName, scopeID, activeNode, len(members))
|
||||
if len(members) == 0 {
|
||||
message := fmt.Sprintf("Cloudflare %s同步完成: %s 无域名成员", scope, scopeName)
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
for index, member := range members {
|
||||
domainName := fmt.Sprintf("zone_domain_id=%d", member.ZoneDomainID)
|
||||
if domain, domainErr := repository.GetZoneDomainByID(ctx, member.ZoneDomainID); domainErr == nil {
|
||||
domainName = domain.Domain
|
||||
}
|
||||
task.AppendLog(ctx, "[%d/%d] 同步域名 %s (member_id=%d proxied=%v)",
|
||||
index+1, len(members), domainName, member.ID, member.Proxied)
|
||||
if err := ReconcileMember(ctx, member.ID); err != nil {
|
||||
task.AppendLog(ctx, "[%d/%d] 失败: domain=%s member_id=%d error=%v",
|
||||
index+1, len(members), domainName, member.ID, err)
|
||||
return nil, err
|
||||
}
|
||||
task.AppendLog(ctx, "[%d/%d] 成功: domain=%s", index+1, len(members), domainName)
|
||||
}
|
||||
return &task.TaskResult{Message: fmt.Sprintf("Cloudflare %s同步完成,共 %d 个域名", scope, len(members))}, nil
|
||||
|
||||
message := fmt.Sprintf("Cloudflare %s同步完成: %s 共 %d 个域名", scope, scopeName, len(members))
|
||||
if activeNode != "" {
|
||||
message = fmt.Sprintf("Cloudflare %s同步完成: %s → %s,共 %d 个域名",
|
||||
scope, scopeName, activeNode, len(members))
|
||||
}
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
func decodePayload(payload []byte, target any) error {
|
||||
|
||||
@@ -85,9 +85,10 @@ type GroupDetail struct {
|
||||
|
||||
// AvailableDomain is a ZoneDomain eligible for pointing.
|
||||
type AvailableDomain struct {
|
||||
ID uint `json:"id"`
|
||||
ZoneID uint `json:"zone_id"`
|
||||
Domain string `json:"domain"`
|
||||
ID uint `json:"id"`
|
||||
ZoneID uint `json:"zone_id"`
|
||||
Domain string `json:"domain"`
|
||||
ZoneDomain string `json:"zone_domain"`
|
||||
}
|
||||
|
||||
// Overview summarizes Cloudflare pointing readiness and sync health.
|
||||
|
||||
@@ -54,8 +54,12 @@ func (TaskExecution) TableName() string {
|
||||
|
||||
// ListTaskExecutionsRequest 查询任务执行记录列表请求
|
||||
type ListTaskExecutionsRequest struct {
|
||||
Status string `form:"status"`
|
||||
TaskType string `form:"task_type"`
|
||||
Page int `form:"page"`
|
||||
PageSize int `form:"page_size"`
|
||||
Status string `form:"status"`
|
||||
TaskType string `form:"task_type"`
|
||||
TaskTypePrefix string `form:"task_type_prefix"`
|
||||
// TaskTypes is a comma-separated list of exact asynq task types (IN filter).
|
||||
// Used when TaskType is empty; takes precedence over TaskTypePrefix.
|
||||
TaskTypes string `form:"task_types"`
|
||||
Page int `form:"page"`
|
||||
PageSize int `form:"page_size"`
|
||||
}
|
||||
|
||||
@@ -146,6 +146,10 @@ func ListTaskExecutions(ctx context.Context, req model.ListTaskExecutionsRequest
|
||||
}
|
||||
if req.TaskType != "" {
|
||||
query = query.Where("task_type = ?", req.TaskType)
|
||||
} else if types := parseTaskTypesFilter(req.TaskTypes); len(types) > 0 {
|
||||
query = query.Where("task_type IN ?", types)
|
||||
} else if req.TaskTypePrefix != "" {
|
||||
query = query.Where("task_type LIKE ?", req.TaskTypePrefix+"%")
|
||||
}
|
||||
|
||||
var total int64
|
||||
@@ -165,6 +169,21 @@ func ListTaskExecutions(ctx context.Context, req model.ListTaskExecutionsRequest
|
||||
return executions, total, nil
|
||||
}
|
||||
|
||||
func parseTaskTypesFilter(raw string) []string {
|
||||
if strings.TrimSpace(raw) == "" {
|
||||
return nil
|
||||
}
|
||||
parts := strings.Split(raw, ",")
|
||||
out := make([]string, 0, len(parts))
|
||||
for _, part := range parts {
|
||||
part = strings.TrimSpace(part)
|
||||
if part != "" {
|
||||
out = append(out, part)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// MarkFailedTaskExecutionsSucceededTx marks failed executions of a task type as succeeded within a transaction.
|
||||
func MarkFailedTaskExecutionsSucceededTx(
|
||||
tx *gorm.DB,
|
||||
|
||||
@@ -415,6 +415,34 @@ func TestListTaskExecutions(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(1), total)
|
||||
assert.Equal(t, "list_001", items[0].TaskID)
|
||||
|
||||
// 按类型前缀筛选
|
||||
items, total, err = ListTaskExecutions(ctx, model.ListTaskExecutionsRequest{TaskTypePrefix: "system:", Page: 1, PageSize: 10})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(3), total)
|
||||
assert.Len(t, items, 3)
|
||||
|
||||
// 按多类型 IN 筛选
|
||||
items, total, err = ListTaskExecutions(ctx, model.ListTaskExecutionsRequest{
|
||||
TaskTypes: "system:cleanup,other:task",
|
||||
Page: 1,
|
||||
PageSize: 10,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(5), total)
|
||||
assert.Len(t, items, 5)
|
||||
|
||||
// 精确类型优先于 task_types / 前缀
|
||||
items, total, err = ListTaskExecutions(ctx, model.ListTaskExecutionsRequest{
|
||||
TaskType: "other:task",
|
||||
TaskTypes: "system:cleanup",
|
||||
TaskTypePrefix: "system:",
|
||||
Page: 1,
|
||||
PageSize: 10,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(2), total)
|
||||
assert.Len(t, items, 2)
|
||||
}
|
||||
|
||||
func TestListTaskExecutionsDefaultPaging(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user