[新增] 实现节点强制同步功能,允许通过 API 请求强制同步配置

This commit is contained in:
ryan
2026-05-29 10:34:25 +08:00
parent 32762fdf3c
commit 32dc7ef68e
9 changed files with 125 additions and 14 deletions
+13
View File
@@ -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()
@@ -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 {
+12
View File
@@ -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)
+24
View File
@@ -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
+6 -4
View File
@@ -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())
+16 -5
View File
@@ -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,
+16
View File
@@ -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 == "" {
@@ -63,6 +63,12 @@ export function requestNodeAgentUpdate(
});
}
export function requestNodeForceSync(id: number) {
return apiRequest<NodeItem>(`/nodes/${id}/force-sync`, {
method: 'POST',
});
}
export function requestNodeOpenrestyRestart(id: number) {
return apiRequest<NodeItem>(`/nodes/${id}/openresty-restart`, {
method: 'POST',
@@ -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 ? '刷新中...' : '刷新'}
</SecondaryButton>
<SecondaryButton
type="button"
onClick={() => {
setFeedback(null);
forceSyncMutation.mutate();
}}
disabled={forceSyncMutation.isPending}
>
{forceSyncMutation.isPending ? '同步中...' : '同步'}
</SecondaryButton>
<PrimaryButton
type="button"
onClick={handleOpenAgentUpdateModal}