Files
2026-07-22 16:51:23 +08:00

196 lines
5.2 KiB
Go

package api
import (
"encoding/json"
"errors"
"net/http"
"strconv"
"github.com/gin-gonic/gin"
_ "oci-portal/internal/model" // swagger 注解引用
"oci-portal/internal/service"
)
// taskHandler 处理后台任务(定时测活、抢机)相关请求。
type taskHandler struct {
svc *service.TaskService
}
type createTaskRequest struct {
Name string `json:"name" binding:"required"`
Type string `json:"type" binding:"required"`
CronExpr string `json:"cronExpr" binding:"required"`
Payload json.RawMessage `json:"payload"`
}
type updateTaskRequest struct {
Name *string `json:"name"`
CronExpr *string `json:"cronExpr"`
Payload json.RawMessage `json:"payload"`
Status *string `json:"status"`
}
// @Summary 创建任务
// @Tags 任务与日志回传
// @Param body body createTaskRequest true "任务定义(类型/cron/payload)"
// @Success 201 {object} model.Task
// @Security BearerAuth
// @Router /api/v1/tasks [post]
func (h *taskHandler) create(c *gin.Context) {
var req createTaskRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
task, err := h.svc.CreateTask(c.Request.Context(), service.CreateTaskInput{
Name: req.Name,
Type: req.Type,
CronExpr: req.CronExpr,
Payload: req.Payload,
})
if err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusCreated, task)
}
// @Summary 任务列表
// @Tags 任务与日志回传
// @Success 200 {array} model.Task
// @Security BearerAuth
// @Router /api/v1/tasks [get]
func (h *taskHandler) list(c *gin.Context) {
tasks, err := h.svc.ListTasks(c.Request.Context())
if err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusOK, tasks)
}
// @Summary 任务详情
// @Tags 任务与日志回传
// @Param id path int true "任务 ID"
// @Success 200 {object} model.Task
// @Security BearerAuth
// @Router /api/v1/tasks/{id} [get]
func (h *taskHandler) get(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
task, err := h.svc.GetTask(c.Request.Context(), id)
if err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusOK, task)
}
// @Summary 更新任务
// @Tags 任务与日志回传
// @Param id path int true "任务 ID"
// @Param body body updateTaskRequest true "可更新字段;抢机 payload 的 count 语义为目标台数"
// @Success 200 {object} model.Task
// @Failure 400 {object} errorResponse "参数非法或抢机目标不大于已完成数量"
// @Failure 409 {object} errorResponse "任务被并发修改,须刷新重试"
// @Security BearerAuth
// @Router /api/v1/tasks/{id} [put]
func (h *taskHandler) update(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
var req updateTaskRequest
if err := c.ShouldBindJSON(&req); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
task, err := h.svc.UpdateTask(c.Request.Context(), id, service.UpdateTaskInput{
Name: req.Name,
CronExpr: req.CronExpr,
Payload: req.Payload,
Status: req.Status,
})
if err != nil {
respondUpdateTaskErr(c, err)
return
}
c.JSON(http.StatusOK, task)
}
// respondUpdateTaskErr 把任务编辑的哨兵错误映射为语义状态码,其余走统一错误边界。
func respondUpdateTaskErr(c *gin.Context, err error) {
switch {
case errors.Is(err, service.ErrTaskConflict):
c.JSON(http.StatusConflict, gin.H{"error": err.Error()})
case errors.Is(err, service.ErrSnatchTargetTooLow):
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
default:
respondError(c, err)
}
}
// @Summary 删除任务
// @Tags 任务与日志回传
// @Param id path int true "任务 ID"
// @Success 204 "无内容"
// @Security BearerAuth
// @Router /api/v1/tasks/{id} [delete]
func (h *taskHandler) remove(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
if err := h.svc.DeleteTask(c.Request.Context(), id); err != nil {
respondError(c, err)
return
}
c.Status(http.StatusNoContent)
}
// @Summary 任务执行日志
// @Tags 任务与日志回传
// @Param id path int true "配置 ID"
// @Success 200 {array} model.TaskLog
// @Security BearerAuth
// @Router /api/v1/tasks/{id}/logs [get]
func (h *taskHandler) logs(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
limit, _ := strconv.Atoi(c.DefaultQuery("limit", "50"))
logs, err := h.svc.TaskLogs(c.Request.Context(), id, limit)
if err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusOK, logs)
}
// @Summary 立即执行任务(异步触发,结果经任务日志轮询获取)
// @Tags 任务与日志回传
// @Param id path int true "配置 ID"
// @Success 202 {object} triggerResponse
// @Failure 409 {object} errorResponse "任务正在执行中"
// @Security BearerAuth
// @Router /api/v1/tasks/{id}/run [post]
func (h *taskHandler) run(c *gin.Context) {
id, ok := pathID(c)
if !ok {
return
}
if err := h.svc.TriggerTask(c.Request.Context(), id); err != nil {
if errors.Is(err, service.ErrTaskRunning) {
c.JSON(http.StatusConflict, gin.H{"error": err.Error()})
return
}
respondError(c, err)
return
}
c.JSON(http.StatusAccepted, gin.H{"triggered": true})
}