fix: resolve task timeout bug and refactor remote executor wait logic
- Prevent hardcoded 30-minute timeout when timeout is set to 0. - Fix zombie processes by actively stopping tasks during deletion. - Refactor remote execution wait logic to properly handle ctx.Done() and agent disconnections.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
// 执行批量删除
|
||||
|
||||
@@ -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 记录的就是完整命令
|
||||
|
||||
@@ -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 离线")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user