feat(team): refine configuration, fix critical bugs, and update documentation

- Fixed dead code in token budget clamping
- Fixed goroutine leak in DAG strategy fast-fail path
- Optimized Evaluator to run without tools in evaluator_optimizer
- Added partial success support to Parallel strategy
- Implemented configurable context truncation via max_context_runes
- Added reviewer_model support for cost-efficient QA
- Added structured logging across all strategies
- Improved ForUser messages with role summaries
- Merged and updated team configuration documentation
This commit is contained in:
Administrator 2026-03-12 10:44:02 +08:00
parent bccdc92ffb
commit f12638ce14
3 changed files with 187 additions and 65 deletions

View file

@ -21,6 +21,9 @@ PicoClaw's tools configuration is located in the `tools` field of `config.json`.
}, },
"skills": { "skills": {
... ...
},
"team": {
...
} }
} }
} }
@ -302,6 +305,44 @@ The skills tool configures skill discovery and installation via registries like
} }
``` ```
## Team Tool
The team tool allows PicoClaw to orchestrate multiple agents using strategies like Sequential, Parallel, DAG, and Evaluator-Optimizer.
### Global Config
| Config | Type | Default | Description |
| :--- | :--- | :--- | :--- |
| `enabled` | bool | false | Enable the Team tool |
| `max_members` | int | 5 | Max agents in a single team |
| `max_team_tokens` | int | 0 | Hard ceiling for total token usage per team call (0 = no limit) |
| `max_evaluator_loops`| int | 5 | Max retries for `evaluator_optimizer` strategy |
| `max_timeout_minutes`| int | 15 | Max execution time for team tasks |
| `max_context_runes` | int | 8000 | Max dependency context size between workers |
| `disable_auto_reviewer` | bool | false | Skip automatic QA reviewer step |
| `reviewer_model` | string | - | Specifically assigned model for the reviewer step |
| `allowed_strategies` | array | [] | Whitelist of allowed strategies (empty = all) |
| `allowed_models` | array | [] | Whitelist of models team members can use |
### Configuration Example
```json
{
"tools": {
"team": {
"enabled": true,
"max_members": 5,
"max_team_tokens": 100000,
"reviewer_model": "gpt-4o-mini",
"allowed_models": [
{ "name": "gpt-4o", "tags": ["vision", "code"] },
{ "name": "claude-3-5-sonnet", "tags": ["precise"] }
]
}
}
}
```
## Environment Variables ## Environment Variables
All configuration options can be overridden via environment variables with the format `PICOCLAW_TOOLS_<SECTION>_<KEY>`: All configuration options can be overridden via environment variables with the format `PICOCLAW_TOOLS_<SECTION>_<KEY>`:
@ -312,6 +353,6 @@ For example:
- `PICOCLAW_TOOLS_EXEC_ENABLE_DENY_PATTERNS=false` - `PICOCLAW_TOOLS_EXEC_ENABLE_DENY_PATTERNS=false`
- `PICOCLAW_TOOLS_CRON_EXEC_TIMEOUT_MINUTES=10` - `PICOCLAW_TOOLS_CRON_EXEC_TIMEOUT_MINUTES=10`
- `PICOCLAW_TOOLS_MCP_ENABLED=true` - `PICOCLAW_TOOLS_MCP_ENABLED=true`
- `PICOCLAW_TOOLS_TEAM_MAX_TEAM_TOKENS=50000`
Note: Nested map-style config (for example `tools.mcp.servers.<name>.*`) is configured in `config.json` rather than Note: Nested map-style config (for example `tools.mcp.servers.<name>.*` or `tools.team.allowed_models`) is configured in `config.json` rather than environment variables.
environment variables.

View file

