nomadcode/services/core/internal/workflow/lifecycle.go
toki 4b6951cfa8 feat: notification 및 workflow 기능 개선
- notification 서비스 테스트 추가
- workflow lifecycle 관리 구현
- config, scheduler, workflow 서비스 개선
2026-05-24 19:36:43 +09:00

213 lines
4.8 KiB
Go

package workflow
import (
"context"
"encoding/json"
"log/slog"
"strings"
"time"
"github.com/nomadcode/nomadcode-core/internal/storage"
)
type Lifecycle struct {
store *storage.Store
logger *slog.Logger
}
func NewLifecycle(store *storage.Store, logger *slog.Logger) *Lifecycle {
return &Lifecycle{
store: store,
logger: logger,
}
}
func canTransition(from, to TaskStatus) bool {
switch from {
case StatusPending, StatusQueued, StatusFailed:
return to == StatusRunning
case StatusRunning:
return to == StatusCompleted || to == StatusFailed || to == StatusCanceled
default:
return false
}
}
func nextTaskAttempt(metadata json.RawMessage) (int, error) {
if len(metadata) == 0 || string(metadata) == "null" {
return 1, nil
}
var currentMap map[string]any
if err := json.Unmarshal(metadata, &currentMap); err != nil {
return 0, ErrInvalidTaskInput
}
attempt := 1
if currentMap != nil {
if val, ok := currentMap[MetadataKeyAttempt]; ok {
switch v := val.(type) {
case float64:
attempt = int(v) + 1
case int:
attempt = v + 1
case int64:
attempt = int(v) + 1
}
}
}
return attempt, nil
}
func mergeTaskMetadata(current json.RawMessage, updates map[string]any) (json.RawMessage, error) {
var currentMap map[string]any
if len(current) > 0 && string(current) != "null" {
if err := json.Unmarshal(current, &currentMap); err != nil {
return nil, ErrInvalidTaskInput
}
} else {
currentMap = make(map[string]any)
}
for k, v := range updates {
if v == nil {
delete(currentMap, k)
} else {
currentMap[k] = v
}
}
return json.Marshal(currentMap)
}
func (l *Lifecycle) StartTask(ctx context.Context, id string) (storage.Task, error) {
task, err := l.store.GetTask(ctx, id)
if err != nil {
return storage.Task{}, err
}
if !canTransition(TaskStatus(task.Status), StatusRunning) {
return storage.Task{}, ErrInvalidTaskTransition
}
attempt, err := nextTaskAttempt(task.Metadata)
if err != nil {
return storage.Task{}, err
}
updates := map[string]any{
MetadataKeyAttempt: attempt,
MetadataKeyAgentRunState: "running",
MetadataKeyLastHeartbeat: time.Now().UTC().Format(time.RFC3339),
}
newMeta, err := mergeTaskMetadata(task.Metadata, updates)
if err != nil {
return storage.Task{}, err
}
task, err = l.store.UpdateMetadata(ctx, id, newMeta)
if err != nil {
return storage.Task{}, err
}
return l.store.UpdateStatus(ctx, id, string(StatusRunning))
}
func (l *Lifecycle) CompleteTask(ctx context.Context, id string, result json.RawMessage) (storage.Task, error) {
task, err := l.store.GetTask(ctx, id)
if err != nil {
return storage.Task{}, err
}
if !canTransition(TaskStatus(task.Status), StatusCompleted) {
return storage.Task{}, ErrInvalidTaskTransition
}
updates := map[string]any{
MetadataKeyAgentRunState: "completed",
MetadataKeyWaitType: nil,
}
var summary string
if len(result) > 0 && string(result) != "null" {
var resMap map[string]interface{}
if err := json.Unmarshal(result, &resMap); err == nil {
if sVal, ok := resMap["summary"].(string); ok {
summary = strings.TrimSpace(sVal)
}
}
}
if summary != "" {
updates[MetadataKeyStatusReason] = summary
}
newMeta, err := mergeTaskMetadata(task.Metadata, updates)
if err != nil {
return storage.Task{}, err
}
task, err = l.store.UpdateMetadata(ctx, id, newMeta)
if err != nil {
return storage.Task{}, err
}
return l.store.CompleteTask(ctx, id, result)
}
func (l *Lifecycle) FailTask(ctx context.Context, id string, message string) (storage.Task, error) {
task, err := l.store.GetTask(ctx, id)
if err != nil {
return storage.Task{}, err
}
if !canTransition(TaskStatus(task.Status), StatusFailed) {
return storage.Task{}, ErrInvalidTaskTransition
}
updates := map[string]any{
MetadataKeyAgentRunState: "failed",
MetadataKeyStatusReason: message,
}
newMeta, err := mergeTaskMetadata(task.Metadata, updates)
if err != nil {
return storage.Task{}, err
}
task, err = l.store.UpdateMetadata(ctx, id, newMeta)
if err != nil {
return storage.Task{}, err
}
return l.store.FailTask(ctx, id, message)
}
func (l *Lifecycle) CancelTask(ctx context.Context, id string, message string) (storage.Task, error) {
task, err := l.store.GetTask(ctx, id)
if err != nil {
return storage.Task{}, err
}
if !canTransition(TaskStatus(task.Status), StatusCanceled) {
return storage.Task{}, ErrInvalidTaskTransition
}
updates := map[string]any{
MetadataKeyAgentRunState: "canceled",
}
if message != "" {
updates[MetadataKeyStatusReason] = message
}
newMeta, err := mergeTaskMetadata(task.Metadata, updates)
if err != nil {
return storage.Task{}, err
}
task, err = l.store.UpdateMetadata(ctx, id, newMeta)
if err != nil {
return storage.Task{}, err
}
return l.store.UpdateStatus(ctx, id, string(StatusCanceled))
}