diff --git a/internal/controllers/task_controller.go b/internal/controllers/task_controller.go index 3c16081..8690800 100644 --- a/internal/controllers/task_controller.go +++ b/internal/controllers/task_controller.go @@ -362,6 +362,7 @@ func (tc *TaskController) DeleteTask(c *gin.Context) { } tc.executorService.RemoveCronTask(id) + tc.executorService.GetScheduler().StopTask(id) success := tc.taskService.DeleteTask(id) if !success { @@ -460,6 +461,7 @@ func (tc *TaskController) BatchDeleteTasks(c *gin.Context) { // 移除 cron 调度 tc.executorService.RemoveCronTask(id) + tc.executorService.GetScheduler().StopTask(id) } // 执行批量删除 @@ -513,6 +515,7 @@ func (tc *TaskController) BatchDeleteByQuery(c *gin.Context) { } // 移除 cron 调度 tc.executorService.RemoveCronTask(task.ID) + tc.executorService.GetScheduler().StopTask(task.ID) } // 执行批量删除 diff --git a/internal/executor/executor.go b/internal/executor/executor.go index 2ec605e..1158c68 100644 --- a/internal/executor/executor.go +++ b/internal/executor/executor.go @@ -112,10 +112,14 @@ func ExecuteWithHooks(ctx context.Context, req Request, stdout, stderr io.Writer // 2. 执行命令 timeout := req.Timeout - if timeout <= 0 { - timeout = 30 + var execCtx context.Context + var cancel context.CancelFunc + + if timeout > 0 { + execCtx, cancel = context.WithTimeout(ctx, time.Duration(timeout)*time.Minute) + } else { + execCtx, cancel = context.WithCancel(ctx) } - execCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Minute) defer cancel() // 如果指定使用 mise,则预先构建好带 mise 的命令,这样 PreExecute 记录的就是完整命令 diff --git a/internal/services/tasks/executor_service.go b/internal/services/tasks/executor_service.go index a2c6c91..ea47359 100644 --- a/internal/services/tasks/executor_service.go +++ b/internal/services/tasks/executor_service.go @@ -490,7 +490,7 @@ func (es *ExecutorService) ExecuteDispatcher(ctx context.Context, req *executor. // 远程任务 if task.AgentID != nil && *task.AgentID != "" { // 将请求中已包含的环境变量(已合并)传递给 Agent - return es.ExecuteRemoteForScheduler(task, req.LogID, executor.FormatEnvVars(req.Envs), req.Secrets) + return es.ExecuteRemoteForScheduler(ctx, task, req.LogID, executor.FormatEnvVars(req.Envs), req.Secrets) } // 本地任务 @@ -1042,7 +1042,7 @@ func (es *ExecutorService) RemoveRunningGo(taskID string, goid int64) { } // ExecuteRemoteForScheduler 供 Scheduler 调用,执行远程任务并等待结果 -func (es *ExecutorService) ExecuteRemoteForScheduler(task *models.Task, logID string, envs string, secrets []string) (*executor.Result, error) { +func (es *ExecutorService) ExecuteRemoteForScheduler(ctx context.Context, task *models.Task, logID string, envs string, secrets []string) (*executor.Result, error) { agentID := *task.AgentID logger.Infof("[Executor] 远程执行任务 #%s: %s (Agent #%s, LogID: %s)", task.ID, task.Name, agentID, logID) @@ -1079,32 +1079,70 @@ func (es *ExecutorService) ExecuteRemoteForScheduler(task *models.Task, logID st // 4. 等待结果或超时 timeout := task.Timeout - if timeout <= 0 { - timeout = 30 - } start := time.Now() - select { - case agentResult := <-resultChan: - return &executor.Result{ - Output: agentResult.Output, - Error: agentResult.Error, - Status: agentResult.Status, - Duration: agentResult.Duration, - ExitCode: agentResult.ExitCode, - StartTime: time.Unix(agentResult.StartTime, 0), - EndTime: time.Unix(agentResult.EndTime, 0), - }, nil - case <-time.After(time.Duration(timeout) * time.Minute): - end := time.Now() - return &executor.Result{ - Status: constant.TaskStatusFailed, - Error: "等待 Agent 结果超时", - Duration: end.Sub(start).Milliseconds(), - ExitCode: -1, - StartTime: start, - EndTime: end, - }, fmt.Errorf("等待 Agent 结果超时") + + var timeoutChan <-chan time.Time + if timeout > 0 { + timeoutChan = time.After(time.Duration(timeout) * time.Minute) + } + + ticker := time.NewTicker(3 * time.Second) + defer ticker.Stop() + + for { + select { + case agentResult := <-resultChan: + return &executor.Result{ + Output: agentResult.Output, + Error: agentResult.Error, + Status: agentResult.Status, + Duration: agentResult.Duration, + ExitCode: agentResult.ExitCode, + StartTime: time.Unix(agentResult.StartTime, 0), + EndTime: time.Unix(agentResult.EndTime, 0), + }, nil + + case <-timeoutChan: + end := time.Now() + return &executor.Result{ + Status: constant.TaskStatusFailed, + Error: "等待 Agent 结果超时", + Duration: end.Sub(start).Milliseconds(), + ExitCode: -1, + StartTime: start, + EndTime: end, + }, fmt.Errorf("等待 Agent 结果超时") + + case <-ctx.Done(): + // 收到取消信号,向 Agent 发送停止指令 + _ = es.agentWSManager.SendToAgent(agentID, constant.WSTypeStop, map[string]interface{}{ + "log_id": logID, + }) + end := time.Now() + return &executor.Result{ + Status: constant.TaskStatusCancelled, + Error: "任务被手动停止或删除", + Duration: end.Sub(start).Milliseconds(), + ExitCode: -1, + StartTime: start, + EndTime: end, + }, fmt.Errorf("任务被主动停止") + + case <-ticker.C: + // 定期检查 Agent 是否在线 + if !es.agentWSManager.IsAgentOnline(agentID) { + end := time.Now() + return &executor.Result{ + Status: constant.TaskStatusFailed, + Error: "Agent 离线,任务被迫终止", + Duration: end.Sub(start).Milliseconds(), + ExitCode: -1, + StartTime: start, + EndTime: end, + }, fmt.Errorf("Agent 离线") + } + } } }