Files
Test 11a1c1196a Add testable task descriptions: T0.1 through T0.9
Each task includes:
- Scope: what to build
- Implementation: code sketches + details
- Verification: concrete test criteria
- Done criteria: acceptance checklist
2026-08-20 22:41:54 -07:00

6.3 KiB

T0.7: Orchestrator Workflow

Scope

Implement statemachine/orchestrator.go with config state, signals, fan-out/fan-in, continue-as-new, and 504 learning.

Implementation

File: statemachine/orchestrator.go

func OrchestratorWorkflow(ctx workflow.Context, in OrchestratorInput) (OrchestratorOutput, error)
    // 1. Mutable config state (not frozen at start):
    //    config := in.Config
    //    skillsHaveChanged := true // first cycle, prep skills
    //
    // 2. Signal handlers (checked each cycle):
    //    - "update-system-prompt": config.SystemPrompt = signalPayload
    //    - "update-skills": config.Skills = signalPayload, skillsHaveChanged = true
    //    - "update-role-prompt": config.RolePrompts[role] = signalPayload
    //    - "update-tuning": config.Tuning = signalPayload
    //    - "pause": wait for "resume" signal
    //    - "abort-task": forward via SignalExternalWorkflow(ctx, "taskunit-"+taskID, "abort", nil)
    //
    // 3. Query handlers:
    //    - "status": return current cycle count, pending tasks
    //    - "current-config": return config
    //
    // 4. Main loop (continues until submilestone complete):
    //    for {
    //        // Check signals (pause, abort, update-*)
    //        selector := workflow.NewSelector(ctx)
    //        // register signal channels
    //
    //        // Prep skills if needed
    //        if skillsHaveChanged {
    //            call PrepareSkillsActivity(ctx, {config.Skills, config.Tuning.PiRetry.StreamTimeout})
    //                wrap in 504-learning loop:
    //                    for {
    //                        err := ExecuteActivity(...)
    //                        if isPiStreamTimeout(err) && config.Tuning.PiRetry.StreamTimeout < config.Tuning.PiRetry.StreamTimeoutMax:
    //                            config.Tuning.PiRetry.StreamTimeout *= 2
    //                            continue
    //                        break
    //                    }
    //            skillsHaveChanged = false
    //        }
    //
    //        // Planning phase 1: decide what to dispatch
    //        planResult := call PlanningActivity(ctx, {config, boardState, milestone})
    //        if submilestoneComplete(planResult):
    //            // All subtasks done, trigger merge
    //            call GitSquashMergeActivity(ctx, {repoBranches, "T0: squash merge subtasks"})
    //            return OrchestratorOutput{MilestoneComplete: true, Done: true}
    //
    //        // Dispatch: fan out TaskUnitWorkflow for each task
    //        taskFutures := []workflow.Future{}
    //        for taskID in planResult.tasksToDispatch:
    //            spec := selectApplicableSpec(config.RolePrompts, taskID)
    //            future := ExecuteChildWorkflow(ctx, TaskUnitWorkflow, TaskUnitInput{
    //                TaskID: taskID,
    //                JudgeSpec: config.RolePrompts["judge"],
    //                ImplementerSpec: config.RolePrompts["implementer"],
    //                BaseTimeout: config.Tuning.ImplementerBaseTimeout, // or override from planner
    //                MaxJudgeRetries: config.Tuning.ImplementerMaxRetries,
    //            })
    //            taskFutures = append(taskFutures, future)
    //
    //        // Await all
    //        results := []TaskUnitOutput{}
    //        for future in taskFutures:
    //            var out TaskUnitOutput
    //            future.Get(ctx, &out)
    //            results = append(results, out)
    //
    //        // Planning phase 2: update board and commit
    //        call PlanningActivity(ctx, {config, results, boardState, milestone}) → UpdateBoardOutput
    //        call GitCommitActivity(ctx, {repoPath, "Update board after cycle"})
    //        call GitPushActivity(ctx, {repoPath})
    //
    //        // Continue-as-new check
    //        in.CycleCount++
    //        if in.CycleCount >= in.MaxCyclesBeforeCAN:
    //            nextInput := OrchestratorInput{
    //                // carry forward all state
    //                CycleCount: 0,
    //                Config: config, // includes mutated Tuning/RolePrompts/Skills
    //            }
    //            return workflow.NewContinueAsNewError(ctx, OrchestratorWorkflow, nextInput)
    //    }

504 Learning Detail

// Wrapping PrepareSkillsActivity for 504 learning:
for {
    r := config.Tuning.PiRetry
    ao := workflow.ActivityOptions{
        ScheduleToCloseTimeout: r.ScheduleToCloseTimeout, // 5m hard cap
        StartToCloseTimeout: r.MaximumInterval,           // per-attempt ceiling
        RetryPolicy: &temporal.RetryPolicy{
            InitialInterval: r.InitialInterval,
            BackoffCoefficient: r.BackoffCoefficient,
            MaximumInterval: r.MaximumInterval,
            NonRetryableErrorTypes: []string{"PiClientError"},
        },
    }
    err := workflow.ExecuteActivity(
        workflow.WithActivityOptions(ctx, ao),
        action.PrepareSkillsActivity,
        action.PrepareSkillsInput{Skills: config.Skills, StreamTimeout: r.StreamTimeout},
    ).Get(ctx, nil)

    var appErr *temporal.ApplicationError
    if errors.As(err, &appErr) && appErr.Type() == "PiStreamTimeout" && r.StreamTimeout < r.StreamTimeoutMax {
        config.Tuning.PiRetry.StreamTimeout = min(r.StreamTimeout*2, r.StreamTimeoutMax)
        continue // ScheduleToCloseTimeout still bounds each attempt
    }
    if err != nil {
        return OrchestratorOutput{}, err
    }
    break
}

Verification

cd /Users/rockliang/workplace/Poimen/workflows
go test -v ./tests -run TestOrchestrator

# Test file: tests/orchestrator_workflow_test.go

Test cases (mocked activities):

  1. Fan-out/fan-in: Dispatch 3 tasks → all complete → results collected
  2. Squash-merge on complete: All tasks pass → GitSquashMergeActivity called
  3. Continue-as-new: CycleCount reaches MaxCyclesBeforeCAN → returns NewContinueAsNewError
  4. Signal mutation: update-role-prompt signal → next dispatch uses new prompt
  5. 504 learning: PrepareSkillsActivity returns PiStreamTimeout → StreamTimeout doubled → next PrepareSkillsActivity call uses doubled value, capped at Max

Done Criteria

  • go test ./tests -run TestOrchestrator passes all 5 cases
  • Signals mutate config without affecting in-flight TaskUnit
  • Continue-as-new preserves OrchestratorConfig across cycles
  • 504 learning loop doesn't exceed ScheduleToCloseTimeout