From a03583159b59f14ae632578e8eabf1b1a5f98552 Mon Sep 17 00:00:00 2001 From: limityan Date: Sun, 2 Aug 2026 12:51:53 +0800 Subject: [PATCH] feat(extensions): delegate OpenCode commands to approved subagents --- .../agent-runtime-services-design.md | 7 +- .../opencode-config-assets-adapter-design.md | 14 +- .../opencode-extension-compatibility.md | 2 +- docs/plans/core-decomposition-plan.md | 2 +- .../rules/source/public-api-rules.mjs | 2 + .../rules/source/required-rules.mjs | 33 +- scripts/core-boundaries/self-test.mjs | 12 +- src/apps/cli/src/agent/runtime_client.rs | 79 ++- src/apps/cli/src/dispatch/worker.rs | 1 + src/apps/cli/src/modes/chat.rs | 6 +- src/apps/cli/src/modes/chat/commands.rs | 30 +- src/apps/cli/src/modes/chat/sessions.rs | 71 +++ src/apps/cli/src/modes/chat/tests.rs | 28 + src/apps/cli/src/peer_host/commands/dialog.rs | 1 + src/apps/desktop/src/api/agentic_api.rs | 11 +- .../src/tests/protocol_contracts.rs | 1 + .../src/tests/shared_controller.rs | 1 + .../claude-code-adapter/src/command_source.rs | 1 + .../opencode-adapter/src/command_source.rs | 36 +- .../tests/opencode_command_adapter.rs | 58 +- .../src/agentic/agents/registry/external.rs | 42 +- .../core/src/agentic/agents/registry/tests.rs | 26 + .../src/agentic/coordination/coordinator.rs | 584 +++++++++++++++++- .../src/agentic/coordination/scheduler.rs | 277 ++++++++- .../src/agentic/deep_review/task_adapter.rs | 25 - .../core/src/agentic/permission_policy.rs | 14 + .../src/agentic/session/session_manager.rs | 82 ++- .../implementations/session_message_tool.rs | 1 + .../tools/implementations/task/execution.rs | 38 +- .../agentic/tools/implementations/task/mod.rs | 2 +- .../assembly/core/src/external_sources.rs | 167 ++++- .../assembly/core/src/external_subagents.rs | 13 + .../assembly/core/src/service/cron/service.rs | 1 + .../core/src/service_agent_runtime.rs | 1 + .../assembly/external-sources/src/lib.rs | 1 + .../tests/coordinator_contracts.rs | 1 + .../product-domains/src/external_sources.rs | 101 ++- .../tests/external_source_contracts.rs | 2 + src/crates/contracts/runtime-ports/src/lib.rs | 60 ++ .../agent-runtime/examples/sdk_minimal.rs | 2 +- .../src/deep_review/task_execution.rs | 45 +- src/crates/execution/agent-runtime/src/lib.rs | 1 + .../execution/agent-runtime/src/runtime.rs | 3 + src/crates/execution/agent-runtime/src/sdk.rs | 19 +- .../agent-runtime/src/subagent_task.rs | 53 ++ .../agent-runtime/tests/sdk_smoke.rs | 2 +- .../interfaces/acp/src/runtime/prompt.rs | 1 + src/crates/interfaces/sdk-host/src/host.rs | 1 + .../src/flow_chat/components/ChatInput.tsx | 14 +- .../src/flow_chat/hooks/useMessageSender.ts | 7 + .../src/flow_chat/services/FlowChatManager.ts | 1 + .../flow-chat-manager/MessageModule.test.ts | 26 + .../flow-chat-manager/MessageModule.ts | 10 + .../api/service-api/AgentAPI.ts | 9 + .../service-api/ExternalSourcesAPI.test.ts | 44 +- .../api/service-api/ExternalSourcesAPI.ts | 66 +- src/web-ui/src/locales/en-US/flow-chat.json | 1 + src/web-ui/src/locales/zh-CN/flow-chat.json | 1 + src/web-ui/src/locales/zh-TW/flow-chat.json | 1 + 59 files changed, 1972 insertions(+), 169 deletions(-) create mode 100644 src/crates/execution/agent-runtime/src/subagent_task.rs diff --git a/docs/architecture/agent-runtime-services-design.md b/docs/architecture/agent-runtime-services-design.md index a9e910b5c4..15cdc03b05 100644 --- a/docs/architecture/agent-runtime-services-design.md +++ b/docs/architecture/agent-runtime-services-design.md @@ -68,7 +68,7 @@ Agent Runtime API 的逻辑归属与物理部署分离:相同归属模块可 私有 SDK Host 或目标机器 Runtime 中。任何 Rust 部署都只管理自己进程树内的服务与 Node/Bun Plugin Host;不能因为多个 GUI/TUI/Remote Client 连接就复制 Runtime 状态模块,或按 Client/Workspace 创建 Plugin Host。 -Rust Runtime SDK 以 `AGENT_RUNTIME_SDK_API_VERSION` 标记兼容边界。当前接口版本为 v2 preview: +Rust Runtime SDK 以 `AGENT_RUNTIME_SDK_API_VERSION` 标记兼容边界。当前接口版本为 v3 preview: 小版本更新允许增加可选 builder hook、有默认实现的端口方法或注册表查询能力,但不得向外部可用 Rust 结构体字面量(struct literal)构造的 DTO 直接增加字段,也不得改变既有端口语义、错误分类、session / turn 标识含义或 默认 feature 依赖。任何需要调用方改写现有嵌入代码的变更,必须提升接口版本并提供兼容迁移路径。 @@ -77,6 +77,9 @@ v2 的迁移只涉及 Rust 错误名词治理:调用方把 `RuntimeBuildError::UnsupportedPluginRuntimeHostBinding` 替换为 `RuntimeBuildError::UnsupportedPluginRuntimeClientAvailability`;错误分类和 builder 行为不变。 +v3 为 `AgentDialogTurnRequest` 增加来源无关的 `execution` 事实。现有 Rust struct literal 调用方迁移时增加 +`execution: AgentDialogTurnExecution::Standard`(或 `Default::default()`);旧 wire payload 缺省为标准执行。 + 只要外部调用方仍必须导入 `bitfun-core`、启用 `product-full`、持有具体服务管理器、读取产品命令 注册表、理解 ACP/内部端口或依赖全局可变状态,公开 SDK 发布边界就不成立。公开 SDK 的完整 术语、能力等价和版本要求以 [`agent-sdk-product-architecture.md`](agent-sdk-product-architecture.md) 为准。 @@ -420,7 +423,7 @@ impl AgentRuntime { 该 Rust 接口是内部产品入口复用的当前形态,不是公开 Python/TypeScript SDK 的目标 API。它必须只接收 已组装的类型化部件,不负责创建 文件系统、终端、MCP、AI 客户端、Remote 提供方或产品命令。 -当前 v2 preview 接口以 message / attachment / metadata 作为最小输入形态;若把 +当前 v3 preview 接口以 message / attachment / metadata 和默认标准执行目标作为最小输入形态;若把 model-round cancellation token、结构化 AgentInput 或更复杂的事件游标纳入公开 SDK, 必须分别评审 Rust Runtime SDK、SDK Host protocol 和公开 SDK API 的版本,并保留旧路径兼容。 diff --git a/docs/architecture/extensions/opencode-config-assets-adapter-design.md b/docs/architecture/extensions/opencode-config-assets-adapter-design.md index 9251572827..3c63dd4d3b 100644 --- a/docs/architecture/extensions/opencode-config-assets-adapter-design.md +++ b/docs/architecture/extensions/opencode-config-assets-adapter-design.md @@ -209,7 +209,7 @@ OpenCode adapter 在来源发现、解析和审批前不 import module、不读 | Agents / Modes | 当前生产 V1 `agent/prompt/disable/permission`、Core V2 `agents/system/disabled/permissions` 输入形状,以及 Markdown、description、mode、model、variant、temperature、top_p、steps、deprecated `maxSteps`、deprecated `tools`、options、hidden、color | Agent 归属模块创建兼容定义和使用范围视图;OpenCode adapter 只翻译来源语义 | 当前支持 Subagent 安全子集、Agent-local 权限约束和不透明 `variant` profile;V1 是生产兼容主路径,Core V2 字段只按已验证安全子集解析;variant 不映射为 reasoning 或请求 options,需显式绑定现有模型配置;首次按行为、来源、模型/profile、工具与权限范围确认,fresh single-run 调用 | primary/mode、options、采样、steps 与续接保持诊断或阻断;root ambient 权限和 V1 嵌套 resource map 尚不激活,不影响其他 Agent。 | | Skills | `.opencode/.claude/.agents` 项目与用户根、`SKILL.md`、`skills.paths/urls` | OpenCode adapter 只由 `bitfun-core/external_sources` 组合并投影有序本地配置根;Skill 归属模块负责有界递归、解析、覆盖与按需加载 | 标准根及 V1 `skills.paths`/当前本地字符串数组可用;项目配置限项目根,用户配置限项目根或用户目录;配置根最多 64 个、每根 512 个 Skill、单文件 256 KiB、可选策略 64 KiB,实际加载再次执行有界非链接读取;配置根在同 scope 覆盖标准 OpenCode 根,但不重排更早的 BitFun/Claude/Codex/Cursor 来源 | URL、下载/缓存、脚本与外部依赖不加载;无效根不影响标准 Skill。 | | References | `references` / 旧 `reference`,本地 path 或 Git repository/branch/description/hidden | OpenCode adapter 输出来源无关的 Reference provider snapshot;Product Assembly 生命周期协调器与 BitFun 原生关联目录合成唯一有效引用目录;关联目录视图和既有目录选择器消费 | 当前支持本地声明路径、description/hidden、异步刷新和 `@alias` 展示;原生关联目录始终在 OpenCode 引用之前,外部引用只读、不自动进入 Prompt 且不改变权限 | Git 引用、Remote 发现和下载/缓存不实现;无效高优先级 entry 阻断同 alias 的旧值并给出诊断,不回退到更宽松来源。 | -| Commands | JSON/JSONC、Markdown、`$ARGUMENTS`、位置参数、`@file`、`!shell`、agent/model/variant/subtask | Prompt Command 专属契约;adapter 提取静态文件引用与 shell 计划,Product Assembly 负责审批指纹和装配,Terminal owner 负责进程执行 | prompt 与静态 workspace 相对 UTF-8 `@file` 可发送;`!shell` 展示精确命令、工作目录与绝对 shell 路径,经重新校验后以不加载 profile 的隔离式 argv 执行,并仅把 stdout 按模板顺序加入 Prompt。为保持 OpenCode 语义,正常退出后的非零退出码仍使用 stdout。静态计划可记住,参数相关计划仅可单次运行;动态/绝对/越界文件及 agent/model/variant/subtask 整体受限 | 任一文件读取、进程启动、超时或超限失败时不发送部分 Prompt;进程副作用不可回滚。最多 8 文件、单文件 64 KiB、文件总量 128 KiB;最多 8 条 shell 指令、单条 64 KiB、总计 128 KiB、每条 stdout 256 Ki 字符、30 秒;最终命令 1 MiB。安全模式禁用,Remote 不回退到本机。 | +| Commands | JSON/JSONC、Markdown、`$ARGUMENTS`、位置参数、`@file`、`!shell`、agent/model/variant/subtask | Prompt Command 专属契约;adapter 提取静态文件引用与 shell 计划,Product Assembly 负责审批指纹和装配,Terminal owner 负责进程执行 | prompt 与静态 workspace 相对 UTF-8 `@file` 可发送;`!shell` 展示精确命令、工作目录与绝对 shell 路径,经重新校验后以不加载 profile 的隔离式 argv 执行,并仅把 stdout 按模板顺序加入 Prompt。为保持 OpenCode 语义,正常退出后的非零退出码仍使用 stdout。静态计划可记住,参数相关计划仅可单次运行;显式 agent 加缺省/true subtask 可走 approved fresh Subagent,其余 agent/model/variant/subtask 组合以及 shell 与委派的组合整体受限 | 任一文件读取、进程启动、超时或超限失败时不发送部分 Prompt;进程副作用不可回滚。最多 8 文件、单文件 64 KiB、文件总量 128 KiB;最多 8 条 shell 指令、单条 64 KiB、总计 128 KiB、每条 stdout 256 Ki 字符、30 秒;最终命令 1 MiB。安全模式禁用,Remote 不回退到本机。 | | MCP | local 的 command/environment/cwd/timeout,remote 的 URL/headers/oauth/timeout,Agent 选择 | MCP 归属模块创建兼容配置视图 | 当前支持 local stdio 和 HTTPS remote 的静态发现、首次/行为变化审批、冲突选择与 workspace 隔离的运行期接纳;C0a 快照导入只复制无 env/cwd 的 local command/args 或无 header/query/fragment 的 HTTPS remote,并保持 disabled | `{env:NAME}` 当前只允许用于运行期兼容来源的 environment/Header 值,不进入 C0a 快照;SSE、OpenCode OAuth client 配置、完整 timeout/Agent 范围与 Remote 执行域保持明确不支持;凭据或网络失败只影响单个 Server。 | | LSP | command、extensions、env、initialization | LSP 归属模块注册兼容实例 | 首次确认外部进程和使用范围后按文件类型启动 | 自定义 Server 缺少 extensions 或启动失败时只禁用该项。 | | Formatters | command、environment、extensions、`$FILE` | **基础能力缺失**:先补文件写入后的 Formatter 执行消费点,再做格式转换 | 首次确认命令后执行匹配 Formatter | 超时后标记未格式化,文件写入结果保留。 | @@ -335,8 +335,16 @@ Agent Prompt,不把配置声明解释成工作区外文件授权。只有用 当前 Prompt Command 子集展开 `$ARGUMENTS` 与 `$1`、`$2` 等位置参数,并支持模板中可静态确认的 workspace 相对 UTF-8 `@file`。OpenCode adapter 只从原模板提取引用,不扫描用户参数;Product Assembly 在 stale/冲突校验后通过共享本地文本服务 -原子读取并追加内容。动态占位、绝对/`~`/URL/越界路径仍进入目录但整体受限。包含 `!shell`、`{env:...}`、`{file:...}`、 -agent/model/variant/subtask 的命令同样保持受限,不能删除不支持的部分后继续发送。 +原子读取并追加内容。动态占位、绝对/`~`/URL/越界路径仍进入目录但整体受限。包含 `{env:...}`、`{file:...}`、 +`model`、`variant`、`subtask: false` 或仅有 `subtask: true` 而没有显式 `agent` 的命令同样保持受限,不能删除不支持的部分后继续发送。 + +显式 `agent` 且 `subtask` 缺省或为 `true` 的命令携带来源无关的 fresh external subagent 执行目标;该字段只描述本次 +Prompt Command 的执行意图,不公开新的 Subagent API。Product Assembly 仅在同 workspace 存在同 OpenCode 生态、逻辑 ID +精确匹配、已审批且当前 generation 有效的 External route 时保留命令可用性。提交后,Scheduler 要求本地 Session 处于 Idle, +Coordinator 复用既有 Task、Subagent Registry、generation lease、权限上限、取消、事件与会话持久化链路创建一次 fresh child; +任何失配或状态竞争都整次失败,不回退到当前 Agent、其他生态同名 Agent、旧 generation、Remote 或本机替代执行。包含 +`!shell` 的委派命令在目录阶段整体受限,避免在 Session/界面准入失败前产生不可回滚的进程副作用。GUI 的附件/ +引用上下文与 Shared TUI 在该路径明确拒绝,因为当前 Task 输入契约不能无损表达它们;普通 inline command 行为不变。 每次调用最多接纳 8 个不同文件,单文件 64 KiB、文件总量 128 KiB、最终命令 1 MiB。共享服务对每级路径执行 workspace 规范化包含校验并拒绝符号链接/reparse point;任一引用缺失、越界、超限或不是 UTF-8 时整次调用失败,不返回 diff --git a/docs/architecture/extensions/opencode-extension-compatibility.md b/docs/architecture/extensions/opencode-extension-compatibility.md index 88cacff002..add24be959 100644 --- a/docs/architecture/extensions/opencode-extension-compatibility.md +++ b/docs/architecture/extensions/opencode-extension-compatibility.md @@ -110,7 +110,7 @@ OpenCode,和 OpenCode 配置/插件进入 BitFun 是两个独立验收方向 | Agents / Modes | 融合现有能力 + 转换参数 | 部分实现:Subagent 安全子集、模型/profile 绑定与 Agent-local 权限约束 | 可主要适配 | OC-R1 | 已支持当前生产 V1 `agent/prompt/disable/permission` 与 Core V2 `agents/system/disabled/permissions` 的已验证安全子集、全局/项目 Markdown 和 JSON/JSONC、subagent/all、description、`Default`/不透明模型引用、不透明 `variant` 意图和工具映射,并接入唯一精确匹配、显式 BitFun 模型绑定、Web/TUI 可见性、审批、冲突、更新、撤下和 fresh single-run Task;仅显式声明 `model` 的 Agent 才保留 `variant`,未声明模型时与 OpenCode 一样不生效。保留的 `variant` 不推断为 reasoning effort,也不生成请求级 override,需显式绑定现有配置。不维护厂商别名、质量推断或自动 fallback。V1 `disable` 保持 deep-merge,V2 `disabled` 保持 remove/re-add,不能混用生命周期语义。有序 V2 permission rules 与 V1 扁平精确 action map 会成为只可收紧的独立约束。primary/mode、root ambient permission、V1 action pattern/嵌套 resource map、跨路径与命令资源域的歧义 pattern、options、采样与续接明确阻断或降级 | [Agents 与 Skills](opencode-config-assets-adapter-design.md#52-agentsmodes-与-skills) | | Skills | 转换参数 | 部分实现:标准根与本地配置根 | 可完整适配 | OC-R2 | 现有 Registry 除标准用户/项目根外,也通过 `bitfun-core/external_sources` 组合边界按 OpenCode 配置来源顺序累加 V1 `skills.paths` 与当前迁移后的本地字符串数组;仅接受项目根/用户目录内的本地目录并做有界递归发现。同 scope 配置根覆盖标准 OpenCode 根,但不重排更早的 BitFun/Claude/Codex/Cursor 来源。URL、下载/缓存、完整 allow/deny/ask 顺序及外部来源策略仍未实现 | [Agents 与 Skills](opencode-config-assets-adapter-design.md#52-agentsmodes-与-skills) | | References | 融合现有能力 + 转换参数 | 部分实现:本地目录与既有 Workspace 消费点 | 可主要适配 | OC-R2 | 已按 OpenCode V2 `e4bd9757` 的独立来源顺序解析 `references`/旧 `reference` 的本地 path、description/hidden,相同 alias 后者覆盖;通过独立生命周期协调器与 BitFun 原生关联目录合成 native-first 有效快照,接入关联目录弹窗和既有 `@` 目录选择器。外部声明不自动进入 Prompt、不授予文件权限;Git、Remote、下载/缓存明确不支持且不做临时实现 | [References](opencode-config-assets-adapter-design.md#521-references) | -| Commands | 补扩展接口 + 转换参数 | 部分实现:prompt、本地文本文件与经审阅的 shell 上下文 | 可完整适配 | OC-R2 | 已支持全局/项目 JSON、JSONC、Markdown 命令、参数展开、动态目录、刷新和显式冲突选择;模板中的静态 workspace 相对 `@file` 可在调用时有界读取,`!shell` 经精确计划审阅后仅把 stdout 加入 Prompt,静态计划可记住、参数相关计划仅可单次运行。动态/绝对/越界文件引用及 Agent/model/variant/subtask 保持受限且不做部分执行;Remote 不回退本机执行 | [Commands](opencode-config-assets-adapter-design.md#53-commands) | +| Commands | 补扩展接口 + 转换参数 | 部分实现:prompt、本地文本文件、经审阅的 shell 上下文与显式 Subagent 委派 | 可完整适配 | OC-R2 | 已支持全局/项目 JSON、JSONC、Markdown 命令、参数展开、动态目录、刷新和显式冲突选择;模板中的静态 workspace 相对 `@file` 可在调用时有界读取,`!shell` 经精确计划审阅后仅把 stdout 加入 Prompt,静态计划可记住、参数相关计划仅可单次运行。仅 `agent` 加缺省/`true` 的 `subtask` 可委派给同 workspace、同 OpenCode 生态、已审批且仍有效的精确 Subagent,并复用现有 fresh Task 生命周期;shell 与委派的组合、`model`、`variant`、`subtask: false`、隐式默认 Agent、Remote 与附件上下文保持受限,不回退到当前 Agent 或本机执行 | [Commands](opencode-config-assets-adapter-design.md#53-commands) | | Models / Providers 配置 | 融合现有能力 | 未实现 | 可主要适配 | OC-R1 | 静态字段进入模型归属模块;动态模型、鉴权和请求头交给插件运行时 | [声明式资产](opencode-config-assets-adapter-design.md#5-声明式资产映射) | | MCP | 转换参数 | 部分实现:local stdio 与 HTTPS remote | 可完整适配 | OC-R2 | 已接入发现、审批、冲突、workspace 隔离、更新和启动反馈;SSE、OAuth、完整 timeout/Agent 范围仍不支持;Remote 不回退本机实例 | [MCP、LSP 与 Formatter](opencode-config-assets-adapter-design.md#54-mcplsp-与-formatter) | | LSP | 转换参数 | 未实现 | 可完整适配 | OC-R2 | R1 解析;R2 转换 command、extensions、env 和 initialization 并由 LSP 归属模块启动 | [MCP、LSP 与 Formatter](opencode-config-assets-adapter-design.md#54-mcplsp-与-formatter) | diff --git a/docs/plans/core-decomposition-plan.md b/docs/plans/core-decomposition-plan.md index 73478da219..fbc9f1b9bd 100644 --- a/docs/plans/core-decomposition-plan.md +++ b/docs/plans/core-decomposition-plan.md @@ -25,7 +25,7 @@ | CLI / Desktop / ACP | 三者仍按需启用 `bitfun-core/product-full`;CLI 与 ACP 已分别提交对应 `DeliveryProfile` 并消费 Runtime Parts/SDK,Desktop 主交互已消费由现有 owner 构造的窄口径 SDK 接口 | 三个入口均复用单一 Core owner;完整 Desktop profile 和剩余兼容操作仍需逐项迁移 | | Server | 当前生产路由只形成 health/info/ping 基线 | 没有插件状态或独立产品组装完整流程 | | Server / Remote / Web / Mobile Web / SDK profile | 当前为空计划、未接入入口或仅有 preview 测试 | 不得据枚举值宣称产品能力已交付 | -| Agent Runtime SDK | 已有无 `bitfun-core` 依赖的 v2 preview 接口和 smoke test | 发布边界仍需真实嵌入方证明 | +| Agent Runtime SDK | 已有无 `bitfun-core` 依赖的 v3 preview 接口和 smoke test | 发布边界仍需真实嵌入方证明 | | 插件运行时 | 现有路径只覆盖 BitFun 原生包和 OpenCode custom tool 静态名称预览 | 不能据通用消息结构或静态候选扩张稳定 ABI | | Relay | room/device 状态、account/sync 存储、asset store 与 HTTP/WebSocket router 已归属 `services/relay-service`,standalone 与 embedded 入口同向消费;embedded bind、静态 fallback 和任务生命周期由 Desktop 窄宿主端口持有 | Cargo metadata 门禁覆盖 workspace、独立 manifest、normal/build/dev 依赖及 optional/target 变体;宿主归位已完成并由生命周期与边界测试保护 | | CLI CI | 独立 Linux job 运行 CLI test,通用三平台 workspace check 覆盖 CLI 编译;Linux PTY 与 Windows ConPTY 有启动页生命周期及本地确定性流式模型夹具驱动的活动 turn 进程测试,发布归档上传前校验 SHA-256 并解压执行 | 参数/序列化/前置失败和组装已有 focused contract;本地模型 HTTP 403 授权拒绝、流中断后的重试失败、Linux PTY/Windows ConPTY Chat resize/取消、`exec` Ctrl+C 及 Patch I/O 失败已有分层回归,真实供应商审批交互、macOS 活动 PTY 与 OS 级终端故障注入仍需补齐 | diff --git a/scripts/core-boundaries/rules/source/public-api-rules.mjs b/scripts/core-boundaries/rules/source/public-api-rules.mjs index 1b928d91fe..32c97d9d0d 100644 --- a/scripts/core-boundaries/rules/source/public-api-rules.mjs +++ b/scripts/core-boundaries/rules/source/public-api-rules.mjs @@ -711,6 +711,7 @@ export const externalSourceContractPublicApiEntries = [ 'ExternalSourceRecord', 'PromptCommandAvailability', 'PromptCommandDefinition', + 'PromptCommandExecutionTarget', 'ExpandedPromptCommand', 'PromptCommandExpansion', 'PromptCommandShellPreference', @@ -1029,6 +1030,7 @@ export const externalSourceCorePublicApiEntries = [ 'PromptCommandAvailability', 'PromptCommandCatalogEntry', 'PromptCommandDefinition', + 'PromptCommandExecutionTarget', 'PromptCommandInvocationOutcome', 'PromptCommandShellReviewDecision', 'PromptCommandShellReviewMode', diff --git a/scripts/core-boundaries/rules/source/required-rules.mjs b/scripts/core-boundaries/rules/source/required-rules.mjs index 6fbf9c640f..ba484e5097 100644 --- a/scripts/core-boundaries/rules/source/required-rules.mjs +++ b/scripts/core-boundaries/rules/source/required-rules.mjs @@ -1094,6 +1094,21 @@ export const requiredContentRules = [ }, ], }, + { + path: 'src/crates/execution/agent-runtime/src/subagent_task.rs', + reason: + 'agent-runtime must own provider-neutral subagent Task completion presentation shared by ordinary Task and product command delegation', + patterns: [ + { + regex: /\bpub struct SubagentTaskCompletionResultInput\b/, + message: 'missing provider-neutral subagent Task completion input', + }, + { + regex: /\bpub fn subagent_task_completion_result\b/, + message: 'missing provider-neutral subagent Task completion formatter', + }, + ], + }, { path: 'src/crates/execution/agent-runtime/src/deep_review/mod.rs', reason: @@ -1180,7 +1195,11 @@ export const requiredContentRules = [ }, { regex: /\bpub fn deep_review_task_completion_result\b/, - message: 'missing DeepReview task completion result presentation owner function', + message: 'missing DeepReview task completion result compatibility wrapper', + }, + { + regex: /crate::subagent_task::subagent_task_completion_result/, + message: 'missing DeepReview delegation to the provider-neutral Task formatter', }, { regex: /\bpub fn deep_review_cancelled_reviewer_result\b/, @@ -2520,10 +2539,6 @@ export const requiredContentRules = [ regex: /runtime_task_execution::DeepReviewReviewerAdmissionQueueRuntime::start/, message: 'missing DeepReview reviewer admission queue runtime delegation', }, - { - regex: /runtime_task_execution::deep_review_task_completion_result/, - message: 'missing DeepReview task completion result runtime delegation', - }, { regex: /runtime_task_execution::deep_review_cancelled_reviewer_result/, message: 'missing DeepReview cancelled reviewer result runtime delegation', @@ -2579,8 +2594,8 @@ export const requiredContentRules = [ message: 'missing TaskTool DeepReview retry guidance facade call', }, { - regex: /deep_review_task_adapter::deep_review_task_completion_result/, - message: 'missing TaskTool DeepReview completion result facade call', + regex: /bitfun_agent_runtime::subagent_task::subagent_task_completion_result/, + message: 'missing TaskTool provider-neutral completion result owner call', }, { regex: /DeepReviewProviderCapacityRetryRuntime::default/, @@ -5708,6 +5723,10 @@ export const requiredContentRules = [ regex: /pub use bitfun_runtime_ports::DialogTriggerSource;/, message: 'missing dialog trigger source compatibility re-export', }, + { + regex: /bitfun_agent_runtime::subagent_task::subagent_task_completion_result/, + message: 'missing delegated command provider-neutral Task result formatting', + }, ], }, { diff --git a/scripts/core-boundaries/self-test.mjs b/scripts/core-boundaries/self-test.mjs index 7901cdd468..ca1342fbbd 100644 --- a/scripts/core-boundaries/self-test.mjs +++ b/scripts/core-boundaries/self-test.mjs @@ -2986,6 +2986,7 @@ export function runManifestParserSelfTest({ 'RemoteControlStatePort', 'generic attachments', 'DialogTriggerSource', + 'bitfun_agent_runtime::subagent_task::subagent_task_completion_result', ], }, { @@ -3630,10 +3631,18 @@ export function runManifestParserSelfTest({ 'render_subagent_line', ], }, + { + path: 'src/crates/execution/agent-runtime/src/subagent_task.rs', + contracts: [ + 'SubagentTaskCompletionResultInput', + 'subagent_task_completion_result', + ], + }, { path: 'src/crates/execution/agent-runtime/src/deep_review/task_execution.rs', contracts: [ 'deep_review_task_completion_result', + 'crate::subagent_task::subagent_task_completion_result', 'deep_review_cancelled_reviewer_result', 'should_emit_deep_review_retry_guidance', 'deep_review_retry_guidance', @@ -3657,7 +3666,6 @@ export function runManifestParserSelfTest({ { path: 'src/crates/assembly/core/src/agentic/deep_review/task_adapter.rs', contracts: [ - 'runtime_task_execution::deep_review_task_completion_result', 'runtime_task_execution::deep_review_cancelled_reviewer_result', 'runtime_task_execution::should_emit_deep_review_retry_guidance', 'runtime_task_execution::deep_review_retry_guidance', @@ -3682,7 +3690,7 @@ export function runManifestParserSelfTest({ path: 'src/crates/assembly/core/src/agentic/tools/implementations/task/execution.rs', contracts: [ 'deep_review_task_adapter::deep_review_retry_guidance', - 'deep_review_task_adapter::deep_review_task_completion_result', + 'bitfun_agent_runtime::subagent_task::subagent_task_completion_result', 'DeepReviewProviderCapacityRetryRuntime::default', 'DeepReviewProviderCapacityRetryDecision::WaitForCapacity', ], diff --git a/src/apps/cli/src/agent/runtime_client.rs b/src/apps/cli/src/agent/runtime_client.rs index 82c7b627f0..8f5e8db9e4 100644 --- a/src/apps/cli/src/agent/runtime_client.rs +++ b/src/apps/cli/src/agent/runtime_client.rs @@ -12,7 +12,7 @@ use std::sync::{Arc, RwLock}; use tokio::sync::{broadcast, Mutex}; use bitfun_agent_runtime::sdk::{ - AgentDialogTurnRequest, AgentEventReceiver, AgentInputAttachment, + AgentDialogTurnExecution, AgentDialogTurnRequest, AgentEventReceiver, AgentInputAttachment, AgentLocalCommandTurnRecordRequest, AgentMessageWorkspaceReferencesRequest, AgentRuntime, AgentSessionCompactionRequest, AgentSessionCreateRequest, AgentSessionDeleteRequest, AgentSessionForkBeforeTurnRequest, AgentSessionForkRequest, AgentSessionForkResult, @@ -1186,6 +1186,58 @@ impl CliAgentRuntimeClient { )); } let session_id = self.ensure_session(agent_type).await?; + self.submit_dialog_turn_request( + session_id, + message, + None, + workspace_references, + attachments, + AgentDialogTurnExecution::Standard, + agent_type, + ) + .await + } + + pub(crate) async fn send_external_subagent_command( + &self, + prompt: String, + original_command: String, + ecosystem_id: String, + logical_id: String, + agent_type: &str, + ) -> Result { + if self.is_shared() { + return Err(anyhow::anyhow!( + "External subagent commands require Embedded TUI; Shared TUI does not transport delegated command submissions" + )); + } + let session_id = self.ensure_session(agent_type).await?; + self.submit_dialog_turn_request( + session_id, + prompt, + Some(original_command), + Vec::new(), + Vec::new(), + AgentDialogTurnExecution::FreshExternalSubagent { + ecosystem_id, + logical_id, + }, + agent_type, + ) + .await + } + + #[allow(clippy::too_many_arguments)] + async fn submit_dialog_turn_request( + &self, + session_id: String, + message: String, + original_message: Option, + workspace_references: Vec, + attachments: Vec, + execution: AgentDialogTurnExecution, + agent_type: &str, + ) -> Result { tracing::info!("Sending message to session {}: {}", session_id, message); // Generate a turn_id @@ -1204,8 +1256,9 @@ impl CliAgentRuntimeClient { let request = AgentDialogTurnRequest { session_id: session_id.clone(), message: message.clone(), - original_message: None, + original_message, turn_id: Some(turn_id.clone()), + execution, agent_type: agent_type.to_string(), // Dialog submission uses this path to locate persisted session // state. Execution still comes from the session's resolved binding. @@ -1871,6 +1924,28 @@ mod tests { assert!(!submission.contains("imagePath")); } + #[test] + fn delegated_external_commands_fail_before_shared_ipc() { + let source = include_str!("runtime_client.rs").replace("\r\n", "\n"); + let submission = source + .split_once("pub(crate) async fn send_external_subagent_command(") + .expect("delegated command submission method") + .1 + .split_once("async fn submit_dialog_turn_request(") + .expect("delegated command submission boundary") + .0; + + let shared_rejection = submission + .find("if self.is_shared()") + .expect("shared runtime rejection"); + let session_creation = submission + .find("let session_id = self.ensure_session") + .expect("session creation"); + assert!(shared_rejection < session_creation); + assert!(submission.contains("AgentDialogTurnExecution::FreshExternalSubagent")); + assert!(!submission.contains("RuntimeIpcOperation::SubmitTurn")); + } + #[test] fn interactive_session_fork_uses_the_same_runtime_boundary_in_both_deployments() { let source = include_str!("runtime_client.rs").replace("\r\n", "\n"); diff --git a/src/apps/cli/src/dispatch/worker.rs b/src/apps/cli/src/dispatch/worker.rs index 1e69de9339..8a03e3ec34 100644 --- a/src/apps/cli/src/dispatch/worker.rs +++ b/src/apps/cli/src/dispatch/worker.rs @@ -178,6 +178,7 @@ async fn run_inner(store: &DispatchStore, job_id: &str) -> Result<()> { message: prompt, original_message: None, turn_id: Some(turn_id.clone()), + execution: Default::default(), agent_type: job.request.agent_type.clone(), workspace_path: Some(workspace_path), remote_connection_id: None, diff --git a/src/apps/cli/src/modes/chat.rs b/src/apps/cli/src/modes/chat.rs index 7c7a0ade0e..c61a2b6c7d 100644 --- a/src/apps/cli/src/modes/chat.rs +++ b/src/apps/cli/src/modes/chat.rs @@ -90,9 +90,9 @@ use bitfun_core::external_sources::{ ExternalSubagentModelBindingMethod, ExternalSubagentModelBindingTarget, ExternalSubagentModelProfileRequest, ExternalSubagentModelRequest, ExternalToolActivationState, ExternalToolCapability, ExternalToolCatalogEntry, ExternalToolRuntimeKind, - NativePromptCommandDescriptor, PromptCommandAvailability, PromptCommandInvocationOutcome, - PromptCommandShellReviewDecision, PromptCommandShellReviewMode, PromptCommandShellReviewPlan, - EXTERNAL_SOURCE_CONTROL_SCHEMA_V1, + NativePromptCommandDescriptor, PromptCommandAvailability, PromptCommandExecutionTarget, + PromptCommandInvocationOutcome, PromptCommandShellReviewDecision, PromptCommandShellReviewMode, + PromptCommandShellReviewPlan, EXTERNAL_SOURCE_CONTROL_SCHEMA_V1, }; use bitfun_core::native_hooks::{ overview as native_hook_overview, NativeHookOverview, NativeHookRuleView, diff --git a/src/apps/cli/src/modes/chat/commands.rs b/src/apps/cli/src/modes/chat/commands.rs index 075461963e..60eb3cdf79 100644 --- a/src/apps/cli/src/modes/chat/commands.rs +++ b/src/apps/cli/src/modes/chat/commands.rs @@ -907,8 +907,34 @@ impl ChatMode { )) }); match expanded { - Ok(PromptCommandInvocationOutcome::Ready { content }) => { - self.send_message_to_agent(content, chat_view, chat_state, rt_handle); + Ok(PromptCommandInvocationOutcome::Ready { + content, + execution_target, + }) => { + match execution_target { + PromptCommandExecutionTarget::Inline => { + self.send_message_to_agent(content, chat_view, chat_state, rt_handle); + } + PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id, + logical_id, + } => { + let original_command = if invocation.arguments.trim().is_empty() { + format!("/{}", invocation.command_name) + } else { + format!("/{} {}", invocation.command_name, invocation.arguments) + }; + self.send_external_subagent_command_to_agent( + content, + original_command, + ecosystem_id.to_string(), + logical_id, + chat_view, + chat_state, + rt_handle, + ); + } + } Ok(None) } Ok(PromptCommandInvocationOutcome::ReviewRequired { review }) => { diff --git a/src/apps/cli/src/modes/chat/sessions.rs b/src/apps/cli/src/modes/chat/sessions.rs index 2d9bfbbd2c..15ccb924fd 100644 --- a/src/apps/cli/src/modes/chat/sessions.rs +++ b/src/apps/cli/src/modes/chat/sessions.rs @@ -209,6 +209,77 @@ impl ChatMode { ); } + #[allow(clippy::too_many_arguments)] + fn send_external_subagent_command_to_agent( + &mut self, + prompt: String, + original_command: String, + ecosystem_id: String, + logical_id: String, + chat_view: &mut ChatView, + chat_state: &mut ChatState, + rt_handle: &tokio::runtime::Handle, + ) { + let submitted_draft = crate::ui::composer::ComposerDraft { + text: original_command.clone(), + ..crate::ui::composer::ComposerDraft::default() + }; + if self.agent.is_shared() { + chat_view.set_status(Some( + "External subagent commands require Embedded TUI; Shared TUI does not transport delegated command submissions" + .to_string(), + )); + chat_view.set_draft(submitted_draft); + return; + } + if self + .pending_session_operation + .as_ref() + .is_some_and(|pending| pending.session_id == chat_state.core_session_id) + { + chat_view.set_status(Some( + "Waiting for the pending Session operation to finish before sending.".to_string(), + )); + chat_view.set_draft(submitted_draft); + return; + } + if chat_state.is_processing { + chat_state.add_system_message("Already processing, please wait.".to_string()); + chat_view.set_draft(submitted_draft); + return; + } + if let Err(error) = self.materialize_requested_worktree(chat_view, chat_state, rt_handle) { + tracing::error!("Failed to prepare worktree for delegated command: {error}"); + chat_view.set_status(Some(format!("Error: {error}"))); + chat_view.set_draft(submitted_draft); + return; + } + + let display_name = agent_display_name(&self.agent_type); + chat_view.set_status(Some(format!("{} is delegating...", display_name))); + let agent = self.agent.clone(); + let agent_type = self.agent_type.clone(); + match tokio::task::block_in_place(|| { + rt_handle.block_on(agent.send_external_subagent_command( + prompt, + original_command, + ecosystem_id, + logical_id, + &agent_type, + )) + }) { + Ok(turn_id) => { + tracing::info!("Started delegated command turn: {}", turn_id); + chat_view.remember_submitted_draft(&chat_state.core_session_id, &submitted_draft); + } + Err(error) => { + tracing::error!("Failed to delegate external command: {}", error); + chat_view.set_status(Some(format!("Error: {error}"))); + chat_view.set_draft(submitted_draft); + } + } + } + fn send_draft_to_agent( &mut self, draft: crate::ui::composer::ComposerDraft, diff --git a/src/apps/cli/src/modes/chat/tests.rs b/src/apps/cli/src/modes/chat/tests.rs index b2ab767765..865faefba8 100644 --- a/src/apps/cli/src/modes/chat/tests.rs +++ b/src/apps/cli/src/modes/chat/tests.rs @@ -1911,6 +1911,34 @@ mod tests { assert!(!source.contains("PendingSessionDelete")); } + #[test] + fn delegated_command_checks_session_state_before_materializing_a_worktree() { + let source = include_str!("sessions.rs").replace("\r\n", "\n"); + let submission = source + .split_once("fn send_external_subagent_command_to_agent(") + .expect("delegated command submission") + .1 + .split_once("fn send_draft_to_agent(") + .expect("delegated command submission boundary") + .0; + + let shared_guard = submission.find("if self.agent.is_shared()").unwrap(); + let pending_guard = submission.find("pending_session_operation").unwrap(); + let busy_guard = submission.find("if chat_state.is_processing").unwrap(); + let worktree_materialization = submission + .find("self.materialize_requested_worktree") + .unwrap(); + assert!(shared_guard < worktree_materialization); + assert!(pending_guard < worktree_materialization); + assert!(busy_guard < worktree_materialization); + assert!( + submission + .matches("chat_view.set_draft(submitted_draft)") + .count() + >= 4 + ); + } + #[test] fn workspace_diff_load_does_not_block_the_tui_event_loop() { let commands = include_str!("commands.rs").replace("\r\n", "\n"); diff --git a/src/apps/cli/src/peer_host/commands/dialog.rs b/src/apps/cli/src/peer_host/commands/dialog.rs index 49a956a239..323c8d0ce5 100644 --- a/src/apps/cli/src/peer_host/commands/dialog.rs +++ b/src/apps/cli/src/peer_host/commands/dialog.rs @@ -70,6 +70,7 @@ pub(crate) async fn start_dialog_turn( message: user_input, original_message: original_user_input, turn_id: Some(turn_id.clone()), + execution: Default::default(), agent_type, workspace_path, remote_connection_id, diff --git a/src/apps/desktop/src/api/agentic_api.rs b/src/apps/desktop/src/api/agentic_api.rs index 819fac723a..ee4a46bf82 100644 --- a/src/apps/desktop/src/api/agentic_api.rs +++ b/src/apps/desktop/src/api/agentic_api.rs @@ -16,9 +16,10 @@ use crate::runtime::{ use crate::startup_trace::DesktopStartupTrace; use bitfun_agent_runtime::deep_review::sanitize_focused_review_public_metadata; use bitfun_agent_runtime::sdk::{ - AgentDialogTurnRequest, AgentInputAttachment, AgentSessionCreateResult, - AgentSessionModelUpdateRequest, AgentSubmissionSource, AgentTurnCancellationRequest, - PermissionAuditRecord, PermissionGrant, PermissionGrantKey, PermissionReply, PermissionRequest, + AgentDialogTurnExecution, AgentDialogTurnRequest, AgentInputAttachment, + AgentSessionCreateResult, AgentSessionModelUpdateRequest, AgentSubmissionSource, + AgentTurnCancellationRequest, PermissionAuditRecord, PermissionGrant, PermissionGrantKey, + PermissionReply, PermissionRequest, }; use bitfun_core::agentic::agents::AgentSource; use bitfun_core::agentic::coordination::{ @@ -260,6 +261,8 @@ pub struct StartDialogTurnRequest { pub remote_ssh_host: Option, pub turn_id: Option, #[serde(default)] + pub execution: AgentDialogTurnExecution, + #[serde(default)] pub image_contexts: Option>, #[serde(default)] pub user_message_metadata: Option, @@ -1800,6 +1803,7 @@ fn desktop_dialog_turn_request( remote_connection_id, remote_ssh_host, turn_id, + execution, image_contexts, user_message_metadata, } = request; @@ -1819,6 +1823,7 @@ fn desktop_dialog_turn_request( message: user_input, original_message: original_user_input, turn_id, + execution, agent_type, workspace_path: project_workspace_path.or(workspace_path), remote_connection_id, diff --git a/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs b/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs index 41159abd7c..93146b4ddc 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs @@ -414,6 +414,7 @@ fn submit_turn_accepts_the_existing_64_kib_tui_paste_contract() { message: "x".repeat(64 * 1024), original_message: None, turn_id: Some("turn-1".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some("D:/workspace/project".to_string()), remote_connection_id: None, diff --git a/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs b/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs index d7c1d005f8..60278e3da0 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs @@ -532,6 +532,7 @@ fn submit_operation(workspace: &Path, session_id: &str, turn_id: &str) -> Runtim message: "hello".to_string(), original_message: None, turn_id: Some(turn_id.to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some(workspace.to_string_lossy().to_string()), remote_connection_id: None, diff --git a/src/crates/adapters/claude-code-adapter/src/command_source.rs b/src/crates/adapters/claude-code-adapter/src/command_source.rs index 404de8f1ad..a3fb55c758 100644 --- a/src/crates/adapters/claude-code-adapter/src/command_source.rs +++ b/src/crates/adapters/claude-code-adapter/src/command_source.rs @@ -749,6 +749,7 @@ fn command_definition( description: input.description, template: input.template, shell_preference, + execution_target: Default::default(), availability, content_version: format!("sha256:{content_version}"), }; diff --git a/src/crates/adapters/opencode-adapter/src/command_source.rs b/src/crates/adapters/opencode-adapter/src/command_source.rs index b62d627417..1a9219477a 100644 --- a/src/crates/adapters/opencode-adapter/src/command_source.rs +++ b/src/crates/adapters/opencode-adapter/src/command_source.rs @@ -6,10 +6,10 @@ use crate::local_source_paths::{ use bitfun_product_domains::external_sources::{ EcosystemId, ExternalSourceAssetKind, ExternalSourceContext, ExternalSourceDiagnostic, ExternalSourceHealth, ExternalSourceProviderError, ExternalSourceRecord, ExternalSourceScope, - ExternalWatchRoot, PromptCommandAvailability, PromptCommandDefinition, PromptCommandExpansion, - PromptCommandProviderIdentity, PromptCommandProviderSnapshot, PromptCommandShellExpansion, - PromptCommandShellInvocation, PromptCommandShellPreference, PromptCommandSourceProvider, - SourceKey, SourceQualifiedCommandId, + ExternalWatchRoot, PromptCommandAvailability, PromptCommandDefinition, + PromptCommandExecutionTarget, PromptCommandExpansion, PromptCommandProviderIdentity, + PromptCommandProviderSnapshot, PromptCommandShellExpansion, PromptCommandShellInvocation, + PromptCommandShellPreference, PromptCommandSourceProvider, SourceKey, SourceQualifiedCommandId, }; pub(crate) use bitfun_services_core::jsonc::strip_jsonc; use bitfun_services_core::markdown::{ @@ -954,7 +954,27 @@ fn command_definition( }) }); let content_version = command_content_version(&name, &input, shell_preference.as_ref()); - if input.agent.is_some() { + let execution_target = match ( + input.agent.as_deref(), + input.subtask, + input.model.as_ref(), + input.variant.as_ref(), + ) { + (Some(agent), None | Some(true), None, None) => { + PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id: EcosystemId::new(ECOSYSTEM_ID).map_err(|error| { + ExternalSourceProviderError::new( + "opencode.command.ecosystem_invalid", + error.to_string(), + false, + ) + })?, + logical_id: agent.to_string(), + } + } + _ => PromptCommandExecutionTarget::Inline, + }; + if input.agent.is_some() && execution_target.is_inline() { required_capabilities.push("command.agent".to_string()); } if input.model.is_some() { @@ -963,9 +983,12 @@ fn command_definition( if input.variant.is_some() { required_capabilities.push("command.variant".to_string()); } - if input.subtask.is_some() { + if input.subtask.is_some() && execution_target.is_inline() { required_capabilities.push("command.subtask".to_string()); } + if !execution_target.is_inline() && shell_preference.is_some() { + required_capabilities.push("command.external_subagent.shell".to_string()); + } if config_variable_regex().is_match(&input.template) { required_capabilities.push("command.config_variable".to_string()); } @@ -997,6 +1020,7 @@ fn command_definition( .unwrap_or_else(|| format!("OpenCode command /{name}")), template: input.template, shell_preference, + execution_target, availability, content_version, }; diff --git a/src/crates/adapters/opencode-adapter/tests/opencode_command_adapter.rs b/src/crates/adapters/opencode-adapter/tests/opencode_command_adapter.rs index 8e36dfb43a..7db1d6c615 100644 --- a/src/crates/adapters/opencode-adapter/tests/opencode_command_adapter.rs +++ b/src/crates/adapters/opencode-adapter/tests/opencode_command_adapter.rs @@ -1,8 +1,8 @@ use bitfun_opencode_adapter::{OpenCodeCommandProvider, OpenCodeCommandProviderOptions}; use bitfun_product_domains::external_sources::{ - ExecutionDomainId, ExternalSourceContext, ExternalSourceHealth, PromptCommandAvailability, - PromptCommandDefinition, PromptCommandProviderSnapshot, PromptCommandShellPreference, - PromptCommandSourceProvider, + EcosystemId, ExecutionDomainId, ExternalSourceContext, ExternalSourceHealth, + PromptCommandAvailability, PromptCommandDefinition, PromptCommandExecutionTarget, + PromptCommandProviderSnapshot, PromptCommandShellPreference, PromptCommandSourceProvider, }; use std::collections::BTreeSet; use std::fs; @@ -346,7 +346,7 @@ fn invalid_source_is_diagnostic_and_does_not_remove_other_valid_sources() { } #[test] -fn unsupported_expansion_features_restrict_the_whole_command() { +fn supported_subagent_delegation_is_typed_and_other_unsupported_features_stay_restricted() { let fixture = Fixture::new(); write( fixture.user_config.join("opencode.json"), @@ -359,7 +359,12 @@ fn unsupported_expansion_features_restrict_the_whole_command() { "unsafe-file": {"template":"Review @/etc/passwd"}, "config-var": {"template":"Review {env:HOME}"}, "agent": {"template":"Delegate this", "agent":"explore"}, - "subtask": {"template":"Delegate this", "subtask":true} + "agent-subtask": {"template":"Delegate this", "agent":"explore", "subtask":true}, + "agent-shell": {"template":"Run !`git status` then delegate", "agent":"explore"}, + "primary-agent": {"template":"Stay primary", "agent":"build", "subtask":false}, + "subtask": {"template":"Delegate this", "subtask":true}, + "model": {"template":"Delegate this", "agent":"explore", "model":"provider/model"}, + "variant": {"template":"Delegate this", "agent":"explore", "variant":"high"} } }"#, ); @@ -374,8 +379,11 @@ fn unsupported_expansion_features_restrict_the_whole_command() { "dynamic-file", "unsafe-file", "config-var", - "agent", + "agent-shell", + "primary-agent", "subtask", + "model", + "variant", ] { let command = snapshot .commands @@ -388,6 +396,44 @@ fn unsupported_expansion_features_restrict_the_whole_command() { )); assert!(provider.expand(&fixture.context(), command, "").is_err()); } + let agent_shell = snapshot + .commands + .iter() + .find(|command| command.name == "agent-shell") + .unwrap(); + let PromptCommandAvailability::Restricted { + required_capabilities, + .. + } = &agent_shell.availability + else { + panic!("shell-backed subagent delegation must be restricted") + }; + assert!(required_capabilities.contains(&"command.external_subagent.shell".to_string())); + for name in ["agent", "agent-subtask"] { + let command = snapshot + .commands + .iter() + .find(|command| command.name == name) + .unwrap(); + assert!(matches!( + command.availability, + PromptCommandAvailability::Available + )); + assert_eq!( + command.execution_target, + PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id: EcosystemId::new("opencode").unwrap(), + logical_id: "explore".to_string(), + } + ); + assert_eq!( + provider + .expand(&fixture.context(), command, "") + .unwrap() + .content, + "Delegate this" + ); + } let shell = snapshot .commands .iter() diff --git a/src/crates/assembly/core/src/agentic/agents/registry/external.rs b/src/crates/assembly/core/src/agentic/agents/registry/external.rs index d9420a220c..07a084ab06 100644 --- a/src/crates/assembly/core/src/agentic/agents/registry/external.rs +++ b/src/crates/assembly/core/src/agentic/agents/registry/external.rs @@ -3,6 +3,7 @@ use super::AgentRegistry; use crate::agentic::agents::{Agent, SubagentVisibilityPolicy}; use bitfun_agent_runtime::prompt_cache::prompt_cache_scope_key; use bitfun_core_types::{SessionContinuationPolicy, SessionModelBindingPolicy}; +use bitfun_product_domains::external_sources::EcosystemId; use std::collections::{BTreeMap, HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::sync::{Arc, RwLock, Weak}; @@ -45,6 +46,7 @@ impl ExternalSubagentModelBinding { pub struct ExternalSubagentRegistration { pub runtime_key: String, pub logical_id: String, + pub ecosystem_id: EcosystemId, pub provider_label: String, pub model_binding: ExternalSubagentModelBinding, pub hidden: bool, @@ -136,9 +138,18 @@ impl ExternalSubagentRegistryState { .retain(|runtime_key, entry| entry.lease_count > 0 || routed.contains(runtime_key)); } - fn acquire(self: &Arc, runtime_key: &str) -> Option { + fn acquire_matching( + self: &Arc, + runtime_key: &str, + expected_ecosystem_id: Option<&EcosystemId>, + ) -> Option { let mut generations = self.write_generations(); let entry = generations.get_mut(runtime_key)?; + if expected_ecosystem_id + .is_some_and(|expected| expected != &entry.registration.ecosystem_id) + { + return None; + } entry.lease_count = entry.lease_count.saturating_add(1); Some(ExternalSubagentInvocationBinding { runtime_agent_key: runtime_key.to_string(), @@ -154,6 +165,10 @@ impl ExternalSubagentRegistryState { }) } + fn acquire(self: &Arc, runtime_key: &str) -> Option { + self.acquire_matching(runtime_key, None) + } + fn release(&self, runtime_key: &str) { if let Some(entry) = self.write_generations().get_mut(runtime_key) { entry.lease_count = entry.lease_count.saturating_sub(1); @@ -341,6 +356,31 @@ impl AgentRegistry { .map(|entry| local_binding(logical_id, entry.agent.id())) } + /// Resolve only the currently approved external route for an exact + /// ecosystem. Command delegation must never fall back to a same-name local + /// agent or cross an ecosystem boundary after the command was expanded. + pub fn resolve_external_subagent_for_fresh_invocation( + &self, + logical_id: &str, + ecosystem_id: &EcosystemId, + workspace_root: Option<&Path>, + ) -> Option { + let workspace_root = workspace_root?; + let logical_key = normalize_external_logical_id(logical_id); + let route = self + .external_subagents + .read_routes() + .get(workspace_root) + .and_then(|routes| routes.get(&logical_key)) + .cloned()?; + match route { + ExternalSubagentRoute::External(runtime_key) => self + .external_subagents + .acquire_matching(&runtime_key, Some(ecosystem_id)), + ExternalSubagentRoute::Local | ExternalSubagentRoute::Unavailable => None, + } + } + pub(super) fn apply_external_routes_to_query( &self, workspace_root: &Path, diff --git a/src/crates/assembly/core/src/agentic/agents/registry/tests.rs b/src/crates/assembly/core/src/agentic/agents/registry/tests.rs index 817732e998..e40c2a1c62 100644 --- a/src/crates/assembly/core/src/agentic/agents/registry/tests.rs +++ b/src/crates/assembly/core/src/agentic/agents/registry/tests.rs @@ -17,6 +17,7 @@ use bitfun_agent_runtime::custom_agent::{ CustomAgentKind, CustomAgentLevel, }; use bitfun_agent_runtime::sdk::{RuntimeAgentRegistry, RuntimeAgentRegistryQuery}; +use bitfun_product_domains::external_sources::EcosystemId; use std::collections::{BTreeMap, HashMap}; use std::path::{Path, PathBuf}; use std::sync::Arc; @@ -1160,6 +1161,7 @@ async fn external_routes_are_workspace_scoped_fail_closed_and_generation_leased( vec![ExternalSubagentRegistration { runtime_key: runtime_v1.to_string(), logical_id: "Explore".to_string(), + ecosystem_id: EcosystemId::new("opencode").unwrap(), provider_label: "OpenCode".to_string(), model_binding: super::ExternalSubagentModelBinding::Fixed { model_id: "inherit".to_string(), @@ -1214,6 +1216,29 @@ async fn external_routes_are_workspace_scoped_fail_closed_and_generation_leased( assert!(registry.is_external_subagent_route("EXPLORE", Some(&workspace))); assert!(!registry.is_external_subagent_route("Explore", None)); assert!(!registry.is_external_subagent_route("Explore", Some(Path::new("C:/workspace/other")))); + assert!(registry + .resolve_external_subagent_for_fresh_invocation( + "Explore", + &EcosystemId::new("claude-code").unwrap(), + Some(&workspace), + ) + .is_none()); + assert!(registry + .resolve_external_subagent_for_fresh_invocation( + "Explore", + &EcosystemId::new("opencode").unwrap(), + None, + ) + .is_none()); + let command_binding = registry + .resolve_external_subagent_for_fresh_invocation( + "Explore", + &EcosystemId::new("opencode").unwrap(), + Some(&workspace), + ) + .expect("command delegation resolves only the exact external ecosystem route"); + assert_eq!(command_binding.runtime_agent_key, runtime_v1); + drop(command_binding); let binding = registry .resolve_subagent_for_fresh_invocation("Explore", Some(&workspace), true) @@ -1240,6 +1265,7 @@ async fn external_routes_are_workspace_scoped_fail_closed_and_generation_leased( vec![ExternalSubagentRegistration { runtime_key: runtime_v2.to_string(), logical_id: "Explore".to_string(), + ecosystem_id: EcosystemId::new("opencode").unwrap(), provider_label: "OpenCode".to_string(), model_binding: super::ExternalSubagentModelBinding::Fixed { model_id: "inherit".to_string(), diff --git a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs index 03a120b4c8..33730b28f4 100644 --- a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs +++ b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs @@ -17,7 +17,7 @@ use crate::agentic::context_profile::ContextProfilePolicy; use crate::agentic::core::{ InternalReminderKind, Message, MessageContent, MessageSemanticKind, ProcessingPhase, Session, SessionConfig, SessionContinuationPolicy, SessionKind, SessionModelBindingPolicy, SessionState, - SessionSummary, TurnStats, + SessionSummary, ToolCall, ToolResult, TurnStats, }; use crate::agentic::events::{ AgenticEvent, DeepReviewQueueState, EventPriority, EventQueue, EventRouter, EventSubscriber, @@ -85,6 +85,8 @@ use bitfun_agent_runtime::remote_file_delivery::{ }; use bitfun_agent_runtime::sdk::PermissionReply; use bitfun_agent_runtime::user_questions::USER_INPUT_AVAILABLE_CONTEXT_KEY; +use bitfun_events::{ToolEventData, ToolEventIdentity}; +use bitfun_product_domains::external_sources::EcosystemId; use bitfun_runtime_ports::{ agent_workspace_references_from_metadata, AgentMessageWorkspaceReferencesRequest, AgentSessionComposerUpdate, AgentSessionWorkspaceBinding, AgentThreadGoalDeliveryKind, @@ -113,6 +115,7 @@ use tokio_util::sync::CancellationToken; const MANUAL_COMPACTION_COMMAND: &str = "/compact"; const CONTEXT_COMPRESSION_TOOL_NAME: &str = "ContextCompression"; +const TASK_TOOL_NAME: &str = "Task"; const DEFAULT_SUBAGENT_MAX_CONCURRENCY: usize = 5; const MAX_SUBAGENT_MAX_CONCURRENCY: usize = 64; const SUBAGENT_TIMEOUT_GRACE_PERIOD: Duration = Duration::from_secs(10); @@ -3163,6 +3166,585 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet .await } + /// Execute a statically discovered external command through the existing + /// fresh-subagent owner while preserving a normal parent UserDialog/Task + /// transcript. The command source selects the target; no model routing or + /// same-name local fallback is performed here. + #[allow(clippy::too_many_arguments)] + pub(crate) fn start_external_subagent_delegation_turn( + self: &Arc, + session_id: String, + prompt: String, + original_user_input: Option, + requested_turn_id: Option, + agent_type: String, + workspace_path: Option, + _submission_policy: DialogSubmissionPolicy, + extra_user_message_metadata: Option, + ecosystem_id: String, + logical_id: String, + ) -> futures::future::BoxFuture<'_, BitFunResult<()>> { + Box::pin(async move { + bitfun_core_types::validate_session_id(&session_id).map_err(BitFunError::Validation)?; + if prompt.trim().is_empty() { + return Err(BitFunError::Validation( + "External subagent delegation prompt must not be empty".to_string(), + )); + } + let ecosystem_id = EcosystemId::new(ecosystem_id).map_err(|error| { + BitFunError::Validation(format!( + "Invalid external subagent delegation ecosystem: {error}" + )) + })?; + let logical_id = logical_id.trim().to_string(); + if logical_id.is_empty() { + return Err(BitFunError::Validation( + "External subagent delegation logical_id must not be empty".to_string(), + )); + } + + let mut session = self + .session_manager + .get_session(&session_id) + .ok_or_else(|| BitFunError::NotFound(format!("Session not found: {session_id}")))?; + self.ensure_session_runtime_ownership(&session_id, None)?; + if session.config.remote_connection_id.is_some() + || session.config.remote_ssh_host.is_some() + { + return Err(BitFunError::NotImplemented( + "External subagent command delegation is unavailable for remote workspaces" + .to_string(), + )); + } + if !matches!(session.state, SessionState::Idle) { + return Err(BitFunError::Validation(format!( + "Session must be idle before external subagent command delegation: {:?}", + session.state + ))); + } + if self + .wait_session_drained(&session_id, Duration::from_millis(800)) + .await + > 0 + { + return Err(BitFunError::Validation(format!( + "Previous dialog turn is still draining: session_id={session_id}" + ))); + } + + let project_workspace_path = session + .config + .project_workspace_path + .clone() + .or_else(|| session.config.workspace_path.clone()) + .or(workspace_path) + .ok_or_else(|| { + BitFunError::Validation(format!( + "Session workspace_path is missing: {session_id}" + )) + })?; + let execution_workspace_path = session + .config + .workspace_path + .clone() + .unwrap_or_else(|| project_workspace_path.clone()); + + let context_messages = self + .session_manager + .get_context_messages(&session_id) + .await?; + if (context_messages.is_empty() + || (context_messages.len() == 1 && !session.dialog_turn_ids.is_empty())) + && !session.dialog_turn_ids.is_empty() + { + let restore_path = + Self::resolve_session_restore_path(&project_workspace_path, None, None).await?; + self.restore_session_from_storage_path(&restore_path, &session_id) + .await?; + session = self + .session_manager + .get_session(&session_id) + .ok_or_else(|| { + BitFunError::NotFound(format!("Session not found: {session_id}")) + })?; + } + + let binding = get_agent_registry() + .resolve_external_subagent_for_fresh_invocation( + &logical_id, + &ecosystem_id, + Some(Path::new(&project_workspace_path)), + ) + .ok_or_else(|| { + BitFunError::Validation(format!( + "candidate_unavailable: approved external subagent {}:{} changed before the command could start", + ecosystem_id, logical_id + )) + })?; + let external_generation_lease = binding.lease.ok_or_else(|| { + BitFunError::Validation( + "Approved external subagent route is missing its generation lease".to_string(), + ) + })?; + + let effective_agent_type = Self::normalize_agent_type(agent_type.trim()); + let permission_runtime_ceiling = + crate::agentic::permission_policy::load_parent_permission_runtime_ceiling(Some( + &effective_agent_type, + )) + .await?; + if session.agent_type != effective_agent_type { + self.session_manager + .update_session_agent_type(&session_id, &effective_agent_type) + .await?; + } + let display_input = original_user_input + .filter(|input| !input.trim().is_empty()) + .unwrap_or_else(|| prompt.clone()); + let mut user_message_metadata = + Self::ensure_user_message_metadata_object(extra_user_message_metadata); + if let Some(metadata) = user_message_metadata.as_object_mut() { + if display_input != prompt { + metadata.insert( + "original_text".to_string(), + serde_json::Value::String(display_input.clone()), + ); + } + metadata.insert( + "externalCommandDelegation".to_string(), + serde_json::json!({ + "ecosystemId": ecosystem_id.as_str(), + "logicalId": logical_id, + }), + ); + } + let turn_index = self.session_manager.get_turn_count(&session_id); + let turn_id = self + .session_manager + .start_dialog_turn( + &session_id, + effective_agent_type.clone(), + prompt.clone(), + requested_turn_id, + None, + Some(user_message_metadata.clone()), + ) + .await?; + let execution_lease = self.register_session_execution(&session_id); + let turn_settlement_registration = self + .turn_settlements + .register_accepted(session_id.clone(), turn_id.clone()); + let cancellation_token = CancellationToken::new(); + self.execution_engine + .register_cancel_token(&turn_id, cancellation_token.clone()); + if let Err(error) = self + .session_manager + .update_session_state_for_turn_if_processing( + &session_id, + &turn_id, + SessionState::Processing { + current_turn_id: turn_id.clone(), + phase: ProcessingPhase::ToolCalling, + }, + ) + .await + { + warn!( + "Failed to persist delegated command ToolCalling phase: session_id={}, turn_id={}, error={}", + session_id, turn_id, error + ); + } + let round_id = format!("{}-round-0", turn_id); + let tool_call_id = format!("task_{}", uuid::Uuid::new_v4()); + let tool_params = serde_json::json!({ + "action": "spawn", + "description": format!("Run external command with {logical_id}"), + "prompt": prompt, + "subagent_type": logical_id, + }); + + self.emit_event(AgenticEvent::DialogTurnStarted { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + turn_index, + user_input: prompt.clone(), + original_user_input: (display_input != prompt).then_some(display_input), + user_message_metadata: Some(user_message_metadata.clone()), + }) + .await; + self.emit_event(AgenticEvent::ModelRoundStarted { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + round_id: round_id.clone(), + round_group_id: None, + round_index: 0, + model_config_id: String::new(), + effective_model_name: String::new(), + }) + .await; + self.emit_event(AgenticEvent::ToolEvent { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + round_id: round_id.clone(), + attempt_id: None, + attempt_index: None, + tool_event: ToolEventData::Started { + identity: ToolEventIdentity::direct(tool_call_id.clone(), TASK_TOOL_NAME), + params: tool_params.clone(), + timeout_seconds: None, + }, + }) + .await; + + let mut child_context = HashMap::new(); + for key in [ + USER_INPUT_AVAILABLE_CONTEXT_KEY, + AUTO_APPROVE_ASK_CONTEXT_KEY, + ] { + if let Some(value) = metadata_bool(Some(&user_message_metadata), key) { + child_context.insert(key.to_string(), value.to_string()); + } + } + let request = SubagentExecutionRequest { + task_description: prompt.clone(), + context_mode: SubagentContextMode::Fresh, + target_session_id: None, + subagent_type: Some(binding.runtime_agent_key), + logical_subagent_type: Some(binding.logical_id), + continuation_policy: binding.continuation_policy, + model_binding_policy: binding.model_binding_policy, + workspace_path: Some(execution_workspace_path), + model_id: None, + inherit_parent_model: false, + subagent_parent_info: SubagentParentInfo { + tool_call_id: tool_call_id.clone(), + session_id: session_id.clone(), + dialog_turn_id: turn_id.clone(), + }, + context: child_context, + permission_runtime_ceiling, + delegation_policy: DelegationPolicy::top_level().spawn_child(), + external_generation_lease: Some(external_generation_lease), + }; + + let coordinator = Arc::clone(self); + tokio::spawn(async move { + let _execution_lease = execution_lease; + let _turn_settlement_registration = turn_settlement_registration; + let _cancel_guard = CancelTokenGuard { + execution_engine: Arc::clone(&coordinator.execution_engine), + dialog_turn_id: turn_id.clone(), + }; + let started_at = Instant::now(); + let execution_result = coordinator + .execute_subagent(request, Some(&cancellation_token), None) + .await; + let duration_ms = started_at.elapsed().as_millis() as u64; + let ( + result_data, + result_for_assistant, + is_error, + cancelled, + child_session_id, + failure_error, + ) = match execution_result { + Ok(result) => { + let child_session_id = result.session_id().map(str::to_string); + let delegate_target_label = format!("subagent '{}'", logical_id); + let (data, assistant_text) = + bitfun_agent_runtime::subagent_task::subagent_task_completion_result( + bitfun_agent_runtime::subagent_task::SubagentTaskCompletionResultInput { + delegate_target_label: &delegate_target_label, + result_text: &result.text, + context_mode: SubagentContextMode::Fresh.as_str(), + duration_ms: duration_ms as u128, + is_partial_timeout: result.is_partial_timeout(), + reason: result.reason.as_deref(), + ledger_event_id: result.ledger_event_id(), + partial_timeout_suffix: "", + }, + ); + coordinator + .emit_event(AgenticEvent::ToolEvent { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + round_id: round_id.clone(), + attempt_id: None, + attempt_index: None, + tool_event: ToolEventData::Completed { + identity: ToolEventIdentity::direct( + tool_call_id.clone(), + TASK_TOOL_NAME, + ), + result: data.clone(), + result_for_assistant: Some(assistant_text.clone()), + image_attachments: None, + duration_ms, + queue_wait_ms: None, + preflight_ms: None, + confirmation_wait_ms: None, + execution_ms: Some(duration_ms), + }, + }) + .await; + (data, assistant_text, false, false, child_session_id, None) + } + Err(error) => { + let cancelled = matches!(error, BitFunError::Cancelled(_)); + let error_text = error.to_string(); + let tool_event = if cancelled { + ToolEventData::Cancelled { + identity: ToolEventIdentity::direct( + tool_call_id.clone(), + TASK_TOOL_NAME, + ), + reason: error_text.clone(), + duration_ms: Some(duration_ms), + queue_wait_ms: None, + preflight_ms: None, + confirmation_wait_ms: None, + execution_ms: Some(duration_ms), + } + } else { + ToolEventData::Failed { + identity: ToolEventIdentity::direct( + tool_call_id.clone(), + TASK_TOOL_NAME, + ), + error: error_text.clone(), + duration_ms: Some(duration_ms), + queue_wait_ms: None, + preflight_ms: None, + confirmation_wait_ms: None, + execution_ms: Some(duration_ms), + } + }; + coordinator + .emit_event(AgenticEvent::ToolEvent { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + round_id: round_id.clone(), + attempt_id: None, + attempt_index: None, + tool_event, + }) + .await; + ( + serde_json::json!({ "error": error_text }), + error_text, + true, + cancelled, + None, + (!cancelled).then_some(error), + ) + } + }; + + let assistant_message = Message::assistant_with_tools( + String::new(), + vec![ToolCall { + tool_id: tool_call_id.clone(), + tool_name: TASK_TOOL_NAME.to_string(), + arguments: tool_params, + raw_arguments: None, + is_error: false, + parse_error: None, + recovered_from_truncation: false, + repair_kind: Default::default(), + }], + ) + .with_turn_id(turn_id.clone()) + .with_round_id(round_id.clone()); + let tool_result_message = Message::tool_result(ToolResult { + tool_id: tool_call_id.clone(), + tool_name: TASK_TOOL_NAME.to_string(), + effective_tool_name: None, + result: result_data, + result_for_assistant: Some(result_for_assistant), + is_error, + duration_ms: Some(duration_ms), + image_attachments: None, + }) + .with_turn_id(turn_id.clone()) + .with_round_id(round_id.clone()); + let new_messages = vec![assistant_message, tool_result_message]; + for message in &new_messages { + if let Err(error) = coordinator + .session_manager + .add_message(&session_id, message.clone()) + .await + { + error!( + "Failed to append delegated command Task message: session_id={}, turn_id={}, error={}", + session_id, turn_id, error + ); + } + } + + let completed_at = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64; + let mut rounds = SessionManager::build_model_rounds_from_messages( + &new_messages, + &turn_id, + completed_at, + ); + let child_dialog_turn_id = child_session_id.as_ref().and_then(|child_session_id| { + coordinator + .session_manager + .get_session(child_session_id) + .and_then(|session| session.dialog_turn_ids.last().cloned()) + }); + let child_model_id = child_session_id.as_ref().and_then(|child_session_id| { + coordinator + .session_manager + .get_session(child_session_id) + .and_then(|session| session.config.model_id) + }); + if let Some(tool_item) = rounds + .iter_mut() + .flat_map(|round| round.tool_items.iter_mut()) + .find(|item| item.id == tool_call_id) + { + tool_item.subagent_session_id = child_session_id; + tool_item.subagent_dialog_turn_id = child_dialog_turn_id; + tool_item.subagent_model_id = child_model_id; + tool_item.duration_ms = Some(duration_ms); + tool_item.execution_ms = Some(duration_ms); + if let Some(result) = tool_item.tool_result.as_mut() { + result.duration_ms = Some(duration_ms); + } + tool_item.status = + Some(if is_error { "error" } else { "completed" }.to_string()); + } + let turn_persistence = if let Some(error) = failure_error.as_ref() { + coordinator + .session_manager + .fail_synthetic_dialog_turn( + &session_id, + &turn_id, + error.to_string(), + rounds, + ) + .await + } else { + coordinator + .session_manager + .complete_synthetic_dialog_turn(&session_id, &turn_id, rounds, duration_ms) + .await + }; + if let Err(error) = turn_persistence { + error!( + "Failed to persist delegated external command turn: session_id={}, turn_id={}, error={}", + session_id, turn_id, error + ); + } + if cancelled { + let _ = coordinator + .session_manager + .cancel_dialog_turn(&session_id, &turn_id) + .await; + } + let final_session_state = if let Some(error) = failure_error.as_ref() { + SessionState::Error { + error: error.to_string(), + recoverable: !matches!( + error, + BitFunError::AIClient(_) | BitFunError::Timeout(_) + ), + } + } else { + SessionState::Idle + }; + let _ = coordinator + .session_manager + .update_session_state_for_turn_if_processing( + &session_id, + &turn_id, + final_session_state, + ) + .await; + coordinator + .emit_event(AgenticEvent::ModelRoundCompleted { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + round_id, + has_tool_calls: true, + duration_ms: Some(duration_ms), + provider_id: None, + model_config_id: String::new(), + effective_model_name: String::new(), + first_chunk_ms: None, + first_visible_output_ms: None, + stream_duration_ms: None, + attempt_count: None, + failure_category: is_error.then_some("tool_error".to_string()), + token_details: None, + }) + .await; + + if cancelled { + coordinator + .emit_event(AgenticEvent::DialogTurnCancelled { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + }) + .await; + } else if let Some(error) = failure_error.as_ref() { + coordinator + .emit_event(AgenticEvent::DialogTurnFailed { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + error: error.to_string(), + error_category: Some(error.error_category()), + error_detail: Some(error.error_detail()), + }) + .await; + } else { + coordinator + .emit_event(AgenticEvent::DialogTurnCompleted { + session_id: session_id.clone(), + turn_id: turn_id.clone(), + total_rounds: 1, + total_tools: 1, + duration_ms, + partial_recovery_reason: None, + success: Some(true), + finish_reason: Some("complete".to_string()), + has_final_response: Some(false), + }) + .await; + } + if let Some(tx) = coordinator.scheduler_notify_tx.get() { + let outcome = if cancelled { + TurnOutcome::Cancelled { + turn_id: turn_id.clone(), + } + } else if let Some(error) = failure_error.as_ref() { + TurnOutcome::Failed { + turn_id: turn_id.clone(), + error: error.to_string(), + } + } else { + TurnOutcome::Completed { + turn_id: turn_id.clone(), + final_response: String::new(), + } + }; + if let Err(error) = tx.try_send((session_id.clone(), outcome)) { + error!( + "Failed to notify scheduler of delegated command settlement: session_id={}, turn_id={}, error={}", + session_id, turn_id, error + ); + } + } + }); + + Ok(()) + }) + } + #[allow(clippy::too_many_arguments)] pub async fn start_dialog_turn_with_prepended_messages( &self, diff --git a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs index 6cebf6d008..1354bc8786 100644 --- a/src/crates/assembly/core/src/agentic/coordination/scheduler.rs +++ b/src/crates/assembly/core/src/agentic/coordination/scheduler.rs @@ -57,11 +57,12 @@ use bitfun_agent_runtime::scheduler::{ }; use bitfun_runtime_ports::{ resolve_dialog_submit_queue_action, AgentBackgroundResultRequest, AgentDialogPrependedReminder, - AgentDialogTurnPort, AgentDialogTurnRequest, AgentInputAttachment, AgentLifecycleDeliveryPort, - AgentThreadGoalDeliveryKind, AgentThreadGoalDeliveryRequest, AgentTurnCancellationPort, - AgentTurnCancellationRequest, AgentTurnCancellationResult, DialogSessionStateFact, - DialogSubmitQueueAction, DialogSubmitQueueFacts, PortError, PortErrorKind, PortResult, - RoundInjection, RoundInjectionKind, SessionStoragePathRequest, SessionStorePort, + AgentDialogTurnExecution, AgentDialogTurnPort, AgentDialogTurnRequest, AgentInputAttachment, + AgentLifecycleDeliveryPort, AgentThreadGoalDeliveryKind, AgentThreadGoalDeliveryRequest, + AgentTurnCancellationPort, AgentTurnCancellationRequest, AgentTurnCancellationResult, + DialogSessionStateFact, DialogSubmitQueueAction, DialogSubmitQueueFacts, PortError, + PortErrorKind, PortResult, RoundInjection, RoundInjectionKind, SessionStoragePathRequest, + SessionStorePort, }; pub use bitfun_runtime_ports::{ AgentSessionReplyRoute, DialogQueuePriority, DialogSteerOutcome, DialogSubmissionPolicy, @@ -101,9 +102,16 @@ impl QueuedTurn { pub(crate) enum QueuedTurnExecution { #[default] Standard, + FreshExternalSubagent(ExternalSubagentDelegationQueuedExecution), HiddenSubagent(HiddenSubagentQueuedExecution), } +#[derive(Debug, Clone)] +pub(crate) struct ExternalSubagentDelegationQueuedExecution { + ecosystem_id: String, + logical_id: String, +} + fn remove_queued_turn_by_id( queues: &DialogTurnQueue, session_id: &str, @@ -1035,6 +1043,16 @@ impl DialogScheduler { }; let queue_has_items = self.queues.has_items(&session_id); + if matches!( + &queued_turn.execution, + QueuedTurnExecution::FreshExternalSubagent(_) + ) && (!matches!(&state_fact, DialogSessionStateFact::Idle) || queue_has_items) + { + return Err(SchedulerSubmitError::Core(BitFunError::Validation( + "External subagent delegation requires an idle session with an empty queue" + .to_string(), + ))); + } let action = resolve_dialog_submit_queue_action(DialogSubmitQueueFacts { session_state: state_fact, queue_has_items, @@ -1141,7 +1159,7 @@ impl DialogScheduler { async fn finish_removed_queued_turn(&self, session_id: &str, removed_turn: QueuedTurn) { match removed_turn.execution { - QueuedTurnExecution::Standard => { + QueuedTurnExecution::Standard | QueuedTurnExecution::FreshExternalSubagent(_) => { if let Some(turn_id) = removed_turn.turn_id { self.coordinator .emit_event(AgenticEvent::DialogTurnCancelled { @@ -1398,7 +1416,7 @@ impl DialogScheduler { let mut retired_turn_ids = Vec::new(); for queued_turn in cleared_turns { match queued_turn.execution { - QueuedTurnExecution::Standard => { + QueuedTurnExecution::Standard | QueuedTurnExecution::FreshExternalSubagent(_) => { if let Some(turn_id) = queued_turn.turn_id { retired_turn_ids.push(turn_id.clone()); self.coordinator @@ -1491,11 +1509,55 @@ impl DialogScheduler { session_id: &str, queued_turn: &QueuedTurn, ) -> Result { - if let QueuedTurnExecution::HiddenSubagent(execution) = &queued_turn.execution { - return self - .start_hidden_subagent_turn(session_id, queued_turn, execution) - .await - .map_err(SchedulerSubmitError::Message); + match &queued_turn.execution { + QueuedTurnExecution::HiddenSubagent(execution) => { + return self + .start_hidden_subagent_turn(session_id, queued_turn, execution) + .await + .map_err(SchedulerSubmitError::Message); + } + QueuedTurnExecution::FreshExternalSubagent(execution) => { + self.coordinator + .start_external_subagent_delegation_turn( + session_id.to_string(), + queued_turn.user_input.clone(), + queued_turn.original_user_input.clone(), + queued_turn.turn_id.clone(), + queued_turn.agent_type.clone(), + queued_turn.workspace_path.clone(), + queued_turn.policy, + queued_turn.user_message_metadata.clone(), + execution.ecosystem_id.clone(), + execution.logical_id.clone(), + ) + .await + .map_err(SchedulerSubmitError::Core)?; + + let resolved = queued_turn.turn_id.clone().ok_or_else(|| { + format!( + "Scheduled external subagent delegation is missing turn_id: session_id={session_id}" + ) + })?; + self.active_turns.insert( + session_id, + ActiveDialogTurn::new( + resolved.clone(), + queued_turn.workspace_path.clone(), + None, + None, + queued_turn.agent_type.clone(), + queued_turn + .original_user_input + .clone() + .unwrap_or_else(|| queued_turn.user_input.clone()), + queued_turn.user_message_metadata.clone(), + queued_turn.policy, + queued_turn.reply_route.clone(), + ), + ); + return Ok(resolved); + } + QueuedTurnExecution::Standard => {} } let images = queued_turn @@ -2167,6 +2229,41 @@ impl DialogScheduler { request: AgentDialogTurnRequest, reject_if_busy: bool, ) -> PortResult { + let (execution, reject_if_busy) = match &request.execution { + AgentDialogTurnExecution::Standard => (QueuedTurnExecution::Standard, reject_if_busy), + AgentDialogTurnExecution::FreshExternalSubagent { + ecosystem_id, + logical_id, + } => { + if ecosystem_id.trim().is_empty() || logical_id.trim().is_empty() { + return Err(PortError::new( + PortErrorKind::InvalidRequest, + "External subagent delegation requires non-empty ecosystem_id and logical_id", + )); + } + if !request.attachments.is_empty() || !request.prepended_reminders.is_empty() { + return Err(PortError::new( + PortErrorKind::InvalidRequest, + "External subagent delegation does not accept attachments or prepended reminders", + )); + } + if request.remote_connection_id.is_some() || request.remote_ssh_host.is_some() { + return Err(PortError::new( + PortErrorKind::NotAvailable, + "External subagent delegation is unavailable for remote workspaces", + )); + } + ( + QueuedTurnExecution::FreshExternalSubagent( + ExternalSubagentDelegationQueuedExecution { + ecosystem_id: ecosystem_id.trim().to_string(), + logical_id: logical_id.trim().to_string(), + }, + ), + true, + ) + } + }; let image_contexts = agent_dialog_turn_image_contexts(&request.attachments)?; let prepended_messages = agent_dialog_turn_prepended_messages(&request.prepended_reminders)?; @@ -2205,7 +2302,7 @@ impl DialogScheduler { image_contexts, enqueued_at: SystemTime::now(), _settlement_registration: Some(settlement_registration), - execution: QueuedTurnExecution::Standard, + execution, }; self.submit_queued_turn( @@ -2780,6 +2877,7 @@ mod tests { message: "hello".to_string(), original_message: None, turn_id: Some("missing-turn".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some(workspace.to_string_lossy().to_string()), remote_connection_id: None, @@ -2844,6 +2942,7 @@ mod tests { message: "queued prompt".to_string(), original_message: None, turn_id: Some(turn_id.to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: None, remote_connection_id: None, @@ -2883,6 +2982,154 @@ mod tests { .expect("cancelled queued turn should settle"); } + #[tokio::test] + async fn delegated_dialog_turn_rejects_instead_of_queueing_behind_an_active_turn() { + let (scheduler, session_manager, _, root) = test_scheduler(); + let session_id = "delegated-busy-session"; + let workspace = root.path().join("workspace"); + std::fs::create_dir_all(&workspace).expect("workspace"); + session_manager + .create_session_with_id( + Some(session_id.to_string()), + "Delegated".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("create delegated session"); + session_manager + .update_session_state( + session_id, + SessionState::Processing { + current_turn_id: "active-turn".to_string(), + phase: ProcessingPhase::Thinking, + }, + ) + .await + .expect("mark active turn"); + + let error = scheduler + .submit_dialog_turn(AgentDialogTurnRequest { + session_id: session_id.to_string(), + message: "expanded command prompt".to_string(), + original_message: Some("/review".to_string()), + turn_id: Some("delegated-turn".to_string()), + execution: bitfun_runtime_ports::AgentDialogTurnExecution::FreshExternalSubagent { + ecosystem_id: "opencode".to_string(), + logical_id: "reviewer".to_string(), + }, + agent_type: "agentic".to_string(), + workspace_path: None, + remote_connection_id: None, + remote_ssh_host: None, + policy: DialogSubmissionPolicy::for_source(DialogTriggerSource::Cli), + reply_route: None, + prepended_reminders: Vec::new(), + attachments: Vec::new(), + metadata: serde_json::Map::new(), + }) + .await + .expect_err("delegated commands must not queue behind another turn"); + + assert_eq!(error.kind, PortErrorKind::InvalidRequest); + assert!(error.message.contains("idle session"), "{error}"); + assert_eq!(scheduler.queue_depth(session_id), 0); + } + + #[tokio::test] + async fn delegated_dialog_turn_does_not_clear_a_queue_from_an_error_session() { + let (scheduler, session_manager, _, root) = test_scheduler(); + let session_id = "delegated-error-session"; + let workspace = root.path().join("workspace"); + std::fs::create_dir_all(&workspace).expect("workspace"); + session_manager + .create_session_with_id( + Some(session_id.to_string()), + "Delegated".to_string(), + "agentic".to_string(), + SessionConfig { + workspace_path: Some(workspace.to_string_lossy().to_string()), + ..Default::default() + }, + ) + .await + .expect("create delegated session"); + session_manager + .update_session_state( + session_id, + SessionState::Processing { + current_turn_id: "active-turn".to_string(), + phase: ProcessingPhase::Thinking, + }, + ) + .await + .expect("mark active turn"); + + scheduler + .submit_dialog_turn(AgentDialogTurnRequest { + session_id: session_id.to_string(), + message: "queued prompt".to_string(), + original_message: None, + turn_id: Some("queued-turn".to_string()), + execution: Default::default(), + agent_type: "agentic".to_string(), + workspace_path: None, + remote_connection_id: None, + remote_ssh_host: None, + policy: DialogSubmissionPolicy::for_source(DialogTriggerSource::Cli), + reply_route: None, + prepended_reminders: Vec::new(), + attachments: Vec::new(), + metadata: serde_json::Map::new(), + }) + .await + .expect("queue standard turn"); + session_manager + .update_session_state( + session_id, + SessionState::Error { + error: "previous turn failed".to_string(), + recoverable: true, + }, + ) + .await + .expect("mark session recoverable error"); + + let error = scheduler + .submit_dialog_turn(AgentDialogTurnRequest { + session_id: session_id.to_string(), + message: "expanded command prompt".to_string(), + original_message: Some("/review".to_string()), + turn_id: Some("delegated-turn".to_string()), + execution: bitfun_runtime_ports::AgentDialogTurnExecution::FreshExternalSubagent { + ecosystem_id: "opencode".to_string(), + logical_id: "reviewer".to_string(), + }, + agent_type: "agentic".to_string(), + workspace_path: None, + remote_connection_id: None, + remote_ssh_host: None, + policy: DialogSubmissionPolicy::for_source(DialogTriggerSource::Cli), + reply_route: None, + prepended_reminders: Vec::new(), + attachments: Vec::new(), + metadata: serde_json::Map::new(), + }) + .await + .expect_err("delegated commands must not replace a queued turn after an error"); + + assert_eq!(error.kind, PortErrorKind::InvalidRequest); + assert!(error.message.contains("idle"), "{error}"); + assert_eq!(scheduler.queue_depth(session_id), 1); + assert!(scheduler + .cancel_queued_or_active_turn(session_id, "queued-turn") + .await + .expect("cancel preserved queued turn")); + } + #[tokio::test] async fn reject_busy_dialog_port_does_not_enqueue_or_replace_the_active_turn() { let (scheduler, session_manager, _, root) = test_scheduler(); @@ -2918,6 +3165,7 @@ mod tests { message: "second prompt".to_string(), original_message: None, turn_id: Some("rejected-turn".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: None, remote_connection_id: None, @@ -2989,6 +3237,7 @@ mod tests { message: "duplicate".to_string(), original_message: None, turn_id: Some(turn_id.to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: None, remote_connection_id: None, @@ -3032,6 +3281,7 @@ mod tests { message: "wrong workspace".to_string(), original_message: None, turn_id: Some(turn_id.to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some(workspace_b.to_string_lossy().to_string()), remote_connection_id: None, @@ -3081,6 +3331,7 @@ mod tests { message: "invalid agent".to_string(), original_message: None, turn_id: Some(turn_id.to_string()), + execution: Default::default(), agent_type: "agent-that-does-not-exist".to_string(), workspace_path: None, remote_connection_id: None, diff --git a/src/crates/assembly/core/src/agentic/deep_review/task_adapter.rs b/src/crates/assembly/core/src/agentic/deep_review/task_adapter.rs index 934764a333..6ecb974896 100644 --- a/src/crates/assembly/core/src/agentic/deep_review/task_adapter.rs +++ b/src/crates/assembly/core/src/agentic/deep_review/task_adapter.rs @@ -101,31 +101,6 @@ pub(crate) fn ensure_deep_review_auto_retry_allowed( runtime_task_execution::ensure_deep_review_auto_retry_allowed(conc_policy, elapsed_seconds) } -#[allow(clippy::too_many_arguments)] -pub(crate) fn deep_review_task_completion_result( - delegate_target_label: &str, - result_text: &str, - context_mode: &str, - duration_ms: u128, - is_partial_timeout: bool, - reason: Option<&str>, - ledger_event_id: Option<&str>, - retry_hint: &str, -) -> (Value, String) { - runtime_task_execution::deep_review_task_completion_result( - runtime_task_execution::DeepReviewTaskCompletionResultInput { - delegate_target_label, - result_text, - context_mode, - duration_ms, - is_partial_timeout, - reason, - ledger_event_id, - retry_hint, - }, - ) -} - pub(crate) fn deep_review_cancelled_reviewer_result( subagent_type: &str, reason: &str, diff --git a/src/crates/assembly/core/src/agentic/permission_policy.rs b/src/crates/assembly/core/src/agentic/permission_policy.rs index c05cef6c92..d402099557 100644 --- a/src/crates/assembly/core/src/agentic/permission_policy.rs +++ b/src/crates/assembly/core/src/agentic/permission_policy.rs @@ -1,4 +1,6 @@ +use crate::service::config::global::GlobalConfigManager; use crate::service::config::types::{AgentProfileConfig, GlobalConfig}; +use crate::util::errors::BitFunResult; use bitfun_runtime_ports::{ resolve_child_permission_policy, resolve_permission_policy, ChildPermissionPolicyLayers, PermissionConstraintLayer, PermissionEffect, PermissionPolicyLayers, PermissionRule, @@ -22,6 +24,18 @@ pub(crate) fn derive_parent_permission_runtime_ceiling( .expect("parent permission ceiling extraction must exclude allow rules") } +pub(crate) async fn load_parent_permission_runtime_ceiling( + agent_type: Option<&str>, +) -> BitFunResult { + let service = GlobalConfigManager::get_service().await?; + let global: GlobalConfig = service.get_config(None).await?; + let profile = agent_type.and_then(|agent_type| { + let profile_id = crate::agentic::agents::resolve_mode_config_profile_id(agent_type); + global.ai.agent_profiles.get(profile_id.as_ref()) + }); + Ok(derive_parent_permission_runtime_ceiling(profile)) +} + pub(crate) fn resolve_effective_permission_policy( global: &GlobalConfig, project_rules: &[PermissionRule], diff --git a/src/crates/assembly/core/src/agentic/session/session_manager.rs b/src/crates/assembly/core/src/agentic/session/session_manager.rs index 35d8b790cf..c9a2f88a36 100644 --- a/src/crates/assembly/core/src/agentic/session/session_manager.rs +++ b/src/crates/assembly/core/src/agentic/session/session_manager.rs @@ -6329,7 +6329,7 @@ impl SessionManager { /// host surface (e.g. CLI) does not persist rounds itself. This ensures /// turn files contain rich conversation data (text, tools, thinking) that /// other surfaces (e.g. Desktop) can render. - fn build_model_rounds_from_messages( + pub(crate) fn build_model_rounds_from_messages( messages: &[Message], turn_id: &str, timestamp: u64, @@ -6792,10 +6792,45 @@ impl SessionManager { turn_id: &str, model_rounds: Vec, duration_ms: u64, + ) -> BitFunResult<()> { + self.complete_turn_with_model_rounds( + session_id, + turn_id, + model_rounds, + duration_ms, + "maintenance_turn_completed", + ) + .await + } + + pub(crate) async fn complete_synthetic_dialog_turn( + &self, + session_id: &str, + turn_id: &str, + model_rounds: Vec, + duration_ms: u64, + ) -> BitFunResult<()> { + self.complete_turn_with_model_rounds( + session_id, + turn_id, + model_rounds, + duration_ms, + "synthetic_dialog_turn_completed", + ) + .await + } + + async fn complete_turn_with_model_rounds( + &self, + session_id: &str, + turn_id: &str, + model_rounds: Vec, + duration_ms: u64, + snapshot_reason: &str, ) -> BitFunResult<()> { if !self.should_persist_session_id(session_id) { debug!( - "Skipping maintenance turn persistence for transient session completion: session_id={}, turn_id={}, rounds={}, duration_ms={}", + "Skipping turn persistence for transient session completion: session_id={}, turn_id={}, rounds={}, duration_ms={}", session_id, turn_id, model_rounds.len(), @@ -6836,7 +6871,7 @@ impl SessionManager { self.persist_context_snapshot_for_turn_best_effort( session_id, turn.turn_index, - "maintenance_turn_completed", + snapshot_reason, ) .await; @@ -6856,10 +6891,45 @@ impl SessionManager { turn_id: &str, error: String, model_rounds: Vec, + ) -> BitFunResult<()> { + self.fail_turn_with_model_rounds( + session_id, + turn_id, + error, + model_rounds, + "maintenance_turn_failed", + ) + .await + } + + pub(crate) async fn fail_synthetic_dialog_turn( + &self, + session_id: &str, + turn_id: &str, + error: String, + model_rounds: Vec, + ) -> BitFunResult<()> { + self.fail_turn_with_model_rounds( + session_id, + turn_id, + error, + model_rounds, + "synthetic_dialog_turn_failed", + ) + .await + } + + async fn fail_turn_with_model_rounds( + &self, + session_id: &str, + turn_id: &str, + error: String, + model_rounds: Vec, + snapshot_reason: &str, ) -> BitFunResult<()> { if !self.should_persist_session_id(session_id) { debug!( - "Skipping maintenance turn persistence for transient session failure: session_id={}, turn_id={}, rounds={}, error={}", + "Skipping turn persistence for transient session failure: session_id={}, turn_id={}, rounds={}, error={}", session_id, turn_id, model_rounds.len(), @@ -6901,7 +6971,7 @@ impl SessionManager { self.persist_context_snapshot_for_turn_best_effort( session_id, turn.turn_index, - "maintenance_turn_failed", + snapshot_reason, ) .await; @@ -6912,7 +6982,7 @@ impl SessionManager { } debug!( - "Maintenance turn marked as failed: turn_id={}, turn_index={}, error={}", + "Turn marked as failed: turn_id={}, turn_index={}, error={}", turn_id, turn.turn_index, error ); diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/session_message_tool.rs b/src/crates/assembly/core/src/agentic/tools/implementations/session_message_tool.rs index a4d4cf636a..c92de53576 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/session_message_tool.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/session_message_tool.rs @@ -689,6 +689,7 @@ Allowed agent types when creating a session: message: forwarded_message, original_message: Some(params.message.clone()), turn_id: None, + execution: Default::default(), agent_type: target_agent_type.clone(), workspace_path: Some(workspace_target.workspace_path.clone()), remote_connection_id: workspace_target.remote_connection_id.clone(), diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/task/execution.rs b/src/crates/assembly/core/src/agentic/tools/implementations/task/execution.rs index e6727d61eb..58b8f1c12f 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/task/execution.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/task/execution.rs @@ -96,17 +96,11 @@ struct BackgroundTaskStartRequest<'a> { impl TaskTool { async fn derive_parent_permission_runtime_ceiling( context: &ToolUseContext, - ) -> PermissionRuntimeCeiling { - let global: GlobalConfig = match GlobalConfigManager::get_service().await { - Ok(service) => service.get_config(None).await.unwrap_or_default(), - Err(_) => GlobalConfig::default(), - }; - let agent_profile = context.agent_type.as_deref().and_then(|agent_type| { - let profile_id = crate::agentic::agents::resolve_mode_config_profile_id(agent_type); - global.ai.agent_profiles.get(profile_id.as_ref()) - }); - - crate::agentic::permission_policy::derive_parent_permission_runtime_ceiling(agent_profile) + ) -> BitFunResult { + crate::agentic::permission_policy::load_parent_permission_runtime_ceiling( + context.agent_type.as_deref(), + ) + .await } pub(super) async fn load_configured_tool_execution_timeout() -> Option { @@ -676,7 +670,7 @@ impl TaskTool { forward_subagent_invocation_context(context, &mut subagent_context); let subagent_context = (!subagent_context.is_empty()).then_some(subagent_context); let permission_runtime_ceiling = - Self::derive_parent_permission_runtime_ceiling(context).await; + Self::derive_parent_permission_runtime_ceiling(context).await?; let prepared_prompt = prompt; if run_in_background { return Self::start_background_task(BackgroundTaskStartRequest { @@ -1156,15 +1150,17 @@ impl TaskTool { }; let (mut data, mut result_for_assistant) = - deep_review_task_adapter::deep_review_task_completion_result( - &delegate_target_label, - &result.text, - context_mode.as_str(), - duration, - result.is_partial_timeout(), - result.reason.as_deref(), - result.ledger_event_id(), - &retry_hint, + bitfun_agent_runtime::subagent_task::subagent_task_completion_result( + bitfun_agent_runtime::subagent_task::SubagentTaskCompletionResultInput { + delegate_target_label: &delegate_target_label, + result_text: &result.text, + context_mode: context_mode.as_str(), + duration_ms: duration, + is_partial_timeout: result.is_partial_timeout(), + reason: result.reason.as_deref(), + ledger_event_id: result.ledger_event_id(), + partial_timeout_suffix: &retry_hint, + }, ); if supports_follow_up { if let Some(subagent_session_id) = result.session_id() { diff --git a/src/crates/assembly/core/src/agentic/tools/implementations/task/mod.rs b/src/crates/assembly/core/src/agentic/tools/implementations/task/mod.rs index a01cf241f2..998a0b6f06 100644 --- a/src/crates/assembly/core/src/agentic/tools/implementations/task/mod.rs +++ b/src/crates/assembly/core/src/agentic/tools/implementations/task/mod.rs @@ -24,7 +24,7 @@ use crate::agentic::tools::framework::{ }; use crate::agentic::tools::pipeline::SubagentParentInfo; use crate::service::config::global::GlobalConfigManager; -use crate::service::config::types::{AIConfig, GlobalConfig}; +use crate::service::config::types::AIConfig; use crate::util::errors::{BitFunError, BitFunResult}; use crate::util::timing::elapsed_ms_u64; use async_trait::async_trait; diff --git a/src/crates/assembly/core/src/external_sources.rs b/src/crates/assembly/core/src/external_sources.rs index 3f5cbea4a3..4dc14279dc 100644 --- a/src/crates/assembly/core/src/external_sources.rs +++ b/src/crates/assembly/core/src/external_sources.rs @@ -28,9 +28,9 @@ pub use bitfun_product_domains::external_sources::{ ExternalToolRuntimeKind, NativePromptCommandConflictProjection, NativePromptCommandConflictSnapshot, NativePromptCommandDescriptor, NativePromptCommandReconfirmationProjection, PromptCommandAvailability, - PromptCommandCatalogEntry, PromptCommandDefinition, PromptCommandInvocationOutcome, - PromptCommandShellReviewDecision, PromptCommandShellReviewMode, PromptCommandShellReviewPlan, - SourceKey, + PromptCommandCatalogEntry, PromptCommandDefinition, PromptCommandExecutionTarget, + PromptCommandInvocationOutcome, PromptCommandShellReviewDecision, PromptCommandShellReviewMode, + PromptCommandShellReviewPlan, SourceKey, }; pub use bitfun_product_domains::external_subagents::{ ExternalSubagentActivationState, ExternalSubagentCompatibilityState, ExternalSubagentConflict, @@ -83,8 +83,9 @@ use bitfun_product_domains::external_integration_policy::{ }; use bitfun_product_domains::external_sources::{ ExecutionDomainId, ExternalMcpRevisionKey, ExternalMcpSourceProvider, ExternalMcpStaticStatus, - ExternalSourceContext, ExternalSourceScope, ExternalToolSourceProvider, PromptCommandExpansion, - PromptCommandShellInvocation, PromptCommandShellPreference, PromptCommandSourceProvider, + ExternalSourceContext, ExternalSourceScope, ExternalToolSourceProvider, PromptCommandConflict, + PromptCommandExpansion, PromptCommandShellInvocation, PromptCommandShellPreference, + PromptCommandSourceProvider, }; use bitfun_product_domains::external_subagents::ExternalSubagentSourceProvider; use bitfun_product_domains::workspace_references::{ @@ -1319,6 +1320,67 @@ fn source_ecosystem_id( }) } +fn restrict_prompt_commands_without_active_subagents( + commands: &mut [PromptCommandCatalogEntry], + conflicts: &mut [PromptCommandConflict], + active_subagents: &BTreeSet<(EcosystemId, String)>, +) { + for command in commands { + if !matches!( + command.definition.availability, + PromptCommandAvailability::Available + ) { + continue; + } + let PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id, + logical_id, + } = &command.definition.execution_target + else { + continue; + }; + let key = (ecosystem_id.clone(), logical_id.to_ascii_lowercase()); + if !active_subagents.contains(&key) { + command.definition.availability = PromptCommandAvailability::Restricted { + reason: format!( + "External command subagent '{}' is not currently approved and available", + logical_id + ), + required_capabilities: vec!["command.external_subagent".to_string()], + }; + } + } + for conflict in conflicts { + for candidate in &mut conflict.candidates { + if !matches!(candidate.availability, PromptCommandAvailability::Available) { + continue; + } + let PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id, + logical_id, + } = &candidate.execution_target + else { + continue; + }; + let key = (ecosystem_id.clone(), logical_id.to_ascii_lowercase()); + if !active_subagents.contains(&key) { + candidate.availability = PromptCommandAvailability::Restricted { + reason: format!( + "External command subagent '{}' is not currently approved and available", + logical_id + ), + required_capabilities: vec!["command.external_subagent".to_string()], + }; + if conflict.selected_candidate_id.as_deref() + == Some(candidate.candidate_id.as_str()) + { + conflict.selected_candidate_id = None; + } + } + } + } +} + fn ensure_source_capability_active( snapshot: &ExternalSourceCatalogSnapshot, source_key: &SourceKey, @@ -1954,6 +2016,21 @@ impl WorkspaceExternalSourceService { &subagent_state, preferences.preference_revision, ); + let active_subagents = subagent_state + .registrations + .iter() + .map(|registration| { + ( + registration.ecosystem_id.clone(), + registration.logical_id.to_ascii_lowercase(), + ) + }) + .collect::>(); + restrict_prompt_commands_without_active_subagents( + &mut snapshot.commands, + &mut snapshot.command_conflicts, + &active_subagents, + ); if let Some(workspace_root) = self.workspace_root.as_deref() { crate::agentic::agents::get_agent_registry().install_external_subagent_routes( workspace_root, @@ -3545,6 +3622,7 @@ impl WorkspaceExternalSourceService { missing_candidate_error(format!("External prompt command '{name}' was not found")) })?; let source_key = selected_command.id.source.clone(); + let execution_target = selected_command.execution_target.clone(); ensure_source_capability_active(&snapshot, &source_key, EXTERNAL_CAPABILITY_COMMAND)?; let source_display_name = snapshot .sources @@ -3575,6 +3653,7 @@ impl WorkspaceExternalSourceService { .await?; return Ok(PromptCommandInvocationOutcome::Ready { content: expanded.content, + execution_target, }); }; if self.safe_mode_enabled() { @@ -3636,6 +3715,7 @@ impl WorkspaceExternalSourceService { finalize_prompt_command_expansion(self.workspace_root.as_deref(), expansion).await?; Ok(PromptCommandInvocationOutcome::Ready { content: expanded.content, + execution_target, }) } @@ -6824,9 +6904,10 @@ mod tests { use crate::service::mcp::{ConfigLocation, MCPServerConfig, MCPServerType}; use bitfun_product_domains::external_sources::{ EcosystemId, ExternalSourceProviderError, ExternalSourceRecord, ExternalSourceScope, - PromptCommandAvailability, PromptCommandDefinition, PromptCommandProviderIdentity, - PromptCommandProviderSnapshot, PromptCommandShellExpansion, PromptCommandShellInvocation, - PromptCommandShellPreference, SourceQualifiedCommandId, + PromptCommandAvailability, PromptCommandCatalogEntry, PromptCommandConflict, + PromptCommandConflictCandidate, PromptCommandDefinition, PromptCommandExecutionTarget, + PromptCommandProviderIdentity, PromptCommandProviderSnapshot, PromptCommandShellExpansion, + PromptCommandShellInvocation, PromptCommandShellPreference, SourceQualifiedCommandId, }; use bitfun_product_domains::workspace_references::ExternalWorkspaceReferenceDefinition; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -7436,6 +7517,7 @@ mod tests { description: "Review changes".to_string(), template: "Review changes".to_string(), shell_preference: None, + execution_target: Default::default(), availability: PromptCommandAvailability::Available, content_version: "external-v1".to_string(), }; @@ -7572,6 +7654,74 @@ mod tests { assert_eq!(changed.conflicts[0].selected_candidate_id, None); } + #[test] + fn delegated_prompt_commands_require_an_active_same_ecosystem_subagent() { + let source = SourceKey::new("opencode.commands", "project").unwrap(); + let definition = |logical_id: &str| PromptCommandDefinition { + id: SourceQualifiedCommandId::new(source.clone(), logical_id).unwrap(), + name: logical_id.to_string(), + description: logical_id.to_string(), + template: "Review changes".to_string(), + shell_preference: None, + execution_target: PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id: EcosystemId::new("opencode").unwrap(), + logical_id: logical_id.to_string(), + }, + availability: PromptCommandAvailability::Available, + content_version: "command-v1".to_string(), + }; + let mut commands = vec![ + PromptCommandCatalogEntry { + definition: definition("reviewer"), + }, + PromptCommandCatalogEntry { + definition: definition("missing"), + }, + ]; + let active = BTreeSet::from([( + EcosystemId::new("opencode").unwrap(), + "reviewer".to_string(), + )]); + let missing_candidate_id = commands[1].definition.id.stable_key(); + let mut conflicts = vec![PromptCommandConflict { + conflict_key: "prompt-command-conflict".to_string(), + command_name: "missing".to_string(), + candidates: vec![PromptCommandConflictCandidate { + candidate_id: missing_candidate_id.clone(), + source: source.clone(), + source_display_name: "OpenCode".to_string(), + ecosystem_id: EcosystemId::new("opencode").unwrap(), + content_version: "command-v1".to_string(), + command_description: "missing".to_string(), + source_scope: ExternalSourceScope::Project, + source_location: ".opencode/commands/missing.md".to_string(), + execution_target: commands[1].definition.execution_target.clone(), + availability: PromptCommandAvailability::Available, + }], + selected_candidate_id: Some(missing_candidate_id), + }]; + + restrict_prompt_commands_without_active_subagents(&mut commands, &mut conflicts, &active); + + assert_eq!( + commands[0].definition.availability, + PromptCommandAvailability::Available + ); + let PromptCommandAvailability::Restricted { + required_capabilities, + .. + } = &commands[1].definition.availability + else { + panic!("missing delegated subagent must restrict the command"); + }; + assert_eq!(required_capabilities, &["command.external_subagent"]); + assert!(matches!( + conflicts[0].candidates[0].availability, + PromptCommandAvailability::Restricted { .. } + )); + assert_eq!(conflicts[0].selected_candidate_id, None); + } + #[test] fn native_prompt_command_choice_does_not_mark_external_candidate_as_self_conflicted() { let mut config = ExternalSourcesConfig::default(); @@ -7977,6 +8127,7 @@ mod tests { description: self.command_name.clone(), template: self.command_name.clone(), shell_preference: None, + execution_target: Default::default(), availability: PromptCommandAvailability::Available, content_version: "command-v1".to_string(), }], diff --git a/src/crates/assembly/core/src/external_subagents.rs b/src/crates/assembly/core/src/external_subagents.rs index 33d062b931..448ffdfa9b 100644 --- a/src/crates/assembly/core/src/external_subagents.rs +++ b/src/crates/assembly/core/src/external_subagents.rs @@ -84,6 +84,7 @@ struct ProductFacts { struct ResolvedExternalCandidate { definition: ExternalSubagentDefinition, + ecosystem_id: Option, provider_label: String, scope: ExternalSourceScope, source_keys: Vec, @@ -960,6 +961,13 @@ fn resolve_external_candidate( } } } + if ecosystem_id.is_none() { + compatibility = ExternalSubagentCompatibilityState::Invalid; + diagnostics.push(ExternalSubagentDiagnosticSummary { + code: "external_subagent.source_unavailable".to_string(), + blocks_activation: true, + }); + } let model_resolution = resolve_model_request( definition, ecosystem_id.as_ref(), @@ -1084,6 +1092,7 @@ fn resolve_external_candidate( .unwrap_or_else(|| "External AI app".to_string()); ResolvedExternalCandidate { definition: definition.clone(), + ecosystem_id, provider_label, scope, source_keys, @@ -1251,6 +1260,9 @@ fn install_active_candidate( candidate: &ResolvedExternalCandidate, state: &mut ExternalSubagentProductState, ) { + let Some(ecosystem_id) = candidate.ecosystem_id.clone() else { + return; + }; let runtime_key = external_subagent_runtime_key(&stable_digest([ candidate.definition.candidate_id.as_str(), candidate.definition.behavior_version.as_str(), @@ -1291,6 +1303,7 @@ fn install_active_candidate( state.registrations.push(ExternalSubagentRegistration { runtime_key: runtime_key.clone(), logical_id: candidate.definition.logical_id.clone(), + ecosystem_id, provider_label: candidate.provider_label.clone(), model_binding, hidden: candidate.definition.hidden, diff --git a/src/crates/assembly/core/src/service/cron/service.rs b/src/crates/assembly/core/src/service/cron/service.rs index 221d3aeb15..e32b370817 100644 --- a/src/crates/assembly/core/src/service/cron/service.rs +++ b/src/crates/assembly/core/src/service/cron/service.rs @@ -566,6 +566,7 @@ impl CronService { message: enqueue_input.user_input.clone(), original_message: Some(enqueue_input.user_input.clone()), turn_id: Some(enqueue_input.turn_id.clone()), + execution: Default::default(), agent_type: resolved.agent_type, workspace_path: Some(resolved.workspace_path), remote_connection_id: resolved.remote_connection_id, diff --git a/src/crates/assembly/core/src/service_agent_runtime.rs b/src/crates/assembly/core/src/service_agent_runtime.rs index 91751b47ee..a5b11bf1d3 100644 --- a/src/crates/assembly/core/src/service_agent_runtime.rs +++ b/src/crates/assembly/core/src/service_agent_runtime.rs @@ -1692,6 +1692,7 @@ impl RemoteDialogRuntimeHost for CoreRemoteDialogRuntimeHost<'_> { message: submission.content, original_message: None, turn_id: Some(submission.turn_id), + execution: Default::default(), agent_type: submission.resolved_agent_type, workspace_path, remote_connection_id, diff --git a/src/crates/assembly/external-sources/src/lib.rs b/src/crates/assembly/external-sources/src/lib.rs index 7f1062ca94..e08b5bb180 100644 --- a/src/crates/assembly/external-sources/src/lib.rs +++ b/src/crates/assembly/external-sources/src/lib.rs @@ -684,6 +684,7 @@ impl ExternalSourceCoordinator { command_description: command.description.clone(), source_scope: source.record.scope, source_location: source.record.location.clone(), + execution_target: command.execution_target.clone(), availability: command.availability.clone(), }) }) diff --git a/src/crates/assembly/external-sources/tests/coordinator_contracts.rs b/src/crates/assembly/external-sources/tests/coordinator_contracts.rs index c2e766d1da..1b6be25c0c 100644 --- a/src/crates/assembly/external-sources/tests/coordinator_contracts.rs +++ b/src/crates/assembly/external-sources/tests/coordinator_contracts.rs @@ -41,6 +41,7 @@ fn command_named( description: format!("Review from {provider_id}"), template: format!("{provider_id}: $ARGUMENTS"), shell_preference: None, + execution_target: Default::default(), availability: PromptCommandAvailability::Available, content_version: format!("command-v{version}"), } diff --git a/src/crates/contracts/product-domains/src/external_sources.rs b/src/crates/contracts/product-domains/src/external_sources.rs index 61d611530d..620963dd62 100644 --- a/src/crates/contracts/product-domains/src/external_sources.rs +++ b/src/crates/contracts/product-domains/src/external_sources.rs @@ -1509,6 +1509,46 @@ impl PromptCommandShellPreference { } } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde( + tag = "kind", + rename_all = "snake_case", + rename_all_fields = "camelCase", + deny_unknown_fields +)] +pub enum PromptCommandExecutionTarget { + Inline, + FreshExternalSubagent { + ecosystem_id: EcosystemId, + logical_id: String, + }, +} + +impl Default for PromptCommandExecutionTarget { + fn default() -> Self { + Self::Inline + } +} + +impl PromptCommandExecutionTarget { + pub fn is_inline(&self) -> bool { + matches!(self, Self::Inline) + } + + fn validate(&self) -> Result<(), ExternalSourceContractError> { + match self { + Self::Inline => Ok(()), + Self::FreshExternalSubagent { + ecosystem_id, + logical_id, + } => { + validate_id(ecosystem_id.as_str(), "prompt command subagent ecosystem")?; + validate_id(logical_id, "prompt command subagent logical id") + } + } + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct PromptCommandDefinition { @@ -1518,6 +1558,11 @@ pub struct PromptCommandDefinition { pub template: String, #[serde(default, skip_serializing_if = "Option::is_none")] pub shell_preference: Option, + #[serde( + default, + skip_serializing_if = "PromptCommandExecutionTarget::is_inline" + )] + pub execution_target: PromptCommandExecutionTarget, pub availability: PromptCommandAvailability, /// Version of this command only. Unrelated edits in the same source must /// not invalidate a remembered conflict choice. @@ -1536,6 +1581,7 @@ impl PromptCommandDefinition { if let Some(shell) = &self.shell_preference { shell.validate()?; } + self.execution_target.validate()?; validate_id(&self.content_version, "command content version") } } @@ -1593,10 +1639,15 @@ impl fmt::Debug for PromptCommandShellReviewPlan { } #[derive(Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(tag = "state", rename_all = "snake_case")] +#[serde( + tag = "state", + rename_all = "snake_case", + rename_all_fields = "camelCase" +)] pub enum PromptCommandInvocationOutcome { Ready { content: String, + execution_target: PromptCommandExecutionTarget, }, ReviewRequired { review: PromptCommandShellReviewPlan, @@ -1606,9 +1657,13 @@ pub enum PromptCommandInvocationOutcome { impl fmt::Debug for PromptCommandInvocationOutcome { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { match self { - Self::Ready { content } => formatter + Self::Ready { + content, + execution_target, + } => formatter .debug_struct("Ready") .field("content_bytes", &content.len()) + .field("execution_target", execution_target) .finish(), Self::ReviewRequired { review } => formatter .debug_struct("ReviewRequired") @@ -2075,6 +2130,11 @@ pub struct PromptCommandConflictCandidate { pub command_description: String, pub source_scope: ExternalSourceScope, pub source_location: String, + #[serde( + default, + skip_serializing_if = "PromptCommandExecutionTarget::is_inline" + )] + pub execution_target: PromptCommandExecutionTarget, pub availability: PromptCommandAvailability, } @@ -2479,10 +2539,43 @@ impl From for ExternalSourcePublicSnapshot { #[cfg(test)] mod prompt_command_shell_review_tests { use super::{ - PromptCommandInvocationOutcome, PromptCommandShellReviewDecision, - PromptCommandShellReviewMode, PromptCommandShellReviewPlan, + EcosystemId, PromptCommandExecutionTarget, PromptCommandInvocationOutcome, + PromptCommandShellReviewDecision, PromptCommandShellReviewMode, + PromptCommandShellReviewPlan, }; + #[test] + fn ready_outcome_preserves_the_typed_execution_target() { + let outcome = PromptCommandInvocationOutcome::Ready { + content: "Review this change".to_string(), + execution_target: PromptCommandExecutionTarget::FreshExternalSubagent { + ecosystem_id: EcosystemId::new("opencode").unwrap(), + logical_id: "reviewer".to_string(), + }, + }; + + assert_eq!( + serde_json::to_value(outcome).unwrap(), + serde_json::json!({ + "state": "ready", + "content": "Review this change", + "executionTarget": { + "kind": "fresh_external_subagent", + "ecosystemId": "opencode", + "logicalId": "reviewer" + } + }) + ); + + let invalid = serde_json::from_value::(serde_json::json!({ + "kind": "fresh_external_subagent", + "ecosystemId": "", + "logicalId": "reviewer" + })) + .expect("open ids validate at their owning contract boundary"); + assert!(invalid.validate().is_err()); + } + #[test] fn review_outcome_uses_a_stable_tagged_transport_shape() { let outcome = PromptCommandInvocationOutcome::ReviewRequired { diff --git a/src/crates/contracts/product-domains/tests/external_source_contracts.rs b/src/crates/contracts/product-domains/tests/external_source_contracts.rs index 11e5c1313e..2fbcd93fcb 100644 --- a/src/crates/contracts/product-domains/tests/external_source_contracts.rs +++ b/src/crates/contracts/product-domains/tests/external_source_contracts.rs @@ -149,6 +149,7 @@ fn command(provider_id: &str, source_id: &str, precedence: i32) -> PromptCommand description: format!("Review from {provider_id}"), template: format!("{provider_id}: $ARGUMENTS"), shell_preference: None, + execution_target: Default::default(), availability: PromptCommandAvailability::Available, content_version: format!("command-v{precedence}"), } @@ -236,6 +237,7 @@ fn prompt_commands_use_a_typed_contract_instead_of_an_arbitrary_asset_payload() description: "Review the current change".to_string(), template: "Review $ARGUMENTS".to_string(), shell_preference: None, + execution_target: Default::default(), availability: PromptCommandAvailability::Restricted { reason: "Shell expansion is not supported yet".to_string(), required_capabilities: vec!["command.shell".to_string()], diff --git a/src/crates/contracts/runtime-ports/src/lib.rs b/src/crates/contracts/runtime-ports/src/lib.rs index 33d6e84308..bd5ad2cdd4 100644 --- a/src/crates/contracts/runtime-ports/src/lib.rs +++ b/src/crates/contracts/runtime-ports/src/lib.rs @@ -1494,6 +1494,27 @@ pub struct AgentSubmissionRequest { pub metadata: serde_json::Map, } +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde( + tag = "kind", + rename_all = "snake_case", + rename_all_fields = "camelCase", + deny_unknown_fields +)] +pub enum AgentDialogTurnExecution { + Standard, + FreshExternalSubagent { + ecosystem_id: String, + logical_id: String, + }, +} + +impl Default for AgentDialogTurnExecution { + fn default() -> Self { + Self::Standard + } +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct AgentDialogTurnRequest { @@ -1503,6 +1524,8 @@ pub struct AgentDialogTurnRequest { pub original_message: Option, #[serde(skip_serializing_if = "Option::is_none")] pub turn_id: Option, + #[serde(default, skip_serializing_if = "AgentDialogTurnExecution::is_standard")] + pub execution: AgentDialogTurnExecution, pub agent_type: String, #[serde(skip_serializing_if = "Option::is_none")] pub workspace_path: Option, @@ -1521,6 +1544,12 @@ pub struct AgentDialogTurnRequest { pub metadata: serde_json::Map, } +impl AgentDialogTurnExecution { + pub fn is_standard(&self) -> bool { + matches!(self, Self::Standard) + } +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct AgentDialogPrependedReminder { @@ -3055,6 +3084,36 @@ mod tests { assert_eq!(sdk_host, serde_json::json!("sdk_host")); } + #[test] + fn delegated_dialog_turn_target_is_typed_and_provider_neutral() { + let target = AgentDialogTurnExecution::FreshExternalSubagent { + ecosystem_id: "opencode".to_string(), + logical_id: "reviewer".to_string(), + }; + + assert_eq!( + serde_json::to_value(target).expect("serialize delegated execution"), + serde_json::json!({ + "kind": "fresh_external_subagent", + "ecosystemId": "opencode", + "logicalId": "reviewer", + }) + ); + assert_eq!( + AgentDialogTurnExecution::default(), + AgentDialogTurnExecution::Standard + ); + assert!( + serde_json::from_value::(serde_json::json!({ + "kind": "fresh_external_subagent", + "ecosystemId": "opencode", + "logicalId": "reviewer", + "model": "provider/model" + })) + .is_err() + ); + } + #[test] fn dialog_submission_policy_preserves_current_surface_queue_defaults() { let remote = DialogSubmissionPolicy::for_source(DialogTriggerSource::RemoteRelay); @@ -3483,6 +3542,7 @@ mod tests { message: "hello".to_string(), original_message: Some("raw hello".to_string()), turn_id: Some("turn_1".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some("/workspace/project".to_string()), remote_connection_id: Some("conn-1".to_string()), diff --git a/src/crates/execution/agent-runtime/examples/sdk_minimal.rs b/src/crates/execution/agent-runtime/examples/sdk_minimal.rs index 7e405cd267..d9c338d61b 100644 --- a/src/crates/execution/agent-runtime/examples/sdk_minimal.rs +++ b/src/crates/execution/agent-runtime/examples/sdk_minimal.rs @@ -49,7 +49,7 @@ impl AgentSubmissionPort for ExampleAgentProvider { #[tokio::main] async fn main() -> Result<(), Box> { let compatibility = AgentRuntimeSdkCompatibility::current(); - assert_eq!(compatibility.api_version, 2); + assert_eq!(compatibility.api_version, 3); let provider = Arc::new(ExampleAgentProvider::default()); let events = AgentEventStream::new(); diff --git a/src/crates/execution/agent-runtime/src/deep_review/task_execution.rs b/src/crates/execution/agent-runtime/src/deep_review/task_execution.rs index 535535f71a..8617a9b83b 100644 --- a/src/crates/execution/agent-runtime/src/deep_review/task_execution.rs +++ b/src/crates/execution/agent-runtime/src/deep_review/task_execution.rs @@ -425,39 +425,18 @@ pub struct DeepReviewTaskCompletionResultInput<'a> { pub fn deep_review_task_completion_result( input: DeepReviewTaskCompletionResultInput<'_>, ) -> (Value, String) { - let status = if input.is_partial_timeout { - "partial_timeout" - } else { - "completed" - }; - let assistant_message = if input.is_partial_timeout { - format!( - "{} timed out with partial result:\n\n{}\n{}", - input.delegate_target_label, input.result_text, input.retry_hint - ) - } else { - format!( - "{} completed successfully with result:\n\n{}\n", - input.delegate_target_label, input.result_text - ) - }; - let mut data = json!({ - "duration": input.duration_ms, - "context_mode": input.context_mode, - "status": status - }); - - if input.is_partial_timeout { - data["partial_output"] = json!(input.result_text); - if let Some(reason) = input.reason { - data["reason"] = json!(reason); - } - if let Some(event_id) = input.ledger_event_id { - data["ledger_event_id"] = json!(event_id); - } - } - - (data, assistant_message) + crate::subagent_task::subagent_task_completion_result( + crate::subagent_task::SubagentTaskCompletionResultInput { + delegate_target_label: input.delegate_target_label, + result_text: input.result_text, + context_mode: input.context_mode, + duration_ms: input.duration_ms, + is_partial_timeout: input.is_partial_timeout, + reason: input.reason, + ledger_event_id: input.ledger_event_id, + partial_timeout_suffix: input.retry_hint, + }, + ) } pub fn deep_review_cancelled_reviewer_result( diff --git a/src/crates/execution/agent-runtime/src/lib.rs b/src/crates/execution/agent-runtime/src/lib.rs index ae0a9ed419..7c4b2e6b62 100644 --- a/src/crates/execution/agent-runtime/src/lib.rs +++ b/src/crates/execution/agent-runtime/src/lib.rs @@ -37,6 +37,7 @@ pub mod session_state_manager; pub mod side_question; pub mod skill_agent_snapshot; pub mod skills; +pub mod subagent_task; pub mod thread_goal; pub mod thread_goal_tools; pub mod turn_cancellation; diff --git a/src/crates/execution/agent-runtime/src/runtime.rs b/src/crates/execution/agent-runtime/src/runtime.rs index 782df7bcd4..55634e7d93 100644 --- a/src/crates/execution/agent-runtime/src/runtime.rs +++ b/src/crates/execution/agent-runtime/src/runtime.rs @@ -2781,6 +2781,7 @@ mod tests { message: "hello".to_string(), original_message: None, turn_id: Some("turn_1".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some("/workspace/project".to_string()), remote_connection_id: None, @@ -2835,6 +2836,7 @@ mod tests { message: "hello".to_string(), original_message: Some("hello".to_string()), turn_id: Some("turn_1".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some("/workspace/project".to_string()), remote_connection_id: None, @@ -2905,6 +2907,7 @@ mod tests { message: "hello".to_string(), original_message: None, turn_id: Some("turn-1".to_string()), + execution: Default::default(), agent_type: "agentic".to_string(), workspace_path: Some("/workspace/project".to_string()), remote_connection_id: None, diff --git a/src/crates/execution/agent-runtime/src/sdk.rs b/src/crates/execution/agent-runtime/src/sdk.rs index c7db2ed7da..383157192f 100644 --- a/src/crates/execution/agent-runtime/src/sdk.rs +++ b/src/crates/execution/agent-runtime/src/sdk.rs @@ -8,7 +8,7 @@ use std::sync::Arc; -pub const AGENT_RUNTIME_SDK_API_VERSION: u32 = 2; +pub const AGENT_RUNTIME_SDK_API_VERSION: u32 = 3; #[derive(Debug, Clone, Copy, PartialEq, Eq)] #[non_exhaustive] @@ -58,14 +58,15 @@ pub use bitfun_harness::{ HarnessRegistry, HarnessWorkflow, }; pub use bitfun_runtime_ports::{ - AgentBackgroundResultRequest, AgentDialogTurnPort, AgentDialogTurnRequest, - AgentInputAttachment, AgentLifecycleDeliveryPort, AgentLocalCommandTurnPort, - AgentLocalCommandTurnRecordRequest, AgentMessageWorkspaceReferencesRequest, - AgentSessionArchiveRequest, AgentSessionArchiveStateRequest, AgentSessionClosePort, - AgentSessionCompactionPort, AgentSessionCompactionRequest, AgentSessionCompactionResult, - AgentSessionComposerUpdate, AgentSessionCreateRequest, AgentSessionCreateResult, - AgentSessionDeleteRequest, AgentSessionForkAtTurnRequest, AgentSessionForkBeforeTurnRequest, - AgentSessionForkPort, AgentSessionForkRequest, AgentSessionForkResult, AgentSessionListRequest, + AgentBackgroundResultRequest, AgentDialogTurnExecution, AgentDialogTurnPort, + AgentDialogTurnRequest, AgentInputAttachment, AgentLifecycleDeliveryPort, + AgentLocalCommandTurnPort, AgentLocalCommandTurnRecordRequest, + AgentMessageWorkspaceReferencesRequest, AgentSessionArchiveRequest, + AgentSessionArchiveStateRequest, AgentSessionClosePort, AgentSessionCompactionPort, + AgentSessionCompactionRequest, AgentSessionCompactionResult, AgentSessionComposerUpdate, + AgentSessionCreateRequest, AgentSessionCreateResult, AgentSessionDeleteRequest, + AgentSessionForkAtTurnRequest, AgentSessionForkBeforeTurnRequest, AgentSessionForkPort, + AgentSessionForkRequest, AgentSessionForkResult, AgentSessionListRequest, AgentSessionManagementPort, AgentSessionModePort, AgentSessionModeUpdateRequest, AgentSessionModelPort, AgentSessionModelUpdateRequest, AgentSessionRenameRequest, AgentSessionRevertPort, AgentSessionRevertRequest, AgentSessionRevertResult, diff --git a/src/crates/execution/agent-runtime/src/subagent_task.rs b/src/crates/execution/agent-runtime/src/subagent_task.rs new file mode 100644 index 0000000000..e9380839eb --- /dev/null +++ b/src/crates/execution/agent-runtime/src/subagent_task.rs @@ -0,0 +1,53 @@ +//! Provider-neutral formatting for completed subagent Task calls. + +use serde_json::{json, Value}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct SubagentTaskCompletionResultInput<'a> { + pub delegate_target_label: &'a str, + pub result_text: &'a str, + pub context_mode: &'a str, + pub duration_ms: u128, + pub is_partial_timeout: bool, + pub reason: Option<&'a str>, + pub ledger_event_id: Option<&'a str>, + pub partial_timeout_suffix: &'a str, +} + +pub fn subagent_task_completion_result( + input: SubagentTaskCompletionResultInput<'_>, +) -> (Value, String) { + let status = if input.is_partial_timeout { + "partial_timeout" + } else { + "completed" + }; + let assistant_message = if input.is_partial_timeout { + format!( + "{} timed out with partial result:\n\n{}\n{}", + input.delegate_target_label, input.result_text, input.partial_timeout_suffix + ) + } else { + format!( + "{} completed successfully with result:\n\n{}\n", + input.delegate_target_label, input.result_text + ) + }; + let mut data = json!({ + "duration": input.duration_ms, + "context_mode": input.context_mode, + "status": status + }); + + if input.is_partial_timeout { + data["partial_output"] = json!(input.result_text); + if let Some(reason) = input.reason { + data["reason"] = json!(reason); + } + if let Some(event_id) = input.ledger_event_id { + data["ledger_event_id"] = json!(event_id); + } + } + + (data, assistant_message) +} diff --git a/src/crates/execution/agent-runtime/tests/sdk_smoke.rs b/src/crates/execution/agent-runtime/tests/sdk_smoke.rs index 1dbc95aa10..090fc73a5a 100644 --- a/src/crates/execution/agent-runtime/tests/sdk_smoke.rs +++ b/src/crates/execution/agent-runtime/tests/sdk_smoke.rs @@ -53,7 +53,7 @@ struct FakeSessionClosePort { fn sdk_facade_exposes_versioned_preview_compatibility_contract() { let compatibility = AgentRuntimeSdkCompatibility::current(); - assert_eq!(compatibility.api_version, 2); + assert_eq!(compatibility.api_version, 3); assert_eq!(compatibility.crate_version, env!("CARGO_PKG_VERSION")); assert_eq!(compatibility.stability, AgentRuntimeSdkStability::Preview); } diff --git a/src/crates/interfaces/acp/src/runtime/prompt.rs b/src/crates/interfaces/acp/src/runtime/prompt.rs index 5b300dd70f..6e4c44cef4 100644 --- a/src/crates/interfaces/acp/src/runtime/prompt.rs +++ b/src/crates/interfaces/acp/src/runtime/prompt.rs @@ -106,6 +106,7 @@ fn dialog_turn_request(session: &AcpSessionState, prompt: ParsedPrompt) -> Agent message: prompt.user_message, original_message: prompt.original_user_message, turn_id: None, + execution: Default::default(), agent_type: session.mode_id.clone(), workspace_path: Some(session.cwd.clone()), remote_connection_id: None, diff --git a/src/crates/interfaces/sdk-host/src/host.rs b/src/crates/interfaces/sdk-host/src/host.rs index a2de12e757..d6164274e1 100644 --- a/src/crates/interfaces/sdk-host/src/host.rs +++ b/src/crates/interfaces/sdk-host/src/host.rs @@ -984,6 +984,7 @@ impl SdkHostConnection { message: params.prompt, original_message: None, turn_id: None, + execution: Default::default(), agent_type, workspace_path: Some(session.workspace_path.clone()), remote_connection_id: session.remote_connection_id.clone(), diff --git a/src/web-ui/src/flow_chat/components/ChatInput.tsx b/src/web-ui/src/flow_chat/components/ChatInput.tsx index 65547bf4ab..cdb34a8fc1 100644 --- a/src/web-ui/src/flow_chat/components/ChatInput.tsx +++ b/src/web-ui/src/flow_chat/components/ChatInput.tsx @@ -3827,6 +3827,11 @@ export const ChatInput: React.FC = ({ ); } if (!submissionTargetIsCurrent()) return true; + const executionTarget = expanded.executionTarget; + if (executionTarget.kind === 'fresh_external_subagent' && contexts.length > 0) { + notificationService.warning(t('chatInput.externalCommandContextUnsupported')); + return true; + } const expandedCharCount = getCharacterCount(expanded.content); if (expandedCharCount > CHAT_INPUT_CONFIG.largePaste.maxMessageChars) { notificationService.error( @@ -3860,7 +3865,12 @@ export const ChatInput: React.FC = ({ setSelectedNonExternalSlashCommand(undefined); setSelectedNonExternalSlashCandidateId(undefined); } - await sendMessage(expanded.content, { displayMessage: originalMessage }); + await sendMessage(expanded.content, { + displayMessage: originalMessage, + ...(executionTarget.kind === 'fresh_external_subagent' + ? { execution: executionTarget } + : {}), + }); if (!submissionTargetIsCurrent()) return true; if (composerCleared && inputValueRef.current === '') { dispatchInput({ type: 'DEACTIVATE' }); @@ -3904,7 +3914,7 @@ export const ChatInput: React.FC = ({ ); } return true; - }, [addToHistory, clearPendingLargePastes, confirmPromptCacheGuardIfNeeded, dispatchInput, effectiveTargetSessionId, externalPromptCommands, externalPromptCommandsIssue, externalPromptCommandsLoading, externalPromptCommandsPending, getSlashPickerItems, refreshExternalPromptCommands, replacePendingLargePastes, selectedExternalPromptCandidateId, selectedNonExternalSlashCandidateId, selectedNonExternalSlashCommand, sendMessage, sessionBoundWorkspacePath, setQueuedInput, t]); + }, [addToHistory, clearPendingLargePastes, confirmPromptCacheGuardIfNeeded, contexts, dispatchInput, effectiveTargetSessionId, externalPromptCommands, externalPromptCommandsIssue, externalPromptCommandsLoading, externalPromptCommandsPending, getSlashPickerItems, refreshExternalPromptCommands, replacePendingLargePastes, selectedExternalPromptCandidateId, selectedNonExternalSlashCandidateId, selectedNonExternalSlashCommand, sendMessage, sessionBoundWorkspacePath, setQueuedInput, t]); const handleCancelCurrentTask = useCallback(async () => { if (effectiveTargetSessionId) { diff --git a/src/web-ui/src/flow_chat/hooks/useMessageSender.ts b/src/web-ui/src/flow_chat/hooks/useMessageSender.ts index 062ee1a7ed..793a924099 100644 --- a/src/web-ui/src/flow_chat/hooks/useMessageSender.ts +++ b/src/web-ui/src/flow_chat/hooks/useMessageSender.ts @@ -24,6 +24,7 @@ import { composerPresentationSessionReferences, type ComposerPresentation, } from '../utils/composerPresentation'; +import type { AgentDialogTurnExecution } from '@/infrastructure/api/service-api/AgentAPI'; const log = createLogger('FlowChat'); @@ -61,6 +62,7 @@ interface UseMessageSenderReturn { options?: { displayMessage?: string; composerPresentation?: ComposerPresentation | null; + execution?: AgentDialogTurnExecution; } ) => Promise; /** Whether a send is in progress */ @@ -84,6 +86,7 @@ export function useMessageSender(props: UseMessageSenderProps): UseMessageSender options?: { displayMessage?: string; composerPresentation?: ComposerPresentation | null; + execution?: AgentDialogTurnExecution; } ) => { if (!message.trim()) { @@ -115,6 +118,9 @@ export function useMessageSender(props: UseMessageSenderProps): UseMessageSender try { const flowChatManager = FlowChatManager.getInstance(); let agentTypeForSend = currentAgentType || 'agentic'; + if (options?.execution?.kind === 'fresh_external_subagent' && contexts.length > 0) { + throw new Error('External subagent command delegation does not accept composer context'); + } if (!sessionId) { const agentType = currentAgentType || 'agentic'; @@ -197,6 +203,7 @@ export function useMessageSender(props: UseMessageSenderProps): UseMessageSender { ...(imagePayload ?? {}), ...(userMessageMetadata ? { userMessageMetadata } : {}), + ...(options?.execution ? { execution: options.execution } : {}), onSessionConflictRetryStart: () => { onSessionConflictRetryStart?.({ sessionId: sessionId!, diff --git a/src/web-ui/src/flow_chat/services/FlowChatManager.ts b/src/web-ui/src/flow_chat/services/FlowChatManager.ts index 7353cf62cd..662e77a9fa 100644 --- a/src/web-ui/src/flow_chat/services/FlowChatManager.ts +++ b/src/web-ui/src/flow_chat/services/FlowChatManager.ts @@ -662,6 +662,7 @@ export class FlowChatManager { imageContexts?: import('@/infrastructure/api/service-api/ImageContextTypes').ImageContextData[]; imageDisplayData?: Array<{ id: string; name: string; dataUrl?: string; imagePath?: string; mimeType?: string }>; userMessageMetadata?: Record; + execution?: import('@/infrastructure/api/service-api/AgentAPI').AgentDialogTurnExecution; turnId?: string; preserveTurnOnStartError?: boolean; onSessionConflictRetryStart?: () => void; diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.test.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.test.ts index a54508c429..f80fcfea3b 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.test.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.test.ts @@ -345,6 +345,32 @@ describe('MessageModule session writer conflict', () => { retryAction.onClick(); expect(mockEnsureBackendSession).toHaveBeenCalledTimes(1); }); + + it('does not turn external command delegation into a pending queue item', async () => { + const sessionId = 'session-delegated-command'; + const { context } = conflictContext(sessionId); + mockGetCurrentState.mockReturnValue('processing'); + + await expect(sendMessage( + context, + 'expanded prompt', + sessionId, + '/review', + undefined, + undefined, + { + execution: { + kind: 'fresh_external_subagent', + ecosystemId: 'opencode', + logicalId: 'reviewer', + }, + }, + )).rejects.toThrow('requires an idle session'); + + expect(mockPendingEnqueue).not.toHaveBeenCalled(); + expect(mockEnsureBackendSession).not.toHaveBeenCalled(); + expect(mockStartDialogTurn).not.toHaveBeenCalled(); + }); }); describe('MessageModule cancellation', () => { diff --git a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts index 29120109e8..cd43344fc9 100644 --- a/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts +++ b/src/web-ui/src/flow_chat/services/flow-chat-manager/MessageModule.ts @@ -213,6 +213,7 @@ export async function sendMessage( */ bypassPendingQueue?: boolean; userMessageMetadata?: Record; + execution?: import('@/infrastructure/api/service-api/AgentAPI').AgentDialogTurnExecution; turnId?: string; preserveTurnOnStartError?: boolean; onSessionConflictRetryStart?: () => void; @@ -360,6 +361,9 @@ export async function sendMessage( const hasPendingQueue = pendingQueueManager.list(sessionId).length > 0; if (sessionBusy || hasPendingQueue) { + if (options?.execution?.kind === 'fresh_external_subagent') { + throw new Error('External subagent command delegation requires an idle session'); + } if (await appendToRunningDispatch()) { return; } @@ -414,6 +418,10 @@ export async function sendMessage( const currentAgentType = (agentType?.trim() || refreshedSession.mode || 'agentic').trim(); const acpClientId = acpClientIdFromMode(currentAgentType); const isDispatched = isNonLocalDispatchTarget(refreshedSession.config.dispatchTarget); + const delegatesExternalSubagent = options?.execution?.kind === 'fresh_external_subagent'; + if (delegatesExternalSubagent && (acpClientId || isDispatched)) { + throw new Error('External subagent command delegation requires the local BitFun runtime'); + } if ( !acpClientId && @@ -700,6 +708,7 @@ export async function sendMessage( remoteSshHost: updatedSession.remoteSshHost, imageContexts: options?.imageContexts, userMessageMetadata: options?.userMessageMetadata, + execution: options?.execution, }); context.flowChatStore.updateSessionLastSubmittedMode(sessionId, currentAgentType); } catch (error: any) { @@ -723,6 +732,7 @@ export async function sendMessage( remoteSshHost: updatedSession.remoteSshHost, imageContexts: options?.imageContexts, userMessageMetadata: options?.userMessageMetadata, + execution: options?.execution, }); context.flowChatStore.updateSessionLastSubmittedMode(sessionId, currentAgentType); } else { diff --git a/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts b/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts index 2479933c00..4dc6192955 100644 --- a/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts +++ b/src/web-ui/src/infrastructure/api/service-api/AgentAPI.ts @@ -84,6 +84,7 @@ export interface StartDialogTurnRequest { userInput: string; originalUserInput?: string; turnId?: string; + execution?: AgentDialogTurnExecution; agentType: string; /** Concrete root where this session executes. */ workspacePath?: string; @@ -96,6 +97,14 @@ export interface StartDialogTurnRequest { userMessageMetadata?: Record; } +export type AgentDialogTurnExecution = + | { kind: 'standard' } + | { + kind: 'fresh_external_subagent'; + ecosystemId: string; + logicalId: string; + }; + export interface StartDialogTurnResponse { success: boolean; message: string; diff --git a/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.test.ts b/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.test.ts index 63dd9d264a..d963b02217 100644 --- a/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.test.ts +++ b/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.test.ts @@ -146,9 +146,17 @@ describe('ExternalSourcesAPI', () => { }); it('expands a prompt command only with the selected candidate and behavior version', async () => { - invokeMock.mockResolvedValueOnce({ state: 'ready', content: 'expanded prompt' }); + invokeMock.mockResolvedValueOnce({ + state: 'ready', + content: 'expanded prompt', + executionTarget: { + kind: 'fresh_external_subagent', + ecosystemId: 'opencode', + logicalId: 'reviewer', + }, + }); - await externalSourcesAPI.expandPromptCommand( + const outcome = await externalSourcesAPI.expandPromptCommand( 'D:/workspace/project', 'review', 'focus on auth', @@ -170,6 +178,15 @@ describe('ExternalSourcesAPI', () => { }, ); + expect(outcome).toMatchObject({ + state: 'ready', + executionTarget: { + kind: 'fresh_external_subagent', + ecosystemId: 'opencode', + logicalId: 'reviewer', + }, + }); + expect(invokeMock).toHaveBeenCalledWith('expand_external_prompt_command_command', { request: { workspacePath: 'D:/workspace/project', @@ -193,6 +210,29 @@ describe('ExternalSourcesAPI', () => { }); }); + it.each([ + ['a missing execution target', { state: 'ready', content: 'expanded prompt' }], + ['an unknown execution target', { + state: 'ready', + content: 'expanded prompt', + executionTarget: { kind: 'future_target' }, + }], + ])('rejects prompt command expansion with %s', async (_label, response) => { + invokeMock.mockResolvedValueOnce(response); + + await expect(externalSourcesAPI.expandPromptCommand( + 'D:/workspace/project', + 'review', + '', + 'opencode.commands:project:review', + 'behavior-v1', + [], + )).rejects.toMatchObject({ + code: 'invalid_response', + retryable: false, + }); + }); + it('projects native prompt command conflicts through the shared control plane', async () => { invokeMock.mockResolvedValueOnce({ preferenceRevision: 3, conflicts: [] }); const nativeCommands = [{ diff --git a/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.ts b/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.ts index 087b5fc266..1237d3ef7a 100644 --- a/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.ts +++ b/src/web-ui/src/infrastructure/api/service-api/ExternalSourcesAPI.ts @@ -166,6 +166,7 @@ export interface ExternalSourceCatalogSnapshot { commandDescription: string; sourceScope: ExternalSourceScope; sourceLocation: string; + executionTarget?: PromptCommandExecutionTarget; availability: PromptCommandAvailability; }>; }>; @@ -212,8 +213,20 @@ export interface PromptCommandShellReviewDecision { expectedPreferenceRevision: number; } +export type PromptCommandExecutionTarget = + | { kind: 'inline' } + | { + kind: 'fresh_external_subagent'; + ecosystemId: string; + logicalId: string; + }; + export type ExternalPromptCommandInvocationOutcome = - | { state: 'ready'; content: string } + | { + state: 'ready'; + content: string; + executionTarget: PromptCommandExecutionTarget; + } | { state: 'review_required'; review: PromptCommandShellReviewPlan }; export type ExternalSubagentActivation = @@ -600,6 +613,52 @@ export class ExternalSourceApiError extends Error { } } +function normalizePromptCommandInvocationOutcome( + value: unknown, +): ExternalPromptCommandInvocationOutcome { + if (!value || typeof value !== 'object') { + throw new ExternalSourceApiError( + 'invalid_response', + 'Prompt command expansion response was invalid', + false, + ); + } + const outcome = value as Record; + if (outcome.state === 'review_required' && outcome.review && typeof outcome.review === 'object') { + return value as ExternalPromptCommandInvocationOutcome; + } + if (outcome.state !== 'ready' || typeof outcome.content !== 'string') { + throw new ExternalSourceApiError( + 'invalid_response', + 'Prompt command expansion response was invalid', + false, + ); + } + const target = outcome.executionTarget; + if (!target || typeof target !== 'object') { + throw new ExternalSourceApiError( + 'invalid_response', + 'Prompt command expansion execution target was invalid', + false, + ); + } + const targetRecord = target as Record; + const validInline = targetRecord.kind === 'inline'; + const validExternalSubagent = targetRecord.kind === 'fresh_external_subagent' + && typeof targetRecord.ecosystemId === 'string' + && targetRecord.ecosystemId.trim().length > 0 + && typeof targetRecord.logicalId === 'string' + && targetRecord.logicalId.trim().length > 0; + if (!validInline && !validExternalSubagent) { + throw new ExternalSourceApiError( + 'invalid_response', + 'Prompt command expansion execution target was invalid', + false, + ); + } + return value as ExternalPromptCommandInvocationOutcome; +} + export interface NativePromptCommandDescriptor { commandName: string; candidateId: string; @@ -1309,7 +1368,7 @@ export const externalSourcesAPI = { }; }, - expandPromptCommand( + async expandPromptCommand( workspacePath: string | undefined, name: string, argumentsText: string, @@ -1322,7 +1381,7 @@ export const externalSourcesAPI = { }, shellReviewDecision?: PromptCommandShellReviewDecision, ) { - return invokeExternalSourceCommand( + const outcome = await invokeExternalSourceCommand( 'expand_external_prompt_command_command', { request: { @@ -1340,6 +1399,7 @@ export const externalSourcesAPI = { }, }, ); + return normalizePromptCommandInvocationOutcome(outcome); }, getNativePromptCommandConflicts( diff --git a/src/web-ui/src/locales/en-US/flow-chat.json b/src/web-ui/src/locales/en-US/flow-chat.json index a54eff624f..19ca9fcc76 100644 --- a/src/web-ui/src/locales/en-US/flow-chat.json +++ b/src/web-ui/src/locales/en-US/flow-chat.json @@ -782,6 +782,7 @@ "currentMode": "Current mode: {{mode}}", "noMatchingMode": "No matching mode", "noMatchingCommand": "No matching command", + "externalCommandContextUnsupported": "Remove composer attachments and references before running this delegated command.", "nativeCommandChoiceNotSaved": "This command will run once, but your conflict choice could not be saved.", "nativeCommandReconfirmationRequired": "The external command you previously chose is no longer available. Choose the BitFun command from the slash menu to run it.", "selectHint": "↑↓ Select · Enter Confirm · Esc Cancel", diff --git a/src/web-ui/src/locales/zh-CN/flow-chat.json b/src/web-ui/src/locales/zh-CN/flow-chat.json index 899aef7195..5acc22cc05 100644 --- a/src/web-ui/src/locales/zh-CN/flow-chat.json +++ b/src/web-ui/src/locales/zh-CN/flow-chat.json @@ -776,6 +776,7 @@ "currentMode": "当前模式: {{mode}}", "noMatchingMode": "没有匹配的模式", "noMatchingCommand": "没有匹配的命令", + "externalCommandContextUnsupported": "运行此委派命令前,请先移除输入框中的附件和引用。", "nativeCommandChoiceNotSaved": "本次仍会执行该命令,但未能保存这次冲突选择。", "nativeCommandReconfirmationRequired": "之前选择的外部命令已不可用。请从斜杠菜单中选择 BitFun 命令后再执行。", "selectHint": "↑↓ 选择 · Enter 确认 · Esc 取消", diff --git a/src/web-ui/src/locales/zh-TW/flow-chat.json b/src/web-ui/src/locales/zh-TW/flow-chat.json index 2d452f79b3..2d3c299f2c 100644 --- a/src/web-ui/src/locales/zh-TW/flow-chat.json +++ b/src/web-ui/src/locales/zh-TW/flow-chat.json @@ -776,6 +776,7 @@ "currentMode": "目前模式: {{mode}}", "noMatchingMode": "沒有匹配的模式", "noMatchingCommand": "沒有匹配的命令", + "externalCommandContextUnsupported": "執行此委派命令前,請先移除輸入框中的附件與引用。", "nativeCommandChoiceNotSaved": "本次仍會執行該命令,但未能儲存這次衝突選擇。", "nativeCommandReconfirmationRequired": "先前選擇的外部命令已無法使用。請從斜線選單中選擇 BitFun 命令後再執行。", "selectHint": "↑↓ 選擇 · Enter 確認 · Esc 取消",