feat: added google reasoning controls mcp stream resumption and tar support
- Added Google provider thinking configuration parameters and force-reasoning-off controls. - Implemented MCP SSE stream resumption using Last-Event-ID and `SSEResumeError`. - Added support for TAR old-GNU sparse extension blocks, path length checks, and archive entry overrides. - Restricted external thinking support to specific models and added semver fallback parsing.
This commit is contained in:
@@ -858,13 +858,18 @@ export function buildGoogleGenerateContentParams<T extends "google-generative-ai
|
||||
config.toolConfig = undefined;
|
||||
}
|
||||
|
||||
if (options.thinking?.enabled && model.reasoning) {
|
||||
const cfg: ThinkingConfig = { includeThoughts: !options.hideThinkingSummary };
|
||||
if (options.thinking.level !== undefined) {
|
||||
// GoogleThinkingLevel mirrors the SDK's `ThinkingLevel` string enum values 1:1.
|
||||
cfg.thinkingLevel = options.thinking.level as ThinkingLevel;
|
||||
} else if (options.thinking.budgetTokens !== undefined) {
|
||||
cfg.thinkingBudget = options.thinking.budgetTokens;
|
||||
const thinking = options.thinking;
|
||||
if (
|
||||
thinking &&
|
||||
model.reasoning &&
|
||||
(thinking.enabled || thinking.level !== undefined || thinking.budgetTokens !== undefined)
|
||||
) {
|
||||
const cfg: ThinkingConfig = { includeThoughts: thinking.enabled && !options.hideThinkingSummary };
|
||||
if (thinking.level !== undefined) {
|
||||
// GoogleThinkingLevel mirrors the SDK's ThinkingLevel string enum values 1:1.
|
||||
cfg.thinkingLevel = thinking.level as ThinkingLevel;
|
||||
} else if (thinking.budgetTokens !== undefined) {
|
||||
cfg.thinkingBudget = thinking.budgetTokens;
|
||||
}
|
||||
config.thinkingConfig = cfg;
|
||||
}
|
||||
|
||||
@@ -1408,6 +1408,17 @@ function resolveOpenAiReasoningEffort<TApi extends Api>(
|
||||
return requireSupportedEffort(model, reasoning);
|
||||
}
|
||||
|
||||
function resolveGoogleThinkingOff<TApi extends Api>(model: Model<TApi>): NonNullable<GoogleOptions["thinking"]> {
|
||||
const thinking: NonNullable<GoogleOptions["thinking"]> = { enabled: false };
|
||||
if (!model.reasoning || !model.thinking) return thinking;
|
||||
if (model.thinking.mode === "budget" && (!model.thinking.requiresEffort || model.thinking.suppressWhenOff)) {
|
||||
thinking.budgetTokens = 0;
|
||||
} else if (model.thinking.mode === "google-level" && model.thinking.suppressWhenOff) {
|
||||
thinking.level = "MINIMAL";
|
||||
}
|
||||
return thinking;
|
||||
}
|
||||
|
||||
const castApi = <TApi extends Api>(api: OptionsForApi<TApi>): OptionsForApi<Api> => api as OptionsForApi<Api>;
|
||||
|
||||
/**
|
||||
@@ -1428,13 +1439,13 @@ function normalizeMandatoryReasoningOptions<TApi extends Api>(
|
||||
!model.reasoning ||
|
||||
!model.thinking?.requiresEffort ||
|
||||
model.thinking.suppressWhenOff ||
|
||||
(options?.reasoning !== undefined && !options.disableReasoning)
|
||||
(options?.reasoning !== undefined && !options.disableReasoning && !options.forceReasoningOff)
|
||||
) {
|
||||
return options;
|
||||
}
|
||||
const floor = minimumSupportedEffort(model);
|
||||
if (floor === undefined) return options;
|
||||
return { ...options, reasoning: floor, disableReasoning: undefined };
|
||||
return { ...options, reasoning: floor, disableReasoning: undefined, forceReasoningOff: undefined };
|
||||
}
|
||||
|
||||
function supportsExplicitOpenAIResponsesPromptCache(compat: unknown): boolean {
|
||||
@@ -1732,7 +1743,7 @@ function mapOptionsForApi<TApi extends Api>(
|
||||
return castApi<"google-generative-ai">({
|
||||
...base,
|
||||
serviceTier: options?.serviceTier,
|
||||
thinking: { enabled: false },
|
||||
thinking: resolveGoogleThinkingOff(model),
|
||||
toolChoice: mapGoogleToolChoice(options?.toolChoice),
|
||||
cachedContent: options?.cachedContent,
|
||||
});
|
||||
@@ -1838,7 +1849,7 @@ function mapOptionsForApi<TApi extends Api>(
|
||||
return castApi<"google-vertex">({
|
||||
...base,
|
||||
serviceTier: options?.serviceTier,
|
||||
thinking: { enabled: false },
|
||||
thinking: resolveGoogleThinkingOff(model),
|
||||
toolChoice: mapGoogleToolChoice(options?.toolChoice),
|
||||
cachedContent: options?.cachedContent,
|
||||
});
|
||||
|
||||
@@ -150,7 +150,7 @@ describe("google-gemini-cli effort-tier variant routing", () => {
|
||||
});
|
||||
});
|
||||
|
||||
it("routes to the explicit off wire shape when an external work log replaces reasoning", async () => {
|
||||
it("routes to the explicit off wire shape when an external scratchpad replaces reasoning", async () => {
|
||||
const off = await captureRequest(collapsedFlashModel(), Effort.High, { forceReasoningOff: true });
|
||||
expect(off.body.model).toBe("gemini-3.5-flash-extra-low");
|
||||
expect(off.body.request?.generationConfig?.thinkingConfig).toEqual({
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import { Effort, type FetchImpl } from "@oh-my-pi/pi-ai";
|
||||
import { streamSimple } from "@oh-my-pi/pi-ai/stream";
|
||||
import type { Context, Model } from "@oh-my-pi/pi-ai/types";
|
||||
import { buildModel } from "@oh-my-pi/pi-catalog/build";
|
||||
|
||||
interface CapturedPayload {
|
||||
config?: {
|
||||
thinkingConfig?: {
|
||||
includeThoughts?: boolean;
|
||||
thinkingBudget?: number;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
const context: Context = {
|
||||
messages: [{ role: "user", content: "hello", timestamp: Date.now() }],
|
||||
};
|
||||
|
||||
const model: Model<"google-generative-ai"> = buildModel({
|
||||
id: "gemini-2.5-flash",
|
||||
name: "Gemini 2.5 Flash",
|
||||
api: "google-generative-ai",
|
||||
provider: "google",
|
||||
baseUrl: "https://generativelanguage.googleapis.com",
|
||||
reasoning: true,
|
||||
thinking: {
|
||||
mode: "budget",
|
||||
efforts: [Effort.Minimal, Effort.Low, Effort.Medium, Effort.High],
|
||||
},
|
||||
input: ["text"],
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
|
||||
contextWindow: 1_000_000,
|
||||
maxTokens: 65_536,
|
||||
});
|
||||
|
||||
async function capturePayload(flag: "disableReasoning" | "forceReasoningOff"): Promise<CapturedPayload> {
|
||||
let captured: CapturedPayload | undefined;
|
||||
const fetchMock: FetchImpl = async () =>
|
||||
new Response("", { status: 200, headers: { "content-type": "text/event-stream" } });
|
||||
|
||||
await streamSimple(model, context, {
|
||||
apiKey: "test-key",
|
||||
reasoning: Effort.High,
|
||||
[flag]: true,
|
||||
fetch: fetchMock,
|
||||
onPayload: payload => {
|
||||
captured = payload as CapturedPayload;
|
||||
},
|
||||
}).result();
|
||||
|
||||
if (!captured) throw new Error("Google request payload was not captured");
|
||||
return captured;
|
||||
}
|
||||
|
||||
describe("Google reasoning disablement", () => {
|
||||
for (const flag of ["disableReasoning", "forceReasoningOff"] as const) {
|
||||
it(`sends an explicit zero thinking budget for ${flag}`, async () => {
|
||||
const payload = await capturePayload(flag);
|
||||
|
||||
expect(payload.config?.thinkingConfig).toEqual({
|
||||
includeThoughts: false,
|
||||
thinkingBudget: 0,
|
||||
});
|
||||
});
|
||||
}
|
||||
});
|
||||
@@ -35,7 +35,7 @@ function openRouterModel(thinking: ThinkingConfig): Model<"openai-completions">
|
||||
|
||||
async function captureBody(
|
||||
model: Model<"openai-completions">,
|
||||
options: { reasoning?: Effort; disableReasoning?: boolean },
|
||||
options: { reasoning?: Effort; disableReasoning?: boolean; forceReasoningOff?: boolean },
|
||||
): Promise<CapturedBody> {
|
||||
let requestBody: string | undefined;
|
||||
const fetchMock: FetchImpl = (_input, init) => {
|
||||
@@ -67,6 +67,14 @@ describe("thinking.requiresEffort clamping", () => {
|
||||
expect(body.reasoning).toEqual({ effort: "minimal" });
|
||||
});
|
||||
|
||||
it("clamps forceReasoningOff when the endpoint cannot disable reasoning", async () => {
|
||||
const body = await captureBody(openRouterModel(MANDATORY_THINKING), {
|
||||
reasoning: Effort.High,
|
||||
forceReasoningOff: true,
|
||||
});
|
||||
expect(body.reasoning).toEqual({ effort: "minimal" });
|
||||
});
|
||||
|
||||
it("keeps explicit efforts untouched", async () => {
|
||||
const body = await captureBody(openRouterModel(MANDATORY_THINKING), { reasoning: Effort.High });
|
||||
expect(body.reasoning).toEqual({ effort: "high" });
|
||||
|
||||
@@ -181,7 +181,9 @@ function createSemVer(major: number, minor: number, patch = 0): SemVer {
|
||||
return { major, minor, patch };
|
||||
}
|
||||
|
||||
// extend this table if we need anything more than 9.10
|
||||
// Fast path for the common 1–2 component versions; anything the table misses
|
||||
// (large minors, 3-part versions) parses dynamically below so no future
|
||||
// version ever classifies as unknown (the failure class #8256 fixed).
|
||||
const precomputeTable: Record<string, SemVer> = {};
|
||||
for (let major = 0; major <= 9; major++) {
|
||||
for (let minor = 0; minor <= 10; minor++) {
|
||||
@@ -192,8 +194,14 @@ for (let major = 0; major <= 9; major++) {
|
||||
precomputeTable[`${major}`] = createSemVer(major, 0, 0);
|
||||
}
|
||||
|
||||
const SEMVER_PATTERN = /^(\d{1,2})(?:[.-](\d{1,2}))?(?:[.-](\d{1,2}))?$/;
|
||||
|
||||
export function parseSemVer(version: string): SemVer | null {
|
||||
return precomputeTable[version] ?? null;
|
||||
const hit = precomputeTable[version];
|
||||
if (hit) return hit;
|
||||
const match = SEMVER_PATTERN.exec(version);
|
||||
if (!match) return null;
|
||||
return createSemVer(Number(match[1]), Number(match[2] ?? 0), Number(match[3] ?? 0));
|
||||
}
|
||||
|
||||
export function semverGte(left: SemVer | string, right: SemVer | string): boolean {
|
||||
|
||||
@@ -76,6 +76,21 @@ describe("parseAnthropicModel", () => {
|
||||
});
|
||||
expect(parseAnthropicModel("anthropic--claude-4.8-haiku")).toBeNull();
|
||||
});
|
||||
|
||||
test("parses versions past the precompute table instead of classifying the model unknown", () => {
|
||||
// The semver precompute table gates parsing; a too-small bound silently
|
||||
// downgraded `claude-opus-5-11`-shaped ids to unknown (#8256 class).
|
||||
expect(parseAnthropicModel("claude-opus-5-11")).toEqual({
|
||||
family: "anthropic",
|
||||
kind: "opus",
|
||||
version: { major: 5, minor: 11, patch: 0 },
|
||||
});
|
||||
expect(parseAnthropicModel("claude-sonnet-4.25")).toEqual({
|
||||
family: "anthropic",
|
||||
kind: "sonnet",
|
||||
version: { major: 4, minor: 25, patch: 0 },
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
describe("supportsAdaptiveThinkingDisplay", () => {
|
||||
|
||||
@@ -13,7 +13,7 @@
|
||||
|
||||
### Changed
|
||||
|
||||
- Renamed the `think` tool's `thoughts` parameter to `notes` and restricted tool availability to models supporting external thinking
|
||||
- Restricted the `think` tool to GPT Responses, Claude, and Gemini transports that can replace native reasoning, and kept its streamed scratchpad input named `thoughts`.
|
||||
- `omp cleanse` default subagent cap raised from 8 to 32 (`--agents`/`-n` still overrides).
|
||||
|
||||
### Fixed
|
||||
|
||||
@@ -368,7 +368,8 @@ export const agenticFixtures: Record<string, GalleryFixture> = {
|
||||
thoughts: "The retry loop re-reads the config after every failure, which explains the doubled latency.",
|
||||
},
|
||||
args: {
|
||||
thoughts: "The retry loop re-reads the config after every failure, which explains the doubled latency. Cache the parsed config outside the loop, then re-check the invalidation path.",
|
||||
thoughts:
|
||||
"The retry loop re-reads the config after every failure, which explains the doubled latency. Cache the parsed config outside the loop, then re-check the invalidation path.",
|
||||
},
|
||||
result: {
|
||||
content: [{ type: "text", text: "------" }],
|
||||
|
||||
@@ -30,6 +30,13 @@ interface SSEResumeState {
|
||||
retryMs: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Failure resuming an accepted request's logical SSE stream. Carries a
|
||||
* never-replay contract: by resume time the server has accepted (and possibly
|
||||
* executed) the originating POST, so auth-retry paths must not re-send it.
|
||||
*/
|
||||
class SSEResumeError extends Error {}
|
||||
|
||||
/** Wait for the server-provided SSE retry interval while remaining abortable. */
|
||||
async function waitForSSERetry(ms: number, signal: AbortSignal): Promise<void> {
|
||||
if (signal.aborted) throw signal.reason;
|
||||
@@ -195,7 +202,7 @@ export class HttpTransport implements MCPTransport {
|
||||
// If the stream ends unexpectedly (server restart, network drop),
|
||||
// fire onClose so the manager can trigger reconnection.
|
||||
const signal = connection.signal;
|
||||
void this.#readSSEStream(response.body!, signal).finally(() => {
|
||||
void this.#runSSEListener(response.body!, signal).finally(() => {
|
||||
const wasConnected = this.#connected;
|
||||
if (this.#sseConnection === connection) this.#sseConnection = null;
|
||||
if (wasConnected) this.onClose?.();
|
||||
@@ -215,6 +222,95 @@ export class HttpTransport implements MCPTransport {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the long-lived GET SSE stream, resuming with `Last-Event-ID` when
|
||||
* the server closes the physical connection mid-stream (2025-11-25 permits
|
||||
* polling-style servers). Returns only when the logical stream ends — the
|
||||
* caller fires `onClose` and the manager's reconnect path takes over. A
|
||||
* resume cycle that delivers no events before dropping again ends the
|
||||
* stream rather than retrying forever against a broken server.
|
||||
*/
|
||||
async #runSSEListener(initialBody: ReadableStream<Uint8Array>, signal: AbortSignal): Promise<void> {
|
||||
const resume: SSEResumeState = { lastEventId: null, retryMs: DEFAULT_SSE_RETRY_MS };
|
||||
let body = initialBody;
|
||||
let progressed = true;
|
||||
for (;;) {
|
||||
try {
|
||||
for await (const event of readSseEvents(body, signal)) {
|
||||
progressed = true;
|
||||
if (event.id !== undefined) resume.lastEventId = event.id || null;
|
||||
if (event.retry !== undefined) resume.retryMs = event.retry;
|
||||
if (event.data === "") continue;
|
||||
if (!this.#connected) return;
|
||||
this.#dispatchSSEMessage(JSON.parse(event.data) as JsonRpcMessage | JsonRpcMessage[]);
|
||||
}
|
||||
} catch (error) {
|
||||
if (error instanceof Error && error.name === "AbortError") return;
|
||||
logger.debug("HTTP SSE stream error", {
|
||||
url: this.config.url,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
if (resume.lastEventId === null) {
|
||||
if (error instanceof Error) this.onError?.(error);
|
||||
return;
|
||||
}
|
||||
}
|
||||
if (!this.#connected || signal.aborted || resume.lastEventId === null || !progressed) return;
|
||||
progressed = false;
|
||||
try {
|
||||
const response = await this.#fetchSSEResume(resume, signal);
|
||||
body = response.body as ReadableStream<Uint8Array>;
|
||||
} catch (error) {
|
||||
if (!(error instanceof Error && error.name === "AbortError")) {
|
||||
logger.debug("HTTP SSE listener resume failed", {
|
||||
url: this.config.url,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
}
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resume a logical SSE stream via GET + `Last-Event-ID`, honoring the
|
||||
* server-provided retry interval and refreshing auth once on 401/403.
|
||||
* Failures throw {@link SSEResumeError} so `request()` never replays the
|
||||
* originating POST in response.
|
||||
*/
|
||||
async #fetchSSEResume(resume: SSEResumeState, signal: AbortSignal): Promise<Response> {
|
||||
if (resume.lastEventId === null) {
|
||||
throw new SSEResumeError("SSE stream ended without a resumable event ID");
|
||||
}
|
||||
await waitForSSERetry(resume.retryMs, signal);
|
||||
const generated: Record<string, string> = {
|
||||
Accept: "text/event-stream",
|
||||
"Last-Event-ID": resume.lastEventId,
|
||||
};
|
||||
if (this.#sessionId) generated["Mcp-Session-Id"] = this.#sessionId;
|
||||
let response = await this.#fetch({ method: "GET", signal }, generated);
|
||||
if (this.onAuthError && (response.status === 401 || response.status === 403)) {
|
||||
await response.body?.cancel();
|
||||
const newHeaders = await this.onAuthError();
|
||||
if (!newHeaders) {
|
||||
throw new SSEResumeError(`HTTP ${response.status} resuming MCP SSE stream: auth refresh failed`);
|
||||
}
|
||||
// Persist refreshed headers so subsequent requests use them directly
|
||||
this.config = { ...this.config, headers: newHeaders };
|
||||
response = await this.#fetch({ method: "GET", signal }, generated);
|
||||
}
|
||||
if (!response.ok) {
|
||||
const text = await response.text().catch(() => "");
|
||||
throw new SSEResumeError(`HTTP ${response.status} resuming MCP SSE stream: ${text}`);
|
||||
}
|
||||
const contentType = response.headers.get("Content-Type") ?? "";
|
||||
if (!contentType.includes("text/event-stream") || !response.body) {
|
||||
await response.body?.cancel();
|
||||
throw new SSEResumeError(`MCP SSE resume returned unsupported Content-Type: ${contentType || "(missing)"}`);
|
||||
}
|
||||
return response;
|
||||
}
|
||||
|
||||
/** Route an SSE message (or batch) to the appropriate handler. */
|
||||
#dispatchSSEMessage(message: JsonRpcMessage | JsonRpcMessage[]): void {
|
||||
if (Array.isArray(message)) {
|
||||
@@ -240,9 +336,12 @@ export class HttpTransport implements MCPTransport {
|
||||
try {
|
||||
return await this.#executeRequest<T>(method, params, options);
|
||||
} catch (error) {
|
||||
// Retry once on auth failure if onAuthError is wired
|
||||
// Retry once on auth failure if onAuthError is wired. Never replay
|
||||
// after an SSE resume failure: the server already accepted the
|
||||
// original POST and may have executed it — replaying could run a
|
||||
// state-changing tool twice.
|
||||
const status = error instanceof Error ? AIError.status(error) : undefined;
|
||||
if (this.onAuthError && (status === 401 || status === 403)) {
|
||||
if (!(error instanceof SSEResumeError) && this.onAuthError && (status === 401 || status === 403)) {
|
||||
const newHeaders = await this.onAuthError();
|
||||
if (newHeaders) {
|
||||
// Persist refreshed headers so subsequent requests use them directly
|
||||
@@ -355,53 +454,49 @@ export class HttpTransport implements MCPTransport {
|
||||
try {
|
||||
for (;;) {
|
||||
if (!current.body) throw new Error("SSE response did not include a body");
|
||||
for await (const event of readSseEvents(current.body, signal)) {
|
||||
if (event.id !== undefined) resume.lastEventId = event.id || null;
|
||||
if (event.retry !== undefined) resume.retryMs = event.retry;
|
||||
if (event.data === "") continue;
|
||||
const raw = JSON.parse(event.data) as JsonRpcMessage | JsonRpcMessage[];
|
||||
const messages = Array.isArray(raw) ? raw : [raw];
|
||||
for (const message of messages) {
|
||||
if (
|
||||
!captured &&
|
||||
"id" in message &&
|
||||
message.id === expectedId &&
|
||||
("result" in message || "error" in message)
|
||||
) {
|
||||
captured = true;
|
||||
operation.clear();
|
||||
if (message.error) {
|
||||
reject(new Error(`MCP error ${message.error.code}: ${message.error.message}`));
|
||||
} else {
|
||||
resolve(message.result as T);
|
||||
try {
|
||||
for await (const event of readSseEvents(current.body, signal)) {
|
||||
if (event.id !== undefined) resume.lastEventId = event.id || null;
|
||||
if (event.retry !== undefined) resume.retryMs = event.retry;
|
||||
if (event.data === "") continue;
|
||||
const raw = JSON.parse(event.data) as JsonRpcMessage | JsonRpcMessage[];
|
||||
const messages = Array.isArray(raw) ? raw : [raw];
|
||||
for (const message of messages) {
|
||||
if (
|
||||
!captured &&
|
||||
"id" in message &&
|
||||
message.id === expectedId &&
|
||||
("result" in message || "error" in message)
|
||||
) {
|
||||
captured = true;
|
||||
operation.clear();
|
||||
if (message.error) {
|
||||
reject(new Error(`MCP error ${message.error.code}: ${message.error.message}`));
|
||||
} else {
|
||||
resolve(message.result as T);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
continue;
|
||||
if (!this.#connected) continue;
|
||||
this.#dispatchSSEMessage(message);
|
||||
}
|
||||
if (!this.#connected) continue;
|
||||
this.#dispatchSSEMessage(message);
|
||||
}
|
||||
} catch (error) {
|
||||
// An abrupt drop (socket reset, body-read failure) is as
|
||||
// resumable as a server-initiated close once an event ID
|
||||
// exists; the request timeout still bounds the total wait.
|
||||
if (captured) return;
|
||||
if (signal.aborted || resume.lastEventId === null) throw error;
|
||||
logger.debug("MCP SSE response stream dropped; resuming", {
|
||||
url: this.config.url,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
}
|
||||
if (captured) return;
|
||||
if (resume.lastEventId === null) {
|
||||
throw new Error(`No response received for request ID ${expectedId}`);
|
||||
}
|
||||
|
||||
await waitForSSERetry(resume.retryMs, signal);
|
||||
const generated: Record<string, string> = {
|
||||
Accept: "text/event-stream",
|
||||
"Last-Event-ID": resume.lastEventId,
|
||||
};
|
||||
if (this.#sessionId) generated["Mcp-Session-Id"] = this.#sessionId;
|
||||
current = await this.#fetch({ method: "GET", signal }, generated);
|
||||
if (!current.ok) {
|
||||
const text = await current.text();
|
||||
throw new Error(`HTTP ${current.status} resuming MCP SSE stream: ${text}`);
|
||||
}
|
||||
const contentType = current.headers.get("Content-Type") ?? "";
|
||||
if (!contentType.includes("text/event-stream")) {
|
||||
await current.body?.cancel();
|
||||
throw new Error(`MCP SSE resume returned unsupported Content-Type: ${contentType || "(missing)"}`);
|
||||
}
|
||||
current = await this.#fetchSSEResume(resume, signal);
|
||||
}
|
||||
} catch (error) {
|
||||
if (captured) return;
|
||||
|
||||
@@ -26,6 +26,18 @@ export class ToolActivityContainer extends Container implements ToolActivityComp
|
||||
this.invalidate();
|
||||
}
|
||||
|
||||
/**
|
||||
* Forward Ctrl+O expansion to wrapped children. The transcript's expansion
|
||||
* traversal only visits top-level children, so the wrapper must proxy or
|
||||
* wrapped renderers would freeze at their insertion-time expansion state.
|
||||
*/
|
||||
setExpanded(expanded: boolean): void {
|
||||
for (const child of this.children) {
|
||||
const expandable = child as Partial<{ setExpanded(expanded: boolean): void }>;
|
||||
if (typeof expandable.setExpanded === "function") expandable.setExpanded(expanded);
|
||||
}
|
||||
}
|
||||
|
||||
override render(width: number): readonly string[] {
|
||||
if (!this.#visible) return [];
|
||||
return super.render(width);
|
||||
|
||||
@@ -7,14 +7,26 @@ import { getMarkdownTheme, type Theme } from "../modes/theme/theme";
|
||||
|
||||
/** Whether a model transport can suppress native reasoning while private scratchpad thoughts are active. */
|
||||
export function supportsExternalThinking(model: Model | null | undefined): boolean {
|
||||
if (!model) return false;
|
||||
const requiresThinking =
|
||||
model.api === "anthropic-messages" &&
|
||||
model.compat !== undefined &&
|
||||
"requiresThinkingEnabled" in model.compat &&
|
||||
model.compat.requiresThinkingEnabled === true;
|
||||
if (
|
||||
model.reasoning &&
|
||||
(requiresThinking || (model.thinking?.requiresEffort && !model.thinking.suppressWhenOff))
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
if (model.api === "google-generative-ai" || model.api === "google-gemini-cli" || model.api === "google-vertex") {
|
||||
return !model.reasoning || model.thinking?.mode === "budget" || model.thinking?.suppressWhenOff === true;
|
||||
}
|
||||
return (
|
||||
model?.api === "openai-responses" ||
|
||||
model?.api === "azure-openai-responses" ||
|
||||
model?.api === "openai-codex-responses" ||
|
||||
model?.api === "anthropic-messages" ||
|
||||
model?.api === "google-generative-ai" ||
|
||||
model?.api === "google-gemini-cli" ||
|
||||
model?.api === "google-vertex"
|
||||
model.api === "openai-responses" ||
|
||||
model.api === "azure-openai-responses" ||
|
||||
model.api === "openai-codex-responses" ||
|
||||
model.api === "anthropic-messages"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -219,11 +219,14 @@ function upsertArchiveEntry(map: Map<string, ArchiveIndexEntry>, entry: ArchiveI
|
||||
return;
|
||||
}
|
||||
|
||||
// Same-kind duplicate: the later record wins (tar append/update semantics,
|
||||
// matching system tar extraction and whole-archive materialization), while
|
||||
// earlier metadata fills any gaps the newer record leaves.
|
||||
map.set(entry.path, {
|
||||
...existing,
|
||||
size: existing.size || entry.size,
|
||||
mtimeMs: existing.mtimeMs ?? entry.mtimeMs,
|
||||
storage: existing.storage ?? entry.storage,
|
||||
...entry,
|
||||
size: entry.size || existing.size,
|
||||
mtimeMs: entry.mtimeMs ?? existing.mtimeMs,
|
||||
storage: entry.storage ?? existing.storage,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -628,6 +631,14 @@ const TAR_LINKNAME_OFFSET = 157;
|
||||
const TAR_LINKNAME_LENGTH = 100;
|
||||
const TAR_PREFIX_OFFSET = 345;
|
||||
const TAR_PREFIX_LENGTH = 155;
|
||||
// Old-GNU sparse header: `isextended` flag inside the main header and inside
|
||||
// each 512-byte sparse-map continuation block that follows it.
|
||||
const TAR_GNU_SPARSE_ISEXTENDED_OFFSET = 482;
|
||||
const TAR_GNU_SPARSE_CONT_ISEXTENDED_OFFSET = 504;
|
||||
// PATH_MAX-style bound on member paths and link targets. Real archives never
|
||||
// exceed it (system tar cannot extract them), and it caps every prefix walk
|
||||
// below so crafted multi-hundred-KiB PAX paths cannot pin the CPU.
|
||||
const TAR_MAX_PATH_BYTES = 4096;
|
||||
const GZIP_MAGIC_0 = 0x1f;
|
||||
const GZIP_MAGIC_1 = 0x8b;
|
||||
const TAR_TEXT_DECODER = new TextDecoder();
|
||||
@@ -692,7 +703,11 @@ function tarChecksumMatches(buffer: Uint8Array, offset: number): boolean {
|
||||
return stored === unsigned || stored === signed;
|
||||
}
|
||||
|
||||
/** Parse a PAX extended-header payload into its `key → value` records. */
|
||||
/**
|
||||
* Parse a PAX extended-header payload into its `key → value` records. Only
|
||||
* keys the indexer consumes are retained; a crafted header packed with
|
||||
* millions of unique throwaway records must not amplify into heap.
|
||||
*/
|
||||
function parsePaxRecords(data: Uint8Array): Map<string, string> {
|
||||
const attrs = new Map<string, string>();
|
||||
let pos = 0;
|
||||
@@ -715,7 +730,10 @@ function parsePaxRecords(data: Uint8Array): Map<string, string> {
|
||||
const record = data.subarray(space + 1, pos + length - 1);
|
||||
const eq = record.indexOf(0x3d);
|
||||
if (eq >= 0) {
|
||||
attrs.set(TAR_TEXT_DECODER.decode(record.subarray(0, eq)), TAR_TEXT_DECODER.decode(record.subarray(eq + 1)));
|
||||
const key = TAR_TEXT_DECODER.decode(record.subarray(0, eq));
|
||||
if (key === "path" || key === "linkpath" || key === "size" || key.startsWith("GNU.sparse.")) {
|
||||
attrs.set(key, TAR_TEXT_DECODER.decode(record.subarray(eq + 1)));
|
||||
}
|
||||
}
|
||||
pos += length;
|
||||
}
|
||||
@@ -826,6 +844,16 @@ function readTarEntries(rawBytes: Uint8Array): ArchiveIndexEntry[] {
|
||||
if (Number.isFinite(parsed) && parsed >= 0) displaySize = parsed;
|
||||
}
|
||||
const sparse = typeFlag === "S" || paxDeclaresSparse(pax);
|
||||
// Old-GNU sparse members chain extra 512-byte sparse-map blocks between
|
||||
// the main header and the stored data; they are not counted in `size`.
|
||||
// Consume the chain so the data offset and the next header line up.
|
||||
if (typeFlag === "S" && buffer[offset - TAR_BLOCK_SIZE + TAR_GNU_SPARSE_ISEXTENDED_OFFSET] === 1) {
|
||||
let extended = true;
|
||||
while (extended && offset + TAR_BLOCK_SIZE <= buffer.length) {
|
||||
extended = buffer[offset + TAR_GNU_SPARSE_CONT_ISEXTENDED_OFFSET] === 1;
|
||||
offset += TAR_BLOCK_SIZE;
|
||||
}
|
||||
}
|
||||
const dataOffset = offset;
|
||||
const memberDataBlocks = Math.ceil(size / TAR_BLOCK_SIZE) * TAR_BLOCK_SIZE;
|
||||
offset += memberDataBlocks;
|
||||
@@ -836,6 +864,9 @@ function readTarEntries(rawBytes: Uint8Array): ArchiveIndexEntry[] {
|
||||
const isDirectory = typeFlag === "5" || name.endsWith("/");
|
||||
const normalizedPath = normalizeArchiveEntryPath(name);
|
||||
if (!normalizedPath) continue;
|
||||
if (normalizedPath.length > TAR_MAX_PATH_BYTES) {
|
||||
throw new ToolError(`Archive member path exceeds ${TAR_MAX_PATH_BYTES} bytes`);
|
||||
}
|
||||
const mtimeMs = mtime > 0 ? mtime * 1000 : undefined;
|
||||
|
||||
if (isDirectory) {
|
||||
@@ -861,7 +892,7 @@ function readTarEntries(rawBytes: Uint8Array): ArchiveIndexEntry[] {
|
||||
size: 0,
|
||||
mtimeMs,
|
||||
};
|
||||
if (targetPath === undefined) {
|
||||
if (targetPath === undefined || targetPath.length > TAR_MAX_PATH_BYTES) {
|
||||
if (kind === "hard link") {
|
||||
throw new ToolError(`Archive hard link '${normalizedPath}' has an invalid target`);
|
||||
}
|
||||
@@ -902,78 +933,110 @@ function readTarEntries(rawBytes: Uint8Array): ArchiveIndexEntry[] {
|
||||
// Link records carry no data. Resolve file targets after all headers are
|
||||
// indexed; directory symlinks remain one alias node and are traversed lazily
|
||||
// by ArchiveReader so N files behind M aliases never inflate the index to
|
||||
// N×M entries during a root listing.
|
||||
// N×M entries during a root listing. Resolution is a work queue keyed on
|
||||
// blocking links (not a rescan-all fixpoint) so crafted archives with huge
|
||||
// link chains stay linear in dependency edges.
|
||||
if (pendingLinks.size > 0) {
|
||||
const entriesByPath = new Map<string, ArchiveIndexEntry>();
|
||||
for (const entry of entries) entriesByPath.set(entry.path, entry);
|
||||
const unresolved = new Set(pendingLinks.keys());
|
||||
|
||||
// True while the target itself or any directory on its path is a link
|
||||
// that has not been classified yet; such targets must wait a pass so a
|
||||
// file symlink routed through a directory alias is not misjudged
|
||||
// dangling before the alias resolves.
|
||||
const dependsOnUnresolvedLink = (targetPath: string): boolean => {
|
||||
const parts = targetPath.split("/");
|
||||
for (let end = parts.length; end > 0; end--) {
|
||||
const prefixEntry = entriesByPath.get(parts.slice(0, end).join("/"));
|
||||
if (prefixEntry && unresolved.has(prefixEntry)) return true;
|
||||
// Every proper ancestor of a member path: O(1) directory-target checks
|
||||
// instead of scanning all entries per unresolved link.
|
||||
const directoryPrefixes = new Set<string>();
|
||||
for (const entry of entries) {
|
||||
const memberPath = entry.path;
|
||||
for (let cut = memberPath.lastIndexOf("/"); cut > 0; cut = memberPath.lastIndexOf("/", cut - 1)) {
|
||||
const prefix = memberPath.slice(0, cut);
|
||||
if (directoryPrefixes.has(prefix)) break;
|
||||
directoryPrefixes.add(prefix);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
const unresolved = new Set(pendingLinks.keys());
|
||||
// Links deferred behind a still-unclassified link, re-queued when it
|
||||
// settles, so a file symlink routed through a directory alias is not
|
||||
// misjudged dangling before the alias resolves.
|
||||
const dependents = new Map<ArchiveIndexEntry, ArchiveIndexEntry[]>();
|
||||
|
||||
// The first still-unresolved link on `targetPath` (the target itself or
|
||||
// any directory on its path), or null when the target is settled.
|
||||
const findUnresolvedBlocker = (targetPath: string): ArchiveIndexEntry | null => {
|
||||
for (let end = targetPath.length; end > 0; end = targetPath.lastIndexOf("/", end - 1)) {
|
||||
const prefixEntry = entriesByPath.get(targetPath.slice(0, end));
|
||||
if (prefixEntry && unresolved.has(prefixEntry)) return prefixEntry;
|
||||
}
|
||||
return null;
|
||||
};
|
||||
|
||||
while (unresolved.size > 0) {
|
||||
let resolved = 0;
|
||||
for (const entry of unresolved) {
|
||||
const pending = pendingLinks.get(entry)!;
|
||||
if (dependsOnUnresolvedLink(pending.targetPath)) continue;
|
||||
const queue = [...unresolved];
|
||||
while (queue.length > 0) {
|
||||
const entry = queue.pop()!;
|
||||
if (!unresolved.has(entry)) continue;
|
||||
const pending = pendingLinks.get(entry)!;
|
||||
|
||||
// Targets may route through directory aliases classified in an
|
||||
// earlier pass; rewrite before the exact-path lookup.
|
||||
let targetPath = pending.targetPath;
|
||||
// Targets may route through directory aliases classified earlier;
|
||||
// rewrite before the exact-path lookup. A cyclic alias chain falls
|
||||
// through to the dangling-symlink path.
|
||||
let blocker = findUnresolvedBlocker(pending.targetPath);
|
||||
let targetPath = pending.targetPath;
|
||||
if (blocker === null) {
|
||||
try {
|
||||
targetPath = resolveDirectoryAliasPath(entriesByPath, targetPath);
|
||||
} catch {
|
||||
// Cyclic alias chain: fall through to the dangling-symlink path.
|
||||
} catch {}
|
||||
if (targetPath !== pending.targetPath) blocker = findUnresolvedBlocker(targetPath);
|
||||
}
|
||||
if (blocker !== null && blocker !== entry) {
|
||||
const waiting = dependents.get(blocker);
|
||||
if (waiting) {
|
||||
waiting.push(entry);
|
||||
} else {
|
||||
dependents.set(blocker, [entry]);
|
||||
}
|
||||
if (targetPath !== pending.targetPath && dependsOnUnresolvedLink(targetPath)) continue;
|
||||
continue;
|
||||
}
|
||||
unresolved.delete(entry);
|
||||
const settled = dependents.get(entry);
|
||||
if (settled) {
|
||||
dependents.delete(entry);
|
||||
queue.push(...settled);
|
||||
}
|
||||
if (blocker === entry) {
|
||||
// The target passes through the link itself (`a -> a/b`):
|
||||
// inherently cyclic, so it can never become a usable alias even
|
||||
// when real members exist beneath the target prefix.
|
||||
if (pending.kind === "hard link") {
|
||||
throw new ToolError(`Archive hard link '${entry.path}' has a cyclic target '${pending.targetPath}'`);
|
||||
}
|
||||
entry.storage = { type: "tar-link", targetPath: pending.targetPath };
|
||||
continue;
|
||||
}
|
||||
|
||||
const target = entriesByPath.get(targetPath);
|
||||
if (target?.storage && !target.isDirectory) {
|
||||
entry.size = target.size;
|
||||
entry.storage = target.storage;
|
||||
unresolved.delete(entry);
|
||||
resolved++;
|
||||
const target = entriesByPath.get(targetPath);
|
||||
if (target?.storage && !target.isDirectory && !unresolved.has(target)) {
|
||||
entry.size = target.size;
|
||||
entry.storage = target.storage;
|
||||
continue;
|
||||
}
|
||||
|
||||
// An empty target is the archive root, which is always a directory.
|
||||
const targetIsDirectory =
|
||||
targetPath === "" || target?.isDirectory === true || directoryPrefixes.has(targetPath);
|
||||
if (!targetIsDirectory) {
|
||||
if (pending.kind === "symlink") {
|
||||
entry.storage = { type: "tar-link", targetPath: pending.targetPath };
|
||||
continue;
|
||||
}
|
||||
|
||||
// An empty target is the archive root, which is always a directory.
|
||||
const targetPrefix = `${targetPath}/`;
|
||||
const targetIsDirectory =
|
||||
targetPath === "" ||
|
||||
target?.isDirectory === true ||
|
||||
entries.some(candidate => candidate.path.startsWith(targetPrefix));
|
||||
if (!targetIsDirectory) {
|
||||
if (pending.kind === "symlink") {
|
||||
entry.storage = { type: "tar-link", targetPath: pending.targetPath };
|
||||
unresolved.delete(entry);
|
||||
resolved++;
|
||||
continue;
|
||||
}
|
||||
const reason = target ? "unreadable member" : "missing member";
|
||||
throw new ToolError(`Archive hard link '${entry.path}' targets ${reason} '${pending.targetPath}'`);
|
||||
}
|
||||
if (pending.kind === "hard link") {
|
||||
throw new ToolError(`Archive hard link '${entry.path}' targets directory '${pending.targetPath}'`);
|
||||
}
|
||||
|
||||
entry.isDirectory = true;
|
||||
entry.storage = { type: "tar-link", targetPath: pending.targetPath };
|
||||
unresolved.delete(entry);
|
||||
resolved++;
|
||||
const reason = target ? "unreadable member" : "missing member";
|
||||
throw new ToolError(`Archive hard link '${entry.path}' targets ${reason} '${pending.targetPath}'`);
|
||||
}
|
||||
if (resolved === 0) {
|
||||
throw new ToolError("Archive contains cyclic or unsupported links");
|
||||
if (pending.kind === "hard link") {
|
||||
throw new ToolError(`Archive hard link '${entry.path}' targets directory '${pending.targetPath}'`);
|
||||
}
|
||||
|
||||
entry.isDirectory = true;
|
||||
entry.storage = { type: "tar-link", targetPath: pending.targetPath };
|
||||
}
|
||||
// Links never dequeued sit in a dependency cycle (a -> b/x, b -> a/y).
|
||||
if (unresolved.size > 0) {
|
||||
throw new ToolError("Archive contains cyclic or unsupported links");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1012,15 +1075,13 @@ const MAX_LINK_RESOLUTION_DEPTH = 40;
|
||||
function resolveDirectoryAliasPath(entries: ReadonlyMap<string, ArchiveIndexEntry>, archivePath: string): string {
|
||||
let resolvedPath = archivePath;
|
||||
const seen = new Set<string>();
|
||||
for (let depth = 0; depth < MAX_LINK_RESOLUTION_DEPTH && !seen.has(resolvedPath); depth++) {
|
||||
for (let rewrites = 0; !seen.has(resolvedPath); ) {
|
||||
seen.add(resolvedPath);
|
||||
const parts = resolvedPath.split("/");
|
||||
let replacement: string | undefined;
|
||||
for (let end = parts.length; end > 0; end--) {
|
||||
const prefix = parts.slice(0, end).join("/");
|
||||
const entry = entries.get(prefix);
|
||||
for (let end = resolvedPath.length; end > 0; end = resolvedPath.lastIndexOf("/", end - 1)) {
|
||||
const entry = entries.get(resolvedPath.slice(0, end));
|
||||
if (!entry?.isDirectory || entry.storage?.type !== "tar-link") continue;
|
||||
const suffix = parts.slice(end).join("/");
|
||||
const suffix = resolvedPath.slice(end + 1);
|
||||
replacement = suffix
|
||||
? entry.storage.targetPath
|
||||
? `${entry.storage.targetPath}/${suffix}`
|
||||
@@ -1029,6 +1090,10 @@ function resolveDirectoryAliasPath(entries: ReadonlyMap<string, ArchiveIndexEntr
|
||||
break;
|
||||
}
|
||||
if (replacement === undefined) return resolvedPath;
|
||||
// The bound counts performed rewrites, so a chain of exactly
|
||||
// MAX_LINK_RESOLUTION_DEPTH aliases still resolves; only needing one
|
||||
// more trips it.
|
||||
if (++rewrites > MAX_LINK_RESOLUTION_DEPTH) break;
|
||||
resolvedPath = replacement;
|
||||
}
|
||||
throw new ToolError(`Archive path '${archivePath}' crosses a cyclic symlink`);
|
||||
|
||||
@@ -73,7 +73,7 @@ function mountNoticesIn(messages: Message[]): string[] {
|
||||
typeof content === "string"
|
||||
? content
|
||||
: content.flatMap(part => (part.type === "text" ? [part.text] : [])).join("");
|
||||
return text.includes("The xd:// device inventory changed.") ? [text] : [];
|
||||
return text.includes("xd:// device inventory changed.") ? [text] : [];
|
||||
});
|
||||
}
|
||||
|
||||
@@ -938,10 +938,10 @@ describe("AgentSession refreshMCPTools rebuild skipping", () => {
|
||||
expect(contexts).toHaveLength(2);
|
||||
const mountNotices = mountNoticesIn(contexts[1]);
|
||||
expect(mountNotices).toHaveLength(1);
|
||||
expect(mountNotices[0]).toContain("became available");
|
||||
expect(mountNotices[0]).toContain("Available tools.");
|
||||
expect(mountNotices[0]).toContain("xd://mcp__nucleus_search");
|
||||
expect(mountNotices[0]).toContain("xd://mcp__nucleus_fetch");
|
||||
expect(mountNotices[0]).not.toContain("No longer mounted");
|
||||
expect(mountNotices[0]).not.toContain("Unmounted; writes fail:");
|
||||
|
||||
// A later unmount is likewise held for the following user prompt.
|
||||
await session.refreshMCPTools([search]);
|
||||
@@ -950,9 +950,9 @@ describe("AgentSession refreshMCPTools rebuild skipping", () => {
|
||||
await session.prompt("third");
|
||||
const allNotices = mountNoticesIn(contexts[2]);
|
||||
expect(allNotices).toHaveLength(2);
|
||||
expect(allNotices[1]).toContain("No longer mounted");
|
||||
expect(allNotices[1]).toContain("Unmounted; writes fail:");
|
||||
expect(allNotices[1]).toContain("xd://mcp__nucleus_fetch");
|
||||
expect(allNotices[1]).not.toContain("became available");
|
||||
expect(allNotices[1]).not.toContain("Available tools.");
|
||||
});
|
||||
|
||||
it("caps dynamic xd:// mount-notice summaries", async () => {
|
||||
@@ -1009,7 +1009,7 @@ describe("AgentSession refreshMCPTools rebuild skipping", () => {
|
||||
expect(notices).toHaveLength(1);
|
||||
expect(notices[0]).toContain("xd://mcp__nucleus_search");
|
||||
expect(notices[0]).not.toContain("mcp__nucleus_fetch");
|
||||
expect(notices[0]).not.toContain("No longer mounted");
|
||||
expect(notices[0]).not.toContain("Unmounted; writes fail:");
|
||||
});
|
||||
|
||||
it.each([
|
||||
|
||||
@@ -213,4 +213,152 @@ describe("MCP Streamable HTTP POST response resumption", () => {
|
||||
expect(observed.protocolVersion).toBe("2025-11-25");
|
||||
expect(observed.resumedAt - observed.postClosedAt).toBeGreaterThanOrEqual(15);
|
||||
});
|
||||
it("refreshes auth on a 401 resume GET without replaying the POST", async () => {
|
||||
const observed = { posts: 0, gets: 0, auth: [] as (string | null)[], lastEventId: null as string | null };
|
||||
server = Bun.serve({
|
||||
port: 0,
|
||||
fetch(req) {
|
||||
if (req.method === "POST") {
|
||||
observed.posts++;
|
||||
return new Response("id: stream-1\nretry: 10\ndata:\n\n", {
|
||||
headers: { "Content-Type": "text/event-stream" },
|
||||
});
|
||||
}
|
||||
observed.gets++;
|
||||
observed.auth.push(req.headers.get("Authorization"));
|
||||
observed.lastEventId = req.headers.get("Last-Event-ID");
|
||||
if (req.headers.get("Authorization") !== "Bearer fresh") {
|
||||
return new Response("expired", { status: 401 });
|
||||
}
|
||||
return new Response(
|
||||
'id: stream-2\ndata: {"jsonrpc":"2.0","id":1,"result":{"tools":[{"name":"resumed","inputSchema":{"type":"object"}}]}}\n\n',
|
||||
{ headers: { "Content-Type": "text/event-stream" } },
|
||||
);
|
||||
},
|
||||
});
|
||||
if (!server) throw new Error("Test server was not started");
|
||||
const transport = new HttpTransport({
|
||||
type: "http",
|
||||
url: `http://127.0.0.1:${server.port}/mcp`,
|
||||
timeout: GUARD_TIMEOUT_MS,
|
||||
headers: { Authorization: "Bearer stale" },
|
||||
});
|
||||
transport.onAuthError = async () => ({ Authorization: "Bearer fresh" });
|
||||
await transport.connect();
|
||||
|
||||
await expect(withPendingGuard(transport.request<ToolList>("tools/list"), "request")).resolves.toEqual({
|
||||
tools: [{ name: "resumed", inputSchema: { type: "object" } }],
|
||||
});
|
||||
// One POST only: replaying it after the server accepted the request
|
||||
// could double-execute a state-changing tool.
|
||||
expect(observed.posts).toBe(1);
|
||||
expect(observed.gets).toBe(2);
|
||||
expect(observed.auth).toEqual(["Bearer stale", "Bearer fresh"]);
|
||||
expect(observed.lastEventId).toBe("stream-1");
|
||||
});
|
||||
|
||||
it("resumes after an abrupt stream drop once an event ID exists", async () => {
|
||||
// Bun.serve cannot produce a genuine mid-body transport failure in-process
|
||||
// (stream errors surface as clean EOF client-side), so speak raw HTTP: a
|
||||
// chunked response without the terminal chunk, closed mid-body, makes the
|
||||
// client's body read throw.
|
||||
const observed = { posts: 0, lastEventId: null as string | null };
|
||||
const sseChunk = (payload: string): string =>
|
||||
`HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n${payload.length.toString(16)}\r\n${payload}\r\n`;
|
||||
const listener = Bun.listen({
|
||||
hostname: "127.0.0.1",
|
||||
port: 0,
|
||||
socket: {
|
||||
data(socket, data) {
|
||||
const request = new TextDecoder().decode(data);
|
||||
if (request.startsWith("POST")) {
|
||||
observed.posts++;
|
||||
// Priming event, then close without the terminal 0-chunk.
|
||||
socket.write(sseChunk("id: stream-1\nretry: 10\ndata:\n\n"));
|
||||
socket.end();
|
||||
return;
|
||||
}
|
||||
if (!request.startsWith("GET")) return;
|
||||
observed.lastEventId = /^Last-Event-ID:\s*(.+)$/im.exec(request)?.[1]?.trim() ?? null;
|
||||
socket.write(
|
||||
`${sseChunk(
|
||||
'id: stream-2\ndata: {"jsonrpc":"2.0","id":1,"result":{"tools":[{"name":"resumed","inputSchema":{"type":"object"}}]}}\n\n',
|
||||
)}0\r\n\r\n`,
|
||||
);
|
||||
socket.end();
|
||||
},
|
||||
},
|
||||
});
|
||||
try {
|
||||
const transport = new HttpTransport({
|
||||
type: "http",
|
||||
url: `http://127.0.0.1:${listener.port}/mcp`,
|
||||
timeout: GUARD_TIMEOUT_MS,
|
||||
});
|
||||
await transport.connect();
|
||||
|
||||
await expect(withPendingGuard(transport.request<ToolList>("tools/list"), "request")).resolves.toEqual({
|
||||
tools: [{ name: "resumed", inputSchema: { type: "object" } }],
|
||||
});
|
||||
expect(observed.posts).toBe(1);
|
||||
expect(observed.lastEventId).toBe("stream-1");
|
||||
} finally {
|
||||
listener.stop(true);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe("MCP Streamable HTTP GET listener resumption", () => {
|
||||
it("resumes the long-lived GET stream with Last-Event-ID instead of reconnecting", async () => {
|
||||
const observed = { gets: 0, lastEventIds: [] as (string | null)[] };
|
||||
server = Bun.serve({
|
||||
port: 0,
|
||||
fetch(req) {
|
||||
if (req.method !== "GET") {
|
||||
return Response.json({ jsonrpc: "2.0", id: 1, result: {} });
|
||||
}
|
||||
observed.gets++;
|
||||
observed.lastEventIds.push(req.headers.get("Last-Event-ID"));
|
||||
if (observed.gets === 1) {
|
||||
// Polling-style server: deliver one notification with an
|
||||
// event ID, then close the physical connection.
|
||||
return new Response(
|
||||
'id: poll-1\nretry: 10\ndata: {"jsonrpc":"2.0","method":"notifications/first"}\n\n',
|
||||
{ headers: { "Content-Type": "text/event-stream" } },
|
||||
);
|
||||
}
|
||||
return new Response(
|
||||
new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
controller.enqueue(
|
||||
encoder.encode('id: poll-2\ndata: {"jsonrpc":"2.0","method":"notifications/second"}\n\n'),
|
||||
);
|
||||
// Held open: the logical stream continues.
|
||||
},
|
||||
}),
|
||||
{ headers: { "Content-Type": "text/event-stream" } },
|
||||
);
|
||||
},
|
||||
});
|
||||
const transport = await connectedTransport();
|
||||
const notifications: string[] = [];
|
||||
let closed = false;
|
||||
const secondNotification = Promise.withResolvers<void>();
|
||||
transport.onNotification = method => {
|
||||
notifications.push(method);
|
||||
if (notifications.length === 2) secondNotification.resolve();
|
||||
};
|
||||
transport.onClose = () => {
|
||||
closed = true;
|
||||
};
|
||||
|
||||
await transport.startSSEListener();
|
||||
await withPendingGuard(secondNotification.promise, "resumed notification");
|
||||
|
||||
expect(notifications).toEqual(["notifications/first", "notifications/second"]);
|
||||
expect(observed.lastEventIds).toEqual([null, "poll-1"]);
|
||||
// The resume replaced the manager-level reconnect: no close fired.
|
||||
expect(closed).toBe(false);
|
||||
await transport.close();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -50,4 +50,19 @@ describe("tool activity visibility", () => {
|
||||
expect(restored).toContain("tool warning");
|
||||
expect(restored).toContain("late activity");
|
||||
});
|
||||
|
||||
it("forwards Ctrl+O expansion to wrapped expandable components", () => {
|
||||
// The transcript expansion traversal only visits top-level children;
|
||||
// without forwarding, a wrapped renderer freezes at insertion-time state.
|
||||
const states: boolean[] = [];
|
||||
class ExpandableText extends Text {
|
||||
setExpanded(expanded: boolean): void {
|
||||
states.push(expanded);
|
||||
}
|
||||
}
|
||||
const wrapper = new ToolActivityContainer(new ExpandableText("activity", 1, 0));
|
||||
wrapper.setExpanded(true);
|
||||
wrapper.setExpanded(false);
|
||||
expect(states).toEqual([true, false]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -109,10 +109,7 @@ describe("createAgentSession defaultInactive tool activation", () => {
|
||||
workspaceTree: { rootPath: tempDir, rendered: "", truncated: false, totalLines: 0, agentsMdFiles: [] },
|
||||
});
|
||||
|
||||
const requireBundledModel = (
|
||||
provider: "anthropic" | "google-antigravity" | "openai" | "xai",
|
||||
id: string,
|
||||
): Model => {
|
||||
const requireBundledModel = (provider: "anthropic" | "google" | "openai" | "xai", id: string): Model => {
|
||||
const bundled = getBundledModel(provider, id);
|
||||
if (!bundled) throw new Error(`Expected ${provider}/${id} model to exist`);
|
||||
return bundled;
|
||||
@@ -190,7 +187,8 @@ describe("createAgentSession defaultInactive tool activation", () => {
|
||||
const unsupported = requireBundledModel("xai", "grok-4");
|
||||
const fable = requireBundledModel("anthropic", "claude-fable-5");
|
||||
const responses = requireBundledModel("openai", "gpt-5");
|
||||
const gemini = requireBundledModel("google-antigravity", "gemini-3.6-flash");
|
||||
const gemini = requireBundledModel("google", "gemini-2.5-flash");
|
||||
const mandatoryGemini = requireBundledModel("google", "gemini-2.5-pro");
|
||||
const { session } = await createAgentSession({
|
||||
...baseOptions(tempDir),
|
||||
settings,
|
||||
@@ -199,7 +197,7 @@ describe("createAgentSession defaultInactive tool activation", () => {
|
||||
const authStorage = session.modelRegistry.authStorage;
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
authStorage.setRuntimeApiKey("openai", "test-key");
|
||||
authStorage.setRuntimeApiKey("google-antigravity", "test-key");
|
||||
authStorage.setRuntimeApiKey("google", "test-key");
|
||||
authStorage.setRuntimeApiKey("xai", "test-key");
|
||||
|
||||
try {
|
||||
@@ -214,6 +212,8 @@ describe("createAgentSession defaultInactive tool activation", () => {
|
||||
expect(session.getActiveToolNames()).toContain("think");
|
||||
await session.setModel(gemini);
|
||||
expect(session.getActiveToolNames()).toContain("think");
|
||||
await session.setModel(mandatoryGemini);
|
||||
expect(session.getActiveToolNames()).not.toContain("think");
|
||||
|
||||
await session.setModel(unsupported);
|
||||
expect(session.getActiveToolNames()).not.toContain("think");
|
||||
@@ -249,7 +249,7 @@ describe("createAgentSession defaultInactive tool activation", () => {
|
||||
fetch: async request => {
|
||||
requestTexts.push(await request.text());
|
||||
if (requestTexts.length === 1) {
|
||||
const argumentsJson = JSON.stringify({ notes: "Checked the request before answering." });
|
||||
const argumentsJson = JSON.stringify({ thoughts: "Checked the request before answering." });
|
||||
return sse([
|
||||
{
|
||||
type: "response.output_item.added",
|
||||
|
||||
@@ -179,6 +179,35 @@ function createSparsePaxTarArchive(realName: string, realSize: number, storedDat
|
||||
return Buffer.concat(parts);
|
||||
}
|
||||
|
||||
/**
|
||||
* Build an old-GNU sparse archive: an `S` member whose header sets the
|
||||
* `isextended` flag (byte 482), followed by one 512-byte sparse-map
|
||||
* continuation block that is not counted in the member's declared size,
|
||||
* then the stored data and a regular member.
|
||||
*/
|
||||
function createOldGnuSparseTarArchive(): Buffer {
|
||||
const storedData = Buffer.from("sparse-extent\n", "utf-8");
|
||||
const header = Buffer.alloc(512, 0);
|
||||
writeTarString(header, 0, 100, "data/old-sparse.bin");
|
||||
writeTarOctal(header, 100, 8, 0o644);
|
||||
writeTarOctal(header, 124, 12, storedData.length);
|
||||
writeTarOctal(header, 136, 12, Math.floor(Date.now() / 1000));
|
||||
header[156] = "S".charCodeAt(0);
|
||||
writeTarString(header, 257, 6, "ustar");
|
||||
writeTarString(header, 263, 2, "00");
|
||||
header[482] = 1; // sparse map continues in extension blocks
|
||||
tarChecksum(header);
|
||||
|
||||
// Final continuation block: its own isextended byte (504) stays 0.
|
||||
const continuation = Buffer.alloc(512, 0);
|
||||
const parts: Buffer[] = [header, continuation, storedData];
|
||||
const remainder = storedData.length % 512;
|
||||
if (remainder !== 0) parts.push(Buffer.alloc(512 - remainder, 0));
|
||||
// createTarArchive appends the member and the end-of-archive terminator.
|
||||
parts.push(createTarArchive([{ path: "data/after.txt", content: "after sparse\n" }]));
|
||||
return Buffer.concat(parts);
|
||||
}
|
||||
|
||||
const CRC32_TABLE = (() => {
|
||||
const table = new Uint32Array(256);
|
||||
for (let index = 0; index < 256; index++) {
|
||||
@@ -940,21 +969,91 @@ describe("Coding Agent Tools", () => {
|
||||
expect(getTextOutput(linkedResult)).toContain("export const linked = true");
|
||||
});
|
||||
|
||||
it("should reject directory symlinks targeting their own subtree instead of looping", async () => {
|
||||
it("should degrade directory symlinks targeting their own subtree to dangling links", async () => {
|
||||
const archivePath = path.join(testDir, "self-cycle-symlink.tar");
|
||||
fs.writeFileSync(
|
||||
archivePath,
|
||||
createTarArchive([
|
||||
{ path: "a/b/f.txt", content: "unreachable\n" },
|
||||
// `a -> a/b` grows the resolved path on every rewrite; the
|
||||
// pre-fix resolver looped forever on this shape.
|
||||
{ path: "a/b/f.txt", content: "still readable\n" },
|
||||
// `a -> a/b` is inherently cyclic; the pre-fix resolver looped
|
||||
// forever growing the rewritten path. The link now dangles
|
||||
// while real members underneath stay readable.
|
||||
{ path: "a", content: "", typeFlag: "2", linkName: "a/b" },
|
||||
]),
|
||||
);
|
||||
|
||||
const memberResult = await readTool.execute("test-call-tar-self-cycle-member", {
|
||||
path: `${archivePath}:a/b/f.txt`,
|
||||
});
|
||||
expect(getTextOutput(memberResult)).toContain("still readable");
|
||||
await expect(readTool.execute("test-call-tar-self-cycle-link", { path: `${archivePath}:a` })).rejects.toThrow(
|
||||
/cannot be materialized/,
|
||||
);
|
||||
});
|
||||
|
||||
it("should resolve alias chains up to the rewrite bound and reject deeper ones", async () => {
|
||||
const buildChain = (length: number): ArchiveFixtureEntry[] => {
|
||||
const chain: ArchiveFixtureEntry[] = [{ path: "real/f.txt", content: "deep\n" }];
|
||||
for (let i = 0; i < length; i++) {
|
||||
chain.push({
|
||||
path: `a${i}`,
|
||||
content: "",
|
||||
typeFlag: "2",
|
||||
linkName: i === length - 1 ? "real" : `a${i + 1}`,
|
||||
});
|
||||
}
|
||||
return chain;
|
||||
};
|
||||
|
||||
const okPath = path.join(testDir, "alias-chain-40.tar");
|
||||
fs.writeFileSync(okPath, createTarArchive(buildChain(40)));
|
||||
const okResult = await readTool.execute("test-call-tar-alias-chain-40", { path: `${okPath}:a0/f.txt` });
|
||||
expect(getTextOutput(okResult)).toContain("deep");
|
||||
|
||||
const deepPath = path.join(testDir, "alias-chain-41.tar");
|
||||
fs.writeFileSync(deepPath, createTarArchive(buildChain(41)));
|
||||
await expect(
|
||||
readTool.execute("test-call-tar-self-cycle", { path: `${archivePath}:a/b/f.txt` }),
|
||||
).rejects.toThrow(/cyclic or unsupported links/);
|
||||
readTool.execute("test-call-tar-alias-chain-41", { path: `${deepPath}:a0/f.txt` }),
|
||||
).rejects.toThrow(/cyclic symlink/);
|
||||
});
|
||||
|
||||
it("should let the later duplicate tar member win", async () => {
|
||||
const archivePath = path.join(testDir, "duplicate-member.tar");
|
||||
// tar -rf append/update semantics: extraction yields the last member.
|
||||
fs.writeFileSync(
|
||||
archivePath,
|
||||
createTarArchive([
|
||||
{ path: "dup/file.txt", content: "first\n" },
|
||||
{ path: "dup/file.txt", content: "second\n" },
|
||||
]),
|
||||
);
|
||||
|
||||
const result = await readTool.execute("test-call-tar-duplicate-member", {
|
||||
path: `${archivePath}:dup/file.txt`,
|
||||
});
|
||||
expect(getTextOutput(result)).toContain("second");
|
||||
expect(getTextOutput(result)).not.toContain("first");
|
||||
});
|
||||
|
||||
it("should skip old-GNU sparse extension blocks between header and data", async () => {
|
||||
const archivePath = path.join(testDir, "old-gnu-sparse.tar");
|
||||
fs.writeFileSync(archivePath, createOldGnuSparseTarArchive());
|
||||
|
||||
// Pre-fix the continuation block was parsed as the next header and
|
||||
// the whole archive rejected as corrupt.
|
||||
const rootResult = await readTool.execute("test-call-tar-old-gnu-sparse-root", {
|
||||
path: `${archivePath}:data`,
|
||||
});
|
||||
expect(getTextOutput(rootResult)).toContain("old-sparse.bin");
|
||||
expect(getTextOutput(rootResult)).toContain("after.txt");
|
||||
|
||||
const afterResult = await readTool.execute("test-call-tar-old-gnu-sparse-after", {
|
||||
path: `${archivePath}:data/after.txt`,
|
||||
});
|
||||
expect(getTextOutput(afterResult)).toContain("after sparse");
|
||||
await expect(
|
||||
readTool.execute("test-call-tar-old-gnu-sparse-member", { path: `${archivePath}:data/old-sparse.bin` }),
|
||||
).rejects.toThrow(/sparse file and cannot be read/);
|
||||
});
|
||||
|
||||
it("should resolve tar symlinks whose target is the archive root", async () => {
|
||||
|
||||
@@ -315,6 +315,12 @@ describe("shape resolution", () => {
|
||||
// high-res tier rather than falling back to the 1568px default (#8256).
|
||||
expect(snapcompact.resolveShape({ api: "anthropic-messages", id: "claude-opus-5" }).frameSize).toBe(1932);
|
||||
expect(snapcompact.resolveShape({ api: "anthropic-messages", id: "claude-opus-6" }).frameSize).toBe(1932);
|
||||
// Mixed-case gateway ids matched the pre-catalog-parser /i regex; the
|
||||
// parser input is normalized so they stay on the high-res tier.
|
||||
expect(snapcompact.resolveShape({ api: "anthropic-messages", id: "CLAUDE-OPUS-5" }).frameSize).toBe(1932);
|
||||
// Minor versions past the old 9.10 semver-table bound must not fall
|
||||
// back to the 1568px default (same staleness class as #8256).
|
||||
expect(snapcompact.resolveShape({ api: "anthropic-messages", id: "claude-opus-5-11" }).frameSize).toBe(1932);
|
||||
// Opus lines below 4.7 downscale, so they keep the safe family default.
|
||||
expect(snapcompact.resolveShape({ api: "anthropic-messages", id: "claude-opus-4-6" })).toBe(
|
||||
snapcompact.SHAPES.anthropic,
|
||||
|
||||
Reference in New Issue
Block a user