feat: add task cancell
This commit is contained in:
@@ -63,6 +63,7 @@ const (
|
||||
WSTypeEnabled = "enabled"
|
||||
WSTypeFetchTasks = "fetch_tasks"
|
||||
WSTypeTaskHeartbeat = "task_heartbeat"
|
||||
WSTypeStop = "stop"
|
||||
)
|
||||
|
||||
// TablePrefix 表前缀,从配置文件读取
|
||||
|
||||
@@ -236,3 +236,19 @@ func (tc *TaskController) DeleteTask(c *gin.Context) {
|
||||
|
||||
utils.SuccessMsg(c, "删除成功")
|
||||
}
|
||||
|
||||
func (tc *TaskController) StopTask(c *gin.Context) {
|
||||
logID, err := strconv.ParseUint(c.Param("logID"), 10, 32)
|
||||
if err != nil {
|
||||
utils.BadRequest(c, "无效的日志ID")
|
||||
return
|
||||
}
|
||||
|
||||
err = tc.executorService.StopTaskExecution(uint(logID))
|
||||
if err != nil {
|
||||
utils.BadRequest(c, err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
utils.SuccessMsg(c, "停止请求已发送")
|
||||
}
|
||||
|
||||
@@ -170,7 +170,8 @@ type Scheduler struct {
|
||||
wg sync.WaitGroup
|
||||
mu sync.RWMutex
|
||||
logger SchedulerLogger
|
||||
runningTasks map[string]context.CancelFunc // 记录运行中的任务,用于停止
|
||||
runningTasks map[string]context.CancelFunc // 记录运行中的任务,用于停止 (TaskID -> CancelFunc)
|
||||
runningExecs map[uint]context.CancelFunc // 记录运行中的执行,用于停止 (LogID -> CancelFunc)
|
||||
}
|
||||
|
||||
// NewScheduler 创建调度器
|
||||
@@ -202,6 +203,7 @@ func NewScheduler(config SchedulerConfig, handler SchedulerEventHandler) *Schedu
|
||||
stopCh: make(chan struct{}),
|
||||
logger: &DefaultLogger{},
|
||||
runningTasks: make(map[string]context.CancelFunc),
|
||||
runningExecs: make(map[uint]context.CancelFunc),
|
||||
}
|
||||
|
||||
return s
|
||||
@@ -377,11 +379,17 @@ func (s *Scheduler) executeTask(req *ExecutionRequest) (*ExecutionResult, error)
|
||||
// 注册到运行中任务
|
||||
s.mu.Lock()
|
||||
s.runningTasks[req.TaskID] = cancel
|
||||
if req.LogID > 0 {
|
||||
s.runningExecs[req.LogID] = cancel
|
||||
}
|
||||
s.mu.Unlock()
|
||||
|
||||
defer func() {
|
||||
s.mu.Lock()
|
||||
delete(s.runningTasks, req.TaskID)
|
||||
if req.LogID > 0 {
|
||||
delete(s.runningExecs, req.LogID)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
}()
|
||||
|
||||
@@ -440,7 +448,7 @@ func (s *Scheduler) executeTask(req *ExecutionRequest) (*ExecutionResult, error)
|
||||
return result, execErr
|
||||
}
|
||||
|
||||
// StopTask 停止正在运行的任务
|
||||
// StopTask 停止正在运行的任务(通过 TaskID,可能会停止多个并发副本)
|
||||
func (s *Scheduler) StopTask(taskID string) bool {
|
||||
s.mu.RLock()
|
||||
cancel, exists := s.runningTasks[taskID]
|
||||
@@ -454,6 +462,20 @@ func (s *Scheduler) StopTask(taskID string) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
// StopLog 停止正在运行的任务(通过 LogID,精确停止单个执行副本)
|
||||
func (s *Scheduler) StopLog(logID uint) bool {
|
||||
s.mu.RLock()
|
||||
cancel, exists := s.runningExecs[logID]
|
||||
s.mu.RUnlock()
|
||||
|
||||
if exists && cancel != nil {
|
||||
cancel()
|
||||
s.logger.Infof("[Scheduler] 已尝试停止任务执行 #%d", logID)
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// GetRunningTaskCount 获取正在运行的任务数量
|
||||
func (s *Scheduler) GetRunningTaskCount() int {
|
||||
s.mu.RLock()
|
||||
|
||||
@@ -119,6 +119,7 @@ func Setup(c *Controllers) *gin.Engine {
|
||||
tasks.GET("/:id", c.Task.GetTask)
|
||||
tasks.PUT("/:id", c.Task.UpdateTask)
|
||||
tasks.DELETE("/:id", c.Task.DeleteTask)
|
||||
tasks.POST("/stop/:logID", c.Task.StopTask)
|
||||
}
|
||||
|
||||
// Task execution routes
|
||||
|
||||
@@ -480,6 +480,39 @@ func (es *ExecutorService) ExecuteTask(taskID int) *executor.ExecutionResult {
|
||||
}
|
||||
}
|
||||
|
||||
// StopTaskExecution stops a running task execution by LogID
|
||||
func (es *ExecutorService) StopTaskExecution(logID uint) error {
|
||||
var taskLog models.TaskLog
|
||||
if err := database.DB.First(&taskLog, logID).Error; err != nil {
|
||||
return fmt.Errorf("日志不存在")
|
||||
}
|
||||
|
||||
if taskLog.Status != "running" {
|
||||
return fmt.Errorf("任务已结束")
|
||||
}
|
||||
|
||||
task := es.taskService.GetTaskByID(int(taskLog.TaskID))
|
||||
if task == nil {
|
||||
return fmt.Errorf("任务不存在")
|
||||
}
|
||||
|
||||
// 远程任务:发送停止指令到 Agent
|
||||
if task.AgentID != nil && *task.AgentID > 0 {
|
||||
logger.Infof("[Executor] 请求停止远程任务 #%d (Agent #%d, LogID: %d)", task.ID, *task.AgentID, logID)
|
||||
return es.agentWSManager.SendToAgent(*task.AgentID, constant.WSTypeStop, map[string]interface{}{
|
||||
"log_id": logID,
|
||||
})
|
||||
}
|
||||
|
||||
// 本地任务:直接停止调度器中的执行实例
|
||||
logger.Infof("[Executor] 请求停止本地任务 #%d (LogID: %d)", task.ID, logID)
|
||||
if es.scheduler.StopLog(logID) {
|
||||
return nil
|
||||
}
|
||||
|
||||
return fmt.Errorf("任务当前不在运行队列中或已完成")
|
||||
}
|
||||
|
||||
// GetRunningCount 获取正在运行任务数量
|
||||
func (es *ExecutorService) GetRunningCount() int {
|
||||
return es.scheduler.GetRunningTaskCount()
|
||||
|
||||
Reference in New Issue
Block a user