Merge branch 'pr-2182'
# Conflicts: # packages/coding-agent/package.json # packages/coding-agent/src/config/model-registry.ts # packages/coding-agent/src/main.ts # packages/coding-agent/src/modes/components/assistant-message.ts # packages/coding-agent/src/sdk.ts # packages/tui/src/components/markdown.ts
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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 && 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:",
|
||||
|
||||
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);
|
||||
@@ -602,6 +602,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[] = [];
|
||||
@@ -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<Api>): string | undefined {
|
||||
return this.#canonicalIndex.bySelector.get(formatCanonicalVariantSelector(model).toLowerCase());
|
||||
return this.#ensureCanonicalIndex().bySelector.get(formatCanonicalVariantSelector(model).toLowerCase());
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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 = {};
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<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() */
|
||||
@@ -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<AgentTool | null> => {
|
||||
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,
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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<object, LRUCache<string, SummaryResult | false>>();
|
||||
function getSummaryParseCache(session: object): LRUCache<string, SummaryResult | false> {
|
||||
let cache = summaryParseCaches.get(session);
|
||||
if (!cache) {
|
||||
cache = new LRUCache<string, SummaryResult | false>({ 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<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 !== 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;
|
||||
}
|
||||
|
||||
+88
@@ -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<string, unknown>;
|
||||
};
|
||||
|
||||
function buildResult(method: string): Record<string, unknown> {
|
||||
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();
|
||||
}
|
||||
+88
@@ -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 <ms>` 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<string, unknown>;
|
||||
};
|
||||
|
||||
function buildResult(method: string): Record<string, unknown> {
|
||||
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();
|
||||
}
|
||||
+182
@@ -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> = {}): 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));
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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
|
||||
|
||||
@@ -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<string, readonly string[]>({ 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[] = [];
|
||||
|
||||
@@ -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));
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user