mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-09 00:56:37 +08:00
feat(cloudflare): add custom targets and node failover (#43)
* feat(cloudflare): add custom targets and node failover * fix(cloudflare): propagate batch sync errors --------- Co-authored-by: OpenAI Codex <codex@openai.com>
This commit is contained in:
@@ -25,6 +25,8 @@ const (
|
||||
SyncGroupTask = "cloudflare:sync_group"
|
||||
// SyncByNodeTask is the Asynq task type for members targeting one node.
|
||||
SyncByNodeTask = "cloudflare:sync_by_node"
|
||||
// FailoverCheckTask checks node availability and updates active Cloudflare targets.
|
||||
FailoverCheckTask = "cloudflare:failover_check"
|
||||
|
||||
// TaskTypeSyncMember is the task metadata type for member synchronization.
|
||||
TaskTypeSyncMember = "of_cloudflare_sync_member"
|
||||
@@ -32,14 +34,29 @@ const (
|
||||
TaskTypeSyncGroup = "of_cloudflare_sync_group"
|
||||
// TaskTypeSyncByNode is the task metadata type for node-triggered synchronization.
|
||||
TaskTypeSyncByNode = "of_cloudflare_sync_by_node"
|
||||
// TaskTypeFailoverCheck is the schedule metadata type for node failover checks.
|
||||
TaskTypeFailoverCheck = "of_cloudflare_failover_check"
|
||||
)
|
||||
|
||||
// FailoverCheckMeta describes the scheduled node failover reconciliation.
|
||||
var FailoverCheckMeta = task.TaskMeta{
|
||||
Type: TaskTypeFailoverCheck,
|
||||
AsynqTask: FailoverCheckTask,
|
||||
Name: "Cloudflare 节点故障回退检查",
|
||||
Description: "根据节点在线和应用状态切换分组主备节点",
|
||||
SupportsTime: false,
|
||||
MaxRetry: 2,
|
||||
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 记录",
|
||||
Description: "同步单个域名的 Cloudflare DNS 记录",
|
||||
SupportsTime: false,
|
||||
MaxRetry: 3,
|
||||
Queue: task.QueueDefault,
|
||||
@@ -169,14 +186,17 @@ func (handler *SyncMemberTaskHandler) Execute(ctx context.Context, payload []byt
|
||||
if loadErr != nil {
|
||||
task.AppendLog(ctx, "加载成员上下文失败: member_id=%d error=%v", input.MemberID, loadErr)
|
||||
} else {
|
||||
target := fmt.Sprintf("%s(%s)", state.Node.Name, strings.TrimSpace(state.Node.IP))
|
||||
if state.Group.TargetMode == model.CFPointingTargetModeCustom {
|
||||
target = fmt.Sprintf("%s %s", state.Group.RecordType, state.Group.RecordContent)
|
||||
}
|
||||
task.AppendLog(ctx,
|
||||
"开始域名同步: domain=%s zone=%s group=%s(#%d) node=%s(%s) proxied=%v member_id=%d",
|
||||
"开始域名同步: domain=%s zone=%s group=%s(#%d) target=%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),
|
||||
target,
|
||||
state.Member.Proxied,
|
||||
input.MemberID,
|
||||
)
|
||||
@@ -194,12 +214,15 @@ func (handler *SyncMemberTaskHandler) Execute(ctx context.Context, payload []byt
|
||||
|
||||
message := "Cloudflare 域名同步成功"
|
||||
if state != nil {
|
||||
ip := strings.TrimSpace(state.Node.IP)
|
||||
target := strings.TrimSpace(state.Node.IP)
|
||||
if state.Group.TargetMode == model.CFPointingTargetModeCustom {
|
||||
target = fmt.Sprintf("%s %s", state.Group.RecordType, state.Group.RecordContent)
|
||||
}
|
||||
message = fmt.Sprintf("Cloudflare 域名同步成功: %s → %s (proxied=%v)",
|
||||
state.Domain.Domain, ip, state.Member.Proxied)
|
||||
state.Domain.Domain, target, 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,
|
||||
"域名同步成功: domain=%s desired_target=%s proxied=%v group=%s target_mode=%s",
|
||||
state.Domain.Domain, target, state.Member.Proxied, state.Group.Name, state.Group.TargetMode,
|
||||
)
|
||||
} else {
|
||||
task.AppendLog(ctx, "域名同步成功: member_id=%d", input.MemberID)
|
||||
@@ -242,14 +265,9 @@ func (handler *SyncGroupTaskHandler) Execute(ctx context.Context, payload []byte
|
||||
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))
|
||||
}
|
||||
activeNode = groupSyncTarget(ctx, group)
|
||||
task.AppendLog(ctx,
|
||||
"准备分组同步: group=%s id=%d enabled=%v active_node=%s default_proxied=%v",
|
||||
"准备分组同步: group=%s id=%d enabled=%v target=%s default_proxied=%v",
|
||||
group.Name, group.ID, group.Enabled, activeNode, group.DefaultProxied,
|
||||
)
|
||||
}
|
||||
@@ -258,6 +276,19 @@ func (handler *SyncGroupTaskHandler) Execute(ctx context.Context, payload []byte
|
||||
return executeBatchSync(ctx, members, err, "分组", scopeName, input.GroupID, activeNode)
|
||||
}
|
||||
|
||||
func groupSyncTarget(ctx context.Context, group *model.CFPointingGroup) string {
|
||||
if group.TargetMode == model.CFPointingTargetModeCustom {
|
||||
return fmt.Sprintf("%s %s", group.RecordType, group.RecordContent)
|
||||
}
|
||||
node, err := repository.GetOpenFlareNodeByID(ctx, group.ActiveNodeID)
|
||||
if err != nil {
|
||||
task.AppendLog(ctx, "加载生效节点失败: group=%s active_node_id=%d error=%v",
|
||||
group.Name, group.ActiveNodeID, err)
|
||||
return ""
|
||||
}
|
||||
return fmt.Sprintf("%s(%s)", node.Name, strings.TrimSpace(node.IP))
|
||||
}
|
||||
|
||||
// SyncByNodeTaskHandler reconciles every member targeting a node.
|
||||
type SyncByNodeTaskHandler struct{}
|
||||
|
||||
@@ -310,7 +341,7 @@ func executeBatchSync(
|
||||
return nil, listErr
|
||||
}
|
||||
|
||||
task.AppendLog(ctx, "开始%s同步: name=%s id=%d active_node=%s 域名数=%d",
|
||||
task.AppendLog(ctx, "开始%s同步: name=%s id=%d target=%s 域名数=%d",
|
||||
scope, scopeName, scopeID, activeNode, len(members))
|
||||
if len(members) == 0 {
|
||||
message := fmt.Sprintf("Cloudflare %s同步完成: %s 无域名成员", scope, scopeName)
|
||||
@@ -364,3 +395,31 @@ func decodePayload(payload []byte, target any) error {
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// FailoverCheckTaskHandler checks enabled node groups for failover and recovery.
|
||||
type FailoverCheckTaskHandler struct{}
|
||||
|
||||
// ValidatePayload validates the empty scheduled-task payload.
|
||||
func (handler *FailoverCheckTaskHandler) ValidatePayload(payload []byte) ([]byte, error) {
|
||||
var input struct{}
|
||||
if err := decodePayload(payload, &input); err != nil {
|
||||
return nil, fmt.Errorf("无效的 Cloudflare 故障回退检查参数: %w", err)
|
||||
}
|
||||
return json.Marshal(input)
|
||||
}
|
||||
|
||||
// Execute switches affected groups and queues DNS reconciliation.
|
||||
func (handler *FailoverCheckTaskHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) {
|
||||
if _, err := handler.ValidatePayload(payload); err != nil {
|
||||
return nil, task.PermanentError(err.Error())
|
||||
}
|
||||
task.AppendLog(ctx, "开始检查 Cloudflare 节点分组主备状态")
|
||||
changed, err := CheckNodeFailover(ctx)
|
||||
if err != nil {
|
||||
task.AppendLog(ctx, "Cloudflare 节点故障回退检查失败: error=%v", err)
|
||||
return nil, err
|
||||
}
|
||||
message := fmt.Sprintf("Cloudflare 节点故障回退检查完成,切换分组 %d 个", changed)
|
||||
task.AppendLog(ctx, "%s", message)
|
||||
return &task.TaskResult{Message: message}, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user