feat: overhaul team delegation, Telegram resilience, and artifact forwarding

Team delegation:
- Unify spawn/subagent/delegate into single spawn tool
- Sibling-aware announce suppression with artifact accumulation
- Fix auto-complete race (isLastDelegation guard)
- Add team tasks list limit (20) with search guidance
- Multi-round orchestration patterns in TEAM.md
- Communication guidance for initial vs follow-up delegations

Telegram resilience:
- Add retrySend wrapper (3 attempts, escalating delay) for network errors
- Fix HTML fallback: strip tags + unescape entities instead of showing raw HTML
- Pre-process HTML tags in LLM output to markdown before conversion pipeline
- Skip caption truncation entirely when > 1024 bytes, send text separately
- Auto-send large images (>5MB) as documents to avoid compression

Artifact forwarding:
- Fix missing ContentType on forwarded media (mimeFromExt for result.Media/ForwardMedia)
- Add deliver parameter to write_file for file attachment delivery
- Extend mimeFromExt with document MIME types

UI: fix regenerate dialog overflow, improve task list layout, delegation detail view

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
viettranxandClaude Opus 4.6 committed 2026-03-01 14:34:36 +07:00
1 parent 7613e07e4f
commit 0e75c21e55
45 files changed
+1246 -678

No files matched your search

