nomadcode/services/core/internal/workflow/service.go
toki e7ddbcdfed feat: add task metadata migration and update workflow service
- Add migration for task metadata column (00003)
- Update workflow model and service with metadata support
- Update database layer (models, queries) for metadata field
- Update storage store to handle metadata
- Update roadmap documents and milestones
2026-05-24 18:10:40 +09:00

181 lines
4.3 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
}
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
}
}