/
AngelCareMe
/
task_queue
Обзор
Документация
Войти
/
AngelCareMe
/
task_queue
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
internal/service/task_service.go
180 строк
5 KB
AngelCareMe
done
28 июл 2025, 21:07
28 июл 2025, 21:07
c70014b
Код
Авторство
О чём код?
package service import ( "context" "fmt" "task-queue/internal/adapter/repository" "task-queue/internal/entity" "task-queue/internal/usecase" "task-queue/pkg/logger" "github.com/sirupsen/logrus" ) type taskService struct { taskUseCase usecase.TaskUseCase logger *logrus.Logger } // NewTaskService creates new task service instance func NewTaskService(taskUseCase usecase.TaskUseCase) TaskService { return &taskService{ taskUseCase: taskUseCase, logger: logger.Get(), } } type TaskService interface { // createTask creates new task through usecase CreateTask(ctx context.Context, req *TaskCreateRequest) (*entity.Task, error) // getTask gets task by id through usecase GetTask(ctx context.Context, id string) (*entity.Task, error) // listTasks lists tasks with filtering through usecase ListTasks(ctx context.Context, req *TaskListRequest) (*TaskListResponse, error) // retryTask retries task through usecase RetryTask(ctx context.Context, id string) error // deleteTask deletes task through usecase DeleteTask(ctx context.Context, id string) error // getTaskStats gets task statistics through usecase GetTaskStats(ctx context.Context) (*TaskStatsResponse, error) } type TaskCreateRequest struct { Name string `json:"name" validate:"required"` Payload string `json:"payload"` MaxRetry int `json:"max_retry"` Delay int `json:"delay"` } type TaskListRequest struct { Status *entity.TaskStatus `json:"status"` Name *string `json:"name"` Limit int `json:"limit"` Offset int `json:"offset"` } type TaskListResponse struct { Tasks []*entity.Task `json:"tasks"` Total int `json:"total"` } type TaskStatsResponse struct { Total int64 `json:"total"` Pending int64 `json:"pending"` Processing int64 `json:"processing"` Success int64 `json:"success"` Failed int64 `json:"failed"` ByDate map[string]int64 `json:"by_date"` } func (s *taskService) CreateTask(ctx context.Context, req *TaskCreateRequest) (*entity.Task, error) { if req == nil { return nil, fmt.Errorf("task create request is required") } task, err := s.taskUseCase.CreateTask(ctx, req.Name, req.Payload, req.MaxRetry, req.Delay) if err != nil { s.logger.WithError(err).Error("failed to create task through usecase") return nil, fmt.Errorf("failed to create task: %w", err) } return task, nil } func (s *taskService) GetTask(ctx context.Context, id string) (*entity.Task, error) { if id == "" { return nil, fmt.Errorf("task id is required") } task, err := s.taskUseCase.GetTask(ctx, id) if err != nil { s.logger.WithError(err).WithField("task_id", id).Error("failed to get task through usecase") return nil, fmt.Errorf("failed to get task: %w", err) } return task, nil } func (s *taskService) ListTasks(ctx context.Context, req *TaskListRequest) (*TaskListResponse, error) { if req == nil { req = &TaskListRequest{} } filter := repository.TaskFilter{ Status: req.Status, Name: req.Name, } tasks, err := s.taskUseCase.ListTasks(ctx, filter, req.Limit, req.Offset) if err != nil { s.logger.WithError(err).Error("failed to list tasks through usecase") return nil, fmt.Errorf("failed to list tasks: %w", err) } // get total count for pagination stats, err := s.taskUseCase.GetStats(ctx) if err != nil { s.logger.WithError(err).Error("failed to get stats for task list") // don't fail completely if stats failed stats = &repository.TaskStats{} } response := &TaskListResponse{ Tasks: tasks, Total: int(stats.Total), } return response, nil } func (s *taskService) RetryTask(ctx context.Context, id string) error { if id == "" { return fmt.Errorf("task id is required") } err := s.taskUseCase.RetryTask(ctx, id) if err != nil { s.logger.WithError(err).WithField("task_id", id).Error("failed to retry task through usecase") return fmt.Errorf("failed to retry task: %w", err) } return nil } func (s *taskService) DeleteTask(ctx context.Context, id string) error { if id == "" { return fmt.Errorf("task id is required") } err := s.taskUseCase.DeleteTask(ctx, id) if err != nil { s.logger.WithError(err).WithField("task_id", id).Error("failed to delete task through usecase") return fmt.Errorf("failed to delete task: %w", err) } return nil } func (s *taskService) GetTaskStats(ctx context.Context) (*TaskStatsResponse, error) { stats, err := s.taskUseCase.GetStats(ctx) if err != nil { s.logger.WithError(err).Error("failed to get task stats through usecase") return nil, fmt.Errorf("failed to get task stats: %w", err) } response := &TaskStatsResponse{ Total: stats.Total, Pending: stats.Pending, Processing: stats.Processing, Success: stats.Success, Failed: stats.Failed, ByDate: stats.ByDate, } return response, nil }