+1 -2
View File
@@ -252,8 +252,7 @@ func runGateway() {
subagentMgr.SetAnnounceQueue(announceQueue)
toolsReg.Register(tools.NewSpawnTool(subagentMgr, "default", 0))
toolsReg.Register(tools.NewSubagentTool(subagentMgr, "default", 0))
slog.Info("subagent system enabled", "tools", []string{"spawn", "subagent"})
slog.Info("subagent system enabled", "tools", []string{"spawn"})
}
// Exec approval system — always active (deny patterns + safe bins + configurable ask mode)
+2 -8
View File
@@ -72,11 +72,8 @@ func builtinToolSeedData() []store.BuiltinToolDef {
Metadata: json.RawMessage(`{"config_hint":"Config → Cron"}`),
},
// subagents
{Name: "spawn", DisplayName: "Spawn Subagent", Description: "Spawn an asynchronous background subagent", Category: "subagents", Enabled: true,
Metadata: json.RawMessage(`{"config_hint":"Config → Agents Defaults"}`),
},
{Name: "subagent", DisplayName: "Subagent", Description: "Run a synchronous subagent and wait for result", Category: "subagents", Enabled: true,
// subagents & delegation (unified spawn tool)
{Name: "spawn", DisplayName: "Spawn / Delegate", Description: "Spawn a subagent or delegate to another agent", Category: "subagents", Enabled: true,
Metadata: json.RawMessage(`{"config_hint":"Config → Agents Defaults"}`),
},
@@ -84,9 +81,6 @@ func builtinToolSeedData() []store.BuiltinToolDef {
{Name: "skill_search", DisplayName: "Skill Search", Description: "Search available skills by keyword or description", Category: "skills", Enabled: true},
// delegation
{Name: "delegate", DisplayName: "Delegate", Description: "Delegate a task to another agent", Category: "delegation", Enabled: true,
Requires: []string{"managed_mode", "agent_links"},
},
{Name: "delegate_search", DisplayName: "Delegate Search", Description: "Search for agents to delegate tasks to", Category: "delegation", Enabled: true,
Requires: []string{"managed_mode", "agent_links"},
},
+45 -9
View File
@@ -152,11 +152,21 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
"- Address the group naturally. If the history shows a multi-person conversation, consider the full context before answering."
}
// Delegation announces carry media as ForwardMedia (not deleted, forwarded to output).
// User-uploaded media goes in Media (loaded as images for LLM, then deleted).
var reqMedia, fwdMedia []string
if msg.Metadata["delegation_id"] != "" || msg.Metadata["subagent_id"] != "" {
fwdMedia = msg.Media
} else {
reqMedia = msg.Media
}
// Schedule through main lane (per-session concurrency controlled by maxConcurrent)
outCh := sched.ScheduleWithOpts(ctx, "main", agent.RunRequest{
SessionKey: sessionKey,
Message: msg.Content,
Media: msg.Media,
Media: reqMedia,
ForwardMedia: fwdMedia,
Channel: msg.Channel,
ChatID: msg.ChatID,
PeerKind: peerKind,
@@ -332,6 +342,7 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
outCh := sched.Schedule(ctx, scheduler.LaneSubagent, agent.RunRequest{
SessionKey: sessionKey,
Message: msg.Content,
ForwardMedia: msg.Media,
Channel: origChannel,
ChatID: msg.ChatID,
PeerKind: origPeerKind,
@@ -356,7 +367,8 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
}
// Suppress empty/NO_REPLY (matching TS normalize-reply.ts / tokens.ts).
if outcome.Result.Content == "" || agent.IsSilentReply(outcome.Result.Content) {
isSilent := outcome.Result.Content == "" || agent.IsSilentReply(outcome.Result.Content)
if isSilent && len(outcome.Result.Media) == 0 {
slog.Info("subagent announce: suppressed silent/empty reply",
"subagent", senderID,
"label", label,
@@ -365,11 +377,22 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
}
// Deliver agent's reformulated response to origin channel.
msgBus.PublishOutbound(bus.OutboundMessage{
announceContent := outcome.Result.Content
if isSilent {
announceContent = "" // suppress NO_REPLY text but still send media
}
outMsg := bus.OutboundMessage{
Channel: origCh,
ChatID: chatID,
Content: outcome.Result.Content,
})
Content: announceContent,
}
for _, mr := range outcome.Result.Media {
outMsg.Media = append(outMsg.Media, bus.MediaAttachment{
URL: mr.Path,
ContentType: mr.ContentType,
})
}
msgBus.PublishOutbound(outMsg)
}(origChannel, msg.ChatID, msg.SenderID, msg.Metadata["subagent_label"])
continue
}
@@ -417,6 +440,7 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
outCh := sched.Schedule(ctx, scheduler.LaneDelegate, agent.RunRequest{
SessionKey: sessionKey,
Message: msg.Content,
ForwardMedia: msg.Media,
Channel: origChannel,
ChatID: msg.ChatID,
PeerKind: origPeerKind,
@@ -438,15 +462,27 @@ func consumeInboundMessages(ctx context.Context, msgBus *bus.MessageBus, agents
})
return
}
if outcome.Result.Content == "" || agent.IsSilentReply(outcome.Result.Content) {
isSilent := outcome.Result.Content == "" || agent.IsSilentReply(outcome.Result.Content)
if isSilent && len(outcome.Result.Media) == 0 {
slog.Info("delegate announce: suppressed silent/empty reply", "delegation", senderID)
return
}
msgBus.PublishOutbound(bus.OutboundMessage{
announceContent := outcome.Result.Content
if isSilent {
announceContent = "" // suppress NO_REPLY text but still send media
}
outMsg := bus.OutboundMessage{
Channel: origCh,
ChatID: chatID,
Content: outcome.Result.Content,
})
Content: announceContent,
}
for _, mr := range outcome.Result.Media {
outMsg.Media = append(outMsg.Media, bus.MediaAttachment{
URL: mr.Path,
ContentType: mr.ContentType,
})
}
msgBus.PublishOutbound(outMsg)
}(origChannel, msg.ChatID, msg.SenderID)
continue
}
+13 -3
View File
@@ -277,10 +277,14 @@ func wireManagedExtras(
if err != nil {
return nil, err
}
return &tools.DelegateRunResult{
dr := &tools.DelegateRunResult{
Content: result.Content,
Iterations: result.Iterations,
}, nil
}
for _, m := range result.Media {
dr.MediaPaths = append(dr.MediaPaths, m.Path)
}
return dr, nil
}
delegateMgr := tools.NewDelegateManager(runAgentFn, stores.AgentLinks, stores.Agents, msgBus)
if stores.Teams != nil {
@@ -309,7 +313,13 @@ func wireManagedExtras(
// Handoff tool (agent-to-agent conversation transfer)
toolsReg.Register(tools.NewHandoffTool(delegateMgr, stores.Teams, stores.Sessions, msgBus))
toolsReg.Register(tools.NewDelegateTool(delegateMgr))
// Inject delegation capability into existing SpawnTool
if st, ok := toolsReg.Get("spawn"); ok {
if spawnTool, ok := st.(*tools.SpawnTool); ok {
spawnTool.SetDelegateManager(delegateMgr)
slog.Info("spawn tool: delegation enabled")
}
}
// Register delegate_search tool (hybrid FTS + semantic agent discovery)
var delegateEmbProvider store.EmbeddingProvider
+30
View File
@@ -235,6 +235,7 @@ type RunRequest struct {
SessionKey string // composite key: agent:{agentId}:{channel}:{peerKind}:{chatId}
Message string // user message
Media []string // local file paths to images (already sanitized)
ForwardMedia []string // media paths to forward to output (not deleted, from delegation results)
Channel string // source channel
ChatID string // source chat ID
PeerKind string // "direct" or "group" (for session key building and tool context)
@@ -670,6 +671,9 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error)
if mr := parseMediaResult(result.ForLLM); mr != nil {
mediaResults = append(mediaResults, *mr)
}
for _, p := range result.Media {
mediaResults = append(mediaResults, MediaResult{Path: p, ContentType: mimeFromExt(filepath.Ext(p))})
}
toolMsg := providers.Message{
Role: "tool",
@@ -778,6 +782,9 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error)
if mr := parseMediaResult(r.result.ForLLM); mr != nil {
mediaResults = append(mediaResults, *mr)
}
for _, p := range r.result.Media {
mediaResults = append(mediaResults, MediaResult{Path: p, ContentType: mimeFromExt(filepath.Ext(p))})
}
toolMsg := providers.Message{
Role: "tool",
@@ -877,6 +884,11 @@ func (l *Loop) runLoop(ctx context.Context, req RunRequest) (*RunResult, error)
// 5. Maybe summarize
l.maybeSummarize(ctx, req.SessionKey)
// Include forwarded media from delegation results (not cleaned up like req.Media)
for _, p := range req.ForwardMedia {
mediaResults = append(mediaResults, MediaResult{Path: p, ContentType: mimeFromExt(filepath.Ext(p))})
}
return &RunResult{
Content: finalContent,
RunID: req.RunID,
@@ -940,6 +952,24 @@ func mimeFromExt(ext string) string {
return "audio/mpeg"
case ".wav":
return "audio/wav"
case ".txt":
return "text/plain"
case ".pdf":
return "application/pdf"
case ".csv":
return "text/csv"
case ".json":
return "application/json"
case ".html", ".htm":
return "text/html"
case ".xml":
return "application/xml"
case ".zip":
return "application/zip"
case ".doc", ".docx":
return "application/msword"
case ".xls", ".xlsx":
return "application/vnd.ms-excel"
default:
return "application/octet-stream"
}
+41 -14
View File
@@ -171,16 +171,18 @@ func NewManagedResolver(deps ResolverDeps) ResolverFunc {
// Inject negative context so the model doesn't waste iterations probing
// unavailable capabilities (team_tasks, delegate_search, etc.).
if !hasTeam || !hasDelegation {
// Note: team agents have delegation targets via team links (TEAM.md),
// so only inject "no delegation" when both hasDelegation and hasTeam are false.
if !hasTeam || (!hasDelegation && !hasTeam) {
var notes []string
if !hasTeam {
notes = append(notes, "You are NOT part of any team. Do not use team_tasks or team_message tools.")
}
if !hasDelegation {
notes = append(notes, "You have NO delegation targets. Do not use delegate or delegate_search tools.")
if !hasDelegation && !hasTeam {
notes = append(notes, "You have NO delegation targets. Do not use spawn with agent parameter or delegate_search tools.")
}
contextFiles = append(contextFiles, bootstrap.ContextFile{
Path: "AVAILABILITY.md",
Path: bootstrap.AvailabilityFile,
Content: strings.Join(notes, "\n"),
})
}
@@ -328,9 +330,9 @@ func filterManualLinks(targets []store.AgentLinkData) []store.AgentLinkData {
func buildDelegateAgentsMD(targets []store.AgentLinkData) string {
var sb strings.Builder
sb.WriteString("# Agent Delegation\n\n")
sb.WriteString("You have the `delegate` tool available. Use it to delegate tasks to other specialized agents.\n")
sb.WriteString("Use `spawn` with the `agent` parameter to delegate tasks to other specialized agents.\n")
sb.WriteString("The agent list below is complete and authoritative — answer questions about available agents directly from it.\n")
sb.WriteString("Only use `delegate` when you need to actually assign work, not to check who is available.\n\n")
sb.WriteString("Only delegate when you need to actually assign work, not to check who is available.\n\n")
sb.WriteString("## Available Agents\n")
for _, t := range targets {
@@ -342,7 +344,7 @@ func buildDelegateAgentsMD(targets []store.AgentLinkData) string {
if t.TargetDescription != "" {
sb.WriteString(t.TargetDescription + "\n")
}
sb.WriteString(fmt.Sprintf("→ `delegate(agent=\"%s\", task=\"describe the task\")`\n", t.TargetAgentKey))
sb.WriteString(fmt.Sprintf("→ `spawn(agent=\"%s\", task=\"describe the task\")`\n", t.TargetAgentKey))
}
sb.WriteString("\n## When to Delegate\n\n")
@@ -358,16 +360,16 @@ func buildDelegateAgentsMD(targets []store.AgentLinkData) string {
func buildDelegateSearchInstruction(targetCount int) string {
return fmt.Sprintf(`# Agent Delegation
You have the `+"`delegate`"+` and `+"`delegate_search`"+` tools available.
You have the `+"`spawn`"+` tool (with `+"`agent`"+` parameter) and `+"`delegate_search`"+` tool available.
Do NOT look for delegation info on disk — it is provided here.
You have access to %d specialized agents. To find the right one:
1. `+"`delegate_search(query=\"your keywords\")`"+` — search agents by expertise
2. `+"`delegate(agent=\"agent-key\", task=\"describe the task\")`"+` — delegate the task
2. `+"`spawn(agent=\"agent-key\", task=\"describe the task\")`"+` — delegate the task
Example:
- User asks about billing → `+"`delegate_search(query=\"billing payment\")`"+` → `+"`delegate(agent=\"billing-agent\", task=\"...\")`"+`
- User asks about billing → `+"`delegate_search(query=\"billing payment\")`"+` → `+"`spawn(agent=\"billing-agent\", task=\"...\")`"+`
Do NOT guess agent keys. Always search first.
`, targetCount)
@@ -398,6 +400,8 @@ func buildTeamMD(team *store.TeamData, members []store.TeamMemberData, selfID uu
for _, m := range members {
if m.AgentID == selfID {
sb.WriteString(fmt.Sprintf("- **you** (%s)", m.Role))
} else if m.DisplayName != "" {
sb.WriteString(fmt.Sprintf("- **%s** `%s` (%s)", m.DisplayName, m.AgentKey, m.Role))
} else {
sb.WriteString(fmt.Sprintf("- **%s** (%s)", m.AgentKey, m.Role))
}
@@ -410,12 +414,33 @@ func buildTeamMD(team *store.TeamData, members []store.TeamMemberData, selfID uu
// Workflow guidance
sb.WriteString("\n## Workflow\n\n")
if selfRole == store.TeamRoleLead {
sb.WriteString("**MANDATORY**: ALWAYS use `team_tasks` to track work. NEVER call `delegate` without a task.\n\n")
sb.WriteString("**MANDATORY**: ALWAYS use `team_tasks` to track work. NEVER delegate without a task.\n\n")
sb.WriteString("**ONE task per ONE delegation.** Each task tracks one unit of work for one agent.\n")
sb.WriteString("When delegating to multiple agents, create a SEPARATE task for each.\n\n")
sb.WriteString("Every delegation MUST follow these 2 steps:\n")
sb.WriteString("1. `team_tasks` action=create, subject=<brief title> → returns task_id\n")
sb.WriteString("2. `delegate` agent=<member>, task=<instructions>, team_task_id=<the task_id from step 1>\n\n")
sb.WriteString("The system ENFORCES this — delegation without team_task_id will be rejected.\n")
sb.WriteString("The task auto-completes when delegation finishes.\n\n")
sb.WriteString("2. `spawn` agent=<member>, task=<instructions>, team_task_id=<the task_id from step 1>\n\n")
sb.WriteString("Example (2 agents):\n")
sb.WriteString("```\n")
sb.WriteString("team_tasks action=create, subject=\"Create illustration\" → task_id=A\n")
sb.WriteString("team_tasks action=create, subject=\"Write caption\" → task_id=B\n")
sb.WriteString("spawn agent=artist, task=\"...\", team_task_id=A\n")
sb.WriteString("spawn agent=writer, task=\"...\", team_task_id=B\n")
sb.WriteString("```\n\n")
sb.WriteString("The system ENFORCES this — spawn with agent but without team_task_id will be rejected.\n")
sb.WriteString("Each task auto-completes when its delegation finishes.\n\n")
sb.WriteString("When multiple delegations run in parallel, the system collects ALL results and delivers\n")
sb.WriteString("them to you in a single combined notification. Do NOT present partial results.\n\n")
sb.WriteString("## Orchestration Patterns\n\n")
sb.WriteString("You can orchestrate multiple rounds — not just one-shot parallel delegation:\n")
sb.WriteString("- **Sequential**: A finishes → review result → delegate to B with A's output as context\n")
sb.WriteString("- **Iterative**: A produces draft → delegate to B for review → delegate back to A with feedback\n")
sb.WriteString("- **Mixed**: A+B in parallel → review both → delegate to C combining their outputs\n\n")
sb.WriteString("After receiving delegation results, decide: present to user (if done) or continue orchestrating.\n\n")
sb.WriteString("**Communication**: When updating the user, distinguish between:\n")
sb.WriteString("- First delegation round → \"assigning to team\" / notifying who is working on what\n")
sb.WriteString("- Follow-up rounds (after receiving results) → \"updating tasks\" / sharing progress and next steps\n")
sb.WriteString("Never repeat the same announcement phrasing for follow-up delegations.\n\n")
sb.WriteString("`team_tasks` actions:\n")
sb.WriteString("- action=list → active tasks (pending/in_progress/blocked), no results shown\n")
sb.WriteString("- action=list, status=all → all tasks including completed\n")
@@ -427,6 +452,8 @@ func buildTeamMD(team *store.TeamData, members []store.TeamMemberData, selfID uu
} else {
sb.WriteString("As a member, when you receive a delegated task, just do the work.\n")
sb.WriteString("Task completion is handled automatically by the system.\n\n")
sb.WriteString("For long-running tasks, send progress updates to your lead:\n")
sb.WriteString("`team_message` action=send, to=<lead_key>, text=<progress update>\n\n")
sb.WriteString("`team_tasks` actions:\n")
sb.WriteString("- action=list → check team task board (active tasks)\n")
sb.WriteString("- action=get, task_id=<id> → read a completed task's full result\n")
+1 -2
View File
@@ -50,8 +50,7 @@ var coreToolSummaries = map[string]string{
"exec": "Run shell commands",
"memory_search": "Search indexed memory files (MEMORY.md + memory/*.md)",
"memory_get": "Read specific sections of memory files",
"spawn": "Spawn a subagent for parallel/background tasks",
"subagent": "List, steer, or kill subagents",
"spawn": "Spawn a subagent or delegate to another agent",
"web_search": "Search the web",
"web_fetch": "Fetch and extract content from a URL",
"cron": "Manage scheduled jobs and reminders",
+2 -2
View File
@@ -106,13 +106,13 @@ func buildProjectContextSection(files []bootstrap.ContextFile) []string {
// During bootstrap (first run), skip delegation/team/availability files — they add noise
// and waste tokens when the agent should only be introducing itself.
if hasBootstrap && (base == bootstrap.DelegationFile || base == bootstrap.TeamFile || base == "AVAILABILITY.md") {
if hasBootstrap && (base == bootstrap.DelegationFile || base == bootstrap.TeamFile || base == bootstrap.AvailabilityFile) {
continue
}
// Virtual files (DELEGATION.md, TEAM.md, AVAILABILITY.md) are system-injected, not on disk.
// Render with <system_context> so the LLM doesn't try to read/write them as files.
if base == bootstrap.DelegationFile || base == bootstrap.TeamFile || base == "AVAILABILITY.md" {
if base == bootstrap.DelegationFile || base == bootstrap.TeamFile || base == bootstrap.AvailabilityFile {
lines = append(lines,
fmt.Sprintf("<system_context name=%q>", base),
f.Content,
+3 -2
View File
@@ -28,8 +28,9 @@ const (
UserFile = "USER.md"
HeartbeatFile = "HEARTBEAT.md"
BootstrapFile = "BOOTSTRAP.md"
DelegationFile = "DELEGATION.md"
TeamFile = "TEAM.md"
DelegationFile = "DELEGATION.md"
TeamFile = "TEAM.md"
AvailabilityFile = "AVAILABILITY.md"
MemoryFile = "MEMORY.md"
MemoryAltFile = "memory.md"
MemoryJSONFile = "MEMORY.json"
+31
View File
@@ -11,11 +11,42 @@ import (
// --- Markdown to Telegram HTML conversion ---
// Adapted from PicoClaw's telegram.go, extended with table support (matching TS "code" mode).
// htmlTagToMarkdown converts common HTML tags in LLM output to markdown equivalents
// so they survive the escapeHTML step and get re-converted by the markdown pipeline.
var htmlToMdReplacers = []struct {
re *regexp.Regexp
repl string
}{
{regexp.MustCompile(`(?i)<br\s*/?>`), "\n"},
{regexp.MustCompile(`(?i)</?p\s*>`), "\n"},
{regexp.MustCompile(`(?i)<b>([\s\S]*?)</b>`), "**$1**"},
{regexp.MustCompile(`(?i)<strong>([\s\S]*?)</strong>`), "**$1**"},
{regexp.MustCompile(`(?i)<i>([\s\S]*?)</i>`), "_$1_"},
{regexp.MustCompile(`(?i)<em>([\s\S]*?)</em>`), "_$1_"},
{regexp.MustCompile(`(?i)<s>([\s\S]*?)</s>`), "~~$1~~"},
{regexp.MustCompile(`(?i)<strike>([\s\S]*?)</strike>`), "~~$1~~"},
{regexp.MustCompile(`(?i)<del>([\s\S]*?)</del>`), "~~$1~~"},
{regexp.MustCompile(`(?i)<code>([\s\S]*?)</code>`), "`$1`"},
{regexp.MustCompile(`(?i)<a\s+href="([^"]+)"[^>]*>([\s\S]*?)</a>`), "[$2]($1)"},
}
func htmlTagToMarkdown(text string) string {
for _, r := range htmlToMdReplacers {
text = r.re.ReplaceAllString(text, r.repl)
}
return text
}
func markdownToTelegramHTML(text string) string {
if text == "" {
return ""
}
// Pre-process: convert any HTML tags in LLM output to markdown equivalents.
// LLMs sometimes output raw HTML (e.g. <b>bold</b>) which would get escaped
// by escapeHTML() and displayed as literal "<b>bold</b>" text.
text = htmlTagToMarkdown(text)
// Extract markdown tables FIRST — uses dedicated \x00TB placeholders.
// Tables render as <pre> (monospace block) WITHOUT <code> wrapper,
// so Telegram shows them as preformatted text, not as "code" with copy button.
+133 -23
View File
@@ -3,10 +3,12 @@ package telegram
import (
"context"
"fmt"
"html"
"log/slog"
"os"
"regexp"
"strings"
"time"
"github.com/mymmrac/telego"
tu "github.com/mymmrac/telego/telegoutil"
@@ -19,8 +21,64 @@ import (
var (
parseErrRe = regexp.MustCompile(`(?i)can't parse entities|parse entities|find end of the entity`)
messageNotModifiedRe = regexp.MustCompile(`(?i)message is not modified`)
htmlTagRe = regexp.MustCompile(`<[^>]*>`)
)
const (
sendMaxRetries = 3
sendRetryDelay = 2 * time.Second
photoSizeThreshold = 5 * 1024 * 1024 // 5 MB — images larger than this are sent as documents to avoid Telegram compression
)
// stripHTML removes HTML tags and unescapes HTML entities for plain-text fallback.
func stripHTML(s string) string {
return html.UnescapeString(htmlTagRe.ReplaceAllString(s, ""))
}
// isRetryableNetworkErr checks if a Telegram API error is a transient network error worth retrying.
func isRetryableNetworkErr(err error) bool {
if err == nil {
return false
}
s := err.Error()
return strings.Contains(s, "timeout") ||
strings.Contains(s, "connection reset") ||
strings.Contains(s, "broken pipe") ||
strings.Contains(s, "EOF") ||
strings.Contains(s, "lookup") // DNS resolution failure
}
// retrySend wraps a Telegram send call with retry logic for transient network errors.
// Parse errors are NOT retried (handled by caller's HTML fallback).
// resetFn is called before each retry (e.g. to seek file handles back to start). Can be nil.
func retrySend(ctx context.Context, name string, resetFn func(), fn func() error) error {
var err error
for attempt := 1; attempt <= sendMaxRetries; attempt++ {
err = fn()
if err == nil {
return nil
}
// Don't retry parse errors — caller handles HTML fallback
if parseErrRe.MatchString(err.Error()) {
return err
}
if !isRetryableNetworkErr(err) || attempt == sendMaxRetries {
return err
}
slog.Warn("telegram send retry",
"func", name, "attempt", attempt, "max", sendMaxRetries, "error", err)
if resetFn != nil {
resetFn()
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(sendRetryDelay * time.Duration(attempt)):
}
}
return err
}
// Send delivers an outbound message to a Telegram chat.
// Supports text-only messages and messages with media attachments.
// Reads metadata for reply-to-message and forum thread routing.
@@ -137,23 +195,34 @@ func (c *Channel) sendMediaMessage(ctx context.Context, chatID int64, msg bus.Ou
msg.Content = "" // only use for first media
}
// Convert caption from markdown to Telegram HTML (same as regular messages)
// Convert caption from markdown to Telegram HTML (same as regular messages).
// If the HTML caption exceeds Telegram's 1024-byte limit, skip caption entirely
// and send the full text as a separate message. Truncating HTML at a byte boundary
// can split tags (e.g. cut inside <code>...</code>) causing parse errors.
var followUpText string
if caption != "" {
caption = markdownToTelegramHTML(caption)
if len(caption) > telegramCaptionMaxLen {
followUpText = caption
caption = ""
}
}
// Split caption if too long (Telegram limit: 1024 chars)
var followUpText string
if len(caption) > telegramCaptionMaxLen {
followUpText = caption[telegramCaptionMaxLen:]
caption = caption[:telegramCaptionMaxLen]
}
// Send based on content type
// Send based on content type.
// Large images (>photoSizeThreshold) are sent as documents to avoid Telegram compression.
ct := strings.ToLower(media.ContentType)
switch {
case strings.HasPrefix(ct, "image/"):
if err := c.sendPhoto(ctx, chatIDObj, media.URL, caption, replyTo, threadID); err != nil {
sendAsDoc := false
if info, statErr := os.Stat(media.URL); statErr == nil && info.Size() > photoSizeThreshold {
sendAsDoc = true
slog.Info("large image, sending as document to preserve quality", "path", media.URL, "size", info.Size())
}
if sendAsDoc {
if err := c.sendDocument(ctx, chatIDObj, media.URL, caption, replyTo, threadID); err != nil {
return err
}
} else if err := c.sendPhoto(ctx, chatIDObj, media.URL, caption, replyTo, threadID); err != nil {
return err
}
case strings.HasPrefix(ct, "video/"):
@@ -200,16 +269,17 @@ func (c *Channel) sendHTML(ctx context.Context, chatID int64, html string, reply
tgMsg.ReplyParameters = &telego.ReplyParameters{MessageID: replyTo}
}
if _, err := c.bot.SendMessage(ctx, tgMsg); err != nil {
if parseErrRe.MatchString(err.Error()) {
slog.Warn("HTML parse failed, falling back to plain text", "error", err)
tgMsg.ParseMode = ""
_, err = c.bot.SendMessage(ctx, tgMsg)
return err
}
return err
err := retrySend(ctx, "sendMessage", nil, func() error {
_, e := c.bot.SendMessage(ctx, tgMsg)
return e
})
if err != nil && parseErrRe.MatchString(err.Error()) {
slog.Warn("HTML parse failed, falling back to plain text", "error", err)
tgMsg.ParseMode = ""
tgMsg.Text = stripHTML(tgMsg.Text)
_, err = c.bot.SendMessage(ctx, tgMsg)
}
return nil
return err
}
// sendPhoto sends a photo message.
@@ -235,7 +305,17 @@ func (c *Channel) sendPhoto(ctx context.Context, chatID telego.ChatID, filePath,
params.ReplyParameters = &telego.ReplyParameters{MessageID: replyTo}
}
_, err = c.bot.SendPhoto(ctx, params)
err = retrySend(ctx, "sendPhoto", func() { file.Seek(0, 0) }, func() error {
_, e := c.bot.SendPhoto(ctx, params)
return e
})
if err != nil && parseErrRe.MatchString(err.Error()) {
slog.Warn("sendPhoto: HTML parse failed, retrying with plain text caption", "error", err)
file.Seek(0, 0)
params.ParseMode = ""
params.Caption = stripHTML(params.Caption)
_, err = c.bot.SendPhoto(ctx, params)
}
return err
}
@@ -262,7 +342,17 @@ func (c *Channel) sendVideo(ctx context.Context, chatID telego.ChatID, filePath,
params.ReplyParameters = &telego.ReplyParameters{MessageID: replyTo}
}
_, err = c.bot.SendVideo(ctx, params)
err = retrySend(ctx, "sendVideo", func() { file.Seek(0, 0) }, func() error {
_, e := c.bot.SendVideo(ctx, params)
return e
})
if err != nil && parseErrRe.MatchString(err.Error()) {
slog.Warn("sendVideo: HTML parse failed, retrying with plain text caption", "error", err)
file.Seek(0, 0)
params.ParseMode = ""
params.Caption = stripHTML(params.Caption)
_, err = c.bot.SendVideo(ctx, params)
}
return err
}
@@ -289,7 +379,17 @@ func (c *Channel) sendAudio(ctx context.Context, chatID telego.ChatID, filePath,
params.ReplyParameters = &telego.ReplyParameters{MessageID: replyTo}
}
_, err = c.bot.SendAudio(ctx, params)
err = retrySend(ctx, "sendAudio", func() { file.Seek(0, 0) }, func() error {
_, e := c.bot.SendAudio(ctx, params)
return e
})
if err != nil && parseErrRe.MatchString(err.Error()) {
slog.Warn("sendAudio: HTML parse failed, retrying with plain text caption", "error", err)
file.Seek(0, 0)
params.ParseMode = ""
params.Caption = stripHTML(params.Caption)
_, err = c.bot.SendAudio(ctx, params)
}
return err
}
@@ -316,7 +416,17 @@ func (c *Channel) sendDocument(ctx context.Context, chatID telego.ChatID, filePa
params.ReplyParameters = &telego.ReplyParameters{MessageID: replyTo}
}
_, err = c.bot.SendDocument(ctx, params)
err = retrySend(ctx, "sendDocument", func() { file.Seek(0, 0) }, func() error {
_, e := c.bot.SendDocument(ctx, params)
return e
})
if err != nil && parseErrRe.MatchString(err.Error()) {
slog.Warn("sendDocument: HTML parse failed, retrying with plain text caption", "error", err)
file.Seek(0, 0)
params.ParseMode = ""
params.Caption = stripHTML(params.Caption)
_, err = c.bot.SendDocument(ctx, params)
}
return err
}
+28 -34
View File
@@ -137,11 +137,10 @@ func (m *TeamsMethods) handleCreate(_ context.Context, client *gateway.Client, r
}
}
// Auto-create bidirectional agent_links between all team members.
// This enables delegation between teammates.
// Auto-create outbound agent_links from lead to each member.
// Only the lead can delegate to members.
if m.linkStore != nil {
allAgents := append([]*store.AgentData{leadAgent}, memberAgents...)
m.autoCreateTeamLinks(ctx, team.ID, allAgents, client.UserID())
m.autoCreateTeamLinks(ctx, team.ID, leadAgent, memberAgents, client.UserID())
}
// Invalidate agent caches so TEAM.md gets injected
@@ -353,18 +352,11 @@ func (m *TeamsMethods) handleAddMember(_ context.Context, client *gateway.Client
return
}
// Auto-create bidirectional links between new member and existing members
// Auto-create outbound link from lead to new member
if m.linkStore != nil {
existingMembers, _ := m.teamStore.ListMembers(ctx, teamID)
for _, member := range existingMembers {
if member.AgentID == ag.ID {
continue
}
memberAgent, err := m.agentStore.GetByID(ctx, member.AgentID)
if err != nil {
continue
}
m.autoCreateTeamLinks(ctx, teamID, []*store.AgentData{ag, memberAgent}, client.UserID())
leadAgent, err := m.agentStore.GetByID(ctx, team.LeadAgentID)
if err == nil {
m.autoCreateTeamLinks(ctx, teamID, leadAgent, []*store.AgentData{ag}, client.UserID())
}
}
@@ -464,25 +456,27 @@ func (m *TeamsMethods) invalidateTeamCaches(ctx context.Context, teamID uuid.UUI
// --- helpers ---
// autoCreateTeamLinks creates bidirectional agent_links between all team members.
// Silently skips existing links (UNIQUE constraint).
func (m *TeamsMethods) autoCreateTeamLinks(ctx context.Context, teamID uuid.UUID, agents []*store.AgentData, createdBy string) {
for i := 0; i < len(agents); i++ {
for j := i + 1; j < len(agents); j++ {
link := &store.AgentLinkData{
SourceAgentID: agents[i].ID,
TargetAgentID: agents[j].ID,
Direction: store.LinkDirectionBidirectional,
TeamID: &teamID,
Description: "auto-created by team",
MaxConcurrent: 3,
Status: store.LinkStatusActive,
CreatedBy: createdBy,
}
if err := m.linkStore.CreateLink(ctx, link); err != nil {
slog.Debug("teams: auto-link already exists or failed",
"source", agents[i].AgentKey, "target", agents[j].AgentKey, "error", err)
}
// autoCreateTeamLinks creates outbound agent_links from lead to each member.
// Only the lead can delegate to members — members cannot delegate back to lead
// or to other members. Silently skips existing links (UNIQUE constraint).
func (m *TeamsMethods) autoCreateTeamLinks(ctx context.Context, teamID uuid.UUID, leadAgent *store.AgentData, members []*store.AgentData, createdBy string) {
for _, member := range members {
if member.ID == leadAgent.ID {
continue
}
link := &store.AgentLinkData{
SourceAgentID: leadAgent.ID,
TargetAgentID: member.ID,
Direction: store.LinkDirectionOutbound,
TeamID: &teamID,
Description: "auto-created by team",
MaxConcurrent: 3,
Status: store.LinkStatusActive,
CreatedBy: createdBy,
}
if err := m.linkStore.CreateLink(ctx, link); err != nil {
slog.Debug("teams: auto-link already exists or failed",
"source", leadAgent.AgentKey, "target", member.AgentKey, "error", err)
}
}
}
+36 -6
View File
@@ -39,6 +39,9 @@ var summoningFiles = []string{
// fileTagRe parses <file name="SOUL.md">content</file> from LLM output.
var fileTagRe = regexp.MustCompile(`(?s)<file\s+name="([^"]+)">\s*(.*?)\s*</file>`)
// identityNameRe extracts the Name field from IDENTITY.md format: - **Name:** value
var identityNameRe = regexp.MustCompile(`(?m)^-\s*\*\*Name:\*\*\s*(.+)$`)
// frontmatterTagRe parses <frontmatter>short expertise summary</frontmatter> from LLM output.
var frontmatterTagRe = regexp.MustCompile(`(?s)<frontmatter>\s*(.*?)\s*</frontmatter>`)
@@ -81,14 +84,21 @@ func (s *AgentSummoner) SummonAgent(agentID uuid.UUID, providerName, model, desc
s.storeFiles(ctx, agentID, files)
// Save frontmatter — use LLM-generated if available, otherwise fallback to truncated description
// Save frontmatter + display_name extracted from IDENTITY.md
updates := map[string]any{}
fm := files[frontmatterKey]
if fm == "" {
fm = truncateUTF8(description, 200)
}
if fm != "" {
if err := s.agents.Update(ctx, agentID, map[string]any{"frontmatter": fm}); err != nil {
slog.Warn("summoning: failed to save frontmatter", "agent", agentID, "error", err)
updates["frontmatter"] = fm
}
if name := extractIdentityName(files[bootstrap.IdentityFile]); name != "" {
updates["display_name"] = name
}
if len(updates) > 0 {
if err := s.agents.Update(ctx, agentID, updates); err != nil {
slog.Warn("summoning: failed to save agent metadata", "agent", agentID, "error", err)
}
}
@@ -129,10 +139,17 @@ func (s *AgentSummoner) RegenerateAgent(agentID uuid.UUID, providerName, model,
s.storeFiles(ctx, agentID, files)
// Update frontmatter if LLM generated one
// Update frontmatter + display_name if IDENTITY.md was regenerated
updates := map[string]any{}
if fm, ok := files[frontmatterKey]; ok && fm != "" {
if err := s.agents.Update(ctx, agentID, map[string]any{"frontmatter": fm}); err != nil {
slog.Warn("summoning: failed to save frontmatter", "agent", agentID, "error", err)
updates["frontmatter"] = fm
}
if name := extractIdentityName(files[bootstrap.IdentityFile]); name != "" {
updates["display_name"] = name
}
if len(updates) > 0 {
if err := s.agents.Update(ctx, agentID, updates); err != nil {
slog.Warn("summoning: failed to save agent metadata", "agent", agentID, "error", err)
}
}
@@ -333,6 +350,19 @@ Output format:
return sb.String()
}
// extractIdentityName extracts the Name field from IDENTITY.md content.
// Matches format: - **Name:** value
func extractIdentityName(content string) string {
if content == "" {
return ""
}
m := identityNameRe.FindStringSubmatch(content)
if len(m) < 2 {
return ""
}
return strings.TrimSpace(m[1])
}
// truncateUTF8 truncates s to at most maxLen runes, appending "…" if truncated.
func truncateUTF8(s string, maxLen int) string {
runes := []rune(s)
+32 -46
View File
@@ -126,12 +126,14 @@ func (s *PGAgentLinkStore) CanDelegate(ctx context.Context, fromAgentID, toAgent
}
func (s *PGAgentLinkStore) DelegateTargets(ctx context.Context, fromAgentID uuid.UUID) ([]store.AgentLinkData, error) {
// CASE expressions ensure "target" columns always refer to the "other" agent,
// regardless of whether fromAgent is source or target side of the link.
rows, err := s.db.QueryContext(ctx,
`SELECT `+linkSelectColsJoined+`,
sa.agent_key AS source_agent_key,
ta.agent_key AS target_agent_key,
COALESCE(ta.display_name, '') AS target_display_name,
COALESCE(ta.frontmatter, '') AS target_description,
CASE WHEN l.source_agent_id = $1 THEN sa.agent_key ELSE ta.agent_key END AS source_agent_key,
CASE WHEN l.source_agent_id = $1 THEN ta.agent_key ELSE sa.agent_key END AS target_agent_key,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.display_name, '') ELSE COALESCE(sa.display_name, '') END AS target_display_name,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.frontmatter, '') ELSE COALESCE(sa.frontmatter, '') END AS target_description,
COALESCE(tm.name, '') AS team_name
FROM agent_links l
JOIN agents sa ON sa.id = l.source_agent_id
@@ -142,38 +144,12 @@ func (s *PGAgentLinkStore) DelegateTargets(ctx context.Context, fromAgentID uuid
OR
(l.target_agent_id = $1 AND l.direction IN ('inbound', 'bidirectional'))
)
ORDER BY ta.agent_key`, fromAgentID)
ORDER BY CASE WHEN l.source_agent_id = $1 THEN ta.agent_key ELSE sa.agent_key END`, fromAgentID)
if err != nil {
return nil, err
}
defer rows.Close()
var links []store.AgentLinkData
for rows.Next() {
var d store.AgentLinkData
var desc sql.NullString
if err := rows.Scan(
&d.ID, &d.SourceAgentID, &d.TargetAgentID, &d.Direction, &d.TeamID, &desc,
&d.MaxConcurrent, &d.Settings, &d.Status, &d.CreatedBy, &d.CreatedAt, &d.UpdatedAt,
&d.SourceAgentKey, &d.TargetAgentKey, &d.TargetDisplayName, &d.TargetDescription,
&d.TeamName,
); err != nil {
return nil, err
}
if desc.Valid {
d.Description = desc.String
}
// For links where this agent is the target (inbound direction),
// swap to show the actual target (the other agent) as the delegate target.
if d.TargetAgentID == fromAgentID {
d.TargetAgentID = d.SourceAgentID
d.TargetAgentKey = d.SourceAgentKey
}
links = append(links, d)
}
return links, rows.Err()
return scanLinkRowsJoined(rows)
}
func (s *PGAgentLinkStore) GetLinkBetween(ctx context.Context, fromAgentID, toAgentID uuid.UUID) (*store.AgentLinkData, error) {
@@ -195,21 +171,27 @@ func (s *PGAgentLinkStore) SearchDelegateTargets(ctx context.Context, fromAgentI
if limit <= 0 {
limit = 5
}
// Handle both directions: when fromAgent is source OR target of a bidirectional link.
// CASE expressions ensure "target" columns always refer to the "other" agent.
rows, err := s.db.QueryContext(ctx,
`SELECT `+linkSelectColsJoined+`,
sa.agent_key AS source_agent_key,
ta.agent_key AS target_agent_key,
COALESCE(ta.display_name, '') AS target_display_name,
COALESCE(ta.frontmatter, '') AS target_description,
CASE WHEN l.source_agent_id = $1 THEN sa.agent_key ELSE ta.agent_key END AS source_agent_key,
CASE WHEN l.source_agent_id = $1 THEN ta.agent_key ELSE sa.agent_key END AS target_agent_key,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.display_name, '') ELSE COALESCE(sa.display_name, '') END AS target_display_name,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.frontmatter, '') ELSE COALESCE(sa.frontmatter, '') END AS target_description,
COALESCE(tm.name, '') AS team_name
FROM agent_links l
JOIN agents sa ON sa.id = l.source_agent_id
JOIN agents ta ON ta.id = l.target_agent_id
LEFT JOIN agent_teams tm ON tm.id = l.team_id
WHERE l.status = 'active'
AND (l.source_agent_id = $1 AND l.direction IN ('outbound', 'bidirectional'))
AND ta.tsv @@ plainto_tsquery('simple', $2)
ORDER BY ts_rank(ta.tsv, plainto_tsquery('simple', $2)) DESC
AND (
(l.source_agent_id = $1 AND l.direction IN ('outbound', 'bidirectional'))
OR
(l.target_agent_id = $1 AND l.direction IN ('inbound', 'bidirectional'))
)
AND CASE WHEN l.source_agent_id = $1 THEN ta.tsv ELSE sa.tsv END @@ plainto_tsquery('simple', $2)
ORDER BY ts_rank(CASE WHEN l.source_agent_id = $1 THEN ta.tsv ELSE sa.tsv END, plainto_tsquery('simple', $2)) DESC
LIMIT $3`, fromAgentID, query, limit)
if err != nil {
return nil, err
@@ -225,19 +207,23 @@ func (s *PGAgentLinkStore) SearchDelegateTargetsByEmbedding(ctx context.Context,
vecStr := vectorToString(embedding)
rows, err := s.db.QueryContext(ctx,
`SELECT `+linkSelectColsJoined+`,
sa.agent_key AS source_agent_key,
ta.agent_key AS target_agent_key,
COALESCE(ta.display_name, '') AS target_display_name,
COALESCE(ta.frontmatter, '') AS target_description,
CASE WHEN l.source_agent_id = $1 THEN sa.agent_key ELSE ta.agent_key END AS source_agent_key,
CASE WHEN l.source_agent_id = $1 THEN ta.agent_key ELSE sa.agent_key END AS target_agent_key,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.display_name, '') ELSE COALESCE(sa.display_name, '') END AS target_display_name,
CASE WHEN l.source_agent_id = $1 THEN COALESCE(ta.frontmatter, '') ELSE COALESCE(sa.frontmatter, '') END AS target_description,
COALESCE(tm.name, '') AS team_name
FROM agent_links l
JOIN agents sa ON sa.id = l.source_agent_id
JOIN agents ta ON ta.id = l.target_agent_id
LEFT JOIN agent_teams tm ON tm.id = l.team_id
WHERE l.status = 'active'
AND (l.source_agent_id = $1 AND l.direction IN ('outbound', 'bidirectional'))
AND ta.embedding IS NOT NULL
ORDER BY ta.embedding <=> $2::vector
AND (
(l.source_agent_id = $1 AND l.direction IN ('outbound', 'bidirectional'))
OR
(l.target_agent_id = $1 AND l.direction IN ('inbound', 'bidirectional'))
)
AND CASE WHEN l.source_agent_id = $1 THEN ta.embedding ELSE sa.embedding END IS NOT NULL
ORDER BY (CASE WHEN l.source_agent_id = $1 THEN ta.embedding ELSE sa.embedding END) <=> $2::vector
LIMIT $3`, fromAgentID, vecStr, limit)
if err != nil {
return nil, err
+10
View File
@@ -141,6 +141,16 @@ func (s *PGBuiltinToolStore) Seed(ctx context.Context, tools []store.BuiltinTool
}
}
// Reconcile: remove stale entries not in the current seed list
names := make([]string, len(tools))
for i, t := range tools {
names[i] = t.Name
}
if _, err := tx.ExecContext(ctx,
`DELETE FROM builtin_tools WHERE name != ALL($1)`, pqStringArray(names)); err != nil {
return fmt.Errorf("reconcile stale builtin tools: %w", err)
}
return tx.Commit()
}
+1 -1
View File
@@ -26,7 +26,7 @@ const teamSelectCols = `id, name, lead_agent_id, description, status, settings,
const taskSelectCols = `id, team_id, subject, description, status, owner_agent_id, blocked_by, priority, result, created_at, updated_at`
const messageSelectCols = `id, team_id, from_agent_id, to_agent_id, content, message_type, read, created_at`
const messageSelectCols = `id, team_id, from_agent_id, to_agent_id, content, message_type, read, task_id, metadata, created_at`
// ============================================================
// Team CRUD
+21 -7
View File
@@ -3,6 +3,7 @@ package pg
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
@@ -22,14 +23,19 @@ func (s *PGTeamStore) SaveDelegationHistory(ctx context.Context, record *store.D
now := time.Now()
record.CreatedAt = now
metadata, _ := json.Marshal(record.Metadata)
if len(metadata) == 0 {
metadata = []byte(`{}`)
}
_, err := s.db.ExecContext(ctx,
`INSERT INTO delegation_history (id, source_agent_id, target_agent_id, team_id, team_task_id, user_id, task, mode, status, result, error, iterations, trace_id, duration_ms, created_at, completed_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)`,
`INSERT INTO delegation_history (id, source_agent_id, target_agent_id, team_id, team_task_id, user_id, task, mode, status, result, error, iterations, trace_id, duration_ms, metadata, created_at, completed_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17)`,
record.ID, record.SourceAgentID, record.TargetAgentID,
record.TeamID, record.TeamTaskID,
record.UserID, record.Task, record.Mode, record.Status,
record.Result, record.Error, record.Iterations,
record.TraceID, record.DurationMS, now, record.CompletedAt,
record.TraceID, record.DurationMS, metadata, now, record.CompletedAt,
)
return err
}
@@ -81,7 +87,7 @@ func (s *PGTeamStore) ListDelegationHistory(ctx context.Context, opts store.Dele
query := fmt.Sprintf(
`SELECT d.id, d.source_agent_id, d.target_agent_id, d.team_id, d.team_task_id,
d.user_id, d.task, d.mode, d.status, d.result, d.error, d.iterations,
d.trace_id, d.duration_ms, d.created_at, d.completed_at,
d.trace_id, d.duration_ms, d.metadata, d.created_at, d.completed_at,
COALESCE(sa.agent_key, '') AS source_agent_key,
COALESCE(ta.agent_key, '') AS target_agent_key
FROM delegation_history d
@@ -103,10 +109,11 @@ func (s *PGTeamStore) ListDelegationHistory(ctx context.Context, opts store.Dele
var d store.DelegationHistoryData
var result, errStr sql.NullString
var completedAt sql.NullTime
var metadata json.RawMessage
if err := rows.Scan(
&d.ID, &d.SourceAgentID, &d.TargetAgentID, &d.TeamID, &d.TeamTaskID,
&d.UserID, &d.Task, &d.Mode, &d.Status, &result, &errStr, &d.Iterations,
&d.TraceID, &d.DurationMS, &d.CreatedAt, &completedAt,
&d.TraceID, &d.DurationMS, &metadata, &d.CreatedAt, &completedAt,
&d.SourceAgentKey, &d.TargetAgentKey,
); err != nil {
return nil, 0, err
@@ -120,6 +127,9 @@ func (s *PGTeamStore) ListDelegationHistory(ctx context.Context, opts store.Dele
if completedAt.Valid {
d.CompletedAt = &completedAt.Time
}
if len(metadata) > 0 && string(metadata) != "{}" {
_ = json.Unmarshal(metadata, &d.Metadata)
}
records = append(records, d)
}
return records, total, rows.Err()
@@ -129,11 +139,12 @@ func (s *PGTeamStore) GetDelegationHistory(ctx context.Context, id uuid.UUID) (*
var d store.DelegationHistoryData
var result, errStr sql.NullString
var completedAt sql.NullTime
var metadata json.RawMessage
err := s.db.QueryRowContext(ctx,
`SELECT d.id, d.source_agent_id, d.target_agent_id, d.team_id, d.team_task_id,
d.user_id, d.task, d.mode, d.status, d.result, d.error, d.iterations,
d.trace_id, d.duration_ms, d.created_at, d.completed_at,
d.trace_id, d.duration_ms, d.metadata, d.created_at, d.completed_at,
COALESCE(sa.agent_key, '') AS source_agent_key,
COALESCE(ta.agent_key, '') AS target_agent_key
FROM delegation_history d
@@ -142,7 +153,7 @@ func (s *PGTeamStore) GetDelegationHistory(ctx context.Context, id uuid.UUID) (*
WHERE d.id = $1`, id).Scan(
&d.ID, &d.SourceAgentID, &d.TargetAgentID, &d.TeamID, &d.TeamTaskID,
&d.UserID, &d.Task, &d.Mode, &d.Status, &result, &errStr, &d.Iterations,
&d.TraceID, &d.DurationMS, &d.CreatedAt, &completedAt,
&d.TraceID, &d.DurationMS, &metadata, &d.CreatedAt, &completedAt,
&d.SourceAgentKey, &d.TargetAgentKey,
)
if err != nil {
@@ -157,5 +168,8 @@ func (s *PGTeamStore) GetDelegationHistory(ctx context.Context, id uuid.UUID) (*
if completedAt.Valid {
d.CompletedAt = &completedAt.Time
}
if len(metadata) > 0 && string(metadata) != "{}" {
_ = json.Unmarshal(metadata, &d.Metadata)
}
return &d, nil
}
+55 -6
View File
@@ -3,6 +3,7 @@ package pg
import (
"context"
"database/sql"
"encoding/json"
"time"
"github.com/google/uuid"
@@ -20,18 +21,23 @@ func (s *PGTeamStore) SendMessage(ctx context.Context, msg *store.TeamMessageDat
}
msg.CreatedAt = time.Now()
metadata, _ := json.Marshal(msg.Metadata)
if len(metadata) == 0 {
metadata = []byte(`{}`)
}
_, err := s.db.ExecContext(ctx,
`INSERT INTO team_messages (id, team_id, from_agent_id, to_agent_id, content, message_type, read, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`,
`INSERT INTO team_messages (id, team_id, from_agent_id, to_agent_id, content, message_type, read, task_id, metadata, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)`,
msg.ID, msg.TeamID, msg.FromAgentID, msg.ToAgentID,
msg.Content, msg.MessageType, false, msg.CreatedAt,
msg.Content, msg.MessageType, false, msg.TaskID, metadata, msg.CreatedAt,
)
return err
}
func (s *PGTeamStore) GetUnread(ctx context.Context, teamID, agentID uuid.UUID) ([]store.TeamMessageData, error) {
rows, err := s.db.QueryContext(ctx,
`SELECT m.id, m.team_id, m.from_agent_id, m.to_agent_id, m.content, m.message_type, m.read, m.created_at,
`SELECT m.id, m.team_id, m.from_agent_id, m.to_agent_id, m.content, m.message_type, m.read, m.task_id, m.metadata, m.created_at,
COALESCE(fa.agent_key, '') AS from_agent_key,
COALESCE(ta.agent_key, '') AS to_agent_key
FROM team_messages m
@@ -52,20 +58,63 @@ func (s *PGTeamStore) MarkRead(ctx context.Context, messageID uuid.UUID) error {
return err
}
func (s *PGTeamStore) ListMessages(ctx context.Context, teamID uuid.UUID, limit, offset int) ([]store.TeamMessageData, int, error) {
if limit <= 0 || limit > 200 {
limit = 50
}
if offset < 0 {
offset = 0
}
// Count total
var total int
if err := s.db.QueryRowContext(ctx,
`SELECT COUNT(*) FROM team_messages WHERE team_id = $1`, teamID).Scan(&total); err != nil {
return nil, 0, err
}
rows, err := s.db.QueryContext(ctx,
`SELECT m.id, m.team_id, m.from_agent_id, m.to_agent_id, m.content, m.message_type, m.read, m.task_id, m.metadata, m.created_at,
COALESCE(fa.agent_key, '') AS from_agent_key,
COALESCE(ta.agent_key, '') AS to_agent_key
FROM team_messages m
LEFT JOIN agents fa ON fa.id = m.from_agent_id
LEFT JOIN agents ta ON ta.id = m.to_agent_id
WHERE m.team_id = $1
ORDER BY m.created_at DESC
LIMIT $2 OFFSET $3`, teamID, limit, offset)
if err != nil {
return nil, 0, err
}
defer rows.Close()
messages, err := scanMessageRowsJoined(rows)
if err != nil {
return nil, 0, err
}
return messages, total, nil
}
func scanMessageRowsJoined(rows *sql.Rows) ([]store.TeamMessageData, error) {
var messages []store.TeamMessageData
for rows.Next() {
var d store.TeamMessageData
var toAgentID *uuid.UUID
var toAgentID, taskID *uuid.UUID
var metadata json.RawMessage
if err := rows.Scan(
&d.ID, &d.TeamID, &d.FromAgentID, &toAgentID,
&d.Content, &d.MessageType, &d.Read, &d.CreatedAt,
&d.Content, &d.MessageType, &d.Read, &taskID, &metadata, &d.CreatedAt,
&d.FromAgentKey, &d.ToAgentKey,
); err != nil {
return nil, err
}
d.ToAgentID = toAgentID
d.TaskID = taskID
if len(metadata) > 0 && string(metadata) != "{}" {
_ = json.Unmarshal(metadata, &d.Metadata)
}
messages = append(messages, d)
}
return messages, rows.Err()
}
+45 -38
View File
@@ -71,14 +71,15 @@ type TeamMemberData struct {
// TeamTaskData represents a task in the team's shared task list.
type TeamTaskData struct {
BaseModel
TeamID uuid.UUID `json:"team_id"`
Subject string `json:"subject"`
Description string `json:"description,omitempty"`
Status string `json:"status"`
OwnerAgentID *uuid.UUID `json:"owner_agent_id,omitempty"`
BlockedBy []uuid.UUID `json:"blocked_by,omitempty"`
Priority int `json:"priority"`
Result *string `json:"result,omitempty"`
TeamID uuid.UUID `json:"team_id"`
Subject string `json:"subject"`
Description string `json:"description,omitempty"`
Status string `json:"status"`
OwnerAgentID *uuid.UUID `json:"owner_agent_id,omitempty"`
BlockedBy []uuid.UUID `json:"blocked_by,omitempty"`
Priority int `json:"priority"`
Result *string `json:"result,omitempty"`
Metadata map[string]interface{} `json:"metadata,omitempty"`
// Joined fields
OwnerAgentKey string `json:"owner_agent_key,omitempty"`
@@ -87,20 +88,21 @@ type TeamTaskData struct {
// DelegationHistoryData represents a persisted delegation record.
type DelegationHistoryData struct {
BaseModel
SourceAgentID uuid.UUID `json:"source_agent_id"`
TargetAgentID uuid.UUID `json:"target_agent_id"`
TeamID *uuid.UUID `json:"team_id,omitempty"`
TeamTaskID *uuid.UUID `json:"team_task_id,omitempty"`
UserID string `json:"user_id,omitempty"`
Task string `json:"task"`
Mode string `json:"mode"`
Status string `json:"status"`
Result *string `json:"result,omitempty"`
Error *string `json:"error,omitempty"`
Iterations int `json:"iterations"`
TraceID *uuid.UUID `json:"trace_id,omitempty"`
DurationMS int `json:"duration_ms"`
CompletedAt *time.Time `json:"completed_at,omitempty"`
SourceAgentID uuid.UUID `json:"source_agent_id"`
TargetAgentID uuid.UUID `json:"target_agent_id"`
TeamID *uuid.UUID `json:"team_id,omitempty"`
TeamTaskID *uuid.UUID `json:"team_task_id,omitempty"`
UserID string `json:"user_id,omitempty"`
Task string `json:"task"`
Mode string `json:"mode"`
Status string `json:"status"`
Result *string `json:"result,omitempty"`
Error *string `json:"error,omitempty"`
Iterations int `json:"iterations"`
TraceID *uuid.UUID `json:"trace_id,omitempty"`
DurationMS int `json:"duration_ms"`
CompletedAt *time.Time `json:"completed_at,omitempty"`
Metadata map[string]interface{} `json:"metadata,omitempty"`
// Joined fields
SourceAgentKey string `json:"source_agent_key,omitempty"`
@@ -120,26 +122,29 @@ type DelegationHistoryListOpts struct {
// HandoffRouteData represents an active routing override for agent handoff.
type HandoffRouteData struct {
ID uuid.UUID `json:"id"`
Channel string `json:"channel"`
ChatID string `json:"chat_id"`
FromAgentKey string `json:"from_agent_key"`
ToAgentKey string `json:"to_agent_key"`
Reason string `json:"reason,omitempty"`
CreatedBy string `json:"created_by"`
CreatedAt time.Time `json:"created_at"`
ID uuid.UUID `json:"id"`
Channel string `json:"channel"`
ChatID string `json:"chat_id"`
FromAgentKey string `json:"from_agent_key"`
ToAgentKey string `json:"to_agent_key"`
Reason string `json:"reason,omitempty"`
CreatedBy string `json:"created_by"`
CreatedAt time.Time `json:"created_at"`
Metadata map[string]interface{} `json:"metadata,omitempty"`
}
// TeamMessageData represents a message in the team mailbox.
type TeamMessageData struct {
ID uuid.UUID `json:"id"`
TeamID uuid.UUID `json:"team_id"`
FromAgentID uuid.UUID `json:"from_agent_id"`
ToAgentID *uuid.UUID `json:"to_agent_id,omitempty"`
Content string `json:"content"`
MessageType string `json:"message_type"`
Read bool `json:"read"`
CreatedAt time.Time `json:"created_at"`
ID uuid.UUID `json:"id"`
TeamID uuid.UUID `json:"team_id"`
FromAgentID uuid.UUID `json:"from_agent_id"`
ToAgentID *uuid.UUID `json:"to_agent_id,omitempty"`
Content string `json:"content"`
MessageType string `json:"message_type"`
Read bool `json:"read"`
TaskID *uuid.UUID `json:"task_id,omitempty"`
Metadata map[string]interface{} `json:"metadata,omitempty"`
CreatedAt time.Time `json:"created_at"`
// Joined fields
FromAgentKey string `json:"from_agent_key,omitempty"`
@@ -195,4 +200,6 @@ type TeamStore interface {
SendMessage(ctx context.Context, msg *TeamMessageData) error
GetUnread(ctx context.Context, teamID, agentID uuid.UUID) ([]TeamMessageData, error)
MarkRead(ctx context.Context, messageID uuid.UUID) error
// ListMessages returns paginated team messages ordered by created_at DESC.
ListMessages(ctx context.Context, teamID uuid.UUID, limit, offset int) ([]TeamMessageData, int, error)
}
+94 -20
View File
@@ -41,7 +41,8 @@ type DelegationTask struct {
OriginTraceID uuid.UUID `json:"-"`
OriginRootSpanID uuid.UUID `json:"-"`
// Team task auto-completion
// Team tracking
TeamID uuid.UUID `json:"-"` // from link.TeamID (for delegation history)
TeamTaskID uuid.UUID `json:"-"`
cancelFunc context.CancelFunc `json:"-"`
@@ -74,6 +75,23 @@ type DelegateRunRequest struct {
type DelegateRunResult struct {
Content string
Iterations int
MediaPaths []string // media file paths from tool results (e.g. generated images)
}
// DelegateArtifacts holds forwarded artifacts from delegation results.
// Used to accumulate artifacts from intermediate completions until the final
// announce fires. New artifact types (files, voice, etc.) should be added here.
type DelegateArtifacts struct {
Media []string // file paths to forward (images, documents, audio, etc.)
Results []DelegateResultSummary // result summaries from completed delegations
}
// DelegateResultSummary is a compact representation of a delegation result
// included in the final announce so the lead has all results in one message.
type DelegateResultSummary struct {
AgentKey string
Content string
HasMedia bool
}
// AgentRunFunc runs an agent by key with the given request.
@@ -84,7 +102,8 @@ type AgentRunFunc func(ctx context.Context, agentKey string, req DelegateRunRequ
type DelegateResult struct {
Content string
Iterations int
DelegationID string // for async: the delegation ID to track/cancel
DelegationID string // for async: the delegation ID to track/cancel
MediaPaths []string // media file paths from delegation result
}
// linkSettings holds per-user restriction rules from agent_links.settings JSONB.
@@ -107,6 +126,7 @@ type DelegateManager struct {
hookEngine *hooks.Engine // optional: quality gate evaluation
active sync.Map // delegationID → *DelegationTask
pendingArtifacts sync.Map // sourceAgentID string → *DelegateArtifacts
completedMu sync.Mutex
completedSessions []string // session keys pending cleanup
}
@@ -190,7 +210,7 @@ func (dm *DelegateManager) Delegate(ctx context.Context, opts DelegateOpts) (*De
dm.saveDelegationHistory(task, result.Content, nil, duration)
slog.Info("delegation completed", "id", task.ID, "target", opts.TargetAgentKey, "iterations", result.Iterations)
return &DelegateResult{Content: result.Content, Iterations: result.Iterations, DelegationID: task.ID}, nil
return &DelegateResult{Content: result.Content, Iterations: result.Iterations, DelegationID: task.ID, MediaPaths: result.MediaPaths}, nil
}
// DelegateAsync spawns a delegation in the background and announces the result back.
@@ -227,25 +247,72 @@ func (dm *DelegateManager) DelegateAsync(ctx context.Context, opts DelegateOpts)
result, runErr := dm.runAgent(taskCtx, opts.TargetAgentKey, runReq)
duration := time.Since(startTime)
// Count sibling delegations still running (exclude self)
siblings := dm.ListActive(task.SourceAgentID)
siblingCount := 0
for _, s := range siblings {
if s.ID != task.ID {
siblingCount++
}
}
isLastDelegation := siblingCount == 0
// Announce result to parent via message bus
if dm.msgBus != nil && task.OriginChannel != "" {
elapsed := time.Since(task.CreatedAt)
dm.msgBus.PublishInbound(bus.InboundMessage{
Channel: "system",
SenderID: fmt.Sprintf("delegate:%s", task.ID),
ChatID: task.OriginChatID,
Content: formatDelegateAnnounce(task, result, runErr, elapsed),
UserID: task.UserID,
Metadata: map[string]string{
"origin_channel": task.OriginChannel,
"origin_peer_kind": task.OriginPeerKind,
"parent_agent": task.SourceAgentKey,
"delegation_id": task.ID,
"target_agent": task.TargetAgentKey,
"origin_trace_id": task.OriginTraceID.String(),
"origin_root_span_id": task.OriginRootSpanID.String(),
},
})
if siblingCount > 0 {
// Intermediate completion: accumulate artifacts + result summary.
// The final announce includes all sibling results so the lead doesn't
// need to call team_tasks to aggregate.
arts := &DelegateArtifacts{}
if result != nil {
arts.Media = result.MediaPaths
arts.Results = []DelegateResultSummary{{
AgentKey: task.TargetAgentKey,
Content: result.Content,
HasMedia: len(result.MediaPaths) > 0,
}}
} else if runErr != nil {
arts.Results = []DelegateResultSummary{{
AgentKey: task.TargetAgentKey,
Content: fmt.Sprintf("[failed] %s", runErr.Error()),
}}
}
dm.accumulateArtifacts(task.SourceAgentID, arts)
slog.Info("delegation announce suppressed (siblings still running)",
"id", task.ID, "target", task.TargetAgentKey, "siblings", siblingCount)
} else {
// Last completion: collect all accumulated artifacts + own result
artifacts := dm.collectArtifacts(task.SourceAgentID)
if result != nil {
artifacts.Media = append(artifacts.Media, result.MediaPaths...)
artifacts.Results = append(artifacts.Results, DelegateResultSummary{
AgentKey: task.TargetAgentKey,
Content: result.Content,
HasMedia: len(result.MediaPaths) > 0,
})
}
announceMsg := bus.InboundMessage{
Channel: "system",
SenderID: fmt.Sprintf("delegate:%s", task.ID),
ChatID: task.OriginChatID,
Content: formatDelegateAnnounce(task, artifacts, runErr, elapsed),
UserID: task.UserID,
Metadata: map[string]string{
"origin_channel": task.OriginChannel,
"origin_peer_kind": task.OriginPeerKind,
"parent_agent": task.SourceAgentKey,
"delegation_id": task.ID,
"target_agent": task.TargetAgentKey,
"origin_trace_id": task.OriginTraceID.String(),
"origin_root_span_id": task.OriginRootSpanID.String(),
},
Media: artifacts.Media,
}
dm.msgBus.PublishInbound(announceMsg)
}
}
if runErr != nil {
@@ -265,7 +332,9 @@ func (dm *DelegateManager) DelegateAsync(ctx context.Context, opts DelegateOpts)
resultContent := ""
if result != nil {
resultContent = result.Content
dm.autoCompleteTeamTask(task, resultContent)
if isLastDelegation {
dm.autoCompleteTeamTask(task, resultContent)
}
}
dm.saveDelegationHistory(task, resultContent, nil, duration)
}
@@ -357,6 +426,11 @@ func (dm *DelegateManager) prepareDelegation(ctx context.Context, opts DelegateO
TeamTaskID: opts.TeamTaskID,
}
// Carry team_id from the link (for delegation history filtering by team)
if link.TeamID != nil {
task.TeamID = *link.TeamID
}
return task, link, nil
}
+41 -17
View File
@@ -11,28 +11,52 @@ func (dm *DelegateManager) emitEvent(name string, task *DelegationTask) {
if dm.msgBus == nil {
return
}
payload := map[string]string{
"delegation_id": task.ID,
"source_agent": task.SourceAgentID.String(),
"target_agent": task.TargetAgentKey,
"user_id": task.UserID,
"mode": task.Mode,
}
if task.TeamID.String() != "00000000-0000-0000-0000-000000000000" {
payload["team_id"] = task.TeamID.String()
}
if task.TeamTaskID.String() != "00000000-0000-0000-0000-000000000000" {
payload["team_task_id"] = task.TeamTaskID.String()
}
dm.msgBus.Broadcast(bus.Event{
Name: name,
Payload: map[string]string{
"delegation_id": task.ID,
"source_agent": task.SourceAgentID.String(),
"target_agent": task.TargetAgentKey,
"user_id": task.UserID,
"mode": task.Mode,
},
Name: name,
Payload: payload,
})
}
func formatDelegateAnnounce(task *DelegationTask, result *DelegateRunResult, err error, elapsed time.Duration) string {
if err != nil {
func formatDelegateAnnounce(task *DelegationTask, artifacts *DelegateArtifacts, err error, elapsed time.Duration) string {
if err != nil && len(artifacts.Results) == 0 {
return fmt.Sprintf(
"[System Message] Delegation to agent %q failed.\n\nError: %s\n\nStats: runtime %s\n\n"+
"Handle the task yourself or try a different agent.",
"[System Message] All delegations finished. The last delegation to agent %q failed.\n\nError: %s\n\nStats: runtime %s\n\n"+
"Handle the failed task yourself or try a different agent.",
task.TargetAgentKey, err.Error(), elapsed.Round(time.Millisecond))
}
return fmt.Sprintf(
"[System Message] Delegation to agent %q completed.\n\nResult:\n%s\n\nStats: runtime %s, iterations %d\n\n"+
"Convert the result above into your normal assistant voice and send that user-facing update now. "+
"Keep internal details private. Reply ONLY: NO_REPLY if this exact result was already delivered to the user.",
task.TargetAgentKey, result.Content, elapsed.Round(time.Millisecond), result.Iterations)
msg := "[System Message] All team delegations completed.\n\n"
// Render each delegation result
for i, r := range artifacts.Results {
msg += fmt.Sprintf("--- Result from %q ---\n%s\n", r.AgentKey, r.Content)
if r.HasMedia {
msg += "[media file(s) attached — will be delivered automatically. Do NOT recreate or call create_image.]\n"
}
if i < len(artifacts.Results)-1 {
msg += "\n"
}
}
msg += fmt.Sprintf("\nStats: total elapsed %s\n\n", elapsed.Round(time.Millisecond))
msg += "Review the results above. You may:\n" +
"- Present a comprehensive summary to the user (if the task is fully done)\n" +
"- Delegate follow-up tasks to refine, combine, or extend these results\n" +
"- Ask a member to revise based on another member's output\n" +
"Any media files attached will be delivered automatically — do NOT recreate them."
return msg
}
+1 -1
View File
@@ -102,7 +102,7 @@ func (t *DelegateSearchTool) Execute(ctx context.Context, args map[string]interf
}, "", " ")
return NewResult(string(data) +
"\n\nUse `delegate(agent=\"<agent_key>\", task=\"your task\")` to delegate to one of these agents.")
"\n\nUse `spawn(agent=\"<agent_key>\", task=\"your task\")` to delegate to one of these agents.")
}
// hybridSearch merges FTS and embedding results with weighted scoring.
+46
View File
@@ -2,6 +2,7 @@ package tools
import (
"context"
"fmt"
"log/slog"
"time"
@@ -68,6 +69,30 @@ func (dm *DelegateManager) ActiveCountForTarget(targetID uuid.UUID) int {
return count
}
// accumulateArtifacts merges new artifacts into the pending set for a source agent.
// Called for intermediate delegation completions (when siblings are still running).
func (dm *DelegateManager) accumulateArtifacts(sourceAgentID uuid.UUID, arts *DelegateArtifacts) {
key := sourceAgentID.String()
existing, _ := dm.pendingArtifacts.Load(key)
var merged DelegateArtifacts
if existing != nil {
merged = *existing.(*DelegateArtifacts)
}
merged.Media = append(merged.Media, arts.Media...)
merged.Results = append(merged.Results, arts.Results...)
dm.pendingArtifacts.Store(key, &merged)
}
// collectArtifacts retrieves and removes all accumulated artifacts for a source agent.
// Called when the last delegation completes (siblingCount == 0).
func (dm *DelegateManager) collectArtifacts(sourceAgentID uuid.UUID) *DelegateArtifacts {
key := sourceAgentID.String()
if pending, ok := dm.pendingArtifacts.LoadAndDelete(key); ok {
return pending.(*DelegateArtifacts)
}
return &DelegateArtifacts{}
}
// trackCompleted records a delegate session key for deferred cleanup.
func (dm *DelegateManager) trackCompleted(task *DelegationTask) {
if dm.sessionStore == nil {
@@ -101,6 +126,7 @@ func (dm *DelegateManager) flushCompletedSessions() {
// autoCompleteTeamTask attempts to claim+complete the associated team task.
// Called after a delegation finishes successfully. Errors are logged but not fatal.
// On success, flushes all tracked delegate sessions (task done = context no longer needed).
// Also persists a team message record for audit trail / visualization.
func (dm *DelegateManager) autoCompleteTeamTask(task *DelegationTask, resultContent string) {
if dm.teamStore == nil || task.TeamTaskID == uuid.Nil {
return
@@ -114,6 +140,23 @@ func (dm *DelegateManager) autoCompleteTeamTask(task *DelegationTask, resultCont
"task_id", task.TeamTaskID, "delegation_id", task.ID)
// Task done — flush delegate sessions
dm.flushCompletedSessions()
// Persist delegation completion as team message for audit trail
if task.TeamID != uuid.Nil {
summary := resultContent
if len(summary) > 500 {
summary = summary[:500] + "..."
}
taskID := task.TeamTaskID
_ = dm.teamStore.SendMessage(context.Background(), &store.TeamMessageData{
TeamID: task.TeamID,
FromAgentID: task.TargetAgentID,
ToAgentID: &task.SourceAgentID,
Content: fmt.Sprintf("[Delegation completed] %s", summary),
MessageType: store.TeamMessageTypeChat,
TaskID: &taskID,
})
}
}
}
@@ -134,6 +177,9 @@ func (dm *DelegateManager) saveDelegationHistory(task *DelegationTask, resultCon
DurationMS: int(duration.Milliseconds()),
}
if task.TeamID != uuid.Nil {
record.TeamID = &task.TeamID
}
if task.TeamTaskID != uuid.Nil {
record.TeamTaskID = &task.TeamTaskID
}
-161
View File
@@ -1,161 +0,0 @@
package tools
import (
"context"
"encoding/json"
"fmt"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
// DelegateTool is a thin wrapper around DelegateManager.
// Supports actions: delegate (default), cancel, list.
type DelegateTool struct {
manager *DelegateManager
}
func NewDelegateTool(manager *DelegateManager) *DelegateTool {
return &DelegateTool{manager: manager}
}
func (t *DelegateTool) Name() string { return "delegate" }
func (t *DelegateTool) Description() string {
return "Delegate a task to another specialized agent, cancel a running delegation, or list active delegations. The target agent runs with its own identity, tools, and expertise. See AGENTS.md for available agents."
}
func (t *DelegateTool) Parameters() map[string]interface{} {
return map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"action": map[string]interface{}{
"type": "string",
"description": "'delegate' (default), 'cancel', or 'list'",
},
"agent": map[string]interface{}{
"type": "string",
"description": "Target agent key (required for action=delegate)",
},
"task": map[string]interface{}{
"type": "string",
"description": "Task description (required for action=delegate)",
},
"context": map[string]interface{}{
"type": "string",
"description": "Optional additional context for the target agent",
},
"mode": map[string]interface{}{
"type": "string",
"description": "'sync' (default, blocks until done) or 'async' (returns immediately, result announced later)",
},
"delegation_id": map[string]interface{}{
"type": "string",
"description": "Delegation ID to cancel (required for action=cancel)",
},
"team_task_id": map[string]interface{}{
"type": "string",
"description": "Team task ID to auto-complete when delegation finishes (optional, for team workflows)",
},
},
"required": []string{},
}
}
func (t *DelegateTool) Execute(ctx context.Context, args map[string]interface{}) *Result {
action, _ := args["action"].(string)
if action == "" {
action = "delegate"
}
switch action {
case "delegate":
return t.executeDelegation(ctx, args)
case "cancel":
return t.executeCancel(args)
case "list":
return t.executeList(ctx)
default:
return ErrorResult(fmt.Sprintf("unknown action: %s (use delegate, cancel, or list)", action))
}
}
func (t *DelegateTool) executeDelegation(ctx context.Context, args map[string]interface{}) *Result {
agentKey, _ := args["agent"].(string)
if agentKey == "" {
return ErrorResult("agent parameter is required for delegation")
}
task, _ := args["task"].(string)
if task == "" {
return ErrorResult("task parameter is required for delegation")
}
extraContext, _ := args["context"].(string)
mode, _ := args["mode"].(string)
if mode == "" {
mode = "sync"
}
var teamTaskID uuid.UUID
if ttID, _ := args["team_task_id"].(string); ttID != "" {
teamTaskID, _ = uuid.Parse(ttID)
}
opts := DelegateOpts{
TargetAgentKey: agentKey,
Task: task,
Context: extraContext,
Mode: mode,
TeamTaskID: teamTaskID,
}
if mode == "async" {
result, err := t.manager.DelegateAsync(ctx, opts)
if err != nil {
return ErrorResult(err.Error())
}
forLLM := fmt.Sprintf(`{"status":"accepted","delegation_id":%q,"target":%q,"mode":"async"}
Delegated to %q (async, id=%s). The result will be announced automatically when done — do NOT wait or poll.
Briefly tell the user what you've delegated and to whom. Be friendly and natural.`,
result.DelegationID, agentKey, agentKey, result.DelegationID)
return AsyncResult(forLLM)
}
// Sync (default)
result, err := t.manager.Delegate(ctx, opts)
if err != nil {
return ErrorResult(err.Error())
}
forLLM := fmt.Sprintf(
"Delegation to %q completed (%d iterations).\n\nResult:\n%s\n\n"+
"Present the information above to the user in YOUR OWN voice and persona. "+
"Do NOT adopt the delegate agent's personality, tone, or self-references. "+
"Rephrase and summarize naturally as yourself.",
agentKey, result.Iterations, result.Content)
return NewResult(forLLM)
}
func (t *DelegateTool) executeCancel(args map[string]interface{}) *Result {
delegationID, _ := args["delegation_id"].(string)
if delegationID == "" {
return ErrorResult("delegation_id is required for cancel action")
}
if t.manager.Cancel(delegationID) {
return NewResult(fmt.Sprintf("Delegation %s cancelled.", delegationID))
}
return ErrorResult(fmt.Sprintf("delegation %s not found or already completed", delegationID))
}
func (t *DelegateTool) executeList(ctx context.Context) *Result {
sourceAgentID := store.AgentIDFromContext(ctx)
tasks := t.manager.ListActive(sourceAgentID)
if len(tasks) == 0 {
return SilentResult(`{"delegations":[],"count":0}`)
}
out, _ := json.Marshal(map[string]interface{}{
"delegations": tasks,
"count": len(tasks),
})
return SilentResult(string(out))
}
+16
View File
@@ -9,9 +9,18 @@ import (
"strings"
"syscall"
"github.com/nextlevelbuilder/goclaw/internal/bootstrap"
"github.com/nextlevelbuilder/goclaw/internal/sandbox"
)
// virtualSystemFiles are files dynamically injected into the system prompt.
// They don't exist on disk — if the model tries to read them, return a hint.
var virtualSystemFiles = map[string]string{
bootstrap.TeamFile: "TEAM.md is already loaded in your system prompt. Refer to the TEAM.md section in your context above for team member information.",
bootstrap.DelegationFile: "DELEGATION.md is already loaded in your system prompt. Refer to the DELEGATION.md section in your context above for delegation instructions and available agents.",
bootstrap.AvailabilityFile: "AVAILABILITY.md is already loaded in your system prompt. Refer to the AVAILABILITY.md section in your context above for agent availability information.",
}
// ReadFileTool reads file contents, optionally through a sandbox container.
type ReadFileTool struct {
workspace string
@@ -89,6 +98,13 @@ func (t *ReadFileTool) Execute(ctx context.Context, args map[string]interface{})
}
}
// Virtual system files: TEAM.md, DELEGATION.md, AVAILABILITY.md are injected
// into the system prompt and don't exist on disk. Return a helpful hint.
baseName := filepath.Base(path)
if hint, ok := virtualSystemFiles[baseName]; ok {
return SilentResult(hint)
}
// Virtual FS: route memory files to DB (managed mode)
if t.memIntc != nil {
if content, handled, err := t.memIntc.ReadFile(ctx, path); handled {
+23 -4
View File
@@ -59,6 +59,10 @@ func (t *WriteFileTool) Parameters() map[string]interface{} {
"type": "string",
"description": "Content to write",
},
"deliver": map[string]interface{}{
"type": "boolean",
"description": "If true, deliver this file to the user as an attachment (image, document, etc.)",
},
},
"required": []string{"path", "content"},
}
@@ -67,6 +71,7 @@ func (t *WriteFileTool) Parameters() map[string]interface{} {
func (t *WriteFileTool) Execute(ctx context.Context, args map[string]interface{}) *Result {
path, _ := args["path"].(string)
content, _ := args["content"].(string)
deliver, _ := args["deliver"].(bool)
if path == "" {
return ErrorResult("path is required")
}
@@ -94,7 +99,7 @@ func (t *WriteFileTool) Execute(ctx context.Context, args map[string]interface{}
// Sandbox routing (sandboxKey from ctx — thread-safe)
sandboxKey := ToolSandboxKeyFromCtx(ctx)
if t.sandboxMgr != nil && sandboxKey != "" {
return t.executeInSandbox(ctx, path, content, sandboxKey)
return t.executeInSandbox(ctx, path, content, sandboxKey, deliver)
}
// Host execution — use per-user workspace from context if available (managed mode)
@@ -118,10 +123,14 @@ func (t *WriteFileTool) Execute(ctx context.Context, args map[string]interface{}
return ErrorResult(fmt.Sprintf("failed to write file: %v", err))
}
return SilentResult(fmt.Sprintf("File written: %s (%d bytes)", path, len(content)))
result := SilentResult(fmt.Sprintf("File written: %s (%d bytes)", path, len(content)))
if deliver {
result.Media = []string{resolved}
}
return result
}
func (t *WriteFileTool) executeInSandbox(ctx context.Context, path, content, sandboxKey string) *Result {
func (t *WriteFileTool) executeInSandbox(ctx context.Context, path, content, sandboxKey string, deliver bool) *Result {
bridge, err := t.getFsBridge(ctx, sandboxKey)
if err != nil {
return ErrorResult(fmt.Sprintf("sandbox error: %v", err))
@@ -131,7 +140,17 @@ func (t *WriteFileTool) executeInSandbox(ctx context.Context, path, content, san
return ErrorResult(fmt.Sprintf("failed to write file: %v", err))
}
return SilentResult(fmt.Sprintf("File written: %s (%d bytes)", path, len(content)))
result := SilentResult(fmt.Sprintf("File written: %s (%d bytes)", path, len(content)))
if deliver {
// Sandbox workspace is bind-mounted — resolve to host path for delivery
workspace := ToolWorkspaceFromCtx(ctx)
if workspace == "" {
workspace = t.workspace
}
hostPath := filepath.Join(workspace, path)
result.Media = []string{hostPath}
}
return result
}
func (t *WriteFileTool) getFsBridge(ctx context.Context, sandboxKey string) (*sandbox.FsBridge, error) {
+2 -2
View File
@@ -14,7 +14,7 @@ var toolGroups = map[string][]string{
"web": {"web_search", "web_fetch"},
"fs": {"read_file", "write_file", "list_files", "edit_file", "search", "glob"},
"runtime": {"exec", "process"},
"sessions": {"sessions_list", "sessions_history", "sessions_send", "sessions_spawn", "subagents", "session_status"},
"sessions": {"sessions_list", "sessions_history", "sessions_send", "sessions_spawn", "session_status"},
"ui": {"browser", "canvas"},
"automation": {"cron", "gateway"},
"messaging": {"message"},
@@ -24,7 +24,7 @@ var toolGroups = map[string][]string{
"goclaw": {
"browser", "canvas", "nodes", "cron", "message", "gateway",
"agents_list", "sessions_list", "sessions_history", "sessions_send",
"sessions_spawn", "subagents", "session_status",
"sessions_spawn", "session_status",
"memory_search", "memory_get", "web_search", "web_fetch", "read_image", "create_image",
},
}
+3
View File
@@ -11,6 +11,9 @@ type Result struct {
Async bool `json:"async"` // running asynchronously
Err error `json:"-"` // internal error (not serialized)
// Media holds file paths to forward as output (e.g. images from delegation).
Media []string `json:"-"`
// Usage holds token usage from tools that make internal LLM calls (e.g. read_image).
// When set, the agent loop records these on the tool span for tracing.
Usage *providers.Usage `json:"-"`
-1
View File
@@ -131,7 +131,6 @@ var SubagentDenyLeaf = []string{
"sessions_history",
"sessions_spawn",
"spawn",
"subagent",
}
// Spawn creates a new subagent task that runs asynchronously.
+340 -33
View File
@@ -2,81 +2,170 @@ package tools
import (
"context"
"encoding/json"
"fmt"
"strings"
"time"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
// SpawnTool is an async tool that spawns a subagent in the background.
// Per-call values (channel, chatID, peerKind, callback) are read from ctx for thread-safety.
// SpawnTool is the unified tool for spawning subagents and delegating to other agents.
// Replaces the old separate spawn, subagent, and delegate tools.
//
// Routing:
// - No agent param (or agent == self): subagent (clone self)
// - agent param set to a different agent: delegation (run target agent)
// - mode="sync": block until done; mode="async" (default): return immediately
type SpawnTool struct {
manager *SubagentManager
parentID string
depth int
subagentMgr *SubagentManager
delegateMgr *DelegateManager // nil in standalone mode; injected via SetDelegateManager
parentID string
depth int
}
func NewSpawnTool(manager *SubagentManager, parentID string, depth int) *SpawnTool {
return &SpawnTool{
manager: manager,
parentID: parentID,
depth: depth,
subagentMgr: manager,
parentID: parentID,
depth: depth,
}
}
func (t *SpawnTool) Name() string { return "spawn" }
// SetDelegateManager injects delegation capability (managed mode only).
func (t *SpawnTool) SetDelegateManager(dm *DelegateManager) { t.delegateMgr = dm }
func (t *SpawnTool) Name() string { return "spawn" }
func (t *SpawnTool) Description() string {
return "Spawn a subagent to handle a task in the background. The subagent runs independently and reports back when done. Use for complex or time-consuming tasks."
if t.delegateMgr != nil {
return "Spawn an agent to handle a task. Omit 'agent' to clone yourself, or specify 'agent' to delegate to a specialized agent. See DELEGATION.md for available agents."
}
return "Spawn a subagent to handle a task in the background. The subagent runs independently and reports back when done."
}
func (t *SpawnTool) Parameters() map[string]interface{} {
return map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"task": map[string]interface{}{
"type": "string",
"description": "The task for the subagent to complete",
},
"label": map[string]interface{}{
"type": "string",
"description": "Short label for the task (for display)",
},
"model": map[string]interface{}{
"type": "string",
"description": "Optional model override for this subagent (e.g. 'anthropic/claude-sonnet-4-5-20250929')",
},
props := map[string]interface{}{
"action": map[string]interface{}{
"type": "string",
"description": "'spawn' (default), 'list', 'cancel', or 'steer'",
},
"required": []string{"task"},
"task": map[string]interface{}{
"type": "string",
"description": "The task to complete (required for action=spawn)",
},
"mode": map[string]interface{}{
"type": "string",
"description": "'async' (default, returns immediately) or 'sync' (blocks until done)",
},
"label": map[string]interface{}{
"type": "string",
"description": "Short label for the task (for display)",
},
"model": map[string]interface{}{
"type": "string",
"description": "Optional model override (e.g. 'anthropic/claude-sonnet-4-5-20250929')",
},
"id": map[string]interface{}{
"type": "string",
"description": "Task ID for cancel/steer. For cancel: use 'all' to cancel all or 'last' for most recent",
},
"message": map[string]interface{}{
"type": "string",
"description": "New instructions (required for action=steer)",
},
}
// Add delegation-specific params when delegate manager is available
if t.delegateMgr != nil {
props["agent"] = map[string]interface{}{
"type": "string",
"description": "Target agent key. Omit to clone yourself, specify to delegate to another agent",
}
props["context"] = map[string]interface{}{
"type": "string",
"description": "Optional additional context for the target agent (used with agent param)",
}
props["team_task_id"] = map[string]interface{}{
"type": "string",
"description": "Team task ID to auto-complete when task finishes (for team workflows)",
}
}
return map[string]interface{}{
"type": "object",
"properties": props,
"required": []string{"task"},
}
}
func (t *SpawnTool) Execute(ctx context.Context, args map[string]interface{}) *Result {
action, _ := args["action"].(string)
if action == "" {
action = "spawn"
}
switch action {
case "list":
return t.executeList(ctx)
case "cancel":
return t.executeCancel(ctx, args)
case "steer":
return t.executeSteer(ctx, args)
default:
return t.executeSpawn(ctx, args)
}
}
// executeSpawn routes to subagent (self-clone) or delegation (different agent).
func (t *SpawnTool) executeSpawn(ctx context.Context, args map[string]interface{}) *Result {
task, _ := args["task"].(string)
if task == "" {
return ErrorResult("task parameter is required")
}
agentKey, _ := args["agent"].(string)
selfKey := ToolAgentKeyFromCtx(ctx)
if selfKey == "" {
selfKey = t.parentID
}
// If agent is specified and different from self → delegation
if agentKey != "" && agentKey != selfKey && t.delegateMgr != nil {
return t.executeDelegation(ctx, args, agentKey, task)
}
// Self-clone path
mode, _ := args["mode"].(string)
if mode == "sync" {
return t.executeSubagentSync(ctx, args, task)
}
return t.executeSubagentAsync(ctx, args, task)
}
// executeSubagentAsync spawns an async self-clone (old SpawnTool behavior).
func (t *SpawnTool) executeSubagentAsync(ctx context.Context, args map[string]interface{}, task string) *Result {
label, _ := args["label"].(string)
modelOverride, _ := args["model"].(string)
// Read per-call values from ctx (thread-safe)
channel := ToolChannelFromCtx(ctx)
chatID := ToolChatIDFromCtx(ctx)
peerKind := ToolPeerKindFromCtx(ctx)
callback := ToolAsyncCBFromCtx(ctx)
// Resolve parent agent from ctx (managed mode) with fallback to construction-time default
parentID := ToolAgentKeyFromCtx(ctx)
if parentID == "" {
parentID = t.parentID
}
msg, err := t.manager.Spawn(ctx, parentID, t.depth, task, label, modelOverride,
msg, err := t.subagentMgr.Spawn(ctx, parentID, t.depth, task, label, modelOverride,
channel, chatID, peerKind, callback)
if err != nil {
return ErrorResult(err.Error())
}
// Match TS pattern: return structured status + instruction for the LLM.
// TS returns {status: "accepted", childSessionKey, runId} and the LLM
// naturally acknowledges. We add a brief instruction since weaker models
// may return empty content after multiple spawn calls.
forLLM := fmt.Sprintf(`{"status":"accepted","label":%q}
%s
After all spawn tool calls in this turn are complete, briefly tell the user what tasks you've started. Subagents will announce results when done — do NOT wait or poll.`, label, msg)
@@ -84,6 +173,189 @@ After all spawn tool calls in this turn are complete, briefly tell the user what
return AsyncResult(forLLM)
}
// executeSubagentSync runs a sync self-clone (old SubagentTool action=run behavior).
func (t *SpawnTool) executeSubagentSync(ctx context.Context, args map[string]interface{}, task string) *Result {
label, _ := args["label"].(string)
if label == "" {
label = truncate(task, 50)
}
channel := ToolChannelFromCtx(ctx)
chatID := ToolChatIDFromCtx(ctx)
parentID := ToolAgentKeyFromCtx(ctx)
if parentID == "" {
parentID = t.parentID
}
result, iterations, err := t.subagentMgr.RunSync(ctx, parentID, t.depth, task, label,
channel, chatID)
if err != nil {
return ErrorResult(fmt.Sprintf("Subagent '%s' failed: %v", label, err))
}
forUser := fmt.Sprintf("Subagent '%s' completed.", label)
if len(result) > 500 {
forUser += "\n" + result[:500] + "..."
} else {
forUser += "\n" + result
}
forLLM := fmt.Sprintf("Subagent '%s' completed in %d iterations.\n\nFull result:\n%s",
label, iterations, result)
return &Result{ForLLM: forLLM, ForUser: forUser}
}
// executeDelegation delegates to a different agent (old DelegateTool behavior).
func (t *SpawnTool) executeDelegation(ctx context.Context, args map[string]interface{}, agentKey, task string) *Result {
extraContext, _ := args["context"].(string)
mode, _ := args["mode"].(string)
if mode == "" {
mode = "async"
}
var teamTaskID uuid.UUID
if ttID, _ := args["team_task_id"].(string); ttID != "" {
teamTaskID, _ = uuid.Parse(ttID)
}
opts := DelegateOpts{
TargetAgentKey: agentKey,
Task: task,
Context: extraContext,
Mode: mode,
TeamTaskID: teamTaskID,
}
if mode == "async" {
result, err := t.delegateMgr.DelegateAsync(ctx, opts)
if err != nil {
return ErrorResult(err.Error())
}
forLLM := fmt.Sprintf(`{"status":"accepted","delegation_id":%q,"target":%q,"mode":"async"}
Delegated to %q (async, id=%s). The result will be announced automatically when done — do NOT wait or poll.
Briefly tell the user what you've delegated and to whom. Be friendly and natural.`,
result.DelegationID, agentKey, agentKey, result.DelegationID)
return AsyncResult(forLLM)
}
// Sync delegation
result, err := t.delegateMgr.Delegate(ctx, opts)
if err != nil {
return ErrorResult(err.Error())
}
mediaNote := ""
if len(result.MediaPaths) > 0 {
mediaNote = fmt.Sprintf("\n\n[%d media file(s) attached — will be delivered automatically. Do NOT recreate or call create_image.]",
len(result.MediaPaths))
}
forLLM := fmt.Sprintf(
"Delegation to %q completed (%d iterations).\n\nResult:\n%s%s\n\n"+
"Present the information above to the user in YOUR OWN voice and persona. "+
"Do NOT adopt the delegate agent's personality, tone, or self-references. "+
"Rephrase and summarize naturally as yourself.",
agentKey, result.Iterations, result.Content, mediaNote)
toolResult := NewResult(forLLM)
if len(result.MediaPaths) > 0 {
toolResult.Media = result.MediaPaths
}
return toolResult
}
// executeList shows active subagents and delegations.
func (t *SpawnTool) executeList(ctx context.Context) *Result {
var sections []string
// Subagent tasks
parentID := ToolAgentKeyFromCtx(ctx)
if parentID == "" {
parentID = t.parentID
}
tasks := t.subagentMgr.ListTasks(parentID)
if len(tasks) > 0 {
var lines []string
running, completed, cancelled := 0, 0, 0
for _, task := range tasks {
switch task.Status {
case "running":
running++
case "completed":
completed++
case "cancelled":
cancelled++
}
line := fmt.Sprintf("- [%s] %s (id=%s, status=%s)", task.Label, truncate(task.Task, 60), task.ID, task.Status)
if task.CompletedAt > 0 {
dur := time.Duration(task.CompletedAt-task.CreatedAt) * time.Millisecond
line += fmt.Sprintf(", took %s", dur.Round(time.Millisecond))
}
lines = append(lines, line)
}
sections = append(sections, fmt.Sprintf("Subagent tasks: %d running, %d completed, %d cancelled\n%s",
running, completed, cancelled, strings.Join(lines, "\n")))
}
// Delegation tasks
if t.delegateMgr != nil {
sourceAgentID := store.AgentIDFromContext(ctx)
delegations := t.delegateMgr.ListActive(sourceAgentID)
if len(delegations) > 0 {
out, _ := json.Marshal(map[string]interface{}{
"delegations": delegations,
"count": len(delegations),
})
sections = append(sections, "Delegations:\n"+string(out))
}
}
if len(sections) == 0 {
return &Result{ForLLM: "No active tasks found."}
}
return &Result{ForLLM: strings.Join(sections, "\n\n")}
}
// executeCancel cancels a subagent or delegation by ID.
func (t *SpawnTool) executeCancel(ctx context.Context, args map[string]interface{}) *Result {
id, _ := args["id"].(string)
if id == "" {
return ErrorResult("id is required for action=cancel")
}
// Try subagent first
if t.subagentMgr.CancelTask(id) {
return &Result{ForLLM: fmt.Sprintf("Task '%s' cancelled.", id)}
}
// Try delegation
if t.delegateMgr != nil && t.delegateMgr.Cancel(id) {
return NewResult(fmt.Sprintf("Delegation '%s' cancelled.", id))
}
return ErrorResult(fmt.Sprintf("Task '%s' not found or not running.", id))
}
// executeSteer redirects a running subagent with new instructions.
func (t *SpawnTool) executeSteer(ctx context.Context, args map[string]interface{}) *Result {
id, _ := args["id"].(string)
if id == "" {
return ErrorResult("id is required for action=steer")
}
message, _ := args["message"].(string)
if message == "" {
return ErrorResult("message is required for action=steer")
}
msg, err := t.subagentMgr.Steer(ctx, id, message, nil)
if err != nil {
return ErrorResult(err.Error())
}
return &Result{ForLLM: msg}
}
// SetContext is a no-op; channel/chatID are now read from ctx (thread-safe).
func (t *SpawnTool) SetContext(channel, chatID string) {}
@@ -92,3 +364,38 @@ func (t *SpawnTool) SetPeerKind(peerKind string) {}
// SetCallback is a no-op; callback is now read from ctx (thread-safe).
func (t *SpawnTool) SetCallback(cb AsyncCallback) {}
// --- Helpers moved from old subagent_tool.go ---
// FilterDenyList returns tool names from the registry excluding denied tools.
func FilterDenyList(reg *Registry, denyList []string) []string {
deny := make(map[string]bool, len(denyList))
for _, n := range denyList {
deny[n] = true
}
var allowed []string
for _, name := range reg.List() {
if !deny[name] {
allowed = append(allowed, name)
}
}
return allowed
}
// IsSubagentDenied checks if a tool name is in the subagent deny list.
func IsSubagentDenied(toolName string, depth, maxDepth int) bool {
for _, d := range SubagentDenyAlways {
if strings.EqualFold(toolName, d) {
return true
}
}
if depth >= maxDepth {
for _, d := range SubagentDenyLeaf {
if strings.EqualFold(toolName, d) {
return true
}
}
}
return false
}
-215
View File
@@ -1,215 +0,0 @@
package tools
import (
"context"
"fmt"
"strings"
"time"
)
// SubagentTool manages subagents: run tasks synchronously, list running tasks, or cancel them.
// Per-call values (channel, chatID, peerKind) are read from ctx for thread-safety.
type SubagentTool struct {
manager *SubagentManager
parentID string
depth int
}
func NewSubagentTool(manager *SubagentManager, parentID string, depth int) *SubagentTool {
return &SubagentTool{
manager: manager,
parentID: parentID,
depth: depth,
}
}
func (t *SubagentTool) Name() string { return "subagent" }
func (t *SubagentTool) Description() string {
return "Manage subagents: run a task synchronously, list active/completed tasks, cancel a running task, or steer (redirect) a running task with new instructions."
}
func (t *SubagentTool) Parameters() map[string]interface{} {
return map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"action": map[string]interface{}{
"type": "string",
"enum": []string{"run", "list", "cancel", "steer"},
"description": "Action: 'run' (default) executes a task synchronously, 'list' shows active/completed subagents, 'cancel' stops a running subagent (use subagent_id='all' or 'last' for bulk), 'steer' cancels and restarts a subagent with new instructions",
},
"task": map[string]interface{}{
"type": "string",
"description": "The task for the subagent to complete (required for action=run)",
},
"label": map[string]interface{}{
"type": "string",
"description": "Short label for the task (for display, used with action=run)",
},
"subagent_id": map[string]interface{}{
"type": "string",
"description": "Subagent ID (required for action=cancel/steer). For cancel: use 'all' to cancel all or 'last' for most recent",
},
"message": map[string]interface{}{
"type": "string",
"description": "New instructions for the subagent (required for action=steer)",
},
},
}
}
func (t *SubagentTool) Execute(ctx context.Context, args map[string]interface{}) *Result {
action, _ := args["action"].(string)
if action == "" {
action = "run"
}
switch action {
case "list":
return t.executeList()
case "cancel":
return t.executeCancel(args)
case "steer":
return t.executeSteer(ctx, args)
default:
return t.executeRun(ctx, args)
}
}
func (t *SubagentTool) executeList() *Result {
tasks := t.manager.ListTasks(t.parentID)
if len(tasks) == 0 {
return &Result{ForLLM: "No subagent tasks found."}
}
var lines []string
running, completed, cancelled := 0, 0, 0
for _, task := range tasks {
status := task.Status
switch status {
case "running":
running++
case "completed":
completed++
case "cancelled":
cancelled++
}
line := fmt.Sprintf("- [%s] %s (id=%s, status=%s)", task.Label, truncate(task.Task, 60), task.ID, status)
if task.CompletedAt > 0 {
dur := time.Duration(task.CompletedAt-task.CreatedAt) * time.Millisecond
line += fmt.Sprintf(", took %s", dur.Round(time.Millisecond))
}
lines = append(lines, line)
}
summary := fmt.Sprintf("Subagent tasks: %d running, %d completed, %d cancelled\n\n%s",
running, completed, cancelled, strings.Join(lines, "\n"))
return &Result{ForLLM: summary}
}
func (t *SubagentTool) executeCancel(args map[string]interface{}) *Result {
id, _ := args["subagent_id"].(string)
if id == "" {
return ErrorResult("subagent_id is required for action=cancel")
}
if t.manager.CancelTask(id) {
return &Result{ForLLM: fmt.Sprintf("Subagent '%s' has been cancelled.", id)}
}
return ErrorResult(fmt.Sprintf("Subagent '%s' not found or not running.", id))
}
func (t *SubagentTool) executeSteer(ctx context.Context, args map[string]interface{}) *Result {
id, _ := args["subagent_id"].(string)
if id == "" {
return ErrorResult("subagent_id is required for action=steer")
}
message, _ := args["message"].(string)
if message == "" {
return ErrorResult("message is required for action=steer")
}
msg, err := t.manager.Steer(ctx, id, message, nil)
if err != nil {
return ErrorResult(err.Error())
}
return &Result{ForLLM: msg}
}
func (t *SubagentTool) executeRun(ctx context.Context, args map[string]interface{}) *Result {
task, _ := args["task"].(string)
if task == "" {
return ErrorResult("task parameter is required for action=run")
}
label, _ := args["label"].(string)
if label == "" {
label = truncate(task, 50)
}
// Read per-call values from ctx (thread-safe)
channel := ToolChannelFromCtx(ctx)
chatID := ToolChatIDFromCtx(ctx)
// Resolve parent agent from ctx (managed mode) with fallback to construction-time default
parentID := ToolAgentKeyFromCtx(ctx)
if parentID == "" {
parentID = t.parentID
}
result, iterations, err := t.manager.RunSync(ctx, parentID, t.depth, task, label,
channel, chatID)
if err != nil {
return ErrorResult(fmt.Sprintf("Subagent '%s' failed: %v", label, err))
}
forUser := fmt.Sprintf("Subagent '%s' completed.", label)
if len(result) > 500 {
forUser += "\n" + result[:500] + "..."
} else {
forUser += "\n" + result
}
forLLM := fmt.Sprintf("Subagent '%s' completed in %d iterations.\n\nFull result:\n%s",
label, iterations, result)
return &Result{ForLLM: forLLM, ForUser: forUser}
}
// SetContext is a no-op; channel/chatID are now read from ctx (thread-safe).
func (t *SubagentTool) SetContext(channel, chatID string) {}
// SetPeerKind is a no-op; peerKind is now read from ctx (thread-safe).
func (t *SubagentTool) SetPeerKind(peerKind string) {}
// --- Helper: filter tools by name ---
// FilterDenyList returns tool names from the registry excluding denied tools.
func FilterDenyList(reg *Registry, denyList []string) []string {
deny := make(map[string]bool, len(denyList))
for _, n := range denyList {
deny[n] = true
}
var allowed []string
for _, name := range reg.List() {
if !deny[name] {
allowed = append(allowed, name)
}
}
return allowed
}
// IsSubagentDenied checks if a tool name is in the subagent deny list.
func IsSubagentDenied(toolName string, depth, maxDepth int) bool {
for _, d := range SubagentDenyAlways {
if strings.EqualFold(toolName, d) {
return true
}
}
if depth >= maxDepth {
for _, d := range SubagentDenyLeaf {
if strings.EqualFold(toolName, d) {
return true
}
}
}
return false
}
+23
View File
@@ -7,6 +7,7 @@ import (
"github.com/nextlevelbuilder/goclaw/internal/bus"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/pkg/protocol"
)
// TeamMessageTool exposes the team mailbox to agents.
@@ -97,6 +98,17 @@ func (t *TeamMessageTool) executeSend(ctx context.Context, args map[string]inter
fromKey := t.manager.agentKeyFromID(ctx, agentID)
t.publishTeammateMessage(fromKey, toKey, text, ctx)
preview := text
if len(preview) > 100 {
preview = preview[:100] + "..."
}
t.manager.broadcastTeamEvent(protocol.EventTeamMessageSent, map[string]string{
"team_id": team.ID.String(),
"from": fromKey,
"to": toKey,
"preview": preview,
})
return NewResult(fmt.Sprintf("Message sent to %s.", toKey))
}
@@ -135,6 +147,17 @@ func (t *TeamMessageTool) executeBroadcast(ctx context.Context, args map[string]
}
}
preview := text
if len(preview) > 100 {
preview = preview[:100] + "..."
}
t.manager.broadcastTeamEvent(protocol.EventTeamMessageSent, map[string]string{
"team_id": team.ID.String(),
"from": fromKey,
"to": "broadcast",
"preview": preview,
})
return NewResult(fmt.Sprintf("Broadcast sent to all teammates."))
}
+31 -2
View File
@@ -8,6 +8,7 @@ import (
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/pkg/protocol"
)
// TeamTasksTool exposes the shared team task list to agents.
@@ -93,6 +94,8 @@ func (t *TeamTasksTool) Execute(ctx context.Context, args map[string]interface{}
}
}
const listTasksLimit = 20
func (t *TeamTasksTool) executeList(ctx context.Context, args map[string]interface{}) *Result {
team, _, err := t.manager.resolveTeam(ctx)
if err != nil {
@@ -111,10 +114,21 @@ func (t *TeamTasksTool) executeList(ctx context.Context, args map[string]interfa
tasks[i].Result = nil
}
out, _ := json.Marshal(map[string]interface{}{
hasMore := len(tasks) > listTasksLimit
if hasMore {
tasks = tasks[:listTasksLimit]
}
resp := map[string]interface{}{
"tasks": tasks,
"count": len(tasks),
})
}
if hasMore {
resp["note"] = fmt.Sprintf("Showing first %d tasks. Use action=search with a query to find older tasks.", listTasksLimit)
resp["has_more"] = true
}
out, _ := json.Marshal(resp)
return SilentResult(string(out))
}
@@ -228,6 +242,13 @@ func (t *TeamTasksTool) executeCreate(ctx context.Context, args map[string]inter
return ErrorResult("failed to create task: " + err.Error())
}
t.manager.broadcastTeamEvent(protocol.EventTeamTaskCreated, map[string]string{
"team_id": team.ID.String(),
"task_id": task.ID.String(),
"subject": subject,
"status": status,
})
return NewResult(fmt.Sprintf("Task created: %s (id=%s, status=%s)", subject, task.ID, status))
}
@@ -282,5 +303,13 @@ func (t *TeamTasksTool) executeComplete(ctx context.Context, args map[string]int
return ErrorResult("failed to complete task: " + err.Error())
}
// Resolve team for event payload
if team, _, teamErr := t.manager.resolveTeam(ctx); teamErr == nil {
t.manager.broadcastTeamEvent(protocol.EventTeamTaskCompleted, map[string]string{
"team_id": team.ID.String(),
"task_id": taskIDStr,
})
}
return NewResult(fmt.Sprintf("Task %s completed. Dependent tasks have been unblocked.", taskIDStr))
}
+11
View File
@@ -58,3 +58,14 @@ func (m *TeamToolManager) agentKeyFromID(ctx context.Context, id uuid.UUID) stri
}
return ag.AgentKey
}
// broadcastTeamEvent sends a real-time event via the message bus for team activity visibility.
func (m *TeamToolManager) broadcastTeamEvent(name string, payload map[string]string) {
if m.msgBus == nil {
return
}
m.msgBus.Broadcast(bus.Event{
Name: name,
Payload: payload,
})
}
+1 -1
View File
@@ -2,4 +2,4 @@ package upgrade
// RequiredSchemaVersion is the schema migration version this binary requires.
// Bump this whenever adding a new SQL migration file.
const RequiredSchemaVersion uint = 6
const RequiredSchemaVersion uint = 7
+5
View File
@@ -0,0 +1,5 @@
ALTER TABLE team_messages DROP COLUMN IF EXISTS task_id;
ALTER TABLE team_messages DROP COLUMN IF EXISTS metadata;
ALTER TABLE team_tasks DROP COLUMN IF EXISTS metadata;
ALTER TABLE delegation_history DROP COLUMN IF EXISTS metadata;
ALTER TABLE handoff_routes DROP COLUMN IF EXISTS metadata;
+9
View File
@@ -0,0 +1,9 @@
-- Add task_id to team_messages for linking messages to tasks
ALTER TABLE team_messages ADD COLUMN IF NOT EXISTS task_id UUID REFERENCES team_tasks(id) ON DELETE SET NULL;
CREATE INDEX IF NOT EXISTS idx_team_messages_task ON team_messages(task_id) WHERE task_id IS NOT NULL;
-- Add metadata JSONB to all team-related tables
ALTER TABLE team_messages ADD COLUMN IF NOT EXISTS metadata JSONB NOT NULL DEFAULT '{}';
ALTER TABLE team_tasks ADD COLUMN IF NOT EXISTS metadata JSONB NOT NULL DEFAULT '{}';
ALTER TABLE delegation_history ADD COLUMN IF NOT EXISTS metadata JSONB NOT NULL DEFAULT '{}';
ALTER TABLE handoff_routes ADD COLUMN IF NOT EXISTS metadata JSONB NOT NULL DEFAULT '{}';
+7
View File
@@ -26,6 +26,13 @@ const (
// Agent handoff event (payload: from_agent, to_agent, reason).
EventHandoff = "handoff"
// Team activity events (real-time team workflow visibility).
EventTeamTaskCreated = "team.task.created"
EventTeamTaskCompleted = "team.task.completed"
EventTeamMessageSent = "team.message.sent"
EventDelegationStarted = "delegation.started"
EventDelegationCompleted = "delegation.completed"
// Cache invalidation events (internal, not forwarded to WS clients).
EventCacheInvalidate = "cache.invalidate"
)
@@ -62,6 +62,12 @@ export function AgentCreateDialog({ open, onOpenChange, onCreate }: AgentCreateD
await verify(selectedProviderId, model.trim());
};
const handleVerifyAndCreate = async () => {
if (!selectedProviderId || !model.trim()) return;
const res = await verify(selectedProviderId, model.trim());
if (res?.valid) await handleCreate();
};
const handleCreate = async () => {
if (!agentKey.trim()) return;
setLoading(true);
@@ -251,6 +257,10 @@ export function AgentCreateDialog({ open, onOpenChange, onCreate }: AgentCreateD
</Button>
{loading ? (
<Button disabled>Creating...</Button>
) : !verifyResult?.valid && selectedProviderId && model.trim() ? (
<Button onClick={handleVerifyAndCreate} disabled={verifying || !displayName.trim() || !agentKey.trim() || !isValidSlug(agentKey)}>
{verifying ? "Checking..." : "Check & Create"}
</Button>
) : (
<Button onClick={handleCreate} disabled={!displayName.trim() || !agentKey.trim() || !isValidSlug(agentKey) || !provider.trim() || !model.trim() || !verifyResult?.valid}>
Create
@@ -1,9 +1,12 @@
import { useState, useEffect, useMemo } from "react";
import { Link } from "react-router";
import { Info } from "lucide-react";
import { ConfirmDialog } from "@/components/shared/confirm-dialog";
import { useAgentLinks } from "../hooks/use-agent-links";
import { useAgents } from "../hooks/use-agents";
import type { AgentLinkData } from "@/types/agent";
import { LinkCreateForm, LinkList, LinkEditDialog, linkTargetName } from "./link-sections";
import { ROUTES } from "@/lib/constants";
interface AgentLinksTabProps {
agentId: string;
@@ -43,6 +46,17 @@ export function AgentLinksTab({ agentId }: AgentLinksTabProps) {
return (
<div className="max-w-4xl space-y-6">
<div className="flex items-start gap-3 rounded-lg border border-blue-200 bg-blue-50 px-4 py-3 dark:border-blue-900 dark:bg-blue-950/30">
<Info className="mt-0.5 h-4 w-4 shrink-0 text-blue-600 dark:text-blue-400" />
<p className="text-sm text-blue-800 dark:text-blue-300">
For multi-agent collaboration with task tracking and quality control, consider using{" "}
<Link to={ROUTES.TEAMS} className="font-medium underline underline-offset-2 hover:text-blue-900 dark:hover:text-blue-200">
Agent Teams
</Link>{" "}
instead. Agent Links are best for simple 1-to-1 delegation between two agents.
</p>
</div>
<LinkCreateForm agentOptions={agentOptions} onSubmit={createLink} />
<LinkList
@@ -38,7 +38,7 @@ export function RegenerateDialog({
return (
<Dialog open={open} onOpenChange={onOpenChange}>
<DialogContent className="sm:max-w-md">
<DialogContent className="sm:max-w-lg">
<DialogHeader>
<DialogTitle className="flex items-center gap-2">
<Sparkles className="h-4 w-4" />
@@ -54,7 +54,7 @@ export function RegenerateDialog({
value={prompt}
onChange={(e) => setPrompt(e.target.value)}
placeholder="e.g. Make the agent more formal, add Vietnamese language support, change the name to Luna..."
className="min-h-[100px]"
className="min-h-[100px] max-h-[300px] resize-none"
/>
</div>
<DialogFooter>
@@ -45,13 +45,6 @@ export function SummoningModal({
}
}, [open]);
// Auto-close modal after completion
useEffect(() => {
if (status === "completed") {
const timer = setTimeout(() => onOpenChange(false), 1500);
return () => clearTimeout(timer);
}
}, [status, onOpenChange]);
const handleSummoningEvent = useCallback(
(payload: unknown) => {
@@ -30,7 +30,7 @@ const CATEGORY_LABELS: Record<string, string> = {
sessions: "Sessions",
messaging: "Messaging",
scheduling: "Scheduling",
subagents: "Subagents",
subagents: "Agents",
skills: "Skills",
delegation: "Delegation",
teams: "Teams",
@@ -1,5 +1,4 @@
import { useState, useEffect } from "react";
import { Link } from "react-router";
import { useState, useEffect, useCallback } from "react";
import {
Dialog,
DialogContent,
@@ -8,7 +7,10 @@ import {
} from "@/components/ui/dialog";
import { Badge } from "@/components/ui/badge";
import { formatDate, formatDuration } from "@/lib/format";
import { useHttp } from "@/hooks/use-ws";
import { TraceDetailDialog } from "@/pages/traces/trace-detail-dialog";
import type { DelegationHistoryRecord } from "@/types/delegation";
import type { TraceData, SpanData } from "@/types/trace";
interface DelegationDetailDialogProps {
delegationId: string;
@@ -19,6 +21,19 @@ interface DelegationDetailDialogProps {
export function DelegationDetailDialog({ delegationId, onClose, getDelegation }: DelegationDetailDialogProps) {
const [record, setRecord] = useState<DelegationHistoryRecord | null>(null);
const [loading, setLoading] = useState(true);
const [viewingTraceId, setViewingTraceId] = useState<string | null>(null);
const http = useHttp();
const getTrace = useCallback(
async (traceId: string): Promise<{ trace: TraceData; spans: SpanData[] } | null> => {
try {
return await http.get<{ trace: TraceData; spans: SpanData[] }>(`/v1/traces/${traceId}`);
} catch {
return null;
}
},
[http],
);
useEffect(() => {
setLoading(true);
@@ -38,7 +53,7 @@ export function DelegationDetailDialog({ delegationId, onClose, getDelegation }:
return (
<Dialog open onOpenChange={() => onClose()}>
<DialogContent className="max-h-[85vh] overflow-y-auto sm:max-w-3xl">
<DialogContent className="max-h-[85vh] overflow-y-auto sm:max-w-4xl">
<DialogHeader>
<DialogTitle>Delegation Detail</DialogTitle>
</DialogHeader>
@@ -90,12 +105,13 @@ export function DelegationDetailDialog({ delegationId, onClose, getDelegation }:
{record.trace_id && (
<div className="text-sm">
<span className="text-muted-foreground">Trace:</span>{" "}
<Link
to={`/traces/${record.trace_id}`}
className="font-mono text-xs text-primary hover:underline"
<button
type="button"
onClick={() => setViewingTraceId(record.trace_id!)}
className="cursor-pointer font-mono text-xs text-primary hover:underline"
>
{record.trace_id.slice(0, 12)}...
</Link>
</button>
</div>
)}
@@ -125,6 +141,15 @@ export function DelegationDetailDialog({ delegationId, onClose, getDelegation }:
</div>
)}
</DialogContent>
{viewingTraceId && (
<TraceDetailDialog
traceId={viewingTraceId}
onClose={() => setViewingTraceId(null)}
getTrace={getTrace}
onNavigateTrace={setViewingTraceId}
/>
)}
</Dialog>
);
}
+4 -1
View File
@@ -1,5 +1,5 @@
import { useState } from "react";
import { Activity, RefreshCw, Search } from "lucide-react";
import { Activity, GitFork, RefreshCw, Search } from "lucide-react";
import { Button } from "@/components/ui/button";
import { Badge } from "@/components/ui/badge";
import { Input } from "@/components/ui/input";
@@ -93,6 +93,9 @@ export function TracesPage() {
onClick={() => setSelectedTraceId(trace.id)}
>
<td className="max-w-[200px] truncate px-4 py-3 font-medium">
{trace.parent_trace_id && (
<GitFork className="mr-1.5 inline-block h-3.5 w-3.5 text-muted-foreground" />
)}
{trace.name || "Unnamed"}
{trace.channel && (
<Badge variant="outline" className="ml-2 text-xs">