diff --git a/action/integration_test.go b/action/integration.go similarity index 100% rename from action/integration_test.go rename to action/integration.go diff --git a/cmd/worker/main.go b/cmd/worker/main.go index 51d5861..0b525fa 100644 --- a/cmd/worker/main.go +++ b/cmd/worker/main.go @@ -64,7 +64,7 @@ func main() { w.RegisterActivity(action.ImplementerActivity) w.RegisterActivity(action.JudgeActivity) // Integration and lessons activities - register when fully tested - // w.RegisterActivity(action.RunIntegrationTestActivity) + w.RegisterActivity(action.RunIntegrationTestActivity) // w.RegisterActivity(action.UpdateLessonsActivity) // w.RegisterActivity(action.ReadLessonsActivity) diff --git a/go.mod b/go.mod index f6a6d0c..88f54e3 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( github.com/stretchr/testify v1.12.1 go.temporal.io/sdk v1.48.0 go.uber.org/zap v1.28.0 + gopkg.in/yaml.v3 v3.0.1 ) require ( @@ -38,5 +39,4 @@ require ( google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 // indirect google.golang.org/grpc v1.82.1 // indirect google.golang.org/protobuf v1.36.12 // indirect - gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index b604d06..dcbe0ad 100644 --- a/go.sum +++ b/go.sum @@ -26,6 +26,10 @@ github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk= github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= @@ -44,6 +48,8 @@ github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+ github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= github.com/robfig/cron v1.2.0 h1:ZjScXvvxeQ63Dbyxy76Fj3AT3Ut0aKsyd2/tl3DTMuQ= github.com/robfig/cron v1.2.0/go.mod h1:JGuDeoQd7Z6yL4zQhZ3OPEVHB7fL6Ka6skscFHfmt2k= +github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M= +github.com/rogpeppe/go-internal v1.11.0/go.mod h1:ddIwULY96R17DhadqLgMfk9H9tvdUzkipdSkR5nkCZA= github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= @@ -131,5 +137,7 @@ google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3 google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/k8s/git-commit.yaml b/k8s/git-commit.yaml index e54d371..cd7a07d 100644 --- a/k8s/git-commit.yaml +++ b/k8s/git-commit.yaml @@ -9,6 +9,6 @@ metadata: app.kubernetes.io/name: poimen app.kubernetes.io/component: orchestrator data: - GIT_COMMIT: "0b6bca7" # Updated automatically by CI/CD + GIT_COMMIT: "9fb58ae" # Updated automatically by CI/CD GIT_BRANCH: "main" DEPLOYMENT_DATE: "2026-08-26" diff --git a/k8s/worker-deployment.yaml b/k8s/worker-deployment.yaml index f5b9924..8b39db9 100644 --- a/k8s/worker-deployment.yaml +++ b/k8s/worker-deployment.yaml @@ -13,7 +13,7 @@ spec: labels: app: poimen-worker annotations: - git-commit: "0b6bca7" # ✅ Updated on each push, triggers rolling restart + git-commit: "9fb58ae" # ✅ Updated on each push, triggers rolling restart deployment-date: "2026-08-26" spec: containers: diff --git a/statemachine/orchestrator.go b/statemachine/orchestrator.go index 5e21440..e63da74 100644 --- a/statemachine/orchestrator.go +++ b/statemachine/orchestrator.go @@ -11,151 +11,168 @@ import ( ) // OrchestratorWorkflow orchestrates multi-agent work on a target repository. +// Reconciliation loop: Plan → Dispatch → Review → Update → Repeat func OrchestratorWorkflow(ctx workflow.Context, in OrchestratorInput) (OrchestratorOutput, error) { + logger := workflow.GetLogger(ctx) output := OrchestratorOutput{ MilestoneComplete: false, Done: false, LastError: "", } - // Step 1: Clone the repository - activityOptions := workflow.ActivityOptions{ - StartToCloseTimeout: 10 * time.Minute, + // Clone repo once at start + activityOpts := workflow.ActivityOptions{ + StartToCloseTimeout: 10 * time.Minute, ScheduleToCloseTimeout: 15 * time.Minute, } - ctxWithOptions := workflow.WithActivityOptions(ctx, activityOptions) - - cloneErr := workflow.ExecuteActivity( - ctxWithOptions, - "CloneRepoActivity", - map[string]interface{}{ - "RemoteURL": in.RemoteURL, - "TargetRepoPath": in.TargetRepoPath, - }, - ).Get(ctx, nil) + ctxWithOpts := workflow.WithActivityOptions(ctx, activityOpts) - if cloneErr != nil { - output.LastError = fmt.Sprintf("Clone failed: %v", cloneErr) + if err := workflow.ExecuteActivity(ctxWithOpts, "CloneRepoActivity", map[string]interface{}{ + "RemoteURL": in.RemoteURL, + "TargetRepoPath": in.TargetRepoPath, + }).Get(ctx, nil); err != nil { + output.LastError = fmt.Sprintf("Clone failed: %v", err) return output, nil } - // Step 2: Read tasks from board.md - tasksToRun, err := readTasksFromBoard(in.TargetRepoPath) - if err != nil { - output.LastError = fmt.Sprintf("Failed to read tasks: %v", err) - return output, nil - } - - if len(tasksToRun) == 0 { - output.LastError = "No tasks found in board.md" - return output, nil - } - - // Step 3: Process each task - completedTasks := 0 - for _, task := range tasksToRun { - taskID := task["id"].(string) - taskDesc := task["description"].(string) - - // taskID will be used for worktree and branch - - // Add worktree - var worktreePath string - wtErr := workflow.ExecuteActivity( - ctxWithOptions, - "GitWorktreeAddActivity", - map[string]interface{}{ - "RepoPath": in.TargetRepoPath, - "TaskID": taskID, - }, - ).Get(ctx, &worktreePath) - - if wtErr != nil { - continue // Skip this task on error + // Prepare skills once + if len(in.Config.Skills) > 0 { + if err := workflow.ExecuteActivity(ctxWithOpts, "PrepareSkillsActivity", map[string]interface{}{ + "Skills": in.Config.Skills, + "StreamTimeout": 30 * time.Second, + "Provider": in.PiProvider, + }).Get(ctx, nil); err != nil { + logger.Info("skill preparation failed (continuing anyway)", "error", err) } + } - // Call implementer to generate code (longer timeout for LLM calls) - implOptions := workflow.ActivityOptions{ - StartToCloseTimeout: 30 * time.Minute, + // Main reconciliation loop + var paused bool + var completedTasks int + cycleCount := 0 + completedBranches := []string{} + + for cycleCount < in.MaxCyclesBeforeCAN { + cycleCount++ + logger.Info("orchestrator cycle", "cycle", cycleCount) + + // Read board state + bordState, _ := readTasksFromBoard(in.TargetRepoPath) + + // Call PlanningActivity to decide what to dispatch + implOpts := workflow.ActivityOptions{ + StartToCloseTimeout: 30 * time.Minute, ScheduleToCloseTimeout: 35 * time.Minute, } - implCtx := workflow.WithActivityOptions(ctx, implOptions) - - var implOutput map[string]interface{} - implErr := workflow.ExecuteActivity( - implCtx, - "ImplementerActivity", - map[string]interface{}{ - "TaskID": taskID, - "Description": taskDesc, - "WorktreePath": worktreePath, - "Prompt": PromptSpec{ - TemplateRef: "implementer/default.tmpl", - Model: ModelSpec{ - ModelID: in.Config.RolePrompts["implementer"].Model.ModelID, - }, - }, - }, - ).Get(ctx, &implOutput) + implCtx := workflow.WithActivityOptions(ctx, implOpts) - if implErr != nil { - continue + var planOutput map[string]interface{} + if err := workflow.ExecuteActivity(implCtx, "PlanningActivity", map[string]interface{}{ + "Config": in.Config, + "BoardState": fmt.Sprintf("%v", bordState), + "RepoPath": in.TargetRepoPath, + "Milestone": in.Milestone, + }).Get(ctx, &planOutput); err != nil { + logger.Info("planning failed", "error", err) + break } - // Commit changes - commitErr := workflow.ExecuteActivity( - ctxWithOptions, - "GitCommitActivity", - map[string]interface{}{ - "WorktreePath": worktreePath, - "Message": fmt.Sprintf("%s: implementation", taskID), - }, - ).Get(ctx, nil) + // Extract tasks to dispatch + var tasksToDispatch []interface{} + if tasks, ok := planOutput["TasksToDispatch"]; ok { + tasksToDispatch = tasks.([]interface{}) + } - if commitErr == nil { - completedTasks++ + if len(tasksToDispatch) == 0 { + logger.Info("no tasks to dispatch, milestone complete") + output.MilestoneComplete = true + break + } + + // Fan-out: Start TaskUnit workflows for each task + logger.Info("dispatching task units", "count", len(tasksToDispatch)) + var childWFs []workflow.Future + + for _, taskIDRaw := range tasksToDispatch { + taskID := taskIDRaw.(string) + + if paused { + logger.Info("skipping dispatch due to pause signal", "task", taskID) + continue + } + + // Child workflow: TaskUnitWorkflow + childOpts := workflow.ChildWorkflowOptions{ + WorkflowID: fmt.Sprintf("%s-%s-cycle%d", in.Milestone, taskID, cycleCount), + } + childCtx := workflow.WithChildOptions(ctx, childOpts) + + taskUnitInput := TaskUnitInput{ + TaskID: taskID, + RemoteURL: in.RemoteURL, + TargetRepoPath: in.TargetRepoPath, + Milestone: in.Milestone, + Config: in.Config, + DryRun: in.DryRun, + } + + future := workflow.ExecuteChildWorkflow(childCtx, TaskUnitWorkflow, taskUnitInput) + childWFs = append(childWFs, future) + } + + // Fan-in: Wait for all TaskUnits to complete + logger.Info("waiting for task units", "count", len(childWFs)) + for _, future := range childWFs { + var taskOutput TaskUnitOutput + if err := future.Get(ctx, &taskOutput); err != nil { + logger.Info("task unit failed", "task", taskOutput.TaskID, "error", err) + } else if taskOutput.Status == "success" { + completedTasks++ + completedBranches = append(completedBranches, taskOutput.Branch) + logger.Info("task unit succeeded", "task", taskOutput.TaskID) + } + } + + // Board update: Call PlanningActivity again to update board + commit + if err := workflow.ExecuteActivity(implCtx, "PlanningActivity", map[string]interface{}{ + "Config": in.Config, + "BoardState": fmt.Sprintf("%v", bordState), + "RepoPath": in.TargetRepoPath, + "Milestone": in.Milestone, + }).Get(ctx, nil); err != nil { + logger.Info("board update failed", "error", err) + } + + // Push completed branches + if err := workflow.ExecuteActivity(ctxWithOpts, "GitPushActivity", map[string]interface{}{ + "RepoPath": in.TargetRepoPath, + }).Get(ctx, nil); err != nil { + logger.Info("push failed", "error", err) + } + + // Squash merge completed branches when milestone ready + if output.MilestoneComplete && len(completedBranches) > 0 { + if err := workflow.ExecuteActivity(ctxWithOpts, "GitSquashMergeActivity", map[string]interface{}{ + "RepoPath": in.TargetRepoPath, + "Branches": completedBranches, + "Message": fmt.Sprintf("%s: squash merge completed tasks", in.Milestone), + }).Get(ctx, nil); err != nil { + logger.Info("squash merge failed", "error", err) + } } } - // Step 4: Push to remote - pushErr := workflow.ExecuteActivity( - ctxWithOptions, - "GitPushActivity", - map[string]interface{}{ - "RepoPath": in.TargetRepoPath, - }, - ).Get(ctx, nil) - - if pushErr != nil { - output.LastError = fmt.Sprintf("Push failed: %v", pushErr) - return output, nil + // Use continue-as-new if hit cycle cap + if cycleCount >= in.MaxCyclesBeforeCAN { + logger.Info("cycle cap reached, continuing as new", "cycles", cycleCount) + nextInput := in + nextInput.CycleCount = cycleCount + return output, workflow.NewContinueAsNewError(ctx, OrchestratorWorkflow, nextInput) } - // Step 5: Squash merge all task branches - branches := make([]string, len(tasksToRun)) - for i, task := range tasksToRun { - branches[i] = fmt.Sprintf("task/%s", task["id"].(string)) - } - - mergeErr := workflow.ExecuteActivity( - ctxWithOptions, - "GitSquashMergeActivity", - map[string]interface{}{ - "RepoPath": in.TargetRepoPath, - "Branches": branches, - "Message": fmt.Sprintf("%s: squash merge all tasks", in.Milestone), - }, - ).Get(ctx, nil) - - if mergeErr != nil { - output.LastError = fmt.Sprintf("Merge failed: %v", mergeErr) - return output, nil - } - - // Success! - output.MilestoneComplete = true + // Success output.Done = true - output.LastError = fmt.Sprintf("Completed %d tasks successfully", completedTasks) + output.LastError = fmt.Sprintf("Completed %d tasks in %d cycles", completedTasks, cycleCount) return output, nil } @@ -192,11 +209,3 @@ func readTasksFromBoard(repoPath string) ([]map[string]interface{}, error) { return tasks, nil } - -// isPiStreamTimeout checks if an error is a 504 stream timeout from Pi command -func isPiStreamTimeout(err error) bool { - if err == nil { - return false - } - return strings.Contains(err.Error(), "PiStreamTimeout") -} diff --git a/statemachine/taskunit.go b/statemachine/taskunit.go index ebad9cf..309ec1a 100644 --- a/statemachine/taskunit.go +++ b/statemachine/taskunit.go @@ -4,161 +4,134 @@ import ( "fmt" "time" - "go.temporal.io/sdk/temporal" "go.temporal.io/sdk/workflow" ) -// TaskUnitWorkflow executes a single task with retry loops, timeout escalation, and lessons injection. +// TaskUnitWorkflow executes a single task with retries and judge review. +// Flow: Worktree → Implementer (retry) → Test → Judge → Commit or Retry func TaskUnitWorkflow(ctx workflow.Context, in TaskUnitInput) (TaskUnitOutput, error) { - // Initialize output + logger := workflow.GetLogger(ctx) output := TaskUnitOutput{ TaskID: in.TaskID, - Verdict: "fail", + Status: "failed", + Reason: "", + Branch: fmt.Sprintf("task/%s", in.TaskID), + Changes: "", } - // 1. Add worktree for isolated work + activityOpts := workflow.ActivityOptions{ + StartToCloseTimeout: 30 * time.Minute, + ScheduleToCloseTimeout: 35 * time.Minute, + } + actCtx := workflow.WithActivityOptions(ctx, activityOpts) + + // Add worktree var worktreePath string - wtErr := workflow.ExecuteActivity( - ctx, - "GitWorktreeAddActivity", - map[string]interface{}{ - "RepoPath": in.TargetRepoPath, - "TaskID": in.TaskID, - }, - ).Get(ctx, &worktreePath) - if wtErr != nil { - output.Critique = fmt.Sprintf("Failed to create worktree: %v", wtErr) + if err := workflow.ExecuteActivity(actCtx, "GitWorktreeAddActivity", map[string]interface{}{ + "RepoPath": in.TargetRepoPath, + "TaskID": in.TaskID, + }).Get(ctx, &worktreePath); err != nil { + output.Reason = fmt.Sprintf("worktree add failed: %v", err) return output, nil } - // 2. Retry loop with separate timeout and judge attempt tracking - timeoutAttempt := 1 - for judgeAttempt := 1; judgeAttempt <= in.MaxJudgeRetries; judgeAttempt++ { - // Calculate timeouts for this attempt - baseTimeout := in.BaseTimeout * time.Duration(timeoutAttempt) - heartbeatTimeout := baseTimeout / 4 + logger.Info("task unit started", "task", in.TaskID, "worktree", worktreePath) - // Prepare activity options with escalating timeout - ao := workflow.ActivityOptions{ - ScheduleToCloseTimeout: baseTimeout, - StartToCloseTimeout: baseTimeout, - HeartbeatTimeout: heartbeatTimeout, - RetryPolicy: &temporal.RetryPolicy{ - MaximumAttempts: 1, // We manage retries in this loop - }, + // Retry loop: implementer + test + judge + maxRetries := 3 + var lessons string + + for attempt := 1; attempt <= maxRetries; attempt++ { + logger.Info("attempt", "task", in.TaskID, "attempt", attempt) + + // Call ImplementerActivity with escalating timeout + timeoutMultiplier := int64(attempt) + implOpts := workflow.ActivityOptions{ + StartToCloseTimeout: time.Duration(timeoutMultiplier*30) * time.Minute, + ScheduleToCloseTimeout: time.Duration(timeoutMultiplier*35) * time.Minute, } - ctxWithOptions := workflow.WithActivityOptions(ctx, ao) + implCtx := workflow.WithActivityOptions(ctx, implOpts) - // Call implementer activity var implOutput map[string]interface{} - implErr := workflow.ExecuteActivity( - ctxWithOptions, - "ImplementerActivity", - map[string]interface{}{ - "TaskID": in.TaskID, - "WorktreePath": worktreePath, - "Prompt": in.ImplementerSpec, - }, - ).Get(ctx, &implOutput) + implErr := workflow.ExecuteActivity(implCtx, "ImplementerActivity", map[string]interface{}{ + "Config": in.Config, + "TaskID": in.TaskID, + "WorktreePath": worktreePath, + "Lessons": lessons, + }).Get(ctx, &implOutput) - // Check if it's a timeout error - if implErr != nil && isStartToCloseTimeout(implErr) { - // Timeout: escalate and retry without consuming judge attempt - timeoutAttempt++ - judgeAttempt-- // Don't consume a judge retry on timeout - continue - } if implErr != nil { - output.Critique = fmt.Sprintf("Implementer failed: %v", implErr) + if attempt < maxRetries { + logger.Info("implementer failed, will retry", "task", in.TaskID, "attempt", attempt, "error", implErr) + continue + } + output.Reason = fmt.Sprintf("implementer exhausted after %d attempts: %v", maxRetries, implErr) return output, nil } - // Call judge activity - judgeTimeout := time.Minute * 5 - judgeCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - ScheduleToCloseTimeout: judgeTimeout, - StartToCloseTimeout: judgeTimeout, - }) - var judgeOutput map[string]interface{} - judgeErr := workflow.ExecuteActivity( - judgeCtx, - "JudgeActivity", - map[string]interface{}{ - "TaskID": in.TaskID, - "WorktreePath": worktreePath, - "Prompt": in.JudgeSpec, - }, - ).Get(ctx, &judgeOutput) - if judgeErr != nil { - output.Critique = fmt.Sprintf("Judge error: %v", judgeErr) - return output, nil - } - - // Check judge verdict - verdict := "" - if judgeOutput != nil { - if v, ok := judgeOutput["Verdict"].(string); ok { - verdict = v + // Run integration test + var testOutput map[string]interface{} + if err := workflow.ExecuteActivity(actCtx, "RunIntegrationTestActivity", map[string]interface{}{ + "WorktreePath": worktreePath, + "TestCmd": "go test ./...", + }).Get(ctx, &testOutput); err != nil { + logger.Info("test failed", "task", in.TaskID, "error", err) + testOutput = map[string]interface{}{ + "success": false, + "logs": fmt.Sprintf("test error: %v", err), } } + // Get diff + var diff string + if err := workflow.ExecuteActivity(actCtx, "GitDiffActivity", map[string]interface{}{ + "WorktreePath": worktreePath, + }).Get(ctx, &diff); err != nil { + logger.Info("diff failed", "task", in.TaskID, "error", err) + } + + // Call JudgeActivity + var judgeOutput map[string]interface{} + if err := workflow.ExecuteActivity(actCtx, "JudgeActivity", map[string]interface{}{ + "Config": in.Config, + "Diff": diff, + "IntegrationTestLogs": fmt.Sprintf("%v", testOutput), + }).Get(ctx, &judgeOutput); err != nil { + logger.Info("judge failed", "task", in.TaskID, "error", err) + } + + verdict, _ := judgeOutput["Verdict"].(string) + critique, _ := judgeOutput["Critique"].(string) + if verdict == "pass" { - // Commit in worktree - commitErr := workflow.ExecuteActivity( - ctx, - "GitCommitActivity", - map[string]interface{}{ - "WorktreePath": worktreePath, - "Message": fmt.Sprintf("%s: implementation", in.TaskID), - }, - ).Get(ctx, nil) - if commitErr != nil { - output.Critique = fmt.Sprintf("Commit failed: %v", commitErr) + // Commit and success + if err := workflow.ExecuteActivity(actCtx, "GitCommitActivity", map[string]interface{}{ + "WorktreePath": worktreePath, + "Message": fmt.Sprintf("%s: implementation", in.TaskID), + }).Get(ctx, nil); err != nil { + output.Reason = fmt.Sprintf("commit failed: %v", err) return output, nil } - // Success! - output.Verdict = "pass" - output.Branch = "task/" + in.TaskID + output.Status = "success" + output.Changes = fmt.Sprintf("completed in %d attempt(s)", attempt) + logger.Info("task success", "task", in.TaskID, "attempt", attempt) return output, nil } - // Judge failed: update lessons and retry - critique := "" - if judgeOutput != nil { - if c, ok := judgeOutput["Critique"].(string); ok { - critique = c - } + // Judge failed: append to lessons and retry + if attempt < maxRetries { + lessons = fmt.Sprintf("%s\nAttempt %d critique: %s", lessons, attempt, critique) + logger.Info("judge rejected, appending to lessons and retrying", "task", in.TaskID, "attempt", attempt) + continue } - updateErr := workflow.ExecuteActivity( - ctx, - "UpdateLessonsActivity", - map[string]interface{}{ - "TargetRepoPath": in.TargetRepoPath, - "TaskID": in.TaskID, - "Attempt": judgeAttempt, - "Critique": critique, - }, - ).Get(ctx, nil) - if updateErr != nil { - output.Critique = fmt.Sprintf("Failed to update lessons: %v", updateErr) - return output, nil - } - - // Continue to next judge attempt with lessons injected + // Exhausted retries after judge failures + output.Reason = fmt.Sprintf("judge rejected after %d attempts. Last critique: %s", maxRetries, critique) + logger.Info("judge exhausted retries", "task", in.TaskID) + return output, nil } - // Retries exhausted - output.Verdict = "fail" - output.Critique = fmt.Sprintf("Exhausted %d judge retries", in.MaxJudgeRetries) + output.Reason = "exhausted all retry attempts" return output, nil } - -// isStartToCloseTimeout checks if an error is a StartToCloseTimeout error -func isStartToCloseTimeout(err error) bool { - if err == nil { - return false - } - return fmt.Sprint(err) == "context deadline exceeded" -} diff --git a/statemachine/types.go b/statemachine/types.go index 93f3b05..cf8b403 100644 --- a/statemachine/types.go +++ b/statemachine/types.go @@ -69,20 +69,23 @@ type OrchestratorOutput struct { // TaskUnitInput is the input to the TaskUnit workflow. type TaskUnitInput struct { - TaskID string - TargetRepoPath string - JudgeSpec PromptSpec - ImplementerSpec PromptSpec - BaseTimeout time.Duration - MaxJudgeRetries int + TaskID string + RemoteURL string + TargetRepoPath string + Milestone string + Config OrchestratorConfig + DryRun bool } // TaskUnitOutput is the output of the TaskUnit workflow. type TaskUnitOutput struct { TaskID string - Verdict string // "pass" or "fail" - Critique string + Status string // "success" or "failed" + Verdict string // "pass" or "fail" from judge + Critique string // feedback from judge Branch string + Reason string // error reason if failed + Changes string // summary of changes } // SkillRef references a skill source.