From 89e124caf9578b1b1fc396695c5b3fa63476cba6 Mon Sep 17 00:00:00 2001 From: Administrator <1280842908@qq.com> Date: Mon, 23 Mar 2026 22:49:46 +0800 Subject: [PATCH] feat(team): enhance team tool with sub-turn spawner integration and improve message handling --- pkg/agent/loop.go | 19 ++--- pkg/agent/subturn.go | 38 ++++++++- pkg/agent/turn.go | 1 + pkg/tools/spawn_sub_agent.go | 124 --------------------------- pkg/tools/subagent.go | 73 ++++++++++++---- pkg/tools/team.go | 158 +++++++++++++++++++++++++++++++---- 6 files changed, 241 insertions(+), 172 deletions(-) delete mode 100644 pkg/tools/spawn_sub_agent.go diff --git a/pkg/agent/loop.go b/pkg/agent/loop.go index f586da631..522997f6c 100644 --- a/pkg/agent/loop.go +++ b/pkg/agent/loop.go @@ -268,22 +268,18 @@ func registerSharedTools( } } - // Team and spawn_sub_agent tools - subagentManager := tools.NewSubagentManager(provider, agent.Model, agent.Candidates, agent.Workspace, cfg.Tools.Team, msgBus) - subagentManager.SetLLMOptions(agent.MaxTokens, agent.Temperature) + // Team tool + teamSubagentManager := tools.NewSubagentManager(provider, agent.Model, agent.Candidates, agent.Workspace, cfg.Tools.Team, msgBus) + teamSubagentManager.SetLLMOptions(agent.MaxTokens, agent.Temperature) - teamTool := tools.NewTeamTool(subagentManager, cfg) + teamTool := tools.NewTeamTool(teamSubagentManager, cfg) if cfg.Tools.IsToolEnabled("team") { + teamTool.SetSpawner(NewSubTurnSpawner(al)) agent.Tools.Register(teamTool) } - spawnSubAgentTool := tools.NewSpawnSubAgentTool(subagentManager) - if cfg.Tools.IsToolEnabled("spawn_sub_agent") { - agent.Tools.Register(spawnSubAgentTool) - } - - // Share the fully-built registry back to subagent manager - subagentManager.SetTools(agent.Tools) + // Share the fully-built registry back to team subagent manager + teamSubagentManager.SetTools(agent.Tools) // Spawn and spawn_status tools share a SubagentManager. // Construct it when either tool is enabled (both require subagent). @@ -2595,6 +2591,7 @@ turnLoop: finalContent: finalContent, status: turnStatus, followUps: append([]bus.InboundMessage(nil), ts.followUps...), + messages: append([]providers.Message(nil), messages...), }, nil } diff --git a/pkg/agent/subturn.go b/pkg/agent/subturn.go index f5ba412ab..7035c8c1e 100644 --- a/pkg/agent/subturn.go +++ b/pkg/agent/subturn.go @@ -102,7 +102,9 @@ type subTurnRuntimeConfig struct { // // Parent turn will poll and process it in a later iteration type SubTurnConfig struct { Model string + Provider providers.LLMProvider // non-nil overrides the child agent's provider Tools []tools.Tool + EmptyTools bool // true: child agent gets an empty ToolRegistry (overrides Tools) SystemPrompt string MaxTokens int @@ -220,7 +222,9 @@ func (s *AgentLoopSpawner) SpawnSubTurn( // Convert tools.SubTurnConfig to agent.SubTurnConfig agentCfg := SubTurnConfig{ Model: cfg.Model, + Provider: cfg.Provider, Tools: cfg.Tools, + EmptyTools: cfg.EmptyTools, SystemPrompt: cfg.SystemPrompt, ActualSystemPrompt: cfg.ActualSystemPrompt, InitialMessages: cfg.InitialMessages, @@ -232,6 +236,17 @@ func (s *AgentLoopSpawner) SpawnSubTurn( MaxContextRunes: cfg.MaxContextRunes, } + // Resolve model → provider when only a model name is given (no explicit provider). + // This enables heterogeneous model routing from tool-layer callers (e.g. subagent tool). + if agentCfg.Provider == nil && agentCfg.Model != "" { + if modelCfg, err := s.al.GetConfig().GetModelConfig(agentCfg.Model); err == nil { + if p, m, err := providers.CreateProviderFromConfig(modelCfg); err == nil { + agentCfg.Provider = p + agentCfg.Model = m + } + } + } + return spawnSubTurn(ctx, s.al, parentTS, agentCfg) } @@ -344,9 +359,25 @@ func spawnSubTurn( ephemeralStore := newEphemeralSession(nil) agent := *baseAgent // shallow copy agent.Sessions = ephemeralStore + // Apply model/provider override for heterogeneous agents. + if cfg.Model != "" { + agent.Model = cfg.Model + } + if cfg.Provider != nil { + agent.Provider = cfg.Provider + } // Clone the tool registry so child turn's tool registrations // don't pollute the parent's registry. - if baseAgent.Tools != nil { + if cfg.EmptyTools { + agent.Tools = tools.NewToolRegistry() + } else if cfg.Tools != nil { + // Tools override will be applied via processOptions below. + // Clone parent registry as base, then replace with cfg.Tools entries. + agent.Tools = tools.NewToolRegistry() + for _, t := range cfg.Tools { + agent.Tools.Register(t) + } + } else if baseAgent.Tools != nil { agent.Tools = baseAgent.Tools.Clone() } @@ -477,8 +508,9 @@ func spawnSubTurn( } } else { result = &tools.ToolResult{ - ForLLM: turnRes.finalContent, - ForUser: turnRes.finalContent, + ForLLM: turnRes.finalContent, + ForUser: turnRes.finalContent, + Messages: turnRes.messages, } } diff --git a/pkg/agent/turn.go b/pkg/agent/turn.go index e4970c519..7545069c1 100644 --- a/pkg/agent/turn.go +++ b/pkg/agent/turn.go @@ -43,6 +43,7 @@ type turnResult struct { finalContent string status TurnEndStatus followUps []bus.InboundMessage + messages []providers.Message // ephemeral session history after execution (for state continuation) } type turnState struct { diff --git a/pkg/tools/spawn_sub_agent.go b/pkg/tools/spawn_sub_agent.go deleted file mode 100644 index ecc71b69f..000000000 --- a/pkg/tools/spawn_sub_agent.go +++ /dev/null @@ -1,124 +0,0 @@ -package tools - -import ( - "context" - "fmt" - "strings" - - "github.com/sipeed/picoclaw/pkg/providers" -) - -// SpawnSubAgentTool executes a customized subagent task synchronously using Anthology-style single worker delegation. -type SpawnSubAgentTool struct { - manager *SubagentManager - originChannel string - originChatID string -} - -func NewSpawnSubAgentTool(manager *SubagentManager) *SpawnSubAgentTool { - return &SpawnSubAgentTool{ - manager: manager, - originChannel: "cli", - originChatID: "direct", - } -} - -func (t *SpawnSubAgentTool) Name() string { - return "spawn_sub_agent" -} - -func (t *SpawnSubAgentTool) Description() string { - base := "Directly delegate a specific task to a new, isolated sub-agent. You (the main agent) should autonomously determine the appropriate expert role and specific task based on the user's high-level request. It will execute independently and return the final result." - if t.manager != nil { - if hint := t.manager.ModelCapabilityHint(); hint != "" { - return base + "\n\n" + hint - } - } - return base -} - -func (t *SpawnSubAgentTool) Parameters() map[string]any { - return map[string]any{ - "type": "object", - "properties": map[string]any{ - "task": map[string]any{ - "type": "string", - "description": "The specific task the sub-agent needs to accomplish.", - }, - "role": map[string]any{ - "type": "string", - "description": "The system prompt/role assignment for the sub-agent (e.g., 'You are an expert code reviewer').", - }, - "model": map[string]any{ - "type": "string", - "description": "Optional specific LLM model ID to route this task to (e.g., 'gpt-4o' for vision, 'claude-3-5-sonnet' for logic). If omitted, inherits the parent's model.", - }, - }, - "required": []string{"task", "role"}, - } -} - -func (t *SpawnSubAgentTool) SetContext(channel, chatID string) { - t.originChannel = channel - t.originChatID = chatID -} - -func (t *SpawnSubAgentTool) Execute(ctx context.Context, args map[string]any) *ToolResult { - task, ok := args["task"].(string) - if !ok || strings.TrimSpace(task) == "" { - return ErrorResult("task is required").WithError(fmt.Errorf("task parameter is required")) - } - - role, ok := args["role"].(string) - if !ok || strings.TrimSpace(role) == "" { - return ErrorResult("role is required").WithError(fmt.Errorf("role parameter is required")) - } - - if t.manager == nil { - return ErrorResult("Subagent manager not configured").WithError(fmt.Errorf("manager is nil")) - } - - // 1. Isolation: Each SubAgent gets a completely fresh message set - messages := []providers.Message{ - { - Role: "system", - Content: role, - }, - { - Role: "user", - Content: task, - }, - } - - // 2. Base Configuration (Timeout & LLM constraints) - config := t.manager.BuildBaseWorkerConfig(ctx) - - // 2.1 Model Override (Heterogeneous Agents) - if modelParam, ok := args["model"].(string); ok && strings.TrimSpace(modelParam) != "" { - requestedModel := strings.TrimSpace(modelParam) - if !t.manager.IsModelAllowed(requestedModel) { - return ErrorResult(fmt.Sprintf("requested model '%s' is not in the allowed fallback candidates list for this agent workspace", requestedModel)).WithError(fmt.Errorf("model %s not allowed", requestedModel)) - } - config.Model = requestedModel - } - - // Note: For MVP, we pass the current ToolRegistry unmodified. - // To enforce strict sandboxing later, we can construct a new ToolRegistry here based on args['allowed_tools']. - - loopResult, err := RunToolLoop(ctx, config, messages, t.originChannel, t.originChatID) - if err != nil { - return ErrorResult(fmt.Sprintf("Subagent execution failed: %v", err)).WithError(err) - } - - // Return full details to LLM - llmContent := fmt.Sprintf("Subagent (Role: %s) task completed:\nIterations: %d\nResult: %s", - role, loopResult.Iterations, loopResult.Content) - - return &ToolResult{ - ForLLM: llmContent, - ForUser: "Sub-agent finished task.", - Silent: false, - IsError: false, - Async: false, - } -} diff --git a/pkg/tools/subagent.go b/pkg/tools/subagent.go index 8b6a0db80..f70b9b182 100644 --- a/pkg/tools/subagent.go +++ b/pkg/tools/subagent.go @@ -22,7 +22,9 @@ type SubTurnSpawner interface { // SubTurnConfig holds configuration for spawning a sub-turn. type SubTurnConfig struct { Model string + Provider providers.LLMProvider // non-nil overrides the child agent's provider Tools []Tool + EmptyTools bool // true: child agent gets an empty ToolRegistry (overrides Tools) SystemPrompt string MaxTokens int Temperature float64 @@ -432,10 +434,12 @@ func (sm *SubagentManager) BuildBaseWorkerConfig(ctx context.Context) ToolLoopCo // SubagentTool executes a subagent task synchronously and returns the result. // It directly calls SubTurnSpawner with Async=false for synchronous execution. type SubagentTool struct { - spawner SubTurnSpawner - defaultModel string - maxTokens int - temperature float64 + spawner SubTurnSpawner + defaultModel string + maxTokens int + temperature float64 + isModelAllowed func(string) bool // nil means no allowlist check + modelHint func() string // nil means no model hint in description } func NewSubagentTool(manager *SubagentManager) *SubagentTool { @@ -443,9 +447,11 @@ func NewSubagentTool(manager *SubagentManager) *SubagentTool { return &SubagentTool{} } return &SubagentTool{ - defaultModel: manager.defaultModel, - maxTokens: manager.maxTokens, - temperature: manager.temperature, + defaultModel: manager.defaultModel, + maxTokens: manager.maxTokens, + temperature: manager.temperature, + isModelAllowed: manager.IsModelAllowed, + modelHint: manager.ModelCapabilityHint, } } @@ -459,7 +465,13 @@ func (t *SubagentTool) Name() string { } func (t *SubagentTool) Description() string { - return "Execute a subagent task synchronously and return the result. Use this for delegating specific tasks to an independent agent instance. Returns execution summary to user and full details to LLM." + base := "Execute a subagent task synchronously and return the result. Use this for delegating specific tasks to an independent agent instance with an optional role (system prompt) and model. Returns execution summary to user and full details to LLM." + if t.modelHint != nil { + if hint := t.modelHint(); hint != "" { + return base + "\n\n" + hint + } + } + return base } func (t *SubagentTool) Parameters() map[string]any { @@ -474,6 +486,14 @@ func (t *SubagentTool) Parameters() map[string]any { "type": "string", "description": "Optional short label for the task (for display)", }, + "role": map[string]any{ + "type": "string", + "description": "Optional system prompt / role assignment for the subagent (e.g. 'You are an expert code reviewer'). If omitted, a default subagent prompt is used.", + }, + "model": map[string]any{ + "type": "string", + "description": "Optional specific LLM model ID to route this task to. If omitted, inherits the parent's model.", + }, }, "required": []string{"task"}, } @@ -486,15 +506,35 @@ func (t *SubagentTool) Execute(ctx context.Context, args map[string]any) *ToolRe } label, _ := args["label"].(string) + role, _ := args["role"].(string) + modelParam, _ := args["model"].(string) + modelParam = strings.TrimSpace(modelParam) - // Build system prompt for subagent + // Validate model against allowlist if provided + if modelParam != "" && t.isModelAllowed != nil && !t.isModelAllowed(modelParam) { + return ErrorResult(fmt.Sprintf("requested model '%s' is not in the allowed models list", modelParam)). + WithError(fmt.Errorf("model %s not allowed", modelParam)) + } + + // Determine the model to use + targetModel := t.defaultModel + if modelParam != "" { + targetModel = modelParam + } + + // Build ActualSystemPrompt: prefer explicit role, fall back to auto-generated prompt + var actualSystemPrompt string + if role != "" { + actualSystemPrompt = role + } + + // Build SystemPrompt (task description, becomes first user message in sub-turn) systemPrompt := fmt.Sprintf( `You are a subagent. Complete the given task independently and provide a clear, concise result. Task: %s`, task, ) - if label != "" { systemPrompt = fmt.Sprintf( `You are a subagent labeled "%s". Complete the given task independently and provide a clear, concise result. @@ -508,12 +548,13 @@ Task: %s`, // Use spawner if available (direct SpawnSubTurn call) if t.spawner != nil { result, err := t.spawner.SpawnSubTurn(ctx, SubTurnConfig{ - Model: t.defaultModel, - Tools: nil, // Will inherit from parent via context - SystemPrompt: systemPrompt, - MaxTokens: t.maxTokens, - Temperature: t.temperature, - Async: false, // Synchronous execution + Model: targetModel, + Tools: nil, // Will inherit from parent via context + SystemPrompt: systemPrompt, + ActualSystemPrompt: actualSystemPrompt, + MaxTokens: t.maxTokens, + Temperature: t.temperature, + Async: false, // Synchronous execution }) if err != nil { return ErrorResult(fmt.Sprintf("Subagent execution failed: %v", err)).WithError(err) diff --git a/pkg/tools/team.go b/pkg/tools/team.go index 8c2223d5e..eb7b8873c 100644 --- a/pkg/tools/team.go +++ b/pkg/tools/team.go @@ -15,6 +15,7 @@ import ( type TeamTool struct { manager *SubagentManager + spawner SubTurnSpawner cfg *config.Config originChannel string originChatID string @@ -38,6 +39,11 @@ func NewTeamTool(manager *SubagentManager, cfg *config.Config) *TeamTool { } } +// SetSpawner sets the SubTurnSpawner used to execute team members as sub-turns. +func (t *TeamTool) SetSpawner(spawner SubTurnSpawner) { + t.spawner = spawner +} + func (t *TeamTool) Name() string { return "team" } @@ -214,11 +220,11 @@ func (t *TeamTool) maybeRunAutoReviewer( "model": teamConfig.ReviewerModel, }) - loopResult, err := RunToolLoop(ctx, reviewerConfig, reviewerMessages, t.originChannel, t.originChatID) + loopContent, _, err := t.spawnWorker(ctx, reviewerConfig, reviewerMessages, nil) if err != nil { return fmt.Sprintf("[Auto-Reviewer] Failed to run: %v", err) } - return "[Auto-Reviewer Result]\n" + loopResult.Content + return "[Auto-Reviewer Result]\n" + loopContent } func (t *TeamTool) Execute(ctx context.Context, args map[string]any) *ToolResult { @@ -412,7 +418,121 @@ func upgradeRegistryForConcurrency(original *ToolRegistry) *ToolRegistry { return upgraded } -// buildWorkerConfig creates a ToolLoopConfig for a specific team member, +// spawnWorker executes a single team member's turn, routing through SubTurnSpawner when available. +// Returns (content, messages, error). The messages slice is non-nil only for stateful workers +// (evaluator_optimizer) and can be passed as InitialMessages for the next iteration. +func (t *TeamTool) spawnWorker(ctx context.Context, cfg ToolLoopConfig, messages []providers.Message, budget *atomic.Int64) (string, []providers.Message, error) { + if t.spawner == nil { + // Fallback: direct RunToolLoop (no turnState integration) + res, err := RunToolLoop(ctx, cfg, messages, t.originChannel, t.originChatID) + if err != nil { + return "", nil, err + } + return res.Content, res.Messages, nil + } + + // Convert ToolLoopConfig + messages into SubTurnConfig for SubTurnSpawner. + var toolSlice []Tool + if cfg.Tools != nil { + for _, name := range cfg.Tools.ListTools() { + if tool, ok := cfg.Tools.Get(name); ok { + toolSlice = append(toolSlice, tool) + } + } + } + + // Extract system prompt and non-system messages + var actualSystemPrompt string + var initialMessages []providers.Message + for _, msg := range messages { + if msg.Role == "system" { + actualSystemPrompt = msg.Content + } else { + initialMessages = append(initialMessages, msg) + } + } + + maxTokens, temperature := getLLMOptionsFromConfig(cfg) + + subCfg := SubTurnConfig{ + Model: cfg.Model, + Provider: cfg.Provider, + Tools: toolSlice, + ActualSystemPrompt: actualSystemPrompt, + InitialMessages: initialMessages, + MaxTokens: maxTokens, + Temperature: temperature, + Async: false, + InitialTokenBudget: budget, + } + + res, err := t.spawner.SpawnSubTurn(ctx, subCfg) + if err != nil { + return "", nil, err + } + return res.ForLLM, res.Messages, nil +} + +// spawnWorkerEmptyTools is like spawnWorker but forces an empty tool registry on the sub-turn. +// Used for the evaluator in evaluator_optimizer to prevent side effects. +func (t *TeamTool) spawnWorkerEmptyTools(ctx context.Context, cfg ToolLoopConfig, messages []providers.Message) (string, error) { + if t.spawner == nil { + // Fallback: direct RunToolLoop with empty registry + emptyConfig := cfg + emptyConfig.Tools = NewToolRegistry() + res, err := RunToolLoop(ctx, emptyConfig, messages, t.originChannel, t.originChatID) + if err != nil { + return "", err + } + return res.Content, nil + } + + var actualSystemPrompt string + var initialMessages []providers.Message + for _, msg := range messages { + if msg.Role == "system" { + actualSystemPrompt = msg.Content + } else { + initialMessages = append(initialMessages, msg) + } + } + + maxTokens, temperature := getLLMOptionsFromConfig(cfg) + + subCfg := SubTurnConfig{ + Model: cfg.Model, + Provider: cfg.Provider, + EmptyTools: true, + ActualSystemPrompt: actualSystemPrompt, + InitialMessages: initialMessages, + MaxTokens: maxTokens, + Temperature: temperature, + Async: false, + } + + res, err := t.spawner.SpawnSubTurn(ctx, subCfg) + if err != nil { + return "", err + } + return res.ForLLM, nil +} + +// getLLMOptionsFromConfig extracts MaxTokens and Temperature from a ToolLoopConfig's LLMOptions map. +func getLLMOptionsFromConfig(cfg ToolLoopConfig) (int, float64) { + var maxTokens int + var temperature float64 + if cfg.LLMOptions != nil { + if v, ok := cfg.LLMOptions["max_tokens"].(int); ok { + maxTokens = v + } + if v, ok := cfg.LLMOptions["temperature"].(float64); ok { + temperature = v + } + } + return maxTokens, temperature +} + + // potentially overriding the model based on the member's definition. func (t *TeamTool) buildWorkerConfig(baseConfig ToolLoopConfig, registry *ToolRegistry, m TeamMember) (ToolLoopConfig, error) { cfg := baseConfig @@ -487,14 +607,14 @@ func (t *TeamTool) executeSequential(ctx context.Context, baseConfig ToolLoopCon return ErrorResult(errStr).WithError(err) } - loopResult, err := RunToolLoop(ctx, workerConfig, messages, t.originChannel, t.originChatID) + content, _, err := t.spawnWorker(ctx, workerConfig, messages, baseConfig.RemainingTokenBudget) if err != nil { errStr := fmt.Sprintf("Phase %d (Role: %s) failed: %v", i+1, m.Role, err) finalOutput.WriteString(errStr + "\n") return ErrorResult(errStr).WithError(err) // Fail fast } - previousResult = loopResult.Content + previousResult = content finalOutput.WriteString(fmt.Sprintf("### Phase %d completed by Role: [%s]\n%s\n\n", i+1, m.Role, previousResult)) } @@ -537,13 +657,13 @@ func (t *TeamTool) executeParallel(ctx context.Context, baseConfig ToolLoopConfi return } - loopResult, err := RunToolLoop(ctx, workerConfig, messages, t.originChannel, t.originChatID) + content, _, err := t.spawnWorker(ctx, workerConfig, messages, baseConfig.RemainingTokenBudget) if err != nil { resultsChan <- workResult{index: index, role: member.Role, err: err} return } - resultsChan <- workResult{index: index, role: member.Role, res: loopResult.Content} + resultsChan <- workResult{index: index, role: member.Role, res: content} logger.InfoCF("team", fmt.Sprintf("[%s] Parallel worker finished", member.Role), map[string]any{ "member_index": index, }) @@ -648,7 +768,7 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too 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) + workerContent, workerMsgs, err := t.spawnWorker(ctx, workerConfig, workerMessages, baseConfig.RemainingTokenBudget) if err != nil { errStr := fmt.Sprintf("Worker failed on attempt %d: %v", attempt, err) finalOutput.WriteString(errStr + "\n") @@ -656,31 +776,33 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too } // Save the worker's cognitive state so it remembers its thought process for the next loop - workerMessages = workerResult.Messages + if workerMsgs != nil { + workerMessages = workerMsgs + } - finalOutput.WriteString(fmt.Sprintf("### Worker Output:\n%s\n\n", workerResult.Content)) + finalOutput.WriteString(fmt.Sprintf("### Worker Output:\n%s\n\n", workerContent)) // 3. Trigger Evaluator (Ephemeral, stateless evaluation) // 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)) + 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(workerContent, contextLimit)) evalMessages := []providers.Message{ {Role: "system", Content: evaluator.Role}, {Role: "user", Content: evalContext}, } - evalResult, err := RunToolLoop(ctx, evalConfig, evalMessages, t.originChannel, t.originChatID) + evalContent, err := t.spawnWorkerEmptyTools(ctx, evalConfig, evalMessages) if err != nil { errStr := fmt.Sprintf("Evaluator failed on attempt %d: %v", attempt, err) finalOutput.WriteString(errStr + "\n") return ErrorResult(errStr).WithError(err) } - finalOutput.WriteString(fmt.Sprintf("### Evaluator Feedback:\n%s\n\n", evalResult.Content)) + finalOutput.WriteString(fmt.Sprintf("### Evaluator Feedback:\n%s\n\n", evalContent)) // 4. Check for PASS condition - if strings.HasPrefix(strings.TrimSpace(evalResult.Content), "[PASS]") { + if strings.HasPrefix(strings.TrimSpace(evalContent), "[PASS]") { finalOutput.WriteString("✅ Evaluation Passed! Loop finished successfully.\n") logger.InfoCF("team", "Evaluator-Optimizer passed", map[string]any{"attempt": attempt}) return &ToolResult{ @@ -693,7 +815,7 @@ func (t *TeamTool) executeEvaluatorOptimizer(ctx context.Context, baseConfig Too // 5. If not passed, and not the last attempt, inject feedback into Worker's stateful memory 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", evalContent) workerMessages = append(workerMessages, providers.Message{ Role: "user", Content: injection, @@ -834,7 +956,7 @@ func (t *TeamTool) executeDAG(ctx context.Context, cancel context.CancelFunc, ba return } - loopResult, err := RunToolLoop(ctx, workerConfig, messages, t.originChannel, t.originChatID) + content, _, err := t.spawnWorker(ctx, workerConfig, messages, baseConfig.RemainingTokenBudget) if err != nil { masterErrMu.Lock() @@ -848,11 +970,11 @@ func (t *TeamTool) executeDAG(ctx context.Context, cancel context.CancelFunc, ba // Store result for final output finalResultsMu.Lock() - finalResults[id] = loopResult.Content + finalResults[id] = content finalResultsMu.Unlock() // Pass result to dependents - resultChan <- nodeResult{id: id, res: loopResult.Content} + resultChan <- nodeResult{id: id, res: content} }(memberID) case res := <-resultChan: