perf(coding-agent): defer boot work and cut streaming/read CPU
- Defer MCP server discovery off the first-paint critical path for UI sessions; tools and slash commands stream in through the existing live-refresh channel once each server connects (non-UI modes keep the blocking path). ~290ms off first paint with MCP servers configured. - Build the model catalog's canonical-equivalence index lazily on first read instead of eagerly in the ModelRegistry constructor. A default interactive launch never reads it pre-paint, moving the ~210ms build (over ~3,200 models) off the critical path: ~244ms (~16%) off cold boot. - Memoize per-session read summaries (tree-sitter parse) on the content hash of the freshly-read bytes; the file is still read fresh each call so results stay correct. Repeat same-file summary read 17ms -> 2.4ms. - Reuse the Markdown subtree across streaming reveal ticks, memoize grapheme counting, and stop re-highlighting finalized thinking blocks. - Attribute the previously-unlabeled synchronous boot region in the PI_TIMING table and add a bench:guard boot-regression target.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -2,6 +2,14 @@
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### 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 MCP server instructions join the system prompt on the next session rather than retroactively. 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.8] - 2026-06-09
|
||||
|
||||
### Added
|
||||
|
||||
@@ -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<unknown>): Promise<number> {
|
||||
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})`);
|
||||
}
|
||||
|
||||
@@ -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",
|
||||
"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:",
|
||||
|
||||
Executable
+71
@@ -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);
|
||||
@@ -899,6 +899,7 @@ function getConfiguredProviderOrderFromSettings(): string[] {
|
||||
export class ModelRegistry {
|
||||
#models: Model<Api>[] = [];
|
||||
#canonicalIndex: CanonicalModelIndex = { records: [], byId: new Map(), bySelector: new Map() };
|
||||
#canonicalIndexDirty: boolean = true;
|
||||
#customProviderApiKeys: Map<string, string> = new Map();
|
||||
#keylessProviders: Set<string> = new Set();
|
||||
#discoverableProviders: DiscoveryProviderConfig[] = [];
|
||||
@@ -2158,10 +2159,25 @@ export class ModelRegistry {
|
||||
this.#rebuildPending = true;
|
||||
return;
|
||||
}
|
||||
this.#canonicalIndex = buildCanonicalModelIndex(this.#models, 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, this.#equivalenceConfig),
|
||||
);
|
||||
this.#canonicalIndexDirty = false;
|
||||
}
|
||||
return this.#canonicalIndex;
|
||||
}
|
||||
|
||||
#suspendRebuild(): void {
|
||||
this.#rebuildSuspended += 1;
|
||||
}
|
||||
@@ -2172,7 +2188,7 @@ export class ModelRegistry {
|
||||
}
|
||||
if (this.#rebuildSuspended === 0 && this.#rebuildPending) {
|
||||
this.#rebuildPending = false;
|
||||
this.#canonicalIndex = buildCanonicalModelIndex(this.#models, this.#equivalenceConfig);
|
||||
this.#canonicalIndexDirty = true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2290,7 +2306,7 @@ export class ModelRegistry {
|
||||
|
||||
getCanonicalModels(options?: CanonicalModelQueryOptions): CanonicalModelRecord[] {
|
||||
const records: CanonicalModelRecord[] = [];
|
||||
for (const record of this.#canonicalIndex.records) {
|
||||
for (const record of this.#ensureCanonicalIndex().records) {
|
||||
const variants = this.#filterCanonicalVariants(record, options);
|
||||
if (variants.length === 0) {
|
||||
continue;
|
||||
@@ -2305,7 +2321,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 [];
|
||||
}
|
||||
@@ -2322,7 +2338,7 @@ export class ModelRegistry {
|
||||
}
|
||||
|
||||
getCanonicalId(model: Model<Api>): string | undefined {
|
||||
return this.#canonicalIndex.bySelector.get(formatCanonicalVariantSelector(model).toLowerCase());
|
||||
return this.#ensureCanonicalIndex().bySelector.get(formatCanonicalVariantSelector(model).toLowerCase());
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -804,7 +804,7 @@ export async function runRootCommand(
|
||||
|
||||
// Create AuthStorage and ModelRegistry upfront
|
||||
const authStorage = await logger.time("discoverModels", 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`);
|
||||
@@ -1045,7 +1045,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 = {};
|
||||
}
|
||||
|
||||
@@ -36,6 +36,11 @@ export class AssistantMessageComponent extends Container {
|
||||
* transcript keeps the error in history.
|
||||
*/
|
||||
#errorPinned = 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,
|
||||
@@ -211,12 +216,107 @@ 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): 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;
|
||||
}
|
||||
// Shape is identical — setText only on Markdown children whose source changed.
|
||||
for (const item of this.#fastPathItems) {
|
||||
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): void {
|
||||
this.#lastMessage = message;
|
||||
|
||||
// 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()) || (c.type === "thinking" && c.thinking.trim()),
|
||||
);
|
||||
@@ -228,7 +328,10 @@ 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
|
||||
this.#contentContainer.addChild(new Markdown(content.text.trim(), 1, 0, getMarkdownTheme()));
|
||||
const trimmed = content.text.trim();
|
||||
const md = new Markdown(trimmed, 1, 0, getMarkdownTheme());
|
||||
this.#contentContainer.addChild(md);
|
||||
captureItems?.push({ md, contentIndex: i, blockType: "text", lastText: trimmed });
|
||||
} else if (content.type === "thinking" && content.thinking.trim()) {
|
||||
// Add spacing only when another visible assistant content block follows.
|
||||
// This avoids a superfluous blank line before separately-rendered tool execution blocks.
|
||||
@@ -245,12 +348,12 @@ export class AssistantMessageComponent extends Container {
|
||||
} else {
|
||||
const thinkingText = content.thinking.trim();
|
||||
// Thinking traces in thinkingText color, italic
|
||||
this.#contentContainer.addChild(
|
||||
new Markdown(thinkingText, 1, 0, getMarkdownTheme(), {
|
||||
color: (text: string) => theme.fg("thinkingText", text),
|
||||
italic: true,
|
||||
}),
|
||||
);
|
||||
const md = new Markdown(thinkingText, 1, 0, getMarkdownTheme(), {
|
||||
color: (text: string) => theme.fg("thinkingText", text),
|
||||
italic: true,
|
||||
});
|
||||
this.#contentContainer.addChild(md);
|
||||
captureItems?.push({ md, contentIndex: i, blockType: "thinking", lastText: thinkingText });
|
||||
this.#appendThinkingExtensions(i, thinkingIndex, thinkingText);
|
||||
thinkingIndex += 1;
|
||||
if (hasVisibleContentAfter) {
|
||||
@@ -299,5 +402,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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<string, number>({ 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -87,7 +87,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 { resolveMemoryBackend } from "./memory-backend";
|
||||
import type { MnemopiSessionState } from "./mnemopi/state";
|
||||
import asyncResultTemplate from "./prompts/tools/async-result.md" with { type: "text" };
|
||||
@@ -141,6 +148,7 @@ import {
|
||||
type DiscoverableTool,
|
||||
filterBySource,
|
||||
formatDiscoverableToolServerSummary,
|
||||
isMCPToolName,
|
||||
selectDiscoverableToolNamesByServer,
|
||||
summarizeDiscoverableTools,
|
||||
} from "./tool-discovery/tool-index";
|
||||
@@ -296,6 +304,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<string>();
|
||||
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() */
|
||||
@@ -1456,46 +1534,89 @@ 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 includeMCPServerInstructionsInPrompt = !deferMCPDiscoveryForUI;
|
||||
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),
|
||||
);
|
||||
applyMCPEnvironment(mcpResult);
|
||||
logMCPLoadErrors(mcpResult.errors);
|
||||
await liveSession.refreshMCPTools(mcpResult.tools, {
|
||||
activateAll: activation.activateAllMCPTools,
|
||||
});
|
||||
if (!activation.mcpDiscoveryEnabled && activation.explicitlyRequestedMCPToolNames.length > 0) {
|
||||
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);
|
||||
|
||||
@@ -1748,6 +1869,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.
|
||||
@@ -1838,8 +1967,13 @@ 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
|
||||
const serverInstructions = mcpManager?.getServerInstructions();
|
||||
// Build combined append prompt: memory instructions + MCP server instructions.
|
||||
// Interactive deferred MCP intentionally omits MCP instructions permanently;
|
||||
// prompt rebuilds caused by later tool refreshes keep cache keys stable by
|
||||
// updating tool metadata only, not server-provided prompt text.
|
||||
const serverInstructions = includeMCPServerInstructionsInPrompt
|
||||
? mcpManager?.getServerInstructions()
|
||||
: undefined;
|
||||
let appendPrompt: string | undefined = memoryInstructions ?? undefined;
|
||||
if (serverInstructions && serverInstructions.size > 0) {
|
||||
const parts: string[] = [];
|
||||
@@ -2202,20 +2336,23 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
|
||||
rebuildSystemPrompt,
|
||||
reloadSshTool,
|
||||
requestedToolNames: requestedToolNameSet,
|
||||
getMcpServerInstructions: mcpManager
|
||||
? () => {
|
||||
const raw = mcpManager.getServerInstructions();
|
||||
if (!raw || raw.size === 0) return raw;
|
||||
const out = new Map<string, string>();
|
||||
for (const [name, text] of raw) {
|
||||
out.set(
|
||||
name,
|
||||
text.length > MAX_MCP_INSTRUCTIONS_LENGTH ? text.slice(0, MAX_MCP_INSTRUCTIONS_LENGTH) : text,
|
||||
);
|
||||
getMcpServerInstructions:
|
||||
includeMCPServerInstructionsInPrompt && mcpManager
|
||||
? () => {
|
||||
const raw = mcpManager.getServerInstructions();
|
||||
if (!raw || raw.size === 0) return raw;
|
||||
const out = new Map<string, string>();
|
||||
for (const [name, text] of raw) {
|
||||
out.set(
|
||||
name,
|
||||
text.length > MAX_MCP_INSTRUCTIONS_LENGTH
|
||||
? text.slice(0, MAX_MCP_INSTRUCTIONS_LENGTH)
|
||||
: text,
|
||||
);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
: undefined,
|
||||
: undefined,
|
||||
mcpDiscoveryEnabled,
|
||||
initialSelectedMCPToolNames,
|
||||
defaultSelectedMCPToolNames,
|
||||
@@ -2349,7 +2486,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 => {
|
||||
@@ -2382,6 +2538,12 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
|
||||
});
|
||||
}
|
||||
|
||||
startDeferredMCPDiscovery?.(session, {
|
||||
mcpDiscoveryEnabled,
|
||||
explicitlyRequestedMCPToolNames,
|
||||
activateAllMCPTools: !mcpDiscoveryEnabled && options.toolNames === undefined,
|
||||
});
|
||||
|
||||
return {
|
||||
session,
|
||||
extensionsResult,
|
||||
|
||||
@@ -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,
|
||||
@@ -99,6 +100,24 @@ 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.
|
||||
const SUMMARY_CACHE_MAX = 48;
|
||||
const summaryParseCaches = new WeakMap<object, LRUCache<string, SummaryResult>>();
|
||||
function getSummaryParseCache(session: object): LRUCache<string, SummaryResult> {
|
||||
let cache = summaryParseCaches.get(session);
|
||||
if (!cache) {
|
||||
cache = new LRUCache<string, SummaryResult>({ 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"]);
|
||||
|
||||
@@ -1512,14 +1531,23 @@ export class ReadTool implements AgentTool<typeof readSchema, ReadToolDetails> {
|
||||
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) return memoized;
|
||||
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,
|
||||
});
|
||||
cache.set(cacheKey, result);
|
||||
return result;
|
||||
} catch {
|
||||
return null;
|
||||
|
||||
+141
@@ -0,0 +1,141 @@
|
||||
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";
|
||||
|
||||
const W = 100;
|
||||
|
||||
function msg(content: AssistantMessage["content"], extra: Partial<AssistantMessage> = {}): 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: #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));
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,95 @@
|
||||
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, getBundledModel } from "@oh-my-pi/pi-ai";
|
||||
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();
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user