From 7da4b72d24751770d35cf3afbb6cadbfb898ac87 Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 11 Jun 2026 16:54:45 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E6=97=A5=E5=BF=97=E4=BC=98?= =?UTF-8?q?=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .agent/skills/new-async-task/SKILL.md | 182 +++++++++++--------------- internal/model/task_execution.go | 116 +++++++++++++++- internal/model/task_execution_test.go | 103 +++++++++++++-- internal/task/executor.go | 26 +++- internal/task/executor_test.go | 38 ++++++ 5 files changed, 340 insertions(+), 125 deletions(-) diff --git a/.agent/skills/new-async-task/SKILL.md b/.agent/skills/new-async-task/SKILL.md index 57b58c3c..567989fb 100644 --- a/.agent/skills/new-async-task/SKILL.md +++ b/.agent/skills/new-async-task/SKILL.md @@ -1,139 +1,103 @@ --- name: "new-async-task" -description: "项目专用:当新增或修改 Asynq 异步任务、后台任务、定时任务、任务元数据、TaskHandler、TaskParam、PayloadValidator、AppendLog、任务重试、任务执行记录或 Admin 任务 API 时必须使用。本技能按项目约束指导常量定义、处理器实现、统一注册、Worker 路由、Cron 配置、Swagger、测试和 code-check 验证。" +description: "Wavelet 项目专用:新增或修改 Asynq 异步任务、后台任务、定时任务、任务元数据、TaskHandler、TaskParam、PayloadValidator、AppendLog、任务重试、任务执行记录或 Admin 任务 API 时必须使用。" --- # 异步任务开发 -本技能只覆盖 Asynq 任务工作流。开始前先读仓库根目录 `AGENTS.md`,并遵守其中的项目级规则:HTTP 路由只在 `internal/router/router.go` 注册、API 变更后运行 `make swagger`、提交前运行 `make code-check`、不要删除 `frontend/node_modules`、`internal/util/` 不引入框架依赖。 +开始前阅读根目录 `AGENTS.md`。只修改任务相关链路,遵守项目路由、日志、数据库迁移和质量门禁要求。 -## 先定位真实链路 +## 开始前 -新增或修改任务前,先快速查看这些文件,确认当前实现没有漂移: +按任务范围检查当前实现: -- `internal/task/handler.go`: `TaskHandler`、`TaskResult`、可选 `PayloadValidator`。 -- `internal/task/constants.go`: 框架通用常量,如 `QueueDefault` 和 `DefaultMaxRetry`。 -- `internal/task/meta.go`: 框架任务元数据结构体(TaskParam、TaskMeta)及全局动态注册与查询接口。 -- `internal/task/executor.go`: `RegisterHandler`、`ValidateAndNormalizePayload`、`DispatchTask`、`RetryTask`、`ProcessTask`、`AppendLog`。 -- `internal/task/handlers/register.go`: 内置 handler 和元数据的统一注册点,Admin API 和 Worker 都依赖它。 -- `internal/task/worker/worker.go`: Asynq mux 动态路由分发和队列配置。 -- `internal/task/scheduler/scheduler.go`: Cron 调度。 -- `internal/apps/admin/task/routers.go`: Admin 下发、查询、详情、重试 API。 -- 现有参考:`internal/apps/upload/tasks.go`(无参数任务)、`internal/apps/user/tasks.go`(带参数任务)。 +- `internal/task/handler.go`:`TaskHandler`、`TaskResult`、`PayloadValidator` +- `internal/task/meta.go`:`TaskMeta`、`TaskParam` +- `internal/task/executor.go`:下发、执行、日志、重试 +- `internal/task/handlers/register.go`:Handler 和元数据注册 +- `internal/task/worker/worker.go`:Worker 路由和队列 +- `internal/task/scheduler/scheduler.go`:定时调度 +- `internal/apps/admin/task/routers.go`:Admin 任务 API +- `internal/model/task_execution.go`:执行记录和日志持久化 -当前任务执行链路: +需要模板时阅读 [references/CODE-EXAMPLES.md](references/CODE-EXAMPLES.md)。 -```text -Admin dispatch -> ValidateAndNormalizePayload -> DispatchTask - -> Asynq Redis queue -> worker mux -> ProcessTask - -> registered TaskHandler.Execute -> TaskExecution status/log/result -``` +## 实现要求 -## 修改检查清单 +### 任务定义 -按任务影响面选择对应步骤。不要只改其中一条链路。 +- 在 `internal/apps//tasks.go` 定义任务类型、Admin 任务类型和 `TaskMeta`。 +- Asynq 任务类型使用 `:` 格式。 +- 完整设置 `Type`、`AsynqTask`、`Name`、`Description`、`MaxRetry`、`Queue`、`Retryable`。 +- 有参数任务必须定义 payload struct。 +- `TaskParam.Name` 必须与 payload JSON tag 一致。 +- `TaskParam` 只描述前端表单,不代替服务端校验。 -> 需要可复制的代码模板时,阅读 [references/CODE-EXAMPLES.md](references/CODE-EXAMPLES.md)。那里包含任务常量、无参数 handler、带参数 `PayloadValidator`、统一注册、Worker 路由、Cron 配置和测试示例。 +### Handler -1. 定义任务元数据与常量。 - - 业务包常量与元数据:在对应业务模块的 `internal/apps//tasks.go` 中定义 Asynq 任务类型常量(如 `CleanupUnusedUploadsTask = "upload:cleanup_unused"`)和 Admin 任务类型常量(如 `TaskTypeCleanupUploads = "cleanup_unused_uploads"`)。 - - 在同一 `tasks.go` 文件中定义该任务的 `TaskMeta` 元数据变量(如 `CleanupUnusedUploadsMeta = task.TaskMeta{...}`),配置 `Type`、`AsynqTask`、`Name`、`Description`、`MaxRetry`、`Queue`、`Retryable` 等字段。 - - 有参数任务在 `Params` 中描述前端表单字段。`TaskParam.Name` 必须与 payload JSON tag 对齐。 +- Handler 必须实现 `task.TaskHandler`。 +- 有参数任务必须实现 `task.PayloadValidator`,负责校验和标准化 Admin 下发参数。 +- `Execute` 必须再次解析 payload;不要假设入口一定经过 Admin 校验。 +- 成功返回 `&task.TaskResult{Message: ..., Detail: ...}`。 +- 失败返回 error,由任务框架处理状态和重试。 +- 不要吞掉关键错误。 +- 复杂 SQL 放到 `internal/model/` 或 `internal/service/`。 +- 新增 Go 文件后检查许可证头,必要时运行 `make license`。 -2. 实现 handler。 - - 优先放在对应业务模块的 `internal/apps//tasks.go`。 - - handler 必须实现 `task.TaskHandler`。 - - 带参数任务定义 payload struct,并实现 `task.PayloadValidator` 做服务端校验和标准化。 - - `Execute` 中仍要解析 payload,因为 Scheduler、重试或其他入口不一定经过 Admin 校验。 - - 不要在 handler 中写复杂 SQL;复杂查询放到 `internal/model/` 或 `internal/service/`。 - - 新增 Go 文件后检查 license header;必要时运行 `make license`。 +### 注册 -3. 统一注册 handler 与元数据。 - - 在 `internal/task/handlers/register.go` 导入业务模块,调用 `task.RegisterHandler(asynqTaskType, handler)` 注册处理器。 - - 同时,在该文件中调用 `task.RegisterTaskMeta(meta)` 注册刚才在业务模块中定义的任务元数据。 - - 这里是 Admin 校验、元数据获取和 Worker 执行共同依赖的注册点。 +- 在 `internal/task/handlers/register.go` 同时注册 Handler 和 `TaskMeta`。 +- 不要在其他位置单独注册任务。 -4. 如需 Cron 调度,系统默认定时任务必须通过 SQL 迁移(goose)初始化。 - - 确保任务已正确注册并载入全局元数据池中。 - - 在 `internal/db/migrator/goose/postgres` 和 `sqlite` 下编写 migration 脚本,使用 `INSERT INTO schedules` 语句初始化任务,指定 `task_type` 和 `cron` 等字段。必须妥善处理冲突(如 `ON CONFLICT DO NOTHING`)以支持幂等。 +## 日志要求 -5. 如改动 Admin API。 - - handler 放在 `internal/apps/admin//` 或现有 Admin task 模块内。 - - 路由只在 `internal/router/router.go` 注册。 - - 响应保持 `{ "error_msg": "", "data": ... }`,分页保持 `{ "total": 0, "results": [] }`。 - - 补完整 Swagger 注释并运行 `make swagger`。 +- 在 `TaskHandler.Execute` 中使用 `task.AppendLog(ctx, format, args...)`。 +- 记录任务开始、参数摘要、批次进度、关键状态、可继续错误和完成摘要。 +- 批量处理按批次记录;禁止为大循环中的每条数据写日志。 +- 不要直接修改任务日志的 Redis key 或 `w_task_executions.log`。 -## Handler 模式 +日志框架约束: -约定: +- 执行状态实时写入数据库:`pending`、`running`、`succeeded`、`failed`。 +- 实时日志写入 Redis,每个任务最多保留最近 1000 行。 +- Redis 日志 TTL 为 24 小时,每次追加时刷新。 +- 查询时优先返回 Redis 日志,Redis 不存在时读取数据库。 +- 任务成功或自动重试耗尽后,将日志写入数据库并删除 Redis 缓冲。 +- 自动重试期间保留同一 taskID 的 Redis 日志。 -- `ValidatePayload` 是 Admin 下发时的服务端校验入口,返回值会作为标准化 payload 存库和入队。 -- `TaskParam` 只是前端表单元数据,不代替服务端校验。 -- 成功返回 `&task.TaskResult{Message: "...", Detail: "..."}`;失败返回 `nil, fmt.Errorf("...")`,由 `ProcessTask` 标记失败并交给 Asynq 重试。 -- 错误要返回给框架,不在 handler 内吞掉;可继续的单条失败可用 `AppendLog` 记录后继续处理。 +## 重试要求 -> 在新增无参数任务、带参数任务或 `PayloadValidator` 时,阅读 [references/CODE-EXAMPLES.md](references/CODE-EXAMPLES.md) 的 Handler 示例。 +- Handler 返回 error 以触发 Asynq 自动重试。 +- 不要在 Handler 内自行实现重复重试循环。 +- Admin 手动重试只允许: + - 原任务状态为 `failed` + - `Retryable=true` + - `RetryCount < MaxRetry` +- 修改重试行为时同时检查: + - `internal/task/executor.go` + - `internal/model/task_execution.go` + - `internal/apps/admin/task/routers.go` + - 前端任务执行列表 -## AppendLog 规则 +## 定时任务 -在 `TaskHandler.Execute` 内用 `task.AppendLog(ctx, format, args...)` 写任务日志。 +- 默认定时任务必须通过 Goose SQL 迁移写入 `schedules`。 +- PostgreSQL 和 SQLite 迁移必须同时提供。 +- 初始化 SQL 必须幂等。 +- 涉及迁移时使用 `database-migration` skill。 -- 任务开始、参数摘要、批次进度、关键状态、可继续错误、完成摘要适合记录。 -- 避免对大循环中的每条记录都写日志;每次 `AppendLog` 都可能触发一次数据库更新。 -- 如果上下文中没有 taskID,`AppendLog` 会降级为普通应用日志,不应额外兜底报错。 +## Admin API -## 重试规则 +- Handler 放在现有 Admin task 模块或 `internal/apps/admin//`。 +- 路由只在 `internal/router/router.go` 注册。 +- 响应保持 `{ "error_msg": "", "data": ... }`。 +- 分页数据保持 `{ "total": 0, "results": [] }`。 +- Swagger 注释必须完整;API 变化后运行 `make swagger`。 -该项目有两层重试: +## 前端 -- Asynq 自动重试:`ProcessTask` 返回 error 后按入队的 `MaxRetry` 处理。 -- Admin 手动重试:`POST /api/v1/admin/tasks/executions/:id/retry` 创建新的 `TaskExecution`,要求原任务状态为 failed、`Retryable=true`、`RetryCount < MaxRetry`。 - -修改重试语义时同时检查 `internal/task/executor.go`、`internal/model/task_execution.go`、`internal/apps/admin/task/routers.go` 和前端任务执行列表。 - -## Frontend/Admin 任务 UI - -只有任务元数据变化时,通常不需要写新页面;现有 Admin UI 会根据 `DispatchableTasks` 和 `Params` 动态渲染。 - -若确实要改前端: - -- 读 shadcn skill。 -- 业务组件优先放在 `frontend/components/common/admin/`。 -- API 访问走 `frontend/lib/services/` 的 service class 和 `services` export。 +- 仅任务元数据变化时,优先复用现有动态任务表单,不新增页面。 +- API 调用必须通过 `frontend/lib/services/`。 +- 修改 shadcn/ui 时使用 `shadcn` skill。 - 不使用 `any`。 -- 页面根容器保持 `w-full`,不要加页面级 `max-w-*`。 - -## 验证 - -根据改动范围运行最小有效验证,最后提交前必须运行项目门禁。 - -- Handler 单测:覆盖成功、失败、日志关键路径;带参数任务覆盖 `ValidatePayload` 成功、空 payload、非法 JSON、缺失必填、标准化。 -- Admin dispatch 单测:合法 payload 返回成功,非法 payload 返回 400 且错误清晰。 -- Retry 单测:failed 且可重试能创建新执行记录;非 failed、`Retryable=false`、超过 `MaxRetry` 都拒绝。 -- 目标包测试示例: - -```bash -go test ./internal/task ./internal/apps/admin/task ./internal/apps/ -``` - -- API 改动后: - -```bash -make swagger -``` - -- 提交前: - -```bash -make code-check -``` - -如涉及前端或整体构建,补跑 `make build-test`。测试中需要 Redis/Asynq 时,可用现有测试模式或 `miniredis`;不要为了测试便利把 `internal/task` 反向塞进通用 util/testhelper,避免 import cycle。 - -## 相关 Skills - -- Go 错误处理:在设计 handler 返回错误、包装底层错误或避免“记录并返回”时,参见 [go-error-handling](../go-error-handling/SKILL.md)。 -- Go 测试:在编写 handler、dispatch 或 retry 单测时,参见 [go-testing](../go-testing/SKILL.md)。 -- Go context:在任务业务逻辑传播取消、超时或 request scoped 值时,参见 [go-context](../go-context/SKILL.md)。 -- Go logging:在决定任务日志、应用日志和日志级别边界时,参见 [go-logging](../go-logging/SKILL.md)。 -- shadcn:在修改 Admin 任务 UI 时,参见 [shadcn](../shadcn/SKILL.md)。 +- 页面根容器使用 `w-full`,不添加页面级 `max-w-*`。 diff --git a/internal/model/task_execution.go b/internal/model/task_execution.go index 8fbe8b79..81a1d193 100644 --- a/internal/model/task_execution.go +++ b/internal/model/task_execution.go @@ -6,12 +6,14 @@ package model import ( "context" + "errors" "fmt" + "strings" "time" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/db/idgen" - "gorm.io/gorm" + "github.com/redis/go-redis/v9" ) // TaskExecutionStatus 任务执行状态 @@ -23,6 +25,10 @@ const ( TaskExecutionStatusRunning TaskExecutionStatus = "running" TaskExecutionStatusSucceeded TaskExecutionStatus = "succeeded" TaskExecutionStatusFailed TaskExecutionStatus = "failed" + + taskExecutionLogRedisKeyPrefix = "task:execution:log:" + taskExecutionLogExpiration = 24 * time.Hour + taskExecutionLogMaxLines = 1000 ) // TaskExecution 任务执行记录 @@ -58,7 +64,7 @@ func CreateTaskExecution(ctx context.Context, execution *TaskExecution) error { return db.DB(ctx).Create(execution).Error } -// UpdateTaskExecution 更新任务执行记录,忽略 log 字段以防覆写正在追加的日志 +// UpdateTaskExecution 更新任务执行记录,忽略由 Redis 缓冲和归档流程管理的 log 字段。 func UpdateTaskExecution(ctx context.Context, execution *TaskExecution) error { return db.DB(ctx).Omit("log").Save(execution).Error } @@ -69,6 +75,9 @@ func GetTaskExecutionByTaskID(ctx context.Context, taskID string) (*TaskExecutio if err := db.DB(ctx).Where("task_id = ?", taskID).First(&execution).Error; err != nil { return nil, err } + if err := loadTaskExecutionLog(ctx, &execution); err != nil { + return nil, err + } return &execution, nil } @@ -78,16 +87,64 @@ func GetTaskExecutionByID(ctx context.Context, id uint64) (*TaskExecution, error if err := db.DB(ctx).Where("id = ?", id).First(&execution).Error; err != nil { return nil, err } + if err := loadTaskExecutionLog(ctx, &execution); err != nil { + return nil, err + } return &execution, nil } -// AppendTaskExecutionLog 追加日志到执行记录 +// AppendTaskExecutionLog 将日志追加到 Redis 缓冲,任务完成后再持久化到数据库。 func AppendTaskExecutionLog(ctx context.Context, taskID string, logLine string) error { + if db.Redis == nil { + return errors.New("redis client is not initialized") + } + now := time.Now().Format("15:04:05") line := fmt.Sprintf("[%s] %s\n", now, logLine) - return db.DB(ctx).Model(&TaskExecution{}). + key := taskExecutionLogRedisKey(taskID) + + _, err := db.Redis.TxPipelined(ctx, func(pipe redis.Pipeliner) error { + pipe.RPush(ctx, key, line) + pipe.LTrim(ctx, key, -taskExecutionLogMaxLines, -1) + pipe.Expire(ctx, key, taskExecutionLogExpiration) + return nil + }) + if err != nil { + return fmt.Errorf("append task execution log to redis: %w", err) + } + return nil +} + +// FlushTaskExecutionLog 将 Redis 中的完整任务日志写入数据库,并在成功后清理缓存。 +func FlushTaskExecutionLog(ctx context.Context, taskID string) error { + if db.Redis == nil { + return errors.New("redis client is not initialized") + } + + key := taskExecutionLogRedisKey(taskID) + logLines, err := db.Redis.LRange(ctx, key, 0, -1).Result() + if err != nil { + return fmt.Errorf("get task execution log from redis: %w", err) + } + if len(logLines) == 0 { + return nil + } + logText := strings.Join(logLines, "") + + result := db.DB(ctx).Model(&TaskExecution{}). Where("task_id = ?", taskID). - Update("log", gorm.Expr("COALESCE(log, '') || ?", line)).Error + Update("log", logText) + if result.Error != nil { + return fmt.Errorf("persist task execution log: %w", result.Error) + } + if result.RowsAffected == 0 { + return fmt.Errorf("persist task execution log: task %q not found", taskID) + } + + if err := db.Redis.Del(ctx, key).Err(); err != nil { + return fmt.Errorf("delete persisted task execution log from redis: %w", err) + } + return nil } // ListTaskExecutionsRequest 查询任务执行记录列表请求 @@ -126,6 +183,55 @@ func ListTaskExecutions(ctx context.Context, req ListTaskExecutionsRequest) ([]T if err := query.Order("id DESC").Offset(offset).Limit(req.PageSize).Find(&executions).Error; err != nil { return nil, 0, err } + if err := loadTaskExecutionLogs(ctx, executions); err != nil { + return nil, 0, err + } return executions, total, nil } + +func taskExecutionLogRedisKey(taskID string) string { + return db.PrefixedKey(taskExecutionLogRedisKeyPrefix + taskID) +} + +func loadTaskExecutionLog(ctx context.Context, execution *TaskExecution) error { + if db.Redis == nil { + return nil + } + + logLines, err := db.Redis.LRange(ctx, taskExecutionLogRedisKey(execution.TaskID), 0, -1).Result() + if err != nil { + return fmt.Errorf("get task execution log from redis: %w", err) + } + if len(logLines) == 0 { + return nil + } + + execution.Log = strings.Join(logLines, "") + return nil +} + +func loadTaskExecutionLogs(ctx context.Context, executions []TaskExecution) error { + if db.Redis == nil || len(executions) == 0 { + return nil + } + + commands := make([]*redis.StringSliceCmd, len(executions)) + _, err := db.Redis.Pipelined(ctx, func(pipe redis.Pipeliner) error { + for i := range executions { + commands[i] = pipe.LRange(ctx, taskExecutionLogRedisKey(executions[i].TaskID), 0, -1) + } + return nil + }) + if err != nil { + return fmt.Errorf("get task execution logs from redis: %w", err) + } + + for i := range executions { + logLines := commands[i].Val() + if len(logLines) > 0 { + executions[i].Log = strings.Join(logLines, "") + } + } + return nil +} diff --git a/internal/model/task_execution_test.go b/internal/model/task_execution_test.go index 458996ab..87c4040b 100644 --- a/internal/model/task_execution_test.go +++ b/internal/model/task_execution_test.go @@ -6,11 +6,14 @@ package model import ( "context" + "fmt" "testing" "time" "github.com/Rain-kl/Wavelet/internal/db" + "github.com/alicebob/miniredis/v2" "github.com/glebarez/sqlite" + "github.com/redis/go-redis/v9" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "gorm.io/gorm" @@ -25,10 +28,18 @@ func setupTaskExecutionTestEnvironment(t *testing.T) func() { err = sqliteDB.AutoMigrate(&TaskExecution{}) require.NoError(t, err) + miniRedis, err := miniredis.Run() + require.NoError(t, err) + redisClient := redis.NewClient(&redis.Options{Addr: miniRedis.Addr()}) + db.SetDB(sqliteDB) + db.Redis = redisClient return func() { + require.NoError(t, redisClient.Close()) + miniRedis.Close() db.SetDB(nil) + db.Redis = nil } } @@ -188,7 +199,7 @@ func TestUpdateTaskExecutionFailed(t *testing.T) { assert.Equal(t, int64(200), found.Duration) } -func TestUpdateTaskExecutionDoesNotOverwriteLog(t *testing.T) { +func TestUpdateTaskExecutionDoesNotPersistBufferedLog(t *testing.T) { cleanup := setupTaskExecutionTestEnvironment(t) defer cleanup() ctx := context.Background() @@ -203,23 +214,25 @@ func TestUpdateTaskExecutionDoesNotOverwriteLog(t *testing.T) { err := CreateTaskExecution(ctx, execution) require.NoError(t, err) - // In a real execution, logs are appended to the DB asynchronously via AppendTaskExecutionLog + // 运行中的日志仅缓存在 Redis。 err = AppendTaskExecutionLog(ctx, "test_omit_log_001", "第一条执行日志") require.NoError(t, err) - // The local struct still has empty Log because it was not reloaded assert.Empty(t, execution.Log) - // Now complete/update the execution (e.g. status, duration) execution.Status = TaskExecutionStatusSucceeded execution.Duration = 100 err = UpdateTaskExecution(ctx, execution) require.NoError(t, err) - // Get the updated execution record and check that the Log was NOT overwritten/wiped + var persisted TaskExecution + err = db.DB(ctx).Where("task_id = ?", "test_omit_log_001").First(&persisted).Error + require.NoError(t, err) + assert.Equal(t, TaskExecutionStatusSucceeded, persisted.Status) + assert.Empty(t, persisted.Log) + found, err := GetTaskExecutionByTaskID(ctx, "test_omit_log_001") require.NoError(t, err) - assert.Equal(t, TaskExecutionStatusSucceeded, found.Status) assert.Contains(t, found.Log, "第一条执行日志") } @@ -248,12 +261,51 @@ func TestAppendTaskExecutionLog(t *testing.T) { err = AppendTaskExecutionLog(ctx, "test_log_001", "清理完成,共删除 42 个文件") require.NoError(t, err) - // 验证日志内容 + // 读取时优先返回 Redis 中的在途日志。 found, err := GetTaskExecutionByTaskID(ctx, "test_log_001") require.NoError(t, err) assert.Contains(t, found.Log, "开始扫描未使用上传文件") assert.Contains(t, found.Log, "本批次找到 42 个待清理文件") assert.Contains(t, found.Log, "清理完成,共删除 42 个文件") + + var persisted TaskExecution + err = db.DB(ctx).Where("task_id = ?", "test_log_001").First(&persisted).Error + require.NoError(t, err) + assert.Empty(t, persisted.Log) + + err = FlushTaskExecutionLog(ctx, "test_log_001") + require.NoError(t, err) + + err = db.DB(ctx).Where("task_id = ?", "test_log_001").First(&persisted).Error + require.NoError(t, err) + assert.Contains(t, persisted.Log, "开始扫描未使用上传文件") + + exists, err := db.Redis.Exists(ctx, taskExecutionLogRedisKey("test_log_001")).Result() + require.NoError(t, err) + assert.Zero(t, exists) +} + +func TestAppendTaskExecutionLogLimitsLinesAndRefreshesTTL(t *testing.T) { + cleanup := setupTaskExecutionTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + const taskID = "limited_log_001" + for i := 0; i < taskExecutionLogMaxLines+5; i++ { + err := AppendTaskExecutionLog(ctx, taskID, fmt.Sprintf("日志-%04d", i)) + require.NoError(t, err) + } + + key := taskExecutionLogRedisKey(taskID) + logLines, err := db.Redis.LRange(ctx, key, 0, -1).Result() + require.NoError(t, err) + assert.Len(t, logLines, taskExecutionLogMaxLines) + assert.Contains(t, logLines[0], "日志-0005") + assert.Contains(t, logLines[len(logLines)-1], "日志-1004") + + ttl, err := db.Redis.TTL(ctx, key).Result() + require.NoError(t, err) + assert.Equal(t, taskExecutionLogExpiration, ttl) } func TestAppendTaskExecutionLogNonExistent(t *testing.T) { @@ -261,10 +313,36 @@ func TestAppendTaskExecutionLogNonExistent(t *testing.T) { defer cleanup() ctx := context.Background() - // 对不存在的 TaskID 追加日志不应报错(COALESCE 处理空值) + // Redis 缓冲不依赖数据库记录是否已经创建。 err := AppendTaskExecutionLog(ctx, "nonexistent_task", "测试日志") - // SQLite 下 COALESCE + || 操作不应报错 assert.NoError(t, err) + + err = FlushTaskExecutionLog(ctx, "nonexistent_task") + assert.Error(t, err) +} + +func TestGetTaskExecutionLogPrefersRedis(t *testing.T) { + cleanup := setupTaskExecutionTestEnvironment(t) + defer cleanup() + ctx := context.Background() + + execution := &TaskExecution{ + TaskID: "redis_priority_001", + TaskType: "upload:cleanup_unused", + TaskName: "清理未使用上传", + Status: TaskExecutionStatusRunning, + Log: "数据库旧日志", + TriggeredBy: "manual", + } + err := CreateTaskExecution(ctx, execution) + require.NoError(t, err) + err = AppendTaskExecutionLog(ctx, execution.TaskID, "Redis 最新日志") + require.NoError(t, err) + + found, err := GetTaskExecutionByID(ctx, execution.ID) + require.NoError(t, err) + assert.Contains(t, found.Log, "Redis 最新日志") + assert.NotContains(t, found.Log, "数据库旧日志") } func TestListTaskExecutions(t *testing.T) { @@ -284,12 +362,19 @@ func TestListTaskExecutions(t *testing.T) { err := CreateTaskExecution(ctx, r) require.NoError(t, err) } + err := AppendTaskExecutionLog(ctx, "list_004", "运行中的 Redis 日志") + require.NoError(t, err) // 查询全部(分页) items, total, err := ListTaskExecutions(ctx, ListTaskExecutionsRequest{Page: 1, PageSize: 10}) require.NoError(t, err) assert.Equal(t, int64(5), total) assert.Len(t, items, 5) + for _, item := range items { + if item.TaskID == "list_004" { + assert.Contains(t, item.Log, "运行中的 Redis 日志") + } + } // 按状态筛选:failed items, total, err = ListTaskExecutions(ctx, ListTaskExecutionsRequest{Status: "failed", Page: 1, PageSize: 10}) diff --git a/internal/task/executor.go b/internal/task/executor.go index 55f466ab..21bff30b 100644 --- a/internal/task/executor.go +++ b/internal/task/executor.go @@ -128,7 +128,9 @@ 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)) + if err := model.AppendTaskExecutionLog(ctx, taskID, fmt.Sprintf("[系统] 任务已成功入队,等待调度执行 (队列: %s, 最大重试次数: %d)", meta.Queue, meta.MaxRetry)); err != nil { + logger.ErrorF(ctx, "[TaskExecutor] 追加入队日志失败 taskID=%s: %v", taskID, err) + } return taskID, nil } @@ -189,7 +191,9 @@ 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)) + if err := model.AppendTaskExecutionLog(ctx, newTaskID, fmt.Sprintf("[系统] 手动触发重试,已重新创建任务并入队 (原任务ID: %s, 重试次数: %d/%d)", execution.TaskID, execution.RetryCount+1, execution.MaxRetry)); err != nil { + logger.ErrorF(ctx, "[TaskExecutor] 追加重试日志失败 taskID=%s: %v", newTaskID, err) + } return newTaskID, nil } @@ -311,6 +315,24 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution, if err := model.UpdateTaskExecution(ctx, execution); err != nil { logger.ErrorF(ctx, "[TaskExecutor] 更新执行记录失败 taskID=%s: %v", execution.TaskID, err) } + if shouldFlushTaskExecutionLog(ctx, execErr) { + if err := model.FlushTaskExecutionLog(ctx, execution.TaskID); err != nil { + logger.ErrorF(ctx, "[TaskExecutor] 持久化任务日志失败 taskID=%s: %v", execution.TaskID, err) + } + } +} + +func shouldFlushTaskExecutionLog(ctx context.Context, execErr error) bool { + if execErr == nil { + return true + } + + retryCount, hasRetryCount := asynq.GetRetryCount(ctx) + maxRetry, hasMaxRetry := asynq.GetMaxRetry(ctx) + if !hasRetryCount || !hasMaxRetry { + return true + } + return retryCount >= maxRetry } func handleFailedTask(ctx context.Context, execution *model.TaskExecution, t *asynq.Task, duration time.Duration, execErr error, span trace.Span) { diff --git a/internal/task/executor_test.go b/internal/task/executor_test.go index 0691f87d..6e0bbab6 100644 --- a/internal/task/executor_test.go +++ b/internal/task/executor_test.go @@ -15,6 +15,7 @@ import ( "github.com/hibiken/asynq" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/trace" ) // mockHandler 用于测试的模拟任务处理器 @@ -209,6 +210,43 @@ func TestProcessTaskFailure(t *testing.T) { assert.Contains(t, found.Log, "开始执行任务") } +func TestCompleteTaskExecutionFlushesLog(t *testing.T) { + cleanup := setupTest(t) + defer cleanup() + ctx := context.Background() + + execution := &model.TaskExecution{ + TaskID: "complete_flush_001", + TaskType: testTaskType, + TaskName: "测试任务", + Status: model.TaskExecutionStatusRunning, + TriggeredBy: "manual", + } + err := model.CreateTaskExecution(ctx, execution) + require.NoError(t, err) + + ctx = withTaskID(ctx, execution.TaskID) + AppendLog(ctx, "任务执行中的日志") + + finishTime := time.Now() + completeTaskExecution( + ctx, + execution, + asynq.NewTask(testTaskType, nil), + 100*time.Millisecond, + finishTime, + &TaskResult{Message: "处理完成"}, + nil, + trace.SpanFromContext(ctx), + ) + + found, err := model.GetTaskExecutionByTaskID(ctx, execution.TaskID) + require.NoError(t, err) + assert.Equal(t, model.TaskExecutionStatusSucceeded, found.Status) + assert.Contains(t, found.Log, "任务执行中的日志") + assert.Contains(t, found.Log, "任务执行成功") +} + func TestRetryTask(t *testing.T) { cleanup := setupTest(t) defer cleanup()