package workflow import ( "context" "encoding/json" "log/slog" "strings" "time" "github.com/nomadcode/nomadcode-core/internal/storage" ) type taskStore interface { CreateTask(context.Context, storage.CreateTaskInput) (storage.Task, error) GetTask(context.Context, string) (storage.Task, error) GetTaskByExternalRef(context.Context, string, string) (storage.Task, error) ListTasks(context.Context, int32) ([]storage.Task, error) UpdateMetadata(context.Context, string, json.RawMessage) (storage.Task, error) UpdateMetadataIfNonTerminal(context.Context, string, json.RawMessage) (storage.Task, bool, error) UpdateStatus(context.Context, string, string) (storage.Task, error) CompleteTask(context.Context, string, json.RawMessage) (storage.Task, error) FailTask(context.Context, string, string) (storage.Task, error) } type Lifecycle struct { store taskStore logger *slog.Logger } func NewLifecycle(store *storage.Store, logger *slog.Logger) *Lifecycle { return &Lifecycle{ store: store, logger: logger, } } func validTaskStatus(status TaskStatus) bool { switch status { case StatusPending, StatusQueued, StatusRunning, StatusCompleted, StatusFailed, StatusCanceled: return true default: return false } } func terminalTaskStatus(status TaskStatus) bool { switch status { case StatusCompleted, StatusFailed, StatusCanceled: return true default: return false } } func canEnqueueStatus(status TaskStatus) bool { switch status { case StatusPending, StatusFailed: return true default: return false } } func canTransition(from, to TaskStatus) bool { switch from { case StatusPending, StatusFailed: return to == StatusRunning case StatusQueued: return to == StatusRunning || to == StatusFailed || to == StatusCanceled 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, ¤tMap); 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 currentTaskAttempt(metadata json.RawMessage) int { if len(metadata) == 0 || string(metadata) == "null" { return 0 } var currentMap map[string]any if err := json.Unmarshal(metadata, ¤tMap); err != nil { return 0 } if currentMap != nil { if val, ok := currentMap[MetadataKeyAttempt]; ok { switch v := val.(type) { case float64: return int(v) case int: return v case int64: return int(v) } } } return 0 } 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, ¤tMap); 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) QueueTask(ctx context.Context, id string) (storage.Task, error) { task, err := l.store.GetTask(ctx, id) if err != nil { return storage.Task{}, err } if !canEnqueueStatus(TaskStatus(task.Status)) { return storage.Task{}, ErrTaskCannotBeEnqueued } return l.store.UpdateStatus(ctx, id, string(StatusQueued)) } 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, MetadataKeyAuthoringFailureCategory: nil, MetadataKeyAuthoringFailureType: 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) } // Promote authoring run state keys from the result payload into // task metadata so the lifecycle layer is the single writer. for _, key := range []string{ MetadataKeyAuthoringRunState, MetadataKeyAuthoringRunUpdatedAt, } { if v, ok := resMap[key]; ok && v != nil { updates[key] = v } } } } 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) } type FailureInput struct { Message string `json:"message"` Type FailureType `json:"type"` ExtraMetadata map[string]any `json:"extra_metadata,omitempty"` } func (l *Lifecycle) FailTaskWithMetadata(ctx context.Context, id string, input FailureInput) (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 } msg := strings.TrimSpace(input.Message) if msg == "" { msg = "unknown failure" } failType := input.Type if failType == "" { failType = FailureTypeExecution } attempt := currentTaskAttempt(task.Metadata) if attempt == 0 { attempt = 1 } retryable := attempt < DefaultTaskMaxAttempts updates := map[string]any{ MetadataKeyAgentRunState: "failed", MetadataKeyStatusReason: msg, MetadataKeyWaitType: nil, MetadataKeyFailureType: string(failType), MetadataKeyFailedAt: time.Now().UTC().Format(time.RFC3339), MetadataKeyRetryable: retryable, } for k, v := range input.ExtraMetadata { updates[k] = v } 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, msg) } func (l *Lifecycle) FailTask(ctx context.Context, id string, message string) (storage.Task, error) { return l.FailTaskWithMetadata(ctx, id, FailureInput{ Message: message, Type: FailureTypeExecution, }) } // MergeTaskMetadata merges the given updates into the task's current metadata // without touching task status. Used by the scheduler to record authoring-specific // state (e.g. authoring_run_state=in_progress) outside of status transitions. func (l *Lifecycle) MergeTaskMetadata(ctx context.Context, id string, updates map[string]any) (storage.Task, error) { task, err := l.store.GetTask(ctx, id) if err != nil { return storage.Task{}, err } if terminalTaskStatus(TaskStatus(task.Status)) { return task, nil } newMeta, err := mergeTaskMetadata(task.Metadata, updates) if err != nil { return storage.Task{}, err } // Apply the write under an atomic terminal guard. If a finalizer transitions // the task to a terminal status between the GetTask above and this update, the // conditional update matches no row and we must not overwrite terminal // metadata. In that case re-read and return the latest terminal task as a // no-op rather than resurrecting stale in_progress/develop_match state. updated, ok, err := l.store.UpdateMetadataIfNonTerminal(ctx, id, newMeta) if err != nil { return storage.Task{}, err } if !ok { return l.store.GetTask(ctx, id) } return updated, nil } 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)) }