git-subtree-dir: services/core git-subtree-mainline:6f5e3a119fgit-subtree-split:6fdbc73753
102 lines
2.4 KiB
Go
102 lines
2.4 KiB
Go
package workflow
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"log/slog"
|
|
"strings"
|
|
|
|
"github.com/nomadcode/nomadcode-core/internal/storage"
|
|
)
|
|
|
|
var (
|
|
ErrInvalidTaskInput = errors.New("invalid task input")
|
|
ErrTaskCannotBeEnqueued = errors.New("task cannot be enqueued in current status")
|
|
)
|
|
|
|
type TaskEnqueuer interface {
|
|
EnqueueTask(ctx context.Context, taskID string) error
|
|
}
|
|
|
|
type Service struct {
|
|
store *storage.Store
|
|
enqueuer TaskEnqueuer
|
|
logger *slog.Logger
|
|
}
|
|
|
|
func NewService(store *storage.Store, enqueuer TaskEnqueuer, logger *slog.Logger) *Service {
|
|
return &Service{
|
|
store: store,
|
|
enqueuer: enqueuer,
|
|
logger: logger,
|
|
}
|
|
}
|
|
|
|
func (s *Service) CreateTask(ctx context.Context, input CreateTaskInput) (storage.Task, error) {
|
|
title := strings.TrimSpace(input.Title)
|
|
source := strings.TrimSpace(input.Source)
|
|
if title == "" || source == "" {
|
|
return storage.Task{}, ErrInvalidTaskInput
|
|
}
|
|
|
|
payload := input.Payload
|
|
if len(payload) == 0 || string(payload) == "null" {
|
|
payload = json.RawMessage(`{}`)
|
|
}
|
|
if !json.Valid(payload) {
|
|
return storage.Task{}, ErrInvalidTaskInput
|
|
}
|
|
|
|
return s.store.CreateTask(ctx, title, source, payload)
|
|
}
|
|
|
|
func (s *Service) GetTask(ctx context.Context, id string) (storage.Task, error) {
|
|
return s.store.GetTask(ctx, id)
|
|
}
|
|
|
|
func (s *Service) ListTasks(ctx context.Context, limit int32) ([]storage.Task, error) {
|
|
if limit <= 0 {
|
|
limit = 20
|
|
}
|
|
if limit > 100 {
|
|
limit = 100
|
|
}
|
|
return s.store.ListTasks(ctx, limit)
|
|
}
|
|
|
|
func (s *Service) EnqueueTask(ctx context.Context, id string) (storage.Task, error) {
|
|
task, err := s.store.GetTask(ctx, id)
|
|
if err != nil {
|
|
return storage.Task{}, err
|
|
}
|
|
if !canEnqueue(task.Status) {
|
|
return storage.Task{}, ErrTaskCannotBeEnqueued
|
|
}
|
|
|
|
queuedTask, err := s.store.UpdateStatus(ctx, id, string(StatusQueued))
|
|
if err != nil {
|
|
return storage.Task{}, err
|
|
}
|
|
|
|
if s.enqueuer == nil {
|
|
return queuedTask, errors.New("task enqueuer is not configured")
|
|
}
|
|
if err := s.enqueuer.EnqueueTask(ctx, id); err != nil {
|
|
if _, failErr := s.store.FailTask(ctx, id, err.Error()); failErr != nil && s.logger != nil {
|
|
s.logger.Error("failed to mark task enqueue error", "task_id", id, "error", failErr)
|
|
}
|
|
return queuedTask, err
|
|
}
|
|
|
|
return queuedTask, nil
|
|
}
|
|
|
|
func canEnqueue(status string) bool {
|
|
switch TaskStatus(status) {
|
|
case StatusPending, StatusFailed:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|