mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
任务日志优化
This commit is contained in:
@@ -1,139 +1,103 @@
|
|||||||
---
|
---
|
||||||
name: "new-async-task"
|
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/handler.go`:`TaskHandler`、`TaskResult`、`PayloadValidator`
|
||||||
- `internal/task/constants.go`: 框架通用常量,如 `QueueDefault` 和 `DefaultMaxRetry`。
|
- `internal/task/meta.go`:`TaskMeta`、`TaskParam`
|
||||||
- `internal/task/meta.go`: 框架任务元数据结构体(TaskParam、TaskMeta)及全局动态注册与查询接口。
|
- `internal/task/executor.go`:下发、执行、日志、重试
|
||||||
- `internal/task/executor.go`: `RegisterHandler`、`ValidateAndNormalizePayload`、`DispatchTask`、`RetryTask`、`ProcessTask`、`AppendLog`。
|
- `internal/task/handlers/register.go`:Handler 和元数据注册
|
||||||
- `internal/task/handlers/register.go`: 内置 handler 和元数据的统一注册点,Admin API 和 Worker 都依赖它。
|
- `internal/task/worker/worker.go`:Worker 路由和队列
|
||||||
- `internal/task/worker/worker.go`: Asynq mux 动态路由分发和队列配置。
|
- `internal/task/scheduler/scheduler.go`:定时调度
|
||||||
- `internal/task/scheduler/scheduler.go`: Cron 调度。
|
- `internal/apps/admin/task/routers.go`:Admin 任务 API
|
||||||
- `internal/apps/admin/task/routers.go`: Admin 下发、查询、详情、重试 API。
|
- `internal/model/task_execution.go`:执行记录和日志持久化
|
||||||
- 现有参考:`internal/apps/upload/tasks.go`(无参数任务)、`internal/apps/user/tasks.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/<module>/tasks.go` 定义任务类型、Admin 任务类型和 `TaskMeta`。
|
||||||
|
- Asynq 任务类型使用 `<module>:<action>` 格式。
|
||||||
|
- 完整设置 `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. 定义任务元数据与常量。
|
- Handler 必须实现 `task.TaskHandler`。
|
||||||
- 业务包常量与元数据:在对应业务模块的 `internal/apps/<module>/tasks.go` 中定义 Asynq 任务类型常量(如 `CleanupUnusedUploadsTask = "upload:cleanup_unused"`)和 Admin 任务类型常量(如 `TaskTypeCleanupUploads = "cleanup_unused_uploads"`)。
|
- 有参数任务必须实现 `task.PayloadValidator`,负责校验和标准化 Admin 下发参数。
|
||||||
- 在同一 `tasks.go` 文件中定义该任务的 `TaskMeta` 元数据变量(如 `CleanupUnusedUploadsMeta = task.TaskMeta{...}`),配置 `Type`、`AsynqTask`、`Name`、`Description`、`MaxRetry`、`Queue`、`Retryable` 等字段。
|
- `Execute` 必须再次解析 payload;不要假设入口一定经过 Admin 校验。
|
||||||
- 有参数任务在 `Params` 中描述前端表单字段。`TaskParam.Name` 必须与 payload JSON tag 对齐。
|
- 成功返回 `&task.TaskResult{Message: ..., Detail: ...}`。
|
||||||
|
- 失败返回 error,由任务框架处理状态和重试。
|
||||||
|
- 不要吞掉关键错误。
|
||||||
|
- 复杂 SQL 放到 `internal/model/` 或 `internal/service/`。
|
||||||
|
- 新增 Go 文件后检查许可证头,必要时运行 `make license`。
|
||||||
|
|
||||||
2. 实现 handler。
|
### 注册
|
||||||
- 优先放在对应业务模块的 `internal/apps/<module>/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` 同时注册 Handler 和 `TaskMeta`。
|
||||||
- 在 `internal/task/handlers/register.go` 导入业务模块,调用 `task.RegisterHandler(asynqTaskType, handler)` 注册处理器。
|
- 不要在其他位置单独注册任务。
|
||||||
- 同时,在该文件中调用 `task.RegisterTaskMeta(meta)` 注册刚才在业务模块中定义的任务元数据。
|
|
||||||
- 这里是 Admin 校验、元数据获取和 Worker 执行共同依赖的注册点。
|
|
||||||
|
|
||||||
4. 如需 Cron 调度,系统默认定时任务必须通过 SQL 迁移(goose)初始化。
|
## 日志要求
|
||||||
- 确保任务已正确注册并载入全局元数据池中。
|
|
||||||
- 在 `internal/db/migrator/goose/postgres` 和 `sqlite` 下编写 migration 脚本,使用 `INSERT INTO schedules` 语句初始化任务,指定 `task_type` 和 `cron` 等字段。必须妥善处理冲突(如 `ON CONFLICT DO NOTHING`)以支持幂等。
|
|
||||||
|
|
||||||
5. 如改动 Admin API。
|
- 在 `TaskHandler.Execute` 中使用 `task.AppendLog(ctx, format, args...)`。
|
||||||
- handler 放在 `internal/apps/admin/<module>/` 或现有 Admin task 模块内。
|
- 记录任务开始、参数摘要、批次进度、关键状态、可继续错误和完成摘要。
|
||||||
- 路由只在 `internal/router/router.go` 注册。
|
- 批量处理按批次记录;禁止为大循环中的每条数据写日志。
|
||||||
- 响应保持 `{ "error_msg": "", "data": ... }`,分页保持 `{ "total": 0, "results": [] }`。
|
- 不要直接修改任务日志的 Redis key 或 `w_task_executions.log`。
|
||||||
- 补完整 Swagger 注释并运行 `make swagger`。
|
|
||||||
|
|
||||||
## 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。
|
||||||
|
|
||||||
- 任务开始、参数摘要、批次进度、关键状态、可继续错误、完成摘要适合记录。
|
## Admin API
|
||||||
- 避免对大循环中的每条记录都写日志;每次 `AppendLog` 都可能触发一次数据库更新。
|
|
||||||
- 如果上下文中没有 taskID,`AppendLog` 会降级为普通应用日志,不应额外兜底报错。
|
|
||||||
|
|
||||||
## 重试规则
|
- Handler 放在现有 Admin task 模块或 `internal/apps/admin/<module>/`。
|
||||||
|
- 路由只在 `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`。
|
- API 调用必须通过 `frontend/lib/services/`。
|
||||||
|
- 修改 shadcn/ui 时使用 `shadcn` skill。
|
||||||
修改重试语义时同时检查 `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。
|
|
||||||
- 不使用 `any`。
|
- 不使用 `any`。
|
||||||
- 页面根容器保持 `w-full`,不要加页面级 `max-w-*`。
|
- 页面根容器使用 `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/<module>
|
|
||||||
```
|
|
||||||
|
|
||||||
- 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)。
|
|
||||||
|
|||||||
@@ -6,12 +6,14 @@ package model
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/db"
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||||
"gorm.io/gorm"
|
"github.com/redis/go-redis/v9"
|
||||||
)
|
)
|
||||||
|
|
||||||
// TaskExecutionStatus 任务执行状态
|
// TaskExecutionStatus 任务执行状态
|
||||||
@@ -23,6 +25,10 @@ const (
|
|||||||
TaskExecutionStatusRunning TaskExecutionStatus = "running"
|
TaskExecutionStatusRunning TaskExecutionStatus = "running"
|
||||||
TaskExecutionStatusSucceeded TaskExecutionStatus = "succeeded"
|
TaskExecutionStatusSucceeded TaskExecutionStatus = "succeeded"
|
||||||
TaskExecutionStatusFailed TaskExecutionStatus = "failed"
|
TaskExecutionStatusFailed TaskExecutionStatus = "failed"
|
||||||
|
|
||||||
|
taskExecutionLogRedisKeyPrefix = "task:execution:log:"
|
||||||
|
taskExecutionLogExpiration = 24 * time.Hour
|
||||||
|
taskExecutionLogMaxLines = 1000
|
||||||
)
|
)
|
||||||
|
|
||||||
// TaskExecution 任务执行记录
|
// TaskExecution 任务执行记录
|
||||||
@@ -58,7 +64,7 @@ func CreateTaskExecution(ctx context.Context, execution *TaskExecution) error {
|
|||||||
return db.DB(ctx).Create(execution).Error
|
return db.DB(ctx).Create(execution).Error
|
||||||
}
|
}
|
||||||
|
|
||||||
// UpdateTaskExecution 更新任务执行记录,忽略 log 字段以防覆写正在追加的日志
|
// UpdateTaskExecution 更新任务执行记录,忽略由 Redis 缓冲和归档流程管理的 log 字段。
|
||||||
func UpdateTaskExecution(ctx context.Context, execution *TaskExecution) error {
|
func UpdateTaskExecution(ctx context.Context, execution *TaskExecution) error {
|
||||||
return db.DB(ctx).Omit("log").Save(execution).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 {
|
if err := db.DB(ctx).Where("task_id = ?", taskID).First(&execution).Error; err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
if err := loadTaskExecutionLog(ctx, &execution); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
return &execution, nil
|
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 {
|
if err := db.DB(ctx).Where("id = ?", id).First(&execution).Error; err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
if err := loadTaskExecutionLog(ctx, &execution); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
return &execution, nil
|
return &execution, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// AppendTaskExecutionLog 追加日志到执行记录
|
// AppendTaskExecutionLog 将日志追加到 Redis 缓冲,任务完成后再持久化到数据库。
|
||||||
func AppendTaskExecutionLog(ctx context.Context, taskID string, logLine string) error {
|
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")
|
now := time.Now().Format("15:04:05")
|
||||||
line := fmt.Sprintf("[%s] %s\n", now, logLine)
|
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).
|
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 查询任务执行记录列表请求
|
// 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 {
|
if err := query.Order("id DESC").Offset(offset).Limit(req.PageSize).Find(&executions).Error; err != nil {
|
||||||
return nil, 0, err
|
return nil, 0, err
|
||||||
}
|
}
|
||||||
|
if err := loadTaskExecutionLogs(ctx, executions); err != nil {
|
||||||
|
return nil, 0, err
|
||||||
|
}
|
||||||
|
|
||||||
return executions, total, nil
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -6,11 +6,14 @@ package model
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/db"
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
|
"github.com/alicebob/miniredis/v2"
|
||||||
"github.com/glebarez/sqlite"
|
"github.com/glebarez/sqlite"
|
||||||
|
"github.com/redis/go-redis/v9"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
@@ -25,10 +28,18 @@ func setupTaskExecutionTestEnvironment(t *testing.T) func() {
|
|||||||
err = sqliteDB.AutoMigrate(&TaskExecution{})
|
err = sqliteDB.AutoMigrate(&TaskExecution{})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
miniRedis, err := miniredis.Run()
|
||||||
|
require.NoError(t, err)
|
||||||
|
redisClient := redis.NewClient(&redis.Options{Addr: miniRedis.Addr()})
|
||||||
|
|
||||||
db.SetDB(sqliteDB)
|
db.SetDB(sqliteDB)
|
||||||
|
db.Redis = redisClient
|
||||||
|
|
||||||
return func() {
|
return func() {
|
||||||
|
require.NoError(t, redisClient.Close())
|
||||||
|
miniRedis.Close()
|
||||||
db.SetDB(nil)
|
db.SetDB(nil)
|
||||||
|
db.Redis = nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -188,7 +199,7 @@ func TestUpdateTaskExecutionFailed(t *testing.T) {
|
|||||||
assert.Equal(t, int64(200), found.Duration)
|
assert.Equal(t, int64(200), found.Duration)
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestUpdateTaskExecutionDoesNotOverwriteLog(t *testing.T) {
|
func TestUpdateTaskExecutionDoesNotPersistBufferedLog(t *testing.T) {
|
||||||
cleanup := setupTaskExecutionTestEnvironment(t)
|
cleanup := setupTaskExecutionTestEnvironment(t)
|
||||||
defer cleanup()
|
defer cleanup()
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
@@ -203,23 +214,25 @@ func TestUpdateTaskExecutionDoesNotOverwriteLog(t *testing.T) {
|
|||||||
err := CreateTaskExecution(ctx, execution)
|
err := CreateTaskExecution(ctx, execution)
|
||||||
require.NoError(t, err)
|
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", "第一条执行日志")
|
err = AppendTaskExecutionLog(ctx, "test_omit_log_001", "第一条执行日志")
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
// The local struct still has empty Log because it was not reloaded
|
|
||||||
assert.Empty(t, execution.Log)
|
assert.Empty(t, execution.Log)
|
||||||
|
|
||||||
// Now complete/update the execution (e.g. status, duration)
|
|
||||||
execution.Status = TaskExecutionStatusSucceeded
|
execution.Status = TaskExecutionStatusSucceeded
|
||||||
execution.Duration = 100
|
execution.Duration = 100
|
||||||
err = UpdateTaskExecution(ctx, execution)
|
err = UpdateTaskExecution(ctx, execution)
|
||||||
require.NoError(t, err)
|
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")
|
found, err := GetTaskExecutionByTaskID(ctx, "test_omit_log_001")
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
assert.Equal(t, TaskExecutionStatusSucceeded, found.Status)
|
|
||||||
assert.Contains(t, found.Log, "第一条执行日志")
|
assert.Contains(t, found.Log, "第一条执行日志")
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -248,12 +261,51 @@ func TestAppendTaskExecutionLog(t *testing.T) {
|
|||||||
err = AppendTaskExecutionLog(ctx, "test_log_001", "清理完成,共删除 42 个文件")
|
err = AppendTaskExecutionLog(ctx, "test_log_001", "清理完成,共删除 42 个文件")
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
// 验证日志内容
|
// 读取时优先返回 Redis 中的在途日志。
|
||||||
found, err := GetTaskExecutionByTaskID(ctx, "test_log_001")
|
found, err := GetTaskExecutionByTaskID(ctx, "test_log_001")
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
assert.Contains(t, found.Log, "开始扫描未使用上传文件")
|
assert.Contains(t, found.Log, "开始扫描未使用上传文件")
|
||||||
assert.Contains(t, found.Log, "本批次找到 42 个待清理文件")
|
assert.Contains(t, found.Log, "本批次找到 42 个待清理文件")
|
||||||
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) {
|
func TestAppendTaskExecutionLogNonExistent(t *testing.T) {
|
||||||
@@ -261,10 +313,36 @@ func TestAppendTaskExecutionLogNonExistent(t *testing.T) {
|
|||||||
defer cleanup()
|
defer cleanup()
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
// 对不存在的 TaskID 追加日志不应报错(COALESCE 处理空值)
|
// Redis 缓冲不依赖数据库记录是否已经创建。
|
||||||
err := AppendTaskExecutionLog(ctx, "nonexistent_task", "测试日志")
|
err := AppendTaskExecutionLog(ctx, "nonexistent_task", "测试日志")
|
||||||
// SQLite 下 COALESCE + || 操作不应报错
|
|
||||||
assert.NoError(t, err)
|
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) {
|
func TestListTaskExecutions(t *testing.T) {
|
||||||
@@ -284,12 +362,19 @@ func TestListTaskExecutions(t *testing.T) {
|
|||||||
err := CreateTaskExecution(ctx, r)
|
err := CreateTaskExecution(ctx, r)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
}
|
}
|
||||||
|
err := AppendTaskExecutionLog(ctx, "list_004", "运行中的 Redis 日志")
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
// 查询全部(分页)
|
// 查询全部(分页)
|
||||||
items, total, err := ListTaskExecutions(ctx, ListTaskExecutionsRequest{Page: 1, PageSize: 10})
|
items, total, err := ListTaskExecutions(ctx, ListTaskExecutionsRequest{Page: 1, PageSize: 10})
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
assert.Equal(t, int64(5), total)
|
assert.Equal(t, int64(5), total)
|
||||||
assert.Len(t, items, 5)
|
assert.Len(t, items, 5)
|
||||||
|
for _, item := range items {
|
||||||
|
if item.TaskID == "list_004" {
|
||||||
|
assert.Contains(t, item.Log, "运行中的 Redis 日志")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// 按状态筛选:failed
|
// 按状态筛选:failed
|
||||||
items, total, err = ListTaskExecutions(ctx, ListTaskExecutionsRequest{Status: "failed", Page: 1, PageSize: 10})
|
items, total, err = ListTaskExecutions(ctx, ListTaskExecutionsRequest{Status: "failed", Page: 1, PageSize: 10})
|
||||||
|
|||||||
@@ -128,7 +128,9 @@ func DispatchTask(ctx context.Context, taskType string, payload []byte, triggere
|
|||||||
return "", fmt.Errorf(errTaskEnqueueFailed, err)
|
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
|
return taskID, nil
|
||||||
}
|
}
|
||||||
@@ -189,7 +191,9 @@ func RetryTask(ctx context.Context, id uint64) (string, error) {
|
|||||||
return "", fmt.Errorf(errRetryTaskEnqueueFailed, err)
|
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
|
return newTaskID, nil
|
||||||
}
|
}
|
||||||
@@ -311,6 +315,24 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution,
|
|||||||
if err := model.UpdateTaskExecution(ctx, execution); err != nil {
|
if err := model.UpdateTaskExecution(ctx, execution); err != nil {
|
||||||
logger.ErrorF(ctx, "[TaskExecutor] 更新执行记录失败 taskID=%s: %v", execution.TaskID, err)
|
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) {
|
func handleFailedTask(ctx context.Context, execution *model.TaskExecution, t *asynq.Task, duration time.Duration, execErr error, span trace.Span) {
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ import (
|
|||||||
"github.com/hibiken/asynq"
|
"github.com/hibiken/asynq"
|
||||||
"github.com/stretchr/testify/assert"
|
"github.com/stretchr/testify/assert"
|
||||||
"github.com/stretchr/testify/require"
|
"github.com/stretchr/testify/require"
|
||||||
|
"go.opentelemetry.io/otel/trace"
|
||||||
)
|
)
|
||||||
|
|
||||||
// mockHandler 用于测试的模拟任务处理器
|
// mockHandler 用于测试的模拟任务处理器
|
||||||
@@ -209,6 +210,43 @@ func TestProcessTaskFailure(t *testing.T) {
|
|||||||
assert.Contains(t, found.Log, "开始执行任务")
|
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) {
|
func TestRetryTask(t *testing.T) {
|
||||||
cleanup := setupTest(t)
|
cleanup := setupTest(t)
|
||||||
defer cleanup()
|
defer cleanup()
|
||||||
|
|||||||
Reference in New Issue
Block a user