From 786fe71778b70a204b4edabbb4c035cdec643215 Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 11 Jun 2026 09:20:28 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96=E5=AE=9A=E6=97=B6=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .agent/skills/database-migration/SKILL.md | 2 + ...0002_alter_schedules_id_auto_increment.sql | 9 +++ ...0002_alter_schedules_id_auto_increment.sql | 32 +++++++++++ internal/model/schedule.go | 2 - internal/task/executor.go | 55 ++++++++++++++----- 5 files changed, 84 insertions(+), 16 deletions(-) create mode 100644 internal/db/migrator/goose/postgres/202606110002_alter_schedules_id_auto_increment.sql create mode 100644 internal/db/migrator/goose/sqlite/202606110002_alter_schedules_id_auto_increment.sql diff --git a/.agent/skills/database-migration/SKILL.md b/.agent/skills/database-migration/SKILL.md index 91e08797..d0a9fc0b 100644 --- a/.agent/skills/database-migration/SKILL.md +++ b/.agent/skills/database-migration/SKILL.md @@ -24,6 +24,8 @@ Wavelet 使用 `github.com/pressly/goose/v3` 执行 SQL 迁移。迁移入口是 ``` - 不要把表结构、默认系统配置、默认模板、默认管理员初始化写回 Go 代码。 +- 编辑表结构(DDL)和插入表数据(DML/Seed)不要放在同一个 SQL 文件里,必须分成两个独立的 SQL 文件完成(例如,先通过一个文件修改表结构,再通过下一个递增版本号的文件插入/初始化数据)。 +- 插入定时任务(schedules 表数据)时绝对不能指定 `id`,必须依靠数据库自增(Identity 或 AUTOINCREMENT)自动分配,防止与用户手动或后续插入的定时任务产生 ID 冲突。 - 不要添加物理外键;关系字段使用显式索引。 - 数据库默认值应匹配 Go model 零值或业务兜底值。 - 系统配置仍然保存字符串值;布尔值写 `"true"` / `"false"`,数字写十进制字符串,复杂结构写合法 JSON 字符串。 diff --git a/internal/db/migrator/goose/postgres/202606110002_alter_schedules_id_auto_increment.sql b/internal/db/migrator/goose/postgres/202606110002_alter_schedules_id_auto_increment.sql new file mode 100644 index 00000000..0797bd64 --- /dev/null +++ b/internal/db/migrator/goose/postgres/202606110002_alter_schedules_id_auto_increment.sql @@ -0,0 +1,9 @@ +-- +goose Up +-- +goose StatementBegin +ALTER TABLE schedules ALTER COLUMN id ADD GENERATED BY DEFAULT AS IDENTITY (START WITH 100); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +ALTER TABLE schedules ALTER COLUMN id DROP IDENTITY IF EXISTS; +-- +goose StatementEnd diff --git a/internal/db/migrator/goose/sqlite/202606110002_alter_schedules_id_auto_increment.sql b/internal/db/migrator/goose/sqlite/202606110002_alter_schedules_id_auto_increment.sql new file mode 100644 index 00000000..4f338526 --- /dev/null +++ b/internal/db/migrator/goose/sqlite/202606110002_alter_schedules_id_auto_increment.sql @@ -0,0 +1,32 @@ +-- +goose Up +-- +goose StatementBegin +-- 1. Rename existing schedules table +ALTER TABLE schedules RENAME TO schedules_old; + +-- 2. Create new schedules table with AUTOINCREMENT +CREATE TABLE schedules ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name VARCHAR(128) NOT NULL, + task_type VARCHAR(64) NOT NULL, + cron VARCHAR(64) NOT NULL, + payload TEXT, + is_active BOOLEAN NOT NULL DEFAULT TRUE, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); + +-- 3. Copy existing data +INSERT INTO schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) +SELECT id, name, task_type, cron, payload, is_active, created_at, updated_at FROM schedules_old; + +-- 4. Drop the old table +DROP TABLE schedules_old; + +-- 5. Recreate index +CREATE INDEX IF NOT EXISTS idx_schedules_is_active ON schedules (is_active); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +-- Re-creating table with AUTOINCREMENT cannot be undone simply without recreating table again. +-- +goose StatementEnd diff --git a/internal/model/schedule.go b/internal/model/schedule.go index f8848cb6..1fd8b7db 100644 --- a/internal/model/schedule.go +++ b/internal/model/schedule.go @@ -8,7 +8,6 @@ import ( "time" "github.com/Rain-kl/Wavelet/internal/db" - "github.com/Rain-kl/Wavelet/internal/db/idgen" ) // Schedule 定时任务配置表 @@ -30,7 +29,6 @@ func (Schedule) TableName() string { // CreateSchedule 创建定时任务 func CreateSchedule(ctx context.Context, schedule *Schedule) error { - schedule.ID = idgen.NextUint64ID() return db.DB(ctx).Create(schedule).Error } diff --git a/internal/task/executor.go b/internal/task/executor.go index 87c0fc66..55f466ab 100644 --- a/internal/task/executor.go +++ b/internal/task/executor.go @@ -128,6 +128,8 @@ func DispatchTask(ctx context.Context, taskType string, payload []byte, triggere return "", fmt.Errorf(errTaskEnqueueFailed, err) } + _ = model.AppendTaskExecutionLog(ctx, taskID, fmt.Sprintf("[系统] 任务已成功入队,等待调度执行 (队列: %s, 最大重试次数: %d)", meta.Queue, meta.MaxRetry)) + return taskID, nil } @@ -187,6 +189,8 @@ func RetryTask(ctx context.Context, id uint64) (string, error) { return "", fmt.Errorf(errRetryTaskEnqueueFailed, err) } + _ = model.AppendTaskExecutionLog(ctx, newTaskID, fmt.Sprintf("[系统] 手动触发重试,已重新创建任务并入队 (原任务ID: %s, 重试次数: %d/%d)", execution.TaskID, execution.RetryCount+1, execution.MaxRetry)) + return newTaskID, nil } @@ -229,6 +233,13 @@ func ProcessTask(ctx context.Context, t *asynq.Task) error { } } + if execution != nil { + AppendLog(ctx, "[系统] 开始执行异步任务 [名称: %s, 类型: %s],重试次数: %d/%d", + execution.TaskName, t.Type(), execution.RetryCount, execution.MaxRetry) + } else { + AppendLog(ctx, "[系统] 开始执行异步任务 [类型: %s]", t.Type()) + } + // 开始计时 start := time.Now() @@ -292,21 +303,9 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution, execution.FinishedAt = &finishTime if execErr != nil { - execution.Status = model.TaskExecutionStatusFailed - execution.ErrorMessage = execErr.Error() - logger.ErrorF(ctx, "[TaskExecutor] 任务处理失败 Type: %s TaskID: %s Duration: %d ms Error: %v", t.Type(), execution.TaskID, duration.Milliseconds(), execErr) - span.SetStatus(codes.Error, execErr.Error()) - span.RecordError(execErr) + handleFailedTask(ctx, execution, t, duration, execErr, span) } else { - execution.Status = model.TaskExecutionStatusSucceeded - execution.ErrorMessage = "" // 清除历史重试失败遗留的错误信息 - if result != nil { - execution.Result = result.Message - if result.Detail != "" { - execution.Result = fmt.Sprintf("%s\n%s", result.Message, result.Detail) - } - } - logger.InfoF(ctx, "[TaskExecutor] 任务处理完成 Type: %s TaskID: %s Duration: %d ms", t.Type(), execution.TaskID, duration.Milliseconds()) + handleSuccessfulTask(ctx, execution, t, duration, result) } if err := model.UpdateTaskExecution(ctx, execution); err != nil { @@ -314,6 +313,34 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution, } } +func handleFailedTask(ctx context.Context, execution *model.TaskExecution, t *asynq.Task, duration time.Duration, execErr error, span trace.Span) { + execution.Status = model.TaskExecutionStatusFailed + execution.ErrorMessage = execErr.Error() + logger.ErrorF(ctx, "[TaskExecutor] 任务处理失败 Type: %s TaskID: %s Duration: %d ms Error: %v", t.Type(), execution.TaskID, duration.Milliseconds(), execErr) + span.SetStatus(codes.Error, execErr.Error()) + span.RecordError(execErr) + + AppendLog(ctx, "[系统] 任务执行失败,耗时: %d ms,错误原因: %v", duration.Milliseconds(), execErr) +} + +func handleSuccessfulTask(ctx context.Context, execution *model.TaskExecution, t *asynq.Task, duration time.Duration, result *TaskResult) { + execution.Status = model.TaskExecutionStatusSucceeded + execution.ErrorMessage = "" // 清除历史重试失败遗留的错误信息 + if result != nil { + execution.Result = result.Message + if result.Detail != "" { + execution.Result = fmt.Sprintf("%s\n%s", result.Message, result.Detail) + } + } + logger.InfoF(ctx, "[TaskExecutor] 任务处理完成 Type: %s TaskID: %s Duration: %d ms", t.Type(), execution.TaskID, duration.Milliseconds()) + + resultMsg := "成功" + if result != nil { + resultMsg = result.Message + } + AppendLog(ctx, "[系统] 任务执行成功,耗时: %d ms,执行结果: %s", duration.Milliseconds(), resultMsg) +} + // generateTaskID 生成任务 ID func generateTaskID(taskType string, triggeredBy string) string { uniqueID := idgen.NextUint64ID()