package task import ( "context" taskbiz "kra/internal/biz/task" "kra/pkg/database/gormkit" "kra/pkg/database/pagination" "time" "gorm.io/gorm" ) type taskPO struct { ID uint `gorm:"primaryKey"` CreatedAt time.Time UpdatedAt time.Time DeletedAt gorm.DeletedAt `gorm:"index"` Name string `gorm:"index"` Description, Spec string WithSeconds bool ExecutorType, MethodName string Params gormkit.JSON HTTPURL string HTTPMethod string HTTPHeader gormkit.JSON HTTPBody string `gorm:"type:text"` HTTPAllowPrivate, Enabled bool } func (taskPO) TableName() string { return "sys_timed_tasks" } type taskLogPO struct { ID uint `gorm:"primaryKey"` CreatedAt time.Time UpdatedAt time.Time DeletedAt gorm.DeletedAt `gorm:"index"` TaskID uint `gorm:"index"` TaskName, TriggerType string StartedAt, FinishedAt time.Time DurationMS int64 Status string ErrorMsg, Output string `gorm:"type:text"` } func (taskLogPO) TableName() string { return "sys_timed_task_logs" } type taskRepo struct{ data Provider } func NewTaskRepo(data Provider) taskbiz.TaskRepo { return &taskRepo{data: data} } func taskToPO(v *taskbiz.TimedTask) taskPO { return taskPO{ID: v.ID, Name: v.Name, Description: v.Description, Spec: v.Spec, WithSeconds: v.WithSeconds, ExecutorType: v.ExecutorType, MethodName: v.MethodName, Params: gormkit.JSON(v.Params), HTTPURL: v.HTTPURL, HTTPMethod: v.HTTPMethod, HTTPHeader: gormkit.JSON(v.HTTPHeader), HTTPBody: v.HTTPBody, HTTPAllowPrivate: v.HTTPAllowPrivate, Enabled: v.Enabled} } func taskFromPO(v taskPO) *taskbiz.TimedTask { return &taskbiz.TimedTask{ID: v.ID, CreatedAt: v.CreatedAt, UpdatedAt: v.UpdatedAt, Name: v.Name, Description: v.Description, Spec: v.Spec, WithSeconds: v.WithSeconds, ExecutorType: v.ExecutorType, MethodName: v.MethodName, Params: []byte(v.Params), HTTPURL: v.HTTPURL, HTTPMethod: v.HTTPMethod, HTTPHeader: []byte(v.HTTPHeader), HTTPBody: v.HTTPBody, HTTPAllowPrivate: v.HTTPAllowPrivate, Enabled: v.Enabled} } func (r *taskRepo) CreateTask(ctx context.Context, v *taskbiz.TimedTask) error { po := taskToPO(v) if err := r.data.DB().WithContext(ctx).Create(&po).Error; err != nil { return err } v.ID = po.ID return nil } func (r *taskRepo) TaskNameExists(ctx context.Context, name string, excludeID uint) (bool, error) { var count int64 db := r.data.DB().WithContext(ctx).Model(&taskPO{}).Where("name = ?", name) if excludeID > 0 { db = db.Where("id <> ?", excludeID) } err := db.Count(&count).Error return count > 0, err } func (r *taskRepo) UpdateTask(ctx context.Context, v *taskbiz.TimedTask) error { po := taskToPO(v) return r.data.DB().WithContext(ctx).Model(&taskPO{}).Where("id = ?", v.ID).Select("name", "description", "spec", "with_seconds", "executor_type", "method_name", "params", "http_url", "http_method", "http_header", "http_body", "http_allow_private", "enabled").Updates(&po).Error } func (r *taskRepo) DeleteTask(ctx context.Context, id uint) error { return r.data.DB().WithContext(ctx).Delete(&taskPO{}, id).Error } func (r *taskRepo) FindTask(ctx context.Context, id uint) (*taskbiz.TimedTask, error) { var po taskPO if err := r.data.DB().WithContext(ctx).First(&po, id).Error; err != nil { return nil, err } return taskFromPO(po), nil } func (r *taskRepo) ListTasks(ctx context.Context, page, size int, q *taskbiz.TimedTask) ([]*taskbiz.TimedTask, int64, error) { // During first-install the data layer intentionally serves a bootstrap // database without system tables. The scheduler starts before /init/initdb // and should remain idle instead of logging a missing-table SQL error. if !r.data.DatabaseReady() { return []*taskbiz.TimedTask{}, 0, nil } db := r.data.DB().WithContext(ctx).Model(&taskPO{}) if q != nil { if q.Name != "" { db = db.Where("name LIKE ?", "%"+q.Name+"%") } if q.ExecutorType != "" { db = db.Where("executor_type = ?", q.ExecutorType) } if q.EnabledFilter != nil { db = db.Where("enabled = ?", *q.EnabledFilter) } } var total int64 if err := db.Count(&total).Error; err != nil { return nil, 0, err } var pos []taskPO if err := pagination.Apply(db.Order("id desc"), page, size, 100).Find(&pos).Error; err != nil { return nil, 0, err } out := make([]*taskbiz.TimedTask, 0, len(pos)) for _, po := range pos { out = append(out, taskFromPO(po)) } return out, total, nil } func (r *taskRepo) ToggleTask(ctx context.Context, id uint, enabled bool) error { return r.data.DB().WithContext(ctx).Model(&taskPO{}).Where("id = ?", id).Update("enabled", enabled).Error } func (r *taskRepo) RecordTaskLog(ctx context.Context, v *taskbiz.TimedTaskLog) error { return r.data.DB().WithContext(ctx).Create(&taskLogPO{TaskID: v.TaskID, TaskName: v.TaskName, TriggerType: v.TriggerType, StartedAt: v.StartedAt, FinishedAt: v.FinishedAt, DurationMS: v.DurationMS, Status: v.Status, ErrorMsg: v.ErrorMsg, Output: v.Output}).Error } func taskLogFromPO(v taskLogPO) *taskbiz.TimedTaskLog { return &taskbiz.TimedTaskLog{ID: v.ID, CreatedAt: v.CreatedAt, UpdatedAt: v.UpdatedAt, TaskID: v.TaskID, TaskName: v.TaskName, TriggerType: v.TriggerType, StartedAt: v.StartedAt, FinishedAt: v.FinishedAt, DurationMS: v.DurationMS, Status: v.Status, ErrorMsg: v.ErrorMsg, Output: v.Output} } func (r *taskRepo) ListTaskLogs(ctx context.Context, page, size int, taskID uint, status string) ([]*taskbiz.TimedTaskLog, int64, error) { db := r.data.DB().WithContext(ctx).Model(&taskLogPO{}) if taskID != 0 { db = db.Where("task_id = ?", taskID) } if status != "" { db = db.Where("status = ?", status) } var total int64 if err := db.Count(&total).Error; err != nil { return nil, 0, err } var pos []taskLogPO if err := pagination.Apply(db.Order("id desc"), page, size, 100).Find(&pos).Error; err != nil { return nil, 0, err } out := make([]*taskbiz.TimedTaskLog, 0, len(pos)) for _, po := range pos { out = append(out, taskLogFromPO(po)) } return out, total, nil } func (r *taskRepo) CleanupTaskLogs(ctx context.Context) error { return r.data.DB().WithContext(ctx).Unscoped(). Where("created_at < ?", time.Now().Add(-720*time.Hour)). Delete(&taskLogPO{}).Error }