From be39f2a03c54620c402b545ed6936ae8134d44af Mon Sep 17 00:00:00 2001 From: Hungerdream <1710233908@qq.com> Date: Tue, 28 Jul 2026 14:29:26 +0800 Subject: [PATCH] feat(delivery): expose rollback retry, release ack and CloudDM retry APIs - POST /delivery/tasks/:id/rollback/retry (platform admin) re-launches the rollback job after a previous cleanup failure - POST /delivery/tasks/:id/rollback/release (platform admin) releases bookkeeping after the operator verified the target host manually - POST /delivery/tasks/:id/clouddm/retry retries only the CloudDM registration step for an already healthy instance - task logs now include rollback AWX job stdout and localized text/classes for the new rollback states --- server/internal/handler/delivery.go | 49 +++++++++++++++++++++++++++++ server/internal/handler/task_log.go | 29 +++++++++++++++-- server/internal/router/router.go | 3 ++ 3 files changed, 78 insertions(+), 3 deletions(-) diff --git a/server/internal/handler/delivery.go b/server/internal/handler/delivery.go index f36b948..34a124a 100644 --- a/server/internal/handler/delivery.go +++ b/server/internal/handler/delivery.go @@ -175,6 +175,55 @@ func (h *DeliveryHandler) Cancel(c *gin.Context) { c.JSON(http.StatusOK, gin.H{"ok": true}) } +// RetryRollback retries the compensating AWX job after a previous cleanup failure. +func (h *DeliveryHandler) RetryRollback(c *gin.Context) { + if !requirePlatformAdmin(c) { + return + } + if err := h.service.RetryRollback(c.Request.Context(), c.Param("id")); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusAccepted, gin.H{"ok": true, "status": model.TaskRollbackPending}) +} + +// AcknowledgeRollbackRelease releases bookkeeping after an administrator has +// independently verified that no instance artifacts remain on the target host. +func (h *DeliveryHandler) AcknowledgeRollbackRelease(c *gin.Context) { + if !requirePlatformAdmin(c) { + return + } + if err := h.service.AcknowledgeRollbackRelease(c.Request.Context(), c.Param("id")); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"ok": true, "status": model.TaskRollbackAck}) +} + +// RetryCloudDMRegistration retries only the CloudDM registration step for an +// already healthy and accounted-for MySQL instance. +func (h *DeliveryHandler) RetryCloudDMRegistration(c *gin.Context) { + claims, ok := CurrentClaims(c) + if !ok { + c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"}) + return + } + taskID := c.Param("id") + if _, _, err := h.service.GetTask(c.Request.Context(), taskID, claims.UserID, claims.IsAdmin); err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + c.JSON(http.StatusNotFound, gin.H{"error": "task not found"}) + return + } + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + if err := h.service.RetryCloudDMRegistration(c.Request.Context(), taskID); err != nil { + c.JSON(http.StatusConflict, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusAccepted, gin.H{"ok": true, "status": model.TaskFinished}) +} + // Targets 获取可用部署目标 // @Summary 获取可用部署目标 // @Description 从 AWX 动态返回可用 Job Template 及其 Inventory hosts diff --git a/server/internal/handler/task_log.go b/server/internal/handler/task_log.go index ee0ad27..617f4a2 100644 --- a/server/internal/handler/task_log.go +++ b/server/internal/handler/task_log.go @@ -138,7 +138,7 @@ func (h *TaskLogHandler) getAWXTask(c *gin.Context, taskID string, userID uint64 lines = append(lines, taskLogLine{Time: formatTaskLogTime(event.CreatedAt), Message: "[" + event.ToState + "] " + event.Message, Class: classForTaskStatus(event.ToState)}) } var execution model.ExecutionJob - if err := h.db.WithContext(c.Request.Context()).Where("task_id = ?", task.ID).First(&execution).Error; err == nil && execution.ExecutorJobID != "" && execution.ExecutorJobID != "pending" { + if err := h.db.WithContext(c.Request.Context()).Where("task_id = ?", task.ID).First(&execution).Error; err == nil && execution.ExecutorJobID != "" && execution.ExecutorJobID != "pending" && !strings.HasPrefix(execution.ExecutorJobID, "pending-") { stdout, stdoutErr := h.delivery.AWXJobStdout(c.Request.Context(), execution.ExecutorJobID) if stdoutErr != nil { lines = append(lines, taskLogLine{Time: formatTaskLogTime(time.Now()), Message: "[awx] stdout fetch failed: " + stdoutErr.Error(), Class: "err"}) @@ -146,6 +146,15 @@ func (h *TaskLogHandler) getAWXTask(c *gin.Context, taskID string, userID uint64 lines = append(lines, splitStdoutLines(stdout)...) } } + var rollback model.RollbackJob + if err := h.db.WithContext(c.Request.Context()).Where("task_id = ?", task.ID).First(&rollback).Error; err == nil && rollback.ExecutorJobID != "" && rollback.ExecutorJobID != "pending" && !strings.HasPrefix(rollback.ExecutorJobID, "pending-") { + stdout, stdoutErr := h.delivery.AWXJobStdout(c.Request.Context(), rollback.ExecutorJobID) + if stdoutErr != nil { + lines = append(lines, taskLogLine{Time: formatTaskLogTime(time.Now()), Message: "[rollback awx] stdout fetch failed: " + stdoutErr.Error(), Class: "err"}) + } else { + lines = append(lines, splitStdoutLines(stdout)...) + } + } c.JSON(http.StatusOK, gin.H{"task": awxTaskSummary(*task), "lines": lines}) } @@ -257,8 +266,18 @@ func textForTaskStatus(status string) string { return "等待" case model.TaskRunning, model.TaskRegistering, model.TaskCanceling: return "执行中" + case model.TaskRollbackPending, model.TaskRollingBack: + return "回退中" case model.TaskFinished: return "成功" + case model.TaskRolledBack: + return "已回退" + case model.TaskRollbackFailed: + return "回退失败" + case model.TaskRollbackAck: + return "已确认释放" + case model.TaskRegisterFailed: + return "注册失败(实例保留)" case model.TaskCanceled: return "已取消" default: @@ -270,10 +289,14 @@ func classForTaskStatus(status string) string { switch status { case model.TaskFinished: return "ok" - case model.TaskExecutionFailed, model.TaskValidationFailed, model.TaskRegisterFailed, model.TaskCanceled: + case model.TaskExecutionFailed, model.TaskValidationFailed, model.TaskCanceled, model.TaskRollbackFailed: return "err" - case model.TaskRunning, model.TaskDispatching, model.TaskRegistering, model.TaskCanceling: + case model.TaskRollbackAck, model.TaskRegisterFailed: return "warn" + case model.TaskRunning, model.TaskDispatching, model.TaskRegistering, model.TaskCanceling, model.TaskRollbackPending, model.TaskRollingBack: + return "warn" + case model.TaskRolledBack: + return "ok" default: return "" } diff --git a/server/internal/router/router.go b/server/internal/router/router.go index a1bb599..a5ddfc7 100644 --- a/server/internal/router/router.go +++ b/server/internal/router/router.go @@ -149,6 +149,9 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { protected.GET("/delivery/tasks", deliveryHandler.List) protected.GET("/delivery/tasks/:id", deliveryHandler.Get) protected.POST("/delivery/tasks/:id/cancel", deliveryHandler.Cancel) + protected.POST("/delivery/tasks/:id/rollback/retry", deliveryHandler.RetryRollback) + protected.POST("/delivery/tasks/:id/rollback/release", deliveryHandler.AcknowledgeRollbackRelease) + protected.POST("/delivery/tasks/:id/clouddm/retry", deliveryHandler.RetryCloudDMRegistration) protected.GET("/task-logs", taskLogHandler.List) protected.GET("/task-logs/:id", taskLogHandler.Get) }