refactor: introduce orch.AgentReporter — decouple Broadcaster from SubagentManager
- Add pkg/orch/reporter.go: AgentReporter interface + Noop singleton
- Broadcaster now implements AgentReporter (ReportSpawn/StateChange/Conversation/GC)
- ToolLoopConfig: replace OnStateChange func with Reporter+AgentID
- SubagentManager: accept AgentReporter in constructor, remove internal Broadcaster
and GetBroadcaster(); all Publish calls replaced with Report* calls
- AgentLoop: add orchBroadcaster/*orchReporter fields, reporter() nil-safe helper,
SetOrchReporter/GetOrchBroadcaster public API
- NewAgentLoop: create struct before registerSharedTools so al.reporter() is
available; auto-detect orchestration from registry config
- runAgentLoop: ReportSpawn on entry, defer ReportGC on exit
- runLLMIteration: ReportStateChange("waiting") before LLM call,
ReportStateChange("toolcall", name) before each tool execution
- cmd_gateway.go: wire GetOrchBroadcaster() → handler.SetOrchBroadcaster()
- Update subagent_tool_test.go for new constructor signature
- Document hierarchy in CLAUDE.md
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
698fe84673
commit
dd81b4b2bf
8 changed files with 218 additions and 80 deletions
71
CLAUDE.md
71
CLAUDE.md
|
|
@ -354,6 +354,77 @@ pkg/agent/
|
||||||
context.go — conductor identity + orchestration guidance 追加
|
context.go — conductor identity + orchestration guidance 追加
|
||||||
```
|
```
|
||||||
|
|
||||||
|
### AgentReporter 抽象化 (実装済み 2026-02-25)
|
||||||
|
|
||||||
|
> branch `sub-agent-technical-breakdown`
|
||||||
|
|
||||||
|
`Broadcaster` を `SubagentManager` 内部で生成する密結合を解消し、
|
||||||
|
`orch.AgentReporter` インターフェースを中心に置くリファクタリングを実施。
|
||||||
|
|
||||||
|
#### オーナーシップ
|
||||||
|
|
||||||
|
```
|
||||||
|
AgentLoop
|
||||||
|
├─ owns: *orch.Broadcaster (orchBroadcaster — nil when disabled)
|
||||||
|
└─ holds: orch.AgentReporter (orchReporter = Broadcaster or Noop)
|
||||||
|
├─ passes to → SubagentManager.reporter
|
||||||
|
│ └─ passes to → ToolLoopConfig.Reporter
|
||||||
|
└─ calls directly for main/heartbeat sessions
|
||||||
|
├─ runAgentLoop: ReportSpawn / ReportGC
|
||||||
|
└─ runLLMIteration: ReportStateChange
|
||||||
|
|
||||||
|
cmd_gateway.go
|
||||||
|
└─ agentLoop.GetOrchBroadcaster() → handler.SetOrchBroadcaster()
|
||||||
|
|
||||||
|
miniapp.Handler
|
||||||
|
└─ borrows *orch.Broadcaster for Subscribe/Snapshot (WS 配信)
|
||||||
|
```
|
||||||
|
|
||||||
|
#### インターフェース (`pkg/orch/reporter.go`)
|
||||||
|
|
||||||
|
```go
|
||||||
|
type AgentReporter interface {
|
||||||
|
ReportSpawn(id, label, task string)
|
||||||
|
ReportStateChange(id, state, tool string)
|
||||||
|
ReportConversation(from, to, text string)
|
||||||
|
ReportGC(id, reason string)
|
||||||
|
}
|
||||||
|
var Noop AgentReporter = &noopReporter{} // nil-free; 全メソッドが no-op
|
||||||
|
```
|
||||||
|
|
||||||
|
`Broadcaster` は `AgentReporter` を満たす (`ReportSpawn` 等が `Publish` のラッパー)。
|
||||||
|
|
||||||
|
#### Noop パターン
|
||||||
|
|
||||||
|
```
|
||||||
|
--orchestration なし: orchReporter = orch.Noop → 全 Report* が空振り (panic なし)
|
||||||
|
--orchestration あり: orchReporter = *Broadcaster → WS 配信
|
||||||
|
```
|
||||||
|
|
||||||
|
呼び出し側は `if reporter != nil` チェック不要。
|
||||||
|
|
||||||
|
#### イベント発火の責任分担
|
||||||
|
|
||||||
|
| 発火元 | イベント | 経由 |
|
||||||
|
|--------|---------|------|
|
||||||
|
| `runAgentLoop` | `ReportSpawn` / `ReportGC` | `al.reporter()` |
|
||||||
|
| `runLLMIteration` | `ReportStateChange("waiting"/"toolcall")` | `al.reporter()` |
|
||||||
|
| `SubagentManager.Spawn` | `ReportSpawn` | `sm.reporter` |
|
||||||
|
| `SubagentManager.runTask` | `ReportConversation` / `ReportGC` | `sm.reporter` |
|
||||||
|
| `RunToolLoop` | `ReportStateChange` | `config.Reporter` |
|
||||||
|
|
||||||
|
main / heartbeat / subagent の全セッションが同一 Broadcaster に発火するため、
|
||||||
|
canvas には全エージェントが統一して表示される。
|
||||||
|
|
||||||
|
#### 変更ファイル
|
||||||
|
|
||||||
|
- `pkg/orch/reporter.go` — **新規** インターフェース + Noop
|
||||||
|
- `pkg/orch/broadcaster.go` — `ReportSpawn/StateChange/Conversation/GC` 追加
|
||||||
|
- `pkg/tools/toolloop.go` — `OnStateChange func` → `Reporter AgentReporter + AgentID`
|
||||||
|
- `pkg/tools/subagent.go` — constructor に `reporter` 受け取り、内部 broadcaster 廃止、`GetBroadcaster()` 削除
|
||||||
|
- `pkg/agent/loop.go` — `orchBroadcaster`/`orchReporter` フィールド追加、`SetOrchReporter`/`GetOrchBroadcaster` 追加、`registerSharedTools` シグネチャに `al *AgentLoop` 追加
|
||||||
|
- `cmd/picoclaw/cmd_gateway.go` — `GetOrchBroadcaster()` → `handler.SetOrchBroadcaster()`
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## Memory Optimization Notes
|
## Memory Optimization Notes
|
||||||
|
|
|
||||||
|
|
@ -243,6 +243,9 @@ func gatewayCmd() {
|
||||||
miniappNotifier = miniapp.NewStateNotifier()
|
miniappNotifier = miniapp.NewStateNotifier()
|
||||||
handler := miniapp.NewHandler(provider, sender, cfg.Channels.Telegram.Token, miniappNotifier, cfg.Channels.Telegram.AllowFrom, cfg.WorkspacePath())
|
handler := miniapp.NewHandler(provider, sender, cfg.Channels.Telegram.Token, miniappNotifier, cfg.Channels.Telegram.AllowFrom, cfg.WorkspacePath())
|
||||||
agentLoop.OnStateChange = miniappNotifier.Notify
|
agentLoop.OnStateChange = miniappNotifier.Notify
|
||||||
|
if b := agentLoop.GetOrchBroadcaster(); b != nil {
|
||||||
|
handler.SetOrchBroadcaster(b)
|
||||||
|
}
|
||||||
handler.RegisterRoutes(healthServer.Mux())
|
handler.RegisterRoutes(healthServer.Mux())
|
||||||
|
|
||||||
// Register dev preview tool for all agents
|
// Register dev preview tool for all agents
|
||||||
|
|
|
||||||
|
|
@ -25,6 +25,7 @@ import (
|
||||||
"github.com/sipeed/picoclaw/pkg/config"
|
"github.com/sipeed/picoclaw/pkg/config"
|
||||||
"github.com/sipeed/picoclaw/pkg/constants"
|
"github.com/sipeed/picoclaw/pkg/constants"
|
||||||
"github.com/sipeed/picoclaw/pkg/logger"
|
"github.com/sipeed/picoclaw/pkg/logger"
|
||||||
|
"github.com/sipeed/picoclaw/pkg/orch"
|
||||||
"github.com/sipeed/picoclaw/pkg/providers"
|
"github.com/sipeed/picoclaw/pkg/providers"
|
||||||
"github.com/sipeed/picoclaw/pkg/providers/protocoltypes"
|
"github.com/sipeed/picoclaw/pkg/providers/protocoltypes"
|
||||||
"github.com/sipeed/picoclaw/pkg/routing"
|
"github.com/sipeed/picoclaw/pkg/routing"
|
||||||
|
|
@ -94,6 +95,8 @@ type AgentLoop struct {
|
||||||
promptDirty atomic.Bool // true = rebuild needed on next GetSystemPrompt read
|
promptDirty atomic.Bool // true = rebuild needed on next GetSystemPrompt read
|
||||||
OnStateChange func() // called on plan/session/skills mutations
|
OnStateChange func() // called on plan/session/skills mutations
|
||||||
OnUserMessage func() // called when a real user message is processed
|
OnUserMessage func() // called when a real user message is processed
|
||||||
|
orchBroadcaster *orch.Broadcaster // nil when --orchestration not set
|
||||||
|
orchReporter orch.AgentReporter // always non-nil (Noop when disabled)
|
||||||
}
|
}
|
||||||
|
|
||||||
// processOptions configures how a message is processed
|
// processOptions configures how a message is processed
|
||||||
|
|
@ -114,9 +117,6 @@ type processOptions struct {
|
||||||
func NewAgentLoop(cfg *config.Config, msgBus *bus.MessageBus, provider providers.LLMProvider, enableStats ...bool) *AgentLoop {
|
func NewAgentLoop(cfg *config.Config, msgBus *bus.MessageBus, provider providers.LLMProvider, enableStats ...bool) *AgentLoop {
|
||||||
registry := NewAgentRegistry(cfg, provider)
|
registry := NewAgentRegistry(cfg, provider)
|
||||||
|
|
||||||
// Register shared tools to all agents
|
|
||||||
registerSharedTools(cfg, msgBus, registry, provider)
|
|
||||||
|
|
||||||
// Set up shared fallback chain
|
// Set up shared fallback chain
|
||||||
cooldown := providers.NewCooldownTracker()
|
cooldown := providers.NewCooldownTracker()
|
||||||
fallbackChain := providers.NewFallbackChain(cooldown)
|
fallbackChain := providers.NewFallbackChain(cooldown)
|
||||||
|
|
@ -136,17 +136,57 @@ func NewAgentLoop(cfg *config.Config, msgBus *bus.MessageBus, provider providers
|
||||||
statsTracker = stats.NewTracker(defaultAgent.Workspace)
|
statsTracker = stats.NewTracker(defaultAgent.Workspace)
|
||||||
}
|
}
|
||||||
|
|
||||||
return &AgentLoop{
|
// Determine if orchestration broadcaster is needed (any agent has subagents enabled).
|
||||||
bus: msgBus,
|
var orchBroadcaster *orch.Broadcaster
|
||||||
cfg: cfg,
|
var orchReporter orch.AgentReporter = orch.Noop
|
||||||
registry: registry,
|
for _, id := range registry.ListAgentIDs() {
|
||||||
state: stateManager,
|
if a, ok := registry.GetAgent(id); ok && a.Subagents != nil && a.Subagents.Enabled {
|
||||||
stats: statsTracker,
|
orchBroadcaster = orch.NewBroadcaster()
|
||||||
summarizing: sync.Map{},
|
orchReporter = orchBroadcaster
|
||||||
fallback: fallbackChain,
|
break
|
||||||
providerCache: providerCache,
|
}
|
||||||
sessions: NewSessionTracker(),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
al := &AgentLoop{
|
||||||
|
bus: msgBus,
|
||||||
|
cfg: cfg,
|
||||||
|
registry: registry,
|
||||||
|
state: stateManager,
|
||||||
|
stats: statsTracker,
|
||||||
|
summarizing: sync.Map{},
|
||||||
|
fallback: fallbackChain,
|
||||||
|
providerCache: providerCache,
|
||||||
|
sessions: NewSessionTracker(),
|
||||||
|
orchBroadcaster: orchBroadcaster,
|
||||||
|
orchReporter: orchReporter,
|
||||||
|
}
|
||||||
|
|
||||||
|
// Register shared tools to all agents (needs al for reporter injection).
|
||||||
|
registerSharedTools(cfg, msgBus, registry, provider, al)
|
||||||
|
|
||||||
|
return al
|
||||||
|
}
|
||||||
|
|
||||||
|
// reporter returns the active AgentReporter (never nil).
|
||||||
|
func (al *AgentLoop) reporter() orch.AgentReporter {
|
||||||
|
if al.orchReporter == nil {
|
||||||
|
return orch.Noop
|
||||||
|
}
|
||||||
|
return al.orchReporter
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetOrchReporter wires a Broadcaster as the active reporter.
|
||||||
|
// Called from cmd_gateway.go when --orchestration is set.
|
||||||
|
// --orchestration なし → 呼ばれない → reporter() は Noop を返す。
|
||||||
|
func (al *AgentLoop) SetOrchReporter(b *orch.Broadcaster) {
|
||||||
|
al.orchBroadcaster = b
|
||||||
|
al.orchReporter = b
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetOrchBroadcaster returns the concrete Broadcaster for miniapp wiring.
|
||||||
|
// Returns nil when orchestration is disabled.
|
||||||
|
func (al *AgentLoop) GetOrchBroadcaster() *orch.Broadcaster {
|
||||||
|
return al.orchBroadcaster
|
||||||
}
|
}
|
||||||
|
|
||||||
func (al *AgentLoop) notifyStateChange() {
|
func (al *AgentLoop) notifyStateChange() {
|
||||||
|
|
@ -162,6 +202,7 @@ func registerSharedTools(
|
||||||
msgBus *bus.MessageBus,
|
msgBus *bus.MessageBus,
|
||||||
registry *AgentRegistry,
|
registry *AgentRegistry,
|
||||||
provider providers.LLMProvider,
|
provider providers.LLMProvider,
|
||||||
|
al *AgentLoop,
|
||||||
) {
|
) {
|
||||||
for _, agentID := range registry.ListAgentIDs() {
|
for _, agentID := range registry.ListAgentIDs() {
|
||||||
agent, ok := registry.GetAgent(agentID)
|
agent, ok := registry.GetAgent(agentID)
|
||||||
|
|
@ -218,7 +259,7 @@ func registerSharedTools(
|
||||||
|
|
||||||
// Spawn tool — only registered when orchestration is explicitly enabled.
|
// Spawn tool — only registered when orchestration is explicitly enabled.
|
||||||
if agent.Subagents != nil && agent.Subagents.Enabled {
|
if agent.Subagents != nil && agent.Subagents.Enabled {
|
||||||
subagentManager := tools.NewSubagentManager(provider, agent.Model, agent.Workspace, msgBus)
|
subagentManager := tools.NewSubagentManager(provider, agent.Model, agent.Workspace, msgBus, al.reporter())
|
||||||
subagentManager.SetLLMOptions(agent.MaxTokens, agent.Temperature)
|
subagentManager.SetLLMOptions(agent.MaxTokens, agent.Temperature)
|
||||||
spawnTool := tools.NewSpawnTool(subagentManager)
|
spawnTool := tools.NewSpawnTool(subagentManager)
|
||||||
currentAgentID := agentID
|
currentAgentID := agentID
|
||||||
|
|
@ -715,6 +756,10 @@ func (al *AgentLoop) runAgentLoop(ctx context.Context, agent *AgentInstance, opt
|
||||||
}
|
}
|
||||||
defer al.releaseSessionLock(opts.SessionKey)
|
defer al.releaseSessionLock(opts.SessionKey)
|
||||||
|
|
||||||
|
// Report session lifecycle to canvas.
|
||||||
|
al.reporter().ReportSpawn(opts.SessionKey, opts.Channel, opts.UserMessage)
|
||||||
|
defer al.reporter().ReportGC(opts.SessionKey, "completed")
|
||||||
|
|
||||||
// -0. Create cancellable child context and register active task
|
// -0. Create cancellable child context and register active task
|
||||||
taskCtx, taskCancel := context.WithCancel(ctx)
|
taskCtx, taskCancel := context.WithCancel(ctx)
|
||||||
defer taskCancel()
|
defer taskCancel()
|
||||||
|
|
@ -1787,6 +1832,9 @@ func (al *AgentLoop) runLLMIteration(
|
||||||
return doCall(ctx, agent.Provider, primaryModel)
|
return doCall(ctx, agent.Provider, primaryModel)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Report waiting state to canvas before each LLM call.
|
||||||
|
al.reporter().ReportStateChange(opts.SessionKey, "waiting", "")
|
||||||
|
|
||||||
// Retry loop for context/token errors
|
// Retry loop for context/token errors
|
||||||
maxRetries := 2
|
maxRetries := 2
|
||||||
for retry := 0; retry <= maxRetries; retry++ {
|
for retry := 0; retry <= maxRetries; retry++ {
|
||||||
|
|
@ -2136,6 +2184,9 @@ func (al *AgentLoop) runLLMIteration(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Report toolcall state to canvas.
|
||||||
|
al.reporter().ReportStateChange(opts.SessionKey, "toolcall", tc.Name)
|
||||||
|
|
||||||
toolStart := time.Now()
|
toolStart := time.Now()
|
||||||
toolCtx := ctx
|
toolCtx := ctx
|
||||||
if wt := agent.GetWorktree(opts.SessionKey); wt != nil {
|
if wt := agent.GetWorktree(opts.SessionKey); wt != nil {
|
||||||
|
|
|
||||||
|
|
@ -83,6 +83,26 @@ func (b *Broadcaster) Snapshot() []AgentInfo {
|
||||||
return out
|
return out
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ReportSpawn implements AgentReporter.
|
||||||
|
func (b *Broadcaster) ReportSpawn(id, label, task string) {
|
||||||
|
b.Publish(Event{Type: "agent_spawn", ID: id, Label: label, Task: task})
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReportStateChange implements AgentReporter.
|
||||||
|
func (b *Broadcaster) ReportStateChange(id, state, tool string) {
|
||||||
|
b.Publish(Event{Type: "agent_state", ID: id, State: state, Tool: tool})
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReportConversation implements AgentReporter.
|
||||||
|
func (b *Broadcaster) ReportConversation(from, to, text string) {
|
||||||
|
b.Publish(Event{Type: "conversation", From: from, To: to, Text: text})
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReportGC implements AgentReporter.
|
||||||
|
func (b *Broadcaster) ReportGC(id, reason string) {
|
||||||
|
b.Publish(Event{Type: "agent_gc", ID: id, Reason: reason})
|
||||||
|
}
|
||||||
|
|
||||||
// Publish updates internal agent state and fans out to all subscribers.
|
// Publish updates internal agent state and fans out to all subscribers.
|
||||||
func (b *Broadcaster) Publish(ev Event) {
|
func (b *Broadcaster) Publish(ev Event) {
|
||||||
if ev.Created == 0 {
|
if ev.Created == 0 {
|
||||||
|
|
|
||||||
21
pkg/orch/reporter.go
Normal file
21
pkg/orch/reporter.go
Normal file
|
|
@ -0,0 +1,21 @@
|
||||||
|
package orch
|
||||||
|
|
||||||
|
// AgentReporter is the interface for reporting agent lifecycle events.
|
||||||
|
// Both Broadcaster (real events) and noopReporter (disabled) implement this.
|
||||||
|
type AgentReporter interface {
|
||||||
|
ReportSpawn(id, label, task string)
|
||||||
|
ReportStateChange(id, state, tool string)
|
||||||
|
ReportConversation(from, to, text string)
|
||||||
|
ReportGC(id, reason string)
|
||||||
|
}
|
||||||
|
|
||||||
|
type noopReporter struct{}
|
||||||
|
|
||||||
|
func (n *noopReporter) ReportSpawn(id, label, task string) {}
|
||||||
|
func (n *noopReporter) ReportStateChange(id, state, tool string) {}
|
||||||
|
func (n *noopReporter) ReportConversation(from, to, text string) {}
|
||||||
|
func (n *noopReporter) ReportGC(id, reason string) {}
|
||||||
|
|
||||||
|
// Noop is the AgentReporter to use when orchestration is disabled.
|
||||||
|
// Allows nil-free code in callers.
|
||||||
|
var Noop AgentReporter = &noopReporter{}
|
||||||
|
|
@ -37,14 +37,18 @@ type SubagentManager struct {
|
||||||
hasMaxTokens bool
|
hasMaxTokens bool
|
||||||
hasTemperature bool
|
hasTemperature bool
|
||||||
nextID int
|
nextID int
|
||||||
broadcaster *orch.Broadcaster
|
reporter orch.AgentReporter
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewSubagentManager(
|
func NewSubagentManager(
|
||||||
provider providers.LLMProvider,
|
provider providers.LLMProvider,
|
||||||
defaultModel, workspace string,
|
defaultModel, workspace string,
|
||||||
bus *bus.MessageBus,
|
bus *bus.MessageBus,
|
||||||
|
reporter orch.AgentReporter,
|
||||||
) *SubagentManager {
|
) *SubagentManager {
|
||||||
|
if reporter == nil {
|
||||||
|
reporter = orch.Noop
|
||||||
|
}
|
||||||
return &SubagentManager{
|
return &SubagentManager{
|
||||||
tasks: make(map[string]*SubagentTask),
|
tasks: make(map[string]*SubagentTask),
|
||||||
provider: provider,
|
provider: provider,
|
||||||
|
|
@ -54,16 +58,10 @@ func NewSubagentManager(
|
||||||
tools: NewToolRegistry(),
|
tools: NewToolRegistry(),
|
||||||
maxIterations: 10,
|
maxIterations: 10,
|
||||||
nextID: 1,
|
nextID: 1,
|
||||||
broadcaster: orch.NewBroadcaster(),
|
reporter: reporter,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetBroadcaster returns the Broadcaster so the miniapp handler can
|
|
||||||
// subscribe to real-time orchestration events.
|
|
||||||
func (sm *SubagentManager) GetBroadcaster() *orch.Broadcaster {
|
|
||||||
return sm.broadcaster
|
|
||||||
}
|
|
||||||
|
|
||||||
// SetLLMOptions sets max tokens and temperature for subagent LLM calls.
|
// SetLLMOptions sets max tokens and temperature for subagent LLM calls.
|
||||||
func (sm *SubagentManager) SetLLMOptions(maxTokens int, temperature float64) {
|
func (sm *SubagentManager) SetLLMOptions(maxTokens int, temperature float64) {
|
||||||
sm.mu.Lock()
|
sm.mu.Lock()
|
||||||
|
|
@ -112,12 +110,7 @@ func (sm *SubagentManager) Spawn(
|
||||||
}
|
}
|
||||||
sm.tasks[taskID] = subagentTask
|
sm.tasks[taskID] = subagentTask
|
||||||
|
|
||||||
sm.broadcaster.Publish(orch.Event{
|
sm.reporter.ReportSpawn(taskID, label, task)
|
||||||
Type: "agent_spawn",
|
|
||||||
ID: taskID,
|
|
||||||
Label: label,
|
|
||||||
Task: task,
|
|
||||||
})
|
|
||||||
|
|
||||||
// Start task in background with context cancellation support
|
// Start task in background with context cancellation support
|
||||||
go sm.runTask(ctx, subagentTask, callback)
|
go sm.runTask(ctx, subagentTask, callback)
|
||||||
|
|
@ -130,7 +123,6 @@ func (sm *SubagentManager) Spawn(
|
||||||
|
|
||||||
func (sm *SubagentManager) runTask(ctx context.Context, task *SubagentTask, callback AsyncCallback) {
|
func (sm *SubagentManager) runTask(ctx context.Context, task *SubagentTask, callback AsyncCallback) {
|
||||||
task.Status = "running"
|
task.Status = "running"
|
||||||
task.Created = time.Now().UnixMilli()
|
|
||||||
|
|
||||||
// Build system prompt for subagent
|
// Build system prompt for subagent
|
||||||
systemPrompt := `You are a subagent. Complete the given task independently and report the result.
|
systemPrompt := `You are a subagent. Complete the given task independently and report the result.
|
||||||
|
|
@ -181,12 +173,7 @@ After completing the task, provide a clear summary of what was done.`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Notify conductor that the subagent is starting
|
// Notify conductor that the subagent is starting
|
||||||
sm.broadcaster.Publish(orch.Event{
|
sm.reporter.ReportConversation("conductor", task.ID, task.Task)
|
||||||
Type: "conversation",
|
|
||||||
From: "conductor",
|
|
||||||
To: task.ID,
|
|
||||||
Text: task.Task,
|
|
||||||
})
|
|
||||||
|
|
||||||
loopResult, err := RunToolLoop(ctx, ToolLoopConfig{
|
loopResult, err := RunToolLoop(ctx, ToolLoopConfig{
|
||||||
Provider: sm.provider,
|
Provider: sm.provider,
|
||||||
|
|
@ -194,14 +181,8 @@ After completing the task, provide a clear summary of what was done.`
|
||||||
Tools: tools,
|
Tools: tools,
|
||||||
MaxIterations: maxIter,
|
MaxIterations: maxIter,
|
||||||
LLMOptions: llmOptions,
|
LLMOptions: llmOptions,
|
||||||
OnStateChange: func(state, tool string) {
|
Reporter: sm.reporter,
|
||||||
sm.broadcaster.Publish(orch.Event{
|
AgentID: task.ID,
|
||||||
Type: "agent_state",
|
|
||||||
ID: task.ID,
|
|
||||||
State: state,
|
|
||||||
Tool: tool,
|
|
||||||
})
|
|
||||||
},
|
|
||||||
}, messages, task.OriginChannel, task.OriginChatID)
|
}, messages, task.OriginChannel, task.OriginChatID)
|
||||||
|
|
||||||
sm.mu.Lock()
|
sm.mu.Lock()
|
||||||
|
|
@ -224,11 +205,7 @@ After completing the task, provide a clear summary of what was done.`
|
||||||
task.Result = "Task cancelled during execution"
|
task.Result = "Task cancelled during execution"
|
||||||
gcReason = "cancelled"
|
gcReason = "cancelled"
|
||||||
}
|
}
|
||||||
sm.broadcaster.Publish(orch.Event{
|
sm.reporter.ReportGC(task.ID, gcReason)
|
||||||
Type: "agent_gc",
|
|
||||||
ID: task.ID,
|
|
||||||
Reason: gcReason,
|
|
||||||
})
|
|
||||||
result = &ToolResult{
|
result = &ToolResult{
|
||||||
ForLLM: task.Result,
|
ForLLM: task.Result,
|
||||||
ForUser: "",
|
ForUser: "",
|
||||||
|
|
@ -241,17 +218,8 @@ After completing the task, provide a clear summary of what was done.`
|
||||||
task.Status = "completed"
|
task.Status = "completed"
|
||||||
task.Result = loopResult.Content
|
task.Result = loopResult.Content
|
||||||
// Notify conductor of the result
|
// Notify conductor of the result
|
||||||
sm.broadcaster.Publish(orch.Event{
|
sm.reporter.ReportConversation(task.ID, "conductor", loopResult.Content)
|
||||||
Type: "conversation",
|
sm.reporter.ReportGC(task.ID, "completed")
|
||||||
From: task.ID,
|
|
||||||
To: "conductor",
|
|
||||||
Text: loopResult.Content,
|
|
||||||
})
|
|
||||||
sm.broadcaster.Publish(orch.Event{
|
|
||||||
Type: "agent_gc",
|
|
||||||
ID: task.ID,
|
|
||||||
Reason: "completed",
|
|
||||||
})
|
|
||||||
result = &ToolResult{
|
result = &ToolResult{
|
||||||
ForLLM: fmt.Sprintf(
|
ForLLM: fmt.Sprintf(
|
||||||
"Subagent '%s' completed (iterations: %d): %s",
|
"Subagent '%s' completed (iterations: %d): %s",
|
||||||
|
|
|
||||||
|
|
@ -6,6 +6,7 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/sipeed/picoclaw/pkg/bus"
|
"github.com/sipeed/picoclaw/pkg/bus"
|
||||||
|
"github.com/sipeed/picoclaw/pkg/orch"
|
||||||
"github.com/sipeed/picoclaw/pkg/providers"
|
"github.com/sipeed/picoclaw/pkg/providers"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -47,7 +48,7 @@ func (m *MockLLMProvider) GetContextWindow() int {
|
||||||
|
|
||||||
func TestSubagentManager_SetLLMOptions_AppliesToRunToolLoop(t *testing.T) {
|
func TestSubagentManager_SetLLMOptions_AppliesToRunToolLoop(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
manager.SetLLMOptions(2048, 0.6)
|
manager.SetLLMOptions(2048, 0.6)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
tool.SetContext("cli", "direct")
|
tool.SetContext("cli", "direct")
|
||||||
|
|
@ -74,7 +75,7 @@ func TestSubagentManager_SetLLMOptions_AppliesToRunToolLoop(t *testing.T) {
|
||||||
// TestSubagentTool_Name verifies tool name
|
// TestSubagentTool_Name verifies tool name
|
||||||
func TestSubagentTool_Name(t *testing.T) {
|
func TestSubagentTool_Name(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
if tool.Name() != "subagent" {
|
if tool.Name() != "subagent" {
|
||||||
|
|
@ -85,7 +86,7 @@ func TestSubagentTool_Name(t *testing.T) {
|
||||||
// TestSubagentTool_Description verifies tool description
|
// TestSubagentTool_Description verifies tool description
|
||||||
func TestSubagentTool_Description(t *testing.T) {
|
func TestSubagentTool_Description(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
desc := tool.Description()
|
desc := tool.Description()
|
||||||
|
|
@ -100,7 +101,7 @@ func TestSubagentTool_Description(t *testing.T) {
|
||||||
// TestSubagentTool_Parameters verifies tool parameters schema
|
// TestSubagentTool_Parameters verifies tool parameters schema
|
||||||
func TestSubagentTool_Parameters(t *testing.T) {
|
func TestSubagentTool_Parameters(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
params := tool.Parameters()
|
params := tool.Parameters()
|
||||||
|
|
@ -150,7 +151,7 @@ func TestSubagentTool_Parameters(t *testing.T) {
|
||||||
// TestSubagentTool_SetContext verifies context setting
|
// TestSubagentTool_SetContext verifies context setting
|
||||||
func TestSubagentTool_SetContext(t *testing.T) {
|
func TestSubagentTool_SetContext(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
tool.SetContext("test-channel", "test-chat")
|
tool.SetContext("test-channel", "test-chat")
|
||||||
|
|
@ -164,7 +165,7 @@ func TestSubagentTool_SetContext(t *testing.T) {
|
||||||
func TestSubagentTool_Execute_Success(t *testing.T) {
|
func TestSubagentTool_Execute_Success(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
msgBus := bus.NewMessageBus()
|
msgBus := bus.NewMessageBus()
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
tool.SetContext("telegram", "chat-123")
|
tool.SetContext("telegram", "chat-123")
|
||||||
|
|
||||||
|
|
@ -220,7 +221,7 @@ func TestSubagentTool_Execute_Success(t *testing.T) {
|
||||||
func TestSubagentTool_Execute_NoLabel(t *testing.T) {
|
func TestSubagentTool_Execute_NoLabel(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
msgBus := bus.NewMessageBus()
|
msgBus := bus.NewMessageBus()
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
@ -243,7 +244,7 @@ func TestSubagentTool_Execute_NoLabel(t *testing.T) {
|
||||||
// TestSubagentTool_Execute_MissingTask tests error handling for missing task
|
// TestSubagentTool_Execute_MissingTask tests error handling for missing task
|
||||||
func TestSubagentTool_Execute_MissingTask(t *testing.T) {
|
func TestSubagentTool_Execute_MissingTask(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", nil, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
@ -294,7 +295,7 @@ func TestSubagentTool_Execute_NilManager(t *testing.T) {
|
||||||
func TestSubagentTool_Execute_ContextPassing(t *testing.T) {
|
func TestSubagentTool_Execute_ContextPassing(t *testing.T) {
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
msgBus := bus.NewMessageBus()
|
msgBus := bus.NewMessageBus()
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
// Set context
|
// Set context
|
||||||
|
|
@ -323,7 +324,7 @@ func TestSubagentTool_ForUserTruncation(t *testing.T) {
|
||||||
// Create a mock provider that returns very long content
|
// Create a mock provider that returns very long content
|
||||||
provider := &MockLLMProvider{}
|
provider := &MockLLMProvider{}
|
||||||
msgBus := bus.NewMessageBus()
|
msgBus := bus.NewMessageBus()
|
||||||
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus)
|
manager := NewSubagentManager(provider, "test-model", "/tmp/test", msgBus, orch.Noop)
|
||||||
tool := NewSubagentTool(manager)
|
tool := NewSubagentTool(manager)
|
||||||
|
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
|
||||||
|
|
@ -12,6 +12,7 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"github.com/sipeed/picoclaw/pkg/logger"
|
"github.com/sipeed/picoclaw/pkg/logger"
|
||||||
|
"github.com/sipeed/picoclaw/pkg/orch"
|
||||||
"github.com/sipeed/picoclaw/pkg/providers"
|
"github.com/sipeed/picoclaw/pkg/providers"
|
||||||
"github.com/sipeed/picoclaw/pkg/utils"
|
"github.com/sipeed/picoclaw/pkg/utils"
|
||||||
)
|
)
|
||||||
|
|
@ -23,11 +24,12 @@ type ToolLoopConfig struct {
|
||||||
Tools *ToolRegistry
|
Tools *ToolRegistry
|
||||||
MaxIterations int
|
MaxIterations int
|
||||||
LLMOptions map[string]any
|
LLMOptions map[string]any
|
||||||
// OnStateChange is an optional hook for UI feedback.
|
// Reporter and AgentID replace the old OnStateChange func.
|
||||||
// Called with ("waiting","") before each LLM call and
|
// Reporter is called with ReportStateChange("waiting","") before each LLM
|
||||||
// ("toolcall", toolName) when each tool starts executing.
|
// call and ReportStateChange("toolcall", toolName) when each tool starts.
|
||||||
// nil is safe to pass.
|
// Pass nil or orch.Noop to disable. nil is treated as orch.Noop internally.
|
||||||
OnStateChange func(state, tool string)
|
Reporter orch.AgentReporter
|
||||||
|
AgentID string
|
||||||
}
|
}
|
||||||
|
|
||||||
// ToolLoopResult contains the result of running the tool loop.
|
// ToolLoopResult contains the result of running the tool loop.
|
||||||
|
|
@ -44,6 +46,11 @@ func RunToolLoop(
|
||||||
messages []providers.Message,
|
messages []providers.Message,
|
||||||
channel, chatID string,
|
channel, chatID string,
|
||||||
) (*ToolLoopResult, error) {
|
) (*ToolLoopResult, error) {
|
||||||
|
reporter := config.Reporter
|
||||||
|
if reporter == nil {
|
||||||
|
reporter = orch.Noop
|
||||||
|
}
|
||||||
|
|
||||||
iteration := 0
|
iteration := 0
|
||||||
var finalContent string
|
var finalContent string
|
||||||
|
|
||||||
|
|
@ -68,9 +75,7 @@ func RunToolLoop(
|
||||||
llmOpts = map[string]any{}
|
llmOpts = map[string]any{}
|
||||||
}
|
}
|
||||||
// 3. Call LLM (hook: waiting for response)
|
// 3. Call LLM (hook: waiting for response)
|
||||||
if config.OnStateChange != nil {
|
reporter.ReportStateChange(config.AgentID, "waiting", "")
|
||||||
config.OnStateChange("waiting", "")
|
|
||||||
}
|
|
||||||
response, err := config.Provider.Chat(ctx, messages, providerToolDefs, config.Model, llmOpts)
|
response, err := config.Provider.Chat(ctx, messages, providerToolDefs, config.Model, llmOpts)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
logger.ErrorCF("toolloop", "LLM call failed",
|
logger.ErrorCF("toolloop", "LLM call failed",
|
||||||
|
|
@ -138,9 +143,7 @@ func RunToolLoop(
|
||||||
"tool": tc.Name,
|
"tool": tc.Name,
|
||||||
"iteration": iteration,
|
"iteration": iteration,
|
||||||
})
|
})
|
||||||
if config.OnStateChange != nil {
|
reporter.ReportStateChange(config.AgentID, "toolcall", tc.Name)
|
||||||
config.OnStateChange("toolcall", tc.Name)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Execute tool (no async callback for subagents - they run independently)
|
// Execute tool (no async callback for subagents - they run independently)
|
||||||
var toolResult *ToolResult
|
var toolResult *ToolResult
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue