feat: refact id column define

This commit is contained in:
engigu
2026-03-03 17:25:38 +08:00
parent 117ae126f7
commit 4589c64923
62 changed files with 881 additions and 454 deletions
+54 -59
View File
@@ -24,9 +24,9 @@ import (
// AgentWSManager 接口定义(避免循环依赖)
type AgentWSManager interface {
RegisterRemoteWaiter(logID uint) chan *models.AgentTaskResult
UnregisterRemoteWaiter(logID uint)
SendToAgent(agentID uint, msgType string, data interface{}) error
RegisterRemoteWaiter(logID string) chan *models.AgentTaskResult
UnregisterRemoteWaiter(logID string)
SendToAgent(agentID string, msgType string, data interface{}) error
}
// SettingsService 接口定义(避免循环依赖)
@@ -115,10 +115,9 @@ func (h *ServerSchedulerHandler) OnTaskScheduled(req *executor.ExecutionRequest)
}
func (h *ServerSchedulerHandler) OnTaskExecuting(req *executor.ExecutionRequest) (io.Writer, io.Writer, error) {
var taskID uint
fmt.Sscanf(req.TaskID, "%d", &taskID)
taskID := req.TaskID
task := h.es.taskService.GetTaskByID(int(taskID))
task := h.es.taskService.GetTaskByID(taskID)
// 系统任务(无 taskID)不记录数据库日志,直接返回空写入器
if task == nil {
return nil, nil, nil
@@ -168,7 +167,7 @@ func (h *ServerSchedulerHandler) OnTaskExecuting(req *executor.ExecutionRequest)
}
func (h *ServerSchedulerHandler) OnTaskHeartbeat(req *executor.ExecutionRequest, duration int64) {
if req.LogID > 0 {
if req.LogID != "" {
h.es.taskLogService.UpdateTaskDuration(req.LogID, duration)
}
@@ -184,14 +183,13 @@ func (h *ServerSchedulerHandler) OnTaskStarted(req *executor.ExecutionRequest) {
}
func (h *ServerSchedulerHandler) OnTaskCompleted(req *executor.ExecutionRequest, result *executor.ExecutionResult) {
if req.LogID == 0 {
if req.LogID == "" {
return
}
var taskID uint
fmt.Sscanf(req.TaskID, "%d", &taskID)
taskID := req.TaskID
task := h.es.taskService.GetTaskByID(int(taskID))
task := h.es.taskService.GetTaskByID(taskID)
if task == nil {
return
}
@@ -204,7 +202,7 @@ func (h *ServerSchedulerHandler) OnTaskCompleted(req *executor.ExecutionRequest,
var err error
output, err = tl.CompressAndCleanup()
if err != nil {
logger.Errorf("[Executor] 压缩任务 #%d 日志失败: %v", task.ID, err)
logger.Errorf("[Executor] 压缩任务 #%s 日志失败: %v", task.ID, err)
output = "[System Error] 日志处理失败: " + err.Error()
}
} else {
@@ -230,7 +228,7 @@ func (h *ServerSchedulerHandler) OnTaskCompleted(req *executor.ExecutionRequest,
}
// 如果有 AgentID,也记录下来
if task.AgentID != nil && *task.AgentID > 0 {
if task.AgentID != nil && *task.AgentID != "" {
agentID := *task.AgentID
taskLog.AgentID = &agentID
}
@@ -251,12 +249,11 @@ func (h *ServerSchedulerHandler) OnTaskCompleted(req *executor.ExecutionRequest,
}
func (h *ServerSchedulerHandler) OnTaskFailed(req *executor.ExecutionRequest, err error) {
if req.LogID == 0 {
if req.LogID == "" {
return
}
var taskID uint
fmt.Sscanf(req.TaskID, "%d", &taskID)
taskID := req.TaskID
// 移除运行记录
if req.Metadata.GoID != 0 {
@@ -288,8 +285,8 @@ func (h *ServerSchedulerHandler) OnTaskFailed(req *executor.ExecutionRequest, er
}
// 补充 AgentID
task := h.es.taskService.GetTaskByID(int(taskID))
if task != nil && task.AgentID != nil && *task.AgentID > 0 {
task := h.es.taskService.GetTaskByID(taskID)
if task != nil && task.AgentID != nil && *task.AgentID != "" {
agentID := *task.AgentID
taskLog.AgentID = &agentID
}
@@ -321,10 +318,10 @@ func (es *ExecutorService) HandleTaskRetry(task *models.Task, req *executor.Exec
if retryIndex < task.RetryCount {
retryIndex++
logger.Infof("[Executor] 任务 #%d 执行失败/出错,将在 %d 秒后进行第 %d/%d 次重试...", task.ID, task.RetryInterval, retryIndex, task.RetryCount)
logger.Infof("[Executor] 任务 #%s 执行失败/出错,将在 %d 秒后进行第 %d/%d 次重试...", task.ID, task.RetryInterval, retryIndex, task.RetryCount)
es.scheduler.EnqueueDelayed(time.Duration(task.RetryInterval)*time.Second, func() *executor.ExecutionRequest {
latestTask := es.taskService.GetTaskByID(int(task.ID))
latestTask := es.taskService.GetTaskByID(task.ID)
if latestTask == nil || !latestTask.Enabled {
return nil
}
@@ -350,8 +347,7 @@ func (es *ExecutorService) HandleTaskRetry(task *models.Task, req *executor.Exec
}
func (h *ServerSchedulerHandler) OnCronNextRun(req *executor.ExecutionRequest, nextRun time.Time) {
var taskID uint
fmt.Sscanf(req.TaskID, "%d", &taskID)
taskID := req.TaskID
// 更新数据库中的下次运行时间
database.DB.Model(&models.Task{}).Where("id = ?", taskID).Update("next_run", nextRun)
}
@@ -359,19 +355,19 @@ func (h *ServerSchedulerHandler) OnCronNextRun(req *executor.ExecutionRequest, n
// LocalTaskHooks 本地任务钩子适配器
type LocalTaskHooks struct {
es *ExecutorService
logID uint
logID string
}
func (h *LocalTaskHooks) PreExecute(ctx context.Context, req executor.Request) (uint, error) {
func (h *LocalTaskHooks) PreExecute(ctx context.Context, req executor.Request) (string, error) {
return h.logID, nil
}
func (h *LocalTaskHooks) PostExecute(ctx context.Context, logID uint, result *executor.Result) error {
func (h *LocalTaskHooks) PostExecute(ctx context.Context, logID string, result *executor.Result) error {
return nil
}
func (h *LocalTaskHooks) OnHeartbeat(ctx context.Context, logID uint, duration int64) error {
if logID > 0 {
func (h *LocalTaskHooks) OnHeartbeat(ctx context.Context, logID string, duration int64) error {
if logID != "" {
return h.es.taskLogService.UpdateTaskDuration(logID, duration)
}
return nil
@@ -379,10 +375,9 @@ func (h *LocalTaskHooks) OnHeartbeat(ctx context.Context, logID uint, duration i
// ExecuteDispatcher 实现任务分发逻辑
func (es *ExecutorService) ExecuteDispatcher(ctx context.Context, req *executor.ExecutionRequest, stdout, stderr io.Writer) (*executor.Result, error) {
var taskID uint
fmt.Sscanf(req.TaskID, "%d", &taskID)
taskID := req.TaskID
task := es.taskService.GetTaskByID(int(taskID))
task := es.taskService.GetTaskByID(taskID)
// 系统任务(无 taskID)直接本地执行
if task == nil {
return executor.Execute(ctx, executor.Request{
@@ -409,7 +404,7 @@ func (es *ExecutorService) ExecuteDispatcher(ctx context.Context, req *executor.
}
// 远程任务
if task.AgentID != nil && *task.AgentID > 0 {
if task.AgentID != nil && *task.AgentID != "" {
return es.ExecuteRemoteForScheduler(task, req.LogID)
}
@@ -467,8 +462,8 @@ func (es *ExecutorService) AddCronTask(task *models.Task) error {
}
// RemoveCronTask 移除计划任务
func (es *ExecutorService) RemoveCronTask(taskID uint) {
es.cronManager.RemoveTask(fmt.Sprintf("%d", taskID))
func (es *ExecutorService) RemoveCronTask(taskID string) {
es.cronManager.RemoveTask(taskID)
}
// ValidateCron 验证 Cron 表达式
@@ -494,10 +489,10 @@ func (es *ExecutorService) loadCronTasks() {
go func(t models.Task) {
// 延迟一点时间再触发,确保系统完全启动
time.Sleep(3 * time.Second)
logger.Infof("[Executor] 触发开机服务启动任务 #%d: %s", t.ID, t.Name)
es.ExecuteTask(int(t.ID), nil)
logger.Infof("[Executor] 触发开机服务启动任务 #%s: %s", t.ID, t.Name)
es.ExecuteTask(t.ID, nil)
}(task)
} else if task.TriggerType == constant.TriggerTypeCron && task.Schedule != "" && (task.AgentID == nil || *task.AgentID == 0) {
} else if task.TriggerType == constant.TriggerTypeCron && task.Schedule != "" && (task.AgentID == nil || *task.AgentID == "") {
// 只调度本地任务(agent_id 为空或 0)的定时任务
err := es.cronManager.AddTask(&task)
if err != nil {
@@ -519,11 +514,11 @@ func (es *ExecutorService) Reload() {
}
// ExecuteTask executes a task by ID(同步执行,供 API 调用)
func (es *ExecutorService) ExecuteTask(taskID int, extraEnvs []string) *executor.ExecutionResult {
func (es *ExecutorService) ExecuteTask(taskID string, extraEnvs []string) *executor.ExecutionResult {
task := es.taskService.GetTaskByID(taskID)
if task == nil {
return &executor.ExecutionResult{
TaskID: fmt.Sprintf("%d", taskID),
TaskID: taskID,
Success: false,
Error: "任务不存在",
StartTime: time.Now(),
@@ -532,9 +527,9 @@ func (es *ExecutorService) ExecuteTask(taskID int, extraEnvs []string) *executor
}
// 1. 检查并发
if err := es.CheckConcurrency(uint(taskID)); err != nil {
if err := es.CheckConcurrency(taskID); err != nil {
return &executor.ExecutionResult{
TaskID: fmt.Sprintf("%d", taskID),
TaskID: taskID,
Success: false,
Error: err.Error(), // 这里会返回 "任务正在运行中,拒绝并行执行"
StartTime: time.Now(),
@@ -548,7 +543,7 @@ func (es *ExecutorService) ExecuteTask(taskID int, extraEnvs []string) *executor
}
req := &executor.ExecutionRequest{
TaskID: fmt.Sprintf("%d", task.ID),
TaskID: task.ID,
Name: task.Name,
Command: task.Command,
WorkDir: task.WorkDir,
@@ -562,7 +557,7 @@ func (es *ExecutorService) ExecuteTask(taskID int, extraEnvs []string) *executor
es.scheduler.EnqueueOrExecute(req)
return &executor.ExecutionResult{
TaskID: fmt.Sprintf("%d", task.ID),
TaskID: task.ID,
Success: true,
Status: constant.TaskStatusQueued,
StartTime: time.Now(),
@@ -570,7 +565,7 @@ func (es *ExecutorService) ExecuteTask(taskID int, extraEnvs []string) *executor
}
// StopTaskExecution stops a running task execution by LogID
func (es *ExecutorService) StopTaskExecution(logID uint) error {
func (es *ExecutorService) StopTaskExecution(logID string) error {
var taskLog models.TaskLog
if err := database.DB.First(&taskLog, logID).Error; err != nil {
return fmt.Errorf("日志不存在")
@@ -580,21 +575,21 @@ func (es *ExecutorService) StopTaskExecution(logID uint) error {
return fmt.Errorf("任务已结束")
}
task := es.taskService.GetTaskByID(int(taskLog.TaskID))
task := es.taskService.GetTaskByID(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)
if task.AgentID != nil && *task.AgentID != "" {
logger.Infof("[Executor] 请求停止远程任务 #%s (Agent #%s, LogID: %s)", 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)
logger.Infof("[Executor] 请求停止本地任务 #%s (LogID: %s)", task.ID, logID)
if es.scheduler.StopLog(logID) {
return nil
}
@@ -658,7 +653,7 @@ func (es *ExecutorService) UpdateResult(res executor.ExecutionResult) {
// 查找是否已存在(通过 LogID)
for i := range es.results {
if es.results[i].LogID == res.LogID && res.LogID != 0 {
if es.results[i].LogID == res.LogID && res.LogID != "" {
es.results[i] = res
return
}
@@ -703,9 +698,9 @@ func (es *ExecutorService) CleanupRunningTasks() error {
}
// CheckConcurrency 检查任务并发限制(只读检查)
func (es *ExecutorService) CheckConcurrency(taskID uint) error {
func (es *ExecutorService) CheckConcurrency(taskID string) error {
var task models.Task
if err := database.DB.Select("config, running_go").First(&task, taskID).Error; err != nil {
if err := database.DB.Select("config, running_go").Where("id = ?", taskID).First(&task).Error; err != nil {
return err
}
var goids []int64
@@ -725,11 +720,11 @@ func (es *ExecutorService) CheckConcurrency(taskID uint) error {
}
// AddRunningGo 添加当前 goroutine ID 到任务的 running_go 字段
func (es *ExecutorService) AddRunningGo(taskID uint) (int64, error) {
func (es *ExecutorService) AddRunningGo(taskID string) (int64, error) {
goid := utils.GetGoroutineID()
err := database.DB.Transaction(func(tx *gorm.DB) error {
var task models.Task
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&task, taskID).Error; err != nil {
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", taskID).First(&task).Error; err != nil {
return err
}
var goids []int64
@@ -756,10 +751,10 @@ func (es *ExecutorService) AddRunningGo(taskID uint) (int64, error) {
}
// RemoveRunningGo 从任务的 running_go 字段移除指定 goroutine ID
func (es *ExecutorService) RemoveRunningGo(taskID uint, goid int64) {
func (es *ExecutorService) RemoveRunningGo(taskID string, goid int64) {
database.DB.Transaction(func(tx *gorm.DB) error {
var task models.Task
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&task, taskID).Error; err != nil {
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("id = ?", taskID).First(&task).Error; err != nil {
return err
}
var goids []int64
@@ -778,17 +773,17 @@ func (es *ExecutorService) RemoveRunningGo(taskID uint, goid int64) {
}
// ExecuteRemoteForScheduler 供 Scheduler 调用,执行远程任务并等待结果
func (es *ExecutorService) ExecuteRemoteForScheduler(task *models.Task, logID uint) (*executor.Result, error) {
func (es *ExecutorService) ExecuteRemoteForScheduler(task *models.Task, logID string) (*executor.Result, error) {
agentID := *task.AgentID
logger.Infof("[Executor] 远程执行任务 #%d: %s (Agent #%d, LogID: %d)", task.ID, task.Name, agentID, logID)
logger.Infof("[Executor] 远程执行任务 #%s: %s (Agent #%s, LogID: %s)", task.ID, task.Name, agentID, logID)
// 1. 检查 Agent 状态
var agent models.Agent
if err := database.DB.First(&agent, agentID).Error; err != nil {
return nil, fmt.Errorf("Agent #%d 不存在", agentID)
if err := database.DB.Where("id = ?", agentID).First(&agent).Error; err != nil {
return nil, fmt.Errorf("Agent #%s 不存在", agentID)
}
if !agent.Enabled {
return nil, fmt.Errorf("Agent #%d 已禁用", agentID)
return nil, fmt.Errorf("Agent #%s 已禁用", agentID)
}
if es.agentWSManager == nil {
return nil, fmt.Errorf("AgentWSManager 未初始化")