nomadcode/services/core/internal/workflow/service.go
toki 6f383f1213 workflow: enqueue failure 시 Lifecycle를 통해 task를 failed 상태로 전환
- enqueuer가 설정되지 않은 경우에도 FailTask를 호출하여 task 상태를 failed로 표시
- scheduler/jobs.go에 책임 경계 주석 추가
- README에 작업 라이프사이클 책임 경계 문서 추가
- EnqueueTask 실패 시 lifecycle를 통한 상태 전이 테스트 추가
2026-06-01 10:48:29 +09:00

175 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
lifecycle *Lifecycle
}
func NewService(store *storage.Store, enqueuer TaskEnqueuer, logger *slog.Logger) *Service {
return &Service{
store: store,
enqueuer: enqueuer,
logger: logger,
lifecycle: NewLifecycle(store, 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) {
queuedTask, err := s.lifecycle.QueueTask(ctx, id)
if err != nil {
return storage.Task{}, err
}
if s.enqueuer == nil {
err := errors.New("task enqueuer is not configured")
if _, failErr := s.lifecycle.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
}
if err := s.enqueuer.EnqueueTask(ctx, id); err != nil {
if _, failErr := s.lifecycle.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 {
return canEnqueueStatus(TaskStatus(status))
}