# 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` ```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 ```go // 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 ```bash 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