diff --git a/docs/tools.md b/docs/tools.md index 8018762..10d48b1 100644 --- a/docs/tools.md +++ b/docs/tools.md @@ -39,7 +39,8 @@ Coding Code 的工具系统是 Agent 与外部世界交互的核心机制。本 | 工具 | 功能 | 关键参数 | |---|---|---| -| `dispatch_agent` | 将任务委派给运行时注册的子智能体 | `agent: string`, `prompt: string` | +| `spawn_agent` | 启动一个后台子智能体并立即返回其会话 id,结果完成后自动注入本会话 | `agentName: string`, `prompt: string`, `model?: string`, `systemPrompt?: string` | +| `wait_agent` | 等待子智能体到达终态,返回 `completed` / `failed` / `timeout` | `sessionId: string`, `timeoutMs?: number`(夹在 `[10000, 3600000]`,默认 `30000`) | --- @@ -50,7 +51,7 @@ Coding Code 的工具系统是 Agent 与外部世界交互的核心机制。本 - **Core 工具**:始终可用,在启动时注册。包括上述所有内置工具。 - **MCP 工具**:从 MCP 服务自动导入和注册。名称空间化为 `serverName:toolName` 格式,避免不同服务间的工具名冲突。 -Agent 在一次运行开始时注册内置工具、项目 MCP 工具和 `dispatch_agent`。plan 模式通过独立的 `PLAN_PROFILE_ALLOWED_TOOLS` 策略过滤工具。 +Agent 在一次运行开始时注册内置工具、项目 MCP 工具和 `spawn_agent` / `wait_agent`。plan 模式通过独立的 `PLAN_PROFILE_ALLOWED_TOOLS` 策略过滤工具。 --- diff --git a/packages/codingcode/src/agent/agent.ts b/packages/codingcode/src/agent/agent.ts index b26a890..ee39672 100644 --- a/packages/codingcode/src/agent/agent.ts +++ b/packages/codingcode/src/agent/agent.ts @@ -7,6 +7,7 @@ import { ApprovalService } from '../approval/port.js'; import { CheckpointService } from '../checkpoint/port.js'; import { ContextService } from '../context/port.js'; import { EventSinkService } from '../sink/port.js'; +import { MailboxService } from '../session/mailbox.js'; import { HookService } from '../hooks/port.js'; import { LLMService } from '../llm/port.js'; import { McpService } from '../mcp/port.js'; @@ -17,7 +18,7 @@ import { SkillService } from '../skills/port.js'; import { TodoService } from '../todo/port.js'; import { ToolExecutorService } from '../tools/port.js'; import { buildSystemPrompt } from './prompt.js'; -import type { FrameBody, FrameError, ResponseMeta, ToolOutcome, Transition } from '../contracts/frame.js'; +import type { EndTransition, FrameBody, FrameError, ResponseMeta, ToolOutcome } from '../contracts/frame.js'; import { isTurnEnd } from '../contracts/frame.js'; import type { SessionRef } from '../contracts/session.js'; import type { ToolCatalog, ToolResult } from '../contracts/tool.js'; @@ -51,6 +52,7 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { const mcp = yield* McpService; const context = yield* ContextService; const sink = yield* EventSinkService; + const mailbox = yield* MailboxService; const memory = yield* MemoryService; const llm = yield* LLMService; const rules = yield* RulesService; @@ -196,7 +198,8 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { const { tools, lookup: toolLookup } = catalog; let ended = false; - const offerEnd = (transition: Extract) => + let deliveryPhase: 'currentTurn' | 'nextTurn' = 'currentTurn'; + const offerEnd = (transition: EndTransition) => Effect.sync(() => { if (ended) return; ended = true; @@ -238,6 +241,14 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { currentTurnId: state.currentTurnId, }; + const mayDrain = deliveryPhase === 'currentTurn' && step > 0; + if (mayDrain) { + for (const item of yield* mailbox.drain(state.sessionId)) { + const ev = yield* session.recordSubagentResult(state, item); + yield* context.absorb(sessionRef, [ev]); + } + } + const history = yield* Effect.either(context.getHistory(sessionRef, model)); if (Either.isLeft(history)) { yield* offerEnd({ to: 'end', reason: 'error', error: toFrameError(history.left) }); @@ -284,6 +295,7 @@ export const AgentLayer = Layer.effect(AgentService, Effect.gen(function* () { if (toolCalls.length === 0) { const assistantEv = yield* session.recordAssistant(state, content, [], responded.usage); yield* context.absorb(sessionRef, [assistantEv]); + deliveryPhase = 'nextTurn'; const stopDecision = yield* hooks.emitDecision('agent.turn.stop', { sessionId: sid, content, turnId: state.currentTurnId, projectPath }); if (stopDecision && stopDecision.decision === 'continue') { diff --git a/packages/codingcode/src/agent/profile.ts b/packages/codingcode/src/agent/profile.ts index e12e9b2..cd8bd09 100644 --- a/packages/codingcode/src/agent/profile.ts +++ b/packages/codingcode/src/agent/profile.ts @@ -27,7 +27,7 @@ export const BUILD_PROMPT = `You are a coding assistant —an AI agent that help 7. For complex or broad tasks (understanding a whole module, cross-file analysis, comprehensive search): a. Briefly assess the task scope using your own reasoning —do not use tools for exploration at this stage, as that would consume your limited context window. b. If you can clearly handle it without extensive file reading or searching, proceed yourself. - c. Otherwise delegate the exploration with dispatch_agent: give the subagent a short agentName and a self-contained prompt. The subagent shares your working directory, so keep the delegated write set disjoint from your own. + c. Otherwise spawn_agent: give the subagent a short agentName, a self-contained prompt, and a write set that does not overlap yours. Then keep working on your own part; call wait_agent only when its result blocks your next step, and give the wait a generous timeout instead of polling. ## Using your tools - **Prefer dedicated tools over shell commands.** Use read_file instead of cat, edit_file instead of sed, search_code instead of grep. Dedicated tools give the user better visibility into your work. @@ -127,7 +127,8 @@ export const BUILD_TOOL_NAMES: readonly string[] = [ 'fetch_url', 'web_search', 'todo_write', - 'dispatch_agent', + 'spawn_agent', + 'wait_agent', ]; export function isPlanProfile(p: { name: string } | null | undefined): boolean { diff --git a/packages/codingcode/src/context/context.ts b/packages/codingcode/src/context/context.ts index 9af0219..d641ece 100644 --- a/packages/codingcode/src/context/context.ts +++ b/packages/codingcode/src/context/context.ts @@ -151,6 +151,9 @@ export function buildContextMessages( case 'summary': messages.push({ role: 'system', name: 'compacted_history', content: event.summaryText }); break; + case 'subagent_result': + messages.push({ role: 'assistant', content: event.content }); + break; } } diff --git a/packages/codingcode/src/contracts/frame.ts b/packages/codingcode/src/contracts/frame.ts index eedd14b..386cbcb 100644 --- a/packages/codingcode/src/contracts/frame.ts +++ b/packages/codingcode/src/contracts/frame.ts @@ -53,6 +53,12 @@ export type RuntimeEvent = readonly id: string; readonly tool: string; readonly args: Readonly>; + } + | { + readonly type: 'subagent_event'; + readonly sessionId: string; + readonly agentName: string; + readonly status: 'spawned' | 'completed' | 'failed'; }; export interface Fatal { @@ -66,8 +72,7 @@ export type FrameBody = | { readonly family: 'fatal'; readonly fatal: Fatal }; export type Frame = Envelope & FrameBody; - -type EndTransition = Extract; +export type EndTransition = Extract; export function isTurnEnd( body: FrameBody diff --git a/packages/codingcode/src/contracts/session.ts b/packages/codingcode/src/contracts/session.ts index 67e5a91..3fe0ffd 100644 --- a/packages/codingcode/src/contracts/session.ts +++ b/packages/codingcode/src/contracts/session.ts @@ -58,6 +58,13 @@ export interface CompactEvent { endTurnId: number; } +export interface SubagentResultEvent { + type: 'subagent_result'; + sessionId: string; + agentName: string; + content: string; +} + export type SessionEvent = | SessionMetaEvent | UserEvent @@ -65,7 +72,8 @@ export type SessionEvent = | ToolResultEvent | SummaryEvent | RollbackEvent - | CompactEvent; + | CompactEvent + | SubagentResultEvent; export interface SessionSummary extends SessionMetaEvent { updatedAt: string; diff --git a/packages/codingcode/src/infra/config.ts b/packages/codingcode/src/infra/config.ts index 6409bce..0d56da1 100644 --- a/packages/codingcode/src/infra/config.ts +++ b/packages/codingcode/src/infra/config.ts @@ -18,6 +18,10 @@ export interface ActiveModelConfig { apiKeyEnv: string; } +export interface SubagentConfig { + maxBackground: number; +} + export interface AppConfig { server: { port: number; @@ -29,6 +33,7 @@ export interface AppConfig { permissionMode: string; context: ContextConfig; memory: MemoryConfig; + subagent: SubagentConfig; } const DEFAULT_CONTEXT: ContextConfig = { @@ -41,6 +46,10 @@ export const DEFAULT_MEMORY: MemoryConfig = { promptMaxBytes: 8192, }; +export const DEFAULT_SUBAGENT: SubagentConfig = { + maxBackground: 4, +}; + export const DEFAULT_CONFIG: AppConfig = { server: { port: 8080, @@ -51,6 +60,7 @@ export const DEFAULT_CONFIG: AppConfig = { permissionMode: 'ask', context: DEFAULT_CONTEXT, memory: DEFAULT_MEMORY, + subagent: DEFAULT_SUBAGENT, }; function deepMerge>(base: T, override: Partial): T { diff --git a/packages/codingcode/src/layer.ts b/packages/codingcode/src/layer.ts index 4c41f06..a7637ef 100644 --- a/packages/codingcode/src/layer.ts +++ b/packages/codingcode/src/layer.ts @@ -16,6 +16,8 @@ import { MemoryLayer } from './memory/memory.js'; import { AgentLayer } from './agent/agent.js'; import { ToolEnvLayer } from './agent/tool-env.js'; import { SubagentRunnerLayer } from './subagent/subagent.js'; +import { SubagentRunRegistryLayer } from './subagent/registry.js'; +import { MailboxLayer } from './session/mailbox.js'; import { SchedulerLayer } from './scheduler/scheduler.js'; // base layers @@ -36,7 +38,7 @@ const MemoryWithDeps = MemoryLayer.pipe(Layer.provide(LlmLayer)); // agent 直接消费的宽服务集合 const AgentServiceLayers = Layer.mergeAll( - InfraLayer, SessionLayer, ToolExecutorWithDeps, ApprovalWithDeps, + InfraLayer, SessionLayer, MailboxLayer, ToolExecutorWithDeps, ApprovalWithDeps, ContextWithDeps, MemoryWithDeps, CheckpointLayer, LlmLayer, ); @@ -48,17 +50,25 @@ const AgentWithDeps = AgentLayer.pipe( // subagent runner (depends on agent) const SubagentWithDeps = SubagentRunnerLayer.pipe(Layer.provide(AgentWithDeps)); +// 运行注册表:要 runner 起子代理、要 mailbox 投递终态、要 sink 发 subagent_event 帧。 +// 不依赖 SessionLayer —— 它不写盘,写盘由父回合循环在 drain 点做。 +const SubagentRunRegistryWithDeps = SubagentRunRegistryLayer.pipe( + Layer.provide(Layer.mergeAll(SubagentWithDeps, MailboxLayer, EventSinkLayer)) +); + export const AppLayer = Layer.mergeAll( InfraLayer, LlmLayer, ApprovalWithDeps, SessionLayer, + MailboxLayer, ToolExecutorWithDeps, ContextWithDeps, MemoryWithDeps, CheckpointLayer, AgentWithDeps, SubagentWithDeps, + SubagentRunRegistryWithDeps, SchedulerLayer, EventSinkLayer, ); diff --git a/packages/codingcode/src/server/routes/sessions.ts b/packages/codingcode/src/server/routes/sessions.ts index 99b0159..ea4212b 100644 --- a/packages/codingcode/src/server/routes/sessions.ts +++ b/packages/codingcode/src/server/routes/sessions.ts @@ -7,6 +7,7 @@ import type { ProfileName } from '../../contracts/types.js'; import { SessionService } from '../../session/port.js'; import { computePaths } from '../../session/paths.js'; import { ContextService } from '../../context/port.js'; +import { MailboxService } from '../../session/mailbox.js'; import { estimatePromptTokensFrom } from '../../context/context.js'; import { CheckpointService } from '../../checkpoint/port.js'; import { activeModelId, setGlobalActive } from '../../infra/models.js'; @@ -143,8 +144,10 @@ export function registerSessionsRoutes(router: Hono, rt: ManagedRt): void { Effect.gen(function* () { const session = yield* SessionService; const context = yield* ContextService; + const mailbox = yield* MailboxService; yield* session.deleteSession(sessionId, cwd); yield* context.dispose(sessionId); + yield* mailbox.dispose(sessionId); }) as any ); return c.json({ ok: true }); diff --git a/packages/codingcode/src/session/mailbox.ts b/packages/codingcode/src/session/mailbox.ts new file mode 100644 index 0000000..a3d44ca --- /dev/null +++ b/packages/codingcode/src/session/mailbox.ts @@ -0,0 +1,48 @@ +import { Chunk, Context, Effect, Layer, Queue } from 'effect'; +import type { SubagentResultEvent } from '../contracts/session.js'; + +/** 入站暂存条目。当前唯一生产者是子代理终态 —— 不预留其它变体 */ +export type MailboxItem = SubagentResultEvent; + +export interface MailboxShape { + /** 投递;该会话还没有队列时自动建 */ + offer(sessionId: string, item: MailboxItem): Effect.Effect; + /** 取走当前全部待处理条目(非阻塞,空则返回空数组) */ + drain(sessionId: string): Effect.Effect>; + /** 会话删除时清空 */ + dispose(sessionId: string): Effect.Effect; +} + +export class MailboxService extends Context.Tag('Mailbox')() {} + +/** + * 会话的易失入站队列,与 transcript(持久出站)对称:同一把键 sessionId。 + * 全局单例,内部按收件人分区 —— 嵌套委派下 B 的终态进 mailbox[A]。 + */ +export const MailboxLayer = Layer.scoped( + MailboxService, + Effect.gen(function* () { + const queues = new Map>(); + const queueFor = (sessionId: string) => { + let q = queues.get(sessionId); + if (!q) { + q = Effect.runSync(Queue.unbounded()); + queues.set(sessionId, q); + } + return q; + }; + + yield* Effect.addFinalizer(() => Effect.sync(() => queues.clear())); + + return { + offer: (sessionId, item) => Queue.offer(queueFor(sessionId), item), + drain: (sessionId) => { + const q = queues.get(sessionId); + return q + ? Queue.takeAll(q).pipe(Effect.map(Chunk.toReadonlyArray)) + : Effect.succeed>([]); + }, + dispose: (sessionId) => Effect.sync(() => { queues.delete(sessionId); }), + }; + }) +); diff --git a/packages/codingcode/src/session/port.ts b/packages/codingcode/src/session/port.ts index 7de3840..a33d723 100644 --- a/packages/codingcode/src/session/port.ts +++ b/packages/codingcode/src/session/port.ts @@ -1,7 +1,7 @@ import { Context } from 'effect'; import type { Effect } from 'effect'; import type { AgentError } from '../core/error.js'; -import type { AssistantEvent, RollbackEvent, SessionCreateOptions, SessionEvent, SessionSummary, SessionStoreState, SummaryEvent, ToolResultEvent, UITurn, UserEvent } from '../contracts/session.js'; +import type { AssistantEvent, RollbackEvent, SessionCreateOptions, SessionEvent, SessionSummary, SessionStoreState, SubagentResultEvent, SummaryEvent, ToolResultEvent, UITurn, UserEvent } from '../contracts/session.js'; import type { TokenUsage, ProfileName } from '../contracts/types.js'; import type { PermissionMode } from '../contracts/permission.js'; @@ -17,6 +17,7 @@ export interface SessionShape { recordSystem(state: SessionStoreState, content: string): Effect.Effect; recordAssistant(state: SessionStoreState, content: string, toolCalls: AssistantEvent['toolCalls'], usage?: TokenUsage): Effect.Effect; recordToolResult(state: SessionStoreState, toolName: string, toolCallId: string, output: string): Effect.Effect; + recordSubagentResult(state: SessionStoreState, result: { sessionId: string; agentName: string; content: string }): Effect.Effect; appendSummary(state: SessionStoreState, summaryText: string, startTurnId: number, endTurnId: number): Effect.Effect; rollbackToTurn(state: SessionStoreState, throughTurnId: number, reason: string): Effect.Effect; readEvents(transcriptPath: string): Effect.Effect; diff --git a/packages/codingcode/src/session/session.ts b/packages/codingcode/src/session/session.ts index aa04436..6cfd58d 100644 --- a/packages/codingcode/src/session/session.ts +++ b/packages/codingcode/src/session/session.ts @@ -5,7 +5,7 @@ import { join, dirname } from 'path'; import { AgentError } from '../core/error.js'; import { encodeProjectPath } from '../core/path.js'; import { computePaths } from './paths.js'; -import type { SessionMetaEvent, UserEvent, AssistantEvent, ToolResultEvent, SummaryEvent, RollbackEvent, SessionEvent, SessionStoreState, SessionSummary, CompactEvent, UITurn } from '../contracts/session.js'; +import type { SessionMetaEvent, UserEvent, AssistantEvent, ToolResultEvent, SubagentResultEvent, SummaryEvent, RollbackEvent, SessionEvent, SessionStoreState, SessionSummary, CompactEvent, UITurn } from '../contracts/session.js'; import type { TokenUsage, ProfileName } from '../contracts/types.js'; import type { PermissionMode } from '../contracts/permission.js'; import { SessionService } from './port.js'; @@ -79,6 +79,7 @@ export function sessionEventsToTurns(events: SessionEvent[]): UITurn[] { for (const event of events) { if (event.type === 'session_meta') continue; if (event.type === 'compact' || event.type === 'rollback') continue; + if (event.type === 'subagent_result') continue; // 结果进模型上下文,不占用户视野 if (event.type === 'summary') { let turn = turnsMap.get(event.endTurnId); @@ -319,6 +320,29 @@ export const SessionLayer = Layer.effect( : new AgentError('SESSION_IO_ERROR', `Session write failed: ${String(e)}`, e), }); + // 由父回合循环在 drain 点调用:终态先入 mailbox,到这里才落盘。 + // 不写 state.usage —— 子代理的用量不算进父会话。 + const recordSubagentResult = ( + state: SessionStoreState, + result: { sessionId: string; agentName: string; content: string } + ): Effect.Effect => + Effect.try({ + try: () => { + const event: SubagentResultEvent = { + type: 'subagent_result', + sessionId: result.sessionId, + agentName: result.agentName, + content: result.content, + }; + appendLine(pathsFromState(state).transcriptPath, event); + return event; + }, + catch: (e) => + e instanceof AgentError + ? e + : new AgentError('SESSION_IO_ERROR', `Session write failed: ${String(e)}`, e), + }); + const appendSummary = ( state: SessionStoreState, summaryText: string, @@ -435,6 +459,7 @@ export const SessionLayer = Layer.effect( recordSystem, recordAssistant, recordToolResult, + recordSubagentResult, appendSummary, rollbackToTurn, diff --git a/packages/codingcode/src/subagent/registry.ts b/packages/codingcode/src/subagent/registry.ts new file mode 100644 index 0000000..972b920 --- /dev/null +++ b/packages/codingcode/src/subagent/registry.ts @@ -0,0 +1,222 @@ +import { Context, Effect, Fiber, Layer, Option, Stream, SubscriptionRef } from 'effect'; +import { AgentError } from '../core/error.js'; +import { isTurnEnd } from '../contracts/frame.js'; +import type { EndTransition, FrameBody } from '../contracts/frame.js'; +import type { ProfileName } from '../contracts/types.js'; +import { estimateTokensForContent } from '../context/tokens.js'; +import { loadConfig } from '../infra/config.js'; +import { MailboxService } from '../session/mailbox.js'; +import { EventSinkService } from '../sink/port.js'; +import { SubagentRunnerService } from './port.js'; + +export type SubagentRunStatus = + | { readonly kind: 'running' } + | { readonly kind: 'ended'; readonly end: EndTransition }; + +/** wait 只回信号、不回内容 —— 内容已由父回合 drain 时写进父 transcript */ +export type WaitOutcome = 'completed' | 'failed' | 'timeout'; + +export const SUBAGENT_WAIT_MIN_MS = 10_000; +export const SUBAGENT_WAIT_DEFAULT_MS = 30_000; +export const SUBAGENT_WAIT_MAX_MS = 3_600_000; + +export interface SpawnOptions { + prompt: string; + agentName: string; + parentSessionId: string; + parentCwd: string; + parentProfile: ProfileName; + /** 已由调用方解析过的模型 id(含清单校验与父模型回退) */ + model: string; + systemPrompt?: string; +} + +/** 阶段二只有两个方法:调用者分别是 spawn.ts 与 wait.ts。 + * 计数 = 模块私有函数 countRunning(只服务 spawn 的配额判断); + * statuses / stopAll 与它们的调用者(HTTP 路由、桌面停止下拉)一起放阶段三。 */ +export interface SubagentRunRegistryShape { + spawn(opts: SpawnOptions): Effect.Effect<{ sessionId: string; agentName: string }, AgentError>; + wait(sessionId: string, timeoutMs: number): Effect.Effect; +} + +export class SubagentRunRegistryService extends Context.Tag('SubagentRunRegistry')< + SubagentRunRegistryService, + SubagentRunRegistryShape +>() {} + +interface SubagentRun { + readonly sessionId: string; + readonly parentSessionId: string; + readonly agentName: string; + readonly status: SubscriptionRef.SubscriptionRef; + readonly abort: AbortController; + fiber?: Fiber.Fiber; +} + +const BODY_TOKEN_BUDGET = 900; + +/** 二分收敛到不超过预算的最长前缀,避免逐字扫描 */ +function truncateToTokens(text: string, budget: number): string { + if (estimateTokensForContent(text) <= budget) return text; + let lo = 0; + let hi = text.length; + while (lo < hi) { + const mid = Math.ceil((lo + hi) / 2); + if (estimateTokensForContent(text.slice(0, mid)) <= budget) lo = mid; + else hi = mid - 1; + } + return `${text.slice(0, lo)}\n…[truncated]`; +} + +function renderResult(run: SubagentRun, outcome: { end: EndTransition; content: string }): string { + const body = outcome.end.reason === 'done' + ? outcome.content + : `${outcome.content}\nThe subagent did not finish. Spawn it again if the task is still needed.`; + return [ + 'Message Type: FINAL_ANSWER', + `Task name: ${run.parentSessionId}`, + `Sender: ${run.sessionId}`, + 'Payload:', + truncateToTokens(body, BODY_TOKEN_BUDGET), + ].join('\n'); +} + +/** 非 done 的终态统一归一成一条 error 帧,形状仍是帧契约的 end */ +const failedEnd = (message: string): EndTransition => + ({ to: 'end', reason: 'error', error: { message, code: 'SUBAGENT_FAILED' } }); + +/** 帧流的唯一读者:攒 text_delta 取最终输出,认 isTurnEnd 拿终态,其余帧丢弃 */ +const consume = (stream: AsyncGenerator) => + Effect.tryPromise({ + try: async (): Promise<{ end: EndTransition; content: string }> => { + let content = ''; + for await (const body of stream) { + if (body.family === 'event') { + if (body.event.type === 'text_delta') content += body.event.text; + continue; + } + if (isTurnEnd(body)) { + const end = body.transition; + if (end.reason === 'done') return { end, content: content || '(subagent completed without output)' }; + if (end.reason === 'error') return { end, content: end.error.message }; + const message = end.reason === 'maxSteps' + ? 'subagent exhausted its step budget before finishing' + : 'subagent was aborted'; + return { end: failedEnd(message), content: message }; + } + } + const message = 'subagent stream ended without a terminal frame'; + return { end: failedEnd(message), content: message }; + }, + catch: (e) => new AgentError('TOOL_EXECUTION_FAILED', e instanceof Error ? e.message : String(e)), + }); + +export const SubagentRunRegistryLayer = Layer.scoped( + SubagentRunRegistryService, + Effect.gen(function* () { + const mailbox = yield* MailboxService; + const runner = yield* SubagentRunnerService; + const sink = yield* EventSinkService; + const runs = new Map(); // 键 = 子会话 sessionId;parentSessionId 只是条目上的字段 + + /** 帧直投父会话的出站队列:EventSink 的键就是收件人会话,不需要任何回调透传 */ + const emitSubagent = ( + parentSessionId: string, sessionId: string, agentName: string, + status: 'spawned' | 'completed' | 'failed', + ) => sink.emit(parentSessionId, { + family: 'event', + event: { type: 'subagent_event', sessionId, agentName, status }, + }); + + /** 同步读当前值:SubscriptionRef 的 get 随时可读、读不走 */ + const currentStatus = (run: SubagentRun): SubagentRunStatus => + Effect.runSync(SubscriptionRef.get(run.status)); + + const countRunning = (parentSessionId: string): number => { + let n = 0; + for (const run of runs.values()) { + if (run.parentSessionId === parentSessionId && currentStatus(run).kind === 'running') n++; + } + return n; + }; + + // 只入队,不写盘:写盘由父回合在自己的 drain 点做(父回合循环手里才有 state) + const drainRun = (run: SubagentRun, stream: AsyncGenerator) => + Effect.gen(function* () { + const settled = yield* Effect.either(consume(stream)); + const outcome = settled._tag === 'Right' + ? settled.right + : { end: failedEnd(settled.left.message), content: settled.left.message }; + + yield* mailbox.offer(run.parentSessionId, { + type: 'subagent_result', + sessionId: run.sessionId, + agentName: run.agentName, + content: renderResult(run, outcome), + }).pipe(Effect.ignore); + + yield* SubscriptionRef.set(run.status, { kind: 'ended', end: outcome.end }); + yield* emitSubagent( + run.parentSessionId, run.sessionId, run.agentName, + outcome.end.reason === 'done' ? 'completed' : 'failed', + ); + }); + + const spawn = (opts: SpawnOptions) => + Effect.gen(function* () { + if (countRunning(opts.parentSessionId) >= loadConfig().subagent.maxBackground) { + return yield* Effect.fail( + new AgentError('TOOL_EXECUTION_FAILED', 'Concurrent subagent limit reached') + ); + } + + const abort = new AbortController(); + const { stream, sessionId } = yield* runner.runSubagent(opts.prompt, { + cwd: opts.parentCwd, + signal: abort.signal, // 子代理自己的 signal,与父回合无关 + activeProfile: opts.parentProfile, + permissionMode: 'bypass', + parentSessionId: opts.parentSessionId, + agentName: opts.agentName, + model: opts.model, + systemPrompt: opts.systemPrompt, + }); + + const status = yield* SubscriptionRef.make({ kind: 'running' }); + const run: SubagentRun = { + sessionId, parentSessionId: opts.parentSessionId, agentName: opts.agentName, + status, abort, + }; + runs.set(sessionId, run); + yield* emitSubagent(opts.parentSessionId, sessionId, opts.agentName, 'spawned'); + + run.fiber = yield* Effect.forkDaemon(drainRun(run, stream)); + return { sessionId, agentName: opts.agentName }; + }); + + const wait = (sessionId: string, timeoutMs: number): Effect.Effect => + Effect.gen(function* () { + const run = runs.get(sessionId); + if (!run) { + return yield* Effect.fail(new AgentError('TOOL_EXECUTION_FAILED', `Unknown subagent: ${sessionId}`)); + } + // 等状态通道的后续值;changes 的首次发射即当前值 ⇒ 已终态的 run 立即返回 + const settled = run.status.changes.pipe( + Stream.filterMap((s) => s.kind === 'ended' ? Option.some(s.end) : Option.none()), + Stream.runHead, + Effect.map((opt): WaitOutcome => + Option.isNone(opt) ? 'timeout' : (opt.value.reason === 'done' ? 'completed' : 'failed') + ) + ); + + return yield* Effect.race( + Effect.sleep(timeoutMs).pipe(Effect.as('timeout')), + settled + ); + }); + + yield* Effect.addFinalizer(() => Effect.sync(() => runs.clear())); + + return { spawn, wait }; + }) +); diff --git a/packages/codingcode/src/tools/catalog.ts b/packages/codingcode/src/tools/catalog.ts index 29b85b2..324d11e 100644 --- a/packages/codingcode/src/tools/catalog.ts +++ b/packages/codingcode/src/tools/catalog.ts @@ -15,7 +15,8 @@ import { globTool } from './domains/fs/glob.js'; import { webFetchTool } from './domains/web/fetch.js'; import { webSearchTool } from './domains/web/search.js'; import { todoWriteTool } from './domains/self/todo-write.js'; -import { dispatchAgentTool } from './domains/subagent/dispatch.js'; +import { spawnAgentTool } from './domains/subagent/spawn.js'; +import { waitAgentTool } from './domains/subagent/wait.js'; import { submitPlanTool } from './domains/subagent/submit-plan.js'; const ALL_TOOLS: ToolDefinition[] = [ @@ -28,7 +29,8 @@ const ALL_TOOLS: ToolDefinition[] = [ webFetchTool, webSearchTool, todoWriteTool, - dispatchAgentTool, + spawnAgentTool, + waitAgentTool, submitPlanTool, ]; diff --git a/packages/codingcode/src/tools/domains/subagent/dispatch.ts b/packages/codingcode/src/tools/domains/subagent/dispatch.ts deleted file mode 100644 index e479802..0000000 --- a/packages/codingcode/src/tools/domains/subagent/dispatch.ts +++ /dev/null @@ -1,101 +0,0 @@ -import { z } from 'zod'; -import { Effect } from 'effect'; -import { AgentError } from '../../../core/error.js'; -import { findModel } from '../../../infra/models.js'; -import type { ToolDefinition } from '../../types.js'; -import { HookService } from '../../../hooks/port.js'; -import { SubagentRunnerService } from '../../../subagent/port.js'; - -export const dispatchAgentTool: ToolDefinition = { - name: 'dispatch_agent', - concurrencySafe: false, - description: - 'Delegate a task to a subagent. The subagent runs in the same working directory as you and returns its final output. ' - + 'Keep the delegated write set disjoint from the files you edit yourself.', - parameters: z.object({ - agentName: z.string().min(1).describe('short nickname for the subagent; used for identification and display'), - prompt: z.string().min(1).describe('task description for the subagent'), - model: z.string().optional().describe('model id for the subagent; must exist in models.json, otherwise the model of the current turn is used'), - systemPrompt: z.string().optional().describe('replaces the middle section of the subagent system prompt; the environment block and system notes are kept'), - }), - execute: (args, ctx) => - Effect.gen(function* () { - const hooks = yield* HookService; - const runner = yield* SubagentRunnerService; - - const { agentName, prompt, model, systemPrompt } = args as { - agentName: string; prompt: string; model?: string; systemPrompt?: string; - }; - const projectPath = ctx?.projectPath || process.cwd(); - - if (!ctx?.activeProfile) { - return yield* Effect.fail( - new AgentError('CONFIG_MISSING', 'dispatch_agent requires the parent session activeProfile') - ); - } - - const parentSessionId = ctx?.sessionId; - // 子代理只跑在模型清单内的模型上,参数空或不在清单里都继承父回合的模型 - const requestedModel = model?.trim(); - const effectiveModel = requestedModel && findModel(requestedModel) ? requestedModel : ctx.model; - const spawnDecision = yield* hooks.emitDecision('agent.subagent.spawn.before', { - agentName, prompt, parentSessionId, projectPath, - }); - if (spawnDecision && spawnDecision.decision === 'deny') { - return yield* Effect.fail( - new AgentError('TOOL_NOT_ALLOWED', `Subagent spawn denied: ${spawnDecision.reason ?? 'no reason'}`) - ); - } - - const { stream, sessionId: childUuid } = yield* runner.runSubagent(prompt, { - cwd: projectPath, - signal: ctx?.signal, - activeProfile: ctx.activeProfile, - permissionMode: 'bypass', - parentSessionId, - agentName, - model: effectiveModel, - systemPrompt, - }); - - yield* hooks.emit('agent.subagent.spawn.after', { - childSessionId: childUuid, agentName, projectPath, - }); - - let didComplete = false; - const finalContent = yield* Effect.async((resume) => { - let content = ''; - (async () => { - try { - for await (const body of stream) { - if (body.family === 'event') { - if (body.event.type === 'text_delta') content += body.event.text; - continue; - } - if ( - body.family === 'transition' && - body.transition.to === 'end' && - body.transition.reason === 'error' - ) { - resume(Effect.fail(new AgentError('TOOL_EXECUTION_FAILED', `Subagent failed: ${body.transition.error.message}`))); - return; - } - } - didComplete = true; - resume(Effect.succeed(content || '(subagent completed without output)')); - } catch (e) { - const msg = e instanceof Error ? e.message : String(e); - resume(Effect.fail(new AgentError('TOOL_EXECUTION_FAILED', msg))); - } - })(); - }); - - if (didComplete) { - yield* hooks.emit('agent.subagent.complete', { - childSessionId: childUuid, agentName, status: 'done', projectPath, - }).pipe(Effect.ignore); - } - - return finalContent; - }), -}; diff --git a/packages/codingcode/src/tools/domains/subagent/spawn.ts b/packages/codingcode/src/tools/domains/subagent/spawn.ts new file mode 100644 index 0000000..2e175e6 --- /dev/null +++ b/packages/codingcode/src/tools/domains/subagent/spawn.ts @@ -0,0 +1,62 @@ +import { z } from 'zod'; +import { Effect } from 'effect'; +import { AgentError } from '../../../core/error.js'; +import { findModel } from '../../../infra/models.js'; +import type { ToolDefinition } from '../../types.js'; +import { HookService } from '../../../hooks/port.js'; +import { SubagentRunRegistryService } from '../../../subagent/registry.js'; + +export const spawnAgentTool: ToolDefinition = { + name: 'spawn_agent', + concurrencySafe: true, + description: + 'Start a subagent and return its session id immediately. The subagent runs in the same working directory ' + + 'and shares your file system: keep the delegated write set disjoint from your own. Its result is appended ' + + 'to this conversation when it finishes.', + parameters: z.object({ + agentName: z.string().min(1).describe('short nickname for the subagent; used for identification and display'), + prompt: z.string().min(1).describe('task description for the subagent'), + model: z.string().optional().describe('set this only when the user asks for a specific model; omit it otherwise and the subagent inherits the parent turn model'), + systemPrompt: z.string().optional().describe('replaces the middle section of the subagent system prompt; the environment block and system notes are kept'), + }), + execute: (args, ctx) => + Effect.gen(function* () { + const hooks = yield* HookService; + const registry = yield* SubagentRunRegistryService; + const { agentName, prompt, model, systemPrompt } = args as { + agentName: string; prompt: string; model?: string; systemPrompt?: string; + }; + const projectPath = ctx?.projectPath || process.cwd(); + + if (!ctx?.activeProfile || !ctx?.sessionId) { + return yield* Effect.fail( + new AgentError('CONFIG_MISSING', 'spawn_agent requires the parent session id and profile') + ); + } + + const requestedModel = model?.trim(); + const effectiveModel = requestedModel && findModel(requestedModel) ? requestedModel : ctx.model; + + const decision = yield* hooks.emitDecision('agent.subagent.spawn.before', { + agentName, prompt, parentSessionId: ctx.sessionId, projectPath, + }); + if (decision && decision.decision === 'deny') { + return yield* Effect.fail( + new AgentError('TOOL_NOT_ALLOWED', `Subagent spawn denied: ${decision.reason ?? 'no reason'}`) + ); + } + + const handle = yield* registry.spawn({ + prompt, agentName, model: effectiveModel, systemPrompt, + parentSessionId: ctx.sessionId, + parentCwd: projectPath, + parentProfile: ctx.activeProfile, + }); + + yield* hooks.emit('agent.subagent.spawn.after', { + childSessionId: handle.sessionId, agentName: handle.agentName, projectPath, + }); + + return `spawned ${handle.agentName} (${handle.sessionId})`; + }), +}; diff --git a/packages/codingcode/src/tools/domains/subagent/wait.ts b/packages/codingcode/src/tools/domains/subagent/wait.ts new file mode 100644 index 0000000..5560fb2 --- /dev/null +++ b/packages/codingcode/src/tools/domains/subagent/wait.ts @@ -0,0 +1,30 @@ +import { z } from 'zod'; +import { Effect } from 'effect'; +import type { ToolDefinition } from '../../types.js'; +import { + SubagentRunRegistryService, + SUBAGENT_WAIT_DEFAULT_MS, + SUBAGENT_WAIT_MIN_MS, + SUBAGENT_WAIT_MAX_MS, +} from '../../../subagent/registry.js'; + +export const waitAgentTool: ToolDefinition = { + name: 'wait_agent', + concurrencySafe: true, + description: + 'Wait until a subagent reaches a terminal state. Returns completed, failed or timeout. ' + + 'The subagent final output is appended to this conversation automatically — never use this tool to fetch text, ' + + 'and only wait when the result blocks your next step. Do not poll with short timeouts.', + parameters: z.object({ + sessionId: z.string().min(1).describe('subagent session id returned by spawn_agent'), + timeoutMs: z.number().int().positive().optional().describe('upper bound for this wait, clamped to [10000, 3600000]; default 30000'), + }), + execute: (args, _ctx) => + Effect.gen(function* () { + const registry = yield* SubagentRunRegistryService; + const { sessionId, timeoutMs } = args as { sessionId: string; timeoutMs?: number }; + const clamped = Math.min(Math.max(timeoutMs ?? SUBAGENT_WAIT_DEFAULT_MS, SUBAGENT_WAIT_MIN_MS), SUBAGENT_WAIT_MAX_MS); + const outcome = yield* registry.wait(sessionId, clamped); + return outcome; + }), +}; diff --git a/packages/codingcode/test/agent/abort.test.ts b/packages/codingcode/test/agent/abort.test.ts index 93a78f1..c6ab72d 100644 --- a/packages/codingcode/test/agent/abort.test.ts +++ b/packages/codingcode/test/agent/abort.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect, vi } from 'vitest'; import { makeState, runAgentTurn, textDeltas } from '../helpers/agent-harness.js'; -import type { FrameBody, Transition } from '../../src/contracts/frame.js'; +import type { EndTransition, FrameBody } from '../../src/contracts/frame.js'; vi.mock('../../src/infra/config.js', () => ({ loadConfig: () => ({ @@ -14,8 +14,6 @@ vi.mock('../../src/infra/config.js', () => ({ const state = makeState({ sessionId: 'abort-sid', cwd: '/tmp', title: 'abort' }); -type EndTransition = Extract; - function endsOf(events: readonly FrameBody[]): EndTransition[] { const out: EndTransition[] = []; for (const b of events) { diff --git a/packages/codingcode/test/approval/pipeline.test.ts b/packages/codingcode/test/approval/pipeline.test.ts index f16a09c..892adfe 100644 --- a/packages/codingcode/test/approval/pipeline.test.ts +++ b/packages/codingcode/test/approval/pipeline.test.ts @@ -30,8 +30,10 @@ const WaitTestLayer = Layer.succeed(ApprovalWaitService, mockApprovalWaitService const SinkTestLayer = Layer.succeed(EventSinkService, mockEventSink as any); const TestLayer = Layer.mergeAll(HookTestLayer, WaitTestLayer, SinkTestLayer); -function runWithLayer(eff: Effect.Effect): Promise { - return Effect.runPromise(eff.pipe(Effect.provide(TestLayer))); +type PipelineEnv = HookService | ApprovalWaitService | EventSinkService; + +function runWithLayer(eff: Effect.Effect): Promise { + return Effect.runPromise(Effect.provide(eff, TestLayer)); } describe('Approval Pipeline — PermissionMode auto-allow (merged from ReadonlyWhitelist + acceptEdits)', () => { diff --git a/packages/codingcode/test/helpers/agent-harness.ts b/packages/codingcode/test/helpers/agent-harness.ts index cdc537d..ae5c4df 100644 --- a/packages/codingcode/test/helpers/agent-harness.ts +++ b/packages/codingcode/test/helpers/agent-harness.ts @@ -14,6 +14,7 @@ import { RulesService } from '../../src/rules/port.js'; import { SessionService } from '../../src/session/port.js'; import { EventSinkService } from '../../src/sink/port.js'; import { EventSinkLayer } from '../../src/sink/sink.js'; +import { MailboxLayer } from '../../src/session/mailbox.js'; import { SkillService } from '../../src/skills/port.js'; import { SubagentRunnerService } from '../../src/subagent/port.js'; import { TodoService } from '../../src/todo/port.js'; @@ -119,7 +120,7 @@ export function hasCompress(events: readonly FrameBody[]): boolean { export interface HarnessMocks { llm: { - completeStream: (params: any, signal?: AbortSignal) => AsyncIterable; + completeStream: (params: any, model: string, signal?: AbortSignal) => AsyncIterable; modelInfo: { maxTokens: number }; }; state?: Partial; @@ -287,6 +288,7 @@ export function makeAgentLayer(mocks: HarnessMocks): Layer.Layer { Layer.succeed(SkillService, skills as any), ContextMockLayer, EventSinkLayer, + MailboxLayer, Layer.succeed(MemoryService, memory as any), Layer.succeed(LLMService, { complete: () => Effect.fail(new Error('complete not implemented in harness')), diff --git a/packages/codingcode/test/plan/gate-pipeline.test.ts b/packages/codingcode/test/plan/gate-pipeline.test.ts index d2107b6..8adfc66 100644 --- a/packages/codingcode/test/plan/gate-pipeline.test.ts +++ b/packages/codingcode/test/plan/gate-pipeline.test.ts @@ -109,10 +109,10 @@ describe('plan profile permission mode (Layer 2)', () => { expect(capturedApproval).toBeNull(); }); - it('plan profile + dispatch_agent: denied by plan mode', async () => { + it('plan profile + spawn_agent: denied by plan mode', async () => { const decision: any = await runPipelineWithMock({ - tool: 'dispatch_agent', - input: { agent: 'build', prompt: 'do something' }, + tool: 'spawn_agent', + input: { agentName: 'build', prompt: 'do something' }, permissionMode: 'ask', sessionId: 's4', profile: 'plan', diff --git a/packages/codingcode/test/server/index.test.ts b/packages/codingcode/test/server/index.test.ts index 7dc1234..63ee532 100644 --- a/packages/codingcode/test/server/index.test.ts +++ b/packages/codingcode/test/server/index.test.ts @@ -2,7 +2,7 @@ import { describe, it, expect, vi } from 'vitest'; import { Effect, Layer, ManagedRuntime } from 'effect'; import { createServer } from '../../src/server/index.js'; import { SessionService } from '../../src/session/port.js'; -import { LLMFactoryService } from '../../src/llm/port.js'; +import { LLMService } from '../../src/llm/port.js'; import { ApprovalService } from '../../src/approval/port.js'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; import { HookService } from '../../src/hooks/port.js'; @@ -37,8 +37,9 @@ const MockSessionLayer = Layer.succeed(SessionService, { }), } as any); -const MockLLMFactoryLayer = Layer.succeed(LLMFactoryService, { - getLLMClient: () => Effect.succeed(null), +const MockLLMLayer = Layer.succeed(LLMService, { + complete: () => Effect.fail(new Error('LLM not available in test')), + completeStream: () => (async function* () {})(), } as any); const MockApprovalLayer = ApprovalLayer.pipe( @@ -103,7 +104,7 @@ const MockCheckpointLayer = Layer.succeed(CheckpointService, { const TestLayer = Layer.mergeAll( MockSessionLayer, - MockLLMFactoryLayer, + MockLLMLayer, MockApprovalLayer, HookLayer, EventSinkLayer, diff --git a/packages/codingcode/test/server/plan-file-route.test.ts b/packages/codingcode/test/server/plan-file-route.test.ts index 608c969..a1ac341 100644 --- a/packages/codingcode/test/server/plan-file-route.test.ts +++ b/packages/codingcode/test/server/plan-file-route.test.ts @@ -8,7 +8,7 @@ import { join, resolve } from 'path'; import { Hono } from 'hono'; import { registerSessionsRoutes } from '../../src/server/routes/sessions.js'; import { SessionService } from '../../src/session/port.js'; -import { LLMFactoryService } from '../../src/llm/port.js'; +import { LLMService } from '../../src/llm/port.js'; import { ApprovalService } from '../../src/approval/port.js'; import { ApprovalWaitService } from '../../src/approval/wait-port.js'; import { HookService } from '../../src/hooks/port.js'; @@ -56,44 +56,9 @@ const MockSessionLayer = Layer.succeed(SessionService, { }), } as any); -const MockLLMFactoryLayer = Layer.succeed(LLMFactoryService, { - findModel: () => - Effect.succeed({ - id: 'deepseek-chat', - model: 'deepseek-chat', - activeProfile: 'build', - permissionMode: 'ask', - provider: 'deepseek', - driver: 'openai', - api_key_env: 'DEEPSEEK_API_KEY', - base_url: 'https://api.deepseek.com', - }), - createClient: () => - Effect.succeed({ - modelInfo: { - provider: 'deepseek', - model: 'deepseek-chat', - activeProfile: 'build', - permissionMode: 'ask', - maxTokens: 64000, - supportsToolCalling: true, - supportsStreaming: true, - }, - }), - getLLMClient: () => Effect.succeed(null), - listModels: () => Effect.succeed([]), - getActiveEntry: () => - Effect.succeed({ - id: 'deepseek-chat', - model: 'deepseek-chat', - activeProfile: 'build', - permissionMode: 'ask', - provider: 'deepseek', - driver: 'openai', - api_key_env: 'DEEPSEEK_API_KEY', - base_url: 'https://api.deepseek.com', - }), - switchModel: () => Effect.fail(new Error('no models')), +const MockLLMLayer = Layer.succeed(LLMService, { + complete: () => Effect.fail(new Error('LLM not available in test')), + completeStream: () => (async function* () {})(), } as any); const MockApprovalLayer = ApprovalLayer.pipe( @@ -158,7 +123,7 @@ const MockCheckpointLayer = Layer.succeed(CheckpointService, { const TestLayer = Layer.mergeAll( MockSessionLayer, - MockLLMFactoryLayer, + MockLLMLayer, MockApprovalLayer, HookLayer, EventSinkLayer, diff --git a/packages/codingcode/test/session/filter-ui.test.ts b/packages/codingcode/test/session/filter-ui.test.ts index d39c6a4..4c19108 100644 --- a/packages/codingcode/test/session/filter-ui.test.ts +++ b/packages/codingcode/test/session/filter-ui.test.ts @@ -9,6 +9,8 @@ function makeBaseEvents(extra: SessionEvent[] = []): SessionEvent[] { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -126,6 +128,8 @@ describe('sessionEventsToTurns with summary', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, diff --git a/packages/codingcode/test/session/mailbox.test.ts b/packages/codingcode/test/session/mailbox.test.ts new file mode 100644 index 0000000..f17b7e5 --- /dev/null +++ b/packages/codingcode/test/session/mailbox.test.ts @@ -0,0 +1,66 @@ +import { describe, it, expect } from 'vitest'; +import { Effect } from 'effect'; +import { MailboxLayer, MailboxService } from '../../src/session/mailbox.js'; + +const item = (sessionId: string, content = 'x') => ({ + type: 'subagent_result' as const, + sessionId, + agentName: 'build', + content, +}); + +const run = (eff: Effect.Effect): Promise => + Effect.runPromise( + eff.pipe(Effect.provide(MailboxLayer)) as unknown as Effect.Effect + ); + +describe('session mailbox', () => { + it('offer / drain:取走当前全部,再 drain 为空', async () => { + const result = await run( + Effect.gen(function* () { + const mb = yield* MailboxService; + yield* mb.offer('s1', item('child-1', 'a')); + yield* mb.offer('s1', item('child-2', 'b')); + return { first: yield* mb.drain('s1'), second: yield* mb.drain('s1') }; + }) + ); + expect(result.first.map((i) => i.content)).toEqual(['a', 'b']); + expect(result.second).toEqual([]); + }); + + it('按收件人分区:两个会话互不可见', async () => { + const result = await run( + Effect.gen(function* () { + const mb = yield* MailboxService; + yield* mb.offer('s1', item('child-1')); + return { s2: yield* mb.drain('s2'), s1: yield* mb.drain('s1') }; + }) + ); + expect(result.s2).toEqual([]); + expect(result.s1).toHaveLength(1); + }); + + it('dispose 只清该会话,其它会话不受影响', async () => { + const result = await run( + Effect.gen(function* () { + const mb = yield* MailboxService; + yield* mb.offer('s1', item('child-1')); + yield* mb.offer('s2', item('child-2')); + yield* mb.dispose('s1'); + return { s1: yield* mb.drain('s1'), s2: yield* mb.drain('s2') }; + }) + ); + expect(result.s1).toEqual([]); + expect(result.s2).toHaveLength(1); + }); + + it('drain 一个从未投递过的会话返回空数组而不报错', async () => { + const result = await run( + Effect.gen(function* () { + const mb = yield* MailboxService; + return yield* mb.drain('never-used'); + }) + ); + expect(result).toEqual([]); + }); +}); diff --git a/packages/codingcode/test/session/store-diff-rebuild.test.ts b/packages/codingcode/test/session/store-diff-rebuild.test.ts index 3b7ced7..bae6bac 100644 --- a/packages/codingcode/test/session/store-diff-rebuild.test.ts +++ b/packages/codingcode/test/session/store-diff-rebuild.test.ts @@ -10,6 +10,8 @@ describe('sessionEventsToTurns', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -60,6 +62,8 @@ describe('sessionEventsToTurns', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -108,6 +112,8 @@ describe('sessionEventsToTurns', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, diff --git a/packages/codingcode/test/session/view-assembly.test.ts b/packages/codingcode/test/session/view-assembly.test.ts index d970b85..a9ea6b0 100644 --- a/packages/codingcode/test/session/view-assembly.test.ts +++ b/packages/codingcode/test/session/view-assembly.test.ts @@ -14,6 +14,8 @@ function makeEvents(extra: SessionEvent[] = []): SessionEvent[] { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -108,6 +110,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -135,6 +139,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -184,6 +190,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -238,6 +246,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -282,6 +292,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, @@ -325,6 +337,8 @@ describe('buildContextMessages', () => { sessionId: 's1', cwd: '/tmp', createdAt: new Date().toISOString(), + model: 'deepseek-chat', + title: '', activeProfile: 'build', permissionMode: 'ask', }, diff --git a/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts b/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts deleted file mode 100644 index 3081a37..0000000 --- a/packages/codingcode/test/subagent/dispatch-end-to-end.test.ts +++ /dev/null @@ -1,253 +0,0 @@ -import { describe, it, expect, beforeEach, afterEach } from 'vitest'; -import { Effect, Layer } from 'effect'; -import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync } from 'fs'; -import { join } from 'path'; -import { tmpdir } from 'os'; -import { useTempHome } from '../helpers/temp-home.js'; -import { AgentLayer } from '../../src/agent/agent.js'; -import { ToolEnvLayer } from '../../src/agent/tool-env.js'; -import { AgentService } from '../../src/agent/port.js'; -import { ApprovalService } from '../../src/approval/port.js'; -import { CheckpointService } from '../../src/checkpoint/port.js'; -import { ContextService } from '../../src/context/port.js'; -import { EventSinkLayer } from '../../src/sink/sink.js'; -import { LLMService } from '../../src/llm/port.js'; -import { MemoryService } from '../../src/memory/port.js'; -import { RulesService } from '../../src/rules/port.js'; -import { SkillService } from '../../src/skills/port.js'; -import { ToolExecutorService } from '../../src/tools/port.js'; -import { SessionLayer } from '../../src/session/session.js'; -import { SessionService } from '../../src/session/port.js'; -import { HookService } from '../../src/hooks/port.js'; -import { McpService } from '../../src/mcp/port.js'; -import { SubagentRunnerService } from '../../src/subagent/port.js'; -import { TodoService } from '../../src/todo/port.js'; -import { readHistory } from '../../src/session/file-ops.js'; -import { encodeProjectPath, normalizePath } from '../../src/core/path.js'; -import { computePaths } from '../../src/session/paths.js'; -import { transcriptPathFor } from '../../src/context/context.js'; -import { projectBaseDir } from '../helpers/project-base.js'; -import type { SessionRef } from '../../src/contracts/session.js'; -import type { Message } from '../../src/contracts/types.js'; -import type { FrameBody } from '../../src/contracts/frame.js'; - -const ANSWER = 'subagent final answer'; - -const LLMMock = Layer.succeed(LLMService, { - complete: () => Effect.succeed({ content: ANSWER }), - completeStream: () => - (async function* () { - yield { type: 'text' as const, text: ANSWER }; - yield { type: 'end' as const }; - })(), -} as any); - -/** Read events back into the message list an LLM would see (like context.assemblePayload). */ -function readMessages(transcriptPath: string): Message[] { - return readHistory(transcriptPath).flatMap((e) => { - if (e.type === 'user') return [{ role: 'user', content: e.content }] as Message[]; - if (e.type === 'assistant') - return [{ role: 'assistant', content: e.content, tool_calls: e.toolCalls }] as Message[]; - if (e.type === 'tool_result') - return [ - { - role: 'tool', - content: e.output ?? '', - tool_call_id: e.toolCallId, - tool_name: e.toolName, - } as Message, - ]; - return []; - }); -} - -/** - * Self-contained runtime used by the end-to-end tests: real AgentLayer + - * real file-backed SessionLayer, everything else mocked. This mirrors how - * the app is wired in layer.ts while keeping each dependency explicit. - */ -// agent 依赖的全部宽服务 stub。 -const McpMock = Layer.succeed(McpService, { - syncConnections: () => Effect.void, - listProjectMcpTools: () => Effect.succeed([]), -} as any); - -const HookMock = Layer.succeed(HookService, { - emit: () => Effect.succeed(undefined), - emitDecision: () => Effect.succeed(null), - reloadUserHooks: () => Effect.succeed(undefined), -} as any); - -const TodoMock = Layer.succeed(TodoService, { read: () => [], write: () => {}, reset: () => {} } as any); - -const SubagentMock = Layer.succeed(SubagentRunnerService, {} as any); - -const AgentDeps = Layer.mergeAll( - SessionLayer, - Layer.succeed(ToolExecutorService, { - executeBatch: () => Effect.succeed([]), - prepare: () => Effect.succeed({ tools: [], lookup: () => undefined }), - } as any), - Layer.succeed(CheckpointService, { - snapshotBaseline: () => Effect.void, - snapshotFinal: () => Effect.void, - } as any), - Layer.succeed(ApprovalService, { - evaluate: () => Effect.succeed({ type: 'allow' }), - } as any), - Layer.succeed(SkillService, { - extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), - } as any), - Layer.succeed(ContextService, { - getHistory: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), - absorb: () => Effect.void, - compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), - dispose: () => Effect.void, - } as any), - EventSinkLayer, - Layer.succeed(MemoryService, { - loadMemoryForPrompt: () => Effect.succeed(''), - flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), - } as any), - LLMMock, - Layer.succeed(RulesService, { - getAllRules: () => Effect.succeed(''), - evictProjectRules: () => Effect.void, - } as any), - HookMock, - McpMock, - TodoMock, - SubagentMock, - // ToolEnvPort 在 getToolEnv 运行时从外层 Runtime 解析具体服务(见 Runtime 定义) - ToolEnvLayer -); - -// Real AgentService built on the real SessionService + stubbed wide services. -const AgentWired = AgentLayer.pipe(Layer.provide(AgentDeps as any)); - -// Runtime exposed to the tests: real AgentService + SessionService, plus the -// full services the dispatch_agent tool's execute pulls from the environment. -const Runtime = Layer.mergeAll( - AgentWired, - SessionLayer, - HookMock, - McpMock, - SubagentMock, - TodoMock -); - -function run(eff: Effect.Effect): Promise { - return Effect.runPromise(eff.pipe(Effect.provide(Runtime as any)) as any); -} - -/** Consume a frame stream to completion (mirrors what dispatch.ts does). */ -function drainStream(stream: AsyncGenerator): Effect.Effect { - return Effect.async((resume) => { - (async () => { - let content = ''; - try { - for await (const body of stream) { - if (body.family === 'event' && body.event.type === 'text_delta') { - content += body.event.text; - } else if ( - body.family === 'transition' && - body.transition.to === 'end' && - body.transition.reason === 'error' - ) { - resume(Effect.fail(new Error(`subagent failed: ${body.transition.error.message}`))); - return; - } - } - resume(Effect.succeed(content)); - } catch (e) { - resume(Effect.fail(e instanceof Error ? e : new Error(String(e)))); - } - })(); - }); -} - -describe('subagent run end-to-end (session transcript is read by the agent loop)', () => { - useTempHome('codingcode-test-e2e-home-'); - - let projectBase: string; - let cwd: string; - - beforeEach(() => { - projectBase = projectBaseDir(); - mkdirSync(projectBase, { recursive: true }); - cwd = mkdtempSync(join(tmpdir(), 'codingcode-test-cwd-')); - }); - - afterEach(() => { - if (existsSync(cwd)) rmSync(cwd, { recursive: true, force: true }); - }); - - it('runSubagent drives the agent loop and persists the transcript it reads', async () => { - const result = await run( - Effect.gen(function* () { - const agent = yield* AgentService; - const session = yield* SessionService; - const { stream, sessionId } = yield* agent.runTurn('analyze this code', { - cwd, - model: 'test-model', - activeProfile: 'build', - permissionMode: 'ask', - }); - const content = yield* drainStream(stream); - const state = yield* session.load(normalizePath(cwd), sessionId); - return { content, sessionId, transcriptPath: computePaths(state.cwd, state.sessionId, state.parentSessionId).transcriptPath }; - }) - ); - - expect(typeof result.sessionId).toBe('string'); - expect(result.content.length).toBeGreaterThan(0); - expect(existsSync(result.transcriptPath)).toBe(true); - - const events = readHistory(result.transcriptPath); - - // First event: session_meta (written by session.create in the runner path) - expect(events[0]!.type).toBe('session_meta'); - - // The user prompt recorded before the agentLoop started. If agentLoop - // read the wrong path, this event is invisible to the LLM and the - // assistant response never lands. - const userEv = events.find((e) => e.type === 'user'); - expect(userEv).toBeDefined(); - if (userEv && userEv.type === 'user') { - expect(userEv.content).toBe('analyze this code'); - } - - // The LLM's reply lands on disk — proof that agentLoop read the jsonl. - const assistantEv = events.find((e) => e.type === 'assistant'); - expect(assistantEv).toBeDefined(); - }, 30_000); - - it('child session created under a parent does NOT produce a flat /.jsonl (old bug regression)', async () => { - const result = await run( - Effect.gen(function* () { - const session = yield* SessionService; - const parent = yield* session.create(cwd, { - model: 'parent-model', - activeProfile: 'build', - permissionMode: 'ask', - }); - const child = yield* session.create( - cwd, - { model: 'child-model', activeProfile: 'build', permissionMode: 'ask' }, - { parentSessionId: parent.sessionId, agentName: 'build' } - ); - return { parentId: parent.sessionId, childId: child.sessionId }; - }) - ); - - const sessionsRoot = join(projectBase, encodeProjectPath(normalizePath(cwd)), 'sessions'); - const subagentDir = join(sessionsRoot, result.parentId, 'subagents'); - const nestedFiles = readdirSync(subagentDir).filter((f) => f.endsWith('.jsonl')); - expect(nestedFiles).toContain(`${result.childId}.jsonl`); - - // The wrong-path location MUST NOT contain the child's jsonl. If it did, - // some code constructed the path without parentSessionId. - const flatChildPath = join(sessionsRoot, `${result.childId}.jsonl`); - expect(existsSync(flatChildPath)).toBe(false); - }, 30_000); -}); diff --git a/packages/codingcode/test/subagent/dispatch-production-path.test.ts b/packages/codingcode/test/subagent/dispatch-production-path.test.ts deleted file mode 100644 index 1e95b6a..0000000 --- a/packages/codingcode/test/subagent/dispatch-production-path.test.ts +++ /dev/null @@ -1,249 +0,0 @@ -import { describe, it, expect, beforeEach, afterEach } from 'vitest'; -import { Effect, Layer } from 'effect'; -import { existsSync, mkdirSync, mkdtempSync, readdirSync, rmSync } from 'fs'; -import { join } from 'path'; -import { tmpdir } from 'os'; -import { useTempHome } from '../helpers/temp-home.js'; -import { AgentLayer } from '../../src/agent/agent.js'; -import { ToolEnvLayer } from '../../src/agent/tool-env.js'; -import { AgentService } from '../../src/agent/port.js'; -import { SubagentRunnerLayer } from '../../src/subagent/subagent.js'; -import { ApprovalService } from '../../src/approval/port.js'; -import { CheckpointService } from '../../src/checkpoint/port.js'; -import { ContextService } from '../../src/context/port.js'; -import { EventSinkLayer } from '../../src/sink/sink.js'; -import { LLMService } from '../../src/llm/port.js'; -import { MemoryService } from '../../src/memory/port.js'; -import { RulesService } from '../../src/rules/port.js'; -import { SkillService } from '../../src/skills/port.js'; -import { ToolExecutorLayer } from '../../src/tools/tools.js'; -import { SessionLayer } from '../../src/session/session.js'; -import { HookService } from '../../src/hooks/port.js'; -import { McpService } from '../../src/mcp/port.js'; -import { TodoService } from '../../src/todo/port.js'; -import { readHistory } from '../../src/session/file-ops.js'; -import { encodeProjectPath, normalizePath } from '../../src/core/path.js'; -import { projectBaseDir } from '../helpers/project-base.js'; -import { transcriptPathFor } from '../../src/context/context.js'; -import type { SessionMetaEvent, SessionRef, ToolResultEvent } from '../../src/contracts/session.js'; -import type { Message } from '../../src/contracts/types.js'; -import type { LLMStreamPart } from '../../src/contracts/provider.js'; -import type { FrameBody, ToolOutcome } from '../../src/contracts/frame.js'; - -/** - * Scripted LLM: the parent turn asks for dispatch_agent, the child turn answers, - * then the parent turn wraps up. Ordering is fixed because the delegated turn - * runs synchronously inside the tool execution. - */ -const script: LLMStreamPart[][] = []; -let call = 0; - -const LLMMock = Layer.succeed(LLMService, { - complete: () => Effect.succeed({ content: '' }), - completeStream: () => - (async function* () { - const parts = script[call++] ?? [{ type: 'end' as const }]; - for (const p of parts) yield p; - })(), -} as any); - -function readMessages(transcriptPath: string): Message[] { - return readHistory(transcriptPath).flatMap((e) => { - if (e.type === 'user') return [{ role: 'user', content: e.content }] as Message[]; - if (e.type === 'assistant') - return [{ role: 'assistant', content: e.content, tool_calls: e.toolCalls }] as Message[]; - if (e.type === 'tool_result') - return [ - { - role: 'tool', - content: e.output ?? '', - tool_call_id: e.toolCallId, - tool_name: e.toolName, - } as Message, - ]; - return []; - }); -} - -interface Draining { - text: string; - toolResults: Array<{ id: string; name: string; outcome: ToolOutcome }>; - endReason: string; -} - -function drain(stream: AsyncGenerator): Effect.Effect { - return Effect.async((resume) => { - (async () => { - const out: Draining = { text: '', toolResults: [], endReason: 'unknown' }; - try { - for await (const body of stream) { - if (body.family === 'event') { - if (body.event.type === 'text_delta') out.text += body.event.text; - if (body.event.type === 'tool_result') { - out.toolResults.push({ - id: body.event.id, - name: body.event.name, - outcome: body.event.outcome, - }); - } - } else if (body.family === 'transition' && body.transition.to === 'end') { - out.endReason = body.transition.reason; - } - } - resume(Effect.succeed(out)); - } catch (e) { - resume(Effect.fail(e instanceof Error ? e : new Error(String(e)))); - } - })(); - }); -} - -const McpMock = Layer.succeed(McpService, { - syncConnections: () => Effect.void, - listProjectMcpTools: () => Effect.succeed([]), -} as any); - -const HookMock = Layer.succeed(HookService, { - emit: () => Effect.succeed(undefined), - emitDecision: () => Effect.succeed(null), - reloadUserHooks: () => Effect.succeed(undefined), -} as any); - -const TodoMock = Layer.succeed(TodoService, { read: () => [], write: () => {}, reset: () => {} } as any); - -// 真实的工具执行层:dispatch_agent 的参数解析与 ctx 装配都走生产代码 -const ToolExecutorWithDeps = ToolExecutorLayer.pipe(Layer.provide(HookMock)); - -const AgentDeps = Layer.mergeAll( - SessionLayer, - ToolExecutorWithDeps, - Layer.succeed(CheckpointService, { - snapshotBaseline: () => Effect.void, - snapshotFinal: () => Effect.void, - } as any), - Layer.succeed(ApprovalService, { - evaluate: () => Effect.succeed({ type: 'allow' }), - } as any), - Layer.succeed(SkillService, { - extractSkill: (_cwd: string, query: string) => Effect.succeed([undefined, query]), - } as any), - Layer.succeed(ContextService, { - getHistory: (ref: SessionRef) => Effect.sync(() => readMessages(transcriptPathFor(ref))), - absorb: () => Effect.void, - compact: () => Effect.succeed({ didCompress: false, released: 0, promptEstimate: 0 }), - dispose: () => Effect.void, - } as any), - EventSinkLayer, - Layer.succeed(MemoryService, { - loadMemoryForPrompt: () => Effect.succeed(''), - flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), - } as any), - LLMMock, - Layer.succeed(RulesService, { - getAllRules: () => Effect.succeed(''), - evictProjectRules: () => Effect.void, - } as any), - HookMock, - McpMock, - TodoMock, - ToolEnvLayer -); - -const AgentWired = AgentLayer.pipe(Layer.provide(AgentDeps as any)); -const SubagentWired = SubagentRunnerLayer.pipe(Layer.provide(AgentWired)); - -const Runtime = Layer.mergeAll( - AgentWired, - SubagentWired, - SessionLayer, - ToolExecutorWithDeps, - HookMock, - McpMock, - TodoMock -); - -function run(eff: Effect.Effect): Promise { - return Effect.runPromise(eff.pipe(Effect.provide(Runtime as any)) as any); -} - -describe('dispatch_agent on the production path (parent turn -> tool -> subagent turn)', () => { - useTempHome('codingcode-test-prod-path-home-'); - - let projectBase: string; - let cwd: string; - - beforeEach(() => { - projectBase = projectBaseDir(); - mkdirSync(projectBase, { recursive: true }); - cwd = mkdtempSync(join(tmpdir(), 'codingcode-test-prod-path-cwd-')); - script.length = 0; - call = 0; - }); - - afterEach(() => { - if (existsSync(cwd)) rmSync(cwd, { recursive: true, force: true }); - }); - - it('delegates without CONFIG_MISSING and lands the child transcript under the parent', async () => { - script.push( - [ - { - type: 'tool_call', - id: 'call-1', - name: 'dispatch_agent', - arguments: { agentName: 'reviewer', prompt: 'inspect the module' }, - }, - { type: 'end' }, - ], - // 子代理的那一回合 - [{ type: 'text', text: 'child answer' }, { type: 'end' }], - // 父代理收尾 - [{ type: 'text', text: 'parent done' }, { type: 'end' }] - ); - - const { drained, parentId } = await run( - Effect.gen(function* () { - const agent = yield* AgentService; - const { stream, sessionId } = yield* agent.runTurn('please delegate', { - cwd, - model: 'test-model', - activeProfile: 'build', - permissionMode: 'ask', - }); - const drained = yield* drain(stream); - return { drained, parentId: sessionId }; - }) - ); - - expect(drained.endReason).toBe('done'); - - const delegated = drained.toolResults.find((r) => r.name === 'dispatch_agent'); - expect(delegated).toBeDefined(); - // 阶段一之前这里必然是 [Error: CONFIG_MISSING],父会话的 activeProfile 到不了 dispatch_agent - expect(delegated!.outcome.status).toBe('ok'); - if (delegated!.outcome.status === 'ok') { - expect(delegated!.outcome.output).toContain('child answer'); - } - - const sessionsRoot = join(projectBase, encodeProjectPath(normalizePath(cwd)), 'sessions'); - const subagentDir = join(sessionsRoot, parentId, 'subagents'); - expect(existsSync(subagentDir)).toBe(true); - const childFiles = readdirSync(subagentDir).filter((f) => f.endsWith('.jsonl')); - expect(childFiles).toHaveLength(1); - - const childEvents = readHistory(join(subagentDir, childFiles[0]!)); - const meta = childEvents[0] as SessionMetaEvent; - expect(meta.type).toBe('session_meta'); - expect(meta.agentName).toBe('reviewer'); - expect(meta.parentSessionId).toBe(parentId); - expect(childEvents.some((e) => e.type === 'assistant')).toBe(true); - - // 父转录里留下了这次委派的调用与结果 - const parentEvents = readHistory(join(sessionsRoot, `${parentId}.jsonl`)); - const toolResult = parentEvents.find( - (e) => e.type === 'tool_result' && e.toolName === 'dispatch_agent' - ) as ToolResultEvent | undefined; - expect(toolResult).toBeDefined(); - expect(toolResult!.output).toContain('child answer'); - }, 30_000); -}); diff --git a/packages/codingcode/test/subagent/dispatch.test.ts b/packages/codingcode/test/subagent/dispatch.test.ts deleted file mode 100644 index bc8052d..0000000 --- a/packages/codingcode/test/subagent/dispatch.test.ts +++ /dev/null @@ -1,215 +0,0 @@ -import { expect, it, describe, beforeEach, vi } from 'vitest'; -import { Effect, Layer } from 'effect'; -import { z } from 'zod'; -import { dispatchAgentTool } from '../../src/tools/domains/subagent/dispatch.js'; -import { HookService } from '../../src/hooks/port.js'; -import { SubagentRunnerService } from '../../src/subagent/port.js'; -import type { ToolExecCtx } from '../../src/contracts/tool.js'; -import type { FrameBody } from '../../src/contracts/frame.js'; - -// 只有清单内的模型能被透传,测试内固定一份最小清单 -vi.mock('../../src/infra/models.js', async (importOriginal) => ({ - ...(await importOriginal()), - findModel: (id: string) => - id === 'child-model' ? ({ id: 'child-model@test', model: 'child-model' } as any) : null, -})); - -const mockHooks = { - emit: vi.fn(() => Effect.succeed(undefined)), - emitDecision: vi.fn(() => Effect.succeed(null)), - reloadUserHooks: () => Effect.succeed(undefined), -}; - -const mockRunner = { - runSubagent: vi.fn((_input: string, _opts: Record) => - Effect.succeed({ stream: makeRunStream(), sessionId: 'child-1' }) - ), -}; - -function makeRunStream(): AsyncGenerator { - return (async function* () { - yield { family: 'event', event: { type: 'text_delta', text: 'done' } }; - yield { family: 'transition', transition: { to: 'end', reason: 'done' } }; - })(); -} - -function makeLayers() { - return Layer.mergeAll( - Layer.succeed(HookService, mockHooks as any), - Layer.succeed(SubagentRunnerService, mockRunner as any) - ); -} - -/** 工具执行上下文:activeProfile / model 由父会话透传,缺前者则派发不出来 */ -function parentCtx(overrides: Partial = {}): ToolExecCtx { - return { - projectPath: '/test', - sessionId: 'parent-1', - activeProfile: 'build', - model: 'parent-model', - ...overrides, - }; -} - -function runTool(args: unknown, ctx: ToolExecCtx): Promise { - return Effect.runPromise( - dispatchAgentTool.execute(args, ctx).pipe(Effect.provide(makeLayers())) - ); -} - -function runToolEither(args: unknown, ctx: ToolExecCtx) { - return Effect.runPromise( - Effect.either(dispatchAgentTool.execute(args, ctx).pipe(Effect.provide(makeLayers()))) - ); -} - -describe('dispatch_agent (runner-based subagent spawn)', () => { - beforeEach(() => { - vi.clearAllMocks(); - }); - - it('exposes agentName/prompt/model/systemPrompt and no dangling catalog pointer', () => { - expect(dispatchAgentTool.description).not.toContain('Available Subagents'); - const schema = z.toJSONSchema(dispatchAgentTool.parameters) as any; - expect(Object.keys(schema.properties).sort()).toEqual([ - 'agentName', - 'model', - 'prompt', - 'systemPrompt', - ]); - expect(schema.required).toEqual(['agentName', 'prompt']); - }); - - it('case 1: dispatches a subagent and returns the runner output', async () => { - const out = await runTool( - { agentName: 'build', prompt: 'go' }, - parentCtx() - ); - expect(out).toBe('done'); - }); - - it('case 2: forwards prompt, cwd, parent session id, profile and bypass to the runner', async () => { - await runTool( - { agentName: 'reviewer', prompt: 'analyze this code' }, - parentCtx() - ); - - expect(mockRunner.runSubagent).toHaveBeenCalledTimes(1); - expect(mockRunner.runSubagent).toHaveBeenCalledWith( - 'analyze this code', - expect.objectContaining({ - cwd: '/test', - parentSessionId: 'parent-1', - activeProfile: 'build', - agentName: 'reviewer', - permissionMode: 'bypass', - }) - ); - }); - - it('case 3: passes model and systemPrompt through, inherits the turn model when model is omitted or not in the catalog', async () => { - await runTool( - { agentName: 'build', prompt: 'go', model: 'child-model', systemPrompt: 'CUSTOM PROMPT' }, - parentCtx() - ); - expect(mockRunner.runSubagent).toHaveBeenLastCalledWith( - 'go', - expect.objectContaining({ model: 'child-model', systemPrompt: 'CUSTOM PROMPT' }) - ); - - mockRunner.runSubagent.mockClear(); - await runTool({ agentName: 'build', prompt: 'go' }, parentCtx()); - const omitted = mockRunner.runSubagent.mock.calls.at(-1)![1]; - expect(omitted.model).toBe('parent-model'); - expect(omitted.systemPrompt).toBeUndefined(); - - mockRunner.runSubagent.mockClear(); - await runTool({ agentName: 'build', prompt: 'go', model: 'not-in-catalog' }, parentCtx()); - expect(mockRunner.runSubagent.mock.calls.at(-1)![1].model).toBe('parent-model'); - }); - - it('case 4: accepts any non-empty agentName (no profile lookup)', async () => { - const out = await runTool( - { agentName: 'custom-name', prompt: 'go' }, - parentCtx() - ); - expect(out).toBe('done'); - expect(mockRunner.runSubagent).toHaveBeenCalledWith( - 'go', - expect.objectContaining({ agentName: 'custom-name' }) - ); - }); - - it('case 5: fails with CONFIG_MISSING when the parent profile is missing', async () => { - const outcome = await runToolEither( - { agentName: 'build', prompt: 'go' }, - { projectPath: '/test', sessionId: 'parent-1', model: 'parent-model' } - ); - expect(outcome._tag).toBe('Left'); - if (outcome._tag === 'Left') { - const err: any = outcome.left; - expect(err.code).toBe('CONFIG_MISSING'); - expect(String(err.message)).toContain('activeProfile'); - } - expect(mockRunner.runSubagent).not.toHaveBeenCalled(); - }); - - it('case 6: spawn.before deny hook blocks the dispatch', async () => { - mockHooks.emitDecision.mockReturnValueOnce( - Effect.succeed({ decision: 'deny' as const, reason: 'policy forbids it' }) as any - ); - const outcome = await runToolEither( - { agentName: 'build', prompt: 'go' }, - parentCtx() - ); - expect(outcome._tag).toBe('Left'); - if (outcome._tag === 'Left') { - const err: any = outcome.left; - expect(err.code).toBe('TOOL_NOT_ALLOWED'); - } - }); - - it('case 7: emits spawn.after and complete carrying the child session id and agentName', async () => { - await runTool( - { agentName: 'reviewer', prompt: 'go' }, - parentCtx() - ); - - expect(mockHooks.emit).toHaveBeenCalledWith( - 'agent.subagent.spawn.after', - expect.objectContaining({ childSessionId: 'child-1', agentName: 'reviewer' }) - ); - expect(mockHooks.emit).toHaveBeenCalledWith( - 'agent.subagent.complete', - expect.objectContaining({ childSessionId: 'child-1', status: 'done' }) - ); - }); - - it('case 8: a stream that ends with error fails the tool and skips complete', async () => { - mockRunner.runSubagent.mockReturnValueOnce( - Effect.succeed({ - stream: (async function* () { - yield { - family: 'transition', - transition: { to: 'end', reason: 'error', error: { message: 'boom' } }, - }; - })() as AsyncGenerator, - sessionId: 'child-2', - }) as any - ); - - const outcome = await runToolEither( - { agentName: 'build', prompt: 'go' }, - parentCtx() - ); - - expect(outcome._tag).toBe('Left'); - if (outcome._tag === 'Left') { - expect(String((outcome.left as any).message)).toContain('Subagent failed: boom'); - } - expect(mockHooks.emit).not.toHaveBeenCalledWith( - 'agent.subagent.complete', - expect.anything() - ); - }); -}); diff --git a/packages/codingcode/test/subagent/registry.test.ts b/packages/codingcode/test/subagent/registry.test.ts new file mode 100644 index 0000000..cf91689 --- /dev/null +++ b/packages/codingcode/test/subagent/registry.test.ts @@ -0,0 +1,205 @@ +import { describe, it, expect } from 'vitest'; +import { Effect, Layer, Queue } from 'effect'; +import { SubagentRunRegistryLayer, SubagentRunRegistryService } from '../../src/subagent/registry.js'; +import { SubagentRunnerService } from '../../src/subagent/port.js'; +import { MailboxLayer, MailboxService } from '../../src/session/mailbox.js'; +import { EventSinkService } from '../../src/sink/port.js'; +import type { FrameBody } from '../../src/contracts/frame.js'; + +const spawnOpts = { + prompt: 'do a thing', + agentName: 'build', + parentSessionId: 'parent-1', + parentCwd: '/tmp', + parentProfile: 'build' as const, + model: 'test-model', +}; + +/** 立刻发一条 text_delta 与终态 */ +function doneStream(text = 'child-done'): AsyncGenerator { + return (async function* () { + yield { family: 'event', event: { type: 'text_delta', text } } as FrameBody; + yield { family: 'transition', transition: { to: 'end', reason: 'done' } } as FrameBody; + })(); +} + +/** 永不结束:用来制造 wait 超时 */ +function hangingStream(): AsyncGenerator { + return (async function* () { + await new Promise((r) => setTimeout(r, 30_000)); + yield { family: 'transition', transition: { to: 'end', reason: 'done' } } as FrameBody; + })(); +} + +/** 以 error 终态收尾 */ +function failStream(): AsyncGenerator { + return (async function* () { + yield { + family: 'transition', + transition: { to: 'end', reason: 'error', error: { message: 'boom', code: 'X' } }, + } as FrameBody; + })(); +} + +function makeHarness(makeStream: () => AsyncGenerator) { + const emitted: Array<{ sessionId: string; body: FrameBody }> = []; + // 每次 spawn 给一条独立流与递增的 sessionId(配额用例会连续 spawn 多次) + let n = 0; + const runner = Layer.succeed(SubagentRunnerService, { + runSubagent: () => Effect.sync(() => ({ stream: makeStream(), sessionId: `child-${++n}` })), + } as any); + const sink = Layer.succeed(EventSinkService, { + attach: () => Effect.succeed(Effect.runSync(Queue.unbounded())), + detach: () => Effect.void, + emit: (sessionId: string, body: FrameBody) => + Effect.sync(() => { + emitted.push({ sessionId, body }); + }), + } as any); + // provideMerge:MailboxService 既注入注册表,也保留在输出里供断言使用 + // (同一次 Layer 构建 ⇒ 同一实例) + const layers = SubagentRunRegistryLayer.pipe( + Layer.provideMerge(Layer.mergeAll(runner, MailboxLayer, sink)) + ); + return { emitted, layers }; +} + +const run = (layers: Layer.Layer, eff: Effect.Effect): Promise => + Effect.runPromise( + eff.pipe(Effect.provide(layers)) as unknown as Effect.Effect + ); + +describe('subagent run registry', () => { + it('spawn 立即返回句柄,不等子代理结束', async () => { + const { layers } = makeHarness(hangingStream); + const result = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + const started = Date.now(); + const handle = yield* reg.spawn(spawnOpts); + return { handle, elapsed: Date.now() - started }; + }) + ); + expect(result.handle).toEqual({ sessionId: 'child-1', agentName: 'build' }); + expect(result.elapsed).toBeLessThan(500); + }); + + it('spawn 时向父会话投递 spawned 帧', async () => { + const { layers, emitted } = makeHarness(hangingStream); + await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + yield* reg.spawn(spawnOpts); + }) + ); + expect(emitted[0]?.sessionId).toBe('parent-1'); + expect(emitted[0]?.body).toMatchObject({ + family: 'event', + event: { type: 'subagent_event', sessionId: 'child-1', agentName: 'build', status: 'spawned' }, + }); + }); + + it('终态到达后:结果进父会话 mailbox,且 completed 帧投给父会话', async () => { + const { layers, emitted } = makeHarness(() => doneStream('hello')); + const result = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + const mb = yield* MailboxService; + yield* reg.spawn(spawnOpts); + // 让后台 drain fiber 跑完 + yield* Effect.sleep(50); + return yield* mb.drain('parent-1'); + }) + ); + expect(result).toHaveLength(1); + expect(result[0]).toMatchObject({ type: 'subagent_result', sessionId: 'child-1', agentName: 'build' }); + expect(result[0]!.content).toContain('hello'); + expect(result[0]!.content).toContain('Message Type: FINAL_ANSWER'); + + const statuses = emitted + .map((e) => (e.body.family === 'event' ? (e.body.event as any).status : undefined)) + .filter(Boolean); + expect(statuses).toEqual(['spawned', 'completed']); + }); + + it('wait 对已终态返回 completed(不阻塞)', async () => { + const { layers } = makeHarness(doneStream); + const outcome = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + yield* reg.spawn(spawnOpts); + yield* Effect.sleep(50); + return yield* reg.wait('child-1', 10_000); + }) + ); + expect(outcome).toBe('completed'); + }); + + it('error 终态:mailbox 正文含失败原因,wait 返回 failed', async () => { + const { layers } = makeHarness(failStream); + const result = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + const mb = yield* MailboxService; + yield* reg.spawn(spawnOpts); + yield* Effect.sleep(50); + const items = yield* mb.drain('parent-1'); + const outcome = yield* reg.wait('child-1', 10_000); + return { items, outcome }; + }) + ); + expect(result.outcome).toBe('failed'); + expect(result.items[0]!.content).toContain('boom'); + expect(result.items[0]!.content).toContain('did not finish'); + }); + + it('wait 超时返回 timeout,且不取消子代理', async () => { + const { layers } = makeHarness(hangingStream); + const outcome = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + yield* reg.spawn(spawnOpts); + return yield* reg.wait('child-1', 60); + }) + ); + expect(outcome).toBe('timeout'); + }); + + it('wait 未知 sessionId 报 Unknown subagent', async () => { + const { layers } = makeHarness(doneStream); + const error = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + return yield* Effect.either(reg.wait('nope', 60)); + }) + ); + expect(error._tag).toBe('Left'); + expect(String((error as any).left.message)).toMatch(/Unknown subagent/); + }); + + it('配额:超过 maxBackground 时拒绝新的 spawn', async () => { + const { layers } = makeHarness(hangingStream); + const result = await run( + layers, + Effect.gen(function* () { + const reg = yield* SubagentRunRegistryService; + const outcomes: string[] = []; + for (let i = 0; i < 6; i++) { + const r = yield* Effect.either(reg.spawn(spawnOpts)); + outcomes.push(r._tag); + } + return outcomes; + }) + ); + // 默认上限 4:前四个成功,之后失败 + expect(result.slice(0, 4)).toEqual(['Right', 'Right', 'Right', 'Right']); + expect(result.slice(4)).toEqual(['Left', 'Left']); + }); +}); diff --git a/packages/codingcode/test/subagent/runner-wiring.test.ts b/packages/codingcode/test/subagent/runner-wiring.test.ts index 80857db..b56967a 100644 --- a/packages/codingcode/test/subagent/runner-wiring.test.ts +++ b/packages/codingcode/test/subagent/runner-wiring.test.ts @@ -12,6 +12,7 @@ import { ApprovalService } from '../../src/approval/port.js'; import { CheckpointService } from '../../src/checkpoint/port.js'; import { ContextService } from '../../src/context/port.js'; import { EventSinkLayer } from '../../src/sink/sink.js'; +import { MailboxLayer } from '../../src/session/mailbox.js'; import { LLMService } from '../../src/llm/port.js'; import { AgentError } from '../../src/core/error.js'; import { MemoryService } from '../../src/memory/port.js'; @@ -139,6 +140,7 @@ const AgentDeps = Layer.mergeAll( dispose: () => Effect.void, } as any), EventSinkLayer, + MailboxLayer, Layer.succeed(MemoryService, { loadMemoryForPrompt: () => Effect.succeed(''), flushSessionToMemory: () => Effect.succeed({ written: false, bytes: 0 }), diff --git a/packages/codingcode/test/tools/catalog.test.ts b/packages/codingcode/test/tools/catalog.test.ts index b5b99fc..3488406 100644 --- a/packages/codingcode/test/tools/catalog.test.ts +++ b/packages/codingcode/test/tools/catalog.test.ts @@ -13,7 +13,8 @@ const BUILD_NAMES = [ 'fetch_url', 'web_search', 'todo_write', - 'dispatch_agent', + 'spawn_agent', + 'wait_agent', ]; describe('createToolCatalog', () => { diff --git a/packages/codingcode/test/tools/executor-concurrency.test.ts b/packages/codingcode/test/tools/executor-concurrency.test.ts index 37dcf6e..ebc2f0d 100644 --- a/packages/codingcode/test/tools/executor-concurrency.test.ts +++ b/packages/codingcode/test/tools/executor-concurrency.test.ts @@ -9,6 +9,7 @@ import { HookService } from '../../src/hooks/port.js'; import { TodoService } from '../../src/todo/port.js'; import { McpService } from '../../src/mcp/port.js'; import { SubagentRunnerService } from '../../src/subagent/port.js'; +import { SubagentRunRegistryService } from '../../src/subagent/registry.js'; import { TOOLS_BY_NAME, createToolCatalog } from '../../src/tools/catalog.js'; import type { ToolCall, TodoItem } from '../../src/contracts/types.js'; import type { ToolResult } from '../../src/contracts/tool.js'; @@ -95,11 +96,18 @@ function makeHarness(): Harness { Effect.succeed({ stream: makeSubagentStream(200), sessionId: 'child-1' }), }; + // spawn_agent 消费注册表(生产装配在 AppLayer,这里只提供最小桩) + const registry = { + spawn: () => Effect.succeed({ sessionId: 'child-1', agentName: 'build' }), + wait: () => Effect.succeed('completed'), + }; + h.layers = Layer.mergeAll( Layer.succeed(HookService, hooks as any), Layer.succeed(TodoService, todo as any), Layer.succeed(McpService, mcp as any), - Layer.succeed(SubagentRunnerService, runner as any) + Layer.succeed(SubagentRunnerService, runner as any), + Layer.succeed(SubagentRunRegistryService, registry as any) ); return h; } @@ -339,17 +347,18 @@ describe('executeBatch 保序波次调度', () => { expect(h.todoStore.get('sid-1')?.[0]?.step).toBe('step one'); }); - it('[dispatch_agent, write_file] 零重叠,且委派与写入的先后由声明顺序决定', async () => { + it('[spawn_agent, write_file] 零重叠:并发安全的 spawn 自成一波,写入独占一波,先后由声明顺序决定', async () => { const results = await runBatch( [ - tc('d1', 'dispatch_agent', { agentName: 'build', prompt: 'go' }), + tc('d1', 'spawn_agent', { agentName: 'build', prompt: 'go' }), tc('w1', 'write_file', { path: 'z.txt', content: 'z' }), ], { projectPath: dir, sessionId: 'sid-1', activeProfile: 'build', model: 'm' }, h ); - expect(okOutput(results[0])).toBe('child-done'); + // spawn 只登记后台任务后立即返回,不等待子代理 + expect(okOutput(results[0])).toBe('spawned build (child-1)'); expect(okOutput(results[1])).toContain('File written'); const d1 = intervalFor(h, 'd1'); const w1 = intervalFor(h, 'w1'); @@ -360,7 +369,7 @@ describe('executeBatch 保序波次调度', () => { const reversed = await runBatch( [ tc('w2', 'write_file', { path: 'z2.txt', content: 'z' }), - tc('d2', 'dispatch_agent', { agentName: 'build', prompt: 'go' }), + tc('d2', 'spawn_agent', { agentName: 'build', prompt: 'go' }), ], { projectPath: dir, sessionId: 'sid-1', activeProfile: 'build', model: 'm' }, h @@ -410,7 +419,7 @@ describe('executeBatch 保序波次调度', () => { expect(h.hookPoints.filter((p) => p === 'tool.execute.before')).toHaveLength(1); }); - it('11 个内置工具的 concurrencySafe 与分类表逐项一致,且无未覆盖工具', () => { + it('12 个内置工具的 concurrencySafe 与分类表逐项一致,且无未覆盖工具', () => { const expected: Record = { read_file: true, search_code: true, @@ -418,10 +427,11 @@ describe('executeBatch 保序波次调度', () => { fetch_url: true, web_search: true, todo_write: true, + spawn_agent: true, + wait_agent: true, write_file: false, edit_file: false, execute_command: false, - dispatch_agent: false, submit_plan: false, }; diff --git a/packages/desktop/src/lib/frame-reducer.ts b/packages/desktop/src/lib/frame-reducer.ts index 8cd99c0..104a583 100644 --- a/packages/desktop/src/lib/frame-reducer.ts +++ b/packages/desktop/src/lib/frame-reducer.ts @@ -116,6 +116,10 @@ export function reduceFrame(frame: Frame, state: StreamState, fx: StreamEffects) status: 'pending', }); return; + case 'subagent_event': + // 后台子代理的进展:阶段二不新增 UI 元素,仅显式消费该帧 + // (不消费也不会报错,只是被静默忽略);展示留给阶段三的状态快照。 + return; case 'tool_result': { if (e.name === 'submit_plan' && e.outcome.status !== 'ok') state.planTitle = null; if (e.outcome.status === 'denied') { diff --git a/packages/sdk/src/decode.ts b/packages/sdk/src/decode.ts index d9ba2d9..9b3e0eb 100644 --- a/packages/sdk/src/decode.ts +++ b/packages/sdk/src/decode.ts @@ -11,7 +11,7 @@ export type DecodeResult = | { readonly ok: false; readonly reason: DecodeFailureReason; readonly raw: unknown }; const TRANSITIONS = new Set(['start', 'executing', 'compress', 'end']); -const EVENTS = new Set(['text_delta', 'tool_call', 'tool_result', 'approval_request']); +const EVENTS = new Set(['text_delta', 'tool_call', 'tool_result', 'approval_request', 'subagent_event']); const END_REASONS = new Set(['done', 'error', 'maxSteps', 'aborted']); const OUTCOMES = new Set(['ok', 'error', 'denied']); @@ -95,6 +95,11 @@ export function decodeFrame(raw: unknown): DecodeResult { if (e.type === 'approval_request' && typeof e.tool !== 'string') { return { ok: false, reason: 'shape', raw }; } + if (e.type === 'subagent_event') { + if (typeof e.sessionId !== 'string' || typeof e.agentName !== 'string' || typeof e.status !== 'string') { + return { ok: false, reason: 'shape', raw }; + } + } return { ok: true, frame: raw as unknown as Frame }; } diff --git a/packages/sdk/src/protocol.ts b/packages/sdk/src/protocol.ts index 852accd..baf2693 100644 --- a/packages/sdk/src/protocol.ts +++ b/packages/sdk/src/protocol.ts @@ -52,6 +52,12 @@ export type RuntimeEvent = readonly id: string; readonly tool: string; readonly args: Readonly>; + } + | { + readonly type: 'subagent_event'; + readonly sessionId: string; + readonly agentName: string; + readonly status: 'spawned' | 'completed' | 'failed'; }; export interface Fatal {