- notification 서비스 테스트 추가 - workflow lifecycle 관리 구현 - config, scheduler, workflow 서비스 개선
182 lines
4.4 KiB
Go
182 lines
4.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")
|
|
ErrInvalidTaskTransition = errors.New("invalid task transition")
|
|
)
|
|
|
|
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
|
|
}
|
|
|
|
metadata, err := NormalizeTaskMetadata(input.Metadata)
|
|
if err != nil {
|
|
return storage.Task{}, err
|
|
}
|
|
|
|
external, err := NormalizeExternalRef(input.External)
|
|
if err != nil {
|
|
return storage.Task{}, err
|
|
}
|
|
|
|
return s.store.CreateTask(ctx, storage.CreateTaskInput{
|
|
Title: title,
|
|
Source: source,
|
|
Payload: payload,
|
|
Metadata: metadata,
|
|
ExternalProvider: external.Provider,
|
|
ExternalID: external.ID,
|
|
ExternalURL: external.URL,
|
|
ExternalMetadata: external.Metadata,
|
|
})
|
|
}
|
|
|
|
func (s *Service) GetTask(ctx context.Context, id string) (storage.Task, error) {
|
|
return s.store.GetTask(ctx, id)
|
|
}
|
|
|
|
type NormalizedExternalRef struct {
|
|
Provider *string
|
|
ID *string
|
|
URL *string
|
|
Metadata json.RawMessage
|
|
}
|
|
|
|
func NormalizeExternalRef(input *ExternalRefInput) (NormalizedExternalRef, error) {
|
|
if input == nil {
|
|
return NormalizedExternalRef{Metadata: json.RawMessage(`{}`)}, nil
|
|
}
|
|
|
|
provider := strings.TrimSpace(input.Provider)
|
|
id := strings.TrimSpace(input.ID)
|
|
externalURL := strings.TrimSpace(input.URL)
|
|
metadata := input.Metadata
|
|
if len(metadata) == 0 || string(metadata) == "null" {
|
|
metadata = json.RawMessage(`{}`)
|
|
}
|
|
if !json.Valid(metadata) {
|
|
return NormalizedExternalRef{}, ErrInvalidTaskInput
|
|
}
|
|
|
|
if provider == "" {
|
|
return NormalizedExternalRef{}, ErrInvalidTaskInput
|
|
}
|
|
|
|
return NormalizedExternalRef{
|
|
Provider: optionalString(provider),
|
|
ID: optionalString(id),
|
|
URL: optionalString(externalURL),
|
|
Metadata: metadata,
|
|
}, nil
|
|
}
|
|
|
|
func optionalString(value string) *string {
|
|
if value == "" {
|
|
return nil
|
|
}
|
|
return &value
|
|
}
|
|
|
|
func NormalizeTaskMetadata(metadata json.RawMessage) (json.RawMessage, error) {
|
|
if len(metadata) == 0 || string(metadata) == "null" {
|
|
return json.RawMessage(`{}`), nil
|
|
}
|
|
if !json.Valid(metadata) {
|
|
return nil, ErrInvalidTaskInput
|
|
}
|
|
return metadata, nil
|
|
}
|
|
|
|
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) UpdateTaskMetadata(ctx context.Context, id string, metadata json.RawMessage) (storage.Task, error) {
|
|
normalized, err := NormalizeTaskMetadata(metadata)
|
|
if err != nil {
|
|
return storage.Task{}, err
|
|
}
|
|
return s.store.UpdateMetadata(ctx, id, normalized)
|
|
}
|
|
|
|
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
|
|
}
|
|
}
|