Files
Test 6a87833c7f
ci / test (push) Failing after 1m2s
feat: implement proper orchestrator workflow with reconciliation loop
Rewrite OrchestratorWorkflow as true reconciliation loop:
- PlanningActivity decides what tasks to dispatch
- Fan-out TaskUnit workflows for parallel execution
- Each TaskUnit runs Implementer → Test → Judge → Commit
- Judge reviews code quality, retries on failure with lessons
- Fan-in waits for all TaskUnits
- Board update and squash merge on success
- continue-as-new for long-running workflows
- Proper error handling and signal support

Key changes:
- statemachine/orchestrator.go: Reconciliation loop (Plan → Dispatch → Review → Repeat)
- statemachine/taskunit.go: Task execution with retry loop & judge review
- statemachine/types.go: Updated TaskUnitInput/Output for new workflow
- cmd/worker/main.go: Register RunIntegrationTestActivity
- action/integration.go: Renamed from integration_test.go (fix Go build issue)

Models:
- Planner: reasoning (OpenAI-compatible from local LLM API)
- Judge: reasoning (reviews diff + tests, gates success)
- Implementer: ornith:35b (executes tasks)

Verification: go build ./cmd/worker ./cmd/starter ✓
2026-08-26 15:00:42 -07:00

138 lines
4.5 KiB
Go

package statemachine
import (
"fmt"
"time"
"go.temporal.io/sdk/workflow"
)
// 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) {
logger := workflow.GetLogger(ctx)
output := TaskUnitOutput{
TaskID: in.TaskID,
Status: "failed",
Reason: "",
Branch: fmt.Sprintf("task/%s", in.TaskID),
Changes: "",
}
activityOpts := workflow.ActivityOptions{
StartToCloseTimeout: 30 * time.Minute,
ScheduleToCloseTimeout: 35 * time.Minute,
}
actCtx := workflow.WithActivityOptions(ctx, activityOpts)
// Add worktree
var worktreePath string
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
}
logger.Info("task unit started", "task", in.TaskID, "worktree", worktreePath)
// 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,
}
implCtx := workflow.WithActivityOptions(ctx, implOpts)
var implOutput map[string]interface{}
implErr := workflow.ExecuteActivity(implCtx, "ImplementerActivity", map[string]interface{}{
"Config": in.Config,
"TaskID": in.TaskID,
"WorktreePath": worktreePath,
"Lessons": lessons,
}).Get(ctx, &implOutput)
if implErr != nil {
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
}
// 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 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
}
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: 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
}
// 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
}
output.Reason = "exhausted all retry attempts"
return output, nil
}