@ -85,7 +85,9 @@ type TeamToolsConfig struct {
MaxTeamTokens int `json:"max_team_tokens" env:"PICOCLAW_TOOLS_TEAM_MAX_TOKENS"` MaxTeamTokens int `json:"max_team_tokens" env:"PICOCLAW_TOOLS_TEAM_MAX_TOKENS"`
MaxEvaluatorLoops int `json:"max_evaluator_loops" env:"PICOCLAW_TOOLS_TEAM_MAX_EVALUATOR_LOOPS"` MaxEvaluatorLoops int `json:"max_evaluator_loops" env:"PICOCLAW_TOOLS_TEAM_MAX_EVALUATOR_LOOPS"`
MaxTimeoutMinutes int `json:"max_timeout_minutes" env:"PICOCLAW_TOOLS_TEAM_MAX_TIMEOUT_MINUTES"` MaxTimeoutMinutes int `json:"max_timeout_minutes" env:"PICOCLAW_TOOLS_TEAM_MAX_TIMEOUT_MINUTES"`
MaxContextRunes int `json:"max_context_runes" env:"PICOCLAW_TOOLS_TEAM_MAX_CONTEXT_RUNES"`
DisableAutoReviewer bool `json:"disable_auto_reviewer" env:"PICOCLAW_TOOLS_TEAM_DISABLE_AUTO_REVIEWER"` DisableAutoReviewer bool `json:"disable_auto_reviewer" env:"PICOCLAW_TOOLS_TEAM_DISABLE_AUTO_REVIEWER"`
ReviewerModel string `json:"reviewer_model" env:"PICOCLAW_TOOLS_TEAM_REVIEWER_MODEL"`
AllowedStrategies []string `json:"allowed_strategies" env:"PICOCLAW_TOOLS_TEAM_ALLOWED_STRATEGIES"` AllowedStrategies []string `json:"allowed_strategies" env:"PICOCLAW_TOOLS_TEAM_ALLOWED_STRATEGIES"`
AllowedModels []TeamModelConfig `json:"allowed_models" env:"-"` AllowedModels []TeamModelConfig `json:"allowed_models" env:"-"`
} }

View file

