diff --git a/CHANGELOG.md b/CHANGELOG.md index 024676b69..23c3a1a3c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,7 @@ ## Unreleased - Stop the `shell_command` prompt from asking for a description of the command, as the tool has no such param and calls with it now fail. +- Allow `spawn_agent` to continue a subagent conversation with optional `chat_id`, keeping its history and (unless overridden) model; each run gets fresh step/time budgets. #614 ## 0.163.0 diff --git a/docs/config/agents.md b/docs/config/agents.md index 3907456d9..d10099337 100644 --- a/docs/config/agents.md +++ b/docs/config/agents.md @@ -105,6 +105,8 @@ The major advantages of subagents are: - __Less context window usage__: Since subagents work as different chats/context/cleaner context, they have their own context window and when done the tools and process done there doesn't affect the primary agent context window, resulting and bigger conversations and less compaction needed. - __Parallel subagents__: subagents are spawned as tools, and ECA supports parallel tool calls if LLM supports, this increase speed of task solution if LLM needs for example to explore 2-3 different things with `explorer` subagent, spawning those in parallel. +Parents can continue a subagent conversation once its prior work has settled by passing its returned `chat_id` to `spawn_agent`, preserving its context. It keeps its model and variant unless `model` or `variant` is given; a new `model` gets the agent's configured variant. Continuation is limited to the same parent chat. Subagents from before an ECA restart cannot be continued, even if their IDs are still in the parent's saved history. Each continued run gets a fresh `maxSteps` and `timeoutSeconds` budget, and its limits, tools and approval rules follow the current config; the system prompt stays as it was unless `chat.autoSyncSystemPrompt` is enabled. A subagent halted by a limit can also be continued after its final summary. While it is continuable, a subagent chat can only be prompted by its parent, through `spawn_agent`. + Subagents can be configured in config or markdown and support/require these fields: - `mode`: set to `"subagent"` (or `["subagent"]`) to restrict an agent to subagent use only. Omit or include `"primary"` to also allow chat use. diff --git a/docs/protocol.md b/docs/protocol.md index 42894a7b8..7c123689a 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -1330,10 +1330,19 @@ interface SubagentDetails { /** * The chatId of this running subagent, useful to link other chat/ContentReceived * messages to this tool call. - * Available from toolCallRun afterwards + * Available from toolCallRun afterwards. + * A subagent continued with spawn_agent's `chat_id` keeps its chatId, so several + * tool calls of the same parent chat can share it. */ subagentChatId?: string; + /** + * The [start, end) indexes of the subagent chat messages that this tool call ran, + * [0, 0] when it did not run. Set on toolCalled. History replay uses it to show + * each call's part of a continued subagent. + */ + subagentMessageRange?: [number, number]; + /** * The model this subagent is using. */ diff --git a/integration-test/integration/chat/subagent_test.clj b/integration-test/integration/chat/subagent_test.clj index 687056da1..d4d45a9f8 100644 --- a/integration-test/integration/chat/subagent_test.clj +++ b/integration-test/integration/chat/subagent_test.clj @@ -115,10 +115,10 @@ :name "spawn_agent" :error false :outputs (m/embeds [{:type "text" - :text #"^## Agent 'explorer' Result"}])} + :text #"^Subagent chat_id: subagent-[^\n]+\n\n## Agent 'explorer' Result"}])} (:content e)))) events) - "Expected toolCalled for spawn_agent with output text starting with \"## Agent 'explorer' Result\""))) + "Expected toolCalled for spawn_agent with the reusable chat ID followed by the result heading"))) (testing "parent receives final assistant text after subagent completes" (is (some (fn [e] diff --git a/resources/prompts/tools/spawn_agent.md b/resources/prompts/tools/spawn_agent.md index 8517d9aef..a4b7bc25f 100644 --- a/resources/prompts/tools/spawn_agent.md +++ b/resources/prompts/tools/spawn_agent.md @@ -1,4 +1,4 @@ -Spawn an isolated sub-agent to handle complex, multi-step tasks without polluting your current context. +Spawn or continue an isolated sub-agent to handle complex, multi-step tasks without polluting your current context. Use for: Codebase exploration, codebase editing and refactoring, focused research, or delegating specialized tasks. Proactive use: If the specific agent's description suggests proactive use, use it whenever the task complexity justifies delegation. @@ -7,5 +7,6 @@ Agent Limits: Sub-agents cannot spawn other agents (no nesting) and have access Strict rules for arguments: - 'task': Provide a highly detailed prompt. Explicitly state whether it should write/edit code or just research, how to verify its work, and exactly what specific information it must return to you. -- 'activity': Must be a concise 3-4 word label for the UI (e.g., "exploring codebase", "refactoring module"). +- 'activity': Optional concise 3-4 word label for the UI (e.g., "exploring codebase", "refactoring module"). +- 'chat_id': Optional returned ID to continue a conversation in the same parent chat. Reuse its 'agent' and supply a new 'task'. It keeps its model and variant unless you override them. If continuation is rejected because the subagent is unavailable, omit 'chat_id' to spawn a new subagent and include the needed context in its 'task'. - 'model' & 'variant': - NEVER include these arguments if the user hasn't explicitly requested a specific model or variant. diff --git a/src/eca/db.clj b/src/eca/db.clj index 2af7996d2..3151fa6f9 100644 --- a/src/eca/db.clj +++ b/src/eca/db.clj @@ -49,6 +49,9 @@ ;; chat ids deleted in this session; excluded from workspace cache writes so ;; the merge-on-write never resurrects them from a shared cache file. :deleted-chat-ids #{:string} + ;; chats only their owner may prompt (spawn_agent for subagents); the owner + ;; passes its :token as :owner-token, :workers counts unwinding prompt workers. + :managed-chats {"" {:token ::object :workers :number :interrupted? :boolean}} :models {"" {:web-search :boolean :tools :boolean :reason? :boolean @@ -176,6 +179,8 @@ :tool-calls {} ;; Chat ids deleted in this session (not cached), see _db-spec. :deleted-chat-ids #{} + ;; Chats only their owner may prompt (not cached), see _db-spec. + :managed-chats {} ;; cacheable; bump `chats-version` when changing :chats shape, `version` ;; when changing :auth/:mcp-auth shape diff --git a/src/eca/features/chat.clj b/src/eca/features/chat.clj index eaef24ae8..404dd8644 100644 --- a/src/eca/features/chat.clj +++ b/src/eca/features/chat.clj @@ -19,6 +19,7 @@ [eca.features.rules :as f.rules] [eca.features.skills :as f.skills] [eca.features.tools :as f.tools] + [eca.features.tools.agent :as f.tools.agent] [eca.features.tools.mcp :as f.mcp] [eca.features.tools.task :as f.tools.task] [eca.llm-api :as llm-api] @@ -408,7 +409,8 @@ subagent-chat-id (when (= "tool_call_output" (:role message)) (get-in message [:content :details :subagent-chat-id])) subagent-messages (when subagent-chat-id - (get-in db [:chats subagent-chat-id :messages]))] + (f.tools.agent/replayed-messages (get-in db [:chats subagent-chat-id :messages]) + (:content message)))] (if (some? subagent-messages) ;; For subagent tool calls: toolCallRun + toolCallRunning, then ;; subagent messages, then toolCalled — matching live execution order. @@ -993,6 +995,18 @@ (string/trim) (as-> t (subs t 0 (min (count t) 40))))))) +(defn ^:private start-prompt-worker! + "Runs `thunk` in a prompt worker thread. For a managed chat (see + `:managed-chats` in the db), counts the worker from before dispatch until its + whole cleanup has unwound, since the chat may publish idle earlier." + [{:keys [db* config chat-id]} thunk] + (if-not (contains? (:managed-chats @db*) chat-id) + (future* config (thunk)) + (do (swap! db* update-in [:managed-chats chat-id :workers] inc) + (future* config + (try (thunk) + (finally (swap! db* update-in [:managed-chats chat-id :workers] dec))))))) + (defn ^:private prompt-messages! "Send user messages to LLM with hook processing. source-type controls hook agent. @@ -1172,7 +1186,8 @@ (if (and (lifecycle/auto-compact? chat-id agent full-model config @db*) (not (:auto-compacted? chat-ctx))) (trigger-auto-compact! chat-ctx all-tools user-messages) - (future* config + (start-prompt-worker! chat-ctx + (fn [] (try (llm-api/sync-or-async-prompt! {:model model @@ -1779,9 +1794,11 @@ ;; Only notify client if finish-chat-prompt! hasn't already run, ;; otherwise the belated statusChanged causes duplicate finished handling. (when-not (get-in @db* [:chats chat-id :prompt-finished?]) + (when (contains? (:managed-chats @db*) chat-id) + (swap! db* assoc-in [:managed-chats chat-id :interrupted?] true)) (messenger/chat-status-changed (:messenger chat-ctx) {:chat-id chat-id :status :idle}) (lifecycle/trigger-chat-status-hook! chat-ctx)) - (db/save-chat! @db* chat-id metrics)))))))))) + (db/save-chat! @db* chat-id metrics))))))))))) (defn ^:private send-mcp-prompt! [{:keys [prompt args] :as _decision} @@ -2112,12 +2129,21 @@ config should pass the map." [{:keys [message agent behavior chat-id contexts variant trust] :as params} db* messenger config metrics] (let [provided-chat-id chat-id - invalid-id-reason (when (and (some? provided-chat-id) - (not (server-managed-subagent-chat-id? @db* provided-chat-id))) + managed (get-in @db* [:managed-chats chat-id]) + invalid-id-reason (cond + ;; Only the run holding the owner token may prompt a managed chat. + (and managed + (not (and (:token managed) + (identical? (:token managed) (:owner-token params))))) + "this chat is managed (e.g. a subagent) and can only be prompted by its owner" + + (and (some? provided-chat-id) + (not (server-managed-subagent-chat-id? @db* provided-chat-id))) (validate-client-chat-id provided-chat-id))] (if invalid-id-reason - (do (logger/warn logger-tag "Rejected chat/prompt with invalid chat-id" - {:chat-id provided-chat-id :reason invalid-id-reason}) + (do (logger/with-chat-context provided-chat-id (db/parent-chat-id @db* provided-chat-id) + (logger/warn logger-tag "Rejected chat/prompt with invalid chat-id" + {:chat-id provided-chat-id :reason invalid-id-reason})) {:chat-id provided-chat-id :model "error" :status :error}) @@ -2464,7 +2490,10 @@ (when (identical? :running (get-in @db* [:chats chat-id :status])) ;; Set :stopping immediately to prevent race with stream callbacks ;; that check status via assert-chat-not-stopped! or cancelled? - (swap! db* assoc-in [:chats chat-id :status] :stopping) + (swap! db* (fn [db] + (cond-> (assoc-in db [:chats chat-id :status] :stopping) + (contains? (:managed-chats db) chat-id) + (assoc-in [:managed-chats chat-id :interrupted?] true)))) (let [chat-ctx {:chat-id chat-id :db* db* :config config diff --git a/src/eca/features/chat/tool_calls.clj b/src/eca/features/chat/tool_calls.clj index d052f31ab..d8eb34c22 100644 --- a/src/eca/features/chat/tool_calls.clj +++ b/src/eca/features/chat/tool_calls.clj @@ -229,6 +229,8 @@ Note: All actions are run in the order specified. Note: The :send-* actions should be last, so that they have the latest values of the state context. Note: The :status is updated before any actions are run, so the actions are in the context of the latest :status. + Note: :finally-actions run after the actions and the status hook, even when they throw. + The future-cleanup promise is delivered there, so a stop joins the post-tool hooks too. Note: all choices (i.e. conditionals) have to be made in code and result in different events being sent to the state machine. @@ -278,7 +280,8 @@ [:executing :execution-end] {:status :cleanup - :actions [:save-execution-result :deliver-future-cleanup-completed :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook]} + :actions [:save-execution-result :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook] + :finally-actions [:deliver-future-cleanup-completed]} [:cleanup :cleanup-finished] {:status :completed @@ -298,7 +301,8 @@ [:stopping :stop-attempted] {:status :cleanup - :actions [:save-execution-result :deliver-future-cleanup-completed :send-toolCallRejected :trigger-post-tool-call-hook]} + :actions [:save-execution-result :send-toolCallRejected :trigger-post-tool-call-hook] + :finally-actions [:deliver-future-cleanup-completed]} ;; And now all the :stop-requested transitions @@ -574,7 +578,8 @@ - event: Event keyword (e.g., :tool-prepare, :tool-run, :user-approve) - event-data: Optional map with event-specific data - Returns: {:status new-status :actions actions-executed} + Returns: the transition, {:status new-status :actions actions-executed} + plus its :finally-actions, if any. Throws: ex-info if the transition is invalid for the current state. @@ -584,7 +589,7 @@ (let [current-state (get-tool-call-state @db* (:chat-id chat-ctx) tool-call-id) current-status (:status current-state :initial) ; Default to :initial if no state transition-key [current-status event] - {:keys [status actions]} (get tool-call-state-machine transition-key)] + {:keys [status actions finally-actions] :as transition} (get tool-call-state-machine transition-key)] (logger/debug logger-tag "Tool call transition" {:tool-call-id tool-call-id :current-status current-status :event event :status status}) @@ -601,13 +606,15 @@ ;; Atomic status update (swap! db* assoc-in [:chats (:chat-id chat-ctx) :tool-calls tool-call-id :status] status) - ;; Execute all actions sequentially - (doseq [action actions] - (execute-action! action db* chat-ctx tool-call-id event-data)) + (try + (doseq [action actions] + (execute-action! action db* chat-ctx tool-call-id event-data)) + (lifecycle/trigger-chat-status-hook! (assoc chat-ctx :db* db*)) + (finally + (doseq [action finally-actions] + (execute-action! action db* chat-ctx tool-call-id event-data)))) - (lifecycle/trigger-chat-status-hook! (assoc chat-ctx :db* db*)) - - {:status status :actions actions})) + transition)) (def ^:private hook-approval-rank {"allow" 1 @@ -890,7 +897,7 @@ config messenger metrics - (partial get-tool-call-state @db* chat-id id) + #(get-tool-call-state @db* chat-id id) (partial transition-tool-call! db* chat-ctx id) {:trust (db/resolve-trust @db* chat-id)}) details (f.tools/tool-call-details-after-invocation name arguments details result @@ -980,7 +987,6 @@ (reduced nil)))) nil tool-calls) - (lifecycle/assert-chat-not-stopped! chat-ctx) (doseq [[tool-call-id state] (get-active-tool-calls @db* chat-id)] (when-let [f (:future state)] (try (deref f) @@ -1011,6 +1017,7 @@ :ex-data (ex-data t) :message (.getMessage ^Throwable t) :cause (.getCause ^Throwable t)}))))))) + (lifecycle/assert-chat-not-stopped! chat-ctx) (f.tools.mcp/await-pending-tools-refresh @db* 5000) ;; Token can expire during long tool calls (e.g. spawn_agent), ;; so renew before any continuation branch. diff --git a/src/eca/features/tools/agent.clj b/src/eca/features/tools/agent.clj index 6b074028a..6a8f25fff 100644 --- a/src/eca/features/tools/agent.clj +++ b/src/eca/features/tools/agent.clj @@ -1,5 +1,5 @@ (ns eca.features.tools.agent - "Tool for spawning subagents to perform focused tasks in isolated context." + "Tool for spawning or continuing subagents to perform focused tasks in isolated context." (:require [clojure.string :as str] [eca.config :as config] @@ -19,8 +19,17 @@ "How often the subagent chat status is checked while waiting for it." 1000) +(def ^:private settle-poll-ms + "How often a stopped or finishing subagent is checked while waiting for it to settle." + 50) + +(def ^:private stop-settle-timeout-ms + "Max time to wait for a subagent stopped with its parent to unwind. Short, as + the parent's stop waits for it, and only replay completeness depends on it." + (* 5 1000)) + (def ^:private summary-turn-timeout-ms - "Max time the subagent has to write its final summary after a timeout or max steps halt." + "Max time to wait for a halted subagent to settle and write its final summary." (* 2 60 1000)) (defn normalize-arguments @@ -57,9 +66,6 @@ [agent-name config parent-agent-name] (first (filter #(= agent-name (:name %)) (all-agents config parent-agent-name)))) -(defn ^:private max-steps [subagent] - (:max-steps subagent)) - (defn ^:private extract-final-assistant-text "Extracts text from the final assistant message, or nil when none exists." [messages] @@ -73,25 +79,21 @@ (str/join "\n") not-empty)) -(defn ^:private extract-final-summary - "Extract the final assistant message as summary from chat messages." - [messages] - (or (extract-final-assistant-text messages) - "Agent completed without producing output.")) - (defn ^:private failure-guidance "Actionable next-step hint for the parent agent based on the error type." [error-type] (when error-type (if (contains? llm-providers.errors/retryable-error-types error-type) - "This is a transient provider error. Prefer spawning this agent again for the same task (optionally with a different `model`) instead of performing the task yourself." - "Retrying this agent the same way is unlikely to help. Consider spawning it again with a different `model` or handling the task yourself."))) + "This is a transient provider error. Continue this agent with the returned `chat_id` and the same agent (optionally with a different `model`) instead of performing the task yourself." + "Retrying this agent the same way is unlikely to help. Consider continuing it with the returned `chat_id` and a different `model`, or handling the task yourself."))) + +(defn ^:private text-result [error? text] + {:error error? + :contents [{:type :text :text text}]}) (defn ^:private failed-agent-result [agent-name prompt-error partial-output] (let [{:keys [message error-type status code request-id response-id rate-limit-resets-at]} prompt-error] - {:error true - :contents [{:type :text - :text (str "## Agent '" agent-name "' Failed\n\n" + (text-result true (str "## Agent '" agent-name "' Failed\n\n" (or message "The sub-agent prompt failed.") (when error-type (str "\n\nError type: " (name error-type))) @@ -108,7 +110,7 @@ (when-let [guidance (failure-guidance error-type)] (str "\n\n" guidance)) (when partial-output - (str "\n\n## Partial result\n\n" partial-output)))}]})) + (str "\n\n## Partial result\n\n" partial-output)))))) (defn ^:private ->subagent-chat-id "Generate a deterministic subagent chat id from the tool-call-id." @@ -117,64 +119,95 @@ (defn ^:private send-step-progress! "Send a toolCallRunning notification with current step progress to the parent chat." - [messenger chat-id tool-call-id agent-name activity subagent-chat-id step max-steps model variant arguments] - (messenger/chat-content-received - messenger - {:chat-id chat-id - :role :assistant - :content {:type :toolCallRunning - :id tool-call-id - :name "spawn_agent" - :server "eca" - :origin "native" - :summary (if activity - (format "%s: %s" agent-name activity) - agent-name) - :arguments arguments - :details (cond-> {:type :subagent - :subagent-chat-id subagent-chat-id - :model model - :agent-name agent-name - :step step - :max-steps max-steps} - variant (assoc :variant variant))}})) + [{:keys [messenger chat-id tool-call-id agent-name subagent subagent-chat-id arguments]} db step] + (let [child (get-in db [:chats subagent-chat-id]) + activity (get arguments "activity")] + (messenger/chat-content-received + messenger + {:chat-id chat-id + :role :assistant + :content {:type :toolCallRunning + :id tool-call-id + :name "spawn_agent" + :server "eca" + :origin "native" + :summary (if activity + (format "%s: %s" agent-name activity) + agent-name) + :arguments arguments + :details (shared/assoc-some {:type :subagent + :subagent-chat-id subagent-chat-id + :model (:model child) + :agent-name agent-name + :step step + :max-steps (:max-steps subagent)} + :variant (:variant child))}}))) (defn ^:private stop-subagent-chat! "Stop a running subagent chat silently (parent already shows 'Prompt stopped')." - [db* messenger config metrics subagent-chat-id agent-name] + [{:keys [db* messenger config metrics subagent-chat-id agent-name]}] (let [prompt-stop (requiring-resolve 'eca.features.chat/prompt-stop)] (try (prompt-stop {:chat-id subagent-chat-id} db* messenger config metrics {:silent? true}) (catch Exception e (logger/warn logger-tag (format "Error stopping subagent '%s': %s" agent-name (.getMessage e))))))) +(defn ^:private prompt-subagent! + "Sends `message` to the subagent chat, claimed by this run's token." + [{:keys [db* messenger config metrics subagent-chat-id agent-name model variant trust token]} message] + ;; Resolved at runtime to avoid a circular dependency with the chat ns. + ((requiring-resolve 'eca.features.chat/prompt) + (shared/assoc-some {:owner-token token + :chat-id subagent-chat-id + :model model + :agent agent-name + :contexts [] + :trust trust + :message message} + :variant variant) + db* messenger config metrics)) + (defn ^:private summary-turn-prompt [reason] (str reason " Tool calls are no longer allowed.\n\n" "Without calling any tools, reply now with your final report: what you found so far, " "with concrete evidence like file paths, what you did not get to, and any open questions.")) +(defn ^:private settled? + "True when the subagent has no running turn and no prompt worker still + unwinding, so nothing else can write to its history." + [db subagent-chat-id] + (and (#{:idle :error} (get-in db [:chats subagent-chat-id :status])) + (zero? (get-in db [:managed-chats subagent-chat-id :workers] 0)))) + +(defn ^:private await-settled! + "Waits until the subagent is `settled?` or `deadline` (epoch ms) passes. + Returns true when it settled." + [db* subagent-chat-id deadline] + (loop [] + (cond + (settled? @db* subagent-chat-id) true + (< (System/currentTimeMillis) (long deadline)) (do (Thread/sleep (long settle-poll-ms)) + (recur)) + :else false))) + (defn ^:private run-summary-turn! "Prompts the subagent one last time, with tool calls refused, so it reports what - it found. Waits while that turn runs, stopping the subagent if it takes longer - than `summary-turn-timeout-ms`. `chat/prompt` marks the chat :running before - returning, so any other status means the turn is over or never started." - [{:keys [db* messenger config metrics chat-id subagent-chat-id agent-name prompt-params]} reason] + it found. A stopped turn may still be unwinding, so waits for the subagent to + settle first, and again after the summary turn, so it can be continued right + away. Both waits share `summary-turn-timeout-ms`; when it passes, the summary is + skipped or the subagent stopped." + [{:keys [db* chat-id subagent-chat-id agent-name] :as run} reason] (logger/with-chat-context subagent-chat-id chat-id - (logger/info logger-tag (format "Requesting final summary from agent '%s'" agent-name)) - (swap! db* assoc-in [:chats subagent-chat-id :summary-requested?] true) - (let [chat-prompt (requiring-resolve 'eca.features.chat/prompt) - deadline (+ (System/currentTimeMillis) (long summary-turn-timeout-ms))] - (chat-prompt (assoc prompt-params :message (summary-turn-prompt reason)) - db* messenger config metrics) - (loop [] - (when (= :running (get-in @db* [:chats subagent-chat-id :status])) - (if (< (System/currentTimeMillis) deadline) - (do - (Thread/sleep (long poll-interval-ms)) - (recur)) - (do - (logger/warn logger-tag (format "Agent '%s' did not finish its final summary in time, stopping it" agent-name)) - (stop-subagent-chat! db* messenger config metrics subagent-chat-id agent-name)))))))) + (let [deadline (+ (System/currentTimeMillis) (long summary-turn-timeout-ms))] + (if-not (await-settled! db* subagent-chat-id deadline) + (logger/warn logger-tag (format "Agent '%s' did not stop in time, skipping its final summary" agent-name)) + (do + (logger/info logger-tag (format "Requesting final summary from agent '%s'" agent-name)) + (swap! db* assoc-in [:chats subagent-chat-id :summary-requested?] true) + (prompt-subagent! run (summary-turn-prompt reason)) + (when-not (await-settled! db* subagent-chat-id deadline) + (logger/warn logger-tag (format "Agent '%s' did not finish its final summary in time, stopping it" agent-name)) + (stop-subagent-chat! run))))))) (defn ^:private available-model-names "Returns a sorted list of available model names from the runtime db." @@ -197,192 +230,293 @@ variants (config/effective-model-variants config provider model model-capabilities user-variants)] (config/selectable-variant-names variants)))))) +(defn ^:private validate-model! [db user-model] + (let [available-models (:models db)] + (when (and user-model + (seq available-models) + (not (contains? available-models user-model))) + (throw (ex-info (format "Model '%s' is not available. Available models: %s" + user-model + (str/join ", " (available-model-names db))) + {:model user-model + :available (available-model-names db)}))))) + +(defn ^:private validate-variant! + "Rejects a variant only when the model has configured variants and it isn't + among them. Models with no configured variants accept any variant (the LLM + API will reject if invalid)." + [config db model user-variant] + (when user-variant + (let [valid-variants (model-variant-names config db model)] + (when (and (seq valid-variants) + (not (some #{user-variant} valid-variants))) + (throw (ex-info (format "Variant '%s' is not available for model '%s'. Available variants: %s" + user-variant model (str/join ", " valid-variants)) + {:variant user-variant + :model model + :available valid-variants})))))) + +(defn ^:private new-model+variant + "[model variant] for a new subagent: explicit args first, then the agent config, + then the parent's model. The agent's :defaultModel may be a bare alias resolved + against the parent's provider; it is kept verbatim if it doesn't resolve." + [db subagent parent-chat-id user-model user-variant] + (let [parent-model (get-in db [:chats parent-chat-id :model]) + parent-provider (some-> parent-model shared/full-model->provider+model first)] + [(or user-model + (when-let [agent-model (:model subagent)] + (or (models/full-model-for db parent-provider agent-model) + agent-model)) + parent-model) + (or user-variant (:variant subagent))])) + +(defn ^:private continued-model+variant + "[model variant] for a continued subagent: explicit args first, else what it + used last. A new model gets the agent's configured variant instead of the + old one, which may not exist for it." + [child subagent user-model user-variant] + (let [model (or user-model (:model child))] + [model (or user-variant + (if (= model (:model child)) + (:variant child) + (:variant subagent)))])) + +(defn ^:private owned-subagent + "The chat of `subagent-chat-id` when `parent-chat-id` spawned it with + `agent-name` in this server session, else nil." + [db subagent-chat-id parent-chat-id agent-name] + (let [child (get-in db [:chats subagent-chat-id])] + (when (and (contains? (:managed-chats db) subagent-chat-id) + (= parent-chat-id (:parent-chat-id child)) + (= agent-name (:agent-name child))) + child))) + +(defn ^:private admit-new + "Registers the chat of a new subagent, claimed by this run's `token`." + [db {:keys [subagent-chat-id chat-id agent-name subagent model variant trust token]}] + (when (contains? (:chats db) subagent-chat-id) + (throw (ex-info "Subagent chat ID already exists." {:chat-id subagent-chat-id}))) + (-> db + (assoc-in [:managed-chats subagent-chat-id] {:token token :workers 0}) + (assoc-in [:chats subagent-chat-id] + (shared/assoc-some {:id subagent-chat-id + :parent-chat-id chat-id + :agent-name agent-name + :subagent subagent + :model model + :trust trust + :current-step 0} + :variant variant + :max-steps (:max-steps subagent))))) + +(defn ^:private admit-resume + "Claims a settled subagent of this parent for another run, with this run's + `token`. Its limits follow the current agent config, with a fresh budget." + [db {:keys [subagent-chat-id chat-id agent-name subagent token]}] + (let [child (owned-subagent db subagent-chat-id chat-id agent-name)] + (when-not (and child + (not (get-in db [:managed-chats subagent-chat-id :token])) + (settled? db subagent-chat-id) + (not-any? #(or (:future %) (seq (:resources %))) (vals (:tool-calls child)))) + ;; The parent LLM only knows its conversation, so say it in those terms. + (throw (ex-info (format (str "chat_id '%s' cannot be continued with agent '%s'. It must come from an earlier " + "spawn_agent result of this agent that has finished, and older subagents may " + "no longer be available. Omit chat_id to spawn a new subagent instead.") + subagent-chat-id agent-name) + {:chat-id subagent-chat-id})))) + (-> db + (update-in [:managed-chats subagent-chat-id] #(-> % (assoc :token token) (dissoc :interrupted?))) + ;; Clear how the previous run ended. + (update-in [:chats subagent-chat-id] #(-> (dissoc % :max-steps-reached? :summary-requested? :prompt-error + :prompt-finished? :follow-up-active?) + (assoc :subagent subagent + :max-steps (:max-steps subagent) + :current-step 0))))) + +(defn ^:private release + "Ends this run's claim on the subagent and records the part of its chat that + this call ran, for replay. A stopped subagent may still append a few messages + after this; they are not replayed." + [db {:keys [chat-id tool-call-id subagent-chat-id start-count]}] + (-> db + (update-in [:managed-chats subagent-chat-id] dissoc :token) + (assoc-in [:chats chat-id :tool-calls tool-call-id :subagent-message-range] + [start-count (count (get-in db [:chats subagent-chat-id :messages]))]))) + +(defn ^:private task-message + "The task as sent to the subagent." + [task max-steps after-summary?] + (cond-> task + ;; The summary turn told it that tools are no longer allowed. + after-summary? + (str "\n\nTool calls are allowed again.") + + max-steps + (str (format "\n\nIMPORTANT: You have a maximum of %d steps to complete this task. Be efficient and provide a clear summary of your findings before reaching the limit." + max-steps)))) + +(defn ^:private start-run! + "Prompts the subagent with its task. A prompt that fails before it starts + marks the subagent as failed, so the wait ends at once." + [{:keys [db* subagent-chat-id] :as run} message] + (when (= :error (:status (prompt-subagent! run message))) + (swap! db* update-in [:chats subagent-chat-id] + #(assoc % :status :error :prompt-error + (or (:prompt-error %) {:message "Subagent prompt setup failed."}))))) + +(defn ^:private run-output + "The final assistant text of this run only, so an earlier run's answer never + leaks into its result." + [{:keys [db* subagent-chat-id start-count]}] + (extract-final-assistant-text + (drop start-count (get-in @db* [:chats subagent-chat-id :messages] [])))) + +(defn ^:private stopped-result [{:keys [db* subagent-chat-id agent-name] :as run}] + (logger/info logger-tag (format "Agent '%s' stopped by parent chat" agent-name)) + (stop-subagent-chat! run) + ;; A cancelled tool still appends its result; wait for it, so this call's replay + ;; range includes it and the subagent can be continued right away. + (await-settled! db* subagent-chat-id (+ (System/currentTimeMillis) (long stop-settle-timeout-ms))) + (text-result true (str (format "Agent '%s' was stopped because the parent chat was stopped." agent-name) + (when-let [output (run-output run)] + (str "\n\n## Partial result\n\n" output))))) + +(defn ^:private timed-out-result [{:keys [agent-name subagent] :as run}] + (let [timeout-seconds (:timeout-seconds subagent)] + (logger/info logger-tag (format "Agent '%s' timed out after %ds" agent-name timeout-seconds)) + (stop-subagent-chat! run) + (run-summary-turn! run (format "You reached your time limit of %d seconds and your work was interrupted." timeout-seconds)) + (text-result true (format "## Agent '%s' Timed out\n\nAgent was stopped because it reached its timeout (%ds). The result below may be incomplete.\n\n%s" + agent-name timeout-seconds + (or (run-output run) "Agent produced no output before timing out."))))) + +(defn ^:private settled-result [{:keys [db* subagent-chat-id agent-name subagent] :as run} step] + (let [db @db* + {:keys [status prompt-error max-steps-reached?]} (get-in db [:chats subagent-chat-id]) + max-steps (:max-steps subagent) + output (run-output run) + failure (cond + (or prompt-error (= :error status)) (or prompt-error {}) + ;; Only interrupted, e.g. the user stopped the subagent chat. + (get-in db [:managed-chats subagent-chat-id :interrupted?]) + {:message "The sub-agent was stopped before it finished."})] + (cond + max-steps-reached? + (do + (logger/info logger-tag (format "Agent '%s' halted after reaching max steps (%d)" agent-name max-steps)) + (run-summary-turn! run (format "You reached your maximum number of steps (%d)." max-steps)) + (text-result true (format "## Agent '%s' Halted\n\nAgent was halted because it reached the maximum number of steps (%d). The result below may be incomplete.\n\n%s" + agent-name max-steps + (or (run-output run) "Agent completed without producing output.")))) + + failure + (do + (logger/warn logger-tag (format "Agent '%s' failed after %d steps: %s" agent-name step (:message failure))) + (failed-agent-result agent-name failure output)) + + :else + (do + (logger/info logger-tag (format "Agent '%s' completed after %d steps" agent-name step)) + (text-result false (format "## Agent '%s' Result\n\n%s" + agent-name (or output "Agent completed without producing output."))))))) + +(defn ^:private await-outcome + "Polls the subagent, sending step progress, until the parent stops, the + subagent settles or its running turn passes the deadline. + Returns [outcome last-step], outcome being :stopped, :settled or :timed-out." + [{:keys [db* subagent-chat-id call-state-fn deadline] :as run}] + (loop [last-step 0] + (let [db @db* + status (get-in db [:chats subagent-chat-id :status]) + step (get-in db [:chats subagent-chat-id :current-step] 0)] + (when (> step last-step) + (send-step-progress! run db step)) + (cond + (= :stopping (:status (call-state-fn))) [:stopped step] + (settled? db subagent-chat-id) [:settled step] + ;; A turn that already finished is only unwinding, so let it complete. + (and deadline + (= :running status) + (>= (System/currentTimeMillis) (long deadline))) [:timed-out step] + :else (do (Thread/sleep (long poll-interval-ms)) + (recur (long (max last-step step)))))))) + +(defn ^:private await-result [{:keys [db* chat-id tool-call-id] :as run}] + (try + (let [[outcome step] (await-outcome run)] + (when-not (= :stopped outcome) + (swap! db* assoc-in [:chats chat-id :tool-calls tool-call-id :subagent-final-step] step)) + (case outcome + :stopped (stopped-result run) + :settled (settled-result run step) + :timed-out (timed-out-result run))) + (catch InterruptedException _ + (stopped-result run)))) + (defn ^:private spawn-agent "Handler for the spawn_agent tool. - Spawns a subagent to perform a focused task and returns the result." - [arguments {:keys [db* config messenger metrics chat-id tool-call-id call-state-fn trust agent]}] + Runs a focused task in a new or existing subagent conversation and returns the result." + [arguments {:keys [db* config chat-id tool-call-id call-state-fn agent] :as ctx}] (let [arguments (normalize-arguments arguments) - agent-name (get arguments "agent") - task (get arguments "task") - activity (get arguments "activity") - parent-agent-name agent + {agent-name "agent" task "task" resume-id "chat_id" + user-model "model" user-variant "variant"} arguments db @db* - - ;; Check for nesting - prevent subagents from spawning other subagents _ (when (get-in db [:chats chat-id :subagent]) (throw (ex-info "Agents cannot spawn other agents (nesting not allowed)" {:agent-name agent-name :parent-chat-id chat-id}))) - - subagent (get-agent agent-name config parent-agent-name) - _ (when-not subagent - (let [available (all-agents config parent-agent-name)] - (throw (ex-info (format "Agent not found or not available. Available agents: %s" - (if (seq available) - (str/join ", " (map :name available)) - "none")) - {:agent-name agent-name - :available (map :name available)})))) - - ;; Create subagent chat session using deterministic id based on tool-call-id - subagent-chat-id (->subagent-chat-id tool-call-id) - - user-model (get arguments "model") - _ (when user-model - (let [available-models (:models db)] - (when (and (seq available-models) - (not (contains? available-models user-model))) - (throw (ex-info (format "Model '%s' is not available. Available models: %s" - user-model - (str/join ", " (available-model-names db))) - {:model user-model - :available (available-model-names db)}))))) - - parent-model (get-in db [:chats chat-id :model]) - parent-provider (some-> parent-model shared/full-model->provider+model first) - ;; The agent's :defaultModel may be a bare alias resolved against the - ;; currently selected (parent) provider; keep it verbatim if it doesn't resolve. - subagent-model (or user-model - (when-let [agent-model (:model subagent)] - (or (models/full-model-for db parent-provider agent-model) - agent-model)) - parent-model) - - ;; Variant validation: reject only when the resolved model has configured - ;; variants and the user-specified one isn't among them. Models with no - ;; configured variants accept any variant (the LLM API will reject if invalid). - user-variant (get arguments "variant") - _ (when user-variant - (let [valid-variants (model-variant-names config db subagent-model)] - (when (and (seq valid-variants) - (not (some #{user-variant} valid-variants))) - (throw (ex-info (format "Variant '%s' is not available for model '%s'. Available variants: %s" - user-variant subagent-model (str/join ", " valid-variants)) - {:variant user-variant - :model subagent-model - :available valid-variants}))))) - variant (or user-variant (:variant subagent))] - - (logger/info logger-tag (format "Spawning agent '%s' for task: %s (model: %s, variant: %s)" agent-name task subagent-model (or variant "default"))) - - (let [max-steps-limit (max-steps subagent) - timeout-seconds (:timeout-seconds subagent) - deadline (when timeout-seconds - (+ (System/currentTimeMillis) (* 1000 (long timeout-seconds)))) - prompt-params (cond-> {:chat-id subagent-chat-id - :model subagent-model - :agent agent-name - :contexts [] - :trust trust} - variant (assoc :variant variant)) - summary-ctx {:db* db* - :messenger messenger - :config config - :metrics metrics - :chat-id chat-id - :subagent-chat-id subagent-chat-id - :agent-name agent-name - :prompt-params prompt-params} - subagent-messages #(get-in @db* [:chats subagent-chat-id :messages] [])] - (swap! db* assoc-in [:chats subagent-chat-id] - (cond-> {:id subagent-chat-id - :parent-chat-id chat-id - :agent-name agent-name - :subagent subagent - :current-step 0} - max-steps-limit (assoc :max-steps max-steps-limit))) - - (try - ;; Require chat ns here to avoid circular dependency - (let [chat-prompt (requiring-resolve 'eca.features.chat/prompt) - task-prompt (if max-steps-limit - (format "%s\n\nIMPORTANT: You have a maximum of %d steps to complete this task. Be efficient and provide a clear summary of your findings before reaching the limit." - task max-steps-limit) - task)] - (chat-prompt (assoc prompt-params :message task-prompt) db* messenger config metrics)) - - ;; Wait for subagent to complete by polling status - (let [stopped-result (fn [] - (logger/info logger-tag (format "Agent '%s' stopped by parent chat" agent-name)) - (stop-subagent-chat! db* messenger config metrics subagent-chat-id agent-name) - {:error true - :contents [{:type :text - :text (str (format "Agent '%s' was stopped because the parent chat was stopped." agent-name) - (when-let [partial-output (extract-final-assistant-text (subagent-messages))] - (str "\n\n## Partial result\n\n" partial-output)))}]})] - (try - (loop [last-step 0] - (let [db @db* - status (get-in db [:chats subagent-chat-id :status]) - current-step (get-in db [:chats subagent-chat-id :current-step] 0)] - ;; Send step progress when step advances - (when (> current-step last-step) - (send-step-progress! messenger chat-id tool-call-id agent-name activity - subagent-chat-id current-step max-steps-limit subagent-model variant arguments)) - (cond - ;; Parent chat stopped — propagate stop to subagent - (= :stopping (:status (call-state-fn))) - (stopped-result) - - ;; Subagent completed - (#{:idle :error} status) - (let [messages (get-in db [:chats subagent-chat-id :messages] []) - summary (extract-final-summary messages) - partial-output (extract-final-assistant-text messages) - prompt-error (get-in db [:chats subagent-chat-id :prompt-error]) - failed? (boolean (or (= :error status) prompt-error)) - max-steps-reached? (get-in db [:chats subagent-chat-id :max-steps-reached?])] - (cond - max-steps-reached? - (logger/info logger-tag (format "Agent '%s' halted after reaching max steps (%d)" agent-name max-steps-limit)) - - failed? - (logger/warn logger-tag (format "Agent '%s' failed after %d steps: %s" - agent-name current-step (:message prompt-error))) - - :else - (logger/info logger-tag (format "Agent '%s' completed after %d steps" agent-name current-step))) - (swap! db* assoc-in [:chats chat-id :tool-calls tool-call-id :subagent-final-step] current-step) - (cond - max-steps-reached? - (do - (run-summary-turn! summary-ctx (format "You reached your maximum number of steps (%d)." max-steps-limit)) - {:error true - :contents [{:type :text - :text (format "## Agent '%s' Halted\n\nAgent was halted because it reached the maximum number of steps (%d). The result below may be incomplete.\n\n%s" - agent-name max-steps-limit (extract-final-summary (subagent-messages)))}]}) - - failed? - (failed-agent-result agent-name prompt-error partial-output) - - :else - {:error false - :contents [{:type :text - :text (format "## Agent '%s' Result\n\n%s" agent-name summary)}]})) - - ;; Subagent ran past its timeout, stop it and ask for a final summary - (and deadline (>= (System/currentTimeMillis) (long deadline))) - (do - (logger/info logger-tag (format "Agent '%s' timed out after %ds" agent-name timeout-seconds)) - (stop-subagent-chat! db* messenger config metrics subagent-chat-id agent-name) - (run-summary-turn! summary-ctx (format "You reached your time limit of %d seconds and your work was interrupted." timeout-seconds)) - (swap! db* assoc-in [:chats chat-id :tool-calls tool-call-id :subagent-final-step] current-step) - {:error true - :contents [{:type :text - :text (format "## Agent '%s' Timed out\n\nAgent was stopped because it reached its timeout (%ds). The result below may be incomplete.\n\n%s" - agent-name timeout-seconds - (or (extract-final-assistant-text (subagent-messages)) - "Agent produced no output before timing out."))}]}) - - ;; Keep waiting - :else - (do - (Thread/sleep (long poll-interval-ms)) - (recur (long (max last-step current-step))))))) - (catch InterruptedException _ - (stopped-result)))) - (catch Exception e - (throw e)))))) + ;; `agent` in the context is the parent's agent. + subagent (or (get-agent agent-name config agent) + (let [available (map :name (all-agents config agent))] + (throw (ex-info (format "Agent not found or not available. Available agents: %s" + (if (seq available) (str/join ", " available) "none")) + {:agent-name agent-name + :available available})))) + subagent-chat-id (or resume-id (->subagent-chat-id tool-call-id)) + child (get-in db [:chats subagent-chat-id]) + _ (validate-model! db user-model) + [model variant] (if resume-id + (continued-model+variant child subagent user-model user-variant) + (new-model+variant db subagent chat-id user-model user-variant)) + _ (validate-variant! config db model user-variant) + run (assoc ctx + :arguments arguments + :agent-name agent-name + :subagent subagent + :subagent-chat-id subagent-chat-id + :model model + :variant variant + :token (Object.) + :deadline (when-let [seconds (:timeout-seconds subagent)] + (+ (System/currentTimeMillis) (* 1000 (long seconds))))) + admitted-db (swap! db* (if resume-id admit-resume admit-new) run) + run (assoc run :start-count (count (get-in admitted-db [:chats subagent-chat-id :messages])))] + (logger/with-chat-context chat-id (get-in db [:chats chat-id :parent-chat-id]) + (-> (try + (logger/info logger-tag (format "Running agent '%s' for task: %s (model: %s, variant: %s)" + agent-name task model (or variant "default"))) + ;; `child` is read before admission, which clears the summary flag. + (start-run! run (task-message task (:max-steps subagent) (:summary-requested? child))) + (await-result run) + (catch Exception e + (when (or (instance? InterruptedException e) + (= :stopping (:status (call-state-fn)))) + (stop-subagent-chat! run)) + (failed-agent-result agent-name {:message (ex-message e)} nil)) + (finally + (swap! db* release run))) + ;; The chat_id goes first, so output truncation keeps it. + (update-in [:contents 0 :text] #(str "Subagent chat_id: " subagent-chat-id "\n\n" %)))))) + +(defn replayed-messages + "The part of a subagent chat's `messages` that one spawn_agent call ran, from + the call's tool-call `details`. A continued subagent runs over several calls, + so each call replays only its own part. Histories from before continuation + have no range and replay all." + [messages {:keys [details]}] + (if-let [[start end] (:subagent-message-range details)] + (take (- end start) (drop start messages)) + messages)) (defn ^:private build-description "Build tool description with available agents and models listed." @@ -404,11 +538,13 @@ {:description (build-description config parent-agent-name) :parameters {:type "object" :properties {"agent" {:type "string" - :description "Name of the agent to spawn"} + :description "Name of the agent to spawn or continue"} "task" {:type "string" :description "The detailed instructions for the agent"} "activity" {:type "string" - :description "Concise label (max 3-4 words) shown in the UI while the agent runs, e.g. \"exploring codebase\", \"reviewing changes\", \"analyzing tests\"."} + :description "Optional concise label (max 3-4 words) shown in the UI while the agent runs, e.g. \"exploring codebase\", \"reviewing changes\", \"analyzing tests\"."} + "chat_id" {:type "string" + :description "Optional chat_id returned by an earlier spawn_agent call of this chat, to continue that subagent conversation. Repeat its agent. It keeps its model and variant unless you override them."} "model" {:type "string" :description "Optional sub-agent model override. Reserved for explicit user override only. Omit unless the user explicitly named a model."} "variant" {:type "string" @@ -425,26 +561,32 @@ (defmethod tools.util/tool-call-details-before-invocation :spawn_agent [_name arguments _server {:keys [db config chat-id tool-call-id]}] (let [agent-name (get arguments "agent") - user-model (get arguments "model") - user-variant (get arguments "variant") - parent-agent-name (get-in db [:chats chat-id :agent]) subagent (when agent-name - (get-agent agent-name config parent-agent-name)) - parent-model (get-in db [:chats chat-id :model]) - subagent-model (or user-model (:model subagent) parent-model) - variant (or user-variant (:variant subagent)) - subagent-chat-id (when tool-call-id - (->subagent-chat-id tool-call-id))] - (cond-> {:type :subagent - :subagent-chat-id subagent-chat-id - :model subagent-model - :agent-name agent-name - :step (get-in db [:chats subagent-chat-id :current-step] 1) - :max-steps (max-steps subagent)} - variant (assoc :variant variant)))) + (get-agent agent-name config (get-in db [:chats chat-id :agent]))) + resume-id (get arguments "chat_id") + resume? (some? resume-id) + child (when (and resume? subagent) + (owned-subagent db resume-id chat-id agent-name)) + [model variant] (if resume? + (continued-model+variant child subagent (get arguments "model") (get arguments "variant")) + (new-model+variant db subagent chat-id (get arguments "model") (get arguments "variant")))] + (shared/assoc-some {:type :subagent + :subagent-chat-id (if resume? + (:id child) + (some-> tool-call-id ->subagent-chat-id)) + :model model + :agent-name agent-name + :step 1 + :max-steps (:max-steps subagent)} + :variant variant))) (defmethod tools.util/tool-call-details-after-invocation :spawn_agent [_name _arguments before-details _result {:keys [db chat-id tool-call-id]}] - (let [final-step (get-in db [:chats chat-id :tool-calls tool-call-id :subagent-final-step] - (or (:step before-details) 1))] - (assoc before-details :step final-step))) + (let [{:keys [subagent-final-step subagent-message-range]} (get-in db [:chats chat-id :tool-calls tool-call-id]) + child (get-in db [:chats (:subagent-chat-id before-details)])] + (cond-> (assoc before-details + :step (or subagent-final-step (:step before-details) 1) + ;; Replay shows only this part of the subagent chat; none when the call did not run. + :subagent-message-range (or subagent-message-range [0 0])) + ;; The subagent ran, so report the model it really used. + subagent-final-step (shared/assoc-some :model (:model child) :variant (:variant child))))) diff --git a/test/eca/features/chat/tool_calls_test.clj b/test/eca/features/chat/tool_calls_test.clj index fb426fd9d..72284825f 100644 --- a/test/eca/features/chat/tool_calls_test.clj +++ b/test/eca/features/chat/tool_calls_test.clj @@ -819,6 +819,27 @@ :post-tool-call-stop-hook-name "guard"} (get-in @db* [:chats "chat-1" :tool-calls "tool-1"])))))) +(deftest cleanup-signal-follows-hooks-test + (doseq [[status event] [[:executing :execution-end] [:stopping :stop-attempted]] + failure [nil :action :status]] + (let [done (promise) + db* (atom {:chats {"child" {:tool-calls {"tool" {:status status + :future-cleanup-complete?* done}}}}}) + execute @#'tc/execute-action!] + (with-redefs [tc/execute-action! + (fn [action & args] + (if (= :deliver-future-cleanup-completed action) + (apply execute action args) + (do (is (not (realized? done)) "Actions must finish before the join is released") + (when (= :action failure) (throw (ex-info "action failed" {})))))) + lifecycle/trigger-chat-status-hook! + (fn [_] + (is (not (realized? done)) "Status hook is part of tool-side work") + (when (= :status failure) (throw (ex-info "status failed" {}))))] + (try (tc/transition-tool-call! db* {:chat-id "child"} "tool" event {}) + (catch clojure.lang.ExceptionInfo e (is failure (ex-message e)))) + (is (realized? done) "Exceptions must also release the join"))))) + (deftest rejected-tool-call-output-contents-test (testing "states the call did not run and made no changes (#507)" (let [text (-> (#'tc/rejected-tool-call-output-contents "Tool call rejected by user choice") diff --git a/test/eca/features/chat_tool_call_state_test.clj b/test/eca/features/chat_tool_call_state_test.clj index 5e21da890..c974d2d0d 100644 --- a/test/eca/features/chat_tool_call_state_test.clj +++ b/test/eca/features/chat_tool_call_state_test.clj @@ -475,7 +475,8 @@ result (#'tc/transition-tool-call! db* chat-ctx tool-call-id :execution-end result-data)] (is (match? {:status :cleanup - :actions [:save-execution-result :deliver-future-cleanup-completed :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook]} + :actions [:save-execution-result :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook] + :finally-actions [:deliver-future-cleanup-completed]} result) "Expected transition to :cleanup with send toolCalled and record metrics actions") @@ -546,7 +547,8 @@ "Expected transition from :executing to :stopping with relevant actions")) (let [result (#'tc/transition-tool-call! db* chat-ctx "tool-executing" :stop-attempted)] (is (match? {:status :cleanup - :actions [:save-execution-result :deliver-future-cleanup-completed :send-toolCallRejected :trigger-post-tool-call-hook]} + :actions [:save-execution-result :send-toolCallRejected :trigger-post-tool-call-hook] + :finally-actions [:deliver-future-cleanup-completed]} result) "Expected transition from :stopping to :cleanup with relevant actions")))) @@ -1025,7 +1027,8 @@ result (#'tc/transition-tool-call! db* chat-ctx tool-call-id :execution-end error-result)] (is (match? {:status :cleanup - :actions [:save-execution-result :deliver-future-cleanup-completed :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook]} + :actions [:save-execution-result :send-toolCalled :log-metrics :send-progress :trigger-post-tool-call-hook] + :finally-actions [:deliver-future-cleanup-completed]} result) "Expected transition to :cleanup with send toolCalled and record metrics actions") diff --git a/test/eca/features/tools/agent_test.clj b/test/eca/features/tools/agent_test.clj index 8e260b242..fb0d9a635 100644 --- a/test/eca/features/tools/agent_test.clj +++ b/test/eca/features/tools/agent_test.clj @@ -4,6 +4,8 @@ [clojure.test :refer [deftest is testing]] [eca.config :as config] [eca.features.chat :as f.chat] + [eca.features.chat.tool-calls :as tool-calls] + [eca.features.hooks :as hooks] [eca.features.tools :as f.tools] [eca.features.tools.agent :as f.tools.agent] [eca.features.tools.util :as tools.util] @@ -349,7 +351,7 @@ :call-state-fn (constantly {:status :executing})})] (is (match? {:error true :contents [{:type :text - :text #"(?s)Failed.*rate limit.*Error type: rate-limited.*Status: 429.*Code: rate_limit_error.*Rate limit resets at: 2025-08-26T10:40:00Z.*transient provider error\. Prefer spawning this agent again"}]} + :text #"(?s)Failed.*rate limit.*Error type: rate-limited.*Status: 429.*Code: rate_limit_error.*Rate limit resets at: 2025-08-26T10:40:00Z.*transient provider error\. Continue this agent with the returned `chat_id`"}]} result)))))) (testing "non-retryable error advises against retrying the same way" @@ -534,7 +536,8 @@ eca.features.chat/prompt-stop (fn [_params _db* _messenger _config _metrics _opts] (swap! stops* inc) - (swap! db* assoc-in [:chats subagent-chat-id :status] :stopping)) + ;; The stopped turn unwinds and settles as idle. + (swap! db* assoc-in [:chats subagent-chat-id :status] :idle)) (clojure.lang.RT/var (namespace sym) (name sym))))] (let [result ((spawn-handler) {"agent" "slow" "task" "review" "activity" "reviewing"} @@ -577,7 +580,8 @@ eca.features.chat/prompt-stop (fn [_params _db* _messenger _config _metrics _opts] (swap! stops* inc) - (swap! db* assoc-in [:chats subagent-chat-id :status] :stopping)) + ;; The stopped turn unwinds and settles as idle. + (swap! db* assoc-in [:chats subagent-chat-id :status] :idle)) (clojure.lang.RT/var (namespace sym) (name sym))))] (let [result ((spawn-handler) {"agent" "slow" "task" "review" "activity" "reviewing"} @@ -595,6 +599,95 @@ (is (= 2 (count @prompts*))) (is (= 2 @stops*))))))) +(defn ^:private spawn-context [db* & {:as overrides}] + (merge {:db* db* :config test-config :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "tc-1" + :call-state-fn (constantly {:status :executing})} + overrides)) + +(defn ^:private result-text [result] + (tools.util/contents->text (:contents result))) + +(deftest spawn-agent-timeout-summary-waits-for-settle-test + (let [db* (atom {:chats {"parent" {:model "test/model"}}}) + id "subagent-tc-1" + context (spawn-context db* :config timeout-test-config) + running! (fn [& _] (swap! db* update-in [:chats id] assoc :status :running :messages [(assistant-msg "Working.")]))] + (testing "the summary turn starts only after the stopped turn has unwound" + (let [status-at-summary* (promise)] + (with-redefs [f.tools.agent/poll-interval-ms 10 + f.chat/prompt (fn [{:keys [message]} & _] + (if (string/includes? message "Without calling any tools") + (do (deliver status-at-summary* (get-in @db* [:chats id :status])) + (swap! db* update-in [:chats id] + #(-> % (assoc :status :idle) + (update :messages conj (assistant-msg "Final report."))))) + (running!))) + f.chat/prompt-stop (fn [& _] + (swap! db* assoc-in [:chats id :status] :stopping) + (future (Thread/sleep 100) + (swap! db* assoc-in [:chats id :status] :idle)))] + (let [result ((spawn-handler) {"agent" "slow" "task" "review"} context)] + (is (= :idle (deref status-at-summary* 0 ::not-called))) + (is (re-find #"(?s)Timed out.*Final report\." (result-text result))))))) + + (testing "the summary turn is skipped when the stopped turn does not unwind in time" + (reset! db* {:chats {"parent" {:model "test/model"}}}) + (let [prompts* (atom 0)] + (with-redefs [f.tools.agent/poll-interval-ms 10 + f.tools.agent/summary-turn-timeout-ms 50 + f.chat/prompt (fn [& args] (swap! prompts* inc) (apply running! args)) + f.chat/prompt-stop (fn [& _] (swap! db* assoc-in [:chats id :status] :stopping))] + (let [result ((spawn-handler) {"agent" "slow" "task" "review"} context)] + (is (= 1 @prompts*)) + (is (re-find #"(?s)Timed out.*Working\." (result-text result))))))))) + +(deftest spawn-agent-finished-before-timeout-test + (testing "a turn that finished but is still unwinding at the deadline completes instead of timing out" + (let [db* (atom {:chats {"parent" {:model "test/model"}}}) + id "subagent-tc-1" + stops* (atom 0)] + (with-redefs [f.tools.agent/poll-interval-ms 10 + f.chat/prompt (fn [& _] + (swap! db* #(-> % + (update-in [:chats id] assoc :status :idle :messages [(assistant-msg "Done.")]) + (assoc-in [:managed-chats id :workers] 1))) + (future (Thread/sleep 1300) + (swap! db* assoc-in [:managed-chats id :workers] 0))) + f.chat/prompt-stop (fn [& _] (swap! stops* inc))] + (let [result ((spawn-handler) {"agent" "slow" "task" "review"} (spawn-context db* :config timeout-test-config))] + (is (zero? @stops*)) + (is (match? {:error false :contents [{:text #"(?s)Result\n\nDone\.$"}]} result))))))) + +(deftest spawn-agent-continued-run-texts-are-run-local-test + (testing "halted, timed out and stopped texts never show an earlier run's answer" + (let [db* (atom {:chats {"parent" {:model "test/model"}}}) + id "subagent-tc-1" + mode* (atom :answer) + run-agent (fn [args tool-call-id parent-status] + (result-text ((spawn-handler) args + (spawn-context db* + :config (assoc-in timeout-test-config [:agent "slow" :maxSteps] 3) + :tool-call-id tool-call-id + :call-state-fn (constantly {:status parent-status})))))] + (with-redefs [f.tools.agent/poll-interval-ms 10 + f.tools.agent/summary-turn-timeout-ms 50 + f.chat/prompt (fn [& _] + (case @mode* + :answer (swap! db* update-in [:chats id] assoc + :status :idle :messages [(assistant-msg "Old answer.")]) + :halt (swap! db* update-in [:chats id] assoc :status :idle :max-steps-reached? true) + :hang (swap! db* assoc-in [:chats id :status] :running))) + f.chat/prompt-stop (fn [& _] (swap! db* assoc-in [:chats id :status] :idle))] + (is (string/includes? (run-agent {"agent" "slow" "task" "first"} "tc-1" :executing) "Old answer.")) + (doseq [[mode tool-call-id parent-status expected] [[:halt "tc-2" :executing #"(?s)Halted.*Agent completed without producing output\."] + [:hang "tc-3" :executing #"(?s)Timed out.*Agent produced no output before timing out\."] + [:hang "tc-4" :stopping #"was stopped"]]] + (reset! mode* mode) + (let [text (run-agent {"agent" "slow" "task" "next" "chat_id" id} tool-call-id parent-status)] + (is (re-find expected text)) + (is (not (string/includes? text "Old answer."))))))))) + (deftest spawn-agent-parent-stop-test (testing "stops subagent when parent chat is stopped" (let [db* (atom {:chats {"chat-1" {:id "chat-1" :model "test/model"}}}) @@ -633,6 +726,20 @@ (testing "preserves subagent chat for resume replay" (is (some? (get-in @db* [:chats subagent-chat-id]))))))))) +(deftest spawn-agent-setup-cancellation-test + (doseq [failure [(InterruptedException.) (ex-info "setup cancelled" {})]] + (let [db* (atom {:chats {"parent" {:model "test/model"}}}) + stopped* (atom false)] + (with-redefs [f.chat/prompt (fn [& _] + (swap! db* assoc-in [:chats "subagent-setup" :status] :running) + (throw failure)) + f.chat/prompt-stop (fn [& _] (reset! stopped* true))] + (is (true? (:error ((spawn-handler) {"agent" "explorer" "task" "work"} + {:db* db* :config test-config :chat-id "parent" :tool-call-id "setup" + :call-state-fn (constantly {:status (if (instance? InterruptedException failure) + :executing :stopping)})})))) + (is @stopped* "Cancellation during synchronous setup must stop the dispatched child"))))) + (deftest spawn-agent-cleanup-on-exception-test (testing "preserves subagent state when chat/prompt throws" (let [db* (atom {:chats {"chat-1" {:id "chat-1" :model "test/model"}}}) @@ -644,7 +751,7 @@ (fn [_params _db* _messenger _config _metrics] (throw (ex-info "LLM provider error" {}))) (clojure.lang.RT/var (namespace sym) (name sym))))] - (is (thrown? Exception + (is (match? {:error true :contents [{:text #"(?s)^Subagent chat_id: subagent-tc-1\n\n.*Failed.*LLM provider error"}]} ((spawn-handler) {"agent" "explorer" "task" "explore" "activity" "exploring"} {:db* db* @@ -1022,36 +1129,35 @@ (is (= "company-litellm/explorer-small" (:model @chat-prompt-called*)) "bare alias should resolve to the parent provider's model")))))) -(deftest extract-final-summary-test +(deftest extract-final-assistant-text-test (testing "extracts text from last assistant message" (is (= "Hello world" - (#'f.tools.agent/extract-final-summary + (#'f.tools.agent/extract-final-assistant-text [{:role "user" :content [{:type :text :text "Hi"}]} {:role "assistant" :content [{:type :text :text "Hello world"}]}])))) (testing "uses last assistant message when multiple exist" (is (= "Final answer" - (#'f.tools.agent/extract-final-summary + (#'f.tools.agent/extract-final-assistant-text [{:role "assistant" :content [{:type :text :text "First response"}]} {:role "user" :content [{:type :text :text "More?"}]} {:role "assistant" :content [{:type :text :text "Final answer"}]}])))) (testing "joins multiple text blocks with newline" (is (= "Part 1\nPart 2" - (#'f.tools.agent/extract-final-summary + (#'f.tools.agent/extract-final-assistant-text [{:role "assistant" :content [{:type :text :text "Part 1"} {:type :text :text "Part 2"}]}])))) (testing "ignores non-text content types" (is (= "Text only" - (#'f.tools.agent/extract-final-summary + (#'f.tools.agent/extract-final-assistant-text [{:role "assistant" :content [{:type :tool-use :text "ignored"} {:type :text :text "Text only"}]}])))) - (testing "returns default when no assistant messages" - (is (= "Agent completed without producing output." - (#'f.tools.agent/extract-final-summary - [{:role "user" :content [{:type :text :text "Hi"}]}]))))) + (testing "returns nil when no assistant messages" + (is (nil? (#'f.tools.agent/extract-final-assistant-text + [{:role "user" :content [{:type :text :text "Hi"}]}]))))) (deftest definitions-test (testing "spawn_agent tool definition has correct structure" @@ -1063,10 +1169,13 @@ :properties {"agent" {:type "string"} "task" {:type "string"} "activity" {:type "string"} + "chat_id" {:type "string"} "model" {:type "string"} "variant" {:type "string"}} :required ["agent" "task"]} - (:parameters tool))))) + (:parameters tool))) + (is (= ["agent" "task" "activity" "chat_id" "model" "variant"] + (vec (keys (get-in tool [:parameters :properties]))))))) (testing "model and variant enums are absent when no models in db" (let [defs (f.tools.agent/definitions test-config {}) @@ -1099,14 +1208,468 @@ (is (= "Spawning agent" (summary-fn {:args {}})))))) +(deftest resume-admission-and-run-isolation-test + (let [db* (atom {:chats {"parent" {:model "openai/gpt-4.1"}}}) + context {:db* db* :config test-config :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "isolation" :trust true + :call-state-fn (constantly {:status :executing})} + handler (spawn-handler) + args {"agent" "explorer" "task" "next" "chat_id" "subagent-isolation"}] + (with-redefs [f.chat/prompt (fn [{:keys [chat-id]} & _] + (swap! db* update-in [:chats chat-id] + #(assoc % :status :idle :current-step 5 :max-steps-reached? true + :messages [{:role "assistant" :content [{:type :text :text "old answer"}]}])))] + (handler {"agent" "explorer" "task" "first" "variant" "high"} context)) + (let [baseline @db*] + (with-redefs [f.chat/prompt (fn [& _] (is false "Rejected resume must not prompt"))] + (doseq [selector ["" " " 42 [] "missing" "parent"]] + (is (thrown? clojure.lang.ExceptionInfo + (handler (assoc args "chat_id" selector) context))) + (is (= baseline @db*)))) + (doseq [[label change changed-args changed-context linked?] + [["foreign parent" identity args (assoc context :chat-id "foreign") false] + ["different agent" identity (assoc args "agent" "general") context false] + ["authorization" identity args (assoc-in context [:config :agent "explorer" :spawnableBy] "other") false] + ["tool future" #(assoc-in % [:chats "subagent-isolation" :tool-calls "old" :future] (delay nil)) args context true] + ["tool resources" #(assoc-in % [:chats "subagent-isolation" :tool-calls "old" :resources] {:process :remaining}) args context true] + ["outstanding worker" #(assoc-in % [:managed-chats "subagent-isolation" :workers] 1) args context true] + ["running" #(assoc-in % [:chats "subagent-isolation" :status] :running) args context true] + ["stopping" #(assoc-in % [:chats "subagent-isolation" :status] :stopping) args context true]]] + (testing label + (reset! db* (change baseline)) + (let [before @db*] + (with-redefs [f.chat/prompt (fn [& _] (is false "Rejected resume must not prompt"))] + ;; An agent the parent may not spawn is rejected before the chat_id is looked at. + (is (thrown-with-msg? clojure.lang.ExceptionInfo + #"Agent not found|(?s)'subagent-isolation' cannot be continued.*Omit chat_id to spawn a new subagent" + (handler changed-args changed-context))) + (is (= before @db*)))) + (testing "tool-call details link only a child of this parent, never disclosing another's settings" + (let [details (tools.util/tool-call-details-before-invocation + :spawn_agent changed-args nil {:db @db* :config (:config changed-context) + :chat-id (:chat-id changed-context) :tool-call-id "x"})] + (if linked? + (is (= "subagent-isolation" (:subagent-chat-id details))) + (is (match? {:subagent-chat-id nil :model nil} details))))))) + (reset! db* baseline) + (testing "settled interruption resets its outcome and ignores unstarted cleanup promises" + (swap! db* assoc-in [:managed-chats "subagent-isolation" :interrupted?] true) + (swap! db* assoc-in [:chats "subagent-isolation" :tool-calls "rejected"] + {:status :rejected :future-cleanup-complete?* (promise)}) + (with-redefs [f.chat/prompt (fn [params & _] + (is (= "high" (:variant params))) + (swap! db* assoc-in [:chats "subagent-isolation" :status] :idle))] + (let [result (handler args (assoc context :tool-call-id "empty"))] + (is (false? (:error result))) + (is (string/ends-with? (get-in result [:contents 0 :text]) + "Agent completed without producing output.")) + (is (not (string/includes? (tools.util/contents->text (:contents result)) "old answer")))))) + (testing "immediate prompt errors finish without polling and allow settled reuse" + (with-redefs [f.chat/prompt (constantly {:status :error})] + (let [run (future (handler args (assoc context :tool-call-id "error")))] + (try + (is (match? {:error true :contents [{:text #"(?s)^Subagent chat_id: subagent-isolation\n\n.*Failed.*setup failed"}]} + (deref run 5000 ::timeout))) + (is (true? (:error (handler args context)))) + (finally (when-not (realized? run) (future-cancel run))))))) + (testing "a run that was only interrupted says so instead of a generic failure" + (reset! db* baseline) + (with-redefs [f.chat/prompt (fn [& _] + (swap! db* #(-> % + (assoc-in [:chats "subagent-isolation" :status] :idle) + (assoc-in [:managed-chats "subagent-isolation" :interrupted?] true))))] + (is (match? {:error true :contents [{:text #"Failed\n\nThe sub-agent was stopped before it finished\."}]} + (handler args (assoc context :tool-call-id "interrupted"))))))))) + +(deftest spawn-agent-replay-per-call-test + (let [id "subagent-first" + answer #(assoc (assistant-msg %) :content-id %) + db {:chats {id {:messages [(answer "run 1") (answer "run 2")]}}} + output (fn [call-id arguments details] + {:role "tool_call_output" + :content {:id call-id :name "spawn_agent" :arguments arguments :error false + :details (merge {:type :subagent :agent-name "explorer" :subagent-chat-id id} details) + :output {:contents [{:type :text :text "Tool result"}]}}}) + child-texts (fn [message] + (->> (f.chat/messages->contents [message] {:chat-id "parent" :db db}) + (filter #(= id (:chat-id %))) + (keep #(get-in % [:content :text])) + (map string/trim)))] + (testing "each call replays only the part of the child chat it ran" + (is (= ["run 1"] (child-texts (output "first" {"agent" "explorer"} {:subagent-message-range [0 1]})))) + (is (= ["run 2"] (child-texts (output "second" {"agent" "explorer" "chat_id" id} {:subagent-message-range [1 2]}))))) + (testing "a call that did not run replays nothing" + (is (empty? (child-texts (output "busy" {"agent" "explorer" "chat_id" id} {:subagent-message-range [0 0]}))))) + (testing "a history from before continuation replays the whole child chat" + (is (= ["run 1" "run 2"] (child-texts (output "old" {"agent" "explorer"} {}))))) + (testing "a range past a rolled-back child chat is cut to what exists" + (is (= ["run 2"] (child-texts (output "rolled" {"agent" "explorer"} {:subagent-message-range [1 5]}))))))) + +(deftest spawn-agent-records-message-range-test + (let [db* (atom {:chats {"parent" {:model "openai/gpt-4.1"}}}) + context {:db* db* :config test-config :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "first" + :call-state-fn (constantly {:status :executing})} + details (fn [tool-call-id arguments] + (tools.util/tool-call-details-after-invocation + :spawn_agent arguments + (tools.util/tool-call-details-before-invocation :spawn_agent arguments nil + {:db @db* :config test-config :chat-id "parent" + :tool-call-id tool-call-id}) + nil {:db @db* :chat-id "parent" :tool-call-id tool-call-id}))] + (with-redefs [f.chat/prompt (fn [{:keys [chat-id message]} & _] + (swap! db* update-in [:chats chat-id] + ;; Like prompt-messages!, it stores the variant even when nil. + #(-> % (assoc :status :idle :variant nil) + (update :messages (fnil conj []) (assistant-msg message)))))] + ((spawn-handler) {"agent" "explorer" "task" "first"} context) + (let [first-details (details "first" {"agent" "explorer"})] + (is (= [0 1] (:subagent-message-range first-details))) + (is (not (contains? first-details :variant)) "No variant, no null on the wire")) + (let [args {"agent" "explorer" "task" "next" "chat_id" "subagent-first"}] + ((spawn-handler) args (assoc context :tool-call-id "next")) + (is (= [1 2] (:subagent-message-range (details "next" args))))) + (testing "a call that did not run gets an empty range" + (is (= [0 0] (:subagent-message-range (details "never" {"agent" "explorer"})))))))) + +(deftest managed-subagent-direct-prompt-test + (testing "a live subagent chat cannot be prompted directly, only by its parent through spawn_agent" + (swap! (h/db*) assoc-in [:chats "subagent-live"] {:id "subagent-live" :subagent {:name "explorer"} :status :idle}) + (swap! (h/db*) assoc-in [:managed-chats "subagent-live"] {:workers 0}) + (let [before @(h/db*)] + (is (match? {:chat-id "subagent-live" :status :error} + (f.chat/prompt {:chat-id "subagent-live" :message "hi"} + (h/db*) (h/messenger) (h/config) (h/metrics)))) + (is (= before @(h/db*)))))) + +(deftest spawn-agent-continue-model-override-test + (let [db* (atom {:chats {"parent" {:model "anthropic/claude-sonnet-4-6"}}}) + prompts* (atom []) + context {:db* db* :config test-config :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "first" + :call-state-fn (constantly {:status :executing})} + continue! (fn [tool-call-id overrides] + ((spawn-handler) (merge {"agent" "explorer" "task" "next" "chat_id" "subagent-first"} overrides) + (assoc context :tool-call-id tool-call-id)))] + (with-redefs [f.chat/prompt (fn [{:keys [chat-id model variant] :as params} & _] + (swap! prompts* conj params) + ;; Like prompt-messages!, the chat keeps the model it last used. + (swap! db* update-in [:chats chat-id] assoc :status :idle :model model :variant variant))] + ((spawn-handler) {"agent" "explorer" "task" "first" "variant" "high"} context) + (testing "without overrides it keeps its model and variant" + (continue! "keep" {}) + (is (match? {:model "anthropic/claude-sonnet-4-6" :variant "high"} (last @prompts*)))) + (testing "a variant override alone keeps the model" + (continue! "variant" {"variant" "low"}) + (is (match? {:model "anthropic/claude-sonnet-4-6" :variant "low"} (last @prompts*)))) + (testing "a new model drops the old variant" + (continue! "model" {"model" "openai/gpt-4.1"}) + (is (= "openai/gpt-4.1" (:model (last @prompts*)))) + (is (nil? (:variant (last @prompts*)))) + (testing "and is kept by later continues" + (continue! "after" {}) + (is (= "openai/gpt-4.1" (:model (last @prompts*)))))) + (testing "the variant is validated against the final model" + (is (thrown-with-msg? clojure.lang.ExceptionInfo #"Variant 'nope' is not available" + (continue! "bad" {"model" "anthropic/claude-opus-4-6" "variant" "nope"}))))))) + +(deftest spawn-agent-stop-waits-for-late-messages-test + (testing "a parent stop waits for the cancelled tool's result, so the call's replay range includes it" + (let [db* (atom {:chats {"parent" {:model "openai/gpt-4.1"}}}) + id "subagent-stop" + call-state* (atom {:status :executing})] + (with-redefs [f.chat/prompt (fn [& _] + (swap! db* update-in [:chats id] assoc :status :running :messages [(assistant-msg "Working.")]) + (reset! call-state* {:status :stopping})) + f.chat/prompt-stop (fn [& _] + (swap! db* assoc-in [:chats id :status] :stopping) + (future (Thread/sleep 100) + (swap! db* update-in [:chats id] #(-> % + (update :messages conj {:role "tool_call_output"}) + (assoc :status :idle)))))] + (let [started (System/currentTimeMillis) + result ((spawn-handler) {"agent" "explorer" "task" "t"} + (spawn-context db* :tool-call-id "stop" :call-state-fn #(deref call-state*)))] + (is (< (- (System/currentTimeMillis) started) 800) "Returns soon after the subagent settles") + (is (string/includes? (result-text result) "was stopped")) + (is (= [0 2] (get-in @db* [:chats "parent" :tool-calls "stop" :subagent-message-range])))))))) + +(deftest spawn-agent-empty-chat-id-spawns-fresh-test + (h/config! {:agent {"explorer" {:mode "subagent" :description "Explorer"}}}) + (doseq [chat-id ["" nil]] + (testing (str "chat_id " (pr-str chat-id) " is treated as absent") + (let [prompted* (promise)] + (with-redefs [f.chat/prompt (fn [{:keys [chat-id]} & _] + (deliver prompted* chat-id) + (swap! (h/db*) assoc-in [:chats chat-id :status] :idle))] + (let [tool-call-id (str "empty-" (count (pr-str chat-id))) + result (f.tools/call-tool! "eca__spawn_agent" {"agent" "explorer" "task" "work" "chat_id" chat-id "model" nil} + "parent" tool-call-id "code" (h/db*) (h/config) (h/messenger) (h/metrics) + (constantly {:status :executing}) (fn [& _]) {})] + (is (false? (:error result))) + (is (= (str "subagent-" tool-call-id) (deref prompted* 0 ::not-prompted))))))))) + +(deftest spawn-agent-resume-uses-current-agent-config-test + (let [db* (atom {:chats {"parent" {:model "openai/gpt-4.1"}}}) + prompts* (atom []) + context {:db* db* :config test-config :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "limits" + :call-state-fn (constantly {:status :executing})}] + (with-redefs [f.chat/prompt (fn [{:keys [chat-id] :as params} & _] + (swap! prompts* conj params) + (swap! db* update-in [:chats chat-id] assoc :status :idle :current-step 4))] + ((spawn-handler) {"agent" "explorer" "task" "first"} context) + (let [config (assoc-in test-config [:agent "explorer" :maxSteps] 2)] + ((spawn-handler) {"agent" "explorer" "task" "next" "chat_id" "subagent-limits"} + (assoc context :config config :tool-call-id "next")) + (is (string/includes? (:message (first @prompts*)) "maximum of 5 steps")) + (is (string/includes? (:message (second @prompts*)) "maximum of 2 steps")) + (is (= 2 (get-in @db* [:chats "subagent-limits" :max-steps]))) + (is (= 2 (get-in @db* [:chats "subagent-limits" :subagent :max-steps])) + "The step check reads the limit from the refreshed agent config"))))) + +(deftest managed-subagent-followup-workers-test + (h/config! {:env "dev" :hooks {"status" {:type "chatStatusChanged"}} + :agent {"explorer" {:mode "subagent" :description "Explorer"}}}) + (swap! (h/db*) assoc-in [:chats "parent"] {:model "openai/gpt-5.2"}) + (let [idle (promise) release-idle (promise) polled (promise) + second-finished (promise) release-worker (promise) + requests* (atom 0) idle-count* (atom 0) + context {:db* (h/db*) :config (h/config) :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "workers" + :call-state-fn (fn [] + (when (= :idle (get-in @(h/db*) [:chats "subagent-workers" :status])) + (deliver polled true)) + {:status :executing})} + handler (spawn-handler)] + (with-redefs [llm-api/sync-prompt! (constantly nil) + config/await-plugins-resolved! (constantly true) + f.tools/all-tools (constantly []) + hooks/trigger-if-matches! + (fn [type data callbacks & _] + (when (and (= :subagentPostRequest type) (not (:follow-up-active data))) + ((:on-after-action callbacks) {:name "follow" :exit 0 :parsed {"followUp" "continue"}})) + (when (and (= :chatStatusChanged type) (= :idle (:status data)) + (= 1 (swap! idle-count* inc))) + (deliver idle true) + (is (= true (deref release-idle 10000 ::timeout))))) + llm-api/sync-or-async-prompt! + (fn [{:keys [on-first-response-received on-message-received]}] + (let [n (swap! requests* inc)] + (on-first-response-received {:type :text :text "answer"}) + (on-message-received {:type :text :text (str "answer " n)}) + (on-message-received {:type :finish}) + (when (= 2 n) + (deliver second-finished true) + (is (= true (deref release-worker 10000 ::timeout))))))] + (let [run (future (handler {"agent" "explorer" "task" "work"} context))] + (try + (is (= true (deref idle 10000 ::timeout))) + (is (= true (deref polled 10000 ::timeout))) + (is (not (realized? run))) + (doseq [args [{"agent" "explorer" "task" "duplicate"} + {"agent" "explorer" "task" "resume" "chat_id" "subagent-workers"}]] + (is (thrown? clojure.lang.ExceptionInfo (handler args context)))) + (let [before @(h/db*)] + (is (match? {:status :error} + (f.chat/prompt {:chat-id "subagent-workers" :message "bypass"} + (h/db*) (h/messenger) (h/config) (h/metrics)))) + (is (= before @(h/db*)))) + (deliver release-idle true) + (is (= true (deref second-finished 10000 ::timeout))) + (is (not (realized? run))) + (deliver release-worker true) + (let [result (deref run 10000 ::timeout)] + (is (map? result)) + (is (string/includes? (tools.util/contents->text (:contents result)) "answer 2"))) + (finally + (deliver release-idle true) + (deliver release-worker true) + (when (= ::timeout (deref run 10000 ::timeout)) (future-cancel run)))))))) + +(deftest spawn-agent-real-max-steps-resume-test + (h/config! {:env "test" + :agent {"explorer" {:mode "subagent" :description "Explorer" :maxSteps 1}}}) + (swap! (h/db*) assoc-in [:chats "parent"] {:model "openai/gpt-5.2"}) + (let [requests* (atom []) + id "subagent-budget" + context {:db* (h/db*) :config (h/config) :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "budget" + :call-state-fn (constantly {:status :executing})}] + (with-redefs [llm-api/sync-prompt! (constantly nil) + config/await-plugins-resolved! (constantly true) + f.tools/all-tools (constantly []) + llm-api/sync-or-async-prompt! + (fn [{:keys [on-first-response-received on-message-received on-tools-called] :as request}] + (swap! requests* conj request) + (on-first-response-received {:type :text :text "Findings"}) + (on-message-received {:type :text :text "Findings"}) + (case (count @requests*) + 1 (on-tools-called [{:id "lookup" :full-name "eca__lookup" :arguments {}}]) + ;; The no-tools summary turn after the max steps halt. + 2 (do + (is (true? (get-in @(h/db*) [:chats id :summary-requested?]))) + (on-message-received {:type :finish})) + (do + (is (= 0 (get-in @(h/db*) [:chats id :current-step]))) + (is (nil? (get-in @(h/db*) [:chats id :max-steps-reached?]))) + (is (nil? (get-in @(h/db*) [:chats id :summary-requested?])) + "A continued child may call tools again") + (is (string/includes? (pr-str (:user-messages request)) "Tool calls are allowed again.")) + (on-message-received {:type :finish}))))] + (let [halted ((spawn-handler) {"agent" "explorer" "task" "find"} context) + history (get-in @(h/db*) [:chats id :messages])] + (is (true? (:error halted))) + (is (re-find #"^Subagent chat_id: subagent-budget\n\n## Agent 'explorer' Halted" + (tools.util/contents->text (:contents halted)))) + (is (true? (get-in @(h/db*) [:chats id :max-steps-reached?]))) + (is (= 1 (get-in @(h/db*) [:chats id :current-step]))) + (let [resumed ((spawn-handler) {"agent" "explorer" "task" "finish" "chat_id" id} + (assoc context :tool-call-id "continued"))] + (is (false? (:error resumed))) + (is (re-find #"^Subagent chat_id: subagent-budget\n\n## Agent 'explorer' Result" + (tools.util/contents->text (:contents resumed)))) + (is (= 3 (count @requests*))) + (is (seq history)) + (is (= history (take (count history) (get-in @(h/db*) [:chats id :messages])))) + (is (some #(= "assistant" (:role %)) (:past-messages (last @requests*))))))))) + +(deftest spawn-agent-provider-failure-resume-test + (h/config! {:env "test" :providers {"openai" {:retry {:maxAutoContinues 0}}} + :agent {"explorer" {:mode "subagent" :description "Explorer"}}}) + (swap! (h/db*) assoc-in [:chats "parent"] {:model "openai/gpt-5.2"}) + (let [requests* (atom []) + context {:db* (h/db*) :config (h/config) :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "network" :call-state-fn (constantly {:status :executing})} + handler (spawn-handler)] + (with-redefs [llm-api/sync-prompt! (constantly nil) + config/await-plugins-resolved! (constantly true) + f.tools/all-tools (constantly [{:name "lookup" :full-name "eca__lookup" + :server {:name "eca"} :origin :native}]) + f.tools/approval (constantly :allow) + f.tools/call-tool! (constantly {:contents [{:type :text :text "earlier tool result"}]}) + llm-api/sync-or-async-prompt! + (fn [{:keys [on-first-response-received on-message-received on-prepare-tool-call + on-tools-called on-error] :as request}] + (swap! requests* conj request) + (on-first-response-received {:type :text :text "started"}) + (case (count @requests*) + 1 (do (on-prepare-tool-call {:id "lookup" :full-name "eca__lookup" :arguments-text "{}"}) + (on-tools-called [{:id "lookup" :full-name "eca__lookup" :arguments {}}]) + (on-message-received {:type :text :text "partial finding"}) + (on-error {:message "Connection lost" :exception (java.net.ConnectException. "Connection refused")})) + (do (is (nil? (get-in @(h/db*) [:chats "subagent-network" :prompt-error]))) + (on-message-received {:type :text :text "recovered answer"}) + (on-message-received {:type :finish}))))] + (let [failed (handler {"agent" "explorer" "task" "find"} context) + history (get-in @(h/db*) [:chats "subagent-network" :messages])] + (is (true? (:error failed))) + (is (string/includes? (tools.util/contents->text (:contents failed)) "returned `chat_id`")) + (let [result (handler {"agent" "explorer" "task" "continue" "chat_id" "subagent-network"} context) + past (:past-messages (last @requests*))] + (is (false? (:error result))) + (is (string/includes? (tools.util/contents->text (:contents result)) "recovered answer")) + (is (= 2 (count @requests*))) + (is (some #(= "earlier tool result" (get-in % [:content :output :contents 0 :text])) past)) + (is (some #(= "partial finding" (get-in % [:content 0 :text])) past)) + (is (= history (take (count history) (get-in @(h/db*) [:chats "subagent-network" :messages]))))))))) + +(deftest spawn-agent-stopped-tool-resume-test + (doseq [phase [:dispatch :uninterruptible]] + (testing (name phase) + (h/reset-components!) + (h/config! {:env "dev" :agent {"explorer" {:mode "subagent" :description "Explorer"}}}) + (swap! (h/db*) assoc-in [:chats "parent"] {:model "openai/gpt-5.2"}) + (let [entered (promise) release (promise) stopped (promise) joining (promise) tool-ended (promise) + requests* (atom []) call-state* (atom {:status :executing}) callback* (atom nil) + id "subagent-cancel" + context {:db* (h/db*) :config (h/config) :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent" :tool-call-id "cancel" :call-state-fn #(deref call-state*)} + handler (spawn-handler) + transition tool-calls/transition-tool-call! + active tool-calls/get-active-tool-calls + ;; Ignore cancellation only until the test explicitly releases old work. + wait! (fn [] (loop [] + (let [result (try (deref release 30000 ::timeout) + (catch InterruptedException _ ::interrupted))] + (if (= ::interrupted result) + (recur) + (is (= true result)))))) + settled? (fn [] (loop [n 1000] + (cond (zero? (get-in @(h/db*) [:managed-chats id :workers] 0)) true + (zero? n) false + :else (do (Thread/sleep 10) (recur (dec n)))))) + stop! (fn [] + (reset! call-state* {:status :stopping}) + (f.chat/prompt-stop {:chat-id id} (h/db*) (h/messenger) (h/config) (h/metrics) {:silent? true}) + (deliver stopped true))] + (with-redefs [f.tools.agent/stop-settle-timeout-ms 100 + llm-api/sync-prompt! (constantly nil) + config/await-plugins-resolved! (constantly true) + f.tools/all-tools (constantly [{:name "lookup" :full-name "eca__lookup" + :server {:name "eca"} :origin :native}]) + f.tools/approval (constantly :allow) + tool-calls/get-active-tool-calls (fn [db chat-id] + (when (realized? stopped) (deliver joining true)) + (active db chat-id)) + tool-calls/transition-tool-call! + (fn [db* ctx tool-id event & data] + (try + (let [result (apply transition db* ctx tool-id event data)] + (when (and (= phase :dispatch) (= event :execution-start)) + (is (= true (deref entered 10000 ::timeout))) + (stop!)) + result) + (finally + (when (#{:execution-end :stop-attempted} event) (deliver tool-ended true))))) + f.tools/call-tool! + (fn [& args] + (reset! callback* (nth args 9)) + (deliver entered true) + (wait!) + {:contents [{:type :text :text "old tool result"}]}) + llm-api/sync-or-async-prompt! + (fn [{:keys [on-first-response-received on-message-received on-prepare-tool-call on-tools-called] + :as request}] + (swap! requests* conj request) + (on-first-response-received {:type :text :text "start"}) + (if (= 1 (count @requests*)) + (do (on-prepare-tool-call {:id "lookup" :full-name "eca__lookup" :arguments-text "{}"}) + (on-tools-called [{:id "lookup" :full-name "eca__lookup" :arguments {}}])) + (do (on-message-received {:type :text :text "corrected answer"}) + (on-message-received {:type :finish}))))] + (let [run (future (handler {"agent" "explorer" "task" "work"} context))] + (try + (is (= true (deref entered 10000 ::timeout))) + (when-not (= phase :dispatch) (stop!)) + (is (= true (deref stopped 10000 ::timeout))) + (is (= :stopping (:status (@callback*))) "Tool observes live state") + (when (= phase :dispatch) (is (= true (deref joining 1000 ::timeout)))) + (is (map? (deref run 10000 ::timeout))) + (is (thrown? clojure.lang.ExceptionInfo + (handler {"agent" "explorer" "task" "too soon" "chat_id" id} context))) + (deliver release true) + (is (map? (deref run 10000 ::timeout))) + (is (settled?)) + (reset! call-state* {:status :executing}) + (let [history (get-in @(h/db*) [:chats id :messages]) + resumed (handler {"agent" "explorer" "task" "correct course" "chat_id" id} context) + past (:past-messages (last @requests*))] + (is (false? (:error resumed))) + (is (string/includes? (tools.util/contents->text (:contents resumed)) "corrected answer")) + (is (= 2 (count @requests*)) "No provider continuation after stop") + (is (every? (set (map :role past)) ["tool_call" "tool_call_output"])) + (is (some #(= "old tool result" (get-in % [:content :output :contents 0 :text])) past)) + (is (= history (take (count history) (get-in @(h/db*) [:chats id :messages]))))) + (finally + (deliver release true) + (is (= true (deref tool-ended 10000 ::timeout))) + (is (not= ::timeout (deref run 10000 ::timeout))) + (is (settled?)))))))))) + (deftest spawn-agent-real-chat-prompt-test - ;; Regression test for v0.133.1 -> v0.133.2: spawn_agent failed end-to-end - ;; because chat/prompt's validate-client-chat-id rejected the deterministic - ;; "subagent-..." chat id used internally by spawn-agent. Every other test - ;; in this namespace mocks chat/prompt via with-redefs of requiring-resolve, - ;; so the validator path was never exercised. This test runs the REAL - ;; chat/prompt and only mocks the LLM transport, asserting the spawn handler - ;; reaches success through the chat layer. + ;; Regression: chat/prompt must accept the server-managed "subagent-..." ID. + ;; Exercise the real chat layer for both fresh and resumed delegation. (testing "spawn handler drives real chat/prompt to success" (h/reset-components!) (h/config! {:env "test" @@ -1117,7 +1680,9 @@ (fn [models] (merge {"openai/gpt-5.2" {:tools true}} (or models {})))) (swap! (h/db*) assoc-in [:chats "parent-1"] {:id "parent-1" :model "openai/gpt-5.2"}) - (let [api-mock (fn [{:keys [on-first-response-received on-message-received]}] + (let [requests* (atom []) + api-mock (fn [{:keys [on-first-response-received on-message-received] :as request}] + (swap! requests* conj request) (on-first-response-received {:type :text :text "Found it"}) (on-message-received {:type :text :text "Found it"}) (on-message-received {:type :finish}))] @@ -1127,16 +1692,8 @@ f.tools/approval (constantly :allow) config/await-plugins-resolved! (constantly true)] (let [handler (get-in (f.tools.agent/definitions (h/config) (h/db)) ["spawn_agent" :handler]) - ;; Run in a future with a timeout so a regression that causes - ;; chat/prompt to reject the subagent chat-id (and thus never - ;; flip the chat to :idle) fails the test instead of hanging - ;; the polling loop forever. The budget is generous because - ;; chat/prompt does its real work in a future* and the agent - ;; polling loop in agent.clj sleeps 1s between status checks, - ;; so on slower CI runners (notably macOS GitHub runners with - ;; cold JIT) a healthy run can still take a few polling - ;; iterations. A genuine regression hangs forever, so 30s is - ;; still a fast failure for that case. + ;; Bound the wait so a stalled prompt fails the test; allow headroom + ;; for slower CI runners. result-fut (future (handler {"agent" "explorer" "task" "find files" "activity" "exploring"} @@ -1154,16 +1711,33 @@ :messages (h/messages)}))] (when (identical? ::timeout result) (future-cancel result-fut)) - (testing "spawn handler completes (regression would hang the polling loop)" + (testing "spawn handler completes within the timeout" (is (not (identical? ::timeout result)) - (str "spawn handler did not complete in 30s — chat/prompt likely rejected the subagent chat-id. " - timeout-details))) + (str "spawn handler did not complete in 30s. " timeout-details))) (when (map? result) - (testing "spawn handler returns success (would be :error true under v0.133.1)" + (testing "spawn handler returns success" (is (match? {:error false :contents [{:type :text - :text #"^## Agent 'explorer' Result"}]} + :text #"^Subagent chat_id: subagent-tc-1\n\n## Agent 'explorer' Result"}]} result))) + (testing "resume retains the real transcript and original selections" + (let [child (get-in @(h/db*) [:chats "subagent-tc-1"]) + history (:messages child) + cache (:prompt-cache child)] + (is (every? (set (map :role history)) ["user" "assistant"])) + (swap! (h/db*) assoc-in [:chats "parent-1" :model] "other/model") + (let [resumed (handler {"agent" "explorer" "task" "Explain that finding" "chat_id" "subagent-tc-1"} + {:db* (h/db*) :config (h/config) :messenger (h/messenger) :metrics (h/metrics) + :chat-id "parent-1" :tool-call-id "tc-2" + :call-state-fn (constantly {:status :executing})})] + (is (false? (:error resumed))) + (is (= history (take (count history) (get-in @(h/db*) [:chats "subagent-tc-1" :messages])))) + (is (every? (set (map :role (:past-messages (last @requests*)))) + ["user" "assistant"])) + (is (= "Explain that finding" + (get-in (last @requests*) [:user-messages 0 :content 0 :text]))) + (is (= "gpt-5.2" (:model (last @requests*)))) + (is (= cache (get-in @(h/db*) [:chats "subagent-tc-1" :prompt-cache])))))) (testing "subagent chat reaches :idle through real chat/prompt" (is (= :idle (get-in @(h/db*) [:chats "subagent-tc-1" :status])))) (testing "subagent chat carries the parent-chat-id" @@ -1173,8 +1747,4 @@ (= "parent-1" (:parent-chat-id m)) (= :assistant (:role m)) (= {:type :text :text "Found it"} (:content m)))) - (:chat-content-received (h/messages))))))))) - ;; Reference f.chat to keep the require non-unused; the actual - ;; eca.features.chat/prompt is invoked indirectly via requiring-resolve - ;; inside spawn-agent. - (is (var? #'f.chat/prompt)))) + (:chat-content-received (h/messages))))))))))) diff --git a/test/eca/features/tools/shell_test.clj b/test/eca/features/tools/shell_test.clj index 29eee4920..1edd36b03 100644 --- a/test/eca/features/tools/shell_test.clj +++ b/test/eca/features/tools/shell_test.clj @@ -4,12 +4,29 @@ [babashka.process :as p] [clojure.test :refer [are deftest is testing]] [eca.features.tools.shell :as f.tools.shell] + [eca.shared :as shared] [eca.test-helper :as h] [matcher-combinators.test :refer [match?]])) (def ^:private call-state-fn (constantly {:status :executing})) (def ^:private state-transition-fn (constantly nil)) +(deftest background-command-is-not-a-tool-resource-test + (testing "a background job is never a tool-call resource, so stopping the chat neither waits for it nor kills it" + (let [events* (atom []) + ;; An existing dir: on Windows a temp dir that the job still uses cannot be deleted. + dir (System/getProperty "user.dir") + result ((get-in f.tools.shell/definitions ["shell_command" :handler]) + {"command" "echo started" "background" "test job"} + {:db {:workspace-folders [{:uri (shared/filename->uri dir) :name "project"}]} + :db* (atom {}) + :chat-id "chat" + :messenger (h/messenger) + :call-state-fn call-state-fn + :state-transition-fn (fn [event & _] (swap! events* conj event))})] + (is (match? {:contents [{:text #"Background job \S+ started"}]} result)) + (is (empty? @events*))))) + (deftest shell-command-test (testing "non-existent working_directory" (is (match?