From aef124defcb09b1d95510b4e4481f30e8ef98ea8 Mon Sep 17 00:00:00 2001 From: engigu Date: Tue, 10 Feb 2026 16:47:18 +0800 Subject: [PATCH] feat: add task cancell --- agent/agent.go | 20 +++ internal/constant/constant.go | 1 + internal/controllers/task_controller.go | 16 ++ internal/executor/scheduler.go | 26 ++- internal/router/router.go | 1 + internal/services/tasks/executor_service.go | 33 ++++ web/src/api/index.ts | 3 +- web/src/views/agents/Agents.vue | 174 +++++++++++++------- web/src/views/history/History.vue | 105 +++++++++--- web/src/views/tasks/Tasks.vue | 67 +++++--- 10 files changed, 335 insertions(+), 111 deletions(-) diff --git a/agent/agent.go b/agent/agent.go index c7e7388..3cc4354 100644 --- a/agent/agent.go +++ b/agent/agent.go @@ -33,6 +33,7 @@ const ( WSTypeTaskLog = "task_log" WSTypeExecute = "execute" WSTypeTaskHeartbeat = "task_heartbeat" + WSTypeStop = "stop" ) type WSMessage struct { @@ -370,6 +371,8 @@ func (a *Agent) handleWSMessage(msg *WSMessage) { a.fetchTasks() case WSTypeExecute: a.handleExecute(msg.Data) + case WSTypeStop: + a.handleStop(msg.Data) } } @@ -505,6 +508,23 @@ func (a *Agent) handleExecute(data json.RawMessage) { a.scheduler.EnqueueOrExecute(execReq) } +func (a *Agent) handleStop(data json.RawMessage) { + var req struct { + LogID uint `json:"log_id"` + } + if err := json.Unmarshal(data, &req); err != nil { + logger.Errorf("解析停止请求失败: %v", err) + return + } + + logger.Infof("[Agent] 收到停止指令 LogID: %d", req.LogID) + if a.scheduler.StopLog(req.LogID) { + logger.Infof("[Agent] 任务执行 #%d 已成功停止", req.LogID) + } else { + logger.Warnf("[Agent] 任务执行 #%d 停止失败(可能已完成或不在运行队列中)", req.LogID) + } +} + // RealTimeLogWriter 实时日志写入器,通过 WebSocket 发送日志 type RealTimeLogWriter struct { agent *Agent diff --git a/internal/constant/constant.go b/internal/constant/constant.go index b82e3dd..b262303 100644 --- a/internal/constant/constant.go +++ b/internal/constant/constant.go @@ -63,6 +63,7 @@ const ( WSTypeEnabled = "enabled" WSTypeFetchTasks = "fetch_tasks" WSTypeTaskHeartbeat = "task_heartbeat" + WSTypeStop = "stop" ) // TablePrefix 表前缀,从配置文件读取 diff --git a/internal/controllers/task_controller.go b/internal/controllers/task_controller.go index d11be16..0d6d3e9 100644 --- a/internal/controllers/task_controller.go +++ b/internal/controllers/task_controller.go @@ -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, "停止请求已发送") +} diff --git a/internal/executor/scheduler.go b/internal/executor/scheduler.go index de56ca6..ddd6fea 100644 --- a/internal/executor/scheduler.go +++ b/internal/executor/scheduler.go @@ -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() diff --git a/internal/router/router.go b/internal/router/router.go index fc82949..495af18 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -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 diff --git a/internal/services/tasks/executor_service.go b/internal/services/tasks/executor_service.go index 0e15332..28be30d 100644 --- a/internal/services/tasks/executor_service.go +++ b/internal/services/tasks/executor_service.go @@ -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() diff --git a/web/src/api/index.ts b/web/src/api/index.ts index d78ac82..032e061 100644 --- a/web/src/api/index.ts +++ b/web/src/api/index.ts @@ -69,7 +69,8 @@ export const api = { create: (data: Partial) => request('/tasks', { method: 'POST', body: JSON.stringify(data) }), update: (id: number, data: Partial) => request(`/tasks/${id}`, { method: 'PUT', body: JSON.stringify(data) }), delete: (id: number) => request(`/tasks/${id}`, { method: 'DELETE' }), - execute: (id: number) => request(`/execute/task/${id}`, { method: 'POST' }) + execute: (id: number) => request(`/execute/task/${id}`, { method: 'POST' }), + stop: (logID: number) => request(`/tasks/stop/${logID}`, { method: 'POST' }) }, scripts: { list: () => request('/scripts'), diff --git a/web/src/views/agents/Agents.vue b/web/src/views/agents/Agents.vue index b00b2be..ecd97ba 100644 --- a/web/src/views/agents/Agents.vue +++ b/web/src/views/agents/Agents.vue @@ -6,7 +6,7 @@ import { Label } from '@/components/ui/label' import { Dialog, DialogContent, DialogHeader, DialogTitle, DialogFooter, DialogDescription } from '@/components/ui/dialog' import { AlertDialog, AlertDialogAction, AlertDialogCancel, AlertDialogContent, AlertDialogDescription, AlertDialogFooter, AlertDialogHeader, AlertDialogTitle } from '@/components/ui/alert-dialog' import { Tabs, TabsContent, TabsList, TabsTrigger } from '@/components/ui/tabs' -import { RefreshCw, Trash2, Edit, Copy, Server, Search, Download, RotateCw, Plus, Ticket, Power, PowerOff, ListTodo, Eye } from 'lucide-vue-next' +import { RefreshCw, Trash2, Edit, Copy, Server, Search, Download, RotateCw, Plus, Ticket, ListTodo, Eye, WifiOff, Zap, Check, X } from 'lucide-vue-next' import { api, type Agent, type AgentToken } from '@/api' import { toast } from 'vue-sonner' import { useRouter } from 'vue-router' @@ -35,8 +35,8 @@ let refreshTimer: ReturnType | null = null const filteredAgents = computed(() => { if (!searchQuery.value) return agents.value const q = searchQuery.value.toLowerCase() - return agents.value.filter(a => - a.name.toLowerCase().includes(q) || + return agents.value.filter(a => + a.name.toLowerCase().includes(q) || a.hostname?.toLowerCase().includes(q) || a.ip?.toLowerCase().includes(q) ) @@ -226,7 +226,8 @@ onUnmounted(() => {
-