From 32dc7ef68ea4a7085018ad38c35c4b02f0eb0801 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 29 May 2026 10:34:25 +0800 Subject: [PATCH] =?UTF-8?q?[=E6=96=B0=E5=A2=9E]=20=E5=AE=9E=E7=8E=B0?= =?UTF-8?q?=E8=8A=82=E7=82=B9=E5=BC=BA=E5=88=B6=E5=90=8C=E6=AD=A5=E5=8A=9F?= =?UTF-8?q?=E8=83=BD=EF=BC=8C=E5=85=81=E8=AE=B8=E9=80=9A=E8=BF=87=20API=20?= =?UTF-8?q?=E8=AF=B7=E6=B1=82=E5=BC=BA=E5=88=B6=E5=90=8C=E6=AD=A5=E9=85=8D?= =?UTF-8?q?=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openflare_agent/internal/agent/runner.go | 13 ++++++++++ .../internal/protocol/agent_api.go | 11 ++++---- openflare_agent/internal/sync/service.go | 12 +++++++++ openflare_server/controller/node.go | 24 +++++++++++++++++ openflare_server/router/api-router.go | 10 ++++--- openflare_server/service/agent_ws.go | 21 +++++++++++---- openflare_server/service/node.go | 16 ++++++++++++ .../web/features/nodes/api/nodes.ts | 6 +++++ .../nodes/components/node-detail-page.tsx | 26 +++++++++++++++++++ 9 files changed, 125 insertions(+), 14 deletions(-) diff --git a/openflare_agent/internal/agent/runner.go b/openflare_agent/internal/agent/runner.go index e8fa5cc4..505a9321 100644 --- a/openflare_agent/internal/agent/runner.go +++ b/openflare_agent/internal/agent/runner.go @@ -23,6 +23,7 @@ type HeartbeatService interface { type SyncService interface { SyncOnStartup(ctx context.Context, target *protocol.ActiveConfigMeta) error SyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error + ForceSyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error } type Updater interface { @@ -294,6 +295,18 @@ func (r *Runner) handleWebSocketMessage(ctx context.Context, message protocol.WS slog.Error("agent ws triggered sync failed", "version", target.Version, "error", err) } return false, nil + case protocol.WSMessageTypeForceSyncConfig: + var target protocol.ActiveConfigMeta + if err := json.Unmarshal(message.Payload, &target); err != nil { + slog.Debug("agent ws force sync config decode failed", "error", err) + return false, nil + } + slog.Debug("agent ws force sync config received", "version", target.Version, "checksum", target.Checksum, "trigger_sync", true) + if err := r.SyncService.ForceSyncOnce(ctx, &target); err != nil { + r.recordSyncError(err) + slog.Error("agent ws triggered force sync failed", "version", target.Version, "error", err) + } + return false, nil case protocol.WSMessageTypePing: slog.Debug("agent ws ping received") return false, conn.SendPong() diff --git a/openflare_agent/internal/protocol/agent_api.go b/openflare_agent/internal/protocol/agent_api.go index 64898141..15846178 100644 --- a/openflare_agent/internal/protocol/agent_api.go +++ b/openflare_agent/internal/protocol/agent_api.go @@ -33,11 +33,12 @@ type AgentSettings struct { } const ( - WSMessageTypeStatus = "status" - WSMessageTypeSettings = "settings" - WSMessageTypeActiveConfig = "active_config" - WSMessageTypePing = "ping" - WSMessageTypePong = "pong" + WSMessageTypeStatus = "status" + WSMessageTypeSettings = "settings" + WSMessageTypeActiveConfig = "active_config" + WSMessageTypeForceSyncConfig = "force_sync_config" + WSMessageTypePing = "ping" + WSMessageTypePong = "pong" ) type WSMessage struct { diff --git a/openflare_agent/internal/sync/service.go b/openflare_agent/internal/sync/service.go index 1c44a497..487e6ca7 100644 --- a/openflare_agent/internal/sync/service.go +++ b/openflare_agent/internal/sync/service.go @@ -137,6 +137,18 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ return s.applyIfNeeded(ctx, mode, startup, snapshot, currentChecksum, target, config) } +func (s *Service) ForceSyncOnce(ctx context.Context, target *protocol.ActiveConfigMeta) error { + snapshot, err := s.stateStore.Load() + if err != nil { + return err + } + if hasBlockedTarget(snapshot) { + clearBlockedTarget(snapshot) + _ = s.stateStore.Save(snapshot) + } + return s.SyncOnce(ctx, target) +} + func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool, snapshot *state.Snapshot, currentChecksum string, target *protocol.ActiveConfigMeta, config *protocol.ActiveConfigResponse) error { if currentChecksum == config.Checksum { slog.Debug("local openresty config already up to date", "mode", mode, "version", config.Version) diff --git a/openflare_server/controller/node.go b/openflare_server/controller/node.go index f11a0063..b8254e32 100644 --- a/openflare_server/controller/node.go +++ b/openflare_server/controller/node.go @@ -188,6 +188,30 @@ func RequestNodeOpenrestyRestart(c *gin.Context) { respondSuccess(c, node) } +// RequestNodeForceSync godoc +// @Summary Request force sync config on node +// @Tags Nodes +// @Produce json +// @Security BearerAuth +// @Param id path int true "Node ID" +// @Success 200 {object} map[string]interface{} +// @Failure 400 {object} map[string]interface{} +// @Router /api/nodes/{id}/force-sync [post] +func RequestNodeForceSync(c *gin.Context) { + id, err := strconv.ParseUint(c.Param("id"), 10, 64) + if err != nil || id == 0 { + respondBadRequest(c, "") + return + } + + node, err := service.RequestNodeForceSync(uint(id)) + if err != nil { + respondFailure(c, err.Error()) + return + } + respondSuccess(c, node) +} + // GetNodeAgentRelease godoc // @Summary Check latest agent release for node // @Tags Nodes diff --git a/openflare_server/router/api-router.go b/openflare_server/router/api-router.go index a1bd1cbb..a0cde9e3 100644 --- a/openflare_server/router/api-router.go +++ b/openflare_server/router/api-router.go @@ -165,12 +165,14 @@ func SetApiRouter(router *gin.Engine) { nodeRoute.GET("/", controller.GetNodes) nodeRoute.POST("/", controller.CreateNode) nodeRoute.GET("/:id/agent-release", controller.GetNodeAgentRelease) - nodeRoute.GET("/:id/observability", controller.GetNodeObservability) - nodeRoute.POST("/:id/observability/cleanup", controller.CleanupNodeHealthEvents) - nodeRoute.POST("/:id/agent-update", controller.RequestNodeAgentUpdate) - nodeRoute.POST("/:id/openresty-restart", controller.RequestNodeOpenrestyRestart) nodeRoute.POST("/:id/update", controller.UpdateNode) nodeRoute.POST("/:id/delete", controller.DeleteNode) + nodeRoute.POST("/:id/agent-update", controller.RequestNodeAgentUpdate) + nodeRoute.POST("/:id/openresty-restart", controller.RequestNodeOpenrestyRestart) + nodeRoute.POST("/:id/force-sync", controller.RequestNodeForceSync) + nodeRoute.GET("/:id/agent-release", controller.GetNodeAgentRelease) + nodeRoute.GET("/:id/observability", controller.GetNodeObservability) + nodeRoute.POST("/:id/observability/cleanup", controller.CleanupNodeHealthEvents) } applyLogRoute := apiRouter.Group("/apply-logs") applyLogRoute.Use(middleware.AdminAuth()) diff --git a/openflare_server/service/agent_ws.go b/openflare_server/service/agent_ws.go index b299ab8a..c8a6a555 100644 --- a/openflare_server/service/agent_ws.go +++ b/openflare_server/service/agent_ws.go @@ -7,11 +7,12 @@ import ( ) const ( - AgentWSMessageTypeStatus = "status" - AgentWSMessageTypeSettings = "settings" - AgentWSMessageTypeActiveConfig = "active_config" - AgentWSMessageTypePing = "ping" - AgentWSMessageTypePong = "pong" + AgentWSMessageTypeStatus = "status" + AgentWSMessageTypeSettings = "settings" + AgentWSMessageTypeActiveConfig = "active_config" + AgentWSMessageTypeForceSyncConfig = "force_sync_config" + AgentWSMessageTypePing = "ping" + AgentWSMessageTypePong = "pong" AgentWSConnectedLastSeenValue = "__OPENFLARE_WS_CONNECTED__" ) @@ -167,6 +168,16 @@ func SendAgentWSActiveConfig(nodeID string, activeConfig *ActiveConfigMeta) bool }) } +func SendAgentWSForceSyncConfig(nodeID string, activeConfig *ActiveConfigMeta) bool { + if activeConfig == nil { + return false + } + return sendAgentWSMessage(nodeID, AgentWSOutboundMessage{ + Type: AgentWSMessageTypeForceSyncConfig, + Payload: activeConfig, + }) +} + func SendAgentWSPong(nodeID string) bool { return sendAgentWSMessage(nodeID, AgentWSOutboundMessage{ Type: AgentWSMessageTypePong, diff --git a/openflare_server/service/node.go b/openflare_server/service/node.go index c3074890..903a3d47 100644 --- a/openflare_server/service/node.go +++ b/openflare_server/service/node.go @@ -189,6 +189,22 @@ func RequestNodeOpenrestyRestart(id uint) (*NodeView, error) { return buildNodeView(node), nil } +func RequestNodeForceSync(id uint) (*NodeView, error) { + node, err := model.GetNodeByID(id) + if err != nil { + return nil, err + } + activeConfig, err := GetActiveConfigMetaForAgent() + if err != nil { + return nil, errors.New("无法获取当前激活的配置版本:" + err.Error()) + } + if !SendAgentWSForceSyncConfig(node.NodeID, activeConfig) { + return nil, errors.New("节点不在线或通过 WebSocket 发送同步指令失败") + } + slog.Info("force sync requested via ws", "node_id", node.NodeID, "name", node.Name) + return buildNodeView(node), nil +} + func AuthenticateAgentToken(token string) (*model.Node, error) { token = strings.TrimSpace(token) if token == "" { diff --git a/openflare_server/web/features/nodes/api/nodes.ts b/openflare_server/web/features/nodes/api/nodes.ts index b2ac965a..162bcb3e 100644 --- a/openflare_server/web/features/nodes/api/nodes.ts +++ b/openflare_server/web/features/nodes/api/nodes.ts @@ -63,6 +63,12 @@ export function requestNodeAgentUpdate( }); } +export function requestNodeForceSync(id: number) { + return apiRequest(`/nodes/${id}/force-sync`, { + method: 'POST', + }); +} + export function requestNodeOpenrestyRestart(id: number) { return apiRequest(`/nodes/${id}/openresty-restart`, { method: 'POST', diff --git a/openflare_server/web/features/nodes/components/node-detail-page.tsx b/openflare_server/web/features/nodes/components/node-detail-page.tsx index 3544dd14..d69eaa2e 100644 --- a/openflare_server/web/features/nodes/components/node-detail-page.tsx +++ b/openflare_server/web/features/nodes/components/node-detail-page.tsx @@ -25,8 +25,10 @@ import { getNodeAgentRelease, getNodeObservability, getNodes, + requestNodeForceSync, requestNodeOpenrestyRestart, requestNodeAgentUpdate, + rotateNodeBootstrapToken, updateNode, } from '@/features/nodes/api/nodes'; import { NodeEditorModal } from '@/features/nodes/components/node-editor-modal'; @@ -365,6 +367,20 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { }, }); + const forceSyncMutation = useMutation({ + mutationFn: () => requestNodeForceSync(Number(nodeId)), + onSuccess: async (updatedNode) => { + setFeedback({ + tone: 'success', + message: `已向节点 ${updatedNode.name} 下发强制同步指令,无视当前错误拦截。`, + }); + await queryClient.invalidateQueries({ queryKey: nodesQueryKey }); + }, + onError: (error) => { + setFeedback({ tone: 'danger', message: getErrorMessage(error) }); + }, + }); + const deleteMutation = useMutation({ mutationFn: () => deleteNode(Number(nodeId)), onSuccess: async () => { @@ -650,6 +666,16 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) { > {isRefreshing ? '刷新中...' : '刷新'} + { + setFeedback(null); + forceSyncMutation.mutate(); + }} + disabled={forceSyncMutation.isPending} + > + {forceSyncMutation.isPending ? '同步中...' : '同步'} +