61eeb0c8ac
- 修复 Wayne 任务状态映射,将数字状态统一为字符串状态 - 添加分页支持,支持 page/page_size 参数 - 添加 AWXJobStdout 内存缓存,减少重复 API 调用 - 修复 SubscribeTask 内存泄漏,自动清理 closed channel - 统一 Stream 端点,AWX 使用 pub/sub,Wayne 使用轮询 背景:task log 接口存在多个潜在问题,包括状态判断错误导致无限轮询、 缺少分页、重复 API 调用性能问题、内存泄漏风险等 关联 commit:fix/logs 分支
701 lines
23 KiB
Go
701 lines
23 KiB
Go
package handler
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"net/http"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/1024XEngineer/xinfra/server/internal/auth"
|
|
"github.com/1024XEngineer/xinfra/server/internal/model"
|
|
"github.com/1024XEngineer/xinfra/server/internal/service"
|
|
"github.com/gin-gonic/gin"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
type TaskLogHandler struct {
|
|
db *gorm.DB
|
|
delivery *service.DeliveryService
|
|
wayne *service.WayneRoleBindingService
|
|
}
|
|
|
|
func NewTaskLogHandler(db *gorm.DB, delivery *service.DeliveryService, wayne *service.WayneRoleBindingService) *TaskLogHandler {
|
|
return &TaskLogHandler{db: db, delivery: delivery, wayne: wayne}
|
|
}
|
|
|
|
type taskLogSummary struct {
|
|
ID string `json:"id"`
|
|
Source string `json:"source"`
|
|
Service string `json:"service"`
|
|
Name string `json:"name"`
|
|
Runner string `json:"runner"`
|
|
Status string `json:"status"`
|
|
StatusText string `json:"status_text"`
|
|
StatusClass string `json:"status_class"`
|
|
BusinessLineID uint64 `json:"business_line_id"`
|
|
ReferenceID string `json:"reference_id"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
UpdatedAt time.Time `json:"updated_at"`
|
|
}
|
|
|
|
type taskLogLine struct {
|
|
Time string `json:"time"`
|
|
Message string `json:"message"`
|
|
Class string `json:"class"`
|
|
}
|
|
|
|
// List 聚合 AWX 交付任务和 Wayne 部署服务任务。
|
|
func (h *TaskLogHandler) List(c *gin.Context) {
|
|
claims, ok := CurrentClaims(c)
|
|
if !ok {
|
|
c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"})
|
|
return
|
|
}
|
|
source := strings.ToLower(strings.TrimSpace(c.DefaultQuery("source", "all")))
|
|
businessLineID, _ := strconv.ParseUint(c.Query("business_line_id"), 10, 64)
|
|
|
|
// 分页参数
|
|
page, _ := strconv.Atoi(c.DefaultQuery("page", "1"))
|
|
pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "20"))
|
|
if page < 1 {
|
|
page = 1
|
|
}
|
|
if pageSize < 1 || pageSize > 100 {
|
|
pageSize = 20
|
|
}
|
|
|
|
items := make([]taskLogSummary, 0)
|
|
if source == "" || source == "all" || source == "awx" {
|
|
awxItems, err := h.listAWXTasks(c, claims.UserID, claims.IsAdmin, businessLineID)
|
|
if err != nil {
|
|
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
|
return
|
|
}
|
|
items = append(items, awxItems...)
|
|
}
|
|
if source == "" || source == "all" || source == "wayne" {
|
|
wayneItems, err := h.listWayneTasks(c, claims.UserID, claims.IsAdmin, businessLineID)
|
|
if err != nil {
|
|
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
|
return
|
|
}
|
|
items = append(items, wayneItems...)
|
|
}
|
|
sortTaskLogSummaries(items)
|
|
|
|
// 计算分页
|
|
total := len(items)
|
|
start := (page - 1) * pageSize
|
|
if start >= total {
|
|
items = make([]taskLogSummary, 0)
|
|
} else {
|
|
end := start + pageSize
|
|
if end > total {
|
|
end = total
|
|
}
|
|
items = items[start:end]
|
|
}
|
|
|
|
c.JSON(http.StatusOK, gin.H{
|
|
"items": items,
|
|
"total": total,
|
|
"page": page,
|
|
"page_size": pageSize,
|
|
"total_pages": (total + pageSize - 1) / pageSize,
|
|
})
|
|
}
|
|
|
|
func (h *TaskLogHandler) Get(c *gin.Context) {
|
|
claims, ok := CurrentClaims(c)
|
|
if !ok {
|
|
c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"})
|
|
return
|
|
}
|
|
id := c.Param("id")
|
|
switch {
|
|
case strings.HasPrefix(id, "awx:"):
|
|
h.getAWXTask(c, strings.TrimPrefix(id, "awx:"), claims.UserID, claims.IsAdmin)
|
|
case strings.HasPrefix(id, "wayne:publish:"):
|
|
h.getWayneTask(c, id, claims.UserID, claims.IsAdmin)
|
|
default:
|
|
c.JSON(http.StatusBadRequest, gin.H{"error": "unknown task log source"})
|
|
}
|
|
}
|
|
|
|
func (h *TaskLogHandler) listAWXTasks(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]taskLogSummary, error) {
|
|
tasks, err := h.delivery.ListTasks(c.Request.Context(), userID, isAdmin, service.DeliveryTaskListFilter{BusinessLineID: businessLineID})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
items := make([]taskLogSummary, 0, len(tasks))
|
|
for _, task := range tasks {
|
|
items = append(items, awxTaskSummary(task))
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
func (h *TaskLogHandler) listWayneTasks(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]taskLogSummary, error) {
|
|
namespaces, err := h.visibleWayneNamespaces(c, userID, isAdmin, businessLineID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
histories, err := h.wayne.ListDeploymentHistories(c.Request.Context(), namespaces, 100)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
items := make([]taskLogSummary, 0, len(histories))
|
|
for _, history := range histories {
|
|
items = append(items, wayneTaskSummary(history))
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
func (h *TaskLogHandler) getAWXTask(c *gin.Context, taskID string, userID uint64, isAdmin bool) {
|
|
task, events, err := h.delivery.GetTask(c.Request.Context(), taskID, userID, isAdmin)
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
c.JSON(http.StatusNotFound, gin.H{"error": "task log not found"})
|
|
return
|
|
}
|
|
if err != nil {
|
|
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
|
return
|
|
}
|
|
lines := make([]taskLogLine, 0, len(events)+16)
|
|
for _, event := range events {
|
|
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" && !strings.HasPrefix(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"})
|
|
} else {
|
|
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})
|
|
}
|
|
|
|
func (h *TaskLogHandler) getWayneTask(c *gin.Context, id string, userID uint64, isAdmin bool) {
|
|
resourceID, historyID, ok := parseWaynePublishTaskID(id)
|
|
if !ok {
|
|
c.JSON(http.StatusBadRequest, gin.H{"error": "invalid Wayne task log id"})
|
|
return
|
|
}
|
|
namespaces, err := h.visibleWayneNamespaces(c, userID, isAdmin, 0)
|
|
if err != nil {
|
|
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
|
|
return
|
|
}
|
|
history, err := h.wayne.GetDeploymentHistory(c.Request.Context(), namespaces, resourceID, historyID)
|
|
if err != nil {
|
|
c.JSON(http.StatusNotFound, gin.H{"error": err.Error()})
|
|
return
|
|
}
|
|
lines := []taskLogLine{
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[wayne] publish history #" + strconv.FormatInt(history.ID, 10), Class: classForWaynePublishStatus(history.Status)},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[deployment] " + history.ResourceName + " resource_id=" + strconv.FormatInt(history.ResourceID, 10), Class: ""},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[cluster] " + history.Cluster + " template_id=" + strconv.FormatInt(history.TemplateID, 10), Class: ""},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[user] " + history.User, Class: ""},
|
|
}
|
|
if strings.TrimSpace(history.Message) != "" {
|
|
lines = append(lines, taskLogLine{Time: formatTaskLogTime(history.CreatedAt), Message: "[message] " + history.Message, Class: classForWaynePublishStatus(history.Status)})
|
|
}
|
|
c.JSON(http.StatusOK, gin.H{"task": wayneTaskSummary(*history), "lines": lines})
|
|
}
|
|
|
|
func awxTaskSummary(task model.DeliveryTask) taskLogSummary {
|
|
return taskLogSummary{
|
|
ID: "awx:" + task.ID,
|
|
Source: "awx",
|
|
Service: "mysql",
|
|
Name: "MySQL 标准化交付 · " + task.InstanceName,
|
|
Runner: "AWX Job Template #" + strconv.FormatUint(task.TargetID, 10),
|
|
Status: task.Status,
|
|
StatusText: textForTaskStatus(task.Status),
|
|
StatusClass: classForTaskStatus(task.Status),
|
|
BusinessLineID: task.BusinessLineID,
|
|
ReferenceID: task.ID,
|
|
CreatedAt: task.CreatedAt,
|
|
UpdatedAt: task.UpdatedAt,
|
|
}
|
|
}
|
|
|
|
func wayneTaskSummary(history service.WayneDeploymentHistory) taskLogSummary {
|
|
return taskLogSummary{
|
|
ID: waynePublishTaskID(history),
|
|
Source: "wayne",
|
|
Service: "wayne-deployment",
|
|
Name: wayneDeploymentTaskName(history),
|
|
Runner: "Wayne Native API",
|
|
Status: wayneStatusToString(history.Status),
|
|
StatusText: textForWaynePublishStatus(history.Status),
|
|
StatusClass: classForWaynePublishStatus(history.Status),
|
|
BusinessLineID: history.BusinessLineID,
|
|
ReferenceID: strconv.FormatInt(history.ID, 10),
|
|
CreatedAt: history.CreatedAt,
|
|
UpdatedAt: history.CreatedAt,
|
|
}
|
|
}
|
|
|
|
func (h *TaskLogHandler) visibleWayneNamespaces(c *gin.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]model.BusinessLineWayneNamespace, error) {
|
|
query := h.db.WithContext(c.Request.Context()).Order("business_line_id ASC, wayne_namespace_id ASC")
|
|
if businessLineID != 0 {
|
|
query = query.Where("business_line_id = ?", businessLineID)
|
|
}
|
|
if !isAdmin {
|
|
query = query.Where("business_line_id IN (?)", h.db.Model(&model.BusinessLineUser{}).Select("business_line_id").Where("user_id = ?", userID))
|
|
}
|
|
var namespaces []model.BusinessLineWayneNamespace
|
|
if err := query.Find(&namespaces).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
return namespaces, nil
|
|
}
|
|
|
|
func waynePublishTaskID(history service.WayneDeploymentHistory) string {
|
|
return "wayne:publish:" + strconv.FormatInt(history.ResourceID, 10) + ":" + strconv.FormatInt(history.ID, 10)
|
|
}
|
|
|
|
func parseWaynePublishTaskID(id string) (int64, int64, bool) {
|
|
parts := strings.Split(id, ":")
|
|
if len(parts) != 4 || parts[0] != "wayne" || parts[1] != "publish" {
|
|
return 0, 0, false
|
|
}
|
|
resourceID, resourceErr := strconv.ParseInt(parts[2], 10, 64)
|
|
historyID, historyErr := strconv.ParseInt(parts[3], 10, 64)
|
|
if resourceErr != nil || historyErr != nil {
|
|
return 0, 0, false
|
|
}
|
|
return resourceID, historyID, true
|
|
}
|
|
|
|
func wayneDeploymentTaskName(history service.WayneDeploymentHistory) string {
|
|
name := strings.TrimSpace(history.ResourceName)
|
|
if name == "" {
|
|
name = strconv.FormatInt(history.ResourceID, 10)
|
|
}
|
|
return "Wayne 服务部署 · " + name
|
|
}
|
|
|
|
func textForTaskStatus(status string) string {
|
|
switch status {
|
|
case model.TaskPending, model.TaskValidating, model.TaskDispatching:
|
|
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:
|
|
return "失败"
|
|
}
|
|
}
|
|
|
|
func classForTaskStatus(status string) string {
|
|
switch status {
|
|
case model.TaskFinished:
|
|
return "ok"
|
|
case model.TaskExecutionFailed, model.TaskValidationFailed, model.TaskCanceled, model.TaskRollbackFailed:
|
|
return "err"
|
|
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 ""
|
|
}
|
|
}
|
|
|
|
func textForWaynePublishStatus(status int) string {
|
|
switch status {
|
|
case 1:
|
|
return "成功"
|
|
case 0:
|
|
return "失败"
|
|
default:
|
|
return "未知"
|
|
}
|
|
}
|
|
|
|
func classForWaynePublishStatus(status int) string {
|
|
switch status {
|
|
case 1:
|
|
return "ok"
|
|
case 0:
|
|
return "err"
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
// wayneStatusToString 将 Wayne 数字状态映射为与 AWX 一致的字符串状态
|
|
func wayneStatusToString(status int) string {
|
|
switch status {
|
|
case 1:
|
|
return model.TaskFinished
|
|
case 0:
|
|
return model.TaskExecutionFailed
|
|
default:
|
|
return "unknown"
|
|
}
|
|
}
|
|
|
|
func splitStdoutLines(stdout string) []taskLogLine {
|
|
lines := make([]taskLogLine, 0)
|
|
for _, line := range strings.Split(stdout, "\n") {
|
|
line = strings.TrimRight(line, "\r")
|
|
if strings.TrimSpace(line) == "" {
|
|
continue
|
|
}
|
|
lines = append(lines, taskLogLine{Time: "", Message: line, Class: classForOutputLine(line)})
|
|
}
|
|
return lines
|
|
}
|
|
|
|
func classForOutputLine(line string) string {
|
|
lower := strings.ToLower(line)
|
|
switch {
|
|
case strings.Contains(lower, "failed") || strings.Contains(lower, "fatal") || strings.Contains(lower, "error"):
|
|
return "err"
|
|
case strings.Contains(lower, "ok:") || strings.Contains(lower, "successful") || strings.Contains(lower, "success"):
|
|
return "ok"
|
|
case strings.Contains(lower, "changed:"):
|
|
return "tag-ok"
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
// Stream SSE 端点,用于实时推送任务日志。
|
|
// AWX 任务使用 pub/sub 模式,Wayne 任务使用轮询模式。
|
|
func (h *TaskLogHandler) Stream(c *gin.Context) {
|
|
claims, ok := CurrentClaims(c)
|
|
if !ok {
|
|
c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"})
|
|
return
|
|
}
|
|
|
|
taskID := c.Query("task_id")
|
|
if taskID == "" {
|
|
c.JSON(http.StatusBadRequest, gin.H{"error": "missing task_id parameter"})
|
|
return
|
|
}
|
|
|
|
// 设置 SSE 响应头
|
|
c.Header("Content-Type", "text/event-stream")
|
|
c.Header("Cache-Control", "no-cache")
|
|
c.Header("Connection", "keep-alive")
|
|
c.Header("X-Accel-Buffering", "no")
|
|
|
|
ctx := c.Request.Context()
|
|
|
|
// AWX 任务使用 pub/sub 模式
|
|
if strings.HasPrefix(taskID, "awx:") {
|
|
h.streamAWXTask(c, taskID, ctx, claims)
|
|
return
|
|
}
|
|
|
|
// Wayne 任务使用轮询模式(因为 Wayne API 不支持推送)
|
|
if strings.HasPrefix(taskID, "wayne:") {
|
|
h.streamWayneTask(c, taskID, ctx, claims)
|
|
return
|
|
}
|
|
|
|
c.SSEvent("message", gin.H{"type": "error", "error": "invalid task_id format"})
|
|
c.Writer.Flush()
|
|
}
|
|
|
|
// streamAWXTask AWX 任务使用 pub/sub 模式
|
|
func (h *TaskLogHandler) streamAWXTask(c *gin.Context, taskID string, ctx context.Context, claims *auth.Claims) {
|
|
actualTaskID := strings.TrimPrefix(taskID, "awx:")
|
|
|
|
// 发送初始日志
|
|
initialLines, taskStatus, err := h.fetchTaskLogLines(ctx, taskID, claims.UserID, claims.IsAdmin)
|
|
if err != nil {
|
|
c.SSEvent("message", gin.H{"type": "error", "error": err.Error()})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
// 发送初始数据
|
|
c.SSEvent("message", gin.H{"type": "init", "lines": initialLines})
|
|
c.Writer.Flush()
|
|
|
|
// 如果任务已完成,直接发送结束事件
|
|
if isTerminalStatus(taskStatus) {
|
|
c.SSEvent("message", gin.H{"type": "finished"})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
// 订阅任务更新(使用 pub/sub 模式)
|
|
updates, cancel := h.delivery.SubscribeTask(actualTaskID)
|
|
defer cancel()
|
|
|
|
heartbeatTicker := time.NewTicker(15 * time.Second)
|
|
defer heartbeatTicker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case _, ok := <-updates:
|
|
if !ok {
|
|
return
|
|
}
|
|
|
|
// 重新获取日志行
|
|
lines, status, err := h.fetchTaskLogLines(ctx, taskID, claims.UserID, claims.IsAdmin)
|
|
if err != nil {
|
|
c.SSEvent("message", gin.H{"type": "error", "error": err.Error()})
|
|
c.Writer.Flush()
|
|
continue
|
|
}
|
|
|
|
// 发送完整日志
|
|
c.SSEvent("message", gin.H{"type": "update", "lines": lines})
|
|
c.Writer.Flush()
|
|
|
|
// 如果任务完成,发送结束事件
|
|
if isTerminalStatus(status) {
|
|
c.SSEvent("message", gin.H{"type": "finished"})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
case <-heartbeatTicker.C:
|
|
c.SSEvent("heartbeat", nil)
|
|
c.Writer.Flush()
|
|
}
|
|
}
|
|
}
|
|
|
|
// streamWayneTask Wayne 任务使用轮询模式
|
|
func (h *TaskLogHandler) streamWayneTask(c *gin.Context, taskID string, ctx context.Context, claims *auth.Claims) {
|
|
// 发送初始日志
|
|
initialLines, taskStatus, err := h.fetchTaskLogLines(ctx, taskID, claims.UserID, claims.IsAdmin)
|
|
if err != nil {
|
|
c.SSEvent("message", gin.H{"type": "error", "error": err.Error()})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
// 发送初始数据
|
|
c.SSEvent("message", gin.H{"type": "init", "lines": initialLines})
|
|
c.Writer.Flush()
|
|
|
|
// 如果任务已完成,直接发送结束事件
|
|
if isTerminalStatus(taskStatus) {
|
|
c.SSEvent("message", gin.H{"type": "finished"})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
// 轮询循环
|
|
ticker := time.NewTicker(5 * time.Second)
|
|
defer ticker.Stop()
|
|
|
|
heartbeatTicker := time.NewTicker(15 * time.Second)
|
|
defer heartbeatTicker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
lines, status, err := h.fetchTaskLogLines(ctx, taskID, claims.UserID, claims.IsAdmin)
|
|
if err != nil {
|
|
c.SSEvent("message", gin.H{"type": "error", "error": err.Error()})
|
|
c.Writer.Flush()
|
|
continue
|
|
}
|
|
|
|
// 发送完整日志
|
|
c.SSEvent("message", gin.H{"type": "update", "lines": lines})
|
|
c.Writer.Flush()
|
|
|
|
// 如果任务完成,发送结束事件
|
|
if isTerminalStatus(status) {
|
|
c.SSEvent("message", gin.H{"type": "finished"})
|
|
c.Writer.Flush()
|
|
return
|
|
}
|
|
|
|
case <-heartbeatTicker.C:
|
|
c.SSEvent("heartbeat", nil)
|
|
c.Writer.Flush()
|
|
}
|
|
}
|
|
}
|
|
|
|
// parseActualTaskID 从 task log ID 解析出实际的任务 ID
|
|
func (h *TaskLogHandler) parseActualTaskID(taskID string) string {
|
|
switch {
|
|
case strings.HasPrefix(taskID, "awx:"):
|
|
return strings.TrimPrefix(taskID, "awx:")
|
|
case strings.HasPrefix(taskID, "wayne:publish:"):
|
|
// Wayne 任务需要使用 resourceID 作为订阅 key
|
|
parts := strings.Split(taskID, ":")
|
|
if len(parts) == 4 {
|
|
return "wayne:" + parts[2]
|
|
}
|
|
return ""
|
|
default:
|
|
return ""
|
|
}
|
|
}
|
|
|
|
// fetchTaskLogLines 获取任务的日志行和状态。
|
|
func (h *TaskLogHandler) fetchTaskLogLines(ctx context.Context, taskID string, userID uint64, isAdmin bool) ([]taskLogLine, string, error) {
|
|
switch {
|
|
case strings.HasPrefix(taskID, "awx:"):
|
|
return h.fetchAWXLogLines(ctx, strings.TrimPrefix(taskID, "awx:"), userID, isAdmin)
|
|
case strings.HasPrefix(taskID, "wayne:publish:"):
|
|
return h.fetchWayneLogLines(ctx, taskID, userID, isAdmin)
|
|
default:
|
|
return nil, "", errors.New("unknown task log source")
|
|
}
|
|
}
|
|
|
|
// fetchAWXLogLines 获取 AWX 任务的日志行。
|
|
func (h *TaskLogHandler) fetchAWXLogLines(ctx context.Context, taskID string, userID uint64, isAdmin bool) ([]taskLogLine, string, error) {
|
|
task, events, err := h.delivery.GetTask(ctx, taskID, userID, isAdmin)
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return nil, "", errors.New("task log not found")
|
|
}
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
|
|
lines := make([]taskLogLine, 0, len(events)+16)
|
|
for _, event := range events {
|
|
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(ctx).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(ctx, execution.ExecutorJobID)
|
|
if stdoutErr != nil {
|
|
lines = append(lines, taskLogLine{Time: formatTaskLogTime(time.Now()), Message: "[awx] stdout fetch failed: " + stdoutErr.Error(), Class: "err"})
|
|
} else {
|
|
lines = append(lines, splitStdoutLines(stdout)...)
|
|
}
|
|
}
|
|
|
|
var rollback model.RollbackJob
|
|
if err := h.db.WithContext(ctx).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(ctx, 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)...)
|
|
}
|
|
}
|
|
|
|
return lines, task.Status, nil
|
|
}
|
|
|
|
// fetchWayneLogLines 获取 Wayne 任务的日志行。
|
|
func (h *TaskLogHandler) fetchWayneLogLines(ctx context.Context, id string, userID uint64, isAdmin bool) ([]taskLogLine, string, error) {
|
|
resourceID, historyID, ok := parseWaynePublishTaskID(id)
|
|
if !ok {
|
|
return nil, "", errors.New("invalid Wayne task log id")
|
|
}
|
|
|
|
namespaces, err := h.visibleWayneNamespacesFromCtx(ctx, userID, isAdmin, 0)
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
|
|
history, err := h.wayne.GetDeploymentHistory(ctx, namespaces, resourceID, historyID)
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
|
|
lines := []taskLogLine{
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[wayne] publish history #" + strconv.FormatInt(history.ID, 10), Class: classForWaynePublishStatus(history.Status)},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[deployment] " + history.ResourceName + " resource_id=" + strconv.FormatInt(history.ResourceID, 10), Class: ""},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[cluster] " + history.Cluster + " template_id=" + strconv.FormatInt(history.TemplateID, 10), Class: ""},
|
|
{Time: formatTaskLogTime(history.CreatedAt), Message: "[user] " + history.User, Class: ""},
|
|
}
|
|
if strings.TrimSpace(history.Message) != "" {
|
|
lines = append(lines, taskLogLine{Time: formatTaskLogTime(history.CreatedAt), Message: "[message] " + history.Message, Class: classForWaynePublishStatus(history.Status)})
|
|
}
|
|
|
|
// 将 Wayne 数字状态映射为与 AWX 一致的字符串状态
|
|
statusText := wayneStatusToString(history.Status)
|
|
return lines, statusText, nil
|
|
}
|
|
|
|
// visibleWayneNamespacesFromCtx 从 context 获取可见的 Wayne 命名空间。
|
|
func (h *TaskLogHandler) visibleWayneNamespacesFromCtx(ctx context.Context, userID uint64, isAdmin bool, businessLineID uint64) ([]model.BusinessLineWayneNamespace, error) {
|
|
query := h.db.WithContext(ctx).Order("business_line_id ASC, wayne_namespace_id ASC")
|
|
if businessLineID != 0 {
|
|
query = query.Where("business_line_id = ?", businessLineID)
|
|
}
|
|
if !isAdmin {
|
|
query = query.Where("business_line_id IN (?)", h.db.Model(&model.BusinessLineUser{}).Select("business_line_id").Where("user_id = ?", userID))
|
|
}
|
|
var namespaces []model.BusinessLineWayneNamespace
|
|
if err := query.Find(&namespaces).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
return namespaces, nil
|
|
}
|
|
|
|
// isTerminalStatus 判断任务状态是否为终态。
|
|
func isTerminalStatus(status string) bool {
|
|
switch status {
|
|
case model.TaskFinished, model.TaskCanceled, model.TaskExecutionFailed,
|
|
model.TaskValidationFailed, model.TaskRegisterFailed,
|
|
model.TaskRolledBack, model.TaskRollbackFailed, model.TaskRollbackAck:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func formatTaskLogTime(t time.Time) string {
|
|
if t.IsZero() {
|
|
return ""
|
|
}
|
|
return t.Format("15:04:05")
|
|
}
|
|
|
|
func sortTaskLogSummaries(items []taskLogSummary) {
|
|
for i := 1; i < len(items); i++ {
|
|
item := items[i]
|
|
j := i - 1
|
|
for j >= 0 && items[j].UpdatedAt.Before(item.UpdatedAt) {
|
|
items[j+1] = items[j]
|
|
j--
|
|
}
|
|
items[j+1] = item
|
|
}
|
|
}
|