@ -8,6 +8,7 @@ import (
"sync/atomic" "sync/atomic"
"time" "time"
"github.com/sipeed/picoclaw/pkg/logger"
"github.com/sipeed/picoclaw/pkg/providers" "github.com/sipeed/picoclaw/pkg/providers"
) )
@ -167,13 +168,12 @@ func (t *TeamTool) maybeRunAutoReviewer(
return "" // Unknown produces types, skip return "" // Unknown produces types, skip
} }
// 3. Skip Auto-Reviewer if disabled by config
sm := t.manager sm := t.manager
sm.mu.RLock() sm.mu.RLock()
disabled := sm.teamConfig.DisableAutoReviewer teamConfig := sm.teamConfig
sm.mu.RUnlock() sm.mu.RUnlock()
if disabled { if teamConfig.DisableAutoReviewer {
return "" return ""
} }
@ -184,7 +184,13 @@ func (t *TeamTool) maybeRunAutoReviewer(
{Role: "user", Content: reviewerTask}, {Role: "user", Content: reviewerTask},
} }
// Use a dedicated reviewer model if configured — typically a cheaper/faster model
// is sufficient for QA review, saving tokens compared to the main worker model.
reviewerConfig := baseConfig reviewerConfig := baseConfig
if teamConfig.ReviewerModel != "" && sm.IsModelAllowed(teamConfig.ReviewerModel) {
reviewerConfig.Model = teamConfig.ReviewerModel
}
loopResult, err := RunToolLoop(ctx, reviewerConfig, reviewerMessages, t.originChannel, t.originChatID) loopResult, err := RunToolLoop(ctx, reviewerConfig, reviewerMessages, t.originChannel, t.originChatID)
if err != nil { if err != nil {
return fmt.Sprintf("[Auto-Reviewer] Failed to run: %v", err) return fmt.Sprintf("[Auto-Reviewer] Failed to run: %v", err)
@ -239,7 +245,7 @@ func (t *TeamTool) Execute(ctx context.Context, args map[string]any) *ToolResult
maxTokensFloat, ok := args["max_team_tokens"].(float64) maxTokensFloat, ok := args["max_team_tokens"].(float64)
// Enforce hard budgets from config // Enforce hard budget from config as the ceiling.
effectiveMaxTokens := int64(0) effectiveMaxTokens := int64(0)
if teamConfig.MaxTeamTokens > 0 { if teamConfig.MaxTeamTokens > 0 {
effectiveMaxTokens = int64(teamConfig.MaxTeamTokens) effectiveMaxTokens = int64(teamConfig.MaxTeamTokens)
@ -247,13 +253,12 @@ func (t *TeamTool) Execute(ctx context.Context, args map[string]any) *ToolResult
if ok && maxTokensFloat > 0 { if ok && maxTokensFloat > 0 {
requestedTokens := int64(maxTokensFloat) requestedTokens := int64(maxTokensFloat)
// If LLM requested tokens but config enforces a smaller hard limit, clamp it
if effectiveMaxTokens > 0 && requestedTokens > effectiveMaxTokens { if effectiveMaxTokens > 0 && requestedTokens > effectiveMaxTokens {
effectiveMaxTokens = requestedTokens // LLM asked for more, but we clamp to config // LLM requested more than config allows: clamp to the hard ceiling.
// Wait, the clamping logic: if requested > max_team_tokens, clamp effectively shrinks it to max_team_tokens. // effectiveMaxTokens already holds the correct ceiling, no change needed.
effectiveMaxTokens = int64(teamConfig.MaxTeamTokens)
} else if effectiveMaxTokens == 0 || requestedTokens < effectiveMaxTokens { } else if effectiveMaxTokens == 0 || requestedTokens < effectiveMaxTokens {
effectiveMaxTokens = requestedTokens // LLM asked for less budget, let them be conservative // LLM asked for less, or there is no hard limit: honour the requested budget.
effectiveMaxTokens = requestedTokens
} }
} }
@ -315,7 +320,8 @@ func (t *TeamTool) Execute(ctx context.Context, args map[string]any) *ToolResult
baseConfig.RemainingTokenBudget = budget baseConfig.RemainingTokenBudget = budget
} }
// Create a new master context for team bounding // Create a cancellable context for team bounding.
// cancel() is always deferred so any spawned goroutines are cleaned up on return.
timeoutDur := 15 * time.Minute timeoutDur := 15 * time.Minute
if teamConfig.MaxTimeoutMinutes > 0 { if teamConfig.MaxTimeoutMinutes > 0 {
timeoutDur = time.Duration(teamConfig.MaxTimeoutMinutes) * time.Minute timeoutDur = time.Duration(teamConfig.MaxTimeoutMinutes) * time.Minute
@ -328,23 +334,29 @@ func (t *TeamTool) Execute(ctx context.Context, args map[string]any) *ToolResult
baseConfig.Tools = upgradeRegistryForConcurrency(baseConfig.Tools) baseConfig.Tools = upgradeRegistryForConcurrency(baseConfig.Tools)
} }
// Resolve max context runes for dependency injection into downstream prompts.
contextLimit := 8000 // default
if teamConfig.MaxContextRunes > 0 {
contextLimit = teamConfig.MaxContextRunes
}
switch strategy { switch strategy {
case "sequential": case "sequential":
result := t.executeSequential(teamCtx, baseConfig, members) result := t.executeSequential(teamCtx, baseConfig, members, contextLimit)
if reviewNote := t.maybeRunAutoReviewer(teamCtx, members, baseConfig, result.ForLLM); reviewNote != "" { if reviewNote := t.maybeRunAutoReviewer(teamCtx, members, baseConfig, result.ForLLM); reviewNote != "" {
result.ForLLM += "\n\n" + reviewNote result.ForLLM += "\n\n" + reviewNote
result.ForUser += "\n\n" + reviewNote result.ForUser += "\n\n" + reviewNote
} }
return result return result
case "dag": case "dag":
result := t.executeDAG(teamCtx, baseConfig, members) result := t.executeDAG(teamCtx, cancel, baseConfig, members, contextLimit)
if reviewNote := t.maybeRunAutoReviewer(teamCtx, members, baseConfig, result.ForLLM); reviewNote != "" { if reviewNote := t.maybeRunAutoReviewer(teamCtx, members, baseConfig, result.ForLLM); reviewNote != "" {
result.ForLLM += "\n\n" + reviewNote result.ForLLM += "\n\n" + reviewNote
result.ForUser += "\n\n" + reviewNote result.ForUser += "\n\n" + reviewNote
} }
return result return result
case "evaluator_optimizer": case "evaluator_optimizer":
return t.executeEvaluatorOptimizer(teamCtx, baseConfig, members) return t.executeEvaluatorOptimizer(teamCtx, baseConfig, members, contextLimit)
} }
// parallel // parallel
result := t.executeParallel(teamCtx, baseConfig, members) result := t.executeParallel(teamCtx, baseConfig, members)
@ -393,7 +405,7 @@ func buildWorkerConfig(baseConfig ToolLoopConfig, registry *ToolRegistry, m Team
return cfg, nil return cfg, nil
} }
func (t *TeamTool) executeSequential(ctx context.Context, baseConfig ToolLoopConfig, members []TeamMember) *ToolResult { func (t *TeamTool) executeSequential(ctx context.Context, baseConfig ToolLoopConfig, members []TeamMember, contextLimit int) *ToolResult {
var finalOutput strings.Builder var finalOutput strings.Builder
finalOutput.WriteString("Team Execution Summary (Sequential):\n\n") finalOutput.WriteString("Team Execution Summary (Sequential):\n\n")
@ -403,7 +415,7 @@ func (t *TeamTool) executeSequential(ctx context.Context, baseConfig ToolLoopCon
// If there is a previous result, we append it to the task so the new agent sees it. // If there is a previous result, we append it to the task so the new agent sees it.
actualTask := m.Task actualTask := m.Task
if i > 0 && previousResult != "" { if i > 0 && previousResult != "" {
actualTask = fmt.Sprintf("%s\n\n--- Context from previous phase ---\n%s", m.Task, truncateContext(previousResult)) actualTask = fmt.Sprintf("%s\n\n--- Context from previous phase ---\n%s", m.Task, truncateContextN(previousResult, contextLimit))
} }
messages := []providers.Message{ messages := []providers.Message{
@ -432,7 +444,7 @@ func (t *TeamTool) executeSequential(ctx context.Context, baseConfig ToolLoopCon
return &ToolResult{ return &ToolResult{
ForLLM: finalOutput.String(), ForLLM: finalOutput.String(),
ForUser: "Team completed sequential execution successfully.", ForUser: buildUserSummary("Sequential", members, nil),
} }
} }
@ -452,6 +464,11 @@ func (t *TeamTool) executeParallel(ctx context.Context, baseConfig ToolLoopConfi
go func(index int, member TeamMember) { go func(index int, member TeamMember) {
defer wg.Done() defer wg.Done()
logger.InfoCF("team", fmt.Sprintf("[%s] Parallel worker starting", member.Role), map[string]any{
"member_index": index,
"model": member.Model,
})
messages := []providers.Message{ messages := []providers.Message{
{Role: "system", Content: member.Role}, {Role: "system", Content: member.Role},
{Role: "user", Content: member.Task}, {Role: "user", Content: member.Task},
@ -470,6 +487,9 @@ func (t *TeamTool) executeParallel(ctx context.Context, baseConfig ToolLoopConfi
return return
} }
resultsChan <- workResult{index: index, role: member.Role, res: loopResult.Content} resultsChan <- workResult{index: index, role: member.Role, res: loopResult.Content}
logger.InfoCF("team", fmt.Sprintf("[%s] Parallel worker finished", member.Role), map[string]any{
"member_index": index,
})
}(i, m) }(i, m)
} }
@ -477,36 +497,54 @@ func (t *TeamTool) executeParallel(ctx context.Context, baseConfig ToolLoopConfi
wg.Wait() wg.Wait()
close(resultsChan) close(resultsChan)
var finalOutput strings.Builder
finalOutput.WriteString("Team Execution Summary (Parallel):\n\n")
// Pre-allocate to maintain order since channels don't guarantee arrival order // Pre-allocate to maintain order since channels don't guarantee arrival order
orderedResults := make([]workResult, len(members)) orderedResults := make([]workResult, len(members))
for res := range resultsChan { for res := range resultsChan {
orderedResults[res.index] = res orderedResults[res.index] = res
} }
hasError := false var successOutput strings.Builder
var failureOutput strings.Builder
successCount, failureCount := 0, 0
successOutput.WriteString("Team Execution Summary (Parallel):\n\n")
for _, res := range orderedResults { for _, res := range orderedResults {
if res.err != nil { if res.err != nil {
hasError = true failureCount++
finalOutput.WriteString(fmt.Sprintf("### Worker [%s] FAILED:\n%v\n\n", res.role, res.err)) failureOutput.WriteString(fmt.Sprintf("### Worker [%s] FAILED:\n%v\n\n", res.role, res.err))
} else { } else {
finalOutput.WriteString(fmt.Sprintf("### Worker [%s] Output:\n%s\n\n", res.role, res.res)) successCount++
successOutput.WriteString(fmt.Sprintf("### Worker [%s] Output:\n%s\n\n", res.role, res.res))
} }
} }
if hasError { if failureCount == 0 {
return ErrorResult("One or more parallel workers failed.\n" + finalOutput.String()) // All workers succeeded
return &ToolResult{
ForLLM: successOutput.String(),
ForUser: buildUserSummary("Parallel", members, nil),
}
}
// Partial failure: preserve successful results and append failure summary.
// This lets the coordinator decide how to handle the partial outcome.
fullOutput := successOutput.String()
if failureCount > 0 {
fullOutput += "---\n## ⚠️ Partial Failures\n\n" + failureOutput.String() +
fmt.Sprintf("\n%d/%d workers succeeded. %d worker(s) failed. The successful results above may still be usable.",
successCount, len(members), failureCount)
} }
return &ToolResult{ return &ToolResult{
ForLLM: finalOutput.String(), ForLLM: fullOutput,
ForUser: "Team completed parallel execution.", ForUser: fmt.Sprintf("⚠️ Parallel execution: %d/%d workers succeeded. %d failed.", successCount, len(members), failureCount),
IsError: failureCount == len(members),
} }
} }
func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig ToolLoopConfig, members []TeamMember) *ToolResult {
func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig ToolLoopConfig, members []TeamMember, contextLimit int) *ToolResult {
if len(members) != 2 { if len(members) != 2 {
return ErrorResult("The evaluator_optimizer strategy requires exactly two members: [0] Worker, [1] Evaluator.") return ErrorResult("The evaluator_optimizer strategy requires exactly two members: [0] Worker, [1] Evaluator.")
} }
@ -533,17 +571,27 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too
maxLoops = teamConfig.MaxEvaluatorLoops maxLoops = teamConfig.MaxEvaluatorLoops
} }
for attempt := 1; attempt <= maxLoops; attempt++ { // Pre-compute both configs once — they don't change between loop iterations.
finalOutput.WriteString(fmt.Sprintf("## Attempt %d\n", attempt))
// 2. Trigger Worker (resumes from its exact previous state!)
workerConfig, err := buildWorkerConfig(baseConfig, baseConfig.Tools, worker, t.manager) workerConfig, err := buildWorkerConfig(baseConfig, baseConfig.Tools, worker, t.manager)
if err != nil { if err != nil {
errStr := fmt.Sprintf("Worker configuration failed on attempt %d: %v", attempt, err) return ErrorResult(fmt.Sprintf("Worker configuration failed: %v", err)).WithError(err)
finalOutput.WriteString(errStr + "\n") }
return ErrorResult(errStr).WithError(err) evalConfig, err := buildWorkerConfig(baseConfig, NewToolRegistry(), evaluator, t.manager)
if err != nil {
return ErrorResult(fmt.Sprintf("Evaluator configuration failed: %v", err)).WithError(err)
} }
logger.InfoCF("team", "Evaluator-Optimizer starting", map[string]any{
"worker": worker.Role,
"evaluator": evaluator.Role,
"max_loops": maxLoops,
})
for attempt := 1; attempt <= maxLoops; attempt++ {
finalOutput.WriteString(fmt.Sprintf("## Attempt %d\n", attempt))
logger.InfoCF("team", fmt.Sprintf("Evaluator-Optimizer attempt %d/%d", attempt, maxLoops), map[string]any{})
// 2. Trigger Worker (resumes from its exact previous state!)
workerResult, err := RunToolLoop(ctx, workerConfig, workerMessages, t.originChannel, t.originChatID) workerResult, err := RunToolLoop(ctx, workerConfig, workerMessages, t.originChannel, t.originChatID)
if err != nil { if err != nil {
errStr := fmt.Sprintf("Worker failed on attempt %d: %v", attempt, err) errStr := fmt.Sprintf("Worker failed on attempt %d: %v", attempt, err)
@ -557,20 +605,15 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too
finalOutput.WriteString(fmt.Sprintf("### Worker Output:\n%s\n\n", workerResult.Content)) finalOutput.WriteString(fmt.Sprintf("### Worker Output:\n%s\n\n", workerResult.Content))
// 3. Trigger Evaluator (Ephemeral, stateless evaluation) // 3. Trigger Evaluator (Ephemeral, stateless evaluation)
evalContext := fmt.Sprintf("%s\n\n--- Worker's Output to Evaluate ---\n%s\n\nIf the output is completely correct and fulfills the task, you MUST reply starting with strictly '[PASS]'. Otherwise, explain the issues in detail.", evaluator.Task, truncateContext(workerResult.Content)) // The evaluator only needs to reason about text — give it no tools to avoid
// unnecessary tool calls, wasted tokens, and potential side effects.
evalContext := fmt.Sprintf("%s\n\n--- Worker's Output to Evaluate ---\n%s\n\nIf the output is completely correct and fulfills the task, you MUST reply starting with strictly '[PASS]'. Otherwise, explain the issues in detail.", evaluator.Task, truncateContextN(workerResult.Content, contextLimit))
evalMessages := []providers.Message{ evalMessages := []providers.Message{
{Role: "system", Content: evaluator.Role}, {Role: "system", Content: evaluator.Role},
{Role: "user", Content: evalContext}, {Role: "user", Content: evalContext},
} }
evalConfig, err := buildWorkerConfig(baseConfig, baseConfig.Tools, evaluator, t.manager)
if err != nil {
errStr := fmt.Sprintf("Evaluator configuration failed on attempt %d: %v", attempt, err)
finalOutput.WriteString(errStr + "\n")
return ErrorResult(errStr).WithError(err)
}
evalResult, err := RunToolLoop(ctx, evalConfig, evalMessages, t.originChannel, t.originChatID) evalResult, err := RunToolLoop(ctx, evalConfig, evalMessages, t.originChannel, t.originChatID)
if err != nil { if err != nil {
errStr := fmt.Sprintf("Evaluator failed on attempt %d: %v", attempt, err) errStr := fmt.Sprintf("Evaluator failed on attempt %d: %v", attempt, err)
@ -583,12 +626,15 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too
// 4. Check for PASS condition // 4. Check for PASS condition
if strings.HasPrefix(strings.TrimSpace(evalResult.Content), "[PASS]") { if strings.HasPrefix(strings.TrimSpace(evalResult.Content), "[PASS]") {
finalOutput.WriteString("✅ Evaluation Passed! Loop finished successfully.\n") finalOutput.WriteString("✅ Evaluation Passed! Loop finished successfully.\n")
logger.InfoCF("team", "Evaluator-Optimizer passed", map[string]any{"attempt": attempt})
return &ToolResult{ return &ToolResult{
ForLLM: finalOutput.String(), ForLLM: finalOutput.String(),
ForUser: "Evaluator-Optimizer loop completed successfully.", ForUser: fmt.Sprintf("✅ Evaluator-Optimizer passed on attempt %d/%d (worker: %s).", attempt, maxLoops, worker.Role),
} }
} }
logger.InfoCF("team", "Evaluator-Optimizer did not pass, retrying", map[string]any{"attempt": attempt, "max_loops": maxLoops})
// 5. If not passed, and not the last attempt, inject feedback into Worker's stateful memory // 5. If not passed, and not the last attempt, inject feedback into Worker's stateful memory
if attempt < maxLoops { if attempt < maxLoops {
injection := fmt.Sprintf("The evaluator rejected your previous attempt. Please fix the issues based on this feedback:\n\n%s", evalResult.Content) injection := fmt.Sprintf("The evaluator rejected your previous attempt. Please fix the issues based on this feedback:\n\n%s", evalResult.Content)
@ -600,13 +646,15 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too
} }
finalOutput.WriteString("❌ Maximum evaluation loops reached without a [PASS]. Returning current state.\n") finalOutput.WriteString("❌ Maximum evaluation loops reached without a [PASS]. Returning current state.\n")
logger.WarnCF("team", "Evaluator-Optimizer exhausted max loops", map[string]any{"max_loops": maxLoops})
return &ToolResult{ return &ToolResult{
ForLLM: finalOutput.String(), ForLLM: finalOutput.String(),
ForUser: "Evaluator-Optimizer loop exhausted maximum attempts.", ForUser: fmt.Sprintf("❌ Evaluator-Optimizer exhausted %d attempts without a [PASS].", maxLoops),
} }
} }
func (t *TeamTool) executeDAG(ctx context.Context, baseConfig ToolLoopConfig, members []TeamMember) *ToolResult { func (t *TeamTool) executeDAG(ctx context.Context, cancel context.CancelFunc, baseConfig ToolLoopConfig, members []TeamMember, contextLimit int) *ToolResult {
logger.InfoCF("team", "DAG execution starting", map[string]any{"member_count": len(members)})
// 1. Build and VALIDATE dependency graph // 1. Build and VALIDATE dependency graph
memberMap := make(map[string]TeamMember) memberMap := make(map[string]TeamMember)
inDegree := make(map[string]int) inDegree := make(map[string]int)
@ -711,7 +759,7 @@ func (t *TeamTool) executeDAG(ctx context.Context, baseConfig ToolLoopConfig, me
contextMu.Unlock() contextMu.Unlock()
if depsContext != "" { if depsContext != "" {
actualTask = fmt.Sprintf("%s\n\n--- Context from dependencies ---\n%s", m.Task, truncateContext(depsContext)) actualTask = fmt.Sprintf("%s\n\n--- Context from dependencies ---\n%s", m.Task, truncateContextN(depsContext, contextLimit))
} }
messages := []providers.Message{ messages := []providers.Message{
@ -753,7 +801,11 @@ func (t *TeamTool) executeDAG(ctx context.Context, baseConfig ToolLoopConfig, me
case res := <-resultChan: case res := <-resultChan:
if res.err != nil { if res.err != nil {
// Fast fail on first error // Fast fail on first error.
// Cancel the team context first so that all in-flight goroutines
// receive the cancellation signal and terminate cleanly.
cancel()
wg.Wait()
return ErrorResult(res.err.Error()) return ErrorResult(res.err.Error())
} }
@ -796,17 +848,44 @@ func (t *TeamTool) executeDAG(ctx context.Context, baseConfig ToolLoopConfig, me
return &ToolResult{ return &ToolResult{
ForLLM: finalOutput.String(), ForLLM: finalOutput.String(),
ForUser: "Team completed DAG execution.", ForUser: buildUserSummary("DAG", members, nil),
} }
} }
// truncateContext prevents Context Window Explosion (Token Bombs) // truncateContextN limits the number of runes in ctx to maxRunes.
// by limiting the size of upstream results injected into downstream prompts. // It prevents Context Window Explosion (Token Bombs) when passing upstream
func truncateContext(ctx string) string { // worker results into downstream agent prompts.
maxRunes := 8000 func truncateContextN(ctx string, maxRunes int) string {
runes := []rune(ctx) runes := []rune(ctx)
if len(runes) > maxRunes { if len(runes) > maxRunes {
return string(runes[:maxRunes]) + "\n...[Context truncated due to length]..." return string(runes[:maxRunes]) + "\n...[Context truncated due to length]..."
} }
return ctx return ctx
} }
// truncateContext is the default wrapper using 8000 runes (≈6000 words).
// Call truncateContextN directly when a configurable limit is needed.
func truncateContext(ctx string) string {
return truncateContextN(ctx, 8000)
}
// buildUserSummary produces a concise human-readable summary for the ForUser field,
// listing each member's role. errors (if any) are appended as a separate section.
func buildUserSummary(strategy string, members []TeamMember, errors []string) string {
var sb strings.Builder
sb.WriteString(fmt.Sprintf("Team (%s) completed — %d member(s):\n", strategy, len(members)))
for i, m := range members {
sb.WriteString(fmt.Sprintf(" [%d] %s", i+1, m.Role))
if m.Model != "" {
sb.WriteString(fmt.Sprintf(" (model: %s)", m.Model))
}
sb.WriteString("\n")
}
if len(errors) > 0 {
sb.WriteString("\n⚠ Failures:\n")
for _, e := range errors {
sb.WriteString(" • " + e + "\n")
}
}
return strings.TrimRight(sb.String(), "\n")
}