Merge pull request #1476 from trheyi/main
Enhance robot execution with language model integration and logging
This commit is contained in:
commit
b5b6118010
16 changed files with 339 additions and 47 deletions
|
|
@ -167,7 +167,7 @@ func loadRobotFromDB(memberID string) (*types.Robot, error) {
|
|||
"id", "member_id", "team_id", "display_name", "bio",
|
||||
"system_prompt", "robot_status", "autonomous_mode",
|
||||
"robot_config", "robot_email", "agents", "mcp_servers",
|
||||
"manager_id",
|
||||
"manager_id", "language_model",
|
||||
},
|
||||
Wheres: []model.QueryWhere{
|
||||
{Column: "member_id", Value: memberID},
|
||||
|
|
@ -230,6 +230,7 @@ func listRobotsFromDB(query *ListQuery) (*ListResult, error) {
|
|||
"id", "member_id", "team_id", "display_name", "bio",
|
||||
"system_prompt", "robot_status", "autonomous_mode",
|
||||
"robot_config", "robot_email", "agents", "mcp_servers",
|
||||
"language_model",
|
||||
},
|
||||
Wheres: wheres,
|
||||
Orders: orders,
|
||||
|
|
|
|||
1
agent/robot/cache/load.go
vendored
1
agent/robot/cache/load.go
vendored
|
|
@ -27,6 +27,7 @@ var memberFields = []interface{}{
|
|||
"agents",
|
||||
"mcp_servers",
|
||||
"manager_id",
|
||||
"language_model",
|
||||
}
|
||||
|
||||
// SetMemberModel sets the member model name
|
||||
|
|
|
|||
|
|
@ -30,6 +30,13 @@ type AgentCaller struct {
|
|||
// ChatID is used for multi-turn conversations to maintain session state
|
||||
// If empty, each call is independent (no history)
|
||||
ChatID string
|
||||
|
||||
// Connector overrides the assistant's default LLM connector (from Robot.LanguageModel).
|
||||
// When non-empty, passed as opts.Connector to ast.Stream so the agent uses the Robot's model.
|
||||
Connector string
|
||||
|
||||
// log is an optional structured logger; when set, Call emits agent-call logs.
|
||||
log *execLogger
|
||||
}
|
||||
|
||||
// NewAgentCaller creates a new AgentCaller with default settings (single-call mode)
|
||||
|
|
@ -78,16 +85,13 @@ func (r *CallResult) GetText() string {
|
|||
if r.Content != "" {
|
||||
return r.Content
|
||||
}
|
||||
// If Next is a string, return it
|
||||
if s, ok := r.Next.(string); ok {
|
||||
return s
|
||||
}
|
||||
// If Next has a "content" field, return it
|
||||
if m, ok := r.Next.(map[string]interface{}); ok {
|
||||
if content, ok := m["content"].(string); ok {
|
||||
return content
|
||||
}
|
||||
// Also check "data" field (common pattern in Next hook)
|
||||
if data, ok := m["data"].(map[string]interface{}); ok {
|
||||
if content, ok := data["content"].(string); ok {
|
||||
return content
|
||||
|
|
@ -103,10 +107,8 @@ func (r *CallResult) GetText() string {
|
|||
// 2. Content parsed using gou/text.ExtractJSON (fault-tolerant)
|
||||
// Returns the parsed data and any error
|
||||
func (r *CallResult) GetJSON() (map[string]interface{}, error) {
|
||||
// Try Next hook data first
|
||||
if r.Next != nil {
|
||||
if m, ok := r.Next.(map[string]interface{}); ok {
|
||||
// Check for "data" wrapper (common in Next hook)
|
||||
if data, ok := m["data"].(map[string]interface{}); ok {
|
||||
return data, nil
|
||||
}
|
||||
|
|
@ -114,7 +116,6 @@ func (r *CallResult) GetJSON() (map[string]interface{}, error) {
|
|||
}
|
||||
}
|
||||
|
||||
// Try parsing Content using gou/text (handles markdown blocks, JSON, YAML)
|
||||
if r.Content != "" {
|
||||
data := text.ExtractJSON(r.Content)
|
||||
if data != nil {
|
||||
|
|
@ -174,6 +175,7 @@ func (c *AgentCaller) Call(ctx *robottypes.Context, assistantID string, messages
|
|||
History: c.SkipHistory,
|
||||
Search: c.SkipSearch,
|
||||
},
|
||||
Connector: c.Connector,
|
||||
}
|
||||
|
||||
// Convert robot context to agent context
|
||||
|
|
@ -203,6 +205,10 @@ func (c *AgentCaller) Call(ctx *robottypes.Context, assistantID string, messages
|
|||
}
|
||||
}
|
||||
|
||||
if c.log != nil {
|
||||
c.log.logAgentCall(assistantID, result)
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -53,6 +53,7 @@ func (e *Executor) RunDelivery(ctx *robottypes.Context, exec *robottypes.Executi
|
|||
|
||||
// Call Delivery Agent
|
||||
caller := NewAgentCaller()
|
||||
caller.Connector = robot.LanguageModel
|
||||
result, err := caller.CallWithMessages(ctx, agentID, userContent)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delivery agent (%s) call failed: %w", agentID, err)
|
||||
|
|
|
|||
|
|
@ -84,6 +84,7 @@ func (e *Executor) RunGoals(ctx *robottypes.Context, exec *robottypes.Execution,
|
|||
|
||||
// Call agent
|
||||
caller := NewAgentCaller()
|
||||
caller.Connector = robot.LanguageModel
|
||||
result, err := caller.CallWithMessages(ctx, agentID, userContent)
|
||||
if err != nil {
|
||||
return fmt.Errorf("goals agent (%s) call failed: %w", agentID, err)
|
||||
|
|
|
|||
|
|
@ -179,6 +179,10 @@ func (f *InputFormatter) FormatAvailableResourcesWithLocale(robot *robottypes.Ro
|
|||
if description != "" {
|
||||
sb.WriteString(fmt.Sprintf(" - %s\n", description))
|
||||
}
|
||||
if ast.Capabilities != "" {
|
||||
capabilities := i18n.Translate(agentID, locale, ast.Capabilities).(string)
|
||||
sb.WriteString(fmt.Sprintf(" - **Capabilities**: %s\n", capabilities))
|
||||
}
|
||||
}
|
||||
sb.WriteString("\n")
|
||||
}
|
||||
|
|
|
|||
|
|
@ -53,6 +53,7 @@ func (e *Executor) RunInspiration(ctx *robottypes.Context, exec *robottypes.Exec
|
|||
|
||||
// Call agent
|
||||
caller := NewAgentCaller()
|
||||
caller.Connector = robot.LanguageModel
|
||||
result, err := caller.CallWithMessages(ctx, agentID, userContent)
|
||||
if err != nil {
|
||||
return fmt.Errorf("inspiration agent (%s) call failed: %w", agentID, err)
|
||||
|
|
|
|||
253
agent/robot/executor/standard/log.go
Normal file
253
agent/robot/executor/standard/log.go
Normal file
|
|
@ -0,0 +1,253 @@
|
|||
package standard
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/yaoapp/kun/log"
|
||||
"github.com/yaoapp/yao/config"
|
||||
|
||||
robottypes "github.com/yaoapp/yao/agent/robot/types"
|
||||
)
|
||||
|
||||
// execLogger provides structured, developer-facing logging for a single Robot execution.
|
||||
// Each execution (Executor.ExecuteWithControl) creates one instance; it is passed to
|
||||
// RunTasks (P2) and Runner (P3) so every log line carries the same identity.
|
||||
//
|
||||
// Output routing:
|
||||
// - development mode (config.IsDevelopment): human-readable console via fmt.Printf
|
||||
// - production mode: structured fields via kun/log
|
||||
type execLogger struct {
|
||||
robot *robottypes.Robot
|
||||
execID string
|
||||
}
|
||||
|
||||
func newExecLogger(robot *robottypes.Robot, execID string) *execLogger {
|
||||
return &execLogger{robot: robot, execID: execID}
|
||||
}
|
||||
|
||||
func (l *execLogger) robotID() string {
|
||||
if l.robot != nil {
|
||||
return l.robot.MemberID
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (l *execLogger) connector() string {
|
||||
if l.robot != nil {
|
||||
return l.robot.LanguageModel
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// P2: Task Overview — called once after RunTasks successfully generates tasks
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (l *execLogger) logTaskOverview(tasks []robottypes.Task) {
|
||||
if config.IsDevelopment() {
|
||||
l.devTaskOverview(tasks)
|
||||
}
|
||||
// Always emit structured log (Info level, hidden in prod unless needed)
|
||||
log.With(log.F{
|
||||
"robot_id": l.robotID(),
|
||||
"execution_id": l.execID,
|
||||
"phase": "tasks",
|
||||
"task_count": len(tasks),
|
||||
"language_model": l.connector(),
|
||||
}).Info("P2 task overview: %d tasks generated", len(tasks))
|
||||
}
|
||||
|
||||
func (l *execLogger) devTaskOverview(tasks []robottypes.Task) {
|
||||
var sb strings.Builder
|
||||
sb.WriteString(fmt.Sprintf("%s ══════ P2: Task Overview ══════\n", l.prefix()))
|
||||
if l.connector() != "" {
|
||||
sb.WriteString(fmt.Sprintf(" Language Model: %s\n", l.connector()))
|
||||
}
|
||||
for i, t := range tasks {
|
||||
desc := t.Description
|
||||
if desc == "" && len(t.Messages) > 0 {
|
||||
if s, ok := t.Messages[0].GetContentAsString(); ok {
|
||||
desc = s
|
||||
}
|
||||
}
|
||||
desc = truncate(desc, 80)
|
||||
sb.WriteString(fmt.Sprintf(" #%d %s [%s:%s] %q\n", i+1, t.ID, t.ExecutorType, t.ExecutorID, desc))
|
||||
}
|
||||
sb.WriteString(fmt.Sprintf(" Total: %d tasks\n", len(tasks)))
|
||||
fmt.Print(sb.String())
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// P3: Task Input — called before each task execution with the full prompt
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (l *execLogger) logTaskInput(task *robottypes.Task, prompt string) {
|
||||
if config.IsDevelopment() {
|
||||
l.devTaskInput(task, prompt)
|
||||
}
|
||||
log.With(log.F{
|
||||
"robot_id": l.robotID(),
|
||||
"execution_id": l.execID,
|
||||
"task_id": task.ID,
|
||||
"executor_type": string(task.ExecutorType),
|
||||
"executor_id": task.ExecutorID,
|
||||
"prompt_len": len(prompt),
|
||||
"language_model": l.connector(),
|
||||
}).Info("Task input: %s [%s]", task.ID, task.ExecutorID)
|
||||
}
|
||||
|
||||
func (l *execLogger) devTaskInput(task *robottypes.Task, prompt string) {
|
||||
sep := strings.Repeat("─", 40)
|
||||
fmt.Printf("%s ▶ Task %s [%s:%s]\n", l.prefix(), task.ID, task.ExecutorType, task.ExecutorID)
|
||||
fmt.Printf(" Prompt (%d chars):\n %s\n%s\n %s\n",
|
||||
len(prompt), sep, indentText(prompt, " "), sep)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// P3: Task Output — called after each task execution with the result
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (l *execLogger) logTaskOutput(task *robottypes.Task, result *robottypes.TaskResult) {
|
||||
if config.IsDevelopment() {
|
||||
l.devTaskOutput(task, result)
|
||||
}
|
||||
|
||||
fields := log.F{
|
||||
"robot_id": l.robotID(),
|
||||
"execution_id": l.execID,
|
||||
"task_id": result.TaskID,
|
||||
"success": result.Success,
|
||||
"duration_ms": result.Duration,
|
||||
"language_model": l.connector(),
|
||||
}
|
||||
if result.Output != nil {
|
||||
fields["output_type"] = fmt.Sprintf("%T", result.Output)
|
||||
fields["output_len"] = outputLen(result.Output)
|
||||
}
|
||||
if result.Error != "" {
|
||||
fields["error"] = result.Error
|
||||
}
|
||||
if result.Success {
|
||||
log.With(fields).Info("Task completed: %s (%dms)", result.TaskID, result.Duration)
|
||||
} else {
|
||||
log.With(fields).Warn("Task failed: %s (%dms) %s", result.TaskID, result.Duration, result.Error)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *execLogger) devTaskOutput(task *robottypes.Task, result *robottypes.TaskResult) {
|
||||
if result.Success {
|
||||
fmt.Printf("%s ✓ Task %s completed (%dms)\n", l.prefix(), result.TaskID, result.Duration)
|
||||
fmt.Printf(" Output: %s\n", outputSummary(result.Output))
|
||||
} else {
|
||||
fmt.Printf("%s ✗ Task %s failed (%dms)\n", l.prefix(), result.TaskID, result.Duration)
|
||||
fmt.Printf(" Error: %s\n", result.Error)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Agent Call — called after every AgentCaller.Call returns
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (l *execLogger) logAgentCall(agentID string, result *CallResult) {
|
||||
if result == nil {
|
||||
return
|
||||
}
|
||||
if config.IsDevelopment() {
|
||||
l.devAgentCall(agentID, result)
|
||||
}
|
||||
|
||||
fields := log.F{
|
||||
"robot_id": l.robotID(),
|
||||
"execution_id": l.execID,
|
||||
"agent_id": agentID,
|
||||
"content_len": len(result.Content),
|
||||
"language_model": l.connector(),
|
||||
}
|
||||
if result.Next != nil {
|
||||
fields["next_type"] = fmt.Sprintf("%T", result.Next)
|
||||
fields["next_len"] = outputLen(result.Next)
|
||||
}
|
||||
log.With(fields).Info("Agent call: %s (content=%d, next=%T)", agentID, len(result.Content), result.Next)
|
||||
}
|
||||
|
||||
func (l *execLogger) devAgentCall(agentID string, result *CallResult) {
|
||||
nextInfo := "<nil>"
|
||||
if result.Next != nil {
|
||||
nextInfo = fmt.Sprintf("%T(len=%d)", result.Next, outputLen(result.Next))
|
||||
}
|
||||
fmt.Printf("%s Agent(%s) → Content(len=%d) Next=%s\n",
|
||||
l.prefix(), agentID, len(result.Content), nextInfo)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Helpers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
func (l *execLogger) prefix() string {
|
||||
if l.connector() != "" {
|
||||
return fmt.Sprintf("[robot:%s|exec:%s|model:%s]", l.robotID(), l.execID, l.connector())
|
||||
}
|
||||
return fmt.Sprintf("[robot:%s|exec:%s]", l.robotID(), l.execID)
|
||||
}
|
||||
|
||||
func truncate(s string, maxLen int) string {
|
||||
if len(s) <= maxLen {
|
||||
return s
|
||||
}
|
||||
return s[:maxLen] + "..."
|
||||
}
|
||||
|
||||
func indentText(s string, prefix string) string {
|
||||
lines := strings.Split(s, "\n")
|
||||
for i, line := range lines {
|
||||
lines[i] = prefix + line
|
||||
}
|
||||
return strings.Join(lines, "\n")
|
||||
}
|
||||
|
||||
func outputSummary(v interface{}) string {
|
||||
if v == nil {
|
||||
return "<nil>"
|
||||
}
|
||||
switch val := v.(type) {
|
||||
case string:
|
||||
if len(val) > 500 {
|
||||
return fmt.Sprintf("string(len=%d) %s...", len(val), val[:500])
|
||||
}
|
||||
return fmt.Sprintf("string(len=%d) %s", len(val), val)
|
||||
case map[string]interface{}:
|
||||
keys := make([]string, 0, len(val))
|
||||
for k := range val {
|
||||
keys = append(keys, k)
|
||||
}
|
||||
return fmt.Sprintf("map{%s}", strings.Join(keys, ", "))
|
||||
default:
|
||||
raw, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
return fmt.Sprintf("%T(marshal-error)", v)
|
||||
}
|
||||
s := string(raw)
|
||||
if len(s) > 500 {
|
||||
return fmt.Sprintf("%T(len=%d) %s...", v, len(s), s[:500])
|
||||
}
|
||||
return fmt.Sprintf("%T %s", v, s)
|
||||
}
|
||||
}
|
||||
|
||||
func outputLen(v interface{}) int {
|
||||
if v == nil {
|
||||
return 0
|
||||
}
|
||||
switch val := v.(type) {
|
||||
case string:
|
||||
return len(val)
|
||||
default:
|
||||
raw, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
return 0
|
||||
}
|
||||
return len(raw)
|
||||
}
|
||||
}
|
||||
|
|
@ -66,7 +66,7 @@ func (e *Executor) RunExecution(ctx *robottypes.Context, exec *robottypes.Execut
|
|||
}
|
||||
|
||||
// Create task runner with execution-level chatID (§8.4)
|
||||
runner := NewRunner(ctx, robot, config, exec.ChatID)
|
||||
runner := NewRunner(ctx, robot, config, exec.ChatID, exec.ID)
|
||||
|
||||
// Execute tasks sequentially from startIndex
|
||||
for i := startIndex; i < len(exec.Tasks); i++ {
|
||||
|
|
|
|||
|
|
@ -18,15 +18,17 @@ type Runner struct {
|
|||
robot *robottypes.Robot
|
||||
config *RunConfig
|
||||
chatID string // execution-level chatID for conversation persistence (§8.4)
|
||||
log *execLogger
|
||||
}
|
||||
|
||||
// NewRunner creates a new task runner
|
||||
func NewRunner(ctx *robottypes.Context, robot *robottypes.Robot, config *RunConfig, chatID string) *Runner {
|
||||
func NewRunner(ctx *robottypes.Context, robot *robottypes.Robot, config *RunConfig, chatID string, execID string) *Runner {
|
||||
return &Runner{
|
||||
ctx: ctx,
|
||||
robot: robot,
|
||||
config: config,
|
||||
chatID: chatID,
|
||||
log: newExecLogger(robot, execID),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -78,12 +80,14 @@ func (r *Runner) ExecuteTask(task *robottypes.Task, taskCtx *RunnerContext) *rob
|
|||
result.Success = false
|
||||
result.Error = fmt.Sprintf("execution failed: %s", err.Error())
|
||||
result.Duration = time.Since(startTime).Milliseconds()
|
||||
r.log.logTaskOutput(task, result)
|
||||
return result
|
||||
}
|
||||
|
||||
result.Output = output
|
||||
result.Success = true
|
||||
result.Duration = time.Since(startTime).Milliseconds()
|
||||
r.log.logTaskOutput(task, result)
|
||||
return result
|
||||
}
|
||||
|
||||
|
|
@ -93,6 +97,7 @@ func (r *Runner) ExecuteTask(task *robottypes.Task, taskCtx *RunnerContext) *rob
|
|||
result.Success = false
|
||||
result.Error = err.Error()
|
||||
result.Duration = time.Since(startTime).Milliseconds()
|
||||
r.log.logTaskOutput(task, result)
|
||||
return result
|
||||
}
|
||||
|
||||
|
|
@ -106,6 +111,7 @@ func (r *Runner) ExecuteTask(task *robottypes.Task, taskCtx *RunnerContext) *rob
|
|||
result.InputQuestion = question
|
||||
}
|
||||
|
||||
r.log.logTaskOutput(task, result)
|
||||
return result
|
||||
}
|
||||
|
||||
|
|
@ -124,15 +130,10 @@ func (r *Runner) executeNonAssistantTask(task *robottypes.Task, taskCtx *RunnerC
|
|||
// executeAssistantTask executes an assistant task with a single conversation turn.
|
||||
// Returns the extracted output, the raw CallResult (for need_input detection), and any error.
|
||||
func (r *Runner) executeAssistantTask(task *robottypes.Task, taskCtx *RunnerContext) (interface{}, *CallResult, error) {
|
||||
chatID := r.chatID
|
||||
if chatID == "" {
|
||||
chatID = fmt.Sprintf("robot-%s-task-%s", r.robot.MemberID, task.ID)
|
||||
}
|
||||
conv := NewConversation(task.ExecutorID, chatID, 1)
|
||||
|
||||
if taskCtx.SystemPrompt != "" {
|
||||
conv.WithSystemPrompt(taskCtx.SystemPrompt)
|
||||
}
|
||||
caller := NewAgentCaller()
|
||||
caller.log = r.log
|
||||
caller.Connector = r.robot.LanguageModel
|
||||
caller.ChatID = r.chatID
|
||||
|
||||
messages := r.BuildAssistantMessages(task, taskCtx)
|
||||
input := r.FormatMessagesAsText(messages)
|
||||
|
|
@ -141,13 +142,19 @@ func (r *Runner) executeAssistantTask(task *robottypes.Task, taskCtx *RunnerCont
|
|||
return nil, nil, fmt.Errorf("no valid input messages for task %s", task.ID)
|
||||
}
|
||||
|
||||
turnResult, err := conv.Turn(r.ctx, input)
|
||||
if taskCtx.SystemPrompt != "" {
|
||||
input = "## Context\n\n" + taskCtx.SystemPrompt + "\n\n## Task\n\n" + input
|
||||
}
|
||||
|
||||
r.log.logTaskInput(task, input)
|
||||
|
||||
result, err := caller.CallWithMessages(r.ctx, task.ExecutorID, input)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("assistant call failed: %w", err)
|
||||
}
|
||||
|
||||
output := r.extractOutput(turnResult.Result)
|
||||
return output, turnResult.Result, nil
|
||||
output := r.extractOutput(result)
|
||||
return output, result, nil
|
||||
}
|
||||
|
||||
// detectNeedMoreInfo checks if the assistant's response signals it needs human input.
|
||||
|
|
@ -179,18 +186,20 @@ func detectNeedMoreInfo(result *CallResult) (bool, string) {
|
|||
}
|
||||
|
||||
// extractOutput extracts the output from a CallResult
|
||||
// Priority: Next hook data > LLM Completion content
|
||||
// Next is the agent's formal A2A output (could be string, map, array, number, etc.)
|
||||
// Content is the raw LLM completion text (fallback only when Next is absent)
|
||||
func (r *Runner) extractOutput(result *CallResult) interface{} {
|
||||
if result == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Try to extract structured JSON output
|
||||
if data, err := result.GetJSON(); err == nil {
|
||||
return data
|
||||
if result.Next != nil {
|
||||
return result.Next
|
||||
}
|
||||
|
||||
// Fall back to text content
|
||||
return result.GetText()
|
||||
if result.Content != "" {
|
||||
return result.Content
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ExecuteMCPTask executes a task using an MCP tool
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ func TestRunnerExecuteTask(t *testing.T) {
|
|||
t.Run("executes assistant task successfully", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-001",
|
||||
|
|
@ -58,7 +58,7 @@ func TestRunnerExecuteTask(t *testing.T) {
|
|||
t.Run("returns success without validation for assistant tasks", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-002",
|
||||
|
|
@ -88,7 +88,7 @@ func TestRunnerExecuteTask(t *testing.T) {
|
|||
t.Run("handles empty messages gracefully", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-003",
|
||||
|
|
@ -123,7 +123,7 @@ func TestRunnerBuildTaskContext(t *testing.T) {
|
|||
t.Run("includes previous results in context", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
exec := &types.Execution{
|
||||
ID: "test-exec",
|
||||
|
|
@ -160,7 +160,7 @@ func TestRunnerBuildTaskContext(t *testing.T) {
|
|||
t.Run("handles first task with no previous results", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
exec := &types.Execution{
|
||||
ID: "test-exec",
|
||||
|
|
@ -182,7 +182,7 @@ func TestRunnerBuildTaskContext(t *testing.T) {
|
|||
t.Run("handles bounds check for task index", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
exec := &types.Execution{
|
||||
ID: "test-exec",
|
||||
|
|
@ -215,7 +215,7 @@ func TestRunnerFormatPreviousResultsAsContext(t *testing.T) {
|
|||
t.Run("formats previous results as markdown", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
results := []types.TaskResult{
|
||||
{
|
||||
|
|
@ -247,7 +247,7 @@ func TestRunnerFormatPreviousResultsAsContext(t *testing.T) {
|
|||
t.Run("returns empty string for no results", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
formatted := runner.FormatPreviousResultsAsContext([]types.TaskResult{})
|
||||
|
||||
|
|
@ -268,7 +268,7 @@ func TestRunnerBuildAssistantMessages(t *testing.T) {
|
|||
t.Run("builds messages with task content", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-001",
|
||||
|
|
@ -300,7 +300,7 @@ func TestRunnerBuildAssistantMessages(t *testing.T) {
|
|||
t.Run("includes previous results in messages", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-002",
|
||||
|
|
@ -340,7 +340,7 @@ func TestRunnerFormatMessagesAsText(t *testing.T) {
|
|||
t.Run("formats string content", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
messages := []agentcontext.Message{
|
||||
{Role: agentcontext.RoleUser, Content: "Hello"},
|
||||
|
|
@ -356,7 +356,7 @@ func TestRunnerFormatMessagesAsText(t *testing.T) {
|
|||
t.Run("handles multipart content", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
messages := []agentcontext.Message{
|
||||
{
|
||||
|
|
@ -377,7 +377,7 @@ func TestRunnerFormatMessagesAsText(t *testing.T) {
|
|||
t.Run("handles map content via JSON", func(t *testing.T) {
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
messages := []agentcontext.Message{
|
||||
{
|
||||
|
|
@ -409,7 +409,7 @@ func TestRunnerExecuteNonAssistantTask(t *testing.T) {
|
|||
ctx := types.NewContext(context.Background(), testAuth())
|
||||
robot := createRunnerTestRobot(t)
|
||||
config := standard.DefaultRunConfig()
|
||||
runner := standard.NewRunner(ctx, robot, config, "")
|
||||
runner := standard.NewRunner(ctx, robot, config, "", "test")
|
||||
|
||||
task := &types.Task{
|
||||
ID: "task-unknown",
|
||||
|
|
|
|||
|
|
@ -53,6 +53,8 @@ func (e *Executor) RunTasks(ctx *robottypes.Context, exec *robottypes.Execution,
|
|||
|
||||
// Call agent
|
||||
caller := NewAgentCaller()
|
||||
caller.log = newExecLogger(robot, exec.ID)
|
||||
caller.Connector = robot.LanguageModel
|
||||
result, err := caller.CallWithMessages(ctx, agentID, userContent)
|
||||
if err != nil {
|
||||
return fmt.Errorf("tasks agent (%s) call failed: %w", agentID, err)
|
||||
|
|
@ -83,6 +85,11 @@ func (e *Executor) RunTasks(ctx *robottypes.Context, exec *robottypes.Execution,
|
|||
}
|
||||
|
||||
exec.Tasks = tasks
|
||||
|
||||
// Log task overview for developer observability
|
||||
el := newExecLogger(robot, exec.ID)
|
||||
el.logTaskOverview(tasks)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -431,6 +431,7 @@ func (v *Validator) validateSemantic(task *robottypes.Task, output interface{})
|
|||
|
||||
// Call validation agent
|
||||
caller := NewAgentCaller()
|
||||
caller.Connector = v.robot.LanguageModel
|
||||
result, err := caller.CallWithMessages(v.ctx, validationAgentID, validationPrompt)
|
||||
if err != nil {
|
||||
return &robottypes.ValidationResult{
|
||||
|
|
@ -637,6 +638,7 @@ func (av *robotAgentValidator) Validate(agentID string, output, input, criteria
|
|||
|
||||
// Call agent
|
||||
caller := NewAgentCaller()
|
||||
caller.Connector = av.v.robot.LanguageModel
|
||||
callResult, err := caller.CallWithMessages(av.v.ctx, agentID, string(inputJSON))
|
||||
if err != nil {
|
||||
result.Passed = false
|
||||
|
|
|
|||
|
|
@ -21,7 +21,8 @@ type Robot struct {
|
|||
SystemPrompt string `json:"system_prompt"`
|
||||
Status RobotStatus `json:"robot_status"`
|
||||
AutonomousMode bool `json:"autonomous_mode"`
|
||||
RobotEmail string `json:"robot_email"` // Robot's email address for sending emails
|
||||
RobotEmail string `json:"robot_email"` // Robot's email address for sending emails
|
||||
LanguageModel string `json:"language_model"` // LLM connector override (from __yao.member.language_model)
|
||||
|
||||
// Manager info (from __yao.member)
|
||||
ManagerID string `json:"manager_id"` // Direct manager user_id (who manages this robot)
|
||||
|
|
@ -499,6 +500,7 @@ func NewRobotFromMap(m map[string]interface{}) (*Robot, error) {
|
|||
RobotEmail: getString(m, "robot_email"),
|
||||
ManagerID: getString(m, "manager_id"),
|
||||
ManagerEmail: getString(m, "manager_email"),
|
||||
LanguageModel: getString(m, "language_model"),
|
||||
}
|
||||
|
||||
// Parse robot_status
|
||||
|
|
|
|||
|
|
@ -153,11 +153,13 @@ func BuildCommandWithContinuation(messages []agentContext.Message, opts *Options
|
|||
}
|
||||
|
||||
// Build claude command with all arguments
|
||||
// Append 2>&1 to the claude command so stderr is merged into stdout;
|
||||
// Docker's stdcopy discards the stderr stream, making errors invisible.
|
||||
bashCmd.WriteString("cat << 'INPUTEOF' | claude -p")
|
||||
for _, arg := range claudeArgs {
|
||||
// Quote arguments that might contain special characters
|
||||
bashCmd.WriteString(fmt.Sprintf(" %q", arg))
|
||||
}
|
||||
bashCmd.WriteString(" 2>&1")
|
||||
bashCmd.WriteString("\n")
|
||||
bashCmd.WriteString(string(inputJSONL))
|
||||
bashCmd.WriteString("\nINPUTEOF")
|
||||
|
|
|
|||
|
|
@ -188,7 +188,8 @@ func (e *Executor) Stream(ctx *agentContext.Context, messages []agentContext.Mes
|
|||
|
||||
// Check if we should skip Claude CLI execution
|
||||
// Skip if no prompts, no skills, and no MCP config
|
||||
if e.shouldSkipClaudeCLI() {
|
||||
skipCLI := e.shouldSkipClaudeCLI()
|
||||
if skipCLI {
|
||||
// Return empty response - hooks can use sandbox API to do their work
|
||||
return &agentContext.CompletionResponse{
|
||||
ID: fmt.Sprintf("sandbox-skip-%d", time.Now().UnixNano()),
|
||||
|
|
@ -1115,8 +1116,9 @@ func (e *Executor) parseStream(ctx *agentContext.Context, reader io.Reader, hand
|
|||
}
|
||||
}
|
||||
|
||||
if err := scanner.Err(); err != nil {
|
||||
return nil, fmt.Errorf("error reading stream: %w", err)
|
||||
scanErr := scanner.Err()
|
||||
if scanErr != nil {
|
||||
return nil, fmt.Errorf("error reading stream: %w", scanErr)
|
||||
}
|
||||
|
||||
// Close the last tool loading message if exists
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue