221 lines
		
	
	
		
			10 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
			
		
		
	
	
			221 lines
		
	
	
		
			10 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
| package repository
 | |
| 
 | |
| import (
 | |
| 	"errors"
 | |
| 	"time"
 | |
| 
 | |
| 	"git.huangwc.com/pig/pig-farm-controller/internal/infra/models"
 | |
| 	"gorm.io/gorm"
 | |
| )
 | |
| 
 | |
| // ExecutionLogRepository 定义了与执行日志交互的接口。
 | |
| // 这为服务层提供了一个清晰的契约,并允许在测试中轻松地进行模拟。
 | |
| type ExecutionLogRepository interface {
 | |
| 	UpdateTaskExecutionLogStatusByIDs(logIDs []uint, status models.ExecutionStatus) error
 | |
| 	UpdateTaskExecutionLogStatus(logID uint, status models.ExecutionStatus) error
 | |
| 	CreateTaskExecutionLog(log *models.TaskExecutionLog) error
 | |
| 	CreatePlanExecutionLog(log *models.PlanExecutionLog) error
 | |
| 	UpdatePlanExecutionLog(log *models.PlanExecutionLog) error
 | |
| 	CreateTaskExecutionLogsInBatch(logs []*models.TaskExecutionLog) error
 | |
| 	UpdateTaskExecutionLog(log *models.TaskExecutionLog) error
 | |
| 	FindTaskExecutionLogByID(id uint) (*models.TaskExecutionLog, error)
 | |
| 	// UpdatePlanExecutionLogStatus 更新计划执行日志的状态
 | |
| 	UpdatePlanExecutionLogStatus(logID uint, status models.ExecutionStatus) error
 | |
| 
 | |
| 	// UpdatePlanExecutionLogsStatusByIDs 批量更新计划执行日志的状态
 | |
| 	UpdatePlanExecutionLogsStatusByIDs(logIDs []uint, status models.ExecutionStatus) error
 | |
| 
 | |
| 	// FindIncompletePlanExecutionLogs 查找所有未完成的计划执行日志
 | |
| 	FindIncompletePlanExecutionLogs() ([]models.PlanExecutionLog, error)
 | |
| 	// FindInProgressPlanExecutionLogByPlanID 根据 PlanID 查找正在进行的计划执行日志
 | |
| 	FindInProgressPlanExecutionLogByPlanID(planID uint) (*models.PlanExecutionLog, error)
 | |
| 	// FindIncompleteTaskExecutionLogsByPlanLogID 根据计划日志ID查找所有未完成的任务日志
 | |
| 	FindIncompleteTaskExecutionLogsByPlanLogID(planLogID uint) ([]models.TaskExecutionLog, error)
 | |
| 
 | |
| 	// FailAllIncompletePlanExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的计划状态都修改为 ExecutionStatusFailed
 | |
| 	FailAllIncompletePlanExecutionLogs() error
 | |
| 	// CancelAllIncompleteTaskExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的任务状态修改为 ExecutionStatusCancelled
 | |
| 	CancelAllIncompleteTaskExecutionLogs() error
 | |
| 
 | |
| 	// FindPlanExecutionLogByID 根据ID查找计划执行日志
 | |
| 	FindPlanExecutionLogByID(id uint) (*models.PlanExecutionLog, error)
 | |
| 
 | |
| 	// CountIncompleteTasksByPlanLogID 计算一个计划执行中未完成的任务数量
 | |
| 	CountIncompleteTasksByPlanLogID(planLogID uint) (int64, error)
 | |
| 
 | |
| 	// FailPlanExecution 将指定的计划执行标记为失败
 | |
| 	FailPlanExecution(planLogID uint, errorMessage string) error
 | |
| 
 | |
| 	// CancelIncompleteTasksByPlanLogID 取消一个计划执行中的所有未完成任务
 | |
| 	CancelIncompleteTasksByPlanLogID(planLogID uint, reason string) error
 | |
| }
 | |
| 
 | |
| // gormExecutionLogRepository 是使用 GORM 的具体实现。
 | |
| type gormExecutionLogRepository struct {
 | |
| 	db *gorm.DB
 | |
| }
 | |
| 
 | |
| // NewGormExecutionLogRepository 创建一个新的执行日志仓库。
 | |
| // 它接收一个 GORM DB 实例作为依赖。
 | |
| func NewGormExecutionLogRepository(db *gorm.DB) ExecutionLogRepository {
 | |
| 	return &gormExecutionLogRepository{db: db}
 | |
| }
 | |
| 
 | |
| func (r *gormExecutionLogRepository) UpdateTaskExecutionLogStatusByIDs(logIDs []uint, status models.ExecutionStatus) error {
 | |
| 	if len(logIDs) == 0 {
 | |
| 		return nil
 | |
| 	}
 | |
| 	return r.db.Model(&models.TaskExecutionLog{}).
 | |
| 		Where("id IN ?", logIDs).
 | |
| 		Update("status", status).Error
 | |
| }
 | |
| 
 | |
| func (r *gormExecutionLogRepository) UpdateTaskExecutionLogStatus(logID uint, status models.ExecutionStatus) error {
 | |
| 	return r.db.Model(&models.TaskExecutionLog{}).Where("id = ?", logID).Update("status", status).Error
 | |
| }
 | |
| 
 | |
| func (r *gormExecutionLogRepository) CreateTaskExecutionLog(log *models.TaskExecutionLog) error {
 | |
| 	return r.db.Create(log).Error
 | |
| }
 | |
| 
 | |
| // CreatePlanExecutionLog 为一次计划执行创建一条新的日志条目。
 | |
| func (r *gormExecutionLogRepository) CreatePlanExecutionLog(log *models.PlanExecutionLog) error {
 | |
| 	return r.db.Create(log).Error
 | |
| }
 | |
| 
 | |
| // UpdatePlanExecutionLog 使用 Updates 方法更新一个计划执行日志。
 | |
| // GORM 的 Updates 传入 struct 时,只会更新非零值字段。
 | |
| // 在这里,我们期望传入的对象一定包含一个有效的 ID。
 | |
| func (r *gormExecutionLogRepository) UpdatePlanExecutionLog(log *models.PlanExecutionLog) error {
 | |
| 	return r.db.Updates(log).Error
 | |
| }
 | |
| 
 | |
| // CreateTaskExecutionLogsInBatch 在一次数据库调用中创建多个任务执行日志条目。
 | |
| // 这是“预写日志”步骤的关键。
 | |
| func (r *gormExecutionLogRepository) CreateTaskExecutionLogsInBatch(logs []*models.TaskExecutionLog) error {
 | |
| 	if len(logs) == 0 {
 | |
| 		return nil
 | |
| 	}
 | |
| 	// GORM 的 Create 传入一个切片指针会执行批量插入。
 | |
| 	return r.db.Create(&logs).Error
 | |
| }
 | |
| 
 | |
| // UpdateTaskExecutionLog 使用 Updates 方法更新一个任务执行日志。
 | |
| // GORM 的 Updates 传入 struct 时,只会更新非零值字段。
 | |
| // 这种方式代码更直观,上层服务可以直接修改模型对象后进行保存。
 | |
| func (r *gormExecutionLogRepository) UpdateTaskExecutionLog(log *models.TaskExecutionLog) error {
 | |
| 	return r.db.Updates(log).Error
 | |
| }
 | |
| 
 | |
| // FindTaskExecutionLogByID 根据 ID 查找单个任务执行日志。
 | |
| // 它会预加载关联的 Task 信息。
 | |
| func (r *gormExecutionLogRepository) FindTaskExecutionLogByID(id uint) (*models.TaskExecutionLog, error) {
 | |
| 	var log models.TaskExecutionLog
 | |
| 	// 使用 Preload("Task") 来确保关联的任务信息被一并加载
 | |
| 	err := r.db.Preload("Task").First(&log, id).Error
 | |
| 	if err != nil {
 | |
| 		return nil, err
 | |
| 	}
 | |
| 	return &log, nil
 | |
| }
 | |
| 
 | |
| // UpdatePlanExecutionLogStatus 更新计划执行日志的状态
 | |
| func (r *gormExecutionLogRepository) UpdatePlanExecutionLogStatus(logID uint, status models.ExecutionStatus) error {
 | |
| 	return r.db.Model(&models.PlanExecutionLog{}).Where("id = ?", logID).Update("status", status).Error
 | |
| }
 | |
| 
 | |
| // UpdatePlanExecutionLogsStatusByIDs 批量更新计划执行日志的状态
 | |
| func (r *gormExecutionLogRepository) UpdatePlanExecutionLogsStatusByIDs(logIDs []uint, status models.ExecutionStatus) error {
 | |
| 	if len(logIDs) == 0 {
 | |
| 		return nil
 | |
| 	}
 | |
| 	return r.db.Model(&models.PlanExecutionLog{}).Where("id IN ?", logIDs).Update("status", status).Error
 | |
| }
 | |
| 
 | |
| // FindIncompletePlanExecutionLogs 查找所有未完成的计划执行日志
 | |
| func (r *gormExecutionLogRepository) FindIncompletePlanExecutionLogs() ([]models.PlanExecutionLog, error) {
 | |
| 	var logs []models.PlanExecutionLog
 | |
| 	err := r.db.Where("status = ? OR status = ?", models.ExecutionStatusStarted, models.ExecutionStatusWaiting).Find(&logs).Error
 | |
| 	return logs, err
 | |
| }
 | |
| 
 | |
| // FindInProgressPlanExecutionLogByPlanID 根据 PlanID 查找正在进行的计划执行日志
 | |
| func (r *gormExecutionLogRepository) FindInProgressPlanExecutionLogByPlanID(planID uint) (*models.PlanExecutionLog, error) {
 | |
| 	var log models.PlanExecutionLog
 | |
| 	err := r.db.Where("plan_id = ? AND status = ?", planID, models.ExecutionStatusStarted).First(&log).Error
 | |
| 	if err != nil {
 | |
| 		if errors.Is(err, gorm.ErrRecordNotFound) {
 | |
| 			// 未找到不是一个需要上报的错误,代表计划当前没有在运行
 | |
| 			return nil, nil
 | |
| 		}
 | |
| 		// 其他数据库错误
 | |
| 		return nil, err
 | |
| 	}
 | |
| 	return &log, nil
 | |
| }
 | |
| 
 | |
| // FindIncompleteTaskExecutionLogsByPlanLogID 根据计划日志ID查找所有未完成的任务日志
 | |
| func (r *gormExecutionLogRepository) FindIncompleteTaskExecutionLogsByPlanLogID(planLogID uint) ([]models.TaskExecutionLog, error) {
 | |
| 	var logs []models.TaskExecutionLog
 | |
| 	err := r.db.Where("plan_execution_log_id = ? AND (status = ? OR status = ?)",
 | |
| 		planLogID, models.ExecutionStatusWaiting, models.ExecutionStatusStarted).Find(&logs).Error
 | |
| 	return logs, err
 | |
| }
 | |
| 
 | |
| // FailAllIncompletePlanExecutionLogs 将所有状态为 ExecutionStatusStarted 和 ExecutionStatusWaiting 的计划状态都修改为 ExecutionStatusFailed
 | |
| func (r *gormExecutionLogRepository) FailAllIncompletePlanExecutionLogs() error {
 | |
| 	return r.db.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
 | |
| func (r *gormExecutionLogRepository) CancelAllIncompleteTaskExecutionLogs() error {
 | |
| 	return r.db.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查找计划执行日志
 | |
| func (r *gormExecutionLogRepository) FindPlanExecutionLogByID(id uint) (*models.PlanExecutionLog, error) {
 | |
| 	var log models.PlanExecutionLog
 | |
| 	err := r.db.First(&log, id).Error
 | |
| 	if err != nil {
 | |
| 		return nil, err
 | |
| 	}
 | |
| 	return &log, nil
 | |
| }
 | |
| 
 | |
| // CountIncompleteTasksByPlanLogID 计算一个计划执行中未完成的任务数量
 | |
| func (r *gormExecutionLogRepository) CountIncompleteTasksByPlanLogID(planLogID uint) (int64, error) {
 | |
| 	var count int64
 | |
| 	err := r.db.Model(&models.TaskExecutionLog{}).
 | |
| 		Where("plan_execution_log_id = ? AND status IN (?, ?)",
 | |
| 			planLogID, models.ExecutionStatusWaiting, models.ExecutionStatusStarted).
 | |
| 		Count(&count).Error
 | |
| 	return count, err
 | |
| }
 | |
| 
 | |
| // FailPlanExecution 将指定的计划执行标记为失败
 | |
| func (r *gormExecutionLogRepository) FailPlanExecution(planLogID uint, errorMessage string) error {
 | |
| 	return r.db.Model(&models.PlanExecutionLog{}).
 | |
| 		Where("id = ?", planLogID).
 | |
| 		Updates(map[string]interface{}{
 | |
| 			"status":   models.ExecutionStatusFailed,
 | |
| 			"error":    errorMessage,
 | |
| 			"ended_at": time.Now(),
 | |
| 		}).Error
 | |
| }
 | |
| 
 | |
| // CancelIncompleteTasksByPlanLogID 取消一个计划执行中的所有未完成任务
 | |
| func (r *gormExecutionLogRepository) CancelIncompleteTasksByPlanLogID(planLogID uint, reason string) error {
 | |
| 	return r.db.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
 | |
| }
 |