feat(subagent): port 3c subagent engine — run_agent loop + spawn/delegate/parallel
This commit is contained in:
@@ -0,0 +1,37 @@
|
||||
/**
|
||||
* Shared config resolver for subagent execution.
|
||||
* Loads Settings + AppConfig from the store directory and resolves
|
||||
* provider, model, baseUrl, and apiKey.
|
||||
*/
|
||||
import type { Settings, AppConfig } from "@zesdex/domain";
|
||||
import { resolveSubagentProvider } from "./provider.ts";
|
||||
|
||||
/** Resolved LLM configuration for a subagent. */
|
||||
export interface ResolvedLlmConfig {
|
||||
baseUrl: string;
|
||||
apiKey: string;
|
||||
model: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve LLM config from the store directory.
|
||||
* Uses JsonSettingsRepository + JsonAppConfigRepository from persistence.
|
||||
*/
|
||||
export async function resolveConfig(): Promise<ResolvedLlmConfig> {
|
||||
const { newStore } = await import("@zesdex/domain");
|
||||
const { JsonSettingsRepository, JsonAppConfigRepository } = await import("../persistence/index.ts");
|
||||
const store = newStore();
|
||||
const settingsRepo = new JsonSettingsRepository();
|
||||
const appConfigRepo = new JsonAppConfigRepository();
|
||||
const settings: Settings = await settingsRepo.load(store.base_dir);
|
||||
const appConfig: AppConfig = await appConfigRepo.load(store.base_dir);
|
||||
const { provider, model } = resolveSubagentProvider(settings, appConfig);
|
||||
const cfg = appConfig.providers[provider] ?? ({ api_base: "" } as unknown as { api_base?: string; api_key_env?: string; default_api_key?: string });
|
||||
const apiKey =
|
||||
settings.api_keys[provider] ??
|
||||
(cfg.api_key_env ? process.env[cfg.api_key_env] : undefined) ??
|
||||
cfg.default_api_key ??
|
||||
"";
|
||||
const baseUrl = cfg.api_base || "https://opencode.ai/zen/v1";
|
||||
return { baseUrl, apiKey, model };
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
/**
|
||||
* Subagent execution context — wraps the shared state needed for a subagent.
|
||||
* Mirrors `apps/infrastructure/src/subagent/context.rs`.
|
||||
*/
|
||||
import type { ToolCtx } from "../tools/context.ts";
|
||||
import type { AccessTier } from "@zesdex/domain";
|
||||
|
||||
/**
|
||||
* Context for a single subagent execution.
|
||||
*
|
||||
* Flow: constructed by the caller (e.g. execute_primitive) with resolved
|
||||
* LLM credentials → passed to `runAgent` → used to create the provider
|
||||
* service for LLM interaction.
|
||||
*/
|
||||
export interface SubagentContext {
|
||||
/** The directive/instruction the subagent should execute. */
|
||||
directive: string;
|
||||
/** Shared tool execution context (workspaces, session, memory paths). */
|
||||
toolCtx: ToolCtx;
|
||||
/** Access tier as a string (used for logging/serialization). */
|
||||
accessTier: string;
|
||||
/** Base URL for the LLM provider API. */
|
||||
baseUrl: string;
|
||||
/** API key for the LLM provider. */
|
||||
apiKey: string;
|
||||
/** Model identifier for the LLM provider. */
|
||||
model: string;
|
||||
}
|
||||
|
||||
/** Create a new subagent context with all required fields. */
|
||||
export function newSubagentContext(
|
||||
directive: string,
|
||||
toolCtx: ToolCtx,
|
||||
accessTier: string,
|
||||
baseUrl: string,
|
||||
apiKey: string,
|
||||
model: string,
|
||||
): SubagentContext {
|
||||
return { directive, toolCtx, accessTier, baseUrl, apiKey, model };
|
||||
}
|
||||
|
||||
/** Convert a string access tier to the typed AccessTier. */
|
||||
export function parseAccessTier(s: string): AccessTier {
|
||||
switch (s.toLowerCase()) {
|
||||
case "read":
|
||||
return "Read";
|
||||
case "write":
|
||||
return "Write";
|
||||
case "full":
|
||||
return "Full";
|
||||
default:
|
||||
return "Read";
|
||||
}
|
||||
}
|
||||
@@ -1,13 +1,49 @@
|
||||
/**
|
||||
* Parallel delegation helper used by parallel_delegate tool. Ported in 3c.
|
||||
* Mirrors `apps/infrastructure/src/subagent/delegate.rs`.
|
||||
*/
|
||||
import type { ToolCtx } from "../tools/mod.ts";
|
||||
import { runAgent } from "../tools/engine.ts";
|
||||
import { resolveConfig } from "./config_resolver.ts";
|
||||
import { buildProviderService } from "./http_provider.ts";
|
||||
|
||||
/**
|
||||
* Run parallel delegation — split a task into sub-tasks that run concurrently
|
||||
* across multiple subagents, then optionally synthesize results.
|
||||
*
|
||||
* Mirrors the Rust `run_parallel_delegation` function.
|
||||
*/
|
||||
export async function runParallelDelegation(
|
||||
_task: string,
|
||||
_directives: Array<{ directive: string; access: string }>,
|
||||
_synthesize: boolean,
|
||||
_ctx: ToolCtx,
|
||||
task: string,
|
||||
directives: Array<{ directive: string; access: string }>,
|
||||
synthesize: boolean,
|
||||
ctx: ToolCtx,
|
||||
): Promise<string> {
|
||||
throw new Error("subagent engine not yet wired in this build");
|
||||
if (directives.length === 0) return "No directives to execute.";
|
||||
|
||||
const { baseUrl, apiKey, model } = await resolveConfig();
|
||||
const svc = buildProviderService(baseUrl, apiKey, model);
|
||||
|
||||
// Run all directives in parallel (bounded concurrency)
|
||||
const results = await Promise.all(
|
||||
directives.map(async (d) => {
|
||||
const access = (d.access as "read" | "write" | "full") ?? "read";
|
||||
return await runAgent(svc, d.directive, access, ctx, undefined);
|
||||
}),
|
||||
);
|
||||
|
||||
let output = `## Parallel Delegation Results\n\n**Task:** ${task}\n**Parallel agents:** ${directives.length}\n\n`;
|
||||
|
||||
results.forEach((result, i) => {
|
||||
const access = directives[i]?.access ?? "read";
|
||||
output += `---\n### Agent ${i}: [${access}]\n\n${result}\n\n`;
|
||||
});
|
||||
|
||||
// Optional synthesis
|
||||
if (synthesize && results.length > 0) {
|
||||
const combined = results.join("\n\n===\n\n");
|
||||
output += `\n## Synthesis\n\nThe following outputs were collected from ${results.length} parallel agents:\n\n${combined}`;
|
||||
}
|
||||
|
||||
return output;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
/**
|
||||
* Shared HTTP provider service builder for subagent execution.
|
||||
* Creates a `SubagentProviderService` that talks to an OpenAI-compatible
|
||||
* /chat/completions endpoint.
|
||||
*/
|
||||
import type { ChatMessage, Role, ToolCall } from "@zesdex/domain";
|
||||
import type { SubagentProviderService } from "../tools/engine.ts";
|
||||
|
||||
interface RawToolCall {
|
||||
id: string;
|
||||
type?: string;
|
||||
function?: { name?: string; arguments?: string };
|
||||
}
|
||||
|
||||
interface RawChoice {
|
||||
message?: { content?: string; role?: string; tool_calls?: RawToolCall[] };
|
||||
}
|
||||
|
||||
interface RawResponse {
|
||||
choices?: RawChoice[];
|
||||
usage?: { prompt_tokens?: number; completion_tokens?: number };
|
||||
}
|
||||
|
||||
/** Build an HTTP-based SubagentProviderService for the given LLM endpoint. */
|
||||
export function buildProviderService(
|
||||
baseUrl: string,
|
||||
apiKey: string,
|
||||
model: string,
|
||||
): SubagentProviderService {
|
||||
return {
|
||||
chat: async (messages, tools, maxTokens, temperature) => {
|
||||
const resp = await fetch(`${baseUrl}/chat/completions`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
"Content-Type": "application/json",
|
||||
...(apiKey ? { Authorization: `Bearer ${apiKey}` } : {}),
|
||||
},
|
||||
body: JSON.stringify({
|
||||
model,
|
||||
messages,
|
||||
max_tokens: maxTokens ?? 4096,
|
||||
temperature: temperature ?? 0.7,
|
||||
tools,
|
||||
stream: false,
|
||||
}),
|
||||
signal: AbortSignal.timeout(600_000),
|
||||
});
|
||||
if (!resp.ok) {
|
||||
const body = await resp.text();
|
||||
throw new Error(`API error ${resp.status}: ${body}`);
|
||||
}
|
||||
const data = (await resp.json()) as RawResponse;
|
||||
const msg = data.choices?.[0]?.message ?? { content: "" };
|
||||
const role = (msg.role ?? "assistant") as Role;
|
||||
const content = msg.content ?? null;
|
||||
const toolCalls: ToolCall[] | undefined = msg.tool_calls?.map((tc) => ({
|
||||
id: tc.id,
|
||||
type: tc.type ?? "function",
|
||||
function: {
|
||||
name: tc.function?.name ?? "",
|
||||
arguments: tc.function?.arguments ?? "",
|
||||
},
|
||||
}));
|
||||
const message: ChatMessage = { role, content };
|
||||
if (toolCalls && toolCalls.length > 0) message.tool_calls = toolCalls;
|
||||
const usage: [number, number] | null = data.usage
|
||||
? ([data.usage.prompt_tokens ?? 0, data.usage.completion_tokens ?? 0] as [number, number])
|
||||
: null;
|
||||
return { message, usage };
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -1,5 +1,10 @@
|
||||
/**
|
||||
* Subagent engine — ported in sub-phase 3c. Placeholder exports so lazy
|
||||
* tool imports resolve cleanly.
|
||||
* Subagent engine — ported in sub-phase 3c.
|
||||
* Mirrors `apps/infrastructure/src/subagent/mod.rs`.
|
||||
*/
|
||||
export const subagentStatus = "not-wired" as const;
|
||||
export { runAgent, type SubagentProviderService, MAX_ITERATIONS, TOOL_OUTPUT_MAX_CHARS } from "../tools/engine.ts";
|
||||
export { toolsFor } from "../tools/division.ts";
|
||||
export { resolveSubagentProvider, resolveSubagentConfig } from "./provider.ts";
|
||||
export { newSubagentContext, type SubagentContext } from "./context.ts";
|
||||
export { buildProviderService } from "./http_provider.ts";
|
||||
export { resolveConfig } from "./config_resolver.ts";
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
/**
|
||||
* Subagent LLM provider — resolves provider/model from settings and wraps
|
||||
* `LlmClient` in a higher-level API for subagent use.
|
||||
* Mirrors `apps/infrastructure/src/subagent/provider.rs`.
|
||||
*
|
||||
* Flow: `resolveSubagentProvider` is called at startup to pick a provider
|
||||
* + model → `SubagentProvider` wraps that pair around a provider service
|
||||
* for use inside the subagent engine loop.
|
||||
*/
|
||||
import type { Settings, AppConfig, ProviderConfig } from "@zesdex/domain";
|
||||
|
||||
/** Provider configuration resolved from settings + app_config. */
|
||||
export interface SubagentProvider {
|
||||
/** The provider name (e.g. "zen", "router", "claude"). */
|
||||
provider: string;
|
||||
/** The resolved model name. */
|
||||
model: string;
|
||||
/** The resolved base URL. */
|
||||
baseUrl: string;
|
||||
/** The resolved API key. */
|
||||
apiKey: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve subagent provider and model from settings + app_config.
|
||||
*
|
||||
* Flow: reads `settings.provider` and `settings.model` → if provider is
|
||||
* empty, falls back to `app_config.default_provider` → resolves model from
|
||||
* provider config or app_config default → finally domain default.
|
||||
*/
|
||||
export function resolveSubagentProvider(
|
||||
settings: Settings,
|
||||
appConfig: AppConfig,
|
||||
): { provider: string; model: string } {
|
||||
const provider = settings.provider === "" ? appConfig.default_provider : settings.provider;
|
||||
|
||||
let model = settings.model;
|
||||
if (model === "") {
|
||||
const cfg: ProviderConfig | undefined = appConfig.providers[provider];
|
||||
model = cfg?.default_model ?? appConfig.default_model ?? "deepseek-v4-flash-free";
|
||||
}
|
||||
|
||||
return { provider, model };
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the full LLM configuration (base_url, api_key, model) for a subagent
|
||||
* from settings + app_config.
|
||||
*
|
||||
* Uses `resolveSubagentProvider` then looks up the provider config for
|
||||
* api_base and api_key from env/default_api_key.
|
||||
*/
|
||||
export function resolveSubagentConfig(
|
||||
settings: Settings,
|
||||
appConfig: AppConfig,
|
||||
): SubagentProvider {
|
||||
const { provider, model } = resolveSubagentProvider(settings, appConfig);
|
||||
const cfg: ProviderConfig = appConfig.providers[provider] ?? { api_base: "" };
|
||||
|
||||
let apiKey = settings.api_keys[provider] ?? "";
|
||||
if (apiKey === "" && cfg.api_key_env) {
|
||||
apiKey = process.env[cfg.api_key_env] ?? cfg.default_api_key ?? "";
|
||||
}
|
||||
|
||||
const baseUrl = cfg.api_base || "https://opencode.ai/zen/v1";
|
||||
|
||||
return { provider, model, baseUrl, apiKey };
|
||||
}
|
||||
@@ -1,12 +1,56 @@
|
||||
/**
|
||||
* Subagent spawning helpers used by the spawn tools. Ported in 3c.
|
||||
* Mirrors `apps/infrastructure/src/subagent/spawn.rs`.
|
||||
*/
|
||||
import type { JsonValue } from "@zesdex/domain";
|
||||
import type { ToolCtx } from "../tools/mod.ts";
|
||||
import { runAgent } from "../tools/engine.ts";
|
||||
import { resolveConfig } from "./config_resolver.ts";
|
||||
import { buildProviderService } from "./http_provider.ts";
|
||||
|
||||
export async function spawnAgents(_agents: JsonValue[], _ctx: ToolCtx): Promise<string> {
|
||||
throw new Error("subagent engine not yet wired in this build");
|
||||
/** Spawn multiple agents in parallel, each with its own directive and access tier. */
|
||||
export async function spawnAgents(
|
||||
agents: Array<{ directive: string; access: string }>,
|
||||
ctx: ToolCtx,
|
||||
): Promise<string> {
|
||||
if (agents.length === 0) return "No agents to spawn.";
|
||||
|
||||
const { baseUrl, apiKey, model } = await resolveConfig();
|
||||
const svc = buildProviderService(baseUrl, apiKey, model);
|
||||
|
||||
const results = await Promise.all(
|
||||
agents.map(async (agent, i) => {
|
||||
const access = (agent.access as "read" | "write" | "full") ?? "read";
|
||||
const out = await runAgent(svc, agent.directive, access, ctx, undefined);
|
||||
return `### Agent ${i}: [${access}]\n\n${out}`;
|
||||
}),
|
||||
);
|
||||
|
||||
return `## Parallel Delegation Results\n\n${results.join("\n\n---\n\n")}`;
|
||||
}
|
||||
export async function spawnPipeline(_stages: JsonValue[], _ctx: ToolCtx): Promise<string> {
|
||||
throw new Error("subagent engine not yet wired in this build");
|
||||
|
||||
/** Spawn a sequential pipeline of agent stages — each stage runs after the previous completes. */
|
||||
export async function spawnPipeline(
|
||||
stages: Array<{ directive: string; access: string }>,
|
||||
ctx: ToolCtx,
|
||||
): Promise<string> {
|
||||
if (stages.length === 0) return "No stages to run.";
|
||||
|
||||
const { baseUrl, apiKey, model } = await resolveConfig();
|
||||
const svc = buildProviderService(baseUrl, apiKey, model);
|
||||
|
||||
const outputs: string[] = [];
|
||||
for (let i = 0; i < stages.length; i++) {
|
||||
const stage = stages[i]!;
|
||||
const access = (stage.access as "read" | "write" | "full") ?? "read";
|
||||
const prevOutput = outputs[outputs.length - 1];
|
||||
const directive = prevOutput
|
||||
? `${stage.directive}\n\nPrevious stage output:\n${prevOutput}`
|
||||
: stage.directive;
|
||||
const out = await runAgent(svc, directive, access, ctx, undefined);
|
||||
outputs.push(out);
|
||||
}
|
||||
|
||||
return `## Sequential Pipeline Results\n\n${outputs
|
||||
.map((o, i) => `### Stage ${i}\n\n${o}`)
|
||||
.join("\n\n---\n\n")}`;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
/**
|
||||
* Subagent division — access-tier tool filtering for subagent permissions.
|
||||
* Mirrors `apps/infrastructure/src/subagent/division.rs`.
|
||||
*
|
||||
* Flow: the calling code picks an `AccessTier` → `toolsFor()` returns the
|
||||
* subset of all built-in tools allowed at that tier → those tools are passed
|
||||
* to the engine for the subagent's tool-execution loop.
|
||||
*/
|
||||
import type { Tool } from "./mod.ts";
|
||||
import { allTools } from "./registry.ts";
|
||||
import type { AccessTier } from "@zesdex/domain";
|
||||
|
||||
/**
|
||||
* Read-only tool names available at the `Read` access tier.
|
||||
* Non-mutating introspection and utility tools only.
|
||||
*/
|
||||
const READ_TIER_TOOLS = new Set([
|
||||
"read",
|
||||
"grep",
|
||||
"glob",
|
||||
"pong",
|
||||
"todowrite",
|
||||
"todofinish",
|
||||
"dir_list",
|
||||
"dir_cache_update",
|
||||
"cd",
|
||||
"remember",
|
||||
"recall",
|
||||
"forget",
|
||||
]);
|
||||
|
||||
/**
|
||||
* Tool names blocked at the `Write` access tier.
|
||||
* Write tier has everything EXCEPT dangerous system/network/process tools.
|
||||
*/
|
||||
const WRITE_TIER_BLOCKED = new Set([
|
||||
"bash",
|
||||
"bash_output",
|
||||
"bash_kill",
|
||||
"git_operator",
|
||||
"git_worktree",
|
||||
"git_cred",
|
||||
"shell",
|
||||
"workflow_run",
|
||||
"note_finding",
|
||||
"read_findings",
|
||||
"hive_mind",
|
||||
"spawn_agents",
|
||||
"spawn_pipeline",
|
||||
"plan_enter",
|
||||
"plan_ready",
|
||||
"sequential_think",
|
||||
]);
|
||||
|
||||
/** Filter the available tools to match the given access tier. */
|
||||
export function toolsFor(access: AccessTier): Tool[] {
|
||||
const all = allTools();
|
||||
switch (access) {
|
||||
case "Read":
|
||||
return all.filter((t) => READ_TIER_TOOLS.has(t.name));
|
||||
case "Write":
|
||||
return all.filter((t) => !WRITE_TIER_BLOCKED.has(t.name));
|
||||
case "Full":
|
||||
return all;
|
||||
default:
|
||||
return all;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,337 @@
|
||||
/**
|
||||
* Subagent engine — runs an LLM-powered agent with tool execution loop.
|
||||
* Mirrors `apps/infrastructure/src/subagent/engine.rs`.
|
||||
*
|
||||
* Flow: construct system message → call LLM → parse tool calls → execute
|
||||
* tools → continue until the model returns a final text response (no more
|
||||
* tool calls) or the iteration limit is reached.
|
||||
*/
|
||||
import {
|
||||
type AgentProgress,
|
||||
type ChatMessage,
|
||||
type JsonValue,
|
||||
type ToolCall,
|
||||
type ToolDef,
|
||||
type TurnEvent,
|
||||
type TurnEventSink,
|
||||
systemMessage,
|
||||
errorRecoveryNote,
|
||||
subagentDirective,
|
||||
Roles,
|
||||
} from "@zesdex/domain";
|
||||
import type { Tool } from "./mod.ts";
|
||||
import type { ToolCtx } from "./context.ts";
|
||||
import { toolDefs, toolIsParallelSafe } from "./registry.ts";
|
||||
import { toolsFor } from "./division.ts";
|
||||
|
||||
/** Maximum number of tool-call iterations before the engine gives up. */
|
||||
export const MAX_ITERATIONS = 25;
|
||||
|
||||
/** A single tool-result message is truncated before entering the subagent's context. */
|
||||
export const TOOL_OUTPUT_MAX_CHARS = 12_000;
|
||||
|
||||
/** Maximum consecutive identical tool errors before the engine injects a recovery note. */
|
||||
export const MAX_CONSECUTIVE_TOOL_ERRORS = 3;
|
||||
|
||||
/** Max read-only tool calls executed concurrently in a single batch. */
|
||||
export const MAX_PARALLEL_TOOLS = 8;
|
||||
|
||||
/** Truncate tool output to TOOL_OUTPUT_MAX_CHARS, char-safe (not byte-safe). */
|
||||
export function truncateToolOutput(output: string): string {
|
||||
if (output.length <= TOOL_OUTPUT_MAX_CHARS) return output;
|
||||
const chars = Array.from(output);
|
||||
const head = chars.slice(0, TOOL_OUTPUT_MAX_CHARS).join("");
|
||||
return `${head}\n...[truncated ${output.length - TOOL_OUTPUT_MAX_CHARS} chars]`;
|
||||
}
|
||||
|
||||
/** Pick a max_tokens budget proportional to the directive's length. */
|
||||
function adaptiveMaxTokens(directiveLen: number): number {
|
||||
if (directiveLen <= 80) return 800;
|
||||
if (directiveLen <= 400) return 1600;
|
||||
return 4096;
|
||||
}
|
||||
|
||||
/** Emit an AgentProgress event onto the turn-event queue if configured. */
|
||||
function reportProgress(toolCtx: ToolCtx, progress: AgentProgress): void {
|
||||
if (toolCtx.turnEvents) {
|
||||
toolCtx.turnEvents.push({ kind: "agent_progress", progress } as TurnEvent);
|
||||
}
|
||||
}
|
||||
|
||||
/** Truncate a string to maxLen characters, appending "…" if truncated. */
|
||||
function truncateStr(s: string, maxLen: number): string {
|
||||
if (s.length <= maxLen) return s;
|
||||
return s.slice(0, maxLen) + "…";
|
||||
}
|
||||
|
||||
// ── Tool execution ────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Execute a single tool call synchronously and return the result string.
|
||||
* Mirrors `run_one_tool` in the Rust engine.
|
||||
*/
|
||||
function runOneTool(tools: Tool[], toolCtx: ToolCtx, toolName: string, args: JsonValue): string {
|
||||
const tool = tools.find((t) => t.name === toolName);
|
||||
if (!tool) return `Unknown tool: ${toolName}`;
|
||||
try {
|
||||
const result = tool.run(toolCtx, args);
|
||||
if (result instanceof Promise) {
|
||||
return `[async] Tool '${toolName}' returned a promise — use parallel_delegate for async`;
|
||||
}
|
||||
return result;
|
||||
} catch (e) {
|
||||
return `Error: ${(e as Error).message}`;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalize tool-call arguments: if arguments is a string, attempt to parse
|
||||
* as JSON; if already an object, return as-is.
|
||||
*/
|
||||
function sanitizeToolArgs(args: JsonValue): JsonValue {
|
||||
if (typeof args === "string") {
|
||||
try {
|
||||
return JSON.parse(args);
|
||||
} catch {
|
||||
return args;
|
||||
}
|
||||
}
|
||||
return args;
|
||||
}
|
||||
|
||||
/** Chunk an array into fixed-size windows. */
|
||||
function chunkArray<T>(arr: T[], size: number): T[][] {
|
||||
const chunks: T[][] = [];
|
||||
for (let i = 0; i < arr.length; i += size) {
|
||||
chunks.push(arr.slice(i, i + size));
|
||||
}
|
||||
return chunks;
|
||||
}
|
||||
|
||||
/**
|
||||
* Execute a batch of tool calls, running read-only tools concurrently when
|
||||
* the whole batch is parallel-safe.
|
||||
*
|
||||
* Returns one result per call IN THE ORIGINAL call order (OpenAI/Anthropic
|
||||
* tool-result ordering contract).
|
||||
*/
|
||||
async function executeToolBatch(
|
||||
tools: Tool[],
|
||||
toolCtx: ToolCtx,
|
||||
toolCalls: ToolCall[],
|
||||
): Promise<Array<{ id: string; name: string; output: string }>> {
|
||||
const parallel =
|
||||
toolCalls.length > 1 && toolCalls.every((tc) => toolIsParallelSafe(tc.function.name));
|
||||
|
||||
if (!parallel) {
|
||||
const results: Array<{ id: string; name: string; output: string }> = [];
|
||||
for (const tc of toolCalls) {
|
||||
const args = sanitizeToolArgs(tc.function.arguments);
|
||||
const output = runOneTool(tools, toolCtx, tc.function.name, args);
|
||||
results.push({ id: tc.id, name: tc.function.name, output });
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
const allResults: Array<{ id: string; name: string; output: string }> = [];
|
||||
const windows = chunkArray(toolCalls, MAX_PARALLEL_TOOLS);
|
||||
|
||||
for (const window of windows) {
|
||||
const windowResults = await Promise.all(
|
||||
window.map(async (tc) => ({
|
||||
id: tc.id,
|
||||
name: tc.function.name,
|
||||
output: runOneTool(tools, toolCtx, tc.function.name, sanitizeToolArgs(tc.function.arguments)),
|
||||
})),
|
||||
);
|
||||
allResults.push(...windowResults);
|
||||
}
|
||||
return allResults;
|
||||
}
|
||||
|
||||
// ── LLM interaction ──────────────────────────────────────────────────────
|
||||
|
||||
/** Provider interface for subagent LLM calls. */
|
||||
export interface SubagentProviderService {
|
||||
chat(
|
||||
messages: ChatMessage[],
|
||||
tools?: ToolDef[],
|
||||
maxTokens?: number,
|
||||
temperature?: number,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Execute the LLM call and return the assembled assistant message.
|
||||
*/
|
||||
async function callLlm(
|
||||
provider: SubagentProviderService,
|
||||
messages: ChatMessage[],
|
||||
defs: ToolDef[],
|
||||
maxTokens: number,
|
||||
_sink: TurnEventSink | null,
|
||||
): Promise<{ message: ChatMessage; usage: [number, number] | null }> {
|
||||
return await provider.chat(messages, defs, maxTokens, 0.2);
|
||||
}
|
||||
|
||||
// ── Main engine loop ─────────────────────────────────────────────────────
|
||||
|
||||
/** Resolve allowed tools for a given access tier string. */
|
||||
function toolsForAccess(access: "read" | "write" | "full"): Tool[] {
|
||||
const tier: "Read" | "Write" | "Full" =
|
||||
access === "read" ? "Read" : access === "write" ? "Write" : "Full";
|
||||
return toolsFor(tier);
|
||||
}
|
||||
|
||||
/** Build the system message for a subagent. */
|
||||
function buildSystemMessage(directive: string, toolCtx: ToolCtx): ChatMessage {
|
||||
const cwd = process.cwd();
|
||||
const wsRoot = toolCtx.workspaces[0] ?? cwd;
|
||||
return systemMessage(subagentDirective(directive, cwd, wsRoot));
|
||||
}
|
||||
|
||||
/**
|
||||
* Run an agent with a directive, access tier, and tool context.
|
||||
*
|
||||
* Flow:
|
||||
* 1. Resolve allowed tools for the given `access` tier.
|
||||
* 2. Build a system prompt from the directive using the domain prompt module.
|
||||
* 3. Loop (up to MAX_ITERATIONS):
|
||||
* a. Call the LLM (non-streaming) with accumulated messages + tool defs.
|
||||
* b. If the response has no tool calls → return the text content.
|
||||
* c. Otherwise execute each tool call and append the result as a
|
||||
* tool-role message.
|
||||
* 4. If the loop exits naturally, return the iteration-limit message.
|
||||
*
|
||||
* Progress: each tool invocation is reported via AgentProgress if a turn-event
|
||||
* queue is available in the ToolCtx.
|
||||
*/
|
||||
export async function runAgent(
|
||||
provider: SubagentProviderService,
|
||||
directive: string,
|
||||
directiveAccess: "read" | "write" | "full",
|
||||
toolCtx: ToolCtx,
|
||||
onProgress?: (progress: AgentProgress) => void,
|
||||
): Promise<string> {
|
||||
const tools = toolsForAccess(directiveAccess);
|
||||
const defs = toolDefs(tools) as unknown as ToolDef[];
|
||||
|
||||
const sysMsg = buildSystemMessage(directive, toolCtx);
|
||||
const messages: ChatMessage[] = [sysMsg];
|
||||
|
||||
const maxTokens = adaptiveMaxTokens(directive.length);
|
||||
|
||||
let consecutiveErrors = 0;
|
||||
let lastTool = "";
|
||||
|
||||
for (let iteration = 0; iteration < MAX_ITERATIONS; iteration++) {
|
||||
if (toolCtx.abortFlag?.aborted) {
|
||||
const prog: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: truncateStr(directive, 40),
|
||||
status: "Cancelled",
|
||||
current_tool: null,
|
||||
steps: null,
|
||||
};
|
||||
reportProgress(toolCtx, prog);
|
||||
onProgress?.(prog);
|
||||
return "Subagent was cancelled.";
|
||||
}
|
||||
|
||||
const runningProg: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: truncateStr(directive, 40),
|
||||
status: "Running",
|
||||
current_tool: null,
|
||||
steps: [iteration, MAX_ITERATIONS],
|
||||
};
|
||||
reportProgress(toolCtx, runningProg);
|
||||
onProgress?.(runningProg);
|
||||
|
||||
let responseMsg: ChatMessage;
|
||||
try {
|
||||
const result = await callLlm(provider, messages, defs, maxTokens, null);
|
||||
responseMsg = result.message;
|
||||
} catch (e) {
|
||||
const msg = (e as Error).message;
|
||||
const failedProg: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: truncateStr(directive, 40),
|
||||
status: "Failed",
|
||||
error: msg,
|
||||
current_tool: null,
|
||||
steps: null,
|
||||
};
|
||||
reportProgress(toolCtx, failedProg);
|
||||
onProgress?.(failedProg);
|
||||
return `Subagent failed: ${msg}`;
|
||||
}
|
||||
|
||||
const content = responseMsg.content ?? "";
|
||||
const toolCalls = responseMsg.tool_calls ?? [];
|
||||
|
||||
if (toolCalls.length === 0) {
|
||||
const doneProg: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: truncateStr(directive, 40),
|
||||
status: "Completed",
|
||||
current_tool: null,
|
||||
steps: null,
|
||||
};
|
||||
reportProgress(toolCtx, doneProg);
|
||||
onProgress?.(doneProg);
|
||||
return content;
|
||||
}
|
||||
|
||||
// Push assistant message BEFORE executing tools (tool-calling contract)
|
||||
messages.push(responseMsg);
|
||||
|
||||
const results = await executeToolBatch(tools, toolCtx, toolCalls);
|
||||
|
||||
for (const { id, name, output } of results) {
|
||||
const runningToolProg: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: `running:${name}`,
|
||||
status: "Running",
|
||||
current_tool: name,
|
||||
steps: null,
|
||||
};
|
||||
reportProgress(toolCtx, runningToolProg);
|
||||
onProgress?.(runningToolProg);
|
||||
|
||||
// Error-recovery: if the same tool keeps failing, inject a system note
|
||||
if (output.startsWith("Error:")) {
|
||||
if (lastTool === name) {
|
||||
consecutiveErrors += 1;
|
||||
} else {
|
||||
consecutiveErrors = 1;
|
||||
lastTool = name;
|
||||
}
|
||||
if (consecutiveErrors >= MAX_CONSECUTIVE_TOOL_ERRORS) {
|
||||
messages.push(systemMessage(errorRecoveryNote(name, output)));
|
||||
consecutiveErrors = 0;
|
||||
}
|
||||
} else {
|
||||
consecutiveErrors = 0;
|
||||
}
|
||||
|
||||
messages.push({
|
||||
role: Roles.Tool,
|
||||
content: truncateToolOutput(output),
|
||||
tool_call_id: id,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
const failProg: AgentProgress = {
|
||||
agent_id: "subagent",
|
||||
agent_name: truncateStr(directive, 40),
|
||||
status: "Failed",
|
||||
error: `iteration limit (${MAX_ITERATIONS})`,
|
||||
current_tool: null,
|
||||
steps: null,
|
||||
};
|
||||
reportProgress(toolCtx, failProg);
|
||||
onProgress?.(failProg);
|
||||
return `Subagent reached iteration limit (${MAX_ITERATIONS})`;
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
/**
|
||||
* Simple sink adapter that wraps a TurnEvent[] (array-based) into the
|
||||
* TurnEventSink interface, providing a `push` method and a no-op `drain`.
|
||||
*/
|
||||
import type { TurnEvent, TurnEventSink } from "@zesdex/domain";
|
||||
|
||||
/**
|
||||
* Create a TurnEventSink backed by a plain array.
|
||||
* `push` appends to the array; `drain` returns a copy and clears the array.
|
||||
*/
|
||||
export function arrayEventSink(events: TurnEvent[]): TurnEventSink {
|
||||
return {
|
||||
push(event: TurnEvent): void {
|
||||
events.push(event);
|
||||
},
|
||||
drain(): TurnEvent[] {
|
||||
const copy = [...events];
|
||||
events.length = 0;
|
||||
return copy;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -34,7 +34,7 @@ export class SpawnAgents implements Tool {
|
||||
const agents = reqArray(args, "agents");
|
||||
try {
|
||||
const { spawnAgents } = await import("../subagent/spawn_tools.ts");
|
||||
return await spawnAgents(agents, ctx);
|
||||
return await spawnAgents(agents as Array<{ directive: string; access: string }>, ctx);
|
||||
} catch (e) {
|
||||
const msg = (e as Error).message;
|
||||
if (msg.includes("not yet") || msg.includes("Cannot find")) {
|
||||
@@ -72,7 +72,7 @@ export class SpawnPipeline implements Tool {
|
||||
const stages = reqArray(args, "stages");
|
||||
try {
|
||||
const { spawnPipeline } = await import("../subagent/spawn_tools.ts");
|
||||
return await spawnPipeline(stages, ctx);
|
||||
return await spawnPipeline(stages as Array<{ directive: string; access: string }>, ctx);
|
||||
} catch (e) {
|
||||
const msg = (e as Error).message;
|
||||
if (msg.includes("not yet") || msg.includes("Cannot find")) {
|
||||
|
||||
Reference in New Issue
Block a user