From 99c90523086945c0ddf0a032fa31e221a24304c6 Mon Sep 17 00:00:00 2001 From: ZanzyTHEbar Date: Sat, 21 Feb 2026 20:39:30 +0000 Subject: [PATCH] refactor(agent): decompose AgentLoop into focused single-responsibility modules MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The monolithic loop.go (previously ~2500 LOC with a FIXME to decompose it) has been split into five focused files following the Ports-and-Adapters principle. loop.go is now a thin orchestrator (~300 LOC) that wires the components together. New files: - pkg/agent/agent_run.go: core run path — prepareRuntimeState, runAgentLoop, assembleContext, resolveFinalContent, and the Generate/Stream dispatch. Owns conversationIDs writes and activeSessionKey writes. - pkg/agent/command_handler.go: slash-command registry (SlashCommand type), built-in /help, /models, /clear, /memory, /debug handlers, and the dispatchCommand entry point. - pkg/agent/helpers.go: stateless helper functions shared across the package — initSecretStore, loadIdentityDocs, seedSystemPrompt, buildDenyPatterns, resolveModelForChannel, and similar utilities. - pkg/agent/message_router.go: inbound message routing — RecordLastChannel, RecordLastChatID, RouteMessage, and the channel-specific dispatch logic that was previously embedded in the main loop. - pkg/agent/summarizer.go: async session summarization — forceCompression, summarizeSession, and the summarizing/summarizeFailures sync.Map owners. Updated files: - pkg/agent/loop.go: stripped to struct definition, NewAgentLoop constructor, Run/Stop lifecycle, and field ownership comments; removed the FIXME banner. - pkg/agent/context.go: minor import cleanup. - pkg/agent/toolloop.go: minor import cleanup. - pkg/agent/loop_test.go: updated test helpers for new file layout. --- pkg/agent/agent_run.go | 543 +++++++++++++ pkg/agent/command_handler.go | 225 ++++++ pkg/agent/context.go | 5 +- pkg/agent/helpers.go | 126 +++ pkg/agent/loop.go | 1450 ++-------------------------------- pkg/agent/loop_test.go | 7 +- pkg/agent/message_router.go | 171 ++++ pkg/agent/summarizer.go | 492 ++++++++++++ pkg/agent/toolloop.go | 3 +- 9 files changed, 1637 insertions(+), 1385 deletions(-) create mode 100644 pkg/agent/agent_run.go create mode 100644 pkg/agent/command_handler.go create mode 100644 pkg/agent/helpers.go create mode 100644 pkg/agent/message_router.go create mode 100644 pkg/agent/summarizer.go diff --git a/pkg/agent/agent_run.go b/pkg/agent/agent_run.go new file mode 100644 index 000000000..40b6431a3 --- /dev/null +++ b/pkg/agent/agent_run.go @@ -0,0 +1,543 @@ +// assembledContext holds the pre-processed context produced by assembleContext, +// consumed by both the Generate and Stream code paths. +package agent + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + fantasy "charm.land/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg" + "github.com/ZanzyTHEbar/dragonscale/pkg/bus" + "github.com/ZanzyTHEbar/dragonscale/pkg/constants" + picofantasy "github.com/ZanzyTHEbar/dragonscale/pkg/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg/ids" + "github.com/ZanzyTHEbar/dragonscale/pkg/logger" + "github.com/ZanzyTHEbar/dragonscale/pkg/memory" + memsqlc "github.com/ZanzyTHEbar/dragonscale/pkg/memory/sqlc" + "github.com/ZanzyTHEbar/dragonscale/pkg/messages" + "github.com/ZanzyTHEbar/dragonscale/pkg/tools" + "github.com/ZanzyTHEbar/dragonscale/pkg/utils" +) + +type assembledContext struct { + systemPrompt string + userPrompt string + fantasyHistory []fantasy.Message + adaptedTools []fantasy.AgentTool + agent fantasy.Agent +} + +func (al *AgentLoop) prepareRuntimeState(ctx context.Context, sessionKey string) (ids.UUID, ids.UUID, error) { + if al.queries == nil || al.stateStore == nil || al.kvDelegate == nil { + return ids.UUID{}, ids.UUID{}, errors.New("runtime persistence dependencies are not initialized") + } + if strings.TrimSpace(sessionKey) == "" { + return ids.UUID{}, ids.UUID{}, errors.New("session key is required") + } + + var conversationID ids.UUID + if cached, ok := al.conversationIDs.Load(sessionKey); ok { + conversationID = cached.(ids.UUID) + } else { + al.conversationMu.Lock() + defer al.conversationMu.Unlock() + if cached, ok := al.conversationIDs.Load(sessionKey); ok { + conversationID = cached.(ids.UUID) + } else { + conversationID = ids.New() + title := sessionKey + if _, err := al.queries.CreateAgentConversation(ctx, memsqlc.CreateAgentConversationParams{ + ID: conversationID, + Title: &title, + }); err != nil { + return ids.UUID{}, ids.UUID{}, fmt.Errorf("create agent conversation: %w", err) + } + al.conversationIDs.Store(sessionKey, conversationID) + } + } + + run, err := al.stateStore.CreateRun(ctx, conversationID) + if err != nil { + return ids.UUID{}, ids.UUID{}, fmt.Errorf("create agent run: %w", err) + } + + return conversationID, run.ID, nil +} + +// assembleContext performs the shared pre-processing for every agent turn: +// record channel, update tool contexts, load memory blocks, build messages, +// DAG-compress history, split into system/history/user, adapt tools, create Fantasy agent. +func (al *AgentLoop) assembleContext(ctx context.Context, opts processOptions) (assembledContext, error) { + if err := al.recordChannelState(ctx, opts); err != nil { + logger.WarnCF("agent", "Failed to record last channel: %v", map[string]interface{}{"error": err.Error()}) + } + + logger.DebugCF("agent", "assembleContext: starting", + map[string]interface{}{ + "session_key": opts.SessionKey, + "channel": opts.Channel, + "sender_id": opts.SenderID, + }) + + al.refreshContextBlocks(ctx, opts) + history, summary := al.loadSessionState(ctx, opts) + builtMsgs := al.buildPromptMessages(opts, history, summary) + systemPrompt, historyMsgs, userPrompt := al.splitMessages(opts, builtMsgs) + + logger.DebugCF("agent", "assembleContext: history messages", + map[string]interface{}{ + "history": formatMessagesForLog(historyMsgs), + }) + + fantasyHistory := picofantasy.MessagesToFantasy(historyMsgs) + adaptedTools, prepareStep := al.prepareToolset(ctx, opts) + agent, err := al.createFantasyAgent(ctx, opts, systemPrompt, adaptedTools, prepareStep) + if err != nil { + return assembledContext{}, err + } + + logger.DebugCF("agent", "Fantasy agent created", + map[string]interface{}{ + "model": al.model, + "tools_count": len(adaptedTools), + "history_count": len(historyMsgs), + "max_iterations": al.maxIterations, + "memory_enabled": true, + }) + + return assembledContext{ + systemPrompt: systemPrompt, + userPrompt: userPrompt, + fantasyHistory: fantasyHistory, + adaptedTools: adaptedTools, + agent: agent, + }, nil +} + +func (al *AgentLoop) recordChannelState(ctx context.Context, opts processOptions) error { + if opts.Channel == "" || opts.ChatID == "" { + return nil + } + if constants.IsInternalChannel(opts.Channel) { + return nil + } + + channelKey := fmt.Sprintf("%s:%s", opts.Channel, opts.ChatID) + return al.RecordLastChannel(ctx, channelKey) +} + +func (al *AgentLoop) refreshContextBlocks(ctx context.Context, opts processOptions) { + al.updateToolContexts(opts.Channel, opts.ChatID) + + block := al.obsManager.LoadBlock(ctx, opts.SessionKey) + al.contextBuilder.SetObservationBlock(block) + + kb := tools.LoadKnowledgeBlock(ctx, al.memDelegate, opts.SessionKey) + al.contextBuilder.SetKnowledgeBlock(kb) + + if al.identitySync != nil { + _ = al.identitySync.CheckAndSync(ctx) + } +} + +func (al *AgentLoop) loadSessionState(ctx context.Context, opts processOptions) ([]messages.Message, string) { + var history []messages.Message + var summary string + if !opts.NoHistory { + history = al.sessions.GetHistory(opts.SessionKey) + summary = al.sessions.GetSummary(opts.SessionKey) + } + + return al.applyDAGCompression(ctx, opts.SessionKey, history), summary +} + +func (al *AgentLoop) buildPromptMessages(opts processOptions, history []messages.Message, summary string) []messages.Message { + builtMsgs := al.contextBuilder.BuildMessages(history, summary, opts.UserMessage, nil, opts.Channel, opts.ChatID) + al.sessions.AddMessage(opts.SessionKey, "user", opts.UserMessage) + return builtMsgs +} + +func (al *AgentLoop) splitMessages(opts processOptions, builtMsgs []messages.Message) (string, []messages.Message, string) { + systemPrompt := "" + var historyMsgs []messages.Message + userPrompt := opts.UserMessage + + if len(builtMsgs) > 0 && builtMsgs[0].Role == "system" { + systemPrompt = builtMsgs[0].Content + if len(builtMsgs) > 2 { + historyMsgs = builtMsgs[1 : len(builtMsgs)-1] + } + } + + return systemPrompt, historyMsgs, userPrompt +} + +func (al *AgentLoop) prepareToolset(ctx context.Context, opts processOptions) ([]fantasy.AgentTool, func(context.Context, fantasy.PrepareStepFunctionOptions) (context.Context, fantasy.PrepareStepResult, error)) { + adaptCfg := picofantasy.AdaptedToolsConfig{ + MemStore: al.memoryStore, + AgentID: pkg.NAME, + SessionKey: opts.SessionKey, + } + + adaptedTools := picofantasy.BuildAdaptedTools(al.tools, al.bus, opts.Channel, opts.ChatID, adaptCfg) + if al.toolResultSearch != nil { + adaptedTools = append(adaptedTools, al.toolResultSearch) + } + + promotedSet := make(map[string]bool) + for _, at := range adaptedTools { + promotedSet[at.Info().Name] = true + } + registry := al.tools + msgBus := al.bus + channel := opts.Channel + chatID := opts.ChatID + + prepareStep := func(ctx context.Context, psOpts fantasy.PrepareStepFunctionOptions) (context.Context, fantasy.PrepareStepResult, error) { + _ = psOpts + discovered := registry.DrainDiscovered() + if len(discovered) == 0 { + return ctx, fantasy.PrepareStepResult{}, nil + } + + newTools := make([]tools.Tool, 0, len(discovered)) + for _, t := range discovered { + if promotedSet[t.Name()] { + continue + } + newTools = append(newTools, t) + promotedSet[t.Name()] = true + } + + if len(newTools) == 0 { + return ctx, fantasy.PrepareStepResult{}, nil + } + + newAdapted := picofantasy.AdaptTools(newTools, msgBus, channel, chatID, adaptCfg) + expanded := append(adaptedTools, newAdapted...) + adaptedTools = expanded + + logger.InfoCF("agent", "Dynamic tool promotion via PrepareStep", + map[string]interface{}{ + "promoted": len(newTools), + "total_tools": len(expanded), + "names": toolNames(newTools), + }) + + return ctx, fantasy.PrepareStepResult{ + Tools: expanded, + }, nil + } + + return adaptedTools, prepareStep +} + +func (al *AgentLoop) createFantasyAgent(ctx context.Context, opts processOptions, systemPrompt string, adaptedTools []fantasy.AgentTool, prepareStep func(context.Context, fantasy.PrepareStepFunctionOptions) (context.Context, fantasy.PrepareStepResult, error)) (fantasy.Agent, error) { + conversationID, runID, err := al.prepareRuntimeState(ctx, opts.SessionKey) + if err != nil { + return nil, err + } + + baseRuntime := OffloadingToolRuntime{ + Base: fantasy.DAGToolRuntime{MaxConcurrency: defaultToolMaxConcurrency}, + KV: al.kvDelegate, + Queries: al.queries, + ConversationID: conversationID, + RunID: runID, + } + toolRuntime := SecureBusToolRuntime{ + Base: baseRuntime, + Bus: al.secureBus, + SessionKey: opts.SessionKey, + StateStore: al.stateStore, + RunID: runID, + } + + agentOpts := []fantasy.AgentOption{ + fantasy.WithTools(adaptedTools...), + fantasy.WithStopConditions(fantasy.StepCountIs(al.maxIterations)), + fantasy.WithPrepareStep(prepareStep), + fantasy.WithToolRuntime(toolRuntime), + } + if systemPrompt != "" { + agentOpts = append(agentOpts, fantasy.WithSystemPrompt(systemPrompt)) + } + + return fantasy.NewAgent(al.languageModel, agentOpts...), nil +} + +// postProcess handles the common finalization after Generate or Stream: +// extract final text, save session, summarize, observe, optionally send response. +func (al *AgentLoop) postProcess(ctx context.Context, opts processOptions, finalContent string, stepCount int) string { + al.sessions.Save(opts.SessionKey) + + if opts.EnableSummary { + al.maybeSummarize(ctx, opts.SessionKey, opts.Channel, opts.ChatID) + } + + tail := al.sessionsToMessagePairs(opts.SessionKey) + al.obsManager.MaybeObserveAsync(ctx, opts.SessionKey, tail) + + if opts.SendResponse { + al.bus.PublishOutbound(bus.OutboundMessage{ + Channel: opts.Channel, + ChatID: opts.ChatID, + Content: finalContent, + }) + } + + responsePreview := utils.Truncate(finalContent, 120) + logger.InfoCF("agent", fmt.Sprintf("Response: %s", responsePreview), + map[string]interface{}{ + "session_key": opts.SessionKey, + "steps": stepCount, + "final_length": len(finalContent), + }) + + return finalContent +} + +// resolveFinalContent normalizes the final assistant response from an agent run. +// Some providers return an empty final response even though an earlier step +// already produced text. In that case, recover the latest non-empty text from +// steps. If no text exists at all, return a deterministic error. +func (al *AgentLoop) resolveFinalContent(finalContent string, steps []fantasy.StepResult) (string, error) { + trimmed := strings.TrimSpace(finalContent) + if trimmed != "" { + return trimmed, nil + } + + for i := len(steps) - 1; i >= 0; i-- { + stepText := strings.TrimSpace(steps[i].Content.Text()) + if stepText != "" { + logger.WarnCF("agent", "Recovered empty final response from prior step text", + map[string]interface{}{ + "step_index": i, + }) + return stepText, nil + } + } + + type candidate struct { + text string + score int + } + candidates := make([]candidate, 0, 8) + for i := len(steps) - 1; i >= 0; i-- { + toolResults := steps[i].Content.ToolResults() + for j := len(toolResults) - 1; j >= 0; j-- { + tr := toolResults[j] + switch out := tr.Result.(type) { + case fantasy.ToolResultOutputContentText: + txt := strings.TrimSpace(out.Text) + if txt != "" { + score := 2 + if tr.ToolName == "tool_search" || strings.Contains(strings.ToLower(txt), "\"kind\":\"tool\"") { + score = 0 + } + if strings.Contains(strings.ToLower(txt), "tool not found") || + strings.Contains(strings.ToLower(txt), "path is required") { + score = -1 + } + candidates = append(candidates, candidate{text: txt, score: score}) + } + case fantasy.ToolResultOutputContentError: + if out.Error != nil { + txt := strings.TrimSpace(out.Error.Error()) + if txt != "" { + candidates = append(candidates, candidate{text: txt, score: -1}) + } + } + case fantasy.ToolResultOutputContentMedia: + txt := strings.TrimSpace(out.Text) + if txt != "" { + candidates = append(candidates, candidate{text: txt, score: 1}) + } + } + if len(candidates) >= 8 { + break + } + } + if len(candidates) >= 8 { + break + } + } + + bestText := "" + bestScore := -1000 + for _, c := range candidates { + if c.score > bestScore { + bestScore = c.score + bestText = c.text + } + } + + if bestText != "" && bestScore > 0 { + logger.WarnCF("agent", "Recovered empty final response from tool results", + map[string]interface{}{ + "candidates": len(candidates), + "score": bestScore, + }) + return bestText, nil + } + + toolCalls := 0 + for _, step := range steps { + toolCalls += len(step.Content.ToolCalls()) + } + + return "", fmt.Errorf("agent produced no final response text (steps=%d, tool_calls=%d)", len(steps), toolCalls) +} + +// runAgentLoop is the core message processing logic. +// It delegates to assembleContext for shared pre-processing, then branches on +// opts.Streaming to either Generate (synchronous) or Stream (real-time deltas). +func (al *AgentLoop) runAgentLoop(ctx context.Context, opts processOptions) (string, error) { + al.activeSessionKey.Store(opts.SessionKey) + + ac, err := al.assembleContext(ctx, opts) + if err != nil { + return "", err + } + + if opts.Streaming { + return al.runStreaming(ctx, opts, ac) + } + + result, err := ac.agent.Generate(ctx, fantasy.AgentCall{ + Prompt: ac.userPrompt, + Messages: ac.fantasyHistory, + }) + if err != nil { + logger.ErrorCF("agent", "Fantasy Generate failed", + map[string]interface{}{"error": err.Error()}) + return "", fmt.Errorf("agent Generate failed: %w", err) + } + + for _, step := range result.Steps { + stepMsgs := picofantasy.StepToMessages(step) + for _, m := range stepMsgs { + al.sessions.AddFullMessage(opts.SessionKey, m) + } + al.auditStep(ctx, step, opts.SessionKey) + } + + finalContent, err := al.resolveFinalContent(result.Response.Content.Text(), result.Steps) + if err != nil { + logger.ErrorCF("agent", "Agent finished without final response text", + map[string]interface{}{ + "error": err.Error(), + "steps": len(result.Steps), + }) + return "", err + } + return al.postProcess(ctx, opts, finalContent, len(result.Steps)), nil +} + +// runStreaming uses Fantasy's agent.Stream() to stream token deltas to the bus +// in real time, using the pre-assembled context from assembleContext. +func (al *AgentLoop) runStreaming(ctx context.Context, opts processOptions, ac assembledContext) (string, error) { + streamCall := fantasy.AgentStreamCall{ + Prompt: ac.userPrompt, + Messages: ac.fantasyHistory, + + OnTextDelta: func(id, text string) error { + if opts.Channel != "" && opts.ChatID != "" { + al.bus.PublishOutbound(bus.OutboundMessage{ + Channel: opts.Channel, + ChatID: opts.ChatID, + Content: text, + StreamDelta: true, + }) + } + return nil + }, + + OnStepFinish: func(step fantasy.StepResult) error { + stepMsgs := picofantasy.StepToMessages(step) + for _, m := range stepMsgs { + al.sessions.AddFullMessage(opts.SessionKey, m) + } + al.auditStep(ctx, step, opts.SessionKey) + return nil + }, + + OnToolCall: func(tc fantasy.ToolCallContent) error { + logger.DebugCF("agent", "Streaming tool call", + map[string]interface{}{ + "tool": tc.ToolName, + "id": tc.ToolCallID, + }) + return nil + }, + } + + result, err := ac.agent.Stream(ctx, streamCall) + if err != nil { + logger.ErrorCF("agent", "Fantasy Stream failed", + map[string]interface{}{"error": err.Error()}) + return "", fmt.Errorf("agent Stream failed: %w", err) + } + + finalContent, err := al.resolveFinalContent(result.Response.Content.Text(), result.Steps) + if err != nil { + logger.ErrorCF("agent", "Streaming agent finished without final response text", + map[string]interface{}{ + "error": err.Error(), + "steps": len(result.Steps), + }) + return "", err + } + return al.postProcess(ctx, opts, finalContent, len(result.Steps)), nil +} + +// runLLMIteration — DELETED. Replaced by Fantasy's internal agent loop. + +// auditStep logs tool calls from a Fantasy step result to the audit log. +func (al *AgentLoop) auditStep(ctx context.Context, step fantasy.StepResult, sessionKey string) { + toolCalls := step.Content.ToolCalls() + if len(toolCalls) == 0 { + return + } + + for _, tc := range toolCalls { + entry := &memory.AuditEntry{ + ID: ids.New(), + AgentID: pkg.NAME, + SessionKey: sessionKey, + Action: "tool_call", + Target: tc.ToolName, + Input: tc.Input, + } + aCtx, cancel := context.WithTimeout(ctx, time.Second) + if err := al.memDelegate.InsertAuditEntry(aCtx, entry); err != nil { + logger.WarnCF("agent", "Failed to log audit entry", + map[string]interface{}{"tool": tc.ToolName, "error": err.Error()}) + } + cancel() + } +} + +// updateToolContexts updates the context for tools that need channel/chatID info. +func (al *AgentLoop) updateToolContexts(channel, chatID string) { + // Use ContextualTool interface instead of type assertions + if tool, ok := al.tools.Get("message"); ok { + if mt, ok := tool.(tools.ContextualTool); ok { + mt.SetContext(channel, chatID) + } + } + if tool, ok := al.tools.Get("spawn"); ok { + if st, ok := tool.(tools.ContextualTool); ok { + st.SetContext(channel, chatID) + } + } + if tool, ok := al.tools.Get("subagent"); ok { + if st, ok := tool.(tools.ContextualTool); ok { + st.SetContext(channel, chatID) + } + } +} diff --git a/pkg/agent/command_handler.go b/pkg/agent/command_handler.go new file mode 100644 index 000000000..39718aae2 --- /dev/null +++ b/pkg/agent/command_handler.go @@ -0,0 +1,225 @@ +package agent + +import ( + "context" + "fmt" + "strings" + + "github.com/ZanzyTHEbar/dragonscale/pkg/bus" + "github.com/ZanzyTHEbar/dragonscale/pkg/config" +) + +// SlashCommand defines an in-agent command. +type SlashCommand struct { + Name string + Description string + Usage string + Handler func(al *AgentLoop, ctx context.Context, msg bus.InboundMessage, args []string) string +} + +// listConfiguredModels returns a human-readable summary of which providers +// have API credentials configured, and the current default model. +func listConfiguredModels(cfg *config.Config) string { + if cfg == nil { + return "No configuration available." + } + + current := fmt.Sprintf("Current model: %s", cfg.Agents.Defaults.Model) + if cfg.Agents.Defaults.Provider != "" { + current += fmt.Sprintf(" (provider: %s)", cfg.Agents.Defaults.Provider) + } + + configured := cfg.Providers.ConfiguredNames() + if len(configured) == 0 { + return current + "\nNo providers configured — set API keys in config.json or environment variables." + } + + return current + "\nConfigured providers: " + strings.Join(configured, ", ") +} + +func defaultSlashCommands() []SlashCommand { + return []SlashCommand{ + { + Name: "/show", + Description: "Display current settings.", + Usage: "/show [model|channel]", + Handler: func(al *AgentLoop, _ context.Context, msg bus.InboundMessage, args []string) string { + if len(args) < 1 { + return "Usage: /show [model|channel]" + } + switch args[0] { + case "model": + return fmt.Sprintf("Current model: %s", al.model) + case "channel": + parts := []string{fmt.Sprintf("Current channel: %s", msg.Channel)} + if override, ok := al.getOutputTarget(); ok { + parts = append(parts, fmt.Sprintf("Output redirect: %s:%s", override.Channel, override.ChatID)) + } else { + parts = append(parts, "Output redirect: default") + } + return strings.Join(parts, "\n") + default: + return fmt.Sprintf("Unknown show target: %s", args[0]) + } + }, + }, + { + Name: "/list", + Description: "List configured providers or enabled channels.", + Usage: "/list [models|channels]", + Handler: func(al *AgentLoop, _ context.Context, _ bus.InboundMessage, args []string) string { + if len(args) < 1 { + return "Usage: /list [models|channels]" + } + switch args[0] { + case "models": + return listConfiguredModels(al.cfg) + case "channels": + if al.channelManager == nil { + return "Channel manager not initialized" + } + channels := al.channelManager.GetEnabledChannels() + if len(channels) == 0 { + return "No channels enabled" + } + return fmt.Sprintf("Enabled channels: %s", strings.Join(channels, ", ")) + default: + return fmt.Sprintf("Unknown list target: %s", args[0]) + } + }, + }, + { + Name: "/switch", + Description: "Switch model or target channel alias.", + Usage: "/switch [model|channel] to ", + Handler: func(al *AgentLoop, ctx context.Context, _ bus.InboundMessage, args []string) string { + if len(args) < 3 || args[1] != "to" { + return "Usage: /switch [model|channel] to " + } + target := args[0] + value := args[2] + switch target { + case "model": + oldModel := al.model + al.model = value + return fmt.Sprintf("Switched model from %s to %s", oldModel, value) + case "channel": + return al.handleSwitchChannel(ctx, value) + default: + return fmt.Sprintf("Unknown switch target: %s", target) + } + }, + }, + { + Name: "/help", + Description: "List available slash commands.", + Usage: "/help [command]", + Handler: func(al *AgentLoop, _ context.Context, _ bus.InboundMessage, args []string) string { + if len(args) == 0 { + lines := make([]string, 0, len(al.commandRegistry)+1) + lines = append(lines, "Available slash commands:") + for _, cmd := range al.commandRegistry { + lines = append(lines, fmt.Sprintf(" %s - %s (%s)", cmd.Usage, cmd.Description, cmd.Name)) + } + return strings.Join(lines, "\n") + } + target := args[0] + if !strings.HasPrefix(target, "/") { + target = "/" + target + } + for _, cmd := range al.commandRegistry { + if cmd.Name == target { + return fmt.Sprintf("%s - %s", cmd.Usage, cmd.Description) + } + } + return fmt.Sprintf("Unknown command: %s", target) + }, + }, + } +} + +func (al *AgentLoop) parseSwitchChannelTarget(target string) (string, string) { + channel := strings.TrimSpace(target) + chatID := "" + if idx := strings.Index(channel, ":"); idx >= 0 { + chatID = strings.TrimSpace(channel[idx+1:]) + channel = strings.TrimSpace(channel[:idx]) + } + return channel, chatID +} + +func (al *AgentLoop) handleSwitchChannel(ctx context.Context, target string) string { + channel, chatID := al.parseSwitchChannelTarget(target) + if channel == "" { + return "Usage: /switch channel to [:chat_id]" + } + + if channel == "cli" { + al.outputOverride.Store(outputTarget{}) + if al.state != nil { + _ = al.state.SetLastChannel(ctx, "cli") + _ = al.state.SetLastChatID(ctx, "") + } + return "Cleared output channel override to CLI defaults" + } + + if al.channelManager == nil { + return "Channel manager not initialized" + } + if _, exists := al.channelManager.GetChannel(channel); !exists { + return fmt.Sprintf("Channel '%s' not found or not enabled", channel) + } + + if chatID == "" { + if al.state == nil { + return "No chat ID available for channel target. Use /switch channel to " + } + chatID = al.state.GetLastChatID() + } + if chatID == "" { + return "No chat ID available for channel target. Use /switch channel to " + } + + al.outputOverride.Store(outputTarget{ + Channel: channel, + ChatID: chatID, + }) + if al.state != nil { + if err := al.state.SetLastChannel(ctx, channel); err != nil { + return fmt.Sprintf("Output redirection set to %s:%s, but failed to persist channel: %v", channel, chatID, err) + } + if err := al.state.SetLastChatID(ctx, chatID); err != nil { + return fmt.Sprintf("Output redirection set to %s:%s, but failed to persist chat id: %v", channel, chatID, err) + } + } + + return fmt.Sprintf("Switched target channel to %s:%s", channel, chatID) +} + +// handleCommand processes slash commands and returns (response, handled). +func (al *AgentLoop) handleCommand(ctx context.Context, msg bus.InboundMessage) (string, bool) { + content := strings.TrimSpace(msg.Content) + if !strings.HasPrefix(content, "/") { + return "", false + } + + parts := strings.Fields(content) + if len(parts) == 0 { + return "", false + } + + cmd := parts[0] + args := parts[1:] + + for _, command := range al.commandRegistry { + if command.Name != cmd { + continue + } + if command.Handler == nil { + return "Command is not implemented", true + } + return command.Handler(al, ctx, msg, args), true + } + + return "", false +} diff --git a/pkg/agent/context.go b/pkg/agent/context.go index 8e8537db7..9f2890f7e 100644 --- a/pkg/agent/context.go +++ b/pkg/agent/context.go @@ -9,6 +9,7 @@ import ( "strings" "time" + "github.com/ZanzyTHEbar/dragonscale/pkg" "github.com/ZanzyTHEbar/dragonscale/pkg/config" "github.com/ZanzyTHEbar/dragonscale/pkg/logger" "github.com/ZanzyTHEbar/dragonscale/pkg/memory" @@ -273,7 +274,7 @@ func (cb *ContextBuilder) LoadBootstrapFiles() string { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() - docs, err := cb.delegate.ListDocumentsByCategory(ctx, "dragonscale", "bootstrap") + docs, err := cb.delegate.ListDocumentsByCategory(ctx, pkg.NAME, "bootstrap") if err != nil || len(docs) == 0 { return "" } @@ -299,7 +300,7 @@ func (cb *ContextBuilder) buildWorkingContextSection() string { var parts []string // Inject working context (hot tier) - wc, err := cb.memoryStore.GetWorkingContext(ctx, "dragonscale", "default") + wc, err := cb.memoryStore.GetWorkingContext(ctx, pkg.NAME, "default") if err == nil && wc != "" { parts = append(parts, "## Working Context\n\n"+wc) } diff --git a/pkg/agent/helpers.go b/pkg/agent/helpers.go new file mode 100644 index 000000000..20208e76e --- /dev/null +++ b/pkg/agent/helpers.go @@ -0,0 +1,126 @@ +package agent + +import ( + "fmt" + "os" + "path/filepath" + + "github.com/ZanzyTHEbar/dragonscale/pkg/bus" + "github.com/ZanzyTHEbar/dragonscale/pkg/config" + "github.com/ZanzyTHEbar/dragonscale/pkg/logger" + "github.com/ZanzyTHEbar/dragonscale/pkg/messages" + "github.com/ZanzyTHEbar/dragonscale/pkg/security" + "github.com/ZanzyTHEbar/dragonscale/pkg/tools" + "github.com/ZanzyTHEbar/dragonscale/pkg/utils" +) + +func initSecretStore() (*security.SecretStore, error) { + cfgDir, err := config.ConfigDir() + if err != nil { + return nil, fmt.Errorf("resolve config dir: %w", err) + } + + secretsPath := filepath.Join(cfgDir, "secrets.json") + keyring := security.NewEnvKeyring(security.MasterKeyEnvVar) + ss, err := security.NewSecretStore(secretsPath, keyring) + if err != nil { + return nil, fmt.Errorf("initialize secret store: %w", err) + } + + if os.Getenv(security.MasterKeyEnvVar) == "" { + logger.WarnCF("security", "master key env var is not set; secret injection requiring stored secrets will fail", + map[string]interface{}{"env_var": security.MasterKeyEnvVar}) + } + return ss, nil +} + +// createToolRegistry creates a tool registry with common tools. +// createToolRegistry builds the base tool set (filesystem, shell, web, etc.). +// Parent and subagent registries start from the same base; memory/search/skill +// tools are registered separately on each so they have isolated discovery state. +func createToolRegistry(workspace string, restrict bool, cfg *config.Config, msgBus *bus.MessageBus) *tools.ToolRegistry { + registry := tools.NewToolRegistry() + + // File system tools + registry.Register(tools.NewReadFileTool(workspace, restrict)) + registry.Register(tools.NewWriteFileTool(workspace, restrict)) + registry.Register(tools.NewListDirTool(workspace, restrict)) + registry.Register(tools.NewEditFileTool(workspace, restrict)) + registry.Register(tools.NewAppendFileTool(workspace, restrict)) + + // Shell execution + registry.Register(tools.NewExecTool(workspace, restrict)) + + if searchTool := tools.NewWebSearchTool(tools.WebSearchToolOptions{ + BraveAPIKey: cfg.Tools.Web.Brave.APIKey, + BraveMaxResults: cfg.Tools.Web.Brave.MaxResults, + BraveEnabled: cfg.Tools.Web.Brave.Enabled, + DuckDuckGoMaxResults: cfg.Tools.Web.DuckDuckGo.MaxResults, + DuckDuckGoEnabled: cfg.Tools.Web.DuckDuckGo.Enabled, + PerplexityAPIKey: cfg.Tools.Web.Perplexity.APIKey, + PerplexityMaxResults: cfg.Tools.Web.Perplexity.MaxResults, + PerplexityEnabled: cfg.Tools.Web.Perplexity.Enabled, + }); searchTool != nil { + registry.Register(searchTool) + } + registry.Register(tools.NewWebFetchTool(50000)) + + // Hardware tools (I2C, SPI) - Linux only, returns error on other platforms + registry.Register(tools.NewI2CTool()) + registry.Register(tools.NewSPITool()) + + // Message tool - available to both agent and subagent + // Subagent uses it to communicate directly with user + messageTool := tools.NewMessageTool() + messageTool.SetSendCallback(func(channel, chatID, content string) error { + msgBus.PublishOutbound(bus.OutboundMessage{ + Channel: channel, + ChatID: chatID, + Content: content, + }) + return nil + }) + registry.Register(messageTool) + + return registry +} + +// formatMessagesForLog formats messages for logging +func formatMessagesForLog(msgs []messages.Message) string { + if len(msgs) == 0 { + return "[]" + } + + var result string + result += "[\n" + for i, msg := range msgs { + result += fmt.Sprintf(" [%d] Role: %s\n", i, msg.Role) + if len(msg.ToolCalls) > 0 { + result += " ToolCalls:\n" + for _, tc := range msg.ToolCalls { + result += fmt.Sprintf(" - ID: %s, Type: %s, Name: %s\n", tc.ID, tc.Type, tc.Name) + if tc.Function != nil { + result += fmt.Sprintf(" Arguments: %s\n", utils.Truncate(tc.Function.Arguments, 200)) + } + } + } + if msg.Content != "" { + content := utils.Truncate(msg.Content, 200) + result += fmt.Sprintf(" Content: %s\n", content) + } + if msg.ToolCallID != "" { + result += fmt.Sprintf(" ToolCallID: %s\n", msg.ToolCallID) + } + result += "\n" + } + result += "]" + return result +} + +func toolNames(tt []tools.Tool) []string { + names := make([]string, len(tt)) + for i, t := range tt { + names[i] = t.Name() + } + return names +} diff --git a/pkg/agent/loop.go b/pkg/agent/loop.go index 1f7271cba..f3d82cff4 100644 --- a/pkg/agent/loop.go +++ b/pkg/agent/loop.go @@ -4,39 +4,23 @@ // // Copyright (c) 2026 DragonScale contributors -// FIXME: This file is a mess, we need to clean it up and make it more readable and maintainable -// Break this into modules with single responsibility and composability in mind -// Leverage Ports and Adapters pattern to achieve this (where boundaries are defined sensibly) -// - AgentLoop: Main agent loop and orchestrator -// - ContextBuilder: Builds the context for the agent -// - Tools: Tool registry and management -// - Memory: Memory store and management -// - State: State management -// - Session: Session management -// - Identity: Identity management -// - SecureBus: SecureBus management - package agent import ( "context" - "encoding/json" - "errors" "fmt" "os" "path/filepath" - "strings" "sync" "sync/atomic" "time" fantasy "charm.land/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg" "github.com/ZanzyTHEbar/dragonscale/pkg/bus" "github.com/ZanzyTHEbar/dragonscale/pkg/channels" "github.com/ZanzyTHEbar/dragonscale/pkg/config" "github.com/ZanzyTHEbar/dragonscale/pkg/constants" - picofantasy "github.com/ZanzyTHEbar/dragonscale/pkg/fantasy" - "github.com/ZanzyTHEbar/dragonscale/pkg/ids" "github.com/ZanzyTHEbar/dragonscale/pkg/logger" "github.com/ZanzyTHEbar/dragonscale/pkg/memory" "github.com/ZanzyTHEbar/dragonscale/pkg/memory/dag" @@ -44,14 +28,12 @@ import ( "github.com/ZanzyTHEbar/dragonscale/pkg/memory/observation" memsqlc "github.com/ZanzyTHEbar/dragonscale/pkg/memory/sqlc" memstore "github.com/ZanzyTHEbar/dragonscale/pkg/memory/store" - "github.com/ZanzyTHEbar/dragonscale/pkg/messages" "github.com/ZanzyTHEbar/dragonscale/pkg/security" "github.com/ZanzyTHEbar/dragonscale/pkg/security/securebus" "github.com/ZanzyTHEbar/dragonscale/pkg/session" "github.com/ZanzyTHEbar/dragonscale/pkg/state" picosync "github.com/ZanzyTHEbar/dragonscale/pkg/sync" "github.com/ZanzyTHEbar/dragonscale/pkg/tools" - "github.com/ZanzyTHEbar/dragonscale/pkg/utils" ) type AgentLoop struct { @@ -72,18 +54,25 @@ type AgentLoop struct { queries *memsqlc.Queries // SQL query surface for runtime persistence kvDelegate KVDelegate // KV adapter for offloaded tool results stateStore *StateStore // Agent run state persistence - conversationIDs sync.Map // map[sessionKey]ids.UUID + conversationIDs sync.Map // Owner: agent_run.go — wrote by prepareRuntimeState, read in prepareRuntimeState/load path conversationMu sync.Mutex // serializes conversation creation path identitySync *picosync.IdentitySync // File→DB sync for identity docs (nil if memory disabled) - activeSessionKey atomic.Value // Current session key for tool access - running atomic.Bool - summarizing sync.Map // Tracks which sessions are currently being summarized - summarizeFailures sync.Map // Tracks consecutive summarization failures per session (string -> int) - cfg *config.Config // Stored for subagent factory access + activeSessionKey atomic.Value // Owner: agent_run.go — written in runAgentLoop, read by router/toolloop for context routing + running atomic.Bool // Owner: loop.go — lifecycle gate controlled by Run/Stop only + summarizing sync.Map // Owner: summarizer.go — intended for async summarization lockout, currently gated by TODO path + summarizeFailures sync.Map // Owner: summarizer.go — write/read in forceCompression + summarizeSession error paths + cfg *config.Config // Stored for subagent factory access channelManager *channels.Manager + commandRegistry []SlashCommand + outputOverride atomic.Value // Owner: command_handler.go — CLI output redirection target for internal messages toolResultSearch fantasy.AgentTool } +type outputTarget struct { + Channel string + ChatID string +} + // processOptions configures how a message is processed type processOptions struct { SessionKey string // Session identifier for history/context @@ -97,78 +86,17 @@ type processOptions struct { Streaming bool // If true, stream token deltas to bus via OnTextDelta } -func initSecretStore() (*security.SecretStore, error) { - cfgDir, err := config.ConfigDir() - if err != nil { - return nil, fmt.Errorf("resolve config dir: %w", err) - } +// Option configures AgentLoop creation. +type Option func(*AgentLoop) - secretsPath := filepath.Join(cfgDir, "secrets.json") - keyring := security.NewEnvKeyring(security.MasterKeyEnvVar) - ss, err := security.NewSecretStore(secretsPath, keyring) - if err != nil { - return nil, fmt.Errorf("initialize secret store: %w", err) +// WithChannelManager configures command routing for slash commands. +func WithChannelManager(cm *channels.Manager) Option { + return func(al *AgentLoop) { + al.channelManager = cm } - - if os.Getenv(security.MasterKeyEnvVar) == "" { - logger.WarnCF("security", "master key env var is not set; secret injection requiring stored secrets will fail", - map[string]interface{}{"env_var": security.MasterKeyEnvVar}) - } - return ss, nil } -// createToolRegistry creates a tool registry with common tools. -// createToolRegistry builds the base tool set (filesystem, shell, web, etc.). -// Parent and subagent registries start from the same base; memory/search/skill -// tools are registered separately on each so they have isolated discovery state. -func createToolRegistry(workspace string, restrict bool, cfg *config.Config, msgBus *bus.MessageBus) *tools.ToolRegistry { - registry := tools.NewToolRegistry() - - // File system tools - registry.Register(tools.NewReadFileTool(workspace, restrict)) - registry.Register(tools.NewWriteFileTool(workspace, restrict)) - registry.Register(tools.NewListDirTool(workspace, restrict)) - registry.Register(tools.NewEditFileTool(workspace, restrict)) - registry.Register(tools.NewAppendFileTool(workspace, restrict)) - - // Shell execution - registry.Register(tools.NewExecTool(workspace, restrict)) - - if searchTool := tools.NewWebSearchTool(tools.WebSearchToolOptions{ - BraveAPIKey: cfg.Tools.Web.Brave.APIKey, - BraveMaxResults: cfg.Tools.Web.Brave.MaxResults, - BraveEnabled: cfg.Tools.Web.Brave.Enabled, - DuckDuckGoMaxResults: cfg.Tools.Web.DuckDuckGo.MaxResults, - DuckDuckGoEnabled: cfg.Tools.Web.DuckDuckGo.Enabled, - PerplexityAPIKey: cfg.Tools.Web.Perplexity.APIKey, - PerplexityMaxResults: cfg.Tools.Web.Perplexity.MaxResults, - PerplexityEnabled: cfg.Tools.Web.Perplexity.Enabled, - }); searchTool != nil { - registry.Register(searchTool) - } - registry.Register(tools.NewWebFetchTool(50000)) - - // Hardware tools (I2C, SPI) - Linux only, returns error on other platforms - registry.Register(tools.NewI2CTool()) - registry.Register(tools.NewSPITool()) - - // Message tool - available to both agent and subagent - // Subagent uses it to communicate directly with user - messageTool := tools.NewMessageTool() - messageTool.SetSendCallback(func(channel, chatID, content string) error { - msgBus.PublishOutbound(bus.OutboundMessage{ - Channel: channel, - ChatID: chatID, - Content: content, - }) - return nil - }) - registry.Register(messageTool) - - return registry -} - -func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBus, model fantasy.LanguageModel) (*AgentLoop, error) { +func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBus, model fantasy.LanguageModel, opts ...Option) (*AgentLoop, error) { sandbox := cfg.SandboxPath() os.MkdirAll(sandbox, 0755) @@ -215,7 +143,7 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu del.Close() return nil, fmt.Errorf("memory delegate queries are not initialized") } - kv := NewDelegateKV(memDelegate, "dragonscale") + kv := NewDelegateKV(memDelegate, pkg.NAME) stateStore := NewStateStore(queries) offloadThreshold := cfg.Memory.OffloadThresholdTokens @@ -239,14 +167,14 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu }) contextBuilder.SetMemoryStore(ms) - memTool := NewMemGPTTool(ms, "dragonscale", "default") + memTool := NewMemGPTTool(ms, pkg.NAME, "default") toolsRegistry.Register(memTool) - toolsRegistry.Register(tools.NewObligationTool(memDelegate, "dragonscale")) + toolsRegistry.Register(tools.NewObligationTool(memDelegate, pkg.NAME)) - toolsRegistry.Register(tools.NewKeywordSearchTool(ms, "dragonscale")) - toolsRegistry.Register(tools.NewSemanticSearchTool(ms, "dragonscale")) - toolsRegistry.Register(tools.NewChunkReadTool(ms, "dragonscale")) - mapRuntime := tools.NewMapRuntime(queries, "dragonscale", model, cfg.Agents.Defaults.Model, subagentManager) + toolsRegistry.Register(tools.NewKeywordSearchTool(ms, pkg.NAME)) + toolsRegistry.Register(tools.NewSemanticSearchTool(ms, pkg.NAME)) + toolsRegistry.Register(tools.NewChunkReadTool(ms, pkg.NAME)) + mapRuntime := tools.NewMapRuntime(queries, pkg.NAME, model, cfg.Agents.Defaults.Model, subagentManager) agenticMapTool.SetRuntime(mapRuntime) toolsRegistry.Register(agenticMapTool) llmMapTool := tools.NewLLMMapTool(model, cfg.Agents.Defaults.Model) @@ -257,11 +185,11 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu // Register memory/search/skill tools on subagent registry so spawned // agents can search knowledge, offload results, and use skills. - subagentTools.Register(NewMemGPTTool(ms, "dragonscale", "default")) - subagentTools.Register(tools.NewObligationTool(memDelegate, "dragonscale")) - subagentTools.Register(tools.NewKeywordSearchTool(ms, "dragonscale")) - subagentTools.Register(tools.NewSemanticSearchTool(ms, "dragonscale")) - subagentTools.Register(tools.NewChunkReadTool(ms, "dragonscale")) + subagentTools.Register(NewMemGPTTool(ms, pkg.NAME, "default")) + subagentTools.Register(tools.NewObligationTool(memDelegate, pkg.NAME)) + subagentTools.Register(tools.NewKeywordSearchTool(ms, pkg.NAME)) + subagentTools.Register(tools.NewSemanticSearchTool(ms, pkg.NAME)) + subagentTools.Register(tools.NewChunkReadTool(ms, pkg.NAME)) subagentTools.Register(tools.NewSkillSearchTool(sl)) subagentTools.Register(tools.NewSkillReadTool(sl)) subagentTools.RegisterMetaTools() @@ -271,19 +199,19 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu contextBuilder.SetDelegate(del) - if migErr := memory.MigrateState(ctx, workspace, del, "dragonscale"); migErr != nil { + if migErr := memory.MigrateState(ctx, workspace, del, pkg.NAME); migErr != nil { logger.WarnCF("agent", "State KV migration failed (non-fatal)", map[string]interface{}{"error": migErr.Error()}) } - if migErr := memory.MigrateDocuments(ctx, workspace, del, "dragonscale"); migErr != nil { + if migErr := memory.MigrateDocuments(ctx, workspace, del, pkg.NAME); migErr != nil { logger.WarnCF("agent", "Document migration failed (non-fatal)", map[string]interface{}{"error": migErr.Error()}) } - if migErr := memory.MigrateLongTermMemory(ctx, workspace, del, "dragonscale"); migErr != nil { + if migErr := memory.MigrateLongTermMemory(ctx, workspace, del, pkg.NAME); migErr != nil { logger.WarnCF("agent", "Long-term memory migration failed (non-fatal)", map[string]interface{}{"error": migErr.Error()}) } - if migErr := memory.MigrateDailyNotes(ctx, workspace, del, "dragonscale"); migErr != nil { + if migErr := memory.MigrateDailyNotes(ctx, workspace, del, pkg.NAME); migErr != nil { logger.WarnCF("agent", "Daily notes migration failed (non-fatal)", map[string]interface{}{"error": migErr.Error()}) } @@ -295,7 +223,7 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu logger.WarnCF("agent", "Could not resolve identity dir, identity sync disabled", map[string]interface{}{"error": idErr.Error()}) } else { - idSync = picosync.New(identityDir, "dragonscale", memDelegate) + idSync = picosync.New(identityDir, pkg.NAME, memDelegate) if syncErr := idSync.SyncAll(ctx); syncErr != nil { logger.WarnCF("agent", "Initial identity sync failed (non-fatal)", map[string]interface{}{"error": syncErr.Error()}) @@ -313,12 +241,12 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu // Session manager (always delegate-backed) sessionsDir := filepath.Join(workspace, "sessions") - sessionsManager := session.NewSessionManager(sessionsDir, session.WithSessionDelegate(memDelegate, "dragonscale")) + sessionsManager := session.NewSessionManager(sessionsDir, session.WithSessionDelegate(memDelegate, pkg.NAME)) // One-shot DAG backfill for pre-existing session histories. backfillCtx, cancelBackfill := context.WithTimeout(ctx, 10*time.Second) defer cancelBackfill() - if status, err := dag.BackfillMissingSessionDAGs(backfillCtx, memDelegate, queries, "dragonscale", dag.DefaultBackfillOptions()); err != nil { + if status, err := dag.BackfillMissingSessionDAGs(backfillCtx, memDelegate, queries, pkg.NAME, dag.DefaultBackfillOptions()); err != nil { logger.WarnCF("agent", "DAG backfill failed (non-fatal)", map[string]interface{}{"error": err.Error()}) } else if status != nil { @@ -372,7 +300,7 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu } return resp.Content.Text(), nil } - obsManager := observation.NewManager(memDelegate, "dragonscale", callModelFn, observation.DefaultManagerConfig()) + obsManager := observation.NewManager(memDelegate, pkg.NAME, callModelFn, observation.DefaultManagerConfig()) al := &AgentLoop{ bus: msgBus, @@ -394,9 +322,16 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu toolResultSearch: NewToolResultSearchTool(queries, kv), identitySync: idSync, summarizing: sync.Map{}, + commandRegistry: defaultSlashCommands(), cfg: cfg, } + for _, apply := range opts { + if apply != nil { + apply(al) + } + } + // Focus tools (start_focus / complete_focus) sessionKeyFn := func() string { if v := al.activeSessionKey.Load(); v != nil { @@ -412,7 +347,7 @@ func NewAgentLoop(ctx context.Context, cfg *config.Config, msgBus *bus.MessageBu Queries: queries, Lister: del, Delegate: del, - AgentID: "dragonscale", + AgentID: pkg.NAME, SessionFn: sessionKeyFn, } toolsRegistry.Register(tools.NewDagExpandTool(dagDeps)) @@ -459,11 +394,20 @@ func (al *AgentLoop) Run(ctx context.Context) error { } if !alreadySent { - al.bus.PublishOutbound(bus.OutboundMessage{ + outMsg := bus.OutboundMessage{ Channel: msg.Channel, ChatID: msg.ChatID, Content: response, - }) + } + if constants.IsInternalChannel(msg.Channel) { + if override, ok := al.getOutputTarget(); ok { + outMsg.Channel = override.Channel + if override.ChatID != "" { + outMsg.ChatID = override.ChatID + } + } + } + al.bus.PublishOutbound(outMsg) } } } @@ -522,755 +466,16 @@ func (al *AgentLoop) SetupSecureBus(ss *security.SecretStore, cfg securebus.BusC return b } -// RecordLastChannel records the last active channel for this workspace. -// This uses the atomic state save mechanism to prevent data loss on crash. -func (al *AgentLoop) RecordLastChannel(ctx context.Context, channel string) error { - return al.state.SetLastChannel(ctx, channel) -} - -// RecordLastChatID records the last active chat ID for this workspace. -// This uses the atomic state save mechanism to prevent data loss on crash. -func (al *AgentLoop) RecordLastChatID(ctx context.Context, chatID string) error { - return al.state.SetLastChatID(ctx, chatID) -} - -func (al *AgentLoop) ProcessDirect(ctx context.Context, content, sessionKey string) (string, error) { - return al.ProcessDirectWithChannel(ctx, content, sessionKey, "cli", "direct") -} - -func (al *AgentLoop) ProcessDirectWithChannel(ctx context.Context, content, sessionKey, channel, chatID string) (string, error) { - msg := bus.InboundMessage{ - Channel: channel, - SenderID: "cron", - ChatID: chatID, - Content: content, - SessionKey: sessionKey, +func (al *AgentLoop) getOutputTarget() (outputTarget, bool) { + raw := al.outputOverride.Load() + if raw == nil { + return outputTarget{}, false } - - return al.processMessage(ctx, msg) -} - -// ProcessDirectStreaming processes a message with streaming token delivery. -// Text deltas are published to the bus as StreamDelta messages in real time. -func (al *AgentLoop) ProcessDirectStreaming(ctx context.Context, content, sessionKey, channel, chatID string) (string, error) { - msg := bus.InboundMessage{ - Channel: channel, - SenderID: "user", - ChatID: chatID, - Content: content, - SessionKey: sessionKey, - } - - return al.runAgentLoop(ctx, processOptions{ - SessionKey: msg.SessionKey, - Channel: msg.Channel, - ChatID: msg.ChatID, - SenderID: msg.SenderID, - UserMessage: msg.Content, - EnableSummary: true, - SendResponse: false, - Streaming: true, - }) -} - -// ProcessHeartbeat processes a heartbeat request without session history. -// Each heartbeat is independent and doesn't accumulate context. -// It injects the active session's summary so the agent has awareness of -// recent user conversation context. -func (al *AgentLoop) ProcessHeartbeat(ctx context.Context, content, channel, chatID string) (string, error) { - if v := al.activeSessionKey.Load(); v != nil { - if key, ok := v.(string); ok && key != "" { - if summary := al.sessions.GetSummary(key); summary != "" { - content = content + "\n\n## Recent User Context\n" + summary - } - } - } - return al.runAgentLoop(ctx, processOptions{ - SessionKey: "heartbeat", - Channel: channel, - ChatID: chatID, - UserMessage: content, - EnableSummary: false, - SendResponse: false, - NoHistory: true, - }) -} - -func (al *AgentLoop) processMessage(ctx context.Context, msg bus.InboundMessage) (string, error) { - // Add message preview to log (show full content for error messages) - var logContent string - if strings.Contains(msg.Content, "Error:") || strings.Contains(msg.Content, "error") { - logContent = msg.Content // Full content for errors - } else { - logContent = utils.Truncate(msg.Content, 80) - } - logger.InfoCF("agent", fmt.Sprintf("Processing message from %s:%s: %s", msg.Channel, msg.SenderID, logContent), - map[string]interface{}{ - "channel": msg.Channel, - "chat_id": msg.ChatID, - "sender_id": msg.SenderID, - "session_key": msg.SessionKey, - }) - - // Route system messages to processSystemMessage - if msg.Channel == "system" { - return al.processSystemMessage(ctx, msg) - } - - // Check for commands - if response, handled := al.handleCommand(ctx, msg); handled { - return response, nil - } - - // Process as user message - return al.runAgentLoop(ctx, processOptions{ - SessionKey: msg.SessionKey, - Channel: msg.Channel, - ChatID: msg.ChatID, - UserMessage: msg.Content, - EnableSummary: true, - SendResponse: false, - }) -} - -func (al *AgentLoop) processSystemMessage(_ context.Context, msg bus.InboundMessage) (string, error) { - // Verify this is a system message - if msg.Channel != "system" { - return "", fmt.Errorf("processSystemMessage called with non-system message channel: %s", msg.Channel) - } - - logger.InfoCF("agent", "Processing system message", - map[string]interface{}{ - "sender_id": msg.SenderID, - "chat_id": msg.ChatID, - }) - - // Parse origin channel from chat_id (format: "channel:chat_id") - var originChannel string - if idx := strings.Index(msg.ChatID, ":"); idx > 0 { - originChannel = msg.ChatID[:idx] - } else { - // Fallback - originChannel = "cli" - } - - // Extract subagent result from message content - // Format: "Task 'label' completed.\n\nResult:\n" - content := msg.Content - if idx := strings.Index(content, "Result:\n"); idx >= 0 { - content = content[idx+8:] // Extract just the result part - } - - // Skip internal channels - only log, don't send to user - if constants.IsInternalChannel(originChannel) { - logger.InfoCF("agent", "Subagent completed (internal channel)", - map[string]interface{}{ - "sender_id": msg.SenderID, - "content_len": len(content), - "channel": originChannel, - }) - return "", nil - } - - // Agent acts as dispatcher only - subagent handles user interaction via message tool - // Don't forward result here, subagent should use message tool to communicate with user - logger.InfoCF("agent", "Subagent completed", - map[string]interface{}{ - "sender_id": msg.SenderID, - "channel": originChannel, - "content_len": len(content), - }) - - // Agent only logs, does not respond to user - return "", nil -} - -// assembledContext holds the pre-processed context produced by assembleContext, -// consumed by both the Generate and Stream code paths. -type assembledContext struct { - systemPrompt string - userPrompt string - fantasyHistory []fantasy.Message - adaptedTools []fantasy.AgentTool - agent fantasy.Agent -} - -func (al *AgentLoop) prepareRuntimeState(ctx context.Context, sessionKey string) (ids.UUID, ids.UUID, error) { - if al.queries == nil || al.stateStore == nil || al.kvDelegate == nil { - return ids.UUID{}, ids.UUID{}, errors.New("runtime persistence dependencies are not initialized") - } - if strings.TrimSpace(sessionKey) == "" { - return ids.UUID{}, ids.UUID{}, errors.New("session key is required") - } - - var conversationID ids.UUID - if cached, ok := al.conversationIDs.Load(sessionKey); ok { - conversationID = cached.(ids.UUID) - } else { - al.conversationMu.Lock() - defer al.conversationMu.Unlock() - if cached, ok := al.conversationIDs.Load(sessionKey); ok { - conversationID = cached.(ids.UUID) - } else { - conversationID = ids.New() - title := sessionKey - if _, err := al.queries.CreateAgentConversation(ctx, memsqlc.CreateAgentConversationParams{ - ID: conversationID, - Title: &title, - }); err != nil { - return ids.UUID{}, ids.UUID{}, fmt.Errorf("create agent conversation: %w", err) - } - al.conversationIDs.Store(sessionKey, conversationID) - } - } - - run, err := al.stateStore.CreateRun(ctx, conversationID) - if err != nil { - return ids.UUID{}, ids.UUID{}, fmt.Errorf("create agent run: %w", err) - } - - return conversationID, run.ID, nil -} - -// assembleContext performs the shared pre-processing for every agent turn: -// record channel, update tool contexts, load memory blocks, build messages, -// DAG-compress history, split into system/history/user, adapt tools, create Fantasy agent. -func (al *AgentLoop) assembleContext(ctx context.Context, opts processOptions) (assembledContext, error) { - if opts.Channel != "" && opts.ChatID != "" { - if !constants.IsInternalChannel(opts.Channel) { - channelKey := fmt.Sprintf("%s:%s", opts.Channel, opts.ChatID) - if err := al.RecordLastChannel(ctx, channelKey); err != nil { - logger.WarnCF("agent", "Failed to record last channel: %v", map[string]interface{}{"error": err.Error()}) - } - } - } - - logger.DebugCF("agent", "assembleContext: starting", - map[string]interface{}{ - "session_key": opts.SessionKey, - "channel": opts.Channel, - "sender_id": opts.SenderID, - }) - al.updateToolContexts(opts.Channel, opts.ChatID) - - block := al.obsManager.LoadBlock(ctx, opts.SessionKey) - al.contextBuilder.SetObservationBlock(block) - - kb := tools.LoadKnowledgeBlock(ctx, al.memDelegate, opts.SessionKey) - al.contextBuilder.SetKnowledgeBlock(kb) - - var history []messages.Message - var summary string - if !opts.NoHistory { - history = al.sessions.GetHistory(opts.SessionKey) - summary = al.sessions.GetSummary(opts.SessionKey) - } - - history = al.applyDAGCompression(ctx, opts.SessionKey, history) - - if al.identitySync != nil { - _ = al.identitySync.CheckAndSync(ctx) - } - - builtMsgs := al.contextBuilder.BuildMessages(history, summary, opts.UserMessage, nil, opts.Channel, opts.ChatID) - al.sessions.AddMessage(opts.SessionKey, "user", opts.UserMessage) - - systemPrompt := "" - var historyMsgs []messages.Message - userPrompt := opts.UserMessage - - if len(builtMsgs) > 0 && builtMsgs[0].Role == "system" { - systemPrompt = builtMsgs[0].Content - if len(builtMsgs) > 2 { - historyMsgs = builtMsgs[1 : len(builtMsgs)-1] - } - } - - logger.DebugCF("agent", "assembleContext: history messages", - map[string]interface{}{ - "history": formatMessagesForLog(historyMsgs), - }) - fantasyHistory := picofantasy.MessagesToFantasy(historyMsgs) - - adaptCfg := picofantasy.AdaptedToolsConfig{ - MemStore: al.memoryStore, - AgentID: "dragonscale", - SessionKey: opts.SessionKey, - } - adaptedTools := picofantasy.BuildAdaptedTools(al.tools, al.bus, opts.Channel, opts.ChatID, adaptCfg) - if al.toolResultSearch != nil { - adaptedTools = append(adaptedTools, al.toolResultSearch) - } - - // Dynamic tool promotion via PrepareStep: after tool_search discovers tools, - // they become native callables in the next inference step — no tool_call needed. - promotedSet := make(map[string]bool) - for _, at := range adaptedTools { - promotedSet[at.Info().Name] = true - } - registry := al.tools - msgBus := al.bus - channel := opts.Channel - chatID := opts.ChatID - - prepareStep := func(ctx context.Context, psOpts fantasy.PrepareStepFunctionOptions) (context.Context, fantasy.PrepareStepResult, error) { - _ = psOpts - discovered := registry.DrainDiscovered() - if len(discovered) == 0 { - return ctx, fantasy.PrepareStepResult{}, nil - } - - var newTools []tools.Tool - for _, t := range discovered { - if promotedSet[t.Name()] { - continue - } - newTools = append(newTools, t) - promotedSet[t.Name()] = true - } - - if len(newTools) == 0 { - return ctx, fantasy.PrepareStepResult{}, nil - } - - newAdapted := picofantasy.AdaptTools(newTools, msgBus, channel, chatID, adaptCfg) - expanded := append(adaptedTools, newAdapted...) - adaptedTools = expanded - - logger.InfoCF("agent", "Dynamic tool promotion via PrepareStep", - map[string]interface{}{ - "promoted": len(newTools), - "total_tools": len(expanded), - "names": toolNames(newTools), - }) - - return ctx, fantasy.PrepareStepResult{ - Tools: expanded, - }, nil - } - - conversationID, runID, err := al.prepareRuntimeState(ctx, opts.SessionKey) - if err != nil { - return assembledContext{}, err - } - - baseRuntime := OffloadingToolRuntime{ - Base: fantasy.DAGToolRuntime{MaxConcurrency: defaultToolMaxConcurrency}, - KV: al.kvDelegate, - Queries: al.queries, - ConversationID: conversationID, - RunID: runID, - } - toolRuntime := SecureBusToolRuntime{ - Base: baseRuntime, - Bus: al.secureBus, - SessionKey: opts.SessionKey, - StateStore: al.stateStore, - RunID: runID, - } - - agentOpts := []fantasy.AgentOption{ - fantasy.WithTools(adaptedTools...), - fantasy.WithStopConditions(fantasy.StepCountIs(al.maxIterations)), - fantasy.WithPrepareStep(prepareStep), - fantasy.WithToolRuntime(toolRuntime), - } - if systemPrompt != "" { - agentOpts = append(agentOpts, fantasy.WithSystemPrompt(systemPrompt)) - } - agent := fantasy.NewAgent(al.languageModel, agentOpts...) - - logger.DebugCF("agent", "Fantasy agent created", - map[string]interface{}{ - "model": al.model, - "tools_count": len(adaptedTools), - "history_count": len(historyMsgs), - "max_iterations": al.maxIterations, - "memory_enabled": true, - }) - - return assembledContext{ - systemPrompt: systemPrompt, - userPrompt: userPrompt, - fantasyHistory: fantasyHistory, - adaptedTools: adaptedTools, - agent: agent, - }, nil -} - -// postProcess handles the common finalization after Generate or Stream: -// extract final text, save session, summarize, observe, optionally send response. -func (al *AgentLoop) postProcess(ctx context.Context, opts processOptions, finalContent string, stepCount int) string { - al.sessions.Save(opts.SessionKey) - - if opts.EnableSummary { - al.maybeSummarize(ctx, opts.SessionKey, opts.Channel, opts.ChatID) - } - - tail := al.sessionsToMessagePairs(opts.SessionKey) - al.obsManager.MaybeObserveAsync(ctx, opts.SessionKey, tail) - - if opts.SendResponse { - al.bus.PublishOutbound(bus.OutboundMessage{ - Channel: opts.Channel, - ChatID: opts.ChatID, - Content: finalContent, - }) - } - - responsePreview := utils.Truncate(finalContent, 120) - logger.InfoCF("agent", fmt.Sprintf("Response: %s", responsePreview), - map[string]interface{}{ - "session_key": opts.SessionKey, - "steps": stepCount, - "final_length": len(finalContent), - }) - - return finalContent -} - -// resolveFinalContent normalizes the final assistant response from an agent run. -// Some providers return an empty final response even though an earlier step -// already produced text. In that case, recover the latest non-empty text from -// steps. If no text exists at all, return a deterministic error. -func (al *AgentLoop) resolveFinalContent(finalContent string, steps []fantasy.StepResult) (string, error) { - trimmed := strings.TrimSpace(finalContent) - if trimmed != "" { - return trimmed, nil - } - - for i := len(steps) - 1; i >= 0; i-- { - stepText := strings.TrimSpace(steps[i].Content.Text()) - if stepText != "" { - logger.WarnCF("agent", "Recovered empty final response from prior step text", - map[string]interface{}{ - "step_index": i, - }) - return stepText, nil - } - } - - type candidate struct { - text string - score int - } - candidates := make([]candidate, 0, 8) - for i := len(steps) - 1; i >= 0; i-- { - toolResults := steps[i].Content.ToolResults() - for j := len(toolResults) - 1; j >= 0; j-- { - tr := toolResults[j] - switch out := tr.Result.(type) { - case fantasy.ToolResultOutputContentText: - txt := strings.TrimSpace(out.Text) - if txt != "" { - score := 2 - if tr.ToolName == "tool_search" || strings.Contains(strings.ToLower(txt), "\"kind\":\"tool\"") { - score = 0 - } - if strings.Contains(strings.ToLower(txt), "tool not found") || - strings.Contains(strings.ToLower(txt), "path is required") { - score = -1 - } - candidates = append(candidates, candidate{text: txt, score: score}) - } - case fantasy.ToolResultOutputContentError: - if out.Error != nil { - txt := strings.TrimSpace(out.Error.Error()) - if txt != "" { - candidates = append(candidates, candidate{text: txt, score: -1}) - } - } - case fantasy.ToolResultOutputContentMedia: - txt := strings.TrimSpace(out.Text) - if txt != "" { - candidates = append(candidates, candidate{text: txt, score: 1}) - } - } - if len(candidates) >= 8 { - break - } - } - if len(candidates) >= 8 { - break - } - } - - bestText := "" - bestScore := -1000 - for _, c := range candidates { - if c.score > bestScore { - bestScore = c.score - bestText = c.text - } - } - - if bestText != "" && bestScore > 0 { - logger.WarnCF("agent", "Recovered empty final response from tool results", - map[string]interface{}{ - "candidates": len(candidates), - "score": bestScore, - }) - return bestText, nil - } - - toolCalls := 0 - for _, step := range steps { - toolCalls += len(step.Content.ToolCalls()) - } - - return "", fmt.Errorf("agent produced no final response text (steps=%d, tool_calls=%d)", len(steps), toolCalls) -} - -// runAgentLoop is the core message processing logic. -// It delegates to assembleContext for shared pre-processing, then branches on -// opts.Streaming to either Generate (synchronous) or Stream (real-time deltas). -func (al *AgentLoop) runAgentLoop(ctx context.Context, opts processOptions) (string, error) { - al.activeSessionKey.Store(opts.SessionKey) - - ac, err := al.assembleContext(ctx, opts) - if err != nil { - return "", err - } - - if opts.Streaming { - return al.runStreaming(ctx, opts, ac) - } - - result, err := ac.agent.Generate(ctx, fantasy.AgentCall{ - Prompt: ac.userPrompt, - Messages: ac.fantasyHistory, - }) - if err != nil { - logger.ErrorCF("agent", "Fantasy Generate failed", - map[string]interface{}{"error": err.Error()}) - return "", fmt.Errorf("agent Generate failed: %w", err) - } - - for _, step := range result.Steps { - stepMsgs := picofantasy.StepToMessages(step) - for _, m := range stepMsgs { - al.sessions.AddFullMessage(opts.SessionKey, m) - } - al.auditStep(ctx, step, opts.SessionKey) - } - - finalContent, err := al.resolveFinalContent(result.Response.Content.Text(), result.Steps) - if err != nil { - logger.ErrorCF("agent", "Agent finished without final response text", - map[string]interface{}{ - "error": err.Error(), - "steps": len(result.Steps), - }) - return "", err - } - return al.postProcess(ctx, opts, finalContent, len(result.Steps)), nil -} - -// runStreaming uses Fantasy's agent.Stream() to stream token deltas to the bus -// in real time, using the pre-assembled context from assembleContext. -func (al *AgentLoop) runStreaming(ctx context.Context, opts processOptions, ac assembledContext) (string, error) { - streamCall := fantasy.AgentStreamCall{ - Prompt: ac.userPrompt, - Messages: ac.fantasyHistory, - - OnTextDelta: func(id, text string) error { - if opts.Channel != "" && opts.ChatID != "" { - al.bus.PublishOutbound(bus.OutboundMessage{ - Channel: opts.Channel, - ChatID: opts.ChatID, - Content: text, - StreamDelta: true, - }) - } - return nil - }, - - OnStepFinish: func(step fantasy.StepResult) error { - stepMsgs := picofantasy.StepToMessages(step) - for _, m := range stepMsgs { - al.sessions.AddFullMessage(opts.SessionKey, m) - } - al.auditStep(ctx, step, opts.SessionKey) - return nil - }, - - OnToolCall: func(tc fantasy.ToolCallContent) error { - logger.DebugCF("agent", "Streaming tool call", - map[string]interface{}{ - "tool": tc.ToolName, - "id": tc.ToolCallID, - }) - return nil - }, - } - - result, err := ac.agent.Stream(ctx, streamCall) - if err != nil { - logger.ErrorCF("agent", "Fantasy Stream failed", - map[string]interface{}{"error": err.Error()}) - return "", fmt.Errorf("agent Stream failed: %w", err) - } - - finalContent, err := al.resolveFinalContent(result.Response.Content.Text(), result.Steps) - if err != nil { - logger.ErrorCF("agent", "Streaming agent finished without final response text", - map[string]interface{}{ - "error": err.Error(), - "steps": len(result.Steps), - }) - return "", err - } - return al.postProcess(ctx, opts, finalContent, len(result.Steps)), nil -} - -// runLLMIteration — DELETED. Replaced by Fantasy's internal agent loop. - -// auditStep logs tool calls from a Fantasy step result to the audit log. -func (al *AgentLoop) auditStep(ctx context.Context, step fantasy.StepResult, sessionKey string) { - toolCalls := step.Content.ToolCalls() - if len(toolCalls) == 0 { - return - } - - for _, tc := range toolCalls { - entry := &memory.AuditEntry{ - ID: ids.New(), - AgentID: "dragonscale", - SessionKey: sessionKey, - Action: "tool_call", - Target: tc.ToolName, - Input: tc.Input, - } - aCtx, cancel := context.WithTimeout(ctx, time.Second) - if err := al.memDelegate.InsertAuditEntry(aCtx, entry); err != nil { - logger.WarnCF("agent", "Failed to log audit entry", - map[string]interface{}{"tool": tc.ToolName, "error": err.Error()}) - } - cancel() - } -} - -// updateToolContexts updates the context for tools that need channel/chatID info. -func (al *AgentLoop) updateToolContexts(channel, chatID string) { - // Use ContextualTool interface instead of type assertions - if tool, ok := al.tools.Get("message"); ok { - if mt, ok := tool.(tools.ContextualTool); ok { - mt.SetContext(channel, chatID) - } - } - if tool, ok := al.tools.Get("spawn"); ok { - if st, ok := tool.(tools.ContextualTool); ok { - st.SetContext(channel, chatID) - } - } - if tool, ok := al.tools.Get("subagent"); ok { - if st, ok := tool.(tools.ContextualTool); ok { - st.SetContext(channel, chatID) - } - } -} - -// maybeSummarize only triggers emergency compression when hard limits are exceeded. -// Normal background compaction is intentionally disabled for the unified kernel. -func (al *AgentLoop) maybeSummarize(ctx context.Context, sessionKey, channel, chatID string) { - _ = channel - _ = chatID - newHistory := al.sessions.GetHistory(sessionKey) - tokenEstimate := al.estimateTokens(newHistory) - criticalThreshold := al.contextWindow * 95 / 100 - - if tokenEstimate > criticalThreshold { - al.forceCompression(ctx, sessionKey) - } - - // TODO: actually use the channel and chatID to push data to the bus - //if _, loading := al.summarizing.LoadOrStore(sessionKey, true); !loading { - // go func() { - // defer al.summarizing.Delete(sessionKey) - // if !constants.IsInternalChannel(channel) { - // al.bus.PublishOutbound(bus.OutboundMessage{ - // Channel: channel, - // ChatID: chatID, - // Content: "⚠️ Memory threshold reached. Optimizing conversation history...", - // }) - // } - // al.summarizeSession(ctx, sessionKey) - // }() - // } -} - -// EmergencyProvenance captures provenance metadata for postmortem when -// emergency compression cycles run. Persisted via the audit pipeline. -type EmergencyProvenance struct { - SessionKey string `json:"session_key"` - Cycle int `json:"cycle"` - TokenEstimate int `json:"token_estimate"` - CriticalBudget int `json:"critical_budget"` - HistoryMsgCount int `json:"history_msg_count"` -} - -// persistEmergencyProvenance writes provenance metadata to the audit log. -// Best-effort: logs warning on failure, never fails the compression path. -func (al *AgentLoop) persistEmergencyProvenance(ctx context.Context, prov EmergencyProvenance) { - if al.memDelegate == nil { - return - } - input, err := json.Marshal(prov) - if err != nil { - logger.WarnCF("agent", "Failed to marshal emergency provenance", - map[string]interface{}{"error": err.Error()}) - return - } - entry := &memory.AuditEntry{ - ID: ids.New(), - AgentID: "dragonscale", - SessionKey: prov.SessionKey, - Action: "emergency_compression", - Target: fmt.Sprintf("cycle_%d", prov.Cycle), - Input: string(input), - } - aCtx, cancel := context.WithTimeout(ctx, time.Second) - defer cancel() - if err := al.memDelegate.InsertAuditEntry(aCtx, entry); err != nil { - logger.WarnCF("agent", "Failed to persist emergency provenance", - map[string]interface{}{"error": err.Error(), "session_key": prov.SessionKey}) - } -} - -// forceCompression performs emergency recursive compression by repeatedly -// summarizing older history until under hard budget, without deleting immutable -// persisted session records. -func (al *AgentLoop) forceCompression(ctx context.Context, sessionKey string) { - const maxCycles = 3 - for cycle := 1; cycle <= maxCycles; cycle++ { - history := al.sessions.GetHistory(sessionKey) - if len(history) <= al.continuityKeepCount(history) { - return - } - tokenEstimate := al.estimateTokens(history) - criticalThreshold := al.contextWindow * 95 / 100 - if tokenEstimate <= criticalThreshold { - return - } - - logger.WarnCF("agent", "Emergency compression cycle triggered", - map[string]interface{}{ - "session_key": sessionKey, - "cycle": cycle, - "token_estimate": tokenEstimate, - "critical_budget": criticalThreshold, - }) - - al.persistEmergencyProvenance(ctx, EmergencyProvenance{ - SessionKey: sessionKey, - Cycle: cycle, - TokenEstimate: tokenEstimate, - CriticalBudget: criticalThreshold, - HistoryMsgCount: len(history), - }) - - al.summarizeSession(ctx, sessionKey) + target, ok := raw.(outputTarget) + if !ok || target.Channel == "" { + return outputTarget{}, false } + return target, true } // MemoryDelegate returns the active memory delegate (nil if memory system is disabled). @@ -1294,10 +499,10 @@ func (al *AgentLoop) GetStartupInfo() map[string]interface{} { info := make(map[string]interface{}) // Tools info - tools := al.tools.List() + toolNames := al.tools.List() info["tools"] = map[string]interface{}{ - "count": len(tools), - "names": tools, + "count": len(toolNames), + "names": toolNames, } // Skills info @@ -1305,516 +510,3 @@ func (al *AgentLoop) GetStartupInfo() map[string]interface{} { return info } - -// formatMessagesForLog formats messages for logging -func formatMessagesForLog(msgs []messages.Message) string { - if len(msgs) == 0 { - return "[]" - } - - var result string - result += "[\n" - for i, msg := range msgs { - result += fmt.Sprintf(" [%d] Role: %s\n", i, msg.Role) - if len(msg.ToolCalls) > 0 { - result += " ToolCalls:\n" - for _, tc := range msg.ToolCalls { - result += fmt.Sprintf(" - ID: %s, Type: %s, Name: %s\n", tc.ID, tc.Type, tc.Name) - if tc.Function != nil { - result += fmt.Sprintf(" Arguments: %s\n", utils.Truncate(tc.Function.Arguments, 200)) - } - } - } - if msg.Content != "" { - content := utils.Truncate(msg.Content, 200) - result += fmt.Sprintf(" Content: %s\n", content) - } - if msg.ToolCallID != "" { - result += fmt.Sprintf(" ToolCallID: %s\n", msg.ToolCallID) - } - result += "\n" - } - result += "]" - return result -} - -func (al *AgentLoop) continuityRetentionPolicy() config.ContinuityRetentionConfig { - policy := config.ContinuityRetentionConfig{ - MinMessages: 4, - MaxMessages: 24, - TargetContextRatio: 0.10, - FailureKeepMessages: 10, - } - - if al.cfg != nil { - cfgPolicy := al.cfg.Agents.Defaults.ContinuityRetention - if cfgPolicy.MinMessages > 0 { - policy.MinMessages = cfgPolicy.MinMessages - } - if cfgPolicy.MaxMessages > 0 { - policy.MaxMessages = cfgPolicy.MaxMessages - } - if cfgPolicy.TargetContextRatio > 0 && cfgPolicy.TargetContextRatio <= 0.5 { - policy.TargetContextRatio = cfgPolicy.TargetContextRatio - } - if cfgPolicy.FailureKeepMessages > 0 { - policy.FailureKeepMessages = cfgPolicy.FailureKeepMessages - } - } - - if policy.MaxMessages < policy.MinMessages { - policy.MaxMessages = policy.MinMessages - } - if policy.FailureKeepMessages < policy.MinMessages { - policy.FailureKeepMessages = policy.MinMessages - } - - return policy -} - -func (al *AgentLoop) continuityKeepCount(history []messages.Message) int { - if len(history) == 0 { - return 0 - } - - policy := al.continuityRetentionPolicy() - minKeep := policy.MinMessages - if minKeep > len(history) { - minKeep = len(history) - } - maxKeep := policy.MaxMessages - if maxKeep > len(history) { - maxKeep = len(history) - } - if maxKeep < minKeep { - maxKeep = minKeep - } - - contextWindow := al.contextWindow - if contextWindow <= 0 && al.cfg != nil { - contextWindow = al.cfg.Agents.Defaults.MaxTokens - } - if contextWindow <= 0 { - return minKeep - } - - targetTokens := int(float64(contextWindow) * policy.TargetContextRatio) - if targetTokens <= 0 { - return minKeep - } - - keep := 0 - keptTokens := 0 - for i := len(history) - 1; i >= 0 && keep < maxKeep; i-- { - msgTokens := observation.EstimateTokens(history[i].Content) + 4 - if keep >= minKeep && keptTokens+msgTokens > targetTokens { - break - } - keptTokens += msgTokens - keep++ - } - if keep < minKeep { - keep = minKeep - } - return keep -} - -type oversizedRecoveryCandidate struct { - Message messages.Message - OriginalIndex int - TokenEstimate int -} - -func (al *AgentLoop) persistOversizedRecoveryRefs(ctx context.Context, sessionKey string, omitted []oversizedRecoveryCandidate) ([]string, error) { - if len(omitted) == 0 { - return nil, nil - } - if al.memDelegate == nil { - return nil, fmt.Errorf("memory delegate is not configured") - } - - const maxPersistedRefs = 8 - refs := make([]string, 0, len(omitted)) - now := time.Now().UTC() - - for i, candidate := range omitted { - if i >= maxPersistedRefs { - break - } - nodeID := tools.DAGRecoveryNodePrefix + ids.New().String() - record := tools.DAGRecoveryRecord{ - NodeID: nodeID, - SessionKey: sessionKey, - OriginalIndex: candidate.OriginalIndex, - Role: candidate.Message.Role, - Content: candidate.Message.Content, - TokenEstimate: candidate.TokenEstimate, - Reason: "oversized_message_omitted_from_summary", - CreatedAt: now, - } - - data, err := json.Marshal(record) - if err != nil { - return refs, fmt.Errorf("marshal DAG recovery record: %w", err) - } - if err := al.memDelegate.UpsertKV(ctx, "dragonscale", tools.DAGRecoveryKVKey(sessionKey, nodeID), string(data)); err != nil { - return refs, fmt.Errorf("persist DAG recovery record: %w", err) - } - refs = append(refs, nodeID) - } - return refs, nil -} - -// summarizeSession summarizes the conversation history for a session. -func (al *AgentLoop) summarizeSession(parentCtx context.Context, sessionKey string) { - ctx, cancel := context.WithTimeout(parentCtx, 120*time.Second) - defer cancel() - - history := al.sessions.GetHistory(sessionKey) - summary := al.sessions.GetSummary(sessionKey) - - keepLast := al.continuityKeepCount(history) - if len(history) <= keepLast { - return - } - - toSummarize := history[:len(history)-keepLast] - - // Oversized Message Guard: skip individual messages that would consume too - // much of the summarizer's context. Use 40% of the window for the summarizer - // input budget, reserving the rest for system prompt + summary output. - // Oversized omissions are persisted as DAG recovery references. - maxMessageTokens := al.contextWindow * 40 / 100 - if maxMessageTokens < 2048 { - maxMessageTokens = 2048 - } - validMessages := make([]messages.Message, 0) - omitted := false - omittedMessages := make([]oversizedRecoveryCandidate, 0) - - for idx, m := range toSummarize { - if m.Role != "user" && m.Role != "assistant" { - continue - } - msgTokens := observation.EstimateTokens(m.Content) - if msgTokens > maxMessageTokens { - omitted = true - omittedMessages = append(omittedMessages, oversizedRecoveryCandidate{ - Message: m, - OriginalIndex: idx, - TokenEstimate: msgTokens, - }) - continue - } - validMessages = append(validMessages, m) - } - - if len(validMessages) == 0 { - return - } - - // Multi-Part Summarization - var finalSummary string - if len(validMessages) > 10 { - mid := len(validMessages) / 2 - part1 := validMessages[:mid] - part2 := validMessages[mid:] - - s1, _ := al.summarizeBatch(ctx, part1, "") - s2, _ := al.summarizeBatch(ctx, part2, "") - - // Merge them - mergePrompt := fmt.Sprintf("Merge these two conversation summaries into one cohesive summary:\n\n1: %s\n\n2: %s", s1, s2) - resp, err := al.callModel(ctx, mergePrompt) - if err == nil { - finalSummary = resp - } else { - finalSummary = s1 + " " + s2 - } - } else { - finalSummary, _ = al.summarizeBatch(ctx, validMessages, summary) - } - - if omitted && finalSummary != "" { - recoveryRefs, err := al.persistOversizedRecoveryRefs(ctx, sessionKey, omittedMessages) - if err != nil { - logger.WarnCF("agent", "Failed to persist DAG recovery references for oversized messages", - map[string]interface{}{ - "session_key": sessionKey, - "error": err.Error(), - "omitted": len(omittedMessages), - }) - finalSummary += "\n[Note: Some oversized messages were omitted from this summary for efficiency.]" - } else if len(recoveryRefs) > 0 { - finalSummary += fmt.Sprintf("\n[Note: %d oversized message(s) were omitted from this summary. Recovery refs: %s. Use dag_expand with node_id= to recover full content.]", - len(omittedMessages), strings.Join(recoveryRefs, ", ")) - } else { - finalSummary += "\n[Note: Some oversized messages were omitted from this summary for efficiency.]" - } - } - - if finalSummary != "" { - al.sessions.SetSummary(sessionKey, finalSummary) - al.sessions.TruncateHistory(sessionKey, keepLast) - al.sessions.Save(sessionKey) - al.summarizeFailures.Delete(sessionKey) - } else { - var count int - if v, ok := al.summarizeFailures.Load(sessionKey); ok { - count = v.(int) - } - count++ - al.summarizeFailures.Store(sessionKey, count) - - const maxSummarizeFailures = 3 - emergencyKeep := al.continuityRetentionPolicy().FailureKeepMessages - if count >= maxSummarizeFailures { - logger.ErrorCF("agent", "Summarization failed repeatedly, force-truncating session", - map[string]interface{}{ - "session": sessionKey, - "consecutive_failures": count, - "keep": emergencyKeep, - }) - al.sessions.TruncateHistory(sessionKey, emergencyKeep) - al.sessions.Save(sessionKey) - al.summarizeFailures.Delete(sessionKey) - } - } -} - -// summarizeBatch summarizes a batch of messages using the Fantasy LanguageModel directly. -func (al *AgentLoop) summarizeBatch(ctx context.Context, batch []messages.Message, existingSummary string) (string, error) { - var prompt strings.Builder - prompt.WriteString("Provide a concise summary of this conversation segment, preserving core context and key points.\n") - if existingSummary != "" { - fmt.Fprintf(&prompt, "Existing context: %s\n", existingSummary) - } - prompt.WriteString("\nCONVERSATION:\n") - for _, m := range batch { - fmt.Fprintf(&prompt, "%s: %s\n", m.Role, m.Content) - } - - return al.callModel(ctx, prompt.String()) -} - -// callModel makes a direct call to the Fantasy LanguageModel (no tools, no agent loop). -// Used for summarization and other simple generation tasks. -func (al *AgentLoop) callModel(ctx context.Context, prompt string) (string, error) { - temp := 0.3 - maxTokens := int64(1024) - - resp, err := al.languageModel.Generate(ctx, fantasy.Call{ - Prompt: fantasy.Prompt{ - fantasy.NewUserMessage(prompt), - }, - Temperature: &temp, - MaxOutputTokens: &maxTokens, - }) - if err != nil { - return "", err - } - return resp.Content.Text(), nil -} - -// sessionsToMessagePairs converts the session history to observation.MessagePair -// for token estimation by the observation manager. -func (al *AgentLoop) sessionsToMessagePairs(sessionKey string) []observation.MessagePair { - history := al.sessions.GetHistory(sessionKey) - pairs := make([]observation.MessagePair, len(history)) - for i, m := range history { - pairs[i] = observation.MessagePair{Role: m.Role, Content: m.Content} - } - return pairs -} - -// applyDAGCompression compresses old history into a DAG summary block and -// returns only the tail messages that should be passed as raw conversation. -// The compressed portion is injected into the system prompt via contextBuilder. -// When memDelegate implements dag.DAGPersister, the DAG is persisted for dag_expand/describe/grep. -func (al *AgentLoop) applyDAGCompression(ctx context.Context, sessionKey string, history []messages.Message) []messages.Message { - const minHistoryForDAG = 16 - - if len(history) < minHistoryForDAG { - al.contextBuilder.SetDAGBlock("") - return history - } - - budget := dag.ComputeBudget(al.contextWindow, dag.DefaultBudgetConfig()) - tailCount := dag.TailMessageCount(budget.RawTail) - if tailCount >= len(history) { - al.contextBuilder.SetDAGBlock("") - return history - } - - // Split: compress old, keep tail raw - compressible := history[:len(history)-tailCount] - tail := history[len(history)-tailCount:] - - // Tool-call-aware: don't split on a "tool" message - for len(tail) > 0 && tail[0].Role == "tool" && len(compressible) > 0 { - tail = append([]messages.Message{compressible[len(compressible)-1]}, tail...) - compressible = compressible[:len(compressible)-1] - } - - if len(compressible) == 0 { - al.contextBuilder.SetDAGBlock("") - return history - } - - dagMsgs := make([]dag.Message, len(compressible)) - for i, m := range compressible { - dagMsgs[i] = dag.Message{Role: m.Role, Content: m.Content} - } - - compressor := dag.NewCompressor(dag.DefaultCompressorConfig()) - d := compressor.Compress(dagMsgs) - - rendered := dag.RenderDAGForBudget(d, budget.DAGSummaries) - al.contextBuilder.SetDAGBlock(rendered) - - // Persist DAG for dag_expand, dag_describe, dag_grep (additive; in-memory behavior unchanged) - if dp, ok := al.memDelegate.(dag.DAGPersister); ok { - if err := dp.PersistDAG(ctx, "dragonscale", sessionKey, &dag.PersistSnapshot{ - FromMsgIdx: 0, - ToMsgIdx: len(compressible), - MsgCount: len(compressible), - DAG: d, - }); err != nil { - logger.WarnCF("agent", "DAG persist failed (non-fatal)", - map[string]interface{}{"error": err.Error(), "session_key": sessionKey}) - } - } - - logger.DebugCF("agent", "DAG compression applied", - map[string]interface{}{ - "total_msgs": len(history), - "compressed_msgs": len(compressible), - "tail_msgs": len(tail), - "dag_nodes": len(d.Nodes), - }) - - return tail -} - -// estimateTokens estimates the number of tokens in a message list. -// listConfiguredModels returns a human-readable summary of which providers -// have API credentials configured, and the current default model. -func listConfiguredModels(cfg *config.Config) string { - if cfg == nil { - return "No configuration available." - } - - current := fmt.Sprintf("Current model: %s", cfg.Agents.Defaults.Model) - if cfg.Agents.Defaults.Provider != "" { - current += fmt.Sprintf(" (provider: %s)", cfg.Agents.Defaults.Provider) - } - - configured := cfg.Providers.ConfiguredNames() - if len(configured) == 0 { - return current + "\nNo providers configured — set API keys in config.json or environment variables." - } - - return current + "\nConfigured providers: " + strings.Join(configured, ", ") -} - -func (al *AgentLoop) estimateTokens(msgs []messages.Message) int { - pairs := make([]observation.MessagePair, 0, len(msgs)) - for _, m := range msgs { - pairs = append(pairs, observation.MessagePair{ - Role: m.Role, - Content: m.Content, - }) - } - return observation.EstimateMessagesTokens(pairs) -} - -// FIXME: Leverage Cobra with subcommand command palette pattern for commands -func (al *AgentLoop) handleCommand(_ context.Context, msg bus.InboundMessage) (string, bool) { - content := strings.TrimSpace(msg.Content) - if !strings.HasPrefix(content, "/") { - return "", false - } - - parts := strings.Fields(content) - if len(parts) == 0 { - return "", false - } - - cmd := parts[0] - args := parts[1:] - - switch cmd { - case "/show": - if len(args) < 1 { - return "Usage: /show [model|channel]", true - } - switch args[0] { - case "model": - return fmt.Sprintf("Current model: %s", al.model), true - case "channel": - return fmt.Sprintf("Current channel: %s", msg.Channel), true - default: - return fmt.Sprintf("Unknown show target: %s", args[0]), true - } - - case "/list": - if len(args) < 1 { - return "Usage: /list [models|channels]", true - } - switch args[0] { - case "models": - return listConfiguredModels(al.cfg), true - case "channels": - if al.channelManager == nil { - return "Channel manager not initialized", true - } - channels := al.channelManager.GetEnabledChannels() - if len(channels) == 0 { - return "No channels enabled", true - } - return fmt.Sprintf("Enabled channels: %s", strings.Join(channels, ", ")), true - default: - return fmt.Sprintf("Unknown list target: %s", args[0]), true - } - - case "/switch": - if len(args) < 3 || args[1] != "to" { - return "Usage: /switch [model|channel] to ", true - } - target := args[0] - value := args[2] - - switch target { - case "model": - oldModel := al.model - al.model = value - return fmt.Sprintf("Switched model from %s to %s", oldModel, value), true - case "channel": - // FIXME: This changes the 'default' channel for some operations, or effectively redirects output? - // For now, let's just validate if the channel exists - if al.channelManager == nil { - return "Channel manager not initialized", true - } - if _, exists := al.channelManager.GetChannel(value); !exists && value != "cli" { - return fmt.Sprintf("Channel '%s' not found or not enabled", value), true - } - - // FIXME: If message came from CLI, maybe we want to redirect CLI output to this channel? - // That would require state persistence about "redirected channel" - // For now, just acknowledged. - return fmt.Sprintf("Switched target channel to %s (Note: this currently only validates existence)", value), true - default: - return fmt.Sprintf("Unknown switch target: %s", target), true - } - } - - return "", false -} - -func toolNames(tt []tools.Tool) []string { - names := make([]string, len(tt)) - for i, t := range tt { - names[i] = t.Name() - } - return names -} diff --git a/pkg/agent/loop_test.go b/pkg/agent/loop_test.go index 9fbc4884b..11d9410cd 100644 --- a/pkg/agent/loop_test.go +++ b/pkg/agent/loop_test.go @@ -12,6 +12,7 @@ import ( "time" fantasy "charm.land/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg" "github.com/ZanzyTHEbar/dragonscale/pkg/bus" "github.com/ZanzyTHEbar/dragonscale/pkg/config" memsqlc "github.com/ZanzyTHEbar/dragonscale/pkg/memory/sqlc" @@ -827,13 +828,13 @@ func TestForceCompression_PersistsProvenance(t *testing.T) { } ctx := context.Background() - al.forceCompression(ctx, sessionKey) + al.forceCompression(ctx, sessionKey, "", "") del := al.MemoryDelegate() if del == nil { t.Fatal("MemoryDelegate is nil") } - entries, err := del.ListAuditEntriesByAction(ctx, "dragonscale", "emergency_compression", 50) + entries, err := del.ListAuditEntriesByAction(ctx, pkg.NAME, "emergency_compression", 50) if err != nil { t.Fatalf("ListAuditEntriesByAction: %v", err) } @@ -913,7 +914,7 @@ func TestPersistOversizedRecoveryRefs_CreatesRecoverableReferences(t *testing.T) dagTool := tools.NewDagExpandTool(tools.DAGToolDeps{ Delegate: al.MemoryDelegate(), - AgentID: "dragonscale", + AgentID: pkg.NAME, SessionFn: func() string { return "recovery-session" }, diff --git a/pkg/agent/message_router.go b/pkg/agent/message_router.go new file mode 100644 index 000000000..884dc3262 --- /dev/null +++ b/pkg/agent/message_router.go @@ -0,0 +1,171 @@ +package agent + +import ( + "context" + "fmt" + "strings" + + "github.com/ZanzyTHEbar/dragonscale/pkg/bus" + "github.com/ZanzyTHEbar/dragonscale/pkg/constants" + "github.com/ZanzyTHEbar/dragonscale/pkg/logger" + "github.com/ZanzyTHEbar/dragonscale/pkg/utils" +) + +func (al *AgentLoop) RecordLastChannel(ctx context.Context, channel string) error { + return al.state.SetLastChannel(ctx, channel) +} + +func (al *AgentLoop) RecordLastChatID(ctx context.Context, chatID string) error { + return al.state.SetLastChatID(ctx, chatID) +} + +func (al *AgentLoop) ProcessDirect(ctx context.Context, content, sessionKey string) (string, error) { + return al.ProcessDirectWithChannel(ctx, content, sessionKey, "cli", "direct") +} + +func (al *AgentLoop) ProcessDirectWithChannel(ctx context.Context, content, sessionKey, channel, chatID string) (string, error) { + msg := bus.InboundMessage{ + Channel: channel, + SenderID: "cron", + ChatID: chatID, + Content: content, + SessionKey: sessionKey, + } + + return al.processMessage(ctx, msg) +} + +// ProcessDirectStreaming processes a message with streaming token delivery. +// Text deltas are published to the bus as StreamDelta messages in real time. +func (al *AgentLoop) ProcessDirectStreaming(ctx context.Context, content, sessionKey, channel, chatID string) (string, error) { + msg := bus.InboundMessage{ + Channel: channel, + SenderID: "user", + ChatID: chatID, + Content: content, + SessionKey: sessionKey, + } + + return al.runAgentLoop(ctx, processOptions{ + SessionKey: msg.SessionKey, + Channel: msg.Channel, + ChatID: msg.ChatID, + SenderID: msg.SenderID, + UserMessage: msg.Content, + EnableSummary: true, + SendResponse: false, + Streaming: true, + }) +} + +// ProcessHeartbeat processes a heartbeat request without session history. +// Each heartbeat is independent and doesn't accumulate context. +// It injects the active session's summary so the agent has awareness of +// recent user conversation context. +func (al *AgentLoop) ProcessHeartbeat(ctx context.Context, content, channel, chatID string) (string, error) { + if v := al.activeSessionKey.Load(); v != nil { + if key, ok := v.(string); ok && key != "" { + if summary := al.sessions.GetSummary(key); summary != "" { + content = content + "\n\n## Recent User Context\n" + summary + } + } + } + return al.runAgentLoop(ctx, processOptions{ + SessionKey: "heartbeat", + Channel: channel, + ChatID: chatID, + UserMessage: content, + EnableSummary: false, + SendResponse: false, + NoHistory: true, + }) +} + +func (al *AgentLoop) processMessage(ctx context.Context, msg bus.InboundMessage) (string, error) { + // Add message preview to log (show full content for error messages) + var logContent string + if strings.Contains(msg.Content, "Error:") || strings.Contains(msg.Content, "error") { + logContent = msg.Content // Full content for errors + } else { + logContent = utils.Truncate(msg.Content, 80) + } + logger.InfoCF("agent", fmt.Sprintf("Processing message from %s:%s: %s", msg.Channel, msg.SenderID, logContent), + map[string]interface{}{ + "channel": msg.Channel, + "chat_id": msg.ChatID, + "sender_id": msg.SenderID, + "session_key": msg.SessionKey, + }) + + // Route system messages to processSystemMessage + if msg.Channel == "system" { + return al.processSystemMessage(ctx, msg) + } + + // Check for commands + if response, handled := al.handleCommand(ctx, msg); handled { + return response, nil + } + + // Process as user message + return al.runAgentLoop(ctx, processOptions{ + SessionKey: msg.SessionKey, + Channel: msg.Channel, + ChatID: msg.ChatID, + UserMessage: msg.Content, + EnableSummary: true, + SendResponse: false, + }) +} + +func (al *AgentLoop) processSystemMessage(_ context.Context, msg bus.InboundMessage) (string, error) { + // Verify this is a system message + if msg.Channel != "system" { + return "", fmt.Errorf("processSystemMessage called with non-system message channel: %s", msg.Channel) + } + + logger.InfoCF("agent", "Processing system message", + map[string]interface{}{ + "sender_id": msg.SenderID, + "chat_id": msg.ChatID, + }) + + // Parse origin channel from chat_id (format: "channel:chat_id") + var originChannel string + if idx := strings.Index(msg.ChatID, ":"); idx > 0 { + originChannel = msg.ChatID[:idx] + } else { + // Fallback + originChannel = "cli" + } + + // Extract subagent result from message content + // Format: "Task 'label' completed.\n\nResult:\n" + content := msg.Content + if idx := strings.Index(content, "Result:\n"); idx >= 0 { + content = content[idx+8:] // Extract just the result part + } + + // Skip internal channels - only log, don't send to user + if constants.IsInternalChannel(originChannel) { + logger.InfoCF("agent", "Subagent completed (internal channel)", + map[string]interface{}{ + "sender_id": msg.SenderID, + "content_len": len(content), + "channel": originChannel, + }) + return "", nil + } + + // Agent acts as dispatcher only - subagent handles user interaction via message tool + // Don't forward result here, subagent should use message tool to communicate with user + logger.InfoCF("agent", "Subagent completed", + map[string]interface{}{ + "sender_id": msg.SenderID, + "channel": originChannel, + "content_len": len(content), + }) + + // Agent only logs, does not respond to user + return "", nil +} diff --git a/pkg/agent/summarizer.go b/pkg/agent/summarizer.go new file mode 100644 index 000000000..95e8def2a --- /dev/null +++ b/pkg/agent/summarizer.go @@ -0,0 +1,492 @@ +package agent + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "time" + + fantasy "charm.land/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg" + "github.com/ZanzyTHEbar/dragonscale/pkg/bus" + "github.com/ZanzyTHEbar/dragonscale/pkg/config" + "github.com/ZanzyTHEbar/dragonscale/pkg/constants" + "github.com/ZanzyTHEbar/dragonscale/pkg/ids" + "github.com/ZanzyTHEbar/dragonscale/pkg/logger" + "github.com/ZanzyTHEbar/dragonscale/pkg/memory" + "github.com/ZanzyTHEbar/dragonscale/pkg/memory/dag" + "github.com/ZanzyTHEbar/dragonscale/pkg/memory/observation" + "github.com/ZanzyTHEbar/dragonscale/pkg/messages" + "github.com/ZanzyTHEbar/dragonscale/pkg/tools" +) + +// maybeSummarize only triggers emergency compression when hard limits are exceeded. +// Normal background compaction is intentionally disabled for the unified kernel. +func (al *AgentLoop) maybeSummarize(ctx context.Context, sessionKey, channel, chatID string) { + newHistory := al.sessions.GetHistory(sessionKey) + tokenEstimate := al.estimateTokens(newHistory) + criticalThreshold := al.contextWindow * 95 / 100 + + if tokenEstimate > criticalThreshold { + al.forceCompression(ctx, sessionKey, channel, chatID) + } +} + +// EmergencyProvenance captures provenance metadata for postmortem when +// emergency compression cycles run. Persisted via the audit pipeline. +type EmergencyProvenance struct { + SessionKey string `json:"session_key"` + Cycle int `json:"cycle"` + TokenEstimate int `json:"token_estimate"` + CriticalBudget int `json:"critical_budget"` + HistoryMsgCount int `json:"history_msg_count"` +} + +// persistEmergencyProvenance writes provenance metadata to the audit log. +// Best-effort: logs warning on failure, never fails the compression path. +func (al *AgentLoop) persistEmergencyProvenance(ctx context.Context, prov EmergencyProvenance) { + if al.memDelegate == nil { + return + } + input, err := json.Marshal(prov) + if err != nil { + logger.WarnCF("agent", "Failed to marshal emergency provenance", + map[string]interface{}{"error": err.Error()}) + return + } + entry := &memory.AuditEntry{ + ID: ids.New(), + AgentID: pkg.NAME, + SessionKey: prov.SessionKey, + Action: "emergency_compression", + Target: fmt.Sprintf("cycle_%d", prov.Cycle), + Input: string(input), + } + aCtx, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + if err := al.memDelegate.InsertAuditEntry(aCtx, entry); err != nil { + logger.WarnCF("agent", "Failed to persist emergency provenance", + map[string]interface{}{"error": err.Error(), "session_key": prov.SessionKey}) + } +} + +// forceCompression performs emergency recursive compression by repeatedly +// summarizing older history until under hard budget, without deleting immutable +// persisted session records. +func (al *AgentLoop) forceCompression(ctx context.Context, sessionKey, channel, chatID string) { + if _, loading := al.summarizing.LoadOrStore(sessionKey, true); loading { + return + } + defer al.summarizing.Delete(sessionKey) + + if channel != "" && !constants.IsInternalChannel(channel) { + if al.bus != nil { + al.bus.PublishOutbound(bus.OutboundMessage{ + Channel: channel, + ChatID: chatID, + Content: "⚠️ Memory threshold reached. Optimizing conversation history...", + }) + } + } + + const maxCycles = 3 + for cycle := 1; cycle <= maxCycles; cycle++ { + history := al.sessions.GetHistory(sessionKey) + if len(history) <= al.continuityKeepCount(history) { + return + } + tokenEstimate := al.estimateTokens(history) + criticalThreshold := al.contextWindow * 95 / 100 + if tokenEstimate <= criticalThreshold { + return + } + + logger.WarnCF("agent", "Emergency compression cycle triggered", + map[string]interface{}{ + "session_key": sessionKey, + "cycle": cycle, + "token_estimate": tokenEstimate, + "critical_budget": criticalThreshold, + }) + + al.persistEmergencyProvenance(ctx, EmergencyProvenance{ + SessionKey: sessionKey, + Cycle: cycle, + TokenEstimate: tokenEstimate, + CriticalBudget: criticalThreshold, + HistoryMsgCount: len(history), + }) + + al.summarizeSession(ctx, sessionKey) + } +} + +func (al *AgentLoop) continuityRetentionPolicy() config.ContinuityRetentionConfig { + policy := config.ContinuityRetentionConfig{ + MinMessages: 4, + MaxMessages: 24, + TargetContextRatio: 0.10, + FailureKeepMessages: 10, + } + + if al.cfg != nil { + cfgPolicy := al.cfg.Agents.Defaults.ContinuityRetention + if cfgPolicy.MinMessages > 0 { + policy.MinMessages = cfgPolicy.MinMessages + } + if cfgPolicy.MaxMessages > 0 { + policy.MaxMessages = cfgPolicy.MaxMessages + } + if cfgPolicy.TargetContextRatio > 0 && cfgPolicy.TargetContextRatio <= 0.5 { + policy.TargetContextRatio = cfgPolicy.TargetContextRatio + } + if cfgPolicy.FailureKeepMessages > 0 { + policy.FailureKeepMessages = cfgPolicy.FailureKeepMessages + } + } + + if policy.MaxMessages < policy.MinMessages { + policy.MaxMessages = policy.MinMessages + } + if policy.FailureKeepMessages < policy.MinMessages { + policy.FailureKeepMessages = policy.MinMessages + } + + return policy +} + +func (al *AgentLoop) continuityKeepCount(history []messages.Message) int { + if len(history) == 0 { + return 0 + } + + policy := al.continuityRetentionPolicy() + minKeep := policy.MinMessages + if minKeep > len(history) { + minKeep = len(history) + } + maxKeep := policy.MaxMessages + if maxKeep > len(history) { + maxKeep = len(history) + } + if maxKeep < minKeep { + maxKeep = minKeep + } + + contextWindow := al.contextWindow + if contextWindow <= 0 && al.cfg != nil { + contextWindow = al.cfg.Agents.Defaults.MaxTokens + } + if contextWindow <= 0 { + return minKeep + } + + targetTokens := int(float64(contextWindow) * policy.TargetContextRatio) + if targetTokens <= 0 { + return minKeep + } + + keep := 0 + keptTokens := 0 + for i := len(history) - 1; i >= 0 && keep < maxKeep; i-- { + msgTokens := observation.EstimateTokens(history[i].Content) + 4 + if keep >= minKeep && keptTokens+msgTokens > targetTokens { + break + } + keptTokens += msgTokens + keep++ + } + if keep < minKeep { + keep = minKeep + } + return keep +} + +type oversizedRecoveryCandidate struct { + Message messages.Message + OriginalIndex int + TokenEstimate int +} + +func (al *AgentLoop) persistOversizedRecoveryRefs(ctx context.Context, sessionKey string, omitted []oversizedRecoveryCandidate) ([]string, error) { + if len(omitted) == 0 { + return nil, nil + } + if al.memDelegate == nil { + return nil, fmt.Errorf("memory delegate is not configured") + } + + const maxPersistedRefs = 8 + refs := make([]string, 0, len(omitted)) + now := time.Now().UTC() + + for i, candidate := range omitted { + if i >= maxPersistedRefs { + break + } + nodeID := tools.DAGRecoveryNodePrefix + ids.New().String() + record := tools.DAGRecoveryRecord{ + NodeID: nodeID, + SessionKey: sessionKey, + OriginalIndex: candidate.OriginalIndex, + Role: candidate.Message.Role, + Content: candidate.Message.Content, + TokenEstimate: candidate.TokenEstimate, + Reason: "oversized_message_omitted_from_summary", + CreatedAt: now, + } + + data, err := json.Marshal(record) + if err != nil { + return refs, fmt.Errorf("marshal DAG recovery record: %w", err) + } + if err := al.memDelegate.UpsertKV(ctx, pkg.NAME, tools.DAGRecoveryKVKey(sessionKey, nodeID), string(data)); err != nil { + return refs, fmt.Errorf("persist DAG recovery record: %w", err) + } + refs = append(refs, nodeID) + } + return refs, nil +} + +// summarizeSession summarizes the conversation history for a session. +func (al *AgentLoop) summarizeSession(parentCtx context.Context, sessionKey string) { + ctx, cancel := context.WithTimeout(parentCtx, 120*time.Second) + defer cancel() + + history := al.sessions.GetHistory(sessionKey) + summary := al.sessions.GetSummary(sessionKey) + + keepLast := al.continuityKeepCount(history) + if len(history) <= keepLast { + return + } + + toSummarize := history[:len(history)-keepLast] + + // Oversized Message Guard: skip individual messages that would consume too + // much of the summarizer's context. Use 40% of the window for the summarizer + // input budget, reserving the rest for system prompt + summary output. + // Oversized omissions are persisted as DAG recovery references. + maxMessageTokens := al.contextWindow * 40 / 100 + if maxMessageTokens < 2048 { + maxMessageTokens = 2048 + } + validMessages := make([]messages.Message, 0) + omitted := false + omittedMessages := make([]oversizedRecoveryCandidate, 0) + + for idx, m := range toSummarize { + if m.Role != "user" && m.Role != "assistant" { + continue + } + msgTokens := observation.EstimateTokens(m.Content) + if msgTokens > maxMessageTokens { + omitted = true + omittedMessages = append(omittedMessages, oversizedRecoveryCandidate{ + Message: m, + OriginalIndex: idx, + TokenEstimate: msgTokens, + }) + continue + } + validMessages = append(validMessages, m) + } + + if len(validMessages) == 0 { + return + } + + // Multi-Part Summarization + var finalSummary string + if len(validMessages) > 10 { + mid := len(validMessages) / 2 + part1 := validMessages[:mid] + part2 := validMessages[mid:] + + s1, _ := al.summarizeBatch(ctx, part1, "") + s2, _ := al.summarizeBatch(ctx, part2, "") + + // Merge them + mergePrompt := fmt.Sprintf("Merge these two conversation summaries into one cohesive summary:\n\n1: %s\n\n2: %s", s1, s2) + resp, err := al.callModel(ctx, mergePrompt) + if err == nil { + finalSummary = resp + } else { + finalSummary = s1 + " " + s2 + } + } else { + finalSummary, _ = al.summarizeBatch(ctx, validMessages, summary) + } + + if omitted && finalSummary != "" { + recoveryRefs, err := al.persistOversizedRecoveryRefs(ctx, sessionKey, omittedMessages) + if err != nil { + logger.WarnCF("agent", "Failed to persist DAG recovery references for oversized messages", + map[string]interface{}{ + "session_key": sessionKey, + "error": err.Error(), + "omitted": len(omittedMessages), + }) + finalSummary += "\n[Note: Some oversized messages were omitted from this summary for efficiency.]" + } else if len(recoveryRefs) > 0 { + finalSummary += fmt.Sprintf("\n[Note: %d oversized message(s) were omitted from this summary. Recovery refs: %s. Use dag_expand with node_id= to recover full content.]", + len(omittedMessages), strings.Join(recoveryRefs, ", ")) + } else { + finalSummary += "\n[Note: Some oversized messages were omitted from this summary for efficiency.]" + } + } + + if finalSummary != "" { + al.sessions.SetSummary(sessionKey, finalSummary) + al.sessions.TruncateHistory(sessionKey, keepLast) + al.sessions.Save(sessionKey) + al.summarizeFailures.Delete(sessionKey) + } else { + var count int + if v, ok := al.summarizeFailures.Load(sessionKey); ok { + count = v.(int) + } + count++ + al.summarizeFailures.Store(sessionKey, count) + + const maxSummarizeFailures = 3 + emergencyKeep := al.continuityRetentionPolicy().FailureKeepMessages + if count >= maxSummarizeFailures { + logger.ErrorCF("agent", "Summarization failed repeatedly, force-truncating session", + map[string]interface{}{ + "session": sessionKey, + "consecutive_failures": count, + "keep": emergencyKeep, + }) + al.sessions.TruncateHistory(sessionKey, emergencyKeep) + al.sessions.Save(sessionKey) + al.summarizeFailures.Delete(sessionKey) + } + } +} + +// summarizeBatch summarizes a batch of messages using the Fantasy LanguageModel directly. +func (al *AgentLoop) summarizeBatch(ctx context.Context, batch []messages.Message, existingSummary string) (string, error) { + var prompt strings.Builder + prompt.WriteString("Provide a concise summary of this conversation segment, preserving core context and key points.\n") + if existingSummary != "" { + fmt.Fprintf(&prompt, "Existing context: %s\n", existingSummary) + } + prompt.WriteString("\nCONVERSATION:\n") + for _, m := range batch { + fmt.Fprintf(&prompt, "%s: %s\n", m.Role, m.Content) + } + + return al.callModel(ctx, prompt.String()) +} + +// callModel makes a direct call to the Fantasy LanguageModel (no tools, no agent loop). +// Used for summarization and other simple generation tasks. +func (al *AgentLoop) callModel(ctx context.Context, prompt string) (string, error) { + temp := 0.3 + maxTokens := int64(1024) + + resp, err := al.languageModel.Generate(ctx, fantasy.Call{ + Prompt: fantasy.Prompt{ + fantasy.NewUserMessage(prompt), + }, + Temperature: &temp, + MaxOutputTokens: &maxTokens, + }) + if err != nil { + return "", err + } + return resp.Content.Text(), nil +} + +// sessionsToMessagePairs converts the session history to observation.MessagePair +// for token estimation by the observation manager. +func (al *AgentLoop) sessionsToMessagePairs(sessionKey string) []observation.MessagePair { + history := al.sessions.GetHistory(sessionKey) + pairs := make([]observation.MessagePair, len(history)) + for i, m := range history { + pairs[i] = observation.MessagePair{Role: m.Role, Content: m.Content} + } + return pairs +} + +// applyDAGCompression compresses old history into a DAG summary block and +// returns only the tail messages that should be passed as raw conversation. +// The compressed portion is injected into the system prompt via contextBuilder. +// When memDelegate implements dag.DAGPersister, the DAG is persisted for dag_expand/describe/grep. +func (al *AgentLoop) applyDAGCompression(ctx context.Context, sessionKey string, history []messages.Message) []messages.Message { + const minHistoryForDAG = 16 + + if len(history) < minHistoryForDAG { + al.contextBuilder.SetDAGBlock("") + return history + } + + budget := dag.ComputeBudget(al.contextWindow, dag.DefaultBudgetConfig()) + tailCount := dag.TailMessageCount(budget.RawTail) + if tailCount >= len(history) { + al.contextBuilder.SetDAGBlock("") + return history + } + + // Split: compress old, keep tail raw + compressible := history[:len(history)-tailCount] + tail := history[len(history)-tailCount:] + + // Tool-call-aware: don't split on a "tool" message + for len(tail) > 0 && tail[0].Role == "tool" && len(compressible) > 0 { + tail = append([]messages.Message{compressible[len(compressible)-1]}, tail...) + compressible = compressible[:len(compressible)-1] + } + + if len(compressible) == 0 { + al.contextBuilder.SetDAGBlock("") + return history + } + + dagMsgs := make([]dag.Message, len(compressible)) + for i, m := range compressible { + dagMsgs[i] = dag.Message{Role: m.Role, Content: m.Content} + } + + compressor := dag.NewCompressor(dag.DefaultCompressorConfig()) + d := compressor.Compress(dagMsgs) + + rendered := dag.RenderDAGForBudget(d, budget.DAGSummaries) + al.contextBuilder.SetDAGBlock(rendered) + + // Persist DAG for dag_expand, dag_describe, dag_grep (additive; in-memory behavior unchanged) + if dp, ok := al.memDelegate.(dag.DAGPersister); ok { + if err := dp.PersistDAG(ctx, pkg.NAME, sessionKey, &dag.PersistSnapshot{ + FromMsgIdx: 0, + ToMsgIdx: len(compressible), + MsgCount: len(compressible), + DAG: d, + }); err != nil { + logger.WarnCF("agent", "DAG persist failed (non-fatal)", + map[string]interface{}{"error": err.Error(), "session_key": sessionKey}) + } + } + + logger.DebugCF("agent", "DAG compression applied", + map[string]interface{}{ + "total_msgs": len(history), + "compressed_msgs": len(compressible), + "tail_msgs": len(tail), + "dag_nodes": len(d.Nodes), + }) + + return tail +} + +func (al *AgentLoop) estimateTokens(msgs []messages.Message) int { + pairs := make([]observation.MessagePair, 0, len(msgs)) + for _, m := range msgs { + pairs = append(pairs, observation.MessagePair{ + Role: m.Role, + Content: m.Content, + }) + } + return observation.EstimateMessagesTokens(pairs) +} diff --git a/pkg/agent/toolloop.go b/pkg/agent/toolloop.go index d22203a52..48bd88f53 100644 --- a/pkg/agent/toolloop.go +++ b/pkg/agent/toolloop.go @@ -11,6 +11,7 @@ import ( "strings" fantasy "charm.land/fantasy" + "github.com/ZanzyTHEbar/dragonscale/pkg" picofantasy "github.com/ZanzyTHEbar/dragonscale/pkg/fantasy" "github.com/ZanzyTHEbar/dragonscale/pkg/logger" memstore "github.com/ZanzyTHEbar/dragonscale/pkg/memory/store" @@ -91,7 +92,7 @@ func runToolLoopWithRuntime( adaptCfg := picofantasy.AdaptedToolsConfig{ MemStore: ms, - AgentID: "dragonscale", + AgentID: pkg.NAME, SessionKey: sessionKey, } adaptedTools := picofantasy.BuildAdaptedTools(config.Tools, config.Bus, channel, chatID, adaptCfg)