diff --git a/.gitignore b/.gitignore index 044b6ea46..0c22cc6eb 100644 --- a/.gitignore +++ b/.gitignore @@ -70,3 +70,6 @@ python/robomp/.cache/ python/robomp/src/robomp/static/ python/robomp/web/dist/ python/robomp/.env + +# Local, machine-specific boot perf baseline (see packages/coding-agent/scripts/bench-guard.ts) +packages/coding-agent/bench/boot-baseline.json diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 00b5d6905..8b6c981f5 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -438,6 +438,14 @@ - Removed the `clearOnShrink` setting and its `PI_CLEAR_ON_SHRINK` environment variable: the rewritten renderer always clears shrunken rows exactly, so the flicker/perf tradeoff the setting controlled no longer exists. Existing config entries are ignored. - Removed the prompt-submit native-scrollback reconciliation checkpoint and the eager streaming render mode from the interactive controllers — the renderer's append-only contract made both obsolete. +### Changed + +- Interactive boot no longer blocks first paint on MCP server discovery. For UI sessions, `createAgentSession` returns immediately and connects to configured MCP servers in the background; their tools and slash commands stream in through the existing live-refresh channel once connected. An explicitly requested MCP tool whose server has not finished connecting resolves to a deterministic "still connecting" placeholder instead of an "unknown tool" error, and each server's instructions join the system prompt once its background connection completes — carried in on the same live-refresh that registers its tools — so server-provided instructions are preserved, not dropped. `tools.discoveryMode: "auto"` is re-resolved once the background connect reports the real MCP tool count, so a toolset large enough to cross the threshold flips discovery on (registering `search_tool_bm25`) instead of force-activating every tool, and a session disposed while servers are still connecting disconnects them instead of resurrecting their tools onto the dead session. Non-UI modes (`print`/`rpc`/`acp`) keep the blocking discovery path. Measured ~290 ms (≈24% of cold boot) off the first-paint critical path with MCP servers configured. +- Assistant-message streaming is cheaper per token. During a stream the component now reuses its `Markdown` subtree across reveal ticks — only the growing block is re-rendered — instead of tearing down and re-lexing every block on each ~30 fps tick, and grapheme counting in the reveal controller is memoized. A completed thinking block that precedes a still-streaming answer is no longer re-highlighted every frame (~66% less render work on think-then-answer streams in benchmarks; single-block streams are unchanged). +- Cold boot no longer builds the model catalog's canonical-equivalence index on the first-paint critical path. The `ModelRegistry` constructor built `buildCanonicalModelIndex` over the entire ~3,200-model catalog synchronously (~210 ms); it is now built lazily on first read (`getCanonicalModels`/`getCanonicalVariants`/`getCanonicalId`, reached by the model picker and by `enabledModels`/default-role pattern resolution), which a default interactive launch never touches before paint. Measured ~244 ms (≈16% of cold-boot wall) off first paint; the picker pays the one-time build on first open. +- Repeat `read` summaries of an unchanged file no longer re-run the tree-sitter parse. The per-session summary is memoized on the content hash of the freshly-read bytes — the file is still read fresh on every call, so results stay correct without a staleness window — dropping a repeated same-file summary read from ~17 ms to ~2.5 ms. +- Attributed the previously-unlabeled synchronous boot region in the `PI_TIMING` startup table with `modelRegistry:init`, `buildCanonicalModelIndex`, and `initTelemetryExport` spans. + ## [15.10.9] - 2026-06-09 ### Fixed diff --git a/packages/coding-agent/bench/rendering.ts b/packages/coding-agent/bench/rendering.ts index e15d82976..8e32428a6 100644 --- a/packages/coding-agent/bench/rendering.ts +++ b/packages/coding-agent/bench/rendering.ts @@ -1,6 +1,18 @@ import { initTheme } from "../src/modes/theme/theme"; import { truncateToVisualLines } from "../src/modes/components/visual-truncate"; import { WelcomeComponent } from "../src/modes/components/welcome"; +import type { AssistantMessage } from "@oh-my-pi/pi-ai"; +import { Editor } from "@oh-my-pi/pi-tui"; +import { AssistantMessageComponent } from "../src/modes/components/assistant-message"; +import { TranscriptContainer } from "../src/modes/components/transcript-container"; +import { Settings } from "../src/config/settings"; +import { getEditorTheme } from "../src/modes/theme/theme"; +import { buildDisplayMessage, nextStep, visibleUnits } from "../src/modes/controllers/streaming-reveal"; +import * as os from "node:os"; +import * as path from "node:path"; +import * as fs from "node:fs"; +import { ReadTool } from "../src/tools/read"; +import type { ToolSession } from "../src/tools"; const ITERATIONS = 500; const WIDTH = 100; @@ -20,6 +32,7 @@ function bench(name: string, fn: () => void): number { return elapsed; } +await Settings.init({ inMemory: true }); await initTheme("dark"); console.log(`Rendering benchmark (${ITERATIONS} iterations)\n`); @@ -39,3 +52,223 @@ const welcome = new WelcomeComponent("8.12.3", "claude-3.7", "anthropic", [ bench("WelcomeComponent.render", () => { welcome.render(WIDTH); }); + +// ── A2: streaming reveal + editor render baselines ────────────────────────── +// +// Diagnostic series, not a fixed-iteration micro-op. `streamingReveal` proves +// or refutes the O(N^2) reveal hypothesis: per-step cost (visibleUnits + +// buildDisplayMessage, the work every stream delta/30fps tick does) is sampled +// at growing revealed lengths. Rising per-step ms => O(N) per tick => O(N^2) +// over the message. Flat per-step => already linear. + +function makeMarkdownCorpus(targetGraphemes: number): string { + const para = + "The quick brown fox jumps over the lazy dog while 🚀 emoji and a `code span` " + + "plus **bold** and _italic_ text exercise the markdown lexer and the grapheme segmenter. "; + const codeBlock = "\n```ts\nconst x: number = compute(a, b) + delta;\nreturn x.toFixed(2);\n```\n\n"; + const list = "\n- first bullet item\n- second bullet item with `inline`\n- third\n\n"; + let out = ""; + let i = 0; + while (out.length < targetGraphemes) { + out += `## Section ${++i}\n\n${para}${para}${codeBlock}${list}`; + } + return out.slice(0, targetGraphemes); +} + +function makeTextMessage(text: string): AssistantMessage { + return { + role: "assistant", + content: [{ type: "text", text }], + api: "anthropic-messages", + provider: "anthropic", + model: "bench", + usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 } }, + stopReason: "stop", + timestamp: 0, + }; +} + +/** Average ms for one call of `fn`, over `reps` repeats. */ +function benchStep(reps: number, fn: () => void): number { + const start = Bun.nanoseconds(); + for (let i = 0; i < reps; i++) fn(); + return (Bun.nanoseconds() - start) / 1e6 / reps; +} + +/** Average ms for one awaited call of `fn`, over `reps` repeats. */ +async function benchStepAsync(reps: number, fn: () => Promise): Promise { + const start = Bun.nanoseconds(); + for (let i = 0; i < reps; i++) await fn(); + return (Bun.nanoseconds() - start) / 1e6 / reps; +} + +const REVEAL_CORPUS = makeMarkdownCorpus(6000); +const REVEAL_CHECKPOINTS = [1000, 2000, 3000, 4000, 5000, 6000]; +const STEP_REPS = 40; + +console.log("\nstreamingReveal (isolated C1: visibleUnits + buildDisplayMessage per delta):"); +for (const n of REVEAL_CHECKPOINTS) { + const msg = makeTextMessage(REVEAL_CORPUS.slice(0, n)); + const revealed = Math.floor(n * 0.9); + const ms = benchStep(STEP_REPS, () => { + visibleUnits(msg, false); + buildDisplayMessage(msg, revealed, false); + }); + console.log(` len=${n}: ${ms.toFixed(4)}ms/step`); +} + +// Real streaming cost: text GROWS every tick, so Markdown's text-keyed cache +// misses each step (the actual interactive path). Total ms to fully reveal an +// N-grapheme message in nextStep increments — the number C1+C2 reduce. +console.log("\nstreamingRevealFull (C1+C2: full incremental reveal, growing text => cache-miss/tick):"); +try { + for (const n of REVEAL_CHECKPOINTS) { + const full = makeTextMessage(REVEAL_CORPUS.slice(0, n)); + const total = visibleUnits(full, false); + const component = new AssistantMessageComponent(); + const start = Bun.nanoseconds(); + let revealed = 0; + let steps = 0; + while (revealed < total) { + revealed = Math.min(total, revealed + nextStep(total - revealed)); + component.updateContent(buildDisplayMessage(full, revealed, false)); + component.render(WIDTH); + steps++; + } + const ms = (Bun.nanoseconds() - start) / 1e6; + console.log(` len=${n}: ${ms.toFixed(2)}ms total over ${steps} steps (${(ms / steps).toFixed(4)}ms/step)`); + } +} catch (err) { + console.log(` (skipped: ${(err as Error).message})`); +} + +// Multi-block variant: a finalized thinking block (stable) precedes the growing +// text block — the shape C2 targets. Current code re-lexes BOTH every tick; +// after C2 the finalized thinking block stays L1-cached and only the tail re-lexes. +function makeThinkingPlusText(thinking: string, text: string): AssistantMessage { + return { ...makeTextMessage(text), content: [{ type: "thinking", thinking }, { type: "text", text }] }; +} +console.log("\nstreamingRevealMultiBlock (C2: finalized thinking block + growing text):"); +try { + const thinking = makeMarkdownCorpus(2500); + for (const n of [2000, 4000, 6000]) { + const full = makeThinkingPlusText(thinking, REVEAL_CORPUS.slice(0, n)); + const total = visibleUnits(full, false); + const component = new AssistantMessageComponent(); + const start = Bun.nanoseconds(); + let revealed = 0; + let steps = 0; + while (revealed < total) { + revealed = Math.min(total, revealed + nextStep(total - revealed)); + component.updateContent(buildDisplayMessage(full, revealed, false)); + component.render(WIDTH); + steps++; + } + const ms = (Bun.nanoseconds() - start) / 1e6; + console.log(` text=${n} (+2500 thinking): ${ms.toFixed(2)}ms total over ${steps} steps (${(ms / steps).toFixed(4)}ms/step)`); + } +} catch (err) { + console.log(` (skipped: ${(err as Error).message})`); +} + +console.log("\neditorKeystroke (C3: layout recompute vs no-mutation render):"); +try { + const buffer = Array.from({ length: 50 }) + .map((_, i) => `Line ${i + 1}: some editor content with words to wrap at width ${WIDTH} and more text here`) + .join("\n"); + + const e1 = new Editor(getEditorTheme()); + e1.setText(buffer); + e1.render(WIDTH); // warm + const noMutMs = benchStep(200, () => { + e1.render(WIDTH); + }); + console.log(` no-mutation render: ${noMutMs.toFixed(4)}ms/render`); + + const e2 = new Editor(getEditorTheme()); + e2.setText(buffer); + e2.render(WIDTH); // warm + const editMs = benchStep(200, () => { + e2.insertText("x"); + e2.render(WIDTH); + }); + console.log(` edit + render: ${editMs.toFixed(4)}ms/op`); +} catch (err) { + console.log(` (skipped: ${(err as Error).message})`); +} + +// ── E3: long-transcript frame cost ────────────────────────────────────────── +// +// E3 root cause: Container.render walks EVERY child and concatenates their line +// arrays on every frame. Finalized messages hit their Markdown L1 cache (no +// re-lex) but still pay the tree walk + line-array rebuild/concat per frame. +// Build N finalized assistant messages (prose + closed code fences) + 1 growing +// tail, then time one render(WIDTH) of the whole tree per streaming frame. +// Rising ms/frame in N => the stable history is re-walked/re-concatenated each +// frame (the cost E3 culls); flat => the walk is already cheap. +console.log("\nlongTranscriptFrame (E3: whole-tree render cost vs transcript length N):"); +try { + const histText = makeMarkdownCorpus(800); + const tailCorpus = makeMarkdownCorpus(1200); + for (const n of [50, 100, 200]) { + const container = new TranscriptContainer(); + for (let i = 0; i < n; i++) { + const c = new AssistantMessageComponent(); + c.updateContent(makeTextMessage(histText)); + container.addChild(c); + } + const tail = new AssistantMessageComponent(); + container.addChild(tail); + let revealed = Math.floor(tailCorpus.length * 0.5); + tail.updateContent(makeTextMessage(tailCorpus.slice(0, revealed))); + container.render(WIDTH); // warm finalized history (L1 caches hot) + const ms = benchStep(60, () => { + revealed += 20; + if (revealed > tailCorpus.length) revealed = Math.floor(tailCorpus.length * 0.5); + tail.updateContent(makeTextMessage(tailCorpus.slice(0, revealed))); + container.render(WIDTH); + }); + console.log(` N=${n}: ${ms.toFixed(4)}ms/frame`); + } +} catch (err) { + console.log(` (skipped: ${(err as Error).message})`); +} + +// ── E4: tool read/parse redundancy ────────────────────────────────────────── +// +// E4 root cause: the read tool re-parses (tree-sitter `summarizeCode`, ~12-18ms +// for a ~1500-line file) on every summary read of the same unchanged file. E4-ii +// memoizes the parse per session keyed on the content hash of the freshly-read +// bytes, so a repeat read of the same file reuses the parse (the file is still +// read fresh, so the result stays correct). A repeated same-session summary read +// should drop from ~17ms to a few ms; a fresh session each call stays full cost. +console.log("\ntoolReadReparse (E4: repeat summary read, memoized parse vs cold):"); +try { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "bench-e4-")); + const file = path.join(dir, "big.ts"); + let src = ""; + for (let i = 0; i < 375; i++) { + src += `export function fn${i}(a: number, b: string): boolean {\n const x = a + ${i};\n return x > 0 && b.length === ${i};\n}\n`; + } + fs.writeFileSync(file, src); + const mkSession = (): ToolSession => + ({ + cwd: dir, + hasUI: false, + getSessionFile: () => path.join(dir, "s.jsonl"), + getSessionSpawns: () => "*", + getArtifactsDir: () => path.join(dir, "sess"), + allocateOutputArtifact: async (t: string) => ({ id: "a", path: path.join(dir, `a.${t}.log`) }), + settings: Settings.isolated(), + }) as unknown as ToolSession; + const sameSession = mkSession(); + const rt = new ReadTool(sameSession); + for (let i = 0; i < 3; i++) await rt.execute("warm", { path: file }); + const repeatMs = await benchStepAsync(20, () => rt.execute("c", { path: file })); + const coldMs = await benchStepAsync(20, () => new ReadTool(mkSession()).execute("c", { path: file })); + console.log(` same-session repeat read: ${repeatMs.toFixed(3)}ms/call (memoized parse)`); + console.log(` fresh-session each read: ${coldMs.toFixed(3)}ms/call (cold parse)`); + fs.rmSync(dir, { recursive: true, force: true }); +} catch (err) { + console.log(` (skipped: ${(err as Error).message})`); +} diff --git a/packages/coding-agent/package.json b/packages/coding-agent/package.json index aa7551664..009694006 100644 --- a/packages/coding-agent/package.json +++ b/packages/coding-agent/package.json @@ -41,7 +41,8 @@ "format-prompts": "bun scripts/format-prompts.ts", "generate-docs-index": "bun scripts/generate-docs-index.ts", "prepack": "bun scripts/generate-docs-index.ts && bun scripts/bundle-dist.ts", - "generate-template": "bun scripts/generate-template.ts" + "generate-template": "bun scripts/generate-template.ts", + "bench:guard": "bun scripts/bench-guard.ts" }, "dependencies": { "@agentclientprotocol/sdk": "catalog:", diff --git a/packages/coding-agent/scripts/bench-guard.ts b/packages/coding-agent/scripts/bench-guard.ts new file mode 100755 index 000000000..2ce365d9f --- /dev/null +++ b/packages/coding-agent/scripts/bench-guard.ts @@ -0,0 +1,71 @@ +#!/usr/bin/env bun +/** + * Boot-time regression guard (Phase A1 of the boot/TUI perf work). + * + * Re-runs the `PI_TIMING=x` cold-boot benchmark under hyperfine and fails when + * the median regresses past `baseline * THRESHOLD`. `PI_TIMING=x` runs the full + * pre-paint chain in `runRootCommand` and then `process.exit(0)`, so the + * never-exiting interactive launch becomes a terminating, benchmarkable boot. + * + * Boot wall-clock is MACHINE-RELATIVE: a baseline captured on one machine is + * meaningless on another (and on CI). This is a LOCAL guard — regenerate the + * baseline on the machine you measure on, then compare on that same machine. + * It is intentionally NOT wired into CI for that reason. + * + * bun scripts/bench-guard.ts --update # capture/refresh the baseline + * bun scripts/bench-guard.ts # measure + compare; exit 1 on regression + * + * Requires `hyperfine` on PATH. + */ +import * as fs from "node:fs"; +import * as path from "node:path"; + +const THRESHOLD = 1.05; // 5% regression budget +const BASELINE_PATH = path.join(import.meta.dir, "..", "bench", "boot-baseline.json"); +const BENCH_COMMAND = "PI_TIMING=x bun src/cli.ts"; +const cwd = path.join(import.meta.dir, ".."); + +function medianOf(hyperfineJson: string): number { + const parsed = JSON.parse(hyperfineJson) as { results: Array<{ mean: number; median?: number }> }; + const result = parsed.results[0]; + if (!result) throw new Error("hyperfine produced no result"); + return result.median ?? result.mean; +} + +async function measure(): Promise<{ seconds: number; raw: string }> { + const tmp = path.join(import.meta.dir, "..", "bench", `.boot-run-${Date.now()}.json`); + const proc = Bun.spawn(["hyperfine", "--warmup", "3", "--min-runs", "10", "--export-json", tmp, BENCH_COMMAND], { + cwd, + stdout: "inherit", + stderr: "inherit", + }); + const code = await proc.exited; + if (code !== 0) throw new Error(`hyperfine exited ${code}`); + const raw = await Bun.file(tmp).text(); + fs.rmSync(tmp, { force: true }); + return { seconds: medianOf(raw), raw }; +} + +const update = process.argv.includes("--update"); +const { seconds, raw } = await measure(); + +if (update) { + fs.mkdirSync(path.dirname(BASELINE_PATH), { recursive: true }); + await Bun.write(BASELINE_PATH, raw); + console.log(`Baseline updated: ${(seconds * 1000).toFixed(0)}ms median -> ${BASELINE_PATH}`); + process.exit(0); +} + +if (!fs.existsSync(BASELINE_PATH)) { + console.error("No baseline found. Run `bun scripts/bench-guard.ts --update` on this machine first."); + process.exit(2); +} + +const baseline = medianOf(await Bun.file(BASELINE_PATH).text()); +const ratio = seconds / baseline; +const verdict = ratio > THRESHOLD ? "REGRESSION" : "ok"; +console.log( + `boot median: ${(seconds * 1000).toFixed(0)}ms vs baseline ${(baseline * 1000).toFixed(0)}ms ` + + `(${((ratio - 1) * 100).toFixed(1)}%, budget ${((THRESHOLD - 1) * 100).toFixed(0)}%) -> ${verdict}`, +); +process.exit(ratio > THRESHOLD ? 1 : 0); diff --git a/packages/coding-agent/src/config/model-registry.ts b/packages/coding-agent/src/config/model-registry.ts index 55acd053d..ac0804753 100644 --- a/packages/coding-agent/src/config/model-registry.ts +++ b/packages/coding-agent/src/config/model-registry.ts @@ -602,6 +602,7 @@ function getConfiguredProviderOrderFromSettings(): string[] { export class ModelRegistry { #models: Model[] = []; #canonicalIndex: CanonicalModelIndex = { records: [], byId: new Map(), bySelector: new Map() }; + #canonicalIndexDirty: boolean = true; #customProviderApiKeys: Map = new Map(); #keylessProviders: Set = new Set(); #discoverableProviders: DiscoveryProviderConfig[] = []; @@ -1519,14 +1520,25 @@ export class ModelRegistry { this.#rebuildPending = true; return; } - this.#canonicalIndex = buildCanonicalModelIndex( - this.#models, - getBundledCanonicalReferenceData(), - this.#equivalenceConfig, - ); + // Defer the catalog-wide index build to first read. Boot model + // resolution reads it only when enabledModels or a default-role pattern + // is configured; the empty interactive launch never reads it pre-paint, + // so the ~200ms build over the full catalog moves off the first-paint + // critical path. + this.#canonicalIndexDirty = true; this.#rebuildPending = false; } + #ensureCanonicalIndex(): CanonicalModelIndex { + if (this.#canonicalIndexDirty) { + this.#canonicalIndex = logger.time("buildCanonicalModelIndex", () => + buildCanonicalModelIndex(this.#models, getBundledCanonicalReferenceData(), this.#equivalenceConfig), + ); + this.#canonicalIndexDirty = false; + } + return this.#canonicalIndex; + } + #suspendRebuild(): void { this.#rebuildSuspended += 1; } @@ -1537,11 +1549,7 @@ export class ModelRegistry { } if (this.#rebuildSuspended === 0 && this.#rebuildPending) { this.#rebuildPending = false; - this.#canonicalIndex = buildCanonicalModelIndex( - this.#models, - getBundledCanonicalReferenceData(), - this.#equivalenceConfig, - ); + this.#canonicalIndexDirty = true; } } @@ -1650,7 +1658,7 @@ export class ModelRegistry { getCanonicalModels(options?: CanonicalModelQueryOptions): CanonicalModelRecord[] { const { candidateKeys, isAvailable } = this.#canonicalQueryFilters(options); const records: CanonicalModelRecord[] = []; - for (const record of this.#canonicalIndex.records) { + for (const record of this.#ensureCanonicalIndex().records) { const variants = this.#filterCanonicalVariants(record, candidateKeys, isAvailable); if (variants.length === 0) { continue; @@ -1676,7 +1684,7 @@ export class ModelRegistry { const candidates = options?.candidates ?? (options?.availableOnly ? this.getAvailable() : this.getAll()); const preferences = this.#variantPreferences(candidates); const selections: CanonicalModelSelection[] = []; - for (const record of this.#canonicalIndex.records) { + for (const record of this.#ensureCanonicalIndex().records) { const variants = this.#filterCanonicalVariants(record, candidateKeys, isAvailable); if (variants.length === 0) { continue; @@ -1694,7 +1702,7 @@ export class ModelRegistry { } getCanonicalVariants(canonicalId: string, options?: CanonicalModelQueryOptions): CanonicalModelVariant[] { - const record = this.#canonicalIndex.byId.get(canonicalId.trim().toLowerCase()); + const record = this.#ensureCanonicalIndex().byId.get(canonicalId.trim().toLowerCase()); if (!record) { return []; } @@ -1712,7 +1720,7 @@ export class ModelRegistry { } getCanonicalId(model: Model): string | undefined { - return this.#canonicalIndex.bySelector.get(formatCanonicalVariantSelector(model).toLowerCase()); + return this.#ensureCanonicalIndex().bySelector.get(formatCanonicalVariantSelector(model).toLowerCase()); } /** diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 7025a9317..2c10b2656 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -897,7 +897,7 @@ export async function runRootCommand( // Create AuthStorage and ModelRegistry upfront const authStorage = await logger.time("discoverAuthStorage", deps.discoverAuthStorage ?? discoverAuthStorage); - const modelRegistry = new ModelRegistry(authStorage); + const modelRegistry = logger.time("modelRegistry:init", () => new ModelRegistry(authStorage)); if (parsedArgs.version) { process.stdout.write(`${VERSION}\n`); @@ -1146,7 +1146,7 @@ export async function runRootCommand( // Both are no-ops when OTEL_EXPORTER_OTLP_ENDPOINT is unset. An empty config // is enough to enable telemetry — content capture is governed by the // standard OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT env var. - await initTelemetryExport(); + await logger.time("initTelemetryExport", initTelemetryExport); if (isTelemetryExportEnabled()) { sessionOptions.telemetry = {}; } diff --git a/packages/coding-agent/src/modes/components/assistant-message.ts b/packages/coding-agent/src/modes/components/assistant-message.ts index 2896f9086..dacca0903 100644 --- a/packages/coding-agent/src/modes/components/assistant-message.ts +++ b/packages/coding-agent/src/modes/components/assistant-message.ts @@ -49,6 +49,11 @@ export class AssistantMessageComponent extends Container { /** Whether the last updateContent carried an in-flight streaming partial; such * renders bypass the markdown module LRU (see Markdown.transientRenderCache). */ #lastUpdateTransient = false; + // Fast-path state: reuse Markdown children when message shape is stable during streaming. + #fastPathKey: string | undefined; + #fastPathItems: + | Array<{ md: Markdown; contentIndex: number; blockType: "text" | "thinking"; lastText: string }> + | undefined; constructor( message?: AssistantMessage, @@ -71,6 +76,12 @@ export class AssistantMessageComponent extends Container { override invalidate(): void { super.invalidate(); + // Theme/symbol changes arrive via invalidate(). Fast-path children captured + // getMarkdownTheme() at construction, so drop them and force the teardown + // path to rebuild with the current theme. Streaming updates call + // updateContent() directly and keep the fast path. + this.#fastPathKey = undefined; + this.#fastPathItems = undefined; if (this.#lastMessage) { this.updateContent(this.#lastMessage, { transient: this.#lastUpdateTransient }); } @@ -228,14 +239,111 @@ export class AssistantMessageComponent extends Container { } } + #computeShapeKey(message: AssistantMessage): string { + const parts: string[] = [`htb:${this.hideThinkingBlock ? 1 : 0}`]; + for (const content of message.content) { + if (content.type === "text") { + parts.push(content.text.trim() ? "T1" : "T0"); + } else if (content.type === "thinking") { + if (!content.thinking.trim()) parts.push("K0"); + else if (this.hideThinkingBlock) parts.push("KH"); + else parts.push("KV"); + } else { + // Non-rendered blocks (toolCall, redactedThinking, …) still occupy a + // content index. Encode their position so an inserted/removed one shifts + // the key and forces the teardown path instead of mis-indexing children. + parts.push(`O:${content.type}`); + } + } + if (settings.get("display.showTokenUsage") && this.#usageInfo) { + const u = this.#usageInfo; + parts.push(`u:${u.input + u.cacheWrite}:${u.output}:${u.cacheRead}`); + } else { + parts.push("u:"); + } + return parts.join("|"); + } + + #canFastPath(message: AssistantMessage): boolean { + for (const content of message.content) { + if (content.type === "toolCall") return false; + } + if (this.#toolImagesByCallId.size > 0) return false; + if (message.stopReason === "aborted" && shouldRenderAbortReason(message.errorMessage)) return false; + if (message.stopReason === "error" && !this.#errorPinned) return false; + if ( + message.errorMessage && + shouldRenderAbortReason(message.errorMessage) && + message.stopReason !== "aborted" && + message.stopReason !== "error" + ) + return false; + // Extension stability: if thinking renderers exist and any tracked thinking + // block's text changed, extensions may produce a different child count. + if (this.thinkingRenderers.length > 0 && this.#fastPathItems) { + for (const item of this.#fastPathItems) { + if (item.blockType === "thinking") { + const content = message.content[item.contentIndex]; + if (content?.type === "thinking" && content.thinking.trim() !== item.lastText) return false; + } + } + } + return true; + } + + #tryFastPathUpdate(message: AssistantMessage, opts?: { transient?: boolean }): boolean { + if (!this.#fastPathKey || !this.#fastPathItems) return false; + if (!this.#canFastPath(message)) { + this.#fastPathKey = undefined; + this.#fastPathItems = undefined; + return false; + } + if (this.#computeShapeKey(message) !== this.#fastPathKey) { + this.#fastPathKey = undefined; + this.#fastPathItems = undefined; + return false; + } + const transient = opts?.transient === true; + // Shape is identical — setText only on Markdown children whose source changed. + for (const item of this.#fastPathItems) { + item.md.transientRenderCache = transient; + const content = message.content[item.contentIndex]; + let newText: string; + if (item.blockType === "text" && content?.type === "text") { + newText = content.text.trim(); + } else if (item.blockType === "thinking" && content?.type === "thinking") { + newText = content.thinking.trim(); + } else { + // Block at this index is gone or changed type (index shift) — fail closed. + this.#fastPathKey = undefined; + this.#fastPathItems = undefined; + return false; + } + if (newText !== item.lastText) { + item.md.setText(newText); + item.lastText = newText; + } + } + return true; + } + updateContent(message: AssistantMessage, opts?: { transient?: boolean }): void { this.#blockVersion++; this.#lastMessage = message; this.#lastUpdateTransient = opts?.transient === true; + // Fast path: reuse Markdown children when shape is stable during streaming + if (this.#tryFastPathUpdate(message)) return; + // Clear content container this.#contentContainer.clear(); + // Determine if we should capture Markdown instances for next fast path + const shouldCapture = this.#canFastPath(message); + const captureItems: + | Array<{ md: Markdown; contentIndex: number; blockType: "text" | "thinking"; lastText: string }> + | undefined = shouldCapture ? [] : undefined; + const hasVisibleContent = message.content.some( c => (c.type === "text" && c.text.trim()) || @@ -249,9 +357,11 @@ export class AssistantMessageComponent extends Container { if (content.type === "text" && content.text.trim()) { // Assistant text messages with no background - trim the text // Set paddingY=0 to avoid extra spacing before tool executions - const markdown = new Markdown(content.text.trim(), 1, 0, getMarkdownTheme()); - markdown.transientRenderCache = this.#lastUpdateTransient; - this.#contentContainer.addChild(markdown); + const trimmed = content.text.trim(); + const md = new Markdown(trimmed, 1, 0, getMarkdownTheme()); + md.transientRenderCache = this.#lastUpdateTransient; + this.#contentContainer.addChild(md); + captureItems?.push({ md, contentIndex: i, blockType: "text", lastText: trimmed }); } else if (content.type === "thinking" && content.thinking.trim()) { if (this.hideThinkingBlock) { thinkingIndex += 1; @@ -265,12 +375,13 @@ export class AssistantMessageComponent extends Container { const thinkingText = content.thinking.trim(); // Thinking traces in thinkingText color, italic - const thinkingMarkdown = new Markdown(thinkingText, 1, 0, getMarkdownTheme(), { + const md = new Markdown(thinkingText, 1, 0, getMarkdownTheme(), { color: (text: string) => theme.fg("thinkingText", text), italic: true, }); - thinkingMarkdown.transientRenderCache = this.#lastUpdateTransient; - this.#contentContainer.addChild(thinkingMarkdown); + md.transientRenderCache = this.#lastUpdateTransient; + this.#contentContainer.addChild(md); + captureItems?.push({ md, contentIndex: i, blockType: "thinking", lastText: thinkingText }); this.#appendThinkingExtensions(i, thinkingIndex, thinkingText); thinkingIndex += 1; if (hasVisibleContentAfter) { @@ -318,5 +429,14 @@ export class AssistantMessageComponent extends Container { this.#contentContainer.addChild(new Spacer(1)); this.#contentContainer.addChild(new Text(theme.fg("dim", parts.join(" ")), 1, 0)); } + + // Store fast-path state for next call + if (shouldCapture) { + this.#fastPathItems = captureItems; + this.#fastPathKey = this.#computeShapeKey(message); + } else { + this.#fastPathKey = undefined; + this.#fastPathItems = undefined; + } } } diff --git a/packages/coding-agent/src/modes/controllers/streaming-reveal.ts b/packages/coding-agent/src/modes/controllers/streaming-reveal.ts index a26024e38..cefb997cb 100644 --- a/packages/coding-agent/src/modes/controllers/streaming-reveal.ts +++ b/packages/coding-agent/src/modes/controllers/streaming-reveal.ts @@ -1,5 +1,6 @@ import type { AssistantMessage } from "@oh-my-pi/pi-ai"; import { getSegmenter } from "@oh-my-pi/pi-tui"; +import { LRUCache } from "lru-cache/raw"; import type { AssistantMessageComponent } from "../components/assistant-message"; export const STREAMING_REVEAL_FRAME_MS = 1000 / 30; @@ -15,11 +16,17 @@ type StreamingRevealControllerOptions = { requestRender(): void; }; +const graphemeCountCache = new LRUCache({ max: 128 }); + function countGraphemes(text: string): number { + if (text.length === 0) return 0; + const cached = graphemeCountCache.get(text); + if (cached !== undefined) return cached; let count = 0; for (const _segment of getSegmenter().segment(text)) { count += 1; } + graphemeCountCache.set(text, count); return count; } diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 14edd8386..1e2bfd46c 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -89,7 +89,14 @@ import { type FileSlashCommand, loadSlashCommands as loadSlashCommandsInternal } import type { HindsightSessionState } from "./hindsight/state"; import { LocalProtocolHandler, type LocalProtocolOptions } from "./internal-urls"; import { LSP_STARTUP_EVENT_CHANNEL, type LspStartupEvent } from "./lsp/startup-events"; -import { discoverAndLoadMCPTools, MCPManager, type MCPToolsLoadResult } from "./mcp"; +import { + discoverAndLoadMCPTools, + type MCPLoadResult, + MCPManager, + MCPToolCache, + type MCPToolsLoadResult, + parseMCPToolName, +} from "./mcp"; import { createSessionMemoryRuntimeContext, resolveMemoryBackend } from "./memory-backend"; import type { MnemopiSessionState } from "./mnemopi/state"; import asyncResultTemplate from "./prompts/tools/async-result.md" with { type: "text" }; @@ -147,6 +154,7 @@ import { type DiscoverableTool, filterBySource, formatDiscoverableToolServerSummary, + isMCPToolName, selectDiscoverableToolNamesByServer, summarizeDiscoverableTools, } from "./tool-discovery/tool-index"; @@ -302,6 +310,76 @@ function buildMcpNotificationBatchMessage(entries: McpNotificationEntry[]): Agen }; } +type DeferredMCPActivation = { + mcpDiscoveryEnabled: boolean; + explicitlyRequestedMCPToolNames: string[]; + activateAllMCPTools: boolean; +}; + +function formatMCPConnectingMessage(serverNames: string[]): string { + return `Connecting to MCP servers: ${serverNames.join(", ")}…`; +} + +function createPendingMCPTool(name: string): Tool { + const parsed = parseMCPToolName(name); + const serverName = parsed?.serverName; + const mcpToolName = parsed?.toolName ?? name; + const label = serverName ? `${serverName}/${mcpToolName}` : name; + const message = serverName + ? `MCP server "${serverName}" is still connecting; tool "${name}" is not yet available. Retry after the MCP connection completes.` + : `MCP discovery is still in progress; tool "${name}" is not yet available. Retry after MCP connection completes.`; + const tool: Tool & { mcpServerName?: string; mcpToolName?: string } = { + name, + label, + description: `Pending MCP tool. ${message}`, + parameters: { + type: "object", + properties: {}, + additionalProperties: true, + }, + approval: "write", + intent: "omit", + mcpServerName: serverName, + mcpToolName, + async execute() { + return { + content: [{ type: "text", text: message }], + details: { serverName, mcpToolName, isError: true }, + isError: true, + }; + }, + }; + return tool; +} + +function collectPendingMCPToolNames( + explicitToolNames: readonly string[] | undefined, + restoredSelectedToolNames: readonly string[], +): string[] { + const names = new Set(); + for (const name of explicitToolNames ?? []) { + const normalized = name.toLowerCase(); + if (isMCPToolName(normalized)) names.add(normalized); + } + for (const name of restoredSelectedToolNames) { + const normalized = name.toLowerCase(); + if (isMCPToolName(normalized)) names.add(normalized); + } + return [...names]; +} + +function logMCPLoadErrors(errors: MCPLoadResult["errors"]): void { + for (const [serverName, error] of errors) { + logger.error("MCP tool load failed", { path: `mcp:${serverName}`, error }); + } +} + +function applyMCPEnvironment(result: { exaApiKeys: string[] }): void { + if (result.exaApiKeys.length > 0 && !$env.EXA_API_KEY) { + Bun.env.EXA_API_KEY = result.exaApiKeys[0]; + } +} + // Types export interface CreateAgentSessionOptions { /** Working directory for project-local discovery. Default: getProjectDir() */ @@ -1518,46 +1596,131 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} let mcpManager: MCPManager | undefined = options.mcpManager; toolSession.mcpManager = mcpManager; const enableMCP = options.enableMCP ?? true; + const deferMCPDiscoveryForUI = enableMCP && !mcpManager && options.hasUI === true; const customTools: CustomTool[] = []; + let startDeferredMCPDiscovery: + | ((liveSession: AgentSession, activation: DeferredMCPActivation) => void) + | undefined; + const onMCPConnecting = (serverNames: string[]) => { + if (!options.hasUI || serverNames.length === 0) return; + process.stderr.write(`${chalk.gray(formatMCPConnectingMessage(serverNames))}\n`); + }; + const mcpDiscoverOptions = { + onConnecting: onMCPConnecting, + enableProjectConfig: settings.get("mcp.enableProjectConfig") ?? true, + // Always filter Exa - we have native integration + filterExa: true, + // Filter browser MCP servers when builtin browser tool is active + filterBrowser: settings.get("browser.enabled") ?? false, + }; if (enableMCP && !mcpManager) { - const mcpResult = await logger.time("discoverAndLoadMCPTools", discoverAndLoadMCPTools, cwd, { - onConnecting: serverNames => { - if (options.hasUI && serverNames.length > 0) { - process.stderr.write(`${chalk.gray(`Connecting to MCP servers: ${serverNames.join(", ")}…`)}\n`); - } - }, - enableProjectConfig: settings.get("mcp.enableProjectConfig") ?? true, - // Always filter Exa - we have native integration - filterExa: true, - // Filter browser MCP servers when builtin browser tool is active - filterBrowser: settings.get("browser.enabled") ?? false, - cacheStorage: settings.getStorage(), - authStorage, - }); - mcpManager = mcpResult.manager; - toolSession.mcpManager = mcpManager; + if (deferMCPDiscoveryForUI) { + const cacheStorage = settings.getStorage(); + mcpManager = new MCPManager(cwd, cacheStorage ? new MCPToolCache(cacheStorage) : null); + mcpManager.setAuthStorage(authStorage); + toolSession.mcpManager = mcpManager; - if (settings.get("mcp.notifications")) { - mcpManager.setNotificationsEnabled(true); - } - // If we extracted Exa API keys from MCP configs and EXA_API_KEY isn't set, use the first one - if (mcpResult.exaApiKeys.length > 0 && !$env.EXA_API_KEY) { - Bun.env.EXA_API_KEY = mcpResult.exaApiKeys[0]; - } + if (settings.get("mcp.notifications")) { + mcpManager.setNotificationsEnabled(true); + } - // Log MCP errors - for (const { path, error } of mcpResult.errors) { - logger.error("MCP tool load failed", { path, error }); - } + const deferredMCPManager = mcpManager; + startDeferredMCPDiscovery = (liveSession, activation) => { + void (async () => { + try { + const mcpResult = await logger.time("discoverAndLoadMCPTools", () => + deferredMCPManager.discoverAndConnect(mcpDiscoverOptions), + ); + // The session can be torn down while servers are still connecting. + // Don't resurrect tools on a disposed session, and don't leak the + // transports/subprocesses the connect just spawned. + if (liveSession.isDisposed) { + await deferredMCPManager.disconnectAll(); + return; + } + applyMCPEnvironment(mcpResult); + logMCPLoadErrors(mcpResult.errors); + // `tools.discoveryMode: "auto"` was resolved against a registry that + // held only built-ins plus persisted placeholder names. Recompute with + // the real MCP tool count: a large toolset must flip discovery on + // BEFORE the refresh, or activateAll would dump every MCP tool into + // the active set with no search_tool_bm25 registered. + let discoveryEnabled = activation.mcpDiscoveryEnabled; + let activateAll = activation.activateAllMCPTools; + if (!discoveryEnabled) { + const nonMCPToolNames = [...toolRegistry.keys()].filter(name => !isMCPToolName(name)); + const projectedMode = resolveEffectiveToolDiscoveryMode( + settings, + countToolsForAutoDiscovery([...nonMCPToolNames, ...mcpResult.tools.map(tool => tool.name)]), + ); + if (projectedMode !== "off") { + effectiveDiscoveryMode = projectedMode; + mcpDiscoveryEnabled = true; + discoveryEnabled = true; + activateAll = false; + liveSession.enableMCPDiscovery(); + if (!toolRegistry.has("search_tool_bm25")) { + const searchTool: Tool = new SearchToolBm25Tool(toolSession); + toolRegistry.set( + searchTool.name, + new ExtensionToolWrapper(wrapToolWithMetaNotice(searchTool), extensionRunner) as Tool, + ); + } + await liveSession.setActiveToolsByName([ + ...liveSession.getActiveToolNames(), + "search_tool_bm25", + ]); + } + } + await liveSession.refreshMCPTools(mcpResult.tools, { activateAll }); + if (activation.explicitlyRequestedMCPToolNames.length > 0) { + if (discoveryEnabled && !activation.mcpDiscoveryEnabled) { + // Discovery flipped on mid-flight: route the explicit request + // through discovery-aware activation so selection persists. + await liveSession.activateDiscoveredMCPTools(activation.explicitlyRequestedMCPToolNames); + } else if (!discoveryEnabled) { + await liveSession.setActiveToolsByName([ + ...liveSession.getActiveToolNames(), + ...activation.explicitlyRequestedMCPToolNames, + ]); + } + } + } catch (error) { + logger.error("MCP tool load failed", { + path: ".mcp.json", + error: error instanceof Error ? error.message : String(error), + }); + } + })(); + }; + } else { + const mcpResult = await logger.time("discoverAndLoadMCPTools", discoverAndLoadMCPTools, cwd, { + ...mcpDiscoverOptions, + cacheStorage: settings.getStorage(), + authStorage, + }); + mcpManager = mcpResult.manager; + toolSession.mcpManager = mcpManager; - if (mcpResult.tools.length > 0) { - // MCP tools are LoadedCustomTool, extract the tool property - customTools.push(...mcpResult.tools.map(loaded => loaded.tool)); + if (settings.get("mcp.notifications")) { + mcpManager.setNotificationsEnabled(true); + } + applyMCPEnvironment(mcpResult); + + // Log MCP errors + for (const { path, error } of mcpResult.errors) { + logger.error("MCP tool load failed", { path, error }); + } + + if (mcpResult.tools.length > 0) { + // MCP tools are LoadedCustomTool, extract the tool property + customTools.push(...mcpResult.tools.map(loaded => loaded.tool)); + } } } // Only top-level sessions own the global MCPManager. Subagents already // receive the parent's manager via `options.mcpManager`, and reassigning - // the singleton to the same value is a no-op \u2014 keep the gate explicit + // the singleton to the same value is a no-op — keep the gate explicit // to mirror the AsyncJobManager ownership rule. if (mcpManager && !options.parentTaskPrefix) MCPManager.setInstance(mcpManager); @@ -1852,6 +2015,14 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} for (const tool of wrappedExtensionTools) { toolRegistry.set(tool.name, tool); } + if (deferMCPDiscoveryForUI && mcpManager) { + for (const name of collectPendingMCPToolNames(options.toolNames, existingSession.selectedMCPToolNames)) { + if (!toolRegistry.has(name)) { + toolRegistry.set(name, createPendingMCPTool(name)); + } + } + } + // Wrap every tool with `ExtensionToolWrapper` so the per-tool approval gate runs on every // call site, regardless of whether any user extensions are loaded. See the runner-construction // comment above for the safety invariant this enforces. @@ -1879,7 +2050,10 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} } } - const effectiveDiscoveryMode = resolveEffectiveToolDiscoveryMode( + // `let`: the deferred MCP discovery closure upgrades these when the real + // MCP tool count pushes `auto` past its threshold; `rebuildSystemPrompt` + // below reads the live bindings. + let effectiveDiscoveryMode = resolveEffectiveToolDiscoveryMode( settings, countToolsForAutoDiscovery(toolRegistry.keys()), ); @@ -1890,7 +2064,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} new ExtensionToolWrapper(wrapToolWithMetaNotice(searchTool), extensionRunner) as Tool, ); } - const mcpDiscoveryEnabled = effectiveDiscoveryMode !== "off"; // back-compat: true when any discovery active + let mcpDiscoveryEnabled = effectiveDiscoveryMode !== "off"; // back-compat: true when any discovery active const reloadSshTool = async (): Promise => { if (!requestedToolNameSet.has("ssh")) return null; @@ -1942,7 +2116,11 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} const memoryBackend = await resolveMemoryBackend(settings); const memoryInstructions = await memoryBackend.buildDeveloperInstructions(agentDir, settings, session); - // Build combined append prompt: memory instructions + MCP server instructions + // Build combined append prompt: memory instructions + MCP server instructions. + // For UI sessions MCP discovery is deferred, so `getServerInstructions()` is + // empty until the background connect completes; the rebuild that + // `refreshMCPTools` triggers post-discovery then picks up the now-connected + // servers' instructions, so they join the prompt for the rest of the session. const serverInstructions = mcpManager?.getServerInstructions(); let appendPrompt: string | undefined = memoryInstructions ?? undefined; if (serverInstructions && serverInstructions.size > 0) { @@ -2494,7 +2672,26 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} // Skip when reusing a parent's manager — the parent owns the callbacks. if (mcpManager && !options.mcpManager) { mcpManager.setOnToolsChanged(tools => { - void session.refreshMCPTools(tools); + void (async () => { + try { + await session.refreshMCPTools( + tools, + deferMCPDiscoveryForUI && !mcpDiscoveryEnabled && options.toolNames === undefined + ? { activateAll: true } + : undefined, + ); + if (deferMCPDiscoveryForUI && !mcpDiscoveryEnabled && explicitlyRequestedMCPToolNames.length > 0) { + await session.setActiveToolsByName([ + ...session.getActiveToolNames(), + ...explicitlyRequestedMCPToolNames, + ]); + } + } catch (error) { + logger.warn("MCP tool refresh failed", { + error: error instanceof Error ? error.message : String(error), + }); + } + })(); }); // Wire prompt refresh → rebuild MCP prompt slash commands mcpManager.setOnPromptsChanged(serverName => { @@ -2527,6 +2724,12 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} }); } + startDeferredMCPDiscovery?.(session, { + mcpDiscoveryEnabled, + explicitlyRequestedMCPToolNames, + activateAllMCPTools: !mcpDiscoveryEnabled && options.toolNames === undefined, + }); + return { session, extensionsResult, diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 58de3025f..e3d13e083 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -3140,6 +3140,12 @@ export class AgentSession { state.resetConversationTracking(); } + /** True once dispose() has begun; deferred background work (e.g. the deferred + * MCP discovery task in sdk.ts) must not touch the session past this point. */ + get isDisposed(): boolean { + return this.#isDisposed; + } + /** * Synchronously mark the session as disposing so new work is rejected * immediately: Python/eval starts throw, queued asides are dropped, and the @@ -3473,6 +3479,17 @@ export class AgentSession { return this.#mcpDiscoveryEnabled; } + /** + * Flip MCP discovery on after deferred discovery learns the real tool count. + * UI sessions resolve `tools.discoveryMode: "auto"` before MCP servers + * connect, so a large MCP toolset discovered later must be able to upgrade + * the session from the force-activate path to the discovery path. One-way: + * discovery is never downgraded mid-session. + */ + enableMCPDiscovery(): void { + this.#mcpDiscoveryEnabled = true; + } + getSelectedMCPToolNames(): string[] { if (!this.#mcpDiscoveryEnabled) { return this.getActiveToolNames().filter(name => isMCPToolName(name) && this.#toolRegistry.has(name)); diff --git a/packages/coding-agent/src/tools/read.ts b/packages/coding-agent/src/tools/read.ts index 285dbd09b..24bef6fe0 100644 --- a/packages/coding-agent/src/tools/read.ts +++ b/packages/coding-agent/src/tools/read.ts @@ -8,6 +8,7 @@ import { glob, type SummaryResult, summarizeCode } from "@oh-my-pi/pi-natives"; import type { Component } from "@oh-my-pi/pi-tui"; import { Text } from "@oh-my-pi/pi-tui"; import { getRemoteDir, logger, prompt, readImageMetadata, untilAborted } from "@oh-my-pi/pi-utils"; +import { LRUCache } from "lru-cache/raw"; import * as z from "zod/v4"; import { canonicalSnapshotKey, @@ -100,6 +101,28 @@ import { import { ToolAbortError, ToolError, throwIfAborted } from "./tool-errors"; import { toolResult } from "./tool-result"; +// Per-session memo for tree-sitter summaries. `summarizeCode` is a pure function +// of (code, path, fold settings) but costs ~12-18ms for a ~1500-line file, and a +// repeat summary read of the same unchanged file re-parses from scratch. Key on +// the content hash of the freshly-read bytes (+ path + fold settings): the file +// is still read fresh on every call, so a hit only reuses the deterministic +// parse — there is no staleness window and no stat guard is needed. Bounded LRU, +// aged out with the session via WeakMap. +// Unusable results (not parsed, or nothing elided) are memoized as `false`: the +// full SummaryResult embeds the whole source in kept segments, and the caller +// only ever renders `parsed && elided` summaries — caching the segments would +// retain up to 48 near-2MiB sources just to remember "no summary". +const SUMMARY_CACHE_MAX = 48; +const summaryParseCaches = new WeakMap>(); +function getSummaryParseCache(session: object): LRUCache { + let cache = summaryParseCaches.get(session); + if (!cache) { + cache = new LRUCache({ max: SUMMARY_CACHE_MAX }); + summaryParseCaches.set(session, cache); + } + return cache; +} + // Document types converted to markdown via markit. const CONVERTIBLE_EXTENSIONS = new Set([".pdf", ".doc", ".docx", ".ppt", ".pptx", ".xls", ".xlsx", ".rtf", ".epub"]); @@ -1614,15 +1637,25 @@ export class ReadTool implements AgentTool { if (lineCount > MAX_SUMMARY_LINES) return null; if (lineCount < this.session.settings.get("read.summarize.minTotalLines")) return null; + const minBodyLines = this.session.settings.get("read.summarize.minBodyLines"); + const minCommentLines = this.session.settings.get("read.summarize.minCommentLines"); + const unfoldUntilLines = this.session.settings.get("read.summarize.unfoldUntil"); + const unfoldLimitLines = this.session.settings.get("read.summarize.unfoldLimit"); + const cache = getSummaryParseCache(this.session); + const cacheKey = `${absolutePath}\0${Bun.hash(code)}\0${minBodyLines},${minCommentLines},${unfoldUntilLines},${unfoldLimitLines}`; + const memoized = cache.get(cacheKey); + if (memoized !== undefined) return memoized || null; const result = summarizeCode({ code, path: absolutePath, - minBodyLines: this.session.settings.get("read.summarize.minBodyLines"), - minCommentLines: this.session.settings.get("read.summarize.minCommentLines"), - unfoldUntilLines: this.session.settings.get("read.summarize.unfoldUntil"), - unfoldLimitLines: this.session.settings.get("read.summarize.unfoldLimit"), + minBodyLines, + minCommentLines, + unfoldUntilLines, + unfoldLimitLines, }); - return result; + const usable = result.parsed && result.elided ? result : false; + cache.set(cacheKey, usable); + return usable || null; } catch { return null; } diff --git a/packages/coding-agent/test/fixtures/instructions-mcp.ts b/packages/coding-agent/test/fixtures/instructions-mcp.ts new file mode 100755 index 000000000..5b468fd8b --- /dev/null +++ b/packages/coding-agent/test/fixtures/instructions-mcp.ts @@ -0,0 +1,88 @@ +#!/usr/bin/env bun +/** + * Test fixture: a minimal, well-behaved stdio MCP server that reports + * server-provided `instructions` on `initialize` and exposes a single tool. + * + * Used by `sdk-mcp-instructions.test.ts` to prove that a deferred interactive + * (`hasUI`) session, whose MCP discovery runs in the background, still folds + * each connected server's instructions into the system prompt once the + * connection completes — see issue: instructions were previously dropped + * permanently for deferred UI sessions. + * + * Speaks newline-delimited JSON-RPC 2.0 (the wire format of `StdioTransport`): + * one JSON object per line on stdin, one JSON response per line on stdout. + * Only requests (objects with an `id`) get a response; notifications are + * dropped. Server-to-client requests are never sent — the client side only + * needs `initialize` + `tools/list` answered to register the tool and capture + * the instructions. + * + * Exported `SERVER_INSTRUCTIONS` is imported by the test for the assertion; + * the server only starts when run as the entry module (`import.meta.main`), so + * importing the constant never spawns a server in the test process. + */ +import * as readline from "node:readline"; + +/** Sentinel the test greps for in the rebuilt system prompt. */ +export const SERVER_INSTRUCTIONS = + "INSTR_FIXTURE_SENTINEL_3f9a2c: when this server is connected, always greet in Latin."; + +/** Single tool advertised by the fixture so `tools/list` is non-empty. */ +export const TOOL_NAME = "do_thing"; + +type JsonRpcRequest = { + jsonrpc: "2.0"; + id?: string | number; + method: string; + params?: Record; +}; + +function buildResult(method: string): Record { + switch (method) { + case "initialize": + return { + protocolVersion: "2025-03-26", + serverInfo: { name: "instr-fixture", version: "1.0.0" }, + // Declare only the tools capability so the client never probes + // resources/list or prompts/list — keeps the fixture minimal. + capabilities: { tools: {} }, + instructions: SERVER_INSTRUCTIONS, + }; + case "tools/list": + return { + tools: [ + { + name: TOOL_NAME, + description: "Fixture tool; never actually called by this test.", + inputSchema: { type: "object", properties: {}, additionalProperties: false }, + }, + ], + }; + default: + // `ping` and any other request: a benign empty result keeps the + // transport happy without modelling methods the test never exercises. + return {}; + } +} + +function startServer(): void { + const rl = readline.createInterface({ input: process.stdin }); + rl.on("line", line => { + const trimmed = line.trim(); + if (trimmed.length === 0) return; + let msg: JsonRpcRequest; + try { + msg = JSON.parse(trimmed) as JsonRpcRequest; + } catch { + return; + } + // Notifications (no `id`) get no response. + if (msg.id === undefined || msg.id === null) return; + const response = { jsonrpc: "2.0" as const, id: msg.id, result: buildResult(msg.method) }; + process.stdout.write(`${JSON.stringify(response)}\n`); + }); + rl.on("close", () => process.exit(0)); +} + +if (import.meta.main) { + startServer(); +} diff --git a/packages/coding-agent/test/fixtures/many-tools-mcp.ts b/packages/coding-agent/test/fixtures/many-tools-mcp.ts new file mode 100755 index 000000000..8f7f7933f --- /dev/null +++ b/packages/coding-agent/test/fixtures/many-tools-mcp.ts @@ -0,0 +1,88 @@ +#!/usr/bin/env bun +/** + * Test fixture: a minimal stdio MCP server that advertises MANY tools, used by + * `sdk-mcp-auto-discovery.test.ts` to prove two deferred-discovery contracts: + * + * 1. `tools.discoveryMode: "auto"` is recomputed once the real MCP tool count + * is known — a toolset this large must flip discovery ON for a session whose + * pre-discovery registry was under the threshold, instead of force-activating + * every tool with no `search_tool_bm25` registered. + * 2. A session disposed while the server is still connecting must disconnect + * the manager and never have MCP tools resurrected onto it. The optional + * `--delay ` argv stalls the `initialize` response so the test can + * deterministically dispose mid-connect. + * + * Speaks newline-delimited JSON-RPC 2.0 (the wire format of `StdioTransport`), + * same shape as `instructions-mcp.ts`. Exported constants are imported by the + * test; the server only starts when run as the entry module (`import.meta.main`). + */ +import * as readline from "node:readline"; + +/** Enough tools to push any small session past TOOL_DISCOVERY_AUTO_THRESHOLD (40). */ +export const MANY_TOOL_COUNT = 45; + +/** Alphabetic names: MCP tool-name sanitization strips digits, so numeric + * suffixes like `tool_01` would all collapse into one colliding name. */ +export function manyToolName(index: number): string { + const hi = String.fromCharCode(97 + Math.floor(index / 26)); + const lo = String.fromCharCode(97 + (index % 26)); + return `tool_${hi}${lo}`; +} + +type JsonRpcRequest = { + jsonrpc: "2.0"; + id?: string | number; + method: string; + params?: Record; +}; + +function buildResult(method: string): Record { + switch (method) { + case "initialize": + return { + protocolVersion: "2025-03-26", + serverInfo: { name: "many-fixture", version: "1.0.0" }, + capabilities: { tools: {} }, + }; + case "tools/list": + return { + tools: Array.from({ length: MANY_TOOL_COUNT }, (_, i) => ({ + name: manyToolName(i), + description: `Fixture tool #${i}; never actually called by the test.`, + inputSchema: { type: "object", properties: {}, additionalProperties: false }, + })), + }; + default: + return {}; + } +} + +function startServer(): void { + const delayIndex = process.argv.indexOf("--delay"); + const initializeDelayMs = delayIndex >= 0 ? Number(process.argv[delayIndex + 1]) || 0 : 0; + const rl = readline.createInterface({ input: process.stdin }); + rl.on("line", line => { + void (async () => { + const trimmed = line.trim(); + if (trimmed.length === 0) return; + let msg: JsonRpcRequest; + try { + msg = JSON.parse(trimmed) as JsonRpcRequest; + } catch { + return; + } + // Notifications (no `id`) get no response. + if (msg.id === undefined || msg.id === null) return; + if (msg.method === "initialize" && initializeDelayMs > 0) { + await Bun.sleep(initializeDelayMs); + } + const response = { jsonrpc: "2.0" as const, id: msg.id, result: buildResult(msg.method) }; + process.stdout.write(`${JSON.stringify(response)}\n`); + })(); + }); + rl.on("close", () => process.exit(0)); +} + +if (import.meta.main) { + startServer(); +} diff --git a/packages/coding-agent/test/modes/components/assistant-message-streaming-fastpath.test.ts b/packages/coding-agent/test/modes/components/assistant-message-streaming-fastpath.test.ts new file mode 100644 index 000000000..f926926de --- /dev/null +++ b/packages/coding-agent/test/modes/components/assistant-message-streaming-fastpath.test.ts @@ -0,0 +1,182 @@ +import { afterEach, beforeAll, beforeEach, describe, expect, it } from "bun:test"; +import type { AssistantMessage } from "@oh-my-pi/pi-ai"; +import { resetSettingsForTest, Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { AssistantMessageComponent } from "@oh-my-pi/pi-coding-agent/modes/components/assistant-message"; +import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme"; +import { type Component, Container, Markdown } from "@oh-my-pi/pi-tui"; + +const W = 100; + +function msg(content: AssistantMessage["content"], extra: Partial = {}): AssistantMessage { + return { + role: "assistant", + content, + api: "anthropic-messages", + provider: "anthropic", + model: "m", + usage: { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 0, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop", + timestamp: 0, + ...extra, + }; +} + +/** Render `m` on a brand-new component, which always takes the teardown path. */ +function teardownRender(m: AssistantMessage): string { + const fresh = new AssistantMessageComponent(); + fresh.updateContent(m); + return fresh.render(W).join("\n"); +} + +beforeAll(async () => { + await initTheme(false); +}); + +beforeEach(async () => { + resetSettingsForTest(); + await Settings.init({ inMemory: true }); +}); + +afterEach(() => { + resetSettingsForTest(); +}); + +// Contract: the streaming fast path (a component reused across updateContent +// calls, which reuses Markdown children via setText) MUST render byte-identical +// output to the teardown path (a fresh component that rebuilds every child) for +// the same message — at every step. If they ever diverge, the optimization +// silently corrupts the transcript. +describe("AssistantMessageComponent streaming fast path", () => { + it("matches teardown output across a growing thinking + text stream", () => { + const reused = new AssistantMessageComponent(); + const thinking = "Reasoning about the **problem** with `code` and a list:\n- a\n- b"; + const steps = [ + "He", + "Hello, ", + "Hello, world.", + "Hello, world.\n\n## Heading\n\nSome `inline` and **bold** text.", + "Hello, world.\n\n## Heading\n\nSome `inline` and **bold** text.\n\n```ts\nconst x = 1;\n```", + ]; + for (const text of steps) { + const m = msg([ + { type: "thinking", thinking }, + { type: "text", text }, + ]); + reused.updateContent(m); + expect(reused.render(W).join("\n")).toBe(teardownRender(m)); + } + }); + + it("matches teardown for a single growing text block", () => { + const reused = new AssistantMessageComponent(); + let text = ""; + for (const chunk of ["The ", "quick ", "brown ", "**fox** ", "jumps."]) { + text += chunk; + const m = msg([{ type: "text", text }]); + reused.updateContent(m); + expect(reused.render(W).join("\n")).toBe(teardownRender(m)); + } + }); + + // Regression: theme/symbol changes reach the component via invalidate() + // (InteractiveMode clears the markdown render cache and invalidates the + // tree). Reused fast-path children captured getMarkdownTheme() at + // construction, so invalidate() MUST drop them and rebuild — otherwise a + // theme switch keeps rendering stale symbols until the message shape + // changes. Child identity is the load-bearing mechanism here: a kept + // instance is exactly a kept stale theme. + it("invalidate() rebuilds Markdown children instead of reusing fast-path state", () => { + const collectMarkdown = (component: Container): Markdown[] => { + const found: Markdown[] = []; + const walk = (node: Component): void => { + if (node instanceof Markdown) found.push(node); + if (node instanceof Container) for (const child of node.children) walk(child); + }; + walk(component); + return found; + }; + + const reused = new AssistantMessageComponent(); + reused.updateContent(msg([{ type: "text", text: "Hello **world**, part one." }])); + reused.updateContent(msg([{ type: "text", text: "Hello **world**, part one and two." }])); + const before = collectMarkdown(reused); + expect(before.length).toBeGreaterThan(0); + + // Sanity: a same-shape streaming update reuses the children (fast path on). + reused.updateContent(msg([{ type: "text", text: "Hello **world**, part one, two, three." }])); + const streamed = collectMarkdown(reused); + expect(streamed.length).toBe(before.length); + for (let i = 0; i < streamed.length; i++) { + expect(streamed[i]).toBe(before[i]); + } + + reused.invalidate(); + const rebuilt = collectMarkdown(reused); + expect(rebuilt.length).toBe(before.length); + for (let i = 0; i < rebuilt.length; i++) { + expect(rebuilt[i]).not.toBe(before[i]); + } + }); + + // Regression: #fastPathItems are keyed by raw content index, but a + // `redactedThinking` block is not rendered. If one appears mid-stream it + // shifts the indices of the visible blocks; the shape key must reflect that + // (or the fast path must fail closed) so children are not mis-targeted. + it("matches teardown when a redactedThinking block shifts indices mid-stream", () => { + const reused = new AssistantMessageComponent(); + const a = msg([ + { type: "thinking", thinking: "step one details here" }, + { type: "text", text: "answer one" }, + ]); + reused.updateContent(a); + expect(reused.render(W).join("\n")).toBe(teardownRender(a)); + + // A redactedThinking block appears at index 0, pushing thinking->1, text->2. + const b = msg([ + { type: "redactedThinking", data: "opaque-blob" }, + { type: "thinking", thinking: "step two with more detail" }, + { type: "text", text: "answer two is longer now" }, + ]); + reused.updateContent(b); + expect(reused.render(W).join("\n")).toBe(teardownRender(b)); + }); + + it("matches teardown when an error trailer appears after streamed text", () => { + const reused = new AssistantMessageComponent(); + const ok = msg([{ type: "text", text: "partial answer in progress" }]); + reused.updateContent(ok); + expect(reused.render(W).join("\n")).toBe(teardownRender(ok)); + + const errored = msg([{ type: "text", text: "partial answer in progress" }], { + stopReason: "error", + errorMessage: "upstream 502", + }); + reused.updateContent(errored); + expect(reused.render(W).join("\n")).toBe(teardownRender(errored)); + }); + + it("matches teardown when a block visibility toggles (empty -> non-empty)", () => { + const reused = new AssistantMessageComponent(); + // First an empty trailing text block (not rendered), then it gains content. + const empty = msg([ + { type: "thinking", thinking: "thinking out loud" }, + { type: "text", text: "" }, + ]); + reused.updateContent(empty); + expect(reused.render(W).join("\n")).toBe(teardownRender(empty)); + + const filled = msg([ + { type: "thinking", thinking: "thinking out loud" }, + { type: "text", text: "now there is an answer" }, + ]); + reused.updateContent(filled); + expect(reused.render(W).join("\n")).toBe(teardownRender(filled)); + }); +}); diff --git a/packages/coding-agent/test/sdk-mcp-auto-discovery.test.ts b/packages/coding-agent/test/sdk-mcp-auto-discovery.test.ts new file mode 100644 index 000000000..260f2a103 --- /dev/null +++ b/packages/coding-agent/test/sdk-mcp-auto-discovery.test.ts @@ -0,0 +1,154 @@ +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, mock, spyOn } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { AuthStorage } from "@oh-my-pi/pi-ai"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { createAgentSession } from "@oh-my-pi/pi-coding-agent/sdk"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { Snowflake } from "@oh-my-pi/pi-utils"; +import { MANY_TOOL_COUNT } from "./fixtures/many-tools-mcp"; + +// Contracts for deferred (hasUI) MCP discovery follow-ups: +// +// 1. `tools.discoveryMode: "auto"` is resolved at session build time against a +// registry that cannot yet contain the deferred MCP tools. When the +// background connect reports a toolset large enough to cross the auto +// threshold, the session MUST upgrade to discovery mode — register and +// activate `search_tool_bm25`, mark discovery enabled, and expose the MCP +// tools as discoverable — instead of force-activating all of them. +// +// 2. A session disposed while servers are still connecting MUST NOT be touched +// by the late discovery result: no tools resurrected onto the disposed +// session and the manager's transports disconnected. +const FIXTURE_PATH = path.join(import.meta.dir, "fixtures", "many-tools-mcp.ts"); + +describe("createAgentSession deferred MCP auto discovery", () => { + let registryDir: string; + let tempDir: string; + let authStorage: AuthStorage; + let modelRegistry: ModelRegistry; + // Discovery resolves user-level MCP config from `os.homedir()`; redirect it + // to an empty dir so the test connects ONLY to the fixture server and never + // spawns the developer's real MCP servers. + let isolatedHome: string; + + beforeAll(async () => { + registryDir = path.join(os.tmpdir(), `pi-sdk-mcp-auto-registry-${Snowflake.next()}`); + fs.mkdirSync(registryDir, { recursive: true }); + isolatedHome = path.join(os.tmpdir(), `pi-sdk-mcp-auto-home-${Snowflake.next()}`); + fs.mkdirSync(isolatedHome, { recursive: true }); + authStorage = await AuthStorage.create(path.join(registryDir, "auth.db")); + modelRegistry = new ModelRegistry(authStorage); + }); + + afterAll(() => { + authStorage.close(); + for (const dir of [registryDir, isolatedHome]) { + if (dir && fs.existsSync(dir)) { + fs.rmSync(dir, { recursive: true, force: true }); + } + } + }); + + beforeEach(() => { + tempDir = path.join(os.tmpdir(), `pi-sdk-mcp-auto-${Snowflake.next()}`); + fs.mkdirSync(tempDir, { recursive: true }); + spyOn(os, "homedir").mockReturnValue(isolatedHome); + }); + + afterEach(() => { + if (tempDir && fs.existsSync(tempDir)) { + fs.rmSync(tempDir, { recursive: true, force: true }); + } + mock.restore(); + }); + + const writeMcpConfig = (extraArgs: string[] = []) => { + fs.writeFileSync( + path.join(tempDir, ".mcp.json"), + JSON.stringify({ + mcpServers: { + many: { type: "stdio", command: process.execPath, args: [FIXTURE_PATH, ...extraArgs] }, + }, + }), + ); + }; + + const baseOptions = () => ({ + cwd: tempDir, + agentDir: tempDir, + modelRegistry, + sessionManager: SessionManager.inMemory(), + settings: Settings.isolated({}), + model: getBundledModel("openai", "gpt-4o-mini"), + disableExtensionDiscovery: true, + skills: [], + contextFiles: [], + promptTemplates: [], + slashCommands: [], + enableLsp: false, + skipPythonPreflight: true, + enableMCP: true, + hasUI: true, + }); + + it("flips auto discovery on when the deferred MCP toolset crosses the threshold", async () => { + writeMcpConfig(); + // A small explicit toolset keeps the pre-discovery registry far below the + // 40-tool auto threshold; the fixture's 45 tools must push it across. + const { session } = await createAgentSession({ ...baseOptions(), toolNames: ["read", "edit", "bash"] }); + try { + // Genuine integration wait: discovery spawns the fixture as a real + // subprocess and connects asynchronously, and the SDK fires that work + // fire-and-forget with no completion promise or event exposed — fake + // timers cannot drive a child process, so poll the live session with + // a generous ceiling, exiting the instant discovery flips on. + const deadline = Date.now() + 12_000; + while (!session.isMCPDiscoveryEnabled() && Date.now() < deadline) { + await Bun.sleep(50); + } + + expect(session.isMCPDiscoveryEnabled()).toBe(true); + const activeNames = session.getActiveToolNames(); + expect(activeNames).toContain("search_tool_bm25"); + // Discovery mode means the MCP tools are searchable, NOT force-activated. + expect(activeNames.filter(name => name.startsWith("mcp__"))).toEqual([]); + const discoverable = session.getDiscoverableTools({ source: "mcp" }); + expect(discoverable.length).toBe(MANY_TOOL_COUNT); + } finally { + await session.dispose(); + } + }, 20_000); + + it("disposing mid-connect disconnects the manager and never resurrects tools", async () => { + // Stall `initialize` in the real fixture subprocess so the connect is + // guaranteed to still be in flight when dispose() runs. Deterministic + // time control cannot order a race against a child process. + writeMcpConfig(["--delay", "750"]); + const { session, mcpManager } = await createAgentSession({ + ...baseOptions(), + toolNames: ["read", "edit", "bash"], + }); + expect(mcpManager).toBeDefined(); + if (!mcpManager) throw new Error("expected deferred session to own an MCPManager"); + const disconnectSpy = spyOn(mcpManager, "disconnectAll"); + + await session.dispose(); + expect(session.isDisposed).toBe(true); + + // Genuine integration wait (see above): the deferred task notices the + // disposed session once the stalled connect resolves and must disconnect + // instead of refreshing tools. Exits the instant the spy fires. + const deadline = Date.now() + 12_000; + while (disconnectSpy.mock.calls.length === 0 && Date.now() < deadline) { + await Bun.sleep(50); + } + expect(disconnectSpy).toHaveBeenCalled(); + expect(session.getActiveToolNames().filter(name => name.startsWith("mcp__"))).toEqual([]); + expect(session.getActiveToolNames()).not.toContain("search_tool_bm25"); + expect(session.isMCPDiscoveryEnabled()).toBe(false); + }, 20_000); +}); diff --git a/packages/coding-agent/test/sdk-mcp-defer.test.ts b/packages/coding-agent/test/sdk-mcp-defer.test.ts new file mode 100644 index 000000000..36a3ae42e --- /dev/null +++ b/packages/coding-agent/test/sdk-mcp-defer.test.ts @@ -0,0 +1,96 @@ +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { AuthStorage } from "@oh-my-pi/pi-ai"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { createAgentSession } from "@oh-my-pi/pi-coding-agent/sdk"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { Snowflake } from "@oh-my-pi/pi-utils"; + +// Contract for B1 (interactive MCP deferral): when `hasUI` is true, MCP +// discovery is deferred off the first-paint path, so an explicitly requested +// MCP tool (e.g. via `--tools`) whose server has not yet connected MUST still +// be a *known* tool — registered as a deterministic "still connecting" +// placeholder — rather than vanishing and surfacing as "unknown tool" if the +// model calls it before the background connection completes. With `hasUI` +// false there is no deferral, so an MCP tool name with no real backing is not +// registered at all (the non-UI paths keep the blocking discover path). +describe("createAgentSession MCP deferral (B1)", () => { + let registryDir: string; + let tempDir: string; + let authStorage: AuthStorage; + let modelRegistry: ModelRegistry; + + const PENDING_MCP_TOOL = "mcp__pending_connectingtool"; + + const baseOptions = () => ({ + cwd: tempDir, + agentDir: tempDir, + modelRegistry, + sessionManager: SessionManager.inMemory(), + settings: Settings.isolated({}), + model: getBundledModel("openai", "gpt-4o-mini"), + disableExtensionDiscovery: true, + skills: [], + contextFiles: [], + promptTemplates: [], + slashCommands: [], + enableLsp: false, + skipPythonPreflight: true, + // No .mcp.json in tempDir, so no real MCP server can ever back this name. + enableMCP: true, + toolNames: ["read", PENDING_MCP_TOOL], + }); + + beforeAll(async () => { + registryDir = path.join(os.tmpdir(), `pi-sdk-mcp-defer-registry-${Snowflake.next()}`); + fs.mkdirSync(registryDir, { recursive: true }); + authStorage = await AuthStorage.create(path.join(registryDir, "auth.db")); + modelRegistry = new ModelRegistry(authStorage); + }); + + afterAll(() => { + authStorage.close(); + if (registryDir && fs.existsSync(registryDir)) { + fs.rmSync(registryDir, { recursive: true, force: true }); + } + }); + + beforeEach(() => { + tempDir = path.join(os.tmpdir(), `pi-sdk-mcp-defer-${Snowflake.next()}`); + fs.mkdirSync(tempDir, { recursive: true }); + }); + + afterEach(() => { + if (tempDir && fs.existsSync(tempDir)) { + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }); + + it("registers a pending placeholder for an explicit MCP tool when hasUI defers discovery", async () => { + const { session } = await createAgentSession({ ...baseOptions(), hasUI: true }); + try { + // The explicitly requested MCP tool is a known, resolvable tool even + // though no server has connected — deterministic, not "unknown tool". + expect(session.getActiveToolNames()).toContain(PENDING_MCP_TOOL); + } finally { + await session.dispose(); + } + }); + + it("does not fabricate the MCP tool in non-UI mode (no deferral, no backing server)", async () => { + const { session } = await createAgentSession({ ...baseOptions(), hasUI: false }); + try { + // Without deferral there is no placeholder; the name has no real + // server backing, so it is simply not a registered tool. + expect(session.getActiveToolNames()).not.toContain(PENDING_MCP_TOOL); + // A normal builtin is unaffected. + expect(session.getActiveToolNames()).toContain("read"); + } finally { + await session.dispose(); + } + }); +}); diff --git a/packages/coding-agent/test/sdk-mcp-instructions.test.ts b/packages/coding-agent/test/sdk-mcp-instructions.test.ts new file mode 100644 index 000000000..21bf0a364 --- /dev/null +++ b/packages/coding-agent/test/sdk-mcp-instructions.test.ts @@ -0,0 +1,116 @@ +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, mock, spyOn } from "bun:test"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { AuthStorage } from "@oh-my-pi/pi-ai"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; +import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; +import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; +import { createAgentSession } from "@oh-my-pi/pi-coding-agent/sdk"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; +import { Snowflake } from "@oh-my-pi/pi-utils"; +import { SERVER_INSTRUCTIONS } from "./fixtures/instructions-mcp"; + +// Contract: a deferred interactive (`hasUI`) session runs MCP discovery off the +// first-paint path, so an MCP server's `instructions` are not available when the +// prompt is first built. Once the background connection completes — and the +// resulting `refreshMCPTools` rebuilds the system prompt — that server's +// instructions MUST join the prompt for the rest of the session. Regression +// guard: a prior version gated instruction inclusion on `!deferMCPDiscoveryForUI`, +// which dropped server instructions permanently for every UI session. +const FIXTURE_PATH = path.join(import.meta.dir, "fixtures", "instructions-mcp.ts"); + +describe("createAgentSession MCP server instructions (deferred UI)", () => { + let registryDir: string; + let tempDir: string; + let authStorage: AuthStorage; + let modelRegistry: ModelRegistry; + // Discovery resolves user-level MCP config from `os.homedir()`; redirect it + // to an empty dir so the test connects ONLY to the fixture server and never + // spawns the developer's real MCP servers. + let isolatedHome: string; + + beforeAll(async () => { + registryDir = path.join(os.tmpdir(), `pi-sdk-mcp-instr-registry-${Snowflake.next()}`); + fs.mkdirSync(registryDir, { recursive: true }); + isolatedHome = path.join(os.tmpdir(), `pi-sdk-mcp-instr-home-${Snowflake.next()}`); + fs.mkdirSync(isolatedHome, { recursive: true }); + authStorage = await AuthStorage.create(path.join(registryDir, "auth.db")); + modelRegistry = new ModelRegistry(authStorage); + }); + + afterAll(() => { + authStorage.close(); + for (const dir of [registryDir, isolatedHome]) { + if (dir && fs.existsSync(dir)) { + fs.rmSync(dir, { recursive: true, force: true }); + } + } + }); + + beforeEach(() => { + tempDir = path.join(os.tmpdir(), `pi-sdk-mcp-instr-${Snowflake.next()}`); + fs.mkdirSync(tempDir, { recursive: true }); + spyOn(os, "homedir").mockReturnValue(isolatedHome); + fs.writeFileSync( + path.join(tempDir, ".mcp.json"), + JSON.stringify({ + mcpServers: { + instr: { type: "stdio", command: process.execPath, args: [FIXTURE_PATH] }, + }, + }), + ); + }); + + afterEach(() => { + if (tempDir && fs.existsSync(tempDir)) { + fs.rmSync(tempDir, { recursive: true, force: true }); + } + mock.restore(); + }); + + it("folds server instructions into the prompt once deferred discovery connects", async () => { + const { session } = await createAgentSession({ + cwd: tempDir, + agentDir: tempDir, + modelRegistry, + sessionManager: SessionManager.inMemory(), + settings: Settings.isolated({}), + model: getBundledModel("openai", "gpt-4o-mini"), + disableExtensionDiscovery: true, + skills: [], + contextFiles: [], + promptTemplates: [], + slashCommands: [], + enableLsp: false, + skipPythonPreflight: true, + enableMCP: true, + hasUI: true, + }); + try { + // First paint: discovery is still in flight, so the server's + // instructions are not yet present. + expect(session.systemPrompt.join("\n")).not.toContain(SERVER_INSTRUCTIONS); + + // Background connect + `refreshMCPTools` rebuild must surface the + // instructions. This is a genuine integration wait: discovery spawns + // the fixture as a real subprocess and connects asynchronously, and + // the SDK fires that work fire-and-forget with no completion promise + // or event exposed to await — so fake timers cannot drive it and we + // poll the live prompt with a generous ceiling, exiting the instant + // the rebuilt prompt carries the instructions. + const deadline = Date.now() + 12_000; + let prompt = session.systemPrompt.join("\n"); + while (!prompt.includes(SERVER_INSTRUCTIONS) && Date.now() < deadline) { + await Bun.sleep(50); + prompt = session.systemPrompt.join("\n"); + } + + expect(prompt).toContain(SERVER_INSTRUCTIONS); + // The instructions are framed under the MCP section, not pasted raw. + expect(prompt).toContain("MCP Server Instructions"); + } finally { + await session.dispose(); + } + }, 20_000); +}); diff --git a/packages/tui/CHANGELOG.md b/packages/tui/CHANGELOG.md index d3bdf91fc..d15ef5710 100644 --- a/packages/tui/CHANGELOG.md +++ b/packages/tui/CHANGELOG.md @@ -127,6 +127,10 @@ - Removed the probe/defer API surface: `TUI.setEagerNativeScrollbackRebuild()`, `TUI.refreshNativeScrollbackIfDirty()`, `TUI.setClearOnShrink()`/`getClearOnShrink()`, `RenderRequestOptions.allowUnknownViewportMutation`, `NativeScrollbackRefreshOptions`, `Terminal.isNativeViewportAtBottom()`, `Terminal.hasEagerEraseScrollbackRisk()`, and the `eagerEraseScrollbackRisk`/`submitPinsViewportToTail` capability fields with their detectors. - Removed the `PI_TUI_ED3_SAFE`, `PI_CLEAR_ON_SHRINK`, and `PI_TUI_DEBUG` environment variables (the levers they tuned no longer exist; `PI_DEBUG_REDRAW` now logs the commit-ledger state per frame). +### Changed + +- Markdown rendering during streaming re-lexes only the grown tail instead of the whole buffer on every reveal tick. marked has no resumable lexer, but block tokenization is local across a blank-line boundary with balanced fences, so the largest blank-line-bounded prefix's block tokens are frozen and reused (`lex(prefix) ++ lex(tail)`), with a full-lex fallback for non-append edits, reference-link definitions, and CRLF input. The output is byte-identical to a full lex (covered by a contract test), turning the O(N²) cost of revealing a long single-block message into O(N): a 6,000-grapheme reveal dropped from ~575 ms to ~89 ms of CPU in benchmarks. + ## [15.10.9] - 2026-06-09 ### Added diff --git a/packages/tui/src/components/markdown.ts b/packages/tui/src/components/markdown.ts index bc0218d7e..aca3b9445 100644 --- a/packages/tui/src/components/markdown.ts +++ b/packages/tui/src/components/markdown.ts @@ -61,6 +61,13 @@ const RENDER_CACHE_MAX = 256; // sane cap: ~256 distinct message × width combos const EMPTY_RENDER_LINES: readonly string[] = []; const renderCache = new LRUCache({ max: RENDER_CACHE_MAX }); +// A reference-link definition (`[label]: dest`) resolves across the whole +// document, so a split lex cannot reproduce it — disable the streaming fast path +// when one is present (rare in streamed output). The label may contain +// backslash-escaped characters (`[a\]b]: x`), so escapes are matched explicitly; +// over-matching is safe (it only costs the fast path), under-matching is not. +const HAS_REF_DEF = /^ {0,3}\[(?:\\.|[^\]\\])+\]:/m; + /** Drop all L2 cache entries. Call on theme change to prevent stale styled output. */ export function clearRenderCache(): void { renderCache.clear(); @@ -297,6 +304,16 @@ export class Markdown implements Component { #cachedLines?: readonly string[]; #transientRenderCache = false; + // Streaming-lex cache: the largest blank-line-bounded prefix of #text whose + // block tokens are frozen, plus those tokens. marked has no resumable lexer, + // but block tokenization is local across a "\n\n" boundary with balanced + // fences, so lex(prefix) ++ lex(tail) === lex(prefix+tail). On append-only + // growth (the streaming path) this re-lexes only the grown tail instead of the + // whole buffer, turning O(N^2) reveal cost into O(N). Width/theme do not affect + // tokenization, so this cache is independent of the render caches above. + #streamPrefixText?: string; + #streamPrefixTokens?: Token[]; + constructor( text: string, paddingX: number, @@ -315,6 +332,13 @@ export class Markdown implements Component { setText(text: string): void { this.#text = text; + if (!text.trim()) { + // Blank replacement: render() early-returns before #lexTokens can see + // the non-append edit, so drop the frozen stream state here or it + // outlives the content it indexed. + this.#streamPrefixText = undefined; + this.#streamPrefixTokens = undefined; + } this.invalidate(); } @@ -334,6 +358,76 @@ export class Markdown implements Component { this.invalidate(); } + // Lex `text` into block tokens, reusing the frozen stable prefix when the text + // only grew (the streaming path). Falls back to a full lex whenever the prefix + // is no longer a prefix (non-append edit), the text carries reference-link + // definitions, or it contains CR (marked normalizes CRLF, which would desync + // raw-span offsets). Every fallback is correctness-preserving — only speed + // differs; the render loop sees the identical token list either way. + #lexTokens(text: string): Token[] { + const canStream = !HAS_REF_DEF.test(text) && !text.includes("\r"); + const prefix = this.#streamPrefixText; + const prefixTokens = this.#streamPrefixTokens; + if ( + canStream && + prefix !== undefined && + prefixTokens !== undefined && + text.length > prefix.length && + text.startsWith(prefix) + ) { + const tailTokens = markdownParser.lexer(text.slice(prefix.length)); + const tokens = [...prefixTokens, ...tailTokens]; + this.#freezeStablePrefix(text, tokens); + return tokens; + } + const tokens = markdownParser.lexer(text); + if (canStream) { + this.#freezeStablePrefix(text, tokens); + } else { + this.#streamPrefixText = undefined; + this.#streamPrefixTokens = undefined; + } + return tokens; + } + + // Freeze the largest run of leading blocks that end on a hard "\n\n" boundary + // (complete and immutable under append-only growth) so the next streaming + // render re-lexes only the unfrozen tail. Caller guarantees no CR / no + // reference definitions, so each token's `raw` is a verbatim slice of `text` + // and the summed offsets address `text` exactly. + #freezeStablePrefix(text: string, tokens: Token[]): void { + let pos = 0; + let frozenEnd = 0; + let frozenCount = 0; + for (let i = 0; i < tokens.length; i++) { + const raw = tokens[i].raw; + const end = pos + raw.length; + // A `space` token ending in "\n\n" closes the preceding block, but a + // `list` before it can still be extended by a following same-marker + // item across the blank line (CommonMark loose-list continuation), + // which marked merges into one renumbered loose list. Freezing across + // such a cut would keep the lists separate. Never freeze right after a + // list — it stays in the re-lexed tail. + if (raw.endsWith("\n\n") && tokens[i - 1]?.type !== "list") { + frozenEnd = end; + frozenCount = i + 1; + } + pos = end; + } + // Freeze only when the tail begins with real block content. If the next + // char is whitespace (an extra blank line, or an indented continuation), + // the block separator straddles the cut and lex(prefix)++lex(tail) would + // desync from a full lex — e.g. a fence followed by "\n\n\n- list". When + // frozenEnd is at end-of-text the next char is unknown, so defer. + if (frozenCount > 0 && frozenEnd < text.length) { + const next = text.charCodeAt(frozenEnd); + if (next !== 0x20 /* space */ && next !== 0x0a /* \n */) { + this.#streamPrefixText = text.slice(0, frozenEnd); + this.#streamPrefixTokens = tokens.slice(0, frozenCount); + } + } + } + render(width: number): readonly string[] { // L1: per-instance cache — fastest path for repeated renders of the same // instance at the same width (e.g. resize debounce, repeated redraws). @@ -385,7 +479,7 @@ export class Markdown implements Component { } // Parse markdown to HTML-like tokens - const tokens = markdownParser.lexer(normalizedText); + const tokens = this.#lexTokens(normalizedText); // Convert tokens to styled terminal output const renderedLines: string[] = []; diff --git a/packages/tui/test/markdown-incremental-lex.test.ts b/packages/tui/test/markdown-incremental-lex.test.ts new file mode 100644 index 000000000..818f503ca --- /dev/null +++ b/packages/tui/test/markdown-incremental-lex.test.ts @@ -0,0 +1,195 @@ +import { describe, expect, it } from "bun:test"; +import { clearRenderCache, Markdown } from "@oh-my-pi/pi-tui/components/markdown"; +import { defaultMarkdownTheme } from "./test-themes.js"; + +// E2 contract: the streaming incremental lexer (lex(prefix) ++ lex(tail), reusing +// frozen blank-line-bounded blocks) must produce BYTE-IDENTICAL output to a fresh +// full lex of the same text at every growth step. A faster-but-divergent render is +// a regression, so this is the gate that keeps E2 honest. +// +// Masking hazard: Markdown's module-level L2 render cache keys on (text, width), +// so a streaming render that produced WRONG lines would cache them and the "fresh" +// oracle would read the same wrong lines back. We `clearRenderCache()` around the +// oracle so it always cold-lexes, and again so the next streaming render cannot +// hit a stale entry — the streaming instance must go through its own incremental +// `#lexTokens` path every step. + +const THEME = defaultMarkdownTheme; + +function renderCold(text: string, width: number): readonly string[] { + clearRenderCache(); + const out = new Markdown(text, 0, 0, THEME).render(width); + clearRenderCache(); + return out; +} + +/** Reveal `full` in `step`-char increments through ONE reused (streaming) instance + * and assert each step matches a cold full-lex render of the same prefix. */ +function assertIdenticalGrowth(full: string, width = 60, step = 13): void { + const streaming = new Markdown("", 0, 0, THEME); + for (let len = 1; len <= full.length; len += step) { + const slice = full.slice(0, len); + clearRenderCache(); + streaming.setText(slice); + const streamLines = streaming.render(width); + const oracle = renderCold(slice, width); + expect(streamLines).toEqual(oracle); + } + clearRenderCache(); + streaming.setText(full); + const streamLines = streaming.render(width); + expect(streamLines).toEqual(renderCold(full, width)); +} + +const PROSE = + "Para one with **bold** and _italic_ words and a `code span` for flavor.\n\n" + + "Para two continues the document with more sentences so the lexer has real\n" + + "block structure to chew on, then a third paragraph grows at the tail end as\n\n" + + "the stream appends additional content token by token over many frames here."; + +const FENCED = + "Intro paragraph before the code block begins streaming in slowly.\n\n" + + "```ts\nconst x: number = compute(a, b) + delta;\nfor (let i = 0; i < x; i++) {\n emit(i);\n}\nreturn x.toFixed(2);\n```\n\n" + + "Trailing prose after the fence keeps growing with more and more sentences."; + +const LIST = + "Lead-in sentence before the list.\n\n" + + "- first bullet item with `inline`\n- second bullet item in **bold**\n- third bullet item\n\n" + + "1. ordered one\n2. ordered two\n3. ordered three\n\n" + + "Closing paragraph that keeps streaming additional words to the very end here."; + +const HEADINGS = + "# Title heading\n\nIntro text under the title with some detail.\n\n" + + "## Section two\n\nBody of section two grows over time.\n\n" + + "### Subsection\n\nDeeper content that streams in at the tail as the reveal advances."; + +const MIXED = (() => { + const para = + "The quick brown fox jumps over the lazy dog while a `code span` and **bold** _italic_ exercise things. "; + const cb = "\n```ts\nconst x = compute(a, b);\nreturn x;\n```\n\n"; + const list = "\n- one\n- two `inline`\n- three\n\n"; + let out = ""; + for (let i = 1; i <= 6; i++) out += `## Section ${i}\n\n${para}${para}${cb}${list}`; + return out; +})(); + +describe("Markdown incremental streaming lex (E2)", () => { + it("prose growth is byte-identical to full lex", () => { + assertIdenticalGrowth(PROSE); + }); + + it("fenced code growth (open then close) is byte-identical", () => { + assertIdenticalGrowth(FENCED); + }); + + it("list growth is byte-identical", () => { + assertIdenticalGrowth(LIST); + }); + + it("heading growth is byte-identical", () => { + assertIdenticalGrowth(HEADINGS); + }); + + it("mixed multi-section corpus growth is byte-identical", () => { + assertIdenticalGrowth(MIXED, 80, 29); + }); + + it("a width change mid-stream still matches a cold render at the new width", () => { + const streaming = new Markdown("", 0, 0, THEME); + // Warm the stream cache at width 80 across the whole message. + for (let len = 1; len <= MIXED.length; len += 41) { + clearRenderCache(); + streaming.setText(MIXED.slice(0, len)); + streaming.render(80); + } + // Now render the full text at a NARROWER width: frozen tokens are width- + // independent, so output must match a cold full lex at the new width. + clearRenderCache(); + streaming.setText(MIXED); + const narrow = streaming.render(40); + expect(narrow).toEqual(renderCold(MIXED, 40)); + // And back to a wider width. + clearRenderCache(); + const wide = streaming.render(100); + expect(wide).toEqual(renderCold(MIXED, 100)); + }); + + it("reference-link definitions (fallback path) still render correctly while growing", () => { + const refDoc = + "See [the docs][d] and [the spec][s] for details on the protocol.\n\n" + + "A middle paragraph with ordinary prose that keeps growing here.\n\n" + + "[d]: https://example.com/docs\n[s]: https://example.com/spec\n\n" + + "Closing paragraph streamed at the tail with extra sentences appended."; + assertIdenticalGrowth(refDoc); + }); + + // Regression: HAS_REF_DEF must also catch labels with backslash-escaped + // brackets (`[a\]b]: …`). marked resolves such a definition document-wide, + // so if the detector misses it the already-frozen paragraph keeps its plain + // text inline tokens while a cold lex rewrites `[a\]b]` into a link. + it("an escaped-bracket reference definition falls back to a correct full render", () => { + const escapedRef = + "See [a\\]b] for details in the long discussion that follows below.\n\n" + + "More prose streams in before the definition finally arrives down here.\n\n" + + "[a\\]b]: https://example.com/escaped\n\n" + + "Trailing paragraph after the definition keeps the stream going on."; + assertIdenticalGrowth(escapedRef, 60, 1); + assertIdenticalGrowth(escapedRef, 60, 13); + }); + + // Regression: marked merges a list with a following same-marker list across a + // blank line into one renumbered loose list (CommonMark loose-list + // continuation). Freezing across that "\n\n" cut keeps them separate and + // renumbers/spaces wrong. These cases must hold at the production reveal + // granularity (MIN_STEP=3) and at step=1 — the divergence is phase-sensitive. + it("two consecutive ordered lists stay merged/renumbered while growing", () => { + const twoLists = "1. a\n2. b\n\n1. c\n2. d"; + assertIdenticalGrowth(twoLists, 60, 1); + assertIdenticalGrowth(twoLists, 60, 3); + }); + + it("a loose ordered list (blank lines between items) stays correct while growing", () => { + const loose = + "Intro line before the numbered list begins here.\n\n" + + "1. First point with enough words to wrap nicely.\n\n" + + "2. Second point also with sufficient words here.\n\n" + + "3. Third and final point streamed at the tail end."; + assertIdenticalGrowth(loose, 60, 1); + assertIdenticalGrowth(loose, 60, 3); + }); + + it("a loose bullet list stays correct while growing", () => { + const loose = + "Lead-in before the bullets.\n\n" + + "- alpha item with several words to wrap\n\n" + + "- beta item with several words to wrap\n\n" + + "- gamma item streamed at the tail end here."; + assertIdenticalGrowth(loose, 60, 1); + assertIdenticalGrowth(loose, 60, 3); + }); + + it("a non-append change (text replaced) falls back to a correct full render", () => { + const streaming = new Markdown("", 0, 0, THEME); + clearRenderCache(); + streaming.setText("# First document\n\nOriginal body paragraph one.\n\nOriginal body paragraph two.\n"); + streaming.render(60); + // Replace with unrelated content that is NOT a prefix-extension. + clearRenderCache(); + streaming.setText("## Different\n\nCompletely new content replacing the old buffer entirely.\n"); + const replaced = streaming.render(60); + expect(replaced).toEqual( + renderCold("## Different\n\nCompletely new content replacing the old buffer entirely.\n", 60), + ); + }); + + it("CRLF text (fallback path) renders identically to a cold lex", () => { + const streaming = new Markdown("", 0, 0, THEME); + const crlf = "Para one with content.\r\n\r\nPara two with `code`.\r\n\r\nPara three tail.\r\n"; + for (let len = 1; len <= crlf.length; len += 11) { + clearRenderCache(); + streaming.setText(crlf.slice(0, len)); + const streamLines = streaming.render(60); + expect(streamLines).toEqual(renderCold(crlf.slice(0, len), 60)); + } + }); +});