refactor(agent): improve SubTurn error handling and logging
- Fix context cancellation check order in concurrency timeout - Add structured logging for panic recovery - Replace println with proper logger for channel full warning - Simplify tool registry initialization logic - Remove unused ErrConcurrencyLimitExceeded error
This commit is contained in:
parent
a26a7db7d2
commit
2fec249be1
1 changed files with 27 additions and 24 deletions
|
|
@ -24,10 +24,9 @@ const (
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
ErrDepthLimitExceeded = errors.New("sub-turn depth limit exceeded")
|
ErrDepthLimitExceeded = errors.New("sub-turn depth limit exceeded")
|
||||||
ErrInvalidSubTurnConfig = errors.New("invalid sub-turn config")
|
ErrInvalidSubTurnConfig = errors.New("invalid sub-turn config")
|
||||||
ErrConcurrencyLimitExceeded = errors.New("sub-turn concurrency limit exceeded")
|
ErrConcurrencyTimeout = errors.New("timeout waiting for concurrency slot")
|
||||||
ErrConcurrencyTimeout = errors.New("timeout waiting for concurrency slot")
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// ====================== SubTurn Config ======================
|
// ====================== SubTurn Config ======================
|
||||||
|
|
@ -57,7 +56,6 @@ var (
|
||||||
// result, err := SpawnSubTurn(ctx, cfg)
|
// result, err := SpawnSubTurn(ctx, cfg)
|
||||||
// // Result also available in parent's pendingResults channel
|
// // Result also available in parent's pendingResults channel
|
||||||
// // Parent turn will poll and process it in a later iteration
|
// // Parent turn will poll and process it in a later iteration
|
||||||
//
|
|
||||||
type SubTurnConfig struct {
|
type SubTurnConfig struct {
|
||||||
Model string
|
Model string
|
||||||
Tools []tools.Tool
|
Tools []tools.Tool
|
||||||
|
|
@ -204,12 +202,13 @@ func spawnSubTurn(ctx context.Context, al *AgentLoop, parentTS *turnState, cfg S
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
case <-timeoutCtx.Done():
|
case <-timeoutCtx.Done():
|
||||||
// Check if it was a timeout or parent context cancellation
|
// Check parent context first - if it was cancelled, propagate that error
|
||||||
if timeoutCtx.Err() == context.DeadlineExceeded {
|
if ctx.Err() != nil {
|
||||||
return nil, fmt.Errorf("%w: all %d slots occupied for %v",
|
return nil, ctx.Err()
|
||||||
ErrConcurrencyTimeout, maxConcurrentSubTurns, concurrencyTimeout)
|
|
||||||
}
|
}
|
||||||
return nil, ctx.Err()
|
// Otherwise it's our timeout
|
||||||
|
return nil, fmt.Errorf("%w: all %d slots occupied for %v",
|
||||||
|
ErrConcurrencyTimeout, maxConcurrentSubTurns, concurrencyTimeout)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -259,6 +258,11 @@ func spawnSubTurn(ctx context.Context, al *AgentLoop, parentTS *turnState, cfg S
|
||||||
defer func() {
|
defer func() {
|
||||||
if r := recover(); r != nil {
|
if r := recover(); r != nil {
|
||||||
err = fmt.Errorf("subturn panicked: %v", r)
|
err = fmt.Errorf("subturn panicked: %v", r)
|
||||||
|
logger.ErrorCF("subturn", "SubTurn panicked", map[string]any{
|
||||||
|
"child_id": childID,
|
||||||
|
"parent_id": parentTS.turnID,
|
||||||
|
"panic": r,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// 7. Result Delivery Strategy (Async vs Sync)
|
// 7. Result Delivery Strategy (Async vs Sync)
|
||||||
|
|
@ -351,7 +355,10 @@ func deliverSubTurnResult(parentTS *turnState, childID string, result *tools.Too
|
||||||
})
|
})
|
||||||
default:
|
default:
|
||||||
// Channel is full - treat as orphan result
|
// Channel is full - treat as orphan result
|
||||||
fmt.Println("[SubTurn] warning: pendingResults channel full")
|
logger.WarnCF("subturn", "pendingResults channel full", map[string]any{
|
||||||
|
"parent_id": parentTS.turnID,
|
||||||
|
"child_id": childID,
|
||||||
|
})
|
||||||
if result != nil {
|
if result != nil {
|
||||||
MockEventBus.Emit(SubTurnOrphanResultEvent{
|
MockEventBus.Emit(SubTurnOrphanResultEvent{
|
||||||
ParentID: parentTS.turnID,
|
ParentID: parentTS.turnID,
|
||||||
|
|
@ -378,20 +385,16 @@ func runTurn(ctx context.Context, al *AgentLoop, ts *turnState, cfg SubTurnConfi
|
||||||
// ephemeral session store and tool registry.
|
// ephemeral session store and tool registry.
|
||||||
parentAgent := al.GetRegistry().GetDefaultAgent()
|
parentAgent := al.GetRegistry().GetDefaultAgent()
|
||||||
|
|
||||||
var toolRegistry *tools.ToolRegistry
|
// Determine which tools to use: explicit config or inherit from parent
|
||||||
if len(cfg.Tools) > 0 {
|
toolRegistry := tools.NewToolRegistry()
|
||||||
// Use explicitly provided tools
|
toolsToRegister := cfg.Tools
|
||||||
toolRegistry = tools.NewToolRegistry()
|
if len(toolsToRegister) == 0 {
|
||||||
for _, t := range cfg.Tools {
|
toolsToRegister = parentAgent.Tools.GetAll()
|
||||||
toolRegistry.Register(t)
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
// Inherit tools from parent agent when cfg.Tools is nil or empty
|
|
||||||
toolRegistry = tools.NewToolRegistry()
|
|
||||||
for _, t := range parentAgent.Tools.GetAll() {
|
|
||||||
toolRegistry.Register(t)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
for _, t := range toolsToRegister {
|
||||||
|
toolRegistry.Register(t)
|
||||||
|
}
|
||||||
|
|
||||||
childAgent := &AgentInstance{
|
childAgent := &AgentInstance{
|
||||||
ID: ts.turnID,
|
ID: ts.turnID,
|
||||||
Model: cfg.Model,
|
Model: cfg.Model,
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue