Files
pig-farm-controller/internal/infra/repository/execution_log_repository.go

352 lines
16 KiB
Go
Raw Normal View History

package repository
import (
2025-11-05 23:00:07 +08:00
"context"
2025-09-22 15:17:54 +08:00
"errors"
"time"
2025-09-22 15:17:54 +08:00
2025-11-05 23:00:07 +08:00
"git.huangwc.com/pig/pig-farm-controller/internal/infra/logs"
"git.huangwc.com/pig/pig-farm-controller/internal/infra/models"
2025-11-05 23:00:07 +08:00
"gorm.io/gorm"
)
2025-10-18 15:36:32 +08:00
// PlanExecutionLogListOptions 定义了查询计划执行日志时的可选参数
type PlanExecutionLogListOptions struct {
2025-11-10 22:23:31 +08:00
PlanID *uint32
2025-10-18 15:36:32 +08:00
Status *models.ExecutionStatus
StartTime *time.Time // 基于 created_at 字段
EndTime *time.Time // 基于 created_at 字段
OrderBy string // 例如 "created_at asc"
}
2025-10-18 15:39:47 +08:00
// TaskExecutionLogListOptions 定义了查询任务执行日志时的可选参数
type TaskExecutionLogListOptions struct {
2025-11-10 22:23:31 +08:00
PlanExecutionLogID *uint32
2025-10-18 15:39:47 +08:00
TaskID *int
Status *models.ExecutionStatus
StartTime *time.Time // 基于 created_at 字段
EndTime *time.Time // 基于 created_at 字段
OrderBy string // 例如 "created_at asc"
}
// ExecutionLogRepository 定义了与执行日志交互的接口。
type ExecutionLogRepository interface {
2025-10-18 15:36:32 +08:00
// --- Existing methods ---
2025-11-10 22:23:31 +08:00
UpdateTaskExecutionLogStatusByIDs(ctx context.Context, logIDs []uint32, status models.ExecutionStatus) error
UpdateTaskExecutionLogStatus(ctx context.Context, logID uint32, status models.ExecutionStatus) error
2025-11-05 23:00:07 +08:00
CreateTaskExecutionLog(ctx context.Context, log *models.TaskExecutionLog) error
CreatePlanExecutionLog(ctx context.Context, log *models.PlanExecutionLog) error
UpdatePlanExecutionLog(ctx context.Context, log *models.PlanExecutionLog) error
CreateTaskExecutionLogsInBatch(ctx context.Context, logs []*models.TaskExecutionLog) error
UpdateTaskExecutionLog(ctx context.Context, log *models.TaskExecutionLog) error
2025-11-10 22:23:31 +08:00
FindTaskExecutionLogByID(ctx context.Context, id uint32) (*models.TaskExecutionLog, error)
2025-09-20 22:41:03 +08:00
// UpdatePlanExecutionLogStatus 更新计划执行日志的状态
2025-11-10 22:23:31 +08:00
UpdatePlanExecutionLogStatus(ctx context.Context, logID uint32, status models.ExecutionStatus) error
2025-09-20 23:50:27 +08:00
// UpdatePlanExecutionLogsStatusByIDs 批量更新计划执行日志的状态
2025-11-10 22:23:31 +08:00
UpdatePlanExecutionLogsStatusByIDs(ctx context.Context, logIDs []uint32, status models.ExecutionStatus) error
2025-09-20 23:50:27 +08:00
// FindIncompletePlanExecutionLogs 查找所有未完成的计划执行日志
2025-11-05 23:00:07 +08:00
FindIncompletePlanExecutionLogs(ctx context.Context) ([]models.PlanExecutionLog, error)
2025-09-22 15:17:54 +08:00
// FindInProgressPlanExecutionLogByPlanID 根据 PlanID 查找正在进行的计划执行日志
2025-11-10 22:23:31 +08:00
FindInProgressPlanExecutionLogByPlanID(ctx context.Context, planID uint32) (*models.PlanExecutionLog, error)
2025-09-22 15:17:54 +08:00
// FindIncompleteTaskExecutionLogsByPlanLogID 根据计划日志ID查找所有未完成的任务日志
2025-11-10 22:23:31 +08:00
FindIncompleteTaskExecutionLogsByPlanLogID(ctx context.Context, planLogID uint32) ([]models.TaskExecutionLog, error)
// FailAllIncompletePlanExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的计划状态都修改为 ExecutionStatusFailed
2025-11-05 23:00:07 +08:00
FailAllIncompletePlanExecutionLogs(ctx context.Context) error
// CancelAllIncompleteTaskExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的任务状态修改为 ExecutionStatusCancelled
2025-11-05 23:00:07 +08:00
CancelAllIncompleteTaskExecutionLogs(ctx context.Context) error
// FindPlanExecutionLogByID 根据ID查找计划执行日志
2025-11-10 22:23:31 +08:00
FindPlanExecutionLogByID(ctx context.Context, id uint32) (*models.PlanExecutionLog, error)
// CountIncompleteTasksByPlanLogID 计算一个计划执行中未完成的任务数量
2025-11-10 22:23:31 +08:00
CountIncompleteTasksByPlanLogID(ctx context.Context, planLogID uint32) (int64, error)
// FailPlanExecution 将指定的计划执行标记为失败
2025-11-10 22:23:31 +08:00
FailPlanExecution(ctx context.Context, planLogID uint32, errorMessage string) error
// CancelIncompleteTasksByPlanLogID 取消一个计划执行中的所有未完成任务
2025-11-10 22:23:31 +08:00
CancelIncompleteTasksByPlanLogID(ctx context.Context, planLogID uint32, reason string) error
2025-10-18 15:36:32 +08:00
2025-10-18 15:39:47 +08:00
// --- New methods ---
2025-11-05 23:00:07 +08:00
ListPlanExecutionLogs(ctx context.Context, opts PlanExecutionLogListOptions, page, pageSize int) ([]models.PlanExecutionLog, int64, error)
ListTaskExecutionLogs(ctx context.Context, opts TaskExecutionLogListOptions, page, pageSize int) ([]models.TaskExecutionLog, int64, error)
}
// gormExecutionLogRepository 是使用 GORM 的具体实现。
type gormExecutionLogRepository struct {
2025-11-05 23:00:07 +08:00
ctx context.Context
db *gorm.DB
}
// NewGormExecutionLogRepository 创建一个新的执行日志仓库。
2025-11-05 23:00:07 +08:00
func NewGormExecutionLogRepository(ctx context.Context, db *gorm.DB) ExecutionLogRepository {
return &gormExecutionLogRepository{ctx: ctx, db: db}
}
2025-10-18 15:36:32 +08:00
// ListPlanExecutionLogs 实现了分页和过滤查询计划执行日志的功能
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) ListPlanExecutionLogs(ctx context.Context, opts PlanExecutionLogListOptions, page, pageSize int) ([]models.PlanExecutionLog, int64, error) {
repoCtx := logs.AddFuncName(ctx, r.ctx, "ListPlanExecutionLogs")
2025-10-18 15:36:32 +08:00
if page <= 0 || pageSize <= 0 {
return nil, 0, ErrInvalidPagination
}
var results []models.PlanExecutionLog
var total int64
2025-11-05 23:00:07 +08:00
query := r.db.WithContext(repoCtx).Model(&models.PlanExecutionLog{})
2025-10-18 15:36:32 +08:00
if opts.PlanID != nil {
query = query.Where("plan_id = ?", *opts.PlanID)
}
if opts.Status != nil {
query = query.Where("status = ?", *opts.Status)
}
if opts.StartTime != nil {
query = query.Where("created_at >= ?", *opts.StartTime)
}
if opts.EndTime != nil {
query = query.Where("created_at <= ?", *opts.EndTime)
}
if err := query.Count(&total).Error; err != nil {
return nil, 0, err
}
orderBy := "created_at DESC"
if opts.OrderBy != "" {
orderBy = opts.OrderBy
}
query = query.Order(orderBy)
offset := (page - 1) * pageSize
err := query.Limit(pageSize).Offset(offset).Find(&results).Error
return results, total, err
}
2025-10-18 15:39:47 +08:00
// ListTaskExecutionLogs 实现了分页和过滤查询任务执行日志的功能
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) ListTaskExecutionLogs(ctx context.Context, opts TaskExecutionLogListOptions, page, pageSize int) ([]models.TaskExecutionLog, int64, error) {
repoCtx := logs.AddFuncName(ctx, r.ctx, "ListTaskExecutionLogs")
2025-10-18 15:39:47 +08:00
if page <= 0 || pageSize <= 0 {
return nil, 0, ErrInvalidPagination
}
var results []models.TaskExecutionLog
var total int64
2025-11-05 23:00:07 +08:00
query := r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{})
2025-10-18 15:39:47 +08:00
if opts.PlanExecutionLogID != nil {
query = query.Where("plan_execution_log_id = ?", *opts.PlanExecutionLogID)
}
if opts.TaskID != nil {
query = query.Where("task_id = ?", *opts.TaskID)
}
if opts.Status != nil {
query = query.Where("status = ?", *opts.Status)
}
if opts.StartTime != nil {
query = query.Where("created_at >= ?", *opts.StartTime)
}
if opts.EndTime != nil {
query = query.Where("created_at <= ?", *opts.EndTime)
}
if err := query.Count(&total).Error; err != nil {
return nil, 0, err
}
orderBy := "created_at DESC"
if opts.OrderBy != "" {
orderBy = opts.OrderBy
}
// 预加载关联的Task信息
query = query.Order(orderBy).Preload("Task")
offset := (page - 1) * pageSize
err := query.Limit(pageSize).Offset(offset).Find(&results).Error
return results, total, err
}
2025-10-18 15:36:32 +08:00
// --- Existing method implementations ---
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) UpdateTaskExecutionLogStatusByIDs(ctx context.Context, logIDs []uint32, status models.ExecutionStatus) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdateTaskExecutionLogStatusByIDs")
2025-09-20 21:14:58 +08:00
if len(logIDs) == 0 {
return nil
}
2025-11-05 23:00:07 +08:00
return r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{}).Where("id IN ?", logIDs).Update("status", status).Error
2025-09-20 21:14:58 +08:00
}
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) UpdateTaskExecutionLogStatus(ctx context.Context, logID uint32, status models.ExecutionStatus) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdateTaskExecutionLogStatus")
return r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{}).Where("id = ?", logID).Update("status", status).Error
2025-09-20 21:14:58 +08:00
}
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) CreateTaskExecutionLog(ctx context.Context, log *models.TaskExecutionLog) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "CreateTaskExecutionLog")
return r.db.WithContext(repoCtx).Create(log).Error
2025-09-20 21:14:58 +08:00
}
// CreatePlanExecutionLog 为一次计划执行创建一条新的日志条目。
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) CreatePlanExecutionLog(ctx context.Context, log *models.PlanExecutionLog) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "CreatePlanExecutionLog")
return r.db.WithContext(repoCtx).Create(log).Error
}
// UpdatePlanExecutionLog 使用 Updates 方法更新一个计划执行日志。
// GORM 的 Updates 传入 struct 时,只会更新非零值字段。
// 在这里,我们期望传入的对象一定包含一个有效的 ID。
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) UpdatePlanExecutionLog(ctx context.Context, log *models.PlanExecutionLog) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdatePlanExecutionLog")
return r.db.WithContext(repoCtx).Updates(log).Error
}
// CreateTaskExecutionLogsInBatch 在一次数据库调用中创建多个任务执行日志条目。
// 这是“预写日志”步骤的关键。
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) CreateTaskExecutionLogsInBatch(ctx context.Context, executionLogs []*models.TaskExecutionLog) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "CreateTaskExecutionLogsInBatch")
if len(executionLogs) == 0 {
2025-09-20 21:14:58 +08:00
return nil
}
2025-10-06 15:08:32 +08:00
// GORM 的 CreateTx 传入一个切片指针会执行批量插入。
2025-11-05 23:00:07 +08:00
return r.db.WithContext(repoCtx).Create(&executionLogs).Error
}
// UpdateTaskExecutionLog 使用 Updates 方法更新一个任务执行日志。
// GORM 的 Updates 传入 struct 时,只会更新非零值字段。
// 这种方式代码更直观,上层服务可以直接修改模型对象后进行保存。
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) UpdateTaskExecutionLog(ctx context.Context, log *models.TaskExecutionLog) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdateTaskExecutionLog")
return r.db.WithContext(repoCtx).Updates(log).Error
}
// FindTaskExecutionLogByID 根据 ID 查找单个任务执行日志。
// 它会预加载关联的 Task 信息。
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) FindTaskExecutionLogByID(ctx context.Context, id uint32) (*models.TaskExecutionLog, error) {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "FindTaskExecutionLogByID")
var log models.TaskExecutionLog
// 使用 Preload("Task") 来确保关联的任务信息被一并加载
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).Preload("Task").First(&log, id).Error
if err != nil {
return nil, err
}
return &log, nil
}
2025-09-20 22:41:03 +08:00
// UpdatePlanExecutionLogStatus 更新计划执行日志的状态
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) UpdatePlanExecutionLogStatus(ctx context.Context, logID uint32, status models.ExecutionStatus) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdatePlanExecutionLogStatus")
return r.db.WithContext(repoCtx).Model(&models.PlanExecutionLog{}).Where("id = ?", logID).Update("status", status).Error
2025-09-20 22:41:03 +08:00
}
2025-09-20 23:50:27 +08:00
// UpdatePlanExecutionLogsStatusByIDs 批量更新计划执行日志的状态
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) UpdatePlanExecutionLogsStatusByIDs(ctx context.Context, logIDs []uint32, status models.ExecutionStatus) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "UpdatePlanExecutionLogsStatusByIDs")
2025-09-20 23:50:27 +08:00
if len(logIDs) == 0 {
return nil
}
2025-11-05 23:00:07 +08:00
return r.db.WithContext(repoCtx).Model(&models.PlanExecutionLog{}).Where("id IN ?", logIDs).Update("status", status).Error
2025-09-20 23:50:27 +08:00
}
// FindIncompletePlanExecutionLogs 查找所有未完成的计划执行日志
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) FindIncompletePlanExecutionLogs(ctx context.Context) ([]models.PlanExecutionLog, error) {
repoCtx := logs.AddFuncName(ctx, r.ctx, "FindIncompletePlanExecutionLogs")
2025-09-20 23:50:27 +08:00
var logs []models.PlanExecutionLog
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).Where("status = ? OR status = ?", models.ExecutionStatusStarted, models.ExecutionStatusWaiting).Find(&logs).Error
2025-09-20 23:50:27 +08:00
return logs, err
}
2025-09-22 15:17:54 +08:00
// FindInProgressPlanExecutionLogByPlanID 根据 PlanID 查找正在进行的计划执行日志
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) FindInProgressPlanExecutionLogByPlanID(ctx context.Context, planID uint32) (*models.PlanExecutionLog, error) {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "FindInProgressPlanExecutionLogByPlanID")
2025-09-22 15:17:54 +08:00
var log models.PlanExecutionLog
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).Where("plan_id = ? AND status = ?", planID, models.ExecutionStatusStarted).First(&log).Error
2025-09-22 15:17:54 +08:00
if err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
// 未找到不是一个需要上报的错误,代表计划当前没有在运行
return nil, nil
}
// 其他数据库错误
return nil, err
}
return &log, nil
}
// FindIncompleteTaskExecutionLogsByPlanLogID 根据计划日志ID查找所有未完成的任务日志
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) FindIncompleteTaskExecutionLogsByPlanLogID(ctx context.Context, planLogID uint32) ([]models.TaskExecutionLog, error) {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "FindIncompleteTaskExecutionLogsByPlanLogID")
2025-09-22 15:17:54 +08:00
var logs []models.TaskExecutionLog
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).Where("plan_execution_log_id = ? AND (status = ? OR status = ?)",
2025-09-22 15:17:54 +08:00
planLogID, models.ExecutionStatusWaiting, models.ExecutionStatusStarted).Find(&logs).Error
return logs, err
}
// FailAllIncompletePlanExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的计划状态都修改为 ExecutionStatusFailed
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) FailAllIncompletePlanExecutionLogs(ctx context.Context) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "FailAllIncompletePlanExecutionLogs")
return r.db.WithContext(repoCtx).Model(&models.PlanExecutionLog{}).
Where("status IN (?, ?)", models.ExecutionStatusStarted, models.ExecutionStatusWaiting).
Updates(map[string]interface{}{"status": models.ExecutionStatusFailed, "ended_at": time.Now(), "error": "系统中断"}).Error
}
// CancelAllIncompleteTaskExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的任务状态修改为 ExecutionStatusCancelled
2025-11-05 23:00:07 +08:00
func (r *gormExecutionLogRepository) CancelAllIncompleteTaskExecutionLogs(ctx context.Context) error {
repoCtx := logs.AddFuncName(ctx, r.ctx, "CancelAllIncompleteTaskExecutionLogs")
return r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{}).
Where("status IN (?, ?)", models.ExecutionStatusStarted, models.ExecutionStatusWaiting).
Updates(map[string]interface{}{"status": models.ExecutionStatusCancelled, "ended_at": time.Now(), "output": "系统中断"}).Error
}
// FindPlanExecutionLogByID 根据ID查找计划执行日志
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) FindPlanExecutionLogByID(ctx context.Context, id uint32) (*models.PlanExecutionLog, error) {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "FindPlanExecutionLogByID")
var log models.PlanExecutionLog
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).First(&log, id).Error
if err != nil {
return nil, err
}
return &log, nil
}
// CountIncompleteTasksByPlanLogID 计算一个计划执行中未完成的任务数量
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) CountIncompleteTasksByPlanLogID(ctx context.Context, planLogID uint32) (int64, error) {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "CountIncompleteTasksByPlanLogID")
var count int64
2025-11-05 23:00:07 +08:00
err := r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{}).
Where("plan_execution_log_id = ? AND status IN (?, ?)",
planLogID, models.ExecutionStatusWaiting, models.ExecutionStatusStarted).
Count(&count).Error
return count, err
}
// FailPlanExecution 将指定的计划执行标记为失败
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) FailPlanExecution(ctx context.Context, planLogID uint32, errorMessage string) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "FailPlanExecution")
return r.db.WithContext(repoCtx).Model(&models.PlanExecutionLog{}).
Where("id = ?", planLogID).
Updates(map[string]interface{}{
"status": models.ExecutionStatusFailed,
"error": errorMessage,
"ended_at": time.Now(),
}).Error
}
// CancelIncompleteTasksByPlanLogID 取消一个计划执行中的所有未完成任务
2025-11-10 22:23:31 +08:00
func (r *gormExecutionLogRepository) CancelIncompleteTasksByPlanLogID(ctx context.Context, planLogID uint32, reason string) error {
2025-11-05 23:00:07 +08:00
repoCtx := logs.AddFuncName(ctx, r.ctx, "CancelIncompleteTasksByPlanLogID")
return r.db.WithContext(repoCtx).Model(&models.TaskExecutionLog{}).
Where("plan_execution_log_id = ? AND status IN (?, ?)",
planLogID, models.ExecutionStatusWaiting, models.ExecutionStatusStarted).
Updates(map[string]interface{}{
"status": models.ExecutionStatusCancelled,
"output": reason,
"ended_at": time.Now(),
}).Error
}