From a14a2275ff90059874166aec8c03bd44e6c7efee Mon Sep 17 00:00:00 2001 From: asepharyana Date: Wed, 2 Sep 2026 19:27:36 +0700 Subject: [PATCH] =?UTF-8?q?feat(subagent):=20port=203c=20subagent=20engine?= =?UTF-8?q?=20=E2=80=94=20run=5Fagent=20loop=20+=20spawn/delegate/parallel?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- TODO.md | 198 ++++++++++ .../src/subagent/config_resolver.ts | 37 ++ .../infrastructure/src/subagent/context.ts | 54 +++ .../infrastructure/src/subagent/delegate.ts | 46 ++- .../src/subagent/http_provider.ts | 72 ++++ .../infrastructure/src/subagent/index.ts | 11 +- .../infrastructure/src/subagent/provider.ts | 68 ++++ .../src/subagent/spawn_tools.ts | 54 ++- .../infrastructure/src/tools/division.ts | 68 ++++ .../infrastructure/src/tools/engine.ts | 337 ++++++++++++++++++ .../infrastructure/src/tools/event_sink.ts | 22 ++ .../infrastructure/src/tools/spawn.ts | 4 +- 12 files changed, 956 insertions(+), 15 deletions(-) create mode 100644 TODO.md create mode 100644 apps/packages/infrastructure/src/subagent/config_resolver.ts create mode 100644 apps/packages/infrastructure/src/subagent/context.ts create mode 100644 apps/packages/infrastructure/src/subagent/http_provider.ts create mode 100644 apps/packages/infrastructure/src/subagent/provider.ts create mode 100644 apps/packages/infrastructure/src/tools/division.ts create mode 100644 apps/packages/infrastructure/src/tools/engine.ts create mode 100644 apps/packages/infrastructure/src/tools/event_sink.ts diff --git a/TODO.md b/TODO.md new file mode 100644 index 0000000..003b6ef --- /dev/null +++ b/TODO.md @@ -0,0 +1,198 @@ +# TODO — Rewrite Zesdex Rust → TypeScript/Bun + +> Plan lengkap migrasi in-place dari Rust (11 crate, clean architecture 4 lapis) ke monorepo +> TypeScript/Bun. Bahasa Indonesia, commit pakai Conventional Commits. +> +> **Status:** Fase 0–2 + 3a + 3b + 3e + 3c **selesai**. Berikutnya: **3d workflow + hive-mind**. +> Base: `bun run check` bersih, 54 test hijau. + +--- + +## Prinsip & peta pemetaan + +| Rust | TS/Bun | +|---|---| +| `serde_json` | `JSON.stringify/parse` + `zod` | +| `chrono` | `Date` / `Date.now()` (epoch ms) | +| `uuid` | `crypto.randomUUID()` | +| `sha2`, `base64 URL_SAFE_NO_PAD` | `node:crypto` / `Buffer.toString('base64url')` | +| `reqwest` | `fetch` / bun HTTP | +| `rusqlite` | `bun:sqlite` | +| `argon2` | `@node-rs/argon2` | +| `jsonwebtoken` | `node:crypto` HMAC HS256 | +| `rmcp` | `@modelcontextprotocol/sdk` | +| `ratatui`/`crossterm` | library TUI Bun / ink | +| `tiktoken-rs` | `@dqbd/tiktoken` (cl100k_base) | +| `clap` | `Bun.argv` / `commander` | +| `tokio` semaphore | Bun async + concurrency pool | +| `libc` PID lock | `node:fs` `O_EXCL` + `process.kill` liveness | +| Rust generics `S` | DI via constructor interface | +| `Arc>` | objek mutable (array/objek di ToolCtx) | + +### Struktur target + +``` +apps/ + packages/ # ✅ sudah ada: domain, application, infrastructure + interfaces/ + cli/ # gateway + TUI [baru] + api/ # REST API [baru] + daemon/ # daemon + IPC [baru] + ws/ # WebSocket server [baru] + grpc/ # stub [baru] + web/ # static server [baru] +``` + +--- + +## ✅ Selesai + +### Fase 0 — Bootstrap monorepo +- [x] `package.json` root (workspaces), `tsconfig.json` strict, `.gitignore` +- [x] Paket `domain`, `application`, `infrastructure` (workspace `@zesdex/*`) +- [x] `bun install` + `bun test` jalan + +### Fase 1 — Domain layer +- [x] `core`: ChatMessage, Conversation, ChatRequest/Response/StreamEvent, SseParser, ToolCall, repairJson, sanitizeToolArguments, ToolCallResult, UsageStats, Store +- [x] `auth`: Session, SessionId (path-traversal), OAuthToken/Config, repo/service trait +- [x] `cms`: AppConfig, ProviderConfig, ModelRole, Settings+resolveEffectiveModel, SettingsPatch, Memory, EditLog, repo/service trait +- [x] `agent`: Origin, Toast, AgentStatus, TurnEvent, SessionRuntime, AgentTurnParams, AgentProgress, ExecutionModel +- [x] `workflow`/`subagent`: WorkflowScript/NodeDirective, AccessTier, prompt builders +- [x] 27 test domain + +### Fase 2 — Application layer +- [x] `ports`: ProviderService, PasswordService, TokenService, AuthService +- [x] `AgentTurnServiceImpl` (loop 50 iter, auto-compact 60k, adaptive token, temp, ErrorTracker, bounded parallel) +- [x] `OAuthUseCase` (PKCE S256 + CSRF state) +- [x] `SessionServiceImpl`, `MemoryServiceImpl`, `SettingsServiceImpl`, `ConversationServiceImpl` +- [x] 11 test application + +### Fase 3a — LLM client +- [x] `LlmClient` (fetch, retry 10x backoff+jitter, abort 401/402/403) +- [x] `resolve_api_key`, `is_auth_error`, `backoff_for_error` + +### Fase 3b — Tool system (`commit a291dd1`) +- [x] `Tool` interface, `ToolCtx`/`ToolCtxBuilder`, `InfrastructureToolExecutor` +- [x] Registry: `allTools` (37), `toolDefs`, `toolIsRisky`, `toolIsParallelSafe` +- [x] FS: read, write, edit, delete · **Guard**: graduated checks, resolve_path sandbox +- [x] Search: grep, glob (matcher ringan) +- [x] Shell: bash (timeout), bash_output, bash_kill · **Filter**: checkGitDestructive, checkCredentialRead +- [x] Git: git_operator, git_worktree, git_cred +- [x] Memory: remember, forget, recall +- [x] Util: cd, dir_list, dir_cache_update, pong, todowrite, todofinish +- [x] Reasoning/Plan: sequential_think, plan_enter, plan_ready +- [x] Web: web_search (async → SearXNG) +- [x] Semantic: semantic_search, rebuild_index, list_symbols (index multi-bahasa) +- [x] Best-practice: best_practice, commit_convention + BestPracticeEngine +- [x] Agent/Workflow: spawn_agents, spawn_pipeline, parallel_delegate, workflow_run, note_finding, read_findings, hive_mind (delegasi ke placeholder 3c/3d) +- [x] 16 test tool system + +### Fase 3e — Persistence +- [x] File repos: settings, app_config (+ auto-detect kredensial Claude), conversation, edits.jsonl, markdown memory, rewind blob, session, oauth, session lock +- [x] `writeJsonAtomic`, `slugify`, `buildWorkspaceTree`, `runCommand` + +### Fase 3c — Subagent engine +**Sumber Rust:** `apps/infrastructure/src/subagent/{engine,division,context,provider,spawn}.rs`, `subagent/auto/*` + +- [x] `AccessTier.toolsFor(access)` — `tools/division.ts`: Read / Write / Full yang mem-filter registry + - Read → tool baca-saja + - Write → baca + tulis + - Full → semua +- [x] `run_agent` loop (`tools/engine.ts`): `MAX_ITERATIONS=25`, `TOOL_OUTPUT_MAX_CHARS=12_000`, + `MAX_CONSECUTIVE_TOOL_ERRORS=3`, `MAX_PARALLEL_TOOLS=8`; recovery note setelah 3 error beruntun; + return pesan limit iterasai saat keluar alami +- [x] `SubagentContext` (`subagent/context.ts`) — directive, toolCtx, access, baseURL, apiKey, model +- [x] `spawn_subagent` — `spawn_tools.ts` — runAgent berbasis HTTP provider + config_resolver +- [x] `resolve_subagent_provider` (`subagent/provider.ts`) — pilih provider/model subagent dari settings + app_config +- [x] `spawn_tools.ts` — isi body `spawnAgents` (paralel) & `spawnPipeline` (berurutan) +- [x] `delegate.ts` — `runParallelDelegation` untuk `parallel_delegate` +- [x] `http_provider.ts` — shared HTTP provider service builder +- [ ] Auto-reviewer (`subagent/auto/*`) — `spawn_background_review`: git diff → LLM review → perbaiki; + wiring `review_enabled` (yang sekarang no-op di `Review` agent) +- [x] Placeholder di `apps/packages/infrastructure/src/subagent/*` diisi penuh + +--- + +## ⬜ Yang belum dikerjakan + +### Fase 3d — Workflow + hive-mind +**Sumber Rust:** `apps/infrastructure/src/workflow/{script,mod,docs}.rs`, `workflow/engine/{execution,phases,primitives}.rs`, `workflow/hive_mind/{complexity,cycle,synthesis,live,types}.rs` + +- [ ] `parse_workflow_script(yaml)` — parser YAML workflow (`script.ts`) +- [ ] `execute_workflow(script, ctx)` — eksekusi multi-phase (`engine/execution.ts`) +- [ ] Phase primitives: run sub-agent per node, gate/condition, paralel (bounded), I/O +- [ ] `execute_cycle(cycle, ctx)` — hive-mind cycle, concurrency 8, kumpulkan NodeOutput +- [ ] `synthesize_consensus(node_outputs, ctx)` — konsensus via LLM +- [ ] `complexity_heuristic` — pilih apakah cukup 1 cycle / butuh banyak +- [ ] `docs.ts` — tulis convergence/decision doc +- [ ] Isi body `WorkflowRun` & `HiveMind` tool (ganti placeholder), hapus `workflow/index.ts` + `hive_mind.ts` placeholder + +### Fase 3f — Auth infrastructure +**Sumber Rust:** `apps/infrastructure/src/auth/{password,jwt,oauth_loopback,mod}.rs` + +- [ ] `password.ts` — Argon2 via `@node-rs/argon2` (hash + verify), port `PasswordService` +- [ ] `jwt.ts` — HMAC HS256 via `node:crypto` (access + refresh token, typ, exp), port `TokenService` +- [ ] `oauth_loopback.ts` — server TCP loopback lokal utk terima redirect OAuth; port `oauth_loopback` +- [ ] `mod.ts` — composisi ke `AuthService` +- [ ] Tambah dep `@node-rs/argon2` di `infrastructure/package.json` + +### Fase 3g — IPC + bgbash +**Sumber Rust:** `apps/infrastructure/src/ipc/{frame,protocol,server,client,conn}.rs`, `apps/infrastructure/src/bgbash/{job,control}.rs` + +- [ ] `ipc/frame.ts` — length-prefixed framing via `node:net` (kirim/terima bingkai biner) +- [ ] `ipc/server.ts` + `ipc/client.ts` — Unix socket server + attach client; protocol message (NLJSON/binary) +- [ ] `bgbash/job.ts` — `BashJob` (pid, output file, log), `spawn_bash_job` +- [ ] `bgbash/control.ts` — registry job global: `register`, `cancel(job_id)`, `is_running` +- [ ] Wire `bash` tool `run_in_background` → `spawn_bash_job` + `register`; `bash_kill` → `cancel` +- [ ] Update `bash_output` untuk baca dari registry/output dir + +### Fase 4 — Gateway + CLI + TUI +**Sumber Rust:** `apps/gateway/src/{lib,main}.rs`, `apps/interfaces/tui/src/**` + +- [ ] `apps/interfaces/cli/` (paket Bun): + - [ ] Parser arg: `--api/--daemon/--attach/--ws/--grpc/--web/--tui`, dispatch mode + - [ ] Logging ke file + - [ ] `run_single_process` — komposisi root: wire semua implementasi konkret sekali +- [ ] TUI mode (library TUI Bun): + - [ ] Event loop adaptive poll + - [ ] `Action` dispatcher + slash commands + - [ ] Transcript render (incremental cache), sidebar, overlays, markdown render + - [ ] Theme +- [ ] Bootstrap seed (`bootstrap`) + +### Fase 5 — Server interface: api / daemon / ws / grpc / web +**Sumber Rust:** `apps/interfaces/{api,daemon,ws,grpc,web}/src/**`, `apps/infrastructure/src/middleware/*` + +- [ ] **API** (`apps/interfaces/api/`): + - [ ] Endpoint: `auth` (login/register/refresh, rate-limited), `sessions` CRUD, `conversations`, `chat/completions` proxy, `health` + - [ ] JWT middleware, CORS + - [ ] DTO (dto/) +- [ ] **daemon** — IPC server + attach TUI client (pakai IPC 3g) +- [ ] **ws** — handler `/ws`, protokol token/connected/done/error +- [ ] **grpc** — stub (health endpoint) +- [ ] **web** — static server + proteksi traversal; frontend placeholder + +### Fase 6 — Bersih-bersih + verifikasi akhir +- [ ] Hapus semua sisa crate Rust: `apps/{domain,application,infrastructure,gateway,bootstrap,interfaces/*}` file `.rs`, `Cargo.toml`, `Cargo.lock`, `.cargo/`, `lefthook.yml` +- [ ] Update `install.sh` → Bun +- [ ] Update CI workflow: `.github/workflows/{ci,deploy,release,flakehub-publish-rolling}.yml` → `bun install`/`bun test`/`bun run cli` +- [ ] Update `Dockerfile` + `flake.nix`/`shell.nix` bila dipakai (atau hapus) +- [ ] Ganti/luruskan tool release (semantic-release) ke pipeline Bun +- [ ] Berkas jadi: hapus folder Rust yang sudah kosong dari `apps/` +- [ ] Verifikasi akhir: `bun run build`, `bun test`, `tsc --noEmit`, `bun run cli --help`, smoke server api/ws/daemon +- [ ] Baca ulang file lama (`settings.json`, `app_config.json`, `conversation.json`, memory `.md`) — pastikan wire-shape terbaca + +--- + +## Verifikasi global (tiap fase) +- [ ] `bun run build` / `tsc --noEmit` tanpa error (strict) +- [ ] `bun test` — port test Rust ke `*.test.ts` (SseParser, repairJson, resolveEffectiveModel, SessionId, turn_service, tool system, subagent, workflow) +- [ ] Smoke: `bun run cli --help`; agent mode TUI dengan model stub/lokal; pastikan loop turn + tool read/write +- [ ] Server: `bun run api` + hit health/auth; `bun run ws` + connect +- [ ] Commit tiap fase dengan message Bahasa Indonesia + Conventional Commits; test hijau sebelum lanjut + +## Catatan +- Registry tool = **38** (terhitung workflow_run, note_finding, read_findings, hive_mind; ExploreCodebase tidak diport). +- `Tool.run` boleh return `string | Promise`. +- `spawn_agents`/`spawn_pipeline`/`parallel_delegate` sudah terisi penuh via 3c. +- Workflow/hive-mind tools masih placeholder — diisi penuh di 3d. diff --git a/apps/packages/infrastructure/src/subagent/config_resolver.ts b/apps/packages/infrastructure/src/subagent/config_resolver.ts new file mode 100644 index 0000000..7dc1e6a --- /dev/null +++ b/apps/packages/infrastructure/src/subagent/config_resolver.ts @@ -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 { + 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 }; +} diff --git a/apps/packages/infrastructure/src/subagent/context.ts b/apps/packages/infrastructure/src/subagent/context.ts new file mode 100644 index 0000000..a254ffe --- /dev/null +++ b/apps/packages/infrastructure/src/subagent/context.ts @@ -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"; + } +} diff --git a/apps/packages/infrastructure/src/subagent/delegate.ts b/apps/packages/infrastructure/src/subagent/delegate.ts index beca755..e096f73 100644 --- a/apps/packages/infrastructure/src/subagent/delegate.ts +++ b/apps/packages/infrastructure/src/subagent/delegate.ts @@ -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 { - 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; } diff --git a/apps/packages/infrastructure/src/subagent/http_provider.ts b/apps/packages/infrastructure/src/subagent/http_provider.ts new file mode 100644 index 0000000..a464cde --- /dev/null +++ b/apps/packages/infrastructure/src/subagent/http_provider.ts @@ -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 }; + }, + }; +} diff --git a/apps/packages/infrastructure/src/subagent/index.ts b/apps/packages/infrastructure/src/subagent/index.ts index 00c8247..82ea20a 100644 --- a/apps/packages/infrastructure/src/subagent/index.ts +++ b/apps/packages/infrastructure/src/subagent/index.ts @@ -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"; diff --git a/apps/packages/infrastructure/src/subagent/provider.ts b/apps/packages/infrastructure/src/subagent/provider.ts new file mode 100644 index 0000000..703f86c --- /dev/null +++ b/apps/packages/infrastructure/src/subagent/provider.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 }; +} diff --git a/apps/packages/infrastructure/src/subagent/spawn_tools.ts b/apps/packages/infrastructure/src/subagent/spawn_tools.ts index cb4eaf6..074b705 100644 --- a/apps/packages/infrastructure/src/subagent/spawn_tools.ts +++ b/apps/packages/infrastructure/src/subagent/spawn_tools.ts @@ -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 { - 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 { + 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 { - 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 { + 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")}`; } diff --git a/apps/packages/infrastructure/src/tools/division.ts b/apps/packages/infrastructure/src/tools/division.ts new file mode 100644 index 0000000..30a3991 --- /dev/null +++ b/apps/packages/infrastructure/src/tools/division.ts @@ -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; + } +} diff --git a/apps/packages/infrastructure/src/tools/engine.ts b/apps/packages/infrastructure/src/tools/engine.ts new file mode 100644 index 0000000..4dab54b --- /dev/null +++ b/apps/packages/infrastructure/src/tools/engine.ts @@ -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(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> { + 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 { + 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})`; +} diff --git a/apps/packages/infrastructure/src/tools/event_sink.ts b/apps/packages/infrastructure/src/tools/event_sink.ts new file mode 100644 index 0000000..48f7b13 --- /dev/null +++ b/apps/packages/infrastructure/src/tools/event_sink.ts @@ -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; + }, + }; +} diff --git a/apps/packages/infrastructure/src/tools/spawn.ts b/apps/packages/infrastructure/src/tools/spawn.ts index f526e19..4b8803c 100644 --- a/apps/packages/infrastructure/src/tools/spawn.ts +++ b/apps/packages/infrastructure/src/tools/spawn.ts @@ -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")) {