Merge branch 'main' into farm/9483b02b/fix-alt-backspace-kitty-ghostty

This commit is contained in:
Can Bölük
2026-06-08 00:45:21 +02:00
committed by GitHub
39 changed files with 2092 additions and 182 deletions
+3
View File
@@ -16,6 +16,9 @@ concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: true
jobs:
# scripts/release.ts pushes the version-bump commit and its `v*` tag
# atomically (`git push --atomic origin main refs/tags/v*`), so a release
+6 -6
View File
@@ -81,13 +81,13 @@ These are consumed via `getEnvApiKey()` (`packages/ai/src/stream.ts`) unless not
| `WAFER_SERVERLESS_API_KEY` | Wafer Serverless auth | Using `wafer-serverless` provider | Pay-as-you-go Wafer SKU; validated against `https://pass.wafer.ai/v1/models` |
| `GITLAB_TOKEN` | GitLab Duo auth | Using `gitlab-duo` provider | |
### GitHub/Copilot token chains
### GitHub/Copilot tokens
| Variable | Used for | Chain |
| ---------------------- | ------------------------------------------------ | ---------------------------------------------------- |
| `COPILOT_GITHUB_TOKEN` | GitHub Copilot provider auth | `COPILOT_GITHUB_TOKEN` → `GH_TOKEN` → `GITHUB_TOKEN` |
| `GH_TOKEN` | Copilot fallback; GitHub API auth in web scraper | In web scraper: `GITHUB_TOKEN` → `GH_TOKEN` |
| `GITHUB_TOKEN` | Copilot fallback; GitHub API auth in web scraper | In web scraper: checked before `GH_TOKEN` |
| Variable | Used for | Notes |
| ---------------------- | ------------------------------------------------ | ------------------------------------------ |
| `COPILOT_GITHUB_TOKEN` | GitHub Copilot provider auth | Generic GitHub tokens are not used here |
| `GH_TOKEN` | GitHub API auth in web scraper | Web scraper fallback after `GITHUB_TOKEN` |
| `GITHUB_TOKEN` | GitHub API auth in web scraper | Web scraper checks this before `GH_TOKEN` |
### Auth broker / auth gateway (remote credential vault)
+10 -4
View File
@@ -47,6 +47,9 @@ Upstream uses different package scopes. Replace them consistently.
- `@mariozechner/pi-agent-core` → `@oh-my-pi/pi-agent-core`
- `@mariozechner/pi-tui` → `@oh-my-pi/pi-tui`
- `@mariozechner/pi-ai` → `@oh-my-pi/pi-ai`
- `@mariozechner/pi-utils` → `@oh-my-pi/pi-utils`
- Some upstream packages publish under the `@earendil-works/*` scope instead of `@mariozechner/*`. Map it the same way (`@earendil-works/pi-coding-agent` → `@oh-my-pi/pi-coding-agent`, and so on).
- The bare `typebox` package is not an `@oh-my-pi/*` scope; do not rewrite it as one. See the Extensions divergence in section 15 for how tool-parameter schemas map.
## 4) Use Bun APIs where they improve on Node
@@ -353,10 +356,13 @@ Our fork has architectural decisions that differ from upstream. **Do not port th
### Extensions
| Upstream | Our Fork |
| ----------------------------- | ------------------------------------------------- |
| `jiti` for TypeScript loading | Native Bun `import()` |
| `pkg.pi` manifest field | `pkg.omp` preferred; fallback to `pkg.pi` remains |
| Upstream | Our Fork |
| ---------------------------------------------------------------- | ------------------------------------------------------------------------------------------------- |
| `jiti` for TypeScript loading | Native Bun `import()` |
| `pkg.pi` manifest field | `pkg.omp` preferred; fallback to `pkg.pi` remains |
| `StringEnum` from `pi-ai` | `Type.Enum` from the `pi.typebox` shim (or author the schema with `pi.zod`); `pi-ai` no longer exports `StringEnum` |
| `formatSize` from `pi-coding-agent` | `formatBytes` from `@oh-my-pi/pi-utils` |
| `DefaultResourceLoader` / `DefaultPackageManager` / `SettingsManager` / `createEventBus` | Capability-based discovery (`loadCapability(...)`) plus the `Settings` singleton and `EventBus` |
### Skip These Upstream Features
+6 -1
View File
@@ -1,7 +1,11 @@
# Changelog
## [Unreleased]
### Added
- Added support for `impersonated_service_account` Application Default Credentials (ADC) in Vertex AI to enable chained impersonation without failing via 401 `invalid_client`.
## [15.10.1] - 2026-06-07
### Breaking Changes
@@ -21,6 +25,7 @@
### Fixed
- Fixed duplicate upstream `tool_call_id` values collapsing distinct tool calls during message transformation, preserving one call/result pairing per emitted tool call before provider replay. ([#2055](https://github.com/can1357/oh-my-pi/issues/2055))
- Fixed streaming auth retries to handle `401` and usage-limit errors before replay-unsafe content is emitted, including failures surfaced only via `errorStatus`
- Fixed tool argument validation to coerce singleton non-string values into arrays when the schema expects an array, preventing Anthropic-compatible models that emit `todo.ops` as an object from getting stuck in repeated validation-error loops. ([#2026](https://github.com/can1357/oh-my-pi/issues/2026))
- Fixed streaming retries to buffer and suppress partial `start` events from failed auth attempts so only clean retried events are delivered
@@ -0,0 +1,144 @@
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import { Buffer } from "node:buffer";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import type { FetchImpl } from "../../types";
import { __resetVertexTokenCache, getVertexAccessToken } from "../google-auth";
const CLOUD_PLATFORM_SCOPE = "https://www.googleapis.com/auth/cloud-platform";
const JWT_BEARER_GRANT = "urn:ietf:params:oauth:grant-type:jwt-bearer";
/** Generate a real RS256 private key so signJwtRs256 / pemToPkcs8 run for real. */
async function generateServiceAccountPem(): Promise<string> {
const keyPair = (await globalThis.crypto.subtle.generateKey(
{ name: "RSASSA-PKCS1-v1_5", modulusLength: 2048, publicExponent: new Uint8Array([1, 0, 1]), hash: "SHA-256" },
true,
["sign", "verify"],
)) as CryptoKeyPair;
const pkcs8 = new Uint8Array(await globalThis.crypto.subtle.exportKey("pkcs8", keyPair.privateKey));
const body = (
Buffer.from(pkcs8)
.toString("base64")
.match(/.{1,64}/g) ?? []
).join("\n");
return `-----BEGIN PRIVATE KEY-----\n${body}\n-----END PRIVATE KEY-----\n`;
}
function urlOf(input: string | URL | Request): string {
if (typeof input === "string") return input;
if (input instanceof URL) return input.toString();
return input.url;
}
describe("getVertexAccessToken impersonated_service_account ADC", () => {
let tmpDir: string;
let originalGac: string | undefined;
beforeEach(async () => {
__resetVertexTokenCache();
originalGac = Bun.env.GOOGLE_APPLICATION_CREDENTIALS;
tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-vertex-adc-"));
});
afterEach(async () => {
__resetVertexTokenCache();
if (originalGac === undefined) delete Bun.env.GOOGLE_APPLICATION_CREDENTIALS;
else Bun.env.GOOGLE_APPLICATION_CREDENTIALS = originalGac;
await fs.rm(tmpDir, { recursive: true, force: true });
});
it("rejects a malformed service_account_impersonation_url before any network call", async () => {
const adcPath = path.join(tmpDir, "impersonated-bad-url.json");
await Bun.write(
adcPath,
JSON.stringify({
type: "impersonated_service_account",
// Missing the trailing ":generateAccessToken" the principal parser requires.
service_account_impersonation_url:
"https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/target@project.iam.gserviceaccount.com",
source_credentials: {
type: "authorized_user",
client_id: "client-id",
client_secret: "client-secret",
refresh_token: "refresh-token",
},
}),
);
Bun.env.GOOGLE_APPLICATION_CREDENTIALS = adcPath;
const calls: string[] = [];
const fetchImpl: FetchImpl = async input => {
calls.push(urlOf(input));
return new Response("{}");
};
// The principal is parsed before the source exchange, so a bad URL must fail
// up front rather than after burning a source-token round trip.
await expect(getVertexAccessToken({ fetch: fetchImpl })).rejects.toBeInstanceOf(RangeError);
expect(calls).toEqual([]);
});
it("signs an RS256 JWT for a service_account source and reconstructs the IAM URL", async () => {
const pem = await generateServiceAccountPem();
const adcPath = path.join(tmpDir, "impersonated-sa.json");
await Bun.write(
adcPath,
JSON.stringify({
type: "impersonated_service_account",
// Non-canonical project segment proves the request URL is rebuilt, not echoed.
service_account_impersonation_url:
"https://iamcredentials.googleapis.com/v1/projects/explicit-proj/serviceAccounts/target@project.iam.gserviceaccount.com:generateAccessToken",
source_credentials: {
type: "service_account",
client_email: "source@project.iam.gserviceaccount.com",
private_key: pem,
private_key_id: "key-1",
},
// delegates intentionally omitted — the IAM body must default to [].
}),
);
Bun.env.GOOGLE_APPLICATION_CREDENTIALS = adcPath;
const calls: { url: string; init?: RequestInit }[] = [];
const fetchImpl: FetchImpl = async (input, init) => {
const url = urlOf(input);
calls.push({ url, init });
if (url === "https://oauth2.googleapis.com/token") {
return new Response(JSON.stringify({ access_token: "sa-source-token", expires_in: 3600 }));
}
if (url.startsWith("https://iamcredentials.googleapis.com/")) {
return new Response(
JSON.stringify({
accessToken: "impersonated-token",
expireTime: new Date(Date.now() + 3_600_000).toISOString(),
}),
);
}
return new Response("unexpected", { status: 404 });
};
const token = await getVertexAccessToken({ fetch: fetchImpl });
expect(token).toBe("impersonated-token");
// Source JWT exchange happens first, then the impersonation exchange.
expect(calls.map(c => c.url)).toEqual([
"https://oauth2.googleapis.com/token",
"https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/target@project.iam.gserviceaccount.com:generateAccessToken",
]);
// Source credential is exchanged via a signed JWT bearer assertion, not a refresh grant.
const sourceBody = new URLSearchParams(String(calls[0].init?.body));
expect(sourceBody.get("grant_type")).toBe(JWT_BEARER_GRANT);
expect((sourceBody.get("assertion") ?? "").split(".")).toHaveLength(3);
// The IAM call carries the source-derived bearer token and defaults delegates to [].
const iamHeaders = calls[1].init?.headers as Record<string, string>;
expect(iamHeaders.Authorization).toBe("Bearer sa-source-token");
expect(JSON.parse(String(calls[1].init?.body))).toEqual({
delegates: [],
scope: [CLOUD_PLATFORM_SCOPE],
lifetime: "3600s",
});
});
});
+54 -5
View File
@@ -42,7 +42,14 @@ interface AuthorizedUserCredentials {
refresh_token: string;
}
type AdcFileCredentials = ServiceAccountCredentials | AuthorizedUserCredentials;
interface ImpersonatedServiceAccountCredentials {
type: "impersonated_service_account";
service_account_impersonation_url: string;
source_credentials: AuthorizedUserCredentials | ServiceAccountCredentials;
delegates?: string[];
}
type AdcFileCredentials = ServiceAccountCredentials | AuthorizedUserCredentials | ImpersonatedServiceAccountCredentials;
interface TokenResponse {
access_token: string;
@@ -196,10 +203,52 @@ async function resolveAccessTokenUncached(
): Promise<{ source: string; token: TokenResponse }> {
const adc = await loadAdcCredentials();
if (adc) {
const token =
adc.creds.type === "service_account"
? await exchangeJwtForToken(adc.creds, signal, fetchImpl)
: await exchangeRefreshToken(adc.creds, signal, fetchImpl);
const creds = adc.creds;
let token: TokenResponse;
if (creds.type === "impersonated_service_account") {
const targetPrincipalMatch = /(?<target>[^/]+):(generateAccessToken|generateIdToken)$/.exec(
creds.service_account_impersonation_url,
);
const targetPrincipal = targetPrincipalMatch?.groups?.target;
if (!targetPrincipal) {
throw new RangeError(`Cannot extract target principal from ${creds.service_account_impersonation_url}`);
}
const sourceToken =
creds.source_credentials.type === "service_account"
? await exchangeJwtForToken(creds.source_credentials, signal, fetchImpl)
: await exchangeRefreshToken(creds.source_credentials, signal, fetchImpl);
const response = await fetchImpl(
`https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/${targetPrincipal}:generateAccessToken`,
{
method: "POST",
headers: {
"Content-Type": "application/json",
Authorization: `Bearer ${sourceToken.access_token}`,
},
body: JSON.stringify({
delegates: creds.delegates ?? [],
scope: [CLOUD_PLATFORM_SCOPE],
lifetime: "3600s",
}),
signal,
},
);
if (!response.ok) {
const detail = await response.text().catch(() => "");
throw new Error(`Google Impersonation token exchange failed (${response.status}): ${detail}`);
}
const data = (await response.json()) as { accessToken: string; expireTime: string };
const expiresIn = Math.max(0, Math.floor((new Date(data.expireTime).getTime() - Date.now()) / 1000));
token = { access_token: data.accessToken, expires_in: expiresIn, token_type: "Bearer" };
} else {
token =
creds.type === "service_account"
? await exchangeJwtForToken(creds, signal, fetchImpl)
: await exchangeRefreshToken(creds, signal, fetchImpl);
}
return { source: adc.source, token };
}
const metadata = await fetchMetadataToken(signal, fetchImpl);
@@ -537,7 +537,6 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
}
stream.push({ type: "start", partial: output });
const parseMiniMaxThinkTags = model.provider === "minimax-code" || model.provider === "minimax-code-cn";
// Some OpenAI-compatible DeepSeek hosts (including NVIDIA NIM and DeepSeek's
// native API) leak chat-template tool-call markers in `delta.content` even
// though tool calls are also surfaced structurally. Strip the leaked markers
@@ -678,9 +677,7 @@ export const streamOpenAICompletions: StreamFunction<"openai-completions"> = (
}
};
const streamMarkupHealingPattern = getStreamMarkupHealingPattern(model.provider, model.id, {
parseThinkingTags: parseMiniMaxThinkTags,
});
const streamMarkupHealingPattern = getStreamMarkupHealingPattern(model.provider, model.id);
const streamMarkupHealing = streamMarkupHealingPattern
? new StreamMarkupHealing({ pattern: streamMarkupHealingPattern })
: undefined;
+201 -107
View File
@@ -17,6 +17,96 @@ const enum ToolCallStatus {
Aborted = 2,
}
/**
* Maximum tool-call id length the strictest replay provider accepts.
*
* Anthropic requires `^[a-zA-Z0-9_-]+$` with a 64-char cap; Google and Codex
* `normalizeToolCallId` implementations cap individual id segments to the same
* 64-char ceiling. Replacement ids minted here flow back through
* `convertAnthropicMessages` (and friends) unchanged, so the `_dupN` suffix
* MUST not push a normalized id past this bound.
*/
const MAX_TOOL_CALL_ID_LENGTH = 64;
function appendDuplicateSuffix(originalId: string, suffix: string): string {
if (originalId.length + suffix.length <= MAX_TOOL_CALL_ID_LENGTH) return `${originalId}${suffix}`;
const prefixBudget = Math.max(0, MAX_TOOL_CALL_ID_LENGTH - suffix.length);
return `${originalId.slice(0, prefixBudget)}${suffix}`;
}
type PendingToolResultRewrite = { replacementId: string } | undefined;
function deduplicateToolCallIds(messages: Message[]): Message[] {
const seenToolCallIds = new Map<string, number>();
const pendingToolResultRewrites = new Map<string, PendingToolResultRewrite[]>();
return messages.map(msg => {
if (msg.role === "toolResult") {
const rewrites = pendingToolResultRewrites.get(msg.toolCallId);
if (!rewrites || rewrites.length === 0) return msg;
const rewrite = rewrites.shift();
if (rewrites.length === 0) pendingToolResultRewrites.delete(msg.toolCallId);
if (rewrite) return { ...msg, toolCallId: rewrite.replacementId };
return msg;
}
if (msg.role !== "assistant") return msg;
const enqueueToolResultRewrite = (id: string, rewrite: PendingToolResultRewrite): void => {
const rewrites = pendingToolResultRewrites.get(id);
if (rewrites) {
rewrites.push(rewrite);
return;
}
pendingToolResultRewrites.set(id, [rewrite]);
};
// Ids this turn has already touched; used to scope the "drop carried-over
// pending rewrites" semantics to the FIRST occurrence per turn so multiple
// blocks of the same id within one turn still accumulate as duplicates.
const idsTouchedInTurn = new Set<string>();
let contentChanged = false;
const content = msg.content.map(block => {
if (block.type !== "toolCall") return block;
// Drop any pending rewrites carried over from a prior assistant turn
// for this id on its first appearance this turn. When a later turn
// re-emits the same id, the older duplicate call's expected result
// never landed in time — the second pass synthesizes
// "No result provided" for it, and the upcoming real result(id) must
// route to one of THIS turn's calls. Without this guard the older
// `_dup` id would steal the next result.
if (!idsTouchedInTurn.has(block.id)) {
pendingToolResultRewrites.delete(block.id);
idsTouchedInTurn.add(block.id);
}
const previousCount = seenToolCallIds.get(block.id) ?? 0;
if (previousCount === 0) {
seenToolCallIds.set(block.id, 1);
enqueueToolResultRewrite(block.id, undefined);
return block;
}
let duplicateIndex = previousCount;
let replacementId = appendDuplicateSuffix(block.id, `_dup${duplicateIndex}`);
while (seenToolCallIds.has(replacementId)) {
duplicateIndex += 1;
replacementId = appendDuplicateSuffix(block.id, `_dup${duplicateIndex}`);
}
seenToolCallIds.set(block.id, duplicateIndex + 1);
seenToolCallIds.set(replacementId, 1);
enqueueToolResultRewrite(block.id, { replacementId });
contentChanged = true;
return { ...block, id: replacementId };
});
if (!contentChanged) return msg;
return { ...msg, content };
});
}
function shouldDropTruncatedThinkingOnlyAssistant(msg: AssistantMessage): boolean {
const isTruncatedStop = msg.stopReason === "length" || msg.stopReason === "error" || msg.stopReason === "aborted";
return isTruncatedStop && !msg.content.some(block => block.type === "toolCall" || block.type === "text");
@@ -52,116 +142,120 @@ export function transformMessages<TApi extends Api>(
const latestSurvivingAssistantIndex = getLatestSurvivingAssistantIndex(messages);
// First pass: transform messages (thinking blocks, tool call ID normalization)
const transformed = messages.map((msg, index) => {
// User and developer messages pass through unchanged
if (msg.role === "user" || msg.role === "developer") {
return msg;
}
const transformed = deduplicateToolCallIds(
messages.map((msg, index) => {
// User and developer messages pass through unchanged
if (msg.role === "user" || msg.role === "developer") {
return msg;
}
// Handle toolResult messages - normalize toolCallId if we have a mapping
if (msg.role === "toolResult") {
const normalizedId = toolCallIdMap.get(msg.toolCallId);
if (normalizedId && normalizedId !== msg.toolCallId) {
return { ...msg, toolCallId: normalizedId };
// Handle toolResult messages - normalize toolCallId if we have a mapping
if (msg.role === "toolResult") {
const normalizedId = toolCallIdMap.get(msg.toolCallId);
if (normalizedId && normalizedId !== msg.toolCallId) {
return { ...msg, toolCallId: normalizedId };
}
return msg;
}
// Assistant messages need transformation check
if (msg.role === "assistant") {
const assistantMsg = msg as AssistantMessage;
const isSameModel =
assistantMsg.provider === model.provider &&
assistantMsg.api === model.api &&
assistantMsg.model === model.id;
const mustPreserveLatestAnthropicThinking =
index === latestSurvivingAssistantIndex &&
model.api === "anthropic-messages" &&
assistantMsg.api === "anthropic-messages";
// Aborted/errored messages may have partially-streamed thinking signatures.
// A partial signature is invalid and will be rejected by the API, so we must
// strip signatures from thinking blocks in these messages.
//
// Abandoned tool-use turns get the same treatment once they are no longer
// the latest assistant message. When a turn carries toolCall blocks but did
// NOT request tool execution (stopReason !== "toolUse" — e.g.
// adaptive-thinking Opus emitting tool calls and then ending the turn on
// `end_turn`/`stop`), the agent loop pairs those calls with placeholder
// tool_results to keep the tool_use/tool_result contract valid. Historical
// abandoned turns cannot safely replay their end_turn-bound signatures in
// that continuation, so stripping downgrades them to plain text downstream.
// Latest abandoned turns are exempt because Anthropic requires thinking
// blocks from its most recent response to remain byte-for-byte unmodified.
const invalidStopReason = assistantMsg.stopReason === "aborted" || assistantMsg.stopReason === "error";
const abandonedToolUse =
!invalidStopReason &&
assistantMsg.stopReason !== "toolUse" &&
assistantMsg.content.some(b => b.type === "toolCall");
const hasInvalidSignatures = invalidStopReason || abandonedToolUse;
const transformedContent = assistantMsg.content.flatMap(block => {
if (block.type === "thinking") {
// Strip untrustworthy signatures so the encoder can downgrade to text.
const sanitized =
hasInvalidSignatures && block.thinkingSignature
? { ...block, thinkingSignature: undefined }
: block;
if (mustPreserveLatestAnthropicThinking) return abandonedToolUse ? block : sanitized;
// For same model: keep thinking blocks with signatures (needed for replay)
// even if the thinking text is empty (OpenAI encrypted reasoning)
if (isSameModel && sanitized.thinkingSignature) return sanitized;
// Skip empty thinking blocks, convert others to plain text
if (!sanitized.thinking || sanitized.thinking.trim() === "") return [];
if (isSameModel) return sanitized;
return {
type: "text" as const,
text: sanitized.thinking,
};
}
if (block.type === "redactedThinking") {
if (mustPreserveLatestAnthropicThinking) return block;
if (isSameModel) return block;
return [];
}
if (block.type === "text") {
if (isSameModel) return block;
return {
type: "text" as const,
text: block.text,
};
}
if (block.type === "toolCall") {
const toolCall = block as ToolCall;
let normalizedToolCall: ToolCall = toolCall;
if (!isSameModel && toolCall.thoughtSignature) {
normalizedToolCall = { ...toolCall };
delete (normalizedToolCall as { thoughtSignature?: string }).thoughtSignature;
}
if (!isSameModel && normalizeToolCallId) {
const normalizedId = normalizeToolCallId(toolCall.id, model, assistantMsg);
if (normalizedId !== toolCall.id) {
toolCallIdMap.set(toolCall.id, normalizedId);
normalizedToolCall = { ...normalizedToolCall, id: normalizedId };
}
}
return normalizedToolCall;
}
return block;
});
return {
...assistantMsg,
content: transformedContent,
};
}
return msg;
}
// Assistant messages need transformation check
if (msg.role === "assistant") {
const assistantMsg = msg as AssistantMessage;
const isSameModel =
assistantMsg.provider === model.provider &&
assistantMsg.api === model.api &&
assistantMsg.model === model.id;
const mustPreserveLatestAnthropicThinking =
index === latestSurvivingAssistantIndex &&
model.api === "anthropic-messages" &&
assistantMsg.api === "anthropic-messages";
// Aborted/errored messages may have partially-streamed thinking signatures.
// A partial signature is invalid and will be rejected by the API, so we must
// strip signatures from thinking blocks in these messages.
//
// Abandoned tool-use turns get the same treatment once they are no longer
// the latest assistant message. When a turn carries toolCall blocks but did
// NOT request tool execution (stopReason !== "toolUse" — e.g.
// adaptive-thinking Opus emitting tool calls and then ending the turn on
// `end_turn`/`stop`), the agent loop pairs those calls with placeholder
// tool_results to keep the tool_use/tool_result contract valid. Historical
// abandoned turns cannot safely replay their end_turn-bound signatures in
// that continuation, so stripping downgrades them to plain text downstream.
// Latest abandoned turns are exempt because Anthropic requires thinking
// blocks from its most recent response to remain byte-for-byte unmodified.
const invalidStopReason = assistantMsg.stopReason === "aborted" || assistantMsg.stopReason === "error";
const abandonedToolUse =
!invalidStopReason &&
assistantMsg.stopReason !== "toolUse" &&
assistantMsg.content.some(b => b.type === "toolCall");
const hasInvalidSignatures = invalidStopReason || abandonedToolUse;
const transformedContent = assistantMsg.content.flatMap(block => {
if (block.type === "thinking") {
// Strip untrustworthy signatures so the encoder can downgrade to text.
const sanitized =
hasInvalidSignatures && block.thinkingSignature ? { ...block, thinkingSignature: undefined } : block;
if (mustPreserveLatestAnthropicThinking) return abandonedToolUse ? block : sanitized;
// For same model: keep thinking blocks with signatures (needed for replay)
// even if the thinking text is empty (OpenAI encrypted reasoning)
if (isSameModel && sanitized.thinkingSignature) return sanitized;
// Skip empty thinking blocks, convert others to plain text
if (!sanitized.thinking || sanitized.thinking.trim() === "") return [];
if (isSameModel) return sanitized;
return {
type: "text" as const,
text: sanitized.thinking,
};
}
if (block.type === "redactedThinking") {
if (mustPreserveLatestAnthropicThinking) return block;
if (isSameModel) return block;
return [];
}
if (block.type === "text") {
if (isSameModel) return block;
return {
type: "text" as const,
text: block.text,
};
}
if (block.type === "toolCall") {
const toolCall = block as ToolCall;
let normalizedToolCall: ToolCall = toolCall;
if (!isSameModel && toolCall.thoughtSignature) {
normalizedToolCall = { ...toolCall };
delete (normalizedToolCall as { thoughtSignature?: string }).thoughtSignature;
}
if (!isSameModel && normalizeToolCallId) {
const normalizedId = normalizeToolCallId(toolCall.id, model, assistantMsg);
if (normalizedId !== toolCall.id) {
toolCallIdMap.set(toolCall.id, normalizedId);
normalizedToolCall = { ...normalizedToolCall, id: normalizedId };
}
}
return normalizedToolCall;
}
return block;
});
return {
...assistantMsg,
content: transformedContent,
};
}
return msg;
});
}),
);
const realToolResultsById = new Map<string, ToolResultMessage>();
for (const msg of transformed) {
if (msg.role === "toolResult" && !realToolResultsById.has(msg.toolCallId)) {
+1 -2
View File
@@ -209,8 +209,7 @@ const serviceProviderMap: Record<string, KeyResolver> = {
tavily: "TAVILY_API_KEY",
parallel: "PARALLEL_API_KEY",
kagi: "KAGI_API_KEY",
// GitHub Copilot uses GitHub personal access token
"github-copilot": () => $pickenv("COPILOT_GITHUB_TOKEN", "GH_TOKEN", "GITHUB_TOKEN"),
"github-copilot": "COPILOT_GITHUB_TOKEN",
// Foundry mode optionally switches Anthropic auth to enterprise gateway credentials.
anthropic: () =>
isFoundryEnabled()
@@ -600,12 +600,17 @@ export function modelMayLeakDsmlToolCalls(provider: string, modelId: string): bo
);
}
/** Cheap model/provider gate for MiniMax plain thinking tag leaks. */
export function modelMayLeakThinkingTags(provider: string, modelId: string): boolean {
return /minimax/i.test(provider) || /minimax/i.test(modelId);
}
export function getStreamMarkupHealingPattern(
provider: string,
modelId: string,
options?: { readonly parseThinkingTags?: boolean },
): StreamMarkupHealingPattern | undefined {
if (options?.parseThinkingTags) return "thinking";
if (options?.parseThinkingTags || modelMayLeakThinkingTags(provider, modelId)) return "thinking";
if (modelMayLeakKimiToolCalls(provider, modelId)) return "kimi";
if (modelMayLeakDsmlToolCalls(provider, modelId)) return "dsml";
return undefined;
@@ -32,6 +32,43 @@ describe("Duplicate Tool Results Regression", () => {
reasoning: true,
};
const makeEvalAssistantMessage = (id: string, timestamp: number): AssistantMessage => ({
role: "assistant",
content: [{ type: "toolCall", id, name: "eval", arguments: {} }],
api: "anthropic-messages",
provider: "anthropic",
model: "claude-3-5-sonnet-20241022",
usage: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "toolUse",
timestamp,
});
const makeEvalToolResult = (id: string, text: string, timestamp: number): ToolResultMessage => ({
role: "toolResult",
toolCallId: id,
toolName: "eval",
content: [{ type: "text", text }],
isError: false,
timestamp,
});
const getAssistantToolIds = (messages: Message[]): string[] =>
messages.flatMap(message =>
message.role === "assistant"
? message.content.filter((block): block is ToolCall => block.type === "toolCall").map(block => block.id)
: [],
);
const getToolResults = (messages: Message[]): ToolResultMessage[] =>
messages.filter((message): message is ToolResultMessage => message.role === "toolResult");
it("should not duplicate tool results for errored messages when results already exist", () => {
const toolCallId = "toolu_019xqMTvqWZiTDy8XxmjxrTo";
@@ -316,6 +353,144 @@ describe("Duplicate Tool Results Regression", () => {
expect(result2.length).toBe(1);
expect(result3.length).toBe(1);
});
it("deduplicates repeated tool call ids and preserves call/result pairing", () => {
const duplicateId = "functions.eval:301";
const distinctId = "functions.eval:302";
const messages: Message[] = [
makeEvalAssistantMessage(duplicateId, 1),
makeEvalToolResult(duplicateId, "first", 2),
makeEvalAssistantMessage(duplicateId, 3),
makeEvalToolResult(duplicateId, "second", 4),
makeEvalAssistantMessage(duplicateId, 5),
makeEvalAssistantMessage(distinctId, 6),
makeEvalToolResult(distinctId, "third", 7),
];
const transformed = transformMessages(messages, model);
const assistantToolIds = getAssistantToolIds(transformed);
const toolResults = getToolResults(transformed);
expect(assistantToolIds).toEqual([duplicateId, `${duplicateId}_dup1`, `${duplicateId}_dup2`, distinctId]);
expect(toolResults.map(result => result.toolCallId)).toEqual([
duplicateId,
`${duplicateId}_dup1`,
`${duplicateId}_dup2`,
distinctId,
]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup1`)?.content).toEqual([
{ type: "text", text: "second" },
]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup2`)?.content).toEqual([
{ type: "text", text: "No result provided" },
]);
});
it("deduplicates repeated ids without colliding with existing generated-looking ids", () => {
const duplicateId = "functions.eval:301";
const generatedLookingId = `${duplicateId}_dup1`;
const messages: Message[] = [
makeEvalAssistantMessage(duplicateId, 1),
makeEvalToolResult(duplicateId, "first", 2),
makeEvalAssistantMessage(generatedLookingId, 3),
makeEvalToolResult(generatedLookingId, "already-used", 4),
makeEvalAssistantMessage(duplicateId, 5),
makeEvalToolResult(duplicateId, "second", 6),
];
const transformed = transformMessages(messages, model);
const assistantToolIds = getAssistantToolIds(transformed);
const toolResults = getToolResults(transformed);
expect(assistantToolIds).toEqual([duplicateId, generatedLookingId, `${duplicateId}_dup2`]);
expect(toolResults.map(result => result.toolCallId)).toEqual([
duplicateId,
generatedLookingId,
`${duplicateId}_dup2`,
]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup2`)?.content).toEqual([
{ type: "text", text: "second" },
]);
});
it("preserves delayed duplicate tool results across message gaps", () => {
const duplicateId = "functions.eval:301";
const developerMessage: DeveloperMessage = { role: "developer", content: "handoff summary", timestamp: 4 };
const messages: Message[] = [
makeEvalAssistantMessage(duplicateId, 1),
makeEvalToolResult(duplicateId, "first", 2),
makeEvalAssistantMessage(duplicateId, 3),
developerMessage,
makeEvalToolResult(duplicateId, "second", 5),
];
const transformed = transformMessages(messages, model);
const toolResults = getToolResults(transformed);
expect(getAssistantToolIds(transformed)).toEqual([duplicateId, `${duplicateId}_dup1`]);
expect(toolResults.map(result => result.toolCallId)).toEqual([duplicateId, `${duplicateId}_dup1`]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup1`)?.content).toEqual([
{ type: "text", text: "second" },
]);
});
it("routes the late result to the most recent duplicate call when a new turn re-emits the id across a gap", () => {
const duplicateId = "functions.eval:301";
const developerMessage: DeveloperMessage = { role: "developer", content: "handoff summary", timestamp: 4 };
const messages: Message[] = [
makeEvalAssistantMessage(duplicateId, 1),
makeEvalToolResult(duplicateId, "first", 2),
makeEvalAssistantMessage(duplicateId, 3),
developerMessage,
makeEvalAssistantMessage(duplicateId, 5),
makeEvalToolResult(duplicateId, "second", 6),
];
const transformed = transformMessages(messages, model);
const toolResults = getToolResults(transformed);
expect(getAssistantToolIds(transformed)).toEqual([duplicateId, `${duplicateId}_dup1`, `${duplicateId}_dup2`]);
expect(toolResults.map(result => result.toolCallId)).toEqual([
duplicateId,
`${duplicateId}_dup1`,
`${duplicateId}_dup2`,
]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup1`)?.content).toEqual([
{ type: "text", text: "No result provided" },
]);
expect(toolResults.find(result => result.toolCallId === `${duplicateId}_dup2`)?.content).toEqual([
{ type: "text", text: "second" },
]);
});
it("keeps duplicate-id rewrites within the 64-char tool-call id limit", () => {
const baseId = `toolu_${"a".repeat(58)}`;
expect(baseId.length).toBe(64);
const messages: Message[] = [
makeEvalAssistantMessage(baseId, 1),
makeEvalToolResult(baseId, "first", 2),
makeEvalAssistantMessage(baseId, 3),
makeEvalToolResult(baseId, "second", 4),
];
const transformed = transformMessages(messages, model);
const assistantToolIds = getAssistantToolIds(transformed);
const toolResults = getToolResults(transformed);
expect(assistantToolIds).toHaveLength(2);
for (const id of assistantToolIds) {
expect(id.length).toBeLessThanOrEqual(64);
expect(id).toMatch(/^[A-Za-z0-9_-]+$/);
}
const rewrittenId = assistantToolIds[1];
expect(rewrittenId).not.toBe(baseId);
expect(rewrittenId.endsWith("_dup1")).toBe(true);
expect(toolResults.map(result => result.toolCallId)).toEqual([baseId, rewrittenId]);
expect(toolResults.find(result => result.toolCallId === rewrittenId)?.content).toEqual([
{ type: "text", text: "second" },
]);
});
});
/**
@@ -148,6 +148,7 @@ describe("StreamMarkupHealing pattern selection", () => {
expect(getStreamMarkupHealingPattern("minimax-code", "MiniMax-M2.5", { parseThinkingTags: true })).toBe(
"thinking",
);
expect(getStreamMarkupHealingPattern("opencode-zen", "minimax-m3")).toBe("thinking");
expect(getStreamMarkupHealingPattern("nanogpt", "deepseek/deepseek-v4-pro")).toBe("dsml");
expect(getStreamMarkupHealingPattern("ollama-cloud", "gpt-oss:120b")).toBeUndefined();
expect(getStreamMarkupHealingPattern("openai", "deepseek-v4-pro")).toBeUndefined();
@@ -582,6 +583,39 @@ describe("Ollama provider DSML envelope healing", () => {
});
});
describe("OpenAI completions MiniMax thinking healing", () => {
it("parses OpenCode Zen MiniMax think tags into a thinking block", async () => {
const model: Model<"openai-completions"> = {
id: "minimax-m3",
name: "MiniMax M3",
api: "openai-completions",
provider: "opencode-zen",
baseUrl: "https://opencode.ai/zen/v1",
reasoning: true,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 200_000,
maxTokens: 8_192,
};
global.fetch = mockFetch([
chunk(model.id, { content: "visible <thin" }),
chunk(model.id, { content: "k>hidden reasoning</think" }),
chunk(model.id, { content: ">" }),
chunk(model.id, { content: " answer" }),
chunk(model.id, {}, "stop"),
"[DONE]",
]);
const result = await streamOpenAICompletions(model, baseContext(), { apiKey: "test-key" }).result();
expect(result.content).toEqual([
{ type: "text", text: "visible " },
{ type: "thinking", thinking: "hidden reasoning", thinkingSignature: undefined },
{ type: "text", text: " answer" },
]);
});
});
describe("OpenAI completions provider DSML envelope healing", () => {
it("heals the envelope into a structured tool call and suppresses leaked text", async () => {
const model: Model<"openai-completions"> = {
+132
View File
@@ -652,6 +652,138 @@ describe("Generate E2E Tests", () => {
else Bun.env.GOOGLE_APPLICATION_CREDENTIALS = originalGac;
}
});
it("routes impersonated_service_account ADC through IAM to the Vertex request", async () => {
const originalProject = Bun.env.GOOGLE_CLOUD_PROJECT;
const originalGcpProject = Bun.env.GCP_PROJECT;
const originalGcloudProject = Bun.env.GCLOUD_PROJECT;
const originalVertexLocation = Bun.env.GOOGLE_VERTEX_LOCATION;
const originalCloudLocation = Bun.env.GOOGLE_CLOUD_LOCATION;
const originalLocation = Bun.env.VERTEX_LOCATION;
const originalApiKey = Bun.env.GOOGLE_CLOUD_API_KEY;
const originalGac = Bun.env.GOOGLE_APPLICATION_CREDENTIALS;
const tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-vertex-impersonation-"));
const adcPath = path.join(tmpDir, "impersonated-adc.json");
await Bun.write(
adcPath,
JSON.stringify({
type: "impersonated_service_account",
service_account_impersonation_url:
"https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/target@project.iam.gserviceaccount.com:generateAccessToken",
source_credentials: {
type: "authorized_user",
client_id: "client-id",
client_secret: "client-secret",
refresh_token: "refresh-token",
},
delegates: ["projects/-/serviceAccounts/delegate@project.iam.gserviceaccount.com"],
}),
);
const model: Model<"anthropic-messages"> = {
id: "claude-sonnet-4@20250514",
name: "Claude Sonnet 4",
api: "anthropic-messages",
provider: "google-vertex",
baseUrl:
"https://{location}-aiplatform.googleapis.com/v1/projects/{project}/locations/{location}/publishers/anthropic/models/claude-sonnet-4@20250514:streamRawPredict",
reasoning: true,
input: ["text", "image"],
cost: { input: 3, output: 15, cacheRead: 0.3, cacheWrite: 3.75 },
contextWindow: 200_000,
maxTokens: 64_000,
};
const callOrder: string[] = [];
let iamRequest: { url: string; authorization: string | null; body: unknown } | undefined;
const captured = Promise.withResolvers<{ url: string; authorization: string | null }>();
try {
__resetVertexTokenCache();
Bun.env.GOOGLE_CLOUD_PROJECT = "vertex-project";
Bun.env.GOOGLE_VERTEX_LOCATION = "global";
delete Bun.env.GCP_PROJECT;
delete Bun.env.GCLOUD_PROJECT;
delete Bun.env.GOOGLE_CLOUD_LOCATION;
delete Bun.env.VERTEX_LOCATION;
delete Bun.env.GOOGLE_CLOUD_API_KEY;
Bun.env.GOOGLE_APPLICATION_CREDENTIALS = adcPath;
const events = stream(
model,
{ messages: [{ role: "user", content: "Hello", timestamp: Date.now() }] },
{
apiKey: "<authenticated>",
fetch: async (input, init) => {
const url = input instanceof Request ? input.url : input.toString();
const headers = input instanceof Request ? input.headers : new Headers(init?.headers);
if (url === "https://oauth2.googleapis.com/token") {
callOrder.push("source");
return new Response(JSON.stringify({ access_token: "source-token", expires_in: 3600 }));
}
if (url.startsWith("https://iamcredentials.googleapis.com/")) {
callOrder.push("iam");
const bodyText =
input instanceof Request ? await input.clone().text() : String(init?.body ?? "");
iamRequest = { url, authorization: headers.get("authorization"), body: JSON.parse(bodyText) };
return new Response(
JSON.stringify({
accessToken: "impersonated-token",
expireTime: new Date(Date.now() + 3_600_000).toISOString(),
}),
);
}
callOrder.push("vertex");
captured.resolve({ url, authorization: headers.get("authorization") });
return new Response(JSON.stringify({ error: { message: "stop after capture" } }), { status: 400 });
},
},
);
for await (const _event of events) {
}
const request = await captured.promise;
// Source refresh, then IAM generateAccessToken, then the actual Vertex call.
expect(callOrder).toEqual(["source", "iam", "vertex"]);
// IAM exchange is authorized by the freshly minted source token, posts the
// reconstructed canonical URL, and forwards the configured delegates verbatim.
expect(iamRequest?.url).toBe(
"https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/target@project.iam.gserviceaccount.com:generateAccessToken",
);
expect(iamRequest?.authorization).toBe("Bearer source-token");
expect(iamRequest?.body).toEqual({
delegates: ["projects/-/serviceAccounts/delegate@project.iam.gserviceaccount.com"],
scope: ["https://www.googleapis.com/auth/cloud-platform"],
lifetime: "3600s",
});
// The impersonated token (not the source token) authorizes the Vertex request.
expect(request.url).toBe(
"https://aiplatform.googleapis.com/v1/projects/vertex-project/locations/global/publishers/anthropic/models/claude-sonnet-4@20250514:streamRawPredict",
);
expect(request.authorization).toBe("Bearer impersonated-token");
} finally {
__resetVertexTokenCache();
await fs.rm(tmpDir, { recursive: true, force: true });
if (originalProject === undefined) delete Bun.env.GOOGLE_CLOUD_PROJECT;
else Bun.env.GOOGLE_CLOUD_PROJECT = originalProject;
if (originalGcpProject === undefined) delete Bun.env.GCP_PROJECT;
else Bun.env.GCP_PROJECT = originalGcpProject;
if (originalGcloudProject === undefined) delete Bun.env.GCLOUD_PROJECT;
else Bun.env.GCLOUD_PROJECT = originalGcloudProject;
if (originalVertexLocation === undefined) delete Bun.env.GOOGLE_VERTEX_LOCATION;
else Bun.env.GOOGLE_VERTEX_LOCATION = originalVertexLocation;
if (originalCloudLocation === undefined) delete Bun.env.GOOGLE_CLOUD_LOCATION;
else Bun.env.GOOGLE_CLOUD_LOCATION = originalCloudLocation;
if (originalLocation === undefined) delete Bun.env.VERTEX_LOCATION;
else Bun.env.VERTEX_LOCATION = originalLocation;
if (originalApiKey === undefined) delete Bun.env.GOOGLE_CLOUD_API_KEY;
else Bun.env.GOOGLE_CLOUD_API_KEY = originalApiKey;
if (originalGac === undefined) delete Bun.env.GOOGLE_APPLICATION_CREDENTIALS;
else Bun.env.GOOGLE_APPLICATION_CREDENTIALS = originalGac;
}
});
});
describe("Google Vertex Provider (gemini-3-flash-preview)", () => {
+29
View File
@@ -1,6 +1,33 @@
# Changelog
## [Unreleased]
### Fixed
- Fixed a flaky JS eval worker startup that intermittently failed unrelated CI runs. The worker-ready wait reused Bun's 5s default per-test timeout as its floor, so a slow cold-start under `--isolate` + high concurrency was aborted mid-init; terminating a still-initializing Bun worker is the documented SIGILL/SIGTRAP crash trigger, which took down the whole test file. Worker init now floors at a fixed 15s infrastructure budget (independent of, and still dominated by, a larger per-cell `timeout`), and the JS eval test suites set a 20s file-local timeout so cold starts complete instead of being torn down.
### Fixed
- Fixed reviewer-style subagent yields crashing the calling eval cell when a caller-supplied output schema declares `additionalProperties: false` without a `findings` property. `normalizeCompleteData` now consults the active validator before splicing collected `report_finding` entries onto the yielded payload, so injection is suppressed when the schema would reject it — keeping the executor's post-mortem validation in lockstep with the in-tool `yield` validation that already accepted the same raw payload ([#2070](https://github.com/can1357/oh-my-pi/issues/2070))
### Fixed
- Fixed Anthropic empty `toolUse` stops without tool calls corrupting session history by retrying them and removing orphaned turns even at the retry cap.
### Fixed
- Fixed MCP tools hanging in non-yolo modes by declaring `approval = "write"` on `MCPTool` and `DeferredMCPTool`, and propagating the `approval` property through `customToolToDefinition()` in `sdk.ts`
### Fixed
- Fixed Kitty OSC 5522 paste rejecting plain text as "no supported text or image data": the listing parser now decodes the `mime="."` DATA payload (whitespace-separated MIME list) Kitty actually sends, in addition to the per-type DATA packets described by the ancillary 5522-mode spec ([#2051](https://github.com/can1357/oh-my-pi/issues/2051))
### Fixed
- Fixed follow-up shortcut submission of builtin slash commands so `/goal set ...` applies goal mode instead of queueing as plain text.
### Fixed
- Fixed Ctrl+Z crashing the agent on Windows with `TypeError: Unknown signal: SIGTSTP`. `InputController.handleCtrlZ` called `process.kill(0, "SIGTSTP")` unconditionally, but `SIGTSTP` is POSIX job-control and Bun/Node on Windows rejects the signal name from the JS side; the throw propagated out of the TUI input dispatcher as an uncaught exception. The handler now no-ops with a "Suspend (Ctrl+Z) is not supported on this platform" status on Windows, and on POSIX wraps `process.kill` in a try/catch that detaches the registered SIGCONT resume hook and re-`start()`s the TUI on failure so a rejected signal can never leave the UI stranded with a leaked listener ([#2036](https://github.com/can1357/oh-my-pi/issues/2036)).
## [15.10.1] - 2026-06-07
@@ -50,6 +77,8 @@
### Fixed
- Fixed session auto-retry for generic `upstream_error: Upstream request failed` gateway failures.
- Fixed inline `find` and `search` result blocks to align with grouped `read` output and render their success headers with the normal tool-title color instead of accent blue.
- Fixed the working-status shimmer to opt into the loader's 30fps animated-message repaint path while keeping both the status spinner and pending bash/eval tool spinners on their normal 80 ms glyph cadence.
+1 -1
View File
@@ -258,7 +258,7 @@ export function getExtraHelpText(): string {
NODE_EXTRA_CA_CERTS - CA bundle path (or inline PEM) for server certificate validation
OPENAI_API_KEY - OpenAI GPT models
GEMINI_API_KEY - Google Gemini models
GITHUB_TOKEN - GitHub Copilot (or GH_TOKEN, COPILOT_GITHUB_TOKEN)
COPILOT_GITHUB_TOKEN - GitHub Copilot
${chalk.dim("# Additional LLM Providers")}
AZURE_OPENAI_API_KEY - Azure OpenAI models
@@ -53,7 +53,13 @@ interface JsSession {
const sessions = new Map<string, JsSession>();
const startingSessions = new Map<string, Promise<JsSession>>();
const resettingSessions = new Set<string>();
const READY_TIMEOUT_MS_DEFAULT = 5_000;
// Worker startup (module-graph import + WorkerCore construction) is infrastructure
// cost, not user compute. Floor it independently of Bun's 5s default per-test timeout
// so a slow cold-start under load isn't aborted mid-init — terminating a still-
// initializing Bun worker triggers the same kind of terminate-race that motivates
// avoiding `vm.runInContext` (see shared/indirect-eval.ts), here surfacing as a
// SIGILL/SIGSEGV. Callers that pass a larger per-cell budget still dominate.
const WORKER_INIT_TIMEOUT_MS = 15_000;
export async function executeInVmContext(options: {
sessionKey: string;
@@ -191,9 +197,9 @@ async function acquireSession(sessionKey: string, snapshot: SessionSnapshot, tim
handleSessionMessage(session, msg);
});
try {
// Cold-start can exceed 5s on slow hosts. Let the caller's per-cell timeout dominate so
// users can grant more headroom when they raise `timeout` on a cell.
const readyTimeoutMs = Math.max(READY_TIMEOUT_MS_DEFAULT, timeoutMs ?? 0);
// Init headroom is the fixed infrastructure floor; the caller's per-cell timeout
// dominates when larger so users can grant more by raising `timeout` on a cell.
const readyTimeoutMs = Math.max(WORKER_INIT_TIMEOUT_MS, timeoutMs ?? 0);
await raceWithTimeout(readyPromise, readyTimeoutMs, "Timed out initializing JS eval worker");
worker.send({ type: "init", snapshot });
sessions.set(sessionKey, session);
@@ -7,7 +7,13 @@
* - Register commands, keyboard shortcuts, and CLI flags
* - Interact with the user via UI primitives
*/
import type { AgentMessage, AgentToolResult, AgentToolUpdateCallback, ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import type {
AgentMessage,
AgentToolResult,
AgentToolUpdateCallback,
ThinkingLevel,
ToolApproval,
} from "@oh-my-pi/pi-agent-core";
import type { CompactionResult } from "@oh-my-pi/pi-agent-core/compaction";
import type {
Api,
@@ -392,6 +398,9 @@ export interface ToolDefinition<TParams extends TSchema = TSchema, TDetails = un
defaultInactive?: boolean;
/** If true, tool may stage deferred changes that require explicit resolve/discard. */
deferrable?: boolean;
/** Tool approval tier. Defaults to `"exec"` when omitted.
* `"read"`: read-only operations. `"write"`: mutations. `"exec"`: code execution. */
approval?: ToolApproval;
/** MCP server name for discovery/search metadata when this tool fronts an MCP server. */
mcpServerName?: string;
/** Original MCP tool name for discovery/search metadata. */
@@ -220,6 +220,7 @@ export class MCPTool implements CustomTool<TSchema, MCPToolDetails> {
readonly mcpToolName: string;
/** Server name */
readonly mcpServerName: string;
readonly approval = "write" as const;
/** Render completed MCP calls with the result header replacing the pending call header. */
readonly mergeCallAndResult = true;
@@ -305,6 +306,7 @@ export class DeferredMCPTool implements CustomTool<TSchema, MCPToolDetails> {
readonly mcpToolName: string;
/** Server name */
readonly mcpServerName: string;
readonly approval = "write" as const;
/** Render completed MCP calls with the result header replacing the pending call header. */
readonly mergeCallAndResult = true;
@@ -499,17 +499,44 @@ export class InputController {
}
handleCtrlZ(): void {
// Set up handler to restore TUI when resumed
process.once("SIGCONT", () => {
// SIGTSTP is POSIX job-control: Windows has no equivalent and
// `process.kill(_, "SIGTSTP")` throws `TypeError: Unknown signal:
// SIGTSTP` there, taking the whole agent down via an uncaught
// exception (issue #2036). No-op on platforms that cannot suspend.
if (process.platform === "win32") {
this.ctx.showStatus("Suspend (Ctrl+Z) is not supported on this platform");
return;
}
// Capture the listener so we can detach it if the signal never
// fires; otherwise a failed suspend would leave a stale SIGCONT
// handler that fires on the next unrelated continue and tries to
// re-`start()` an already-running TUI.
const onResume = (): void => {
this.ctx.ui.start();
this.ctx.ui.requestRender(true);
});
};
process.once("SIGCONT", onResume);
// Stop the TUI (restore terminal to normal mode)
// Stop the TUI (restore terminal to normal mode) before sending the
// signal so the parent shell sees a sane terminal state.
this.ctx.ui.stop();
// Send SIGTSTP to process group (pid=0 means all processes in group)
process.kill(0, "SIGTSTP");
try {
// pid=0 → entire foreground process group; the shell receives
// SIGTSTP and parks the job.
process.kill(0, "SIGTSTP");
} catch (err) {
// Either the runtime refused the signal or the kernel rejected
// it (some sandboxes block sending to pid=0). Tear the resume
// hook down and bring the TUI back so the user is not stranded
// on a frozen prompt.
process.removeListener("SIGCONT", onResume);
this.ctx.ui.start();
this.ctx.ui.requestRender(true);
const reason = err instanceof Error ? err.message : String(err);
this.ctx.showError(`Failed to suspend: ${reason}`);
}
}
handleDequeue(): void {
@@ -589,7 +616,7 @@ export class InputController {
/** Send editor text as a follow-up message (queued behind current stream). */
async handleFollowUp(): Promise<void> {
const text = this.ctx.editor.getText().trim();
let text = this.ctx.editor.getText().trim();
if (!text) return;
// Compaction first: while compacting, free text gets queued via
@@ -603,6 +630,16 @@ export class InputController {
return;
}
const slashResult = await executeBuiltinSlashCommand(text, {
ctx: this.ctx,
});
if (slashResult === true) {
return;
}
if (typeof slashResult === "string") {
text = slashResult;
}
// Skill commands invoke through the custom-message path regardless of
// which keybinding submitted them. Enter routes them as `steer`;
// Ctrl+Enter (this handler) routes them as `followUp`.
+1
View File
@@ -709,6 +709,7 @@ function customToolToDefinition(tool: CustomTool): ToolDefinition {
parameters: tool.parameters,
hidden: tool.hidden,
deferrable: tool.deferrable,
approval: typeof tool.approval === "function" ? tool.approval.bind(tool) : tool.approval,
mcpServerName: tool.mcpServerName,
mcpToolName: tool.mcpToolName,
execute: (toolCallId, params, signal, onUpdate, ctx) =>
@@ -283,6 +283,11 @@ export type AgentSessionEventListener = (event: AgentSessionEvent) => void;
export type AsyncJobSnapshotItem = Pick<AsyncJob, "id" | "type" | "status" | "label" | "startTime">;
const EMPTY_STOP_MAX_RETRIES = 3;
const NON_WHITESPACE_RE = /\S/;
function hasNonWhitespace(value: string): boolean {
return NON_WHITESPACE_RE.test(value);
}
export interface AsyncJobSnapshot {
running: AsyncJobSnapshotItem[];
@@ -6539,9 +6544,13 @@ export class AgentSession {
this.#retryAttempt = 0;
}
this.#resolveRetry();
// Tool-use orphans corrupt Anthropic message history (tool_result without
// matching tool_use). Always remove them even when the retry cap is hit.
if (assistantMessage.stopReason === "toolUse") {
this.#removeEmptyStopFromActiveContext(assistantMessage);
}
return true;
}
this.#removeEmptyStopFromActiveContext(assistantMessage);
this.agent.appendMessage({
role: "developer",
@@ -6554,12 +6563,26 @@ export class AgentSession {
}
#isEmptyAssistantStop(assistantMessage: AssistantMessage): boolean {
if (assistantMessage.stopReason !== "stop") return false;
return !assistantMessage.content.some(content => {
if (content.type === "text") return content.text.trim().length > 0;
if (content.type === "thinking") return content.thinking.trim().length > 0;
return content.type === "toolCall";
});
switch (assistantMessage.stopReason) {
case "stop":
for (const content of assistantMessage.content) {
if (content.type === "toolCall") return false;
if (content.type === "text" && hasNonWhitespace(content.text)) return false;
if (content.type === "thinking" && hasNonWhitespace(content.thinking)) return false;
}
return true;
case "toolUse":
// An orphaned toolUse stop (no tool_use block) corrupts Anthropic history:
// a later tool_result has nothing to anchor to. Thinking alone cannot anchor
// a tool_result, so it does not rescue a toolUse stop here.
for (const content of assistantMessage.content) {
if (content.type === "toolCall") return false;
if (content.type === "text" && hasNonWhitespace(content.text)) return false;
}
return true;
default:
return false;
}
}
#emptyStopRetryReminder(): string {
@@ -7874,11 +7897,12 @@ export class AgentSession {
#isTransientTransportErrorMessage(errorMessage: string): boolean {
// Match: overloaded_error, provider returned error, rate limit, 429, 500, 502, 503, 504,
// service unavailable, provider-suggested retry, network/connection/socket errors, fetch failed,
// terminated, retry delay exceeded, Bun HTTP/2 stream resets (RST_STREAM / REFUSED_STREAM /
// ENHANCE_YOUR_CALM, surfaced verbatim from src/http/h2_client/dispatch.zig)
// gateway upstream failures, terminated, retry delay exceeded, Bun HTTP/2 stream resets
// (RST_STREAM / REFUSED_STREAM / ENHANCE_YOUR_CALM, surfaced verbatim from
// src/http/h2_client/dispatch.zig)
return (
isUnexpectedSocketCloseMessage(errorMessage) ||
/overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|other side closed|fetch failed|upstream.?connect|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)/i.test(
/overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)/i.test(
errorMessage,
)
);
+38 -15
View File
@@ -34,7 +34,11 @@ import { SessionManager } from "../session/session-manager";
import { truncateTail } from "../session/streaming-output";
import type { ContextFileEntry } from "../tools";
import { normalizeSchema } from "../tools/jtd-to-json-schema";
import { buildOutputValidator, summarizeValidationFailure } from "../tools/output-schema-validator";
import {
buildOutputValidator,
type OutputValidator,
summarizeValidationFailure,
} from "../tools/output-schema-validator";
import { type ReportFindingDetails, toReviewFinding } from "../tools/review";
import { ToolAbortError } from "../tools/tool-errors";
@@ -256,21 +260,40 @@ function extractCompletionData(parsed: unknown): unknown {
return parsed;
}
function normalizeCompleteData(data: unknown, reportFindings?: ReviewFinding[]): unknown {
let normalized = parseStringifiedJson(data ?? null);
/**
* Resolve the final yielded payload, optionally splicing collected
* `report_finding` entries into a top-level `findings` array.
*
* Injection is suppressed when an active validator would reject the augmented
* payload (e.g. a caller-supplied schema with `additionalProperties: false`
* that does not declare `findings`). That keeps the in-tool yield validator
* (which only sees the raw, pre-injection data) in lockstep with this
* post-mortem validator — honoring the "accepted in-tool ⇒ accepted
* post-mortem" guarantee documented in `output-schema-validator.ts`. The
* dropped findings are still preserved verbatim in the agent's progress
* stream and JSONL artifact, so no information is lost when injection is
* suppressed.
*/
function normalizeCompleteData(
data: unknown,
reportFindings: ReviewFinding[] | undefined,
validator: OutputValidator | undefined,
): unknown {
const normalized = parseStringifiedJson(data ?? null);
if (
Array.isArray(reportFindings) &&
reportFindings.length > 0 &&
normalized &&
typeof normalized === "object" &&
!Array.isArray(normalized)
!Array.isArray(reportFindings) ||
reportFindings.length === 0 ||
!normalized ||
typeof normalized !== "object" ||
Array.isArray(normalized)
) {
const record = normalized as Record<string, unknown>;
if (!("findings" in record)) {
normalized = { ...record, findings: reportFindings };
}
return normalized;
}
return normalized;
const record = normalized as Record<string, unknown>;
if ("findings" in record) return normalized;
const injected = { ...record, findings: reportFindings };
if (validator && !validator.validate(injected).success) return normalized;
return injected;
}
function resolveFallbackCompletion(rawOutput: string, outputSchema: unknown): { data: unknown } | null {
@@ -360,13 +383,13 @@ export function finalizeSubprocessOutput(args: FinalizeSubprocessOutputArgs): Fi
if (submitData === null || submitData === undefined) {
rawOutput = rawOutput ? `${SUBAGENT_WARNING_NULL_YIELD}\n\n${rawOutput}` : SUBAGENT_WARNING_NULL_YIELD;
} else {
const completeData = normalizeCompleteData(submitData, reportFindings);
const { validator, error: schemaError } = buildOutputValidator(outputSchema);
if (schemaError) {
rawOutput = `{"error":"schema_violation","message":"invalid output schema: ${schemaError.replace(/"/g, '\\"')}"}`;
stderr = `schema_violation: invalid output schema: ${schemaError}`;
exitCode = 1;
} else {
const completeData = normalizeCompleteData(submitData, reportFindings, validator);
const result = validator?.validate(completeData) ?? { success: true as const };
if (!result.success) {
const summary = summarizeValidationFailure(result, completeData, validator?.requiredFields ?? []);
@@ -393,8 +416,8 @@ export function finalizeSubprocessOutput(args: FinalizeSubprocessOutputArgs): Fi
const hasOutputSchema = normalizedSchema !== undefined && !schemaError;
const fallback = allowFallback ? resolveFallbackCompletion(rawOutput, outputSchema) : null;
if (fallback) {
const completeData = normalizeCompleteData(fallback.data, reportFindings);
const { validator } = buildOutputValidator(outputSchema);
const completeData = normalizeCompleteData(fallback.data, reportFindings, validator);
const result = validator?.validate(completeData) ?? { success: true as const };
if (!result.success) {
const summary = summarizeValidationFailure(result, completeData, validator?.requiredFields ?? []);
@@ -7,6 +7,8 @@ const PASTE_EVENT_NAME_BASE64 = Buffer.from("Paste event", "utf8").toString("bas
const IMAGE_MIME_PRIORITY = ["image/png", "image/jpeg", "image/webp", "image/gif"] as const;
const TEXT_MIME_TYPE = "text/plain";
/** Kitty's "give me the list of available MIME types" sentinel — see `TARGETS_MIME` in `kitty/clipboard.py`. */
const MIME_LISTING_TARGET = ".";
type PasteReadKind = "image" | "text";
@@ -144,6 +146,24 @@ export class EnhancedPasteController {
if (!mimeType) return;
if (state.phase === "listing") {
// Kitty (as of writing) implements the "list available MIME types"
// response shape by sending a single DATA packet with `mime="."` and
// the available types packed into the payload as a whitespace-
// separated list (see `fulfill_read_request` in
// kovidgoyal/kitty:kitty/clipboard.py). The 5522-mode ancillary
// spec instead encodes each type as its own DATA packet with an
// empty payload. Support both — fall through to the per-packet
// form when the dot sentinel has no payload, or when the packet
// already names a concrete MIME type.
if (mimeType === MIME_LISTING_TARGET) {
if (!packet.payload) return;
const listing = decodeBase64Utf8(packet.payload);
if (!listing) return;
for (const candidate of listing.split(/\s+/)) {
if (candidate && candidate !== MIME_LISTING_TARGET) state.mimes.push(candidate);
}
return;
}
state.mimes.push(mimeType);
return;
}
@@ -192,11 +212,12 @@ export class EnhancedPasteController {
chunks: [],
};
const metadata = [`type=read`, `mime=${Buffer.from(selected.mimeType, "utf8").toString("base64")}`];
const encodedMime = Buffer.from(selected.mimeType, "utf8").toString("base64");
const metadata = ["type=read"];
if (state.loc) metadata.push(`loc=${state.loc}`);
if (state.pw) {
metadata.push(`pw=${state.pw}`, `name=${PASTE_EVENT_NAME_BASE64}`);
}
this.#handlers.write(`${OSC5522_PREFIX}${metadata.join(":")}${OSC_TERMINATOR_ST}`);
this.#handlers.write(`${OSC5522_PREFIX}${metadata.join(":")};${encodedMime}${OSC_TERMINATOR_ST}`);
}
}
@@ -51,6 +51,14 @@ function emptyStop(): MockResponse {
};
}
function orphanedToolUseStop(): MockResponse {
return {
content: [{ type: "thinking", thinking: "I should call a tool next." }],
stopReason: "toolUse",
usage: { output: 1, cacheRead: 100 },
};
}
async function createHarness(
responses: MockResponse[],
settingsOverrides: SettingsOverrides = {},
@@ -177,6 +185,45 @@ describe("AgentSession empty stop guard", () => {
).toHaveLength(1);
});
it("retries a tool-use stop that has no tool call or text", async () => {
const { session, mock } = await createHarness([
recordCall("orphan", "call-record-orphan"),
orphanedToolUseStop(),
{ content: ["finished after orphaned tool-use retry"], stopReason: "stop" },
]);
await session.prompt("record orphan");
await session.waitForIdle();
expect(mock.calls).toHaveLength(3);
expect(assistantText(session.agent.state.messages)).toContain("finished after orphaned tool-use retry");
expect(reminderMessages(session.agent.state.messages)).toHaveLength(1);
});
it("removes orphaned tool-use stops even when retry cap is hit", async () => {
const { session, mock } = await createHarness([
recordCall("gamma", "call-record-gamma"),
orphanedToolUseStop(),
orphanedToolUseStop(),
orphanedToolUseStop(),
orphanedToolUseStop(),
]);
await session.prompt("record gamma");
await session.waitForIdle();
expect(mock.calls).toHaveLength(5);
expect(reminderMessages(session.agent.state.messages)).toHaveLength(3);
const activeBranchMessages = session.sessionManager
.getBranch()
.filter(entry => entry.type === "message")
.map(entry => entry.message as AgentMessage);
const orphanedToolUseStops = activeBranchMessages.filter(
message =>
message.role === "assistant" &&
message.stopReason === "toolUse" &&
!message.content.some(content => content.type === "toolCall"),
);
expect(orphanedToolUseStops).toHaveLength(0);
});
it("caps empty stop retries at three attempts", async () => {
const { session, mock } = await createHarness([
recordCall("beta", "call-record-beta"),
@@ -336,4 +336,60 @@ describe("AgentSession retry delay cap", () => {
const last = lastAssistant(session);
expect(last.stopReason).toBe("stop");
});
it("retries generic upstream_error gateway failures", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!model) {
throw new Error("Expected bundled Anthropic test model to exist");
}
const mock = createMockModel({
responses: [
{ throw: "upstream_error: Upstream request failed" },
{ content: ["recovered after generic gateway upstream error"] },
],
});
const agent = new Agent({
getApiKey: provider => `${provider}-test-key`,
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
messages: [],
},
streamFn: mock.stream,
});
const settings = Settings.isolated({
"compaction.enabled": false,
"retry.baseDelayMs": 5,
"retry.maxDelayMs": 5_000,
"retry.maxRetries": 1,
});
settings.setModelRole("default", `${model.provider}/${model.id}`);
session = new AgentSession({
agent,
sessionManager: SessionManager.inMemory(),
settings,
modelRegistry,
});
vi.spyOn(scheduler, "wait").mockResolvedValue(undefined);
const retryStartEvents: AutoRetryStartEvent[] = [];
const retryEndEvents: AutoRetryEndEvent[] = [];
session.subscribe(event => {
if (event.type === "auto_retry_start") retryStartEvents.push(event);
if (event.type === "auto_retry_end") retryEndEvents.push(event);
});
await session.prompt("Trigger generic upstream_error");
await session.waitForIdle();
expect(retryStartEvents).toHaveLength(1);
expect(retryEndEvents).toHaveLength(1);
expect(retryEndEvents[0]).toMatchObject({ success: true, attempt: 1 });
const last = lastAssistant(session);
expect(last.stopReason).toBe("stop");
expect(last.content).toContainEqual({ type: "text", text: "recovered after generic gateway upstream error" });
});
});
@@ -1,4 +1,4 @@
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test";
import { afterAll, afterEach, beforeAll, describe, expect, it, setDefaultTimeout, vi } from "bun:test";
import * as path from "node:path";
import type { AgentTool, AgentToolResult } from "@oh-my-pi/pi-agent-core";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
@@ -8,6 +8,11 @@ import * as z from "zod/v4";
import { disposeAllVmContexts } from "../../src/eval/js/context-manager";
import { executeJs, type JsResult } from "../../src/eval/js/executor";
// JS eval cold-starts a Bun worker; under --isolate + high CI concurrency that startup
// can exceed Bun's 5s default per-test timeout, flaking the suite. Give the worker-backed
// tests headroom above the worker-init floor (context-manager WORKER_INIT_TIMEOUT_MS).
setDefaultTimeout(20_000);
function createTool(
name: string,
execute: (toolCallId: string, args: unknown, signal?: AbortSignal) => Promise<AgentToolResult>,
@@ -1,4 +1,4 @@
import { afterAll, beforeAll, describe, expect, it } from "bun:test";
import { afterAll, beforeAll, describe, expect, it, setDefaultTimeout } from "bun:test";
import * as path from "node:path";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
@@ -6,6 +6,11 @@ import { TempDir } from "@oh-my-pi/pi-utils";
import { disposeAllVmContexts } from "../../src/eval/js/context-manager";
import { executeJs, type JsResult } from "../../src/eval/js/executor";
// JS eval cold-starts a Bun worker; under --isolate + high CI concurrency that startup
// can exceed Bun's 5s default per-test timeout, flaking the suite. Give the worker-backed
// tests headroom above the worker-init floor (context-manager WORKER_INIT_TIMEOUT_MS).
setDefaultTimeout(20_000);
function statusEvents(result: JsResult) {
return result.displayOutputs.filter(
(output): output is Extract<JsResult["displayOutputs"][number], { type: "status" }> => output.type === "status",
@@ -71,6 +71,8 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
// Annotate parameters so `mock.calls[N]` is typed as a tuple (not `[]`) and
// `message` carries required skill prompt details for assertion below.
const promptCustomMessage = vi.fn(async (_message: { details: SkillPromptDetails }, _options?: unknown) => {});
const prompt = vi.fn(async (_text: string, _options?: unknown) => {});
const handleGoalModeCommand = vi.fn(async (_rest?: string) => {});
const updatePendingMessagesDisplay = vi.fn();
const requestRender = vi.fn();
const showError = vi.fn();
@@ -86,9 +88,12 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
isEvalRunning: false,
extensionRunner: undefined,
enqueueCustomMessageDisplay,
prompt,
promptCustomMessage,
},
showError,
handleGoalModeCommand,
goalModeEnabled: false,
updatePendingMessagesDisplay,
// Defaults that InputController touches on submit but don't matter here.
isBashMode: false,
@@ -101,7 +106,7 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
withLocalSubmission: async (_text: string, fn: () => unknown) => fn(),
} as unknown as InteractiveModeContext;
return { ctx, editor, enqueueCustomMessageDisplay, promptCustomMessage };
return { ctx, editor, enqueueCustomMessageDisplay, prompt, promptCustomMessage, handleGoalModeCommand };
}
describe("InputController #invokeSkillCommand (E1-E3)", () => {
@@ -166,6 +171,21 @@ describe("InputController #invokeSkillCommand (E1-E3)", () => {
expect(messageArg.details.__pendingDisplayTag).toBe("sk-test-0");
});
it("E2b: streaming follow-up applies builtin slash commands instead of queueing them", async () => {
const { ctx, editor, prompt, handleGoalModeCommand } = createStubInputControllerContext({
skillCommands,
isStreaming: true,
});
const controller = new InputController(ctx);
editor.setText("/goal set Ship the release");
await controller.handleFollowUp();
expect(handleGoalModeCommand).toHaveBeenCalledWith("set Ship the release");
expect(prompt).not.toHaveBeenCalled();
expect(editor.getText()).toBe("");
});
it("E3: not streaming -> enqueueCustomMessageDisplay NOT called and tag absent", async () => {
const { ctx, editor, enqueueCustomMessageDisplay, promptCustomMessage } = createStubInputControllerContext({
skillCommands,
@@ -0,0 +1,128 @@
import { afterEach, describe, expect, it, type Mock, vi } from "bun:test";
import { InputController } from "../src/modes/controllers/input-controller";
import type { InteractiveModeContext } from "../src/modes/types";
interface SuspendCtx {
ctx: InteractiveModeContext;
ui: {
start: Mock<() => void>;
stop: Mock<() => void>;
requestRender: Mock<(force?: boolean) => void>;
};
showStatus: Mock<(message: string) => void>;
showError: Mock<(message: string) => void>;
}
function createCtx(): SuspendCtx {
const ui = {
start: vi.fn(),
stop: vi.fn(),
requestRender: vi.fn(),
};
const showStatus = vi.fn();
const showError = vi.fn();
const ctx = {
ui: ui as unknown as InteractiveModeContext["ui"],
showStatus,
showError,
} as unknown as InteractiveModeContext;
return { ctx, ui, showStatus, showError };
}
const originalPlatform = process.platform;
function setPlatform(value: NodeJS.Platform): void {
Object.defineProperty(process, "platform", { value, configurable: true, writable: true });
}
afterEach(() => {
Object.defineProperty(process, "platform", { value: originalPlatform, configurable: true, writable: true });
vi.restoreAllMocks();
// Drop any SIGCONT listener a passing test left behind so a later test
// (or the next file) doesn't get spurious callbacks.
process.removeAllListeners("SIGCONT");
});
describe("InputController.handleCtrlZ", () => {
it("no-ops on Windows so the unsupported SIGTSTP signal can't crash the process (#2036)", () => {
setPlatform("win32");
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => {
throw new Error("process.kill must not be called on win32");
});
const onceSpy = vi.spyOn(process, "once");
const { ctx, ui, showStatus, showError } = createCtx();
const controller = new InputController(ctx);
expect(() => controller.handleCtrlZ()).not.toThrow();
expect(killSpy).not.toHaveBeenCalled();
expect(onceSpy).not.toHaveBeenCalledWith("SIGCONT", expect.anything());
expect(ui.stop).not.toHaveBeenCalled();
expect(ui.start).not.toHaveBeenCalled();
expect(showStatus).toHaveBeenCalledTimes(1);
expect(showStatus.mock.calls[0]?.[0]).toMatch(/not supported/i);
expect(showError).not.toHaveBeenCalled();
});
it("sends SIGTSTP to the process group and registers a SIGCONT resume hook on POSIX", () => {
setPlatform("linux");
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => true);
const onceSpy = vi.spyOn(process, "once");
const { ctx, ui, showError } = createCtx();
const controller = new InputController(ctx);
controller.handleCtrlZ();
// Resume hook registered BEFORE the signal is sent so a same-tick
// SIGCONT delivery can't race past us.
expect(onceSpy).toHaveBeenCalledWith("SIGCONT", expect.any(Function));
const sigcontOrder = onceSpy.mock.invocationCallOrder[0] ?? Infinity;
const stopOrder = ui.stop.mock.invocationCallOrder[0] ?? Infinity;
const killOrder = killSpy.mock.invocationCallOrder[0] ?? Infinity;
expect(sigcontOrder).toBeLessThan(stopOrder);
expect(stopOrder).toBeLessThan(killOrder);
expect(killSpy).toHaveBeenCalledTimes(1);
expect(killSpy).toHaveBeenCalledWith(0, "SIGTSTP");
expect(ui.start).not.toHaveBeenCalled();
expect(showError).not.toHaveBeenCalled();
// Simulating the kernel-delivered SIGCONT drives the TUI back up.
const resume = onceSpy.mock.calls.find(([sig]) => sig === "SIGCONT")?.[1] as (() => void) | undefined;
expect(resume).toBeDefined();
resume?.();
expect(ui.start).toHaveBeenCalledTimes(1);
expect(ui.requestRender).toHaveBeenCalledWith(true);
});
it("restores the TUI and drops the SIGCONT listener when process.kill rejects the signal", () => {
setPlatform("linux");
const killSpy = vi.spyOn(process, "kill").mockImplementation(() => {
throw new Error("Unknown signal: SIGTSTP");
});
const onceSpy = vi.spyOn(process, "once");
const removeSpy = vi.spyOn(process, "removeListener");
const { ctx, ui, showError, showStatus } = createCtx();
const controller = new InputController(ctx);
// Critical contract: the failure must not bubble up to the caller —
// otherwise the TUI's stdin reader (which invoked us) crashes the
// whole process via `[Uncaught Exception]`.
expect(() => controller.handleCtrlZ()).not.toThrow();
// The exact listener we registered for SIGCONT is the one we
// remove; otherwise a leaked handler would fire on the next
// unrelated continue and re-`start()` an already-running TUI.
const registered = onceSpy.mock.calls.find(([sig]) => sig === "SIGCONT")?.[1];
expect(registered).toBeDefined();
expect(removeSpy).toHaveBeenCalledWith("SIGCONT", registered);
expect(killSpy).toHaveBeenCalledTimes(1);
expect(ui.stop).toHaveBeenCalledTimes(1);
expect(ui.start).toHaveBeenCalledTimes(1);
expect(ui.requestRender).toHaveBeenCalledWith(true);
expect(showError).toHaveBeenCalledTimes(1);
expect(showError.mock.calls[0]?.[0]).toMatch(/Failed to suspend/);
expect(showStatus).not.toHaveBeenCalled();
});
});
@@ -115,6 +115,17 @@ describe("MCP fallback and prompt formatting", () => {
expect(resolveApproval(subject, {}, "yolo")).toMatchObject({ policy: "allow", tier: "exec" });
});
it("allows MCP tools with write approval in write mode", () => {
const subject = tool("mcp__server__safe", "write");
expect(resolveApproval(subject, {}, "write")).toMatchObject({ policy: "allow", tier: "write" });
expect(resolveApproval(subject, {}, "yolo")).toMatchObject({ policy: "allow", tier: "write" });
});
it("prompts for MCP tools with write approval in always-ask mode", () => {
const subject = tool("mcp__server__safe", "write");
expect(resolveApproval(subject, {}, "always-ask")).toMatchObject({ policy: "prompt", tier: "write" });
});
it("formats MCP origin, reason, and per-tool details", () => {
const subject = tool("mcp__server__dangerous", undefined, () => ["Path: /tmp/out", "Content:\nhello"]);
expect(formatApprovalPrompt(subject, {}, "Needs confirmation").split("\n")).toEqual([
@@ -132,3 +132,170 @@ describe("toReviewFinding", () => {
expect(parsed.findings[0].priority).toBe(2);
});
});
describe("findings injection respects active output schema", () => {
const finding = toReviewFinding({
title: "[P0] Example finding",
body: "Details",
priority: "P0",
confidence: 0.95,
file_path: "/tmp/example.ts",
line_start: 10,
line_end: 12,
});
// Reproduces #2070: a caller-supplied JSON Schema with
// `additionalProperties: false` and no `findings` property is silently
// rejected post-mortem after the in-tool yield accepted it, because the
// executor auto-injects `findings` from `report_finding`. The injection
// must respect the active schema so that "accepted in-tool ⇒ accepted
// post-mortem" is honored.
it("suppresses findings injection when the schema forbids additional properties", () => {
const callerSchema = {
type: "object",
additionalProperties: false,
required: ["verdict", "acceptance_summary"],
properties: {
verdict: { enum: ["accept", "needs-work", "honest-stop"] },
acceptance_summary: { type: "string" },
},
};
const data = { verdict: "accept", acceptance_summary: "All good." };
const result = finalizeSubprocessOutput({
rawOutput: "",
exitCode: 0,
stderr: "",
doneAborted: false,
signalAborted: false,
yieldItems: [{ status: "success", data }],
reportFindings: [finding],
outputSchema: callerSchema,
});
expect(result.exitCode).toBe(0);
expect(result.stderr).toBe("");
const parsed = JSON.parse(result.rawOutput) as Record<string, unknown>;
expect(parsed).toEqual(data);
expect("findings" in parsed).toBe(false);
});
it("still injects findings when the schema declares them (bundled reviewer JTD)", () => {
// Mirrors the bundled reviewer agent's output: findings live under
// `optionalProperties`, so injection must continue to flow through.
const reviewerSchema = {
properties: {
overall_correctness: { enum: ["correct", "incorrect"] },
explanation: { type: "string" },
confidence: { type: "number" },
},
optionalProperties: {
findings: {
elements: {
properties: {
title: { type: "string" },
body: { type: "string" },
priority: { type: "number" },
confidence: { type: "number" },
file_path: { type: "string" },
line_start: { type: "number" },
line_end: { type: "number" },
},
},
},
},
};
const data = {
overall_correctness: "incorrect",
explanation: "Found one bug",
confidence: 0.9,
};
const result = finalizeSubprocessOutput({
rawOutput: "",
exitCode: 0,
stderr: "",
doneAborted: false,
signalAborted: false,
yieldItems: [{ status: "success", data }],
reportFindings: [finding],
outputSchema: reviewerSchema,
});
expect(result.exitCode).toBe(0);
expect(result.stderr).toBe("");
const parsed = JSON.parse(result.rawOutput) as { findings: Array<{ priority: number }> };
expect(parsed.findings).toHaveLength(1);
expect(parsed.findings[0].priority).toBe(0);
});
it("still injects findings when no schema is declared (legacy free-form)", () => {
const data = { note: "freeform" };
const result = finalizeSubprocessOutput({
rawOutput: "",
exitCode: 0,
stderr: "",
doneAborted: false,
signalAborted: false,
yieldItems: [{ status: "success", data }],
reportFindings: [finding],
outputSchema: undefined,
});
expect(result.exitCode).toBe(0);
const parsed = JSON.parse(result.rawOutput) as { findings?: unknown[] };
expect(parsed.findings).toHaveLength(1);
});
it("still injects findings when additionalProperties is open", () => {
const openSchema = {
type: "object",
required: ["verdict"],
properties: { verdict: { type: "string" } },
};
const data = { verdict: "accept" };
const result = finalizeSubprocessOutput({
rawOutput: "",
exitCode: 0,
stderr: "",
doneAborted: false,
signalAborted: false,
yieldItems: [{ status: "success", data }],
reportFindings: [finding],
outputSchema: openSchema,
});
expect(result.exitCode).toBe(0);
const parsed = JSON.parse(result.rawOutput) as { findings?: unknown[]; verdict: string };
expect(parsed.verdict).toBe("accept");
expect(parsed.findings).toHaveLength(1);
});
it("suppresses injection on the fallback (no-yield) path when the schema forbids it", () => {
const callerSchema = {
type: "object",
additionalProperties: false,
required: ["verdict"],
properties: { verdict: { type: "string" } },
};
const rawJson = JSON.stringify({ data: { verdict: "accept" } });
const result = finalizeSubprocessOutput({
rawOutput: rawJson,
exitCode: 0,
stderr: "",
doneAborted: false,
signalAborted: false,
yieldItems: undefined,
reportFindings: [finding],
outputSchema: callerSchema,
});
expect(result.exitCode).toBe(0);
expect(result.stderr).toBe("");
const parsed = JSON.parse(result.rawOutput) as Record<string, unknown>;
expect(parsed).toEqual({ verdict: "accept" });
});
});
@@ -34,7 +34,7 @@ describe("EnhancedPasteController", () => {
controller.handleInput(packet("type=read:status=DONE"));
const pasteEventName = Buffer.from("Paste event", "utf8").toString("base64");
expect(writes.at(-1)).toBe(`${OSC}type=read:mime=${imageMime}:pw=${password}:name=${pasteEventName}${ST}`);
expect(writes.at(-1)).toBe(`${OSC}type=read:pw=${password}:name=${pasteEventName};${imageMime}${ST}`);
controller.handleInput(packet("type=read:status=OK"));
controller.handleInput(
@@ -68,15 +68,13 @@ describe("EnhancedPasteController", () => {
const textMime = Buffer.from("text/plain", "utf8").toString("base64");
const password = Buffer.from("secret456", "utf8").toString("base64");
const pasteEventName = Buffer.from("Paste event", "utf8").toString("base64");
expect(controller.handleInput("plain text")).toBe(false);
controller.handleInput(packet(`type=read:status=OK:loc=primary:pw=${password}`));
controller.handleInput(packet(`type=read:status=DATA:mime=${textMime}`));
controller.handleInput(packet("type=read:status=DONE"));
expect(writes).toHaveLength(1);
expect(writes[0]).toContain(`mime=${textMime}`);
expect(writes[0]).toContain("loc=primary");
expect(writes[0]).toContain(`pw=${password}`);
expect(writes).toEqual([`${OSC}type=read:loc=primary:pw=${password}:name=${pasteEventName};${textMime}${ST}`]);
controller.handleInput(packet("type=read:status=OK"));
controller.handleInput(
@@ -106,4 +104,73 @@ describe("EnhancedPasteController", () => {
expect(statuses).toEqual(["Clipboard paste has no supported text or image data"]);
});
it("decodes Kitty's dot-listing DATA payload to discover plain-text and request it", () => {
const writes: string[] = [];
const pastedText: string[] = [];
const controller = new EnhancedPasteController({
write: data => writes.push(data),
pasteText: text => pastedText.push(text),
pasteImage: () => {
throw new Error("unexpected image paste");
},
showStatus: message => pastedText.push(`status:${message}`),
});
const dot = Buffer.from(".", "utf8").toString("base64");
const textMime = Buffer.from("text/plain", "utf8").toString("base64");
const password = Buffer.from("secret-token-123", "utf8").toString("base64");
const pasteEventName = Buffer.from("Paste event", "utf8").toString("base64");
// Kitty bundles the available MIME types into a single DATA packet
// whose `mime` field is the literal `.` and whose payload carries a
// whitespace-separated, base64-encoded list (e.g. "text/plain\n").
controller.handleInput(packet(`type=read:status=OK:pw=${password}`));
controller.handleInput(
packet(
`type=read:status=DATA:mime=${dot}:pw=${password}`,
Buffer.from("text/plain\n", "utf8").toString("base64"),
),
);
controller.handleInput(packet(`type=read:status=DONE:pw=${password}`));
expect(writes.at(-1)).toBe(`${OSC}type=read:pw=${password}:name=${pasteEventName};${textMime}${ST}`);
controller.handleInput(packet("type=read:status=OK"));
controller.handleInput(
packet(`type=read:status=DATA:mime=${textMime}`, Buffer.from("hello", "utf8").toString("base64")),
);
controller.handleInput(
packet(`type=read:status=DATA:mime=${textMime}`, Buffer.from(" world", "utf8").toString("base64")),
);
controller.handleInput(packet("type=read:status=DONE"));
expect(pastedText).toEqual(["hello world"]);
});
it("prefers images when Kitty's dot-listing payload advertises multiple MIME types", () => {
const writes: string[] = [];
const controller = new EnhancedPasteController({
write: data => writes.push(data),
pasteText: () => {
throw new Error("unexpected text paste");
},
pasteImage: () => {},
showStatus: () => {},
});
const dot = Buffer.from(".", "utf8").toString("base64");
const imageMime = Buffer.from("image/png", "utf8").toString("base64");
controller.handleInput(packet("type=read:status=OK"));
controller.handleInput(
packet(
`type=read:status=DATA:mime=${dot}`,
Buffer.from("text/plain image/png text/html\n", "utf8").toString("base64"),
),
);
controller.handleInput(packet("type=read:status=DONE"));
expect(writes.at(-1)).toBe(`${OSC}type=read;${imageMime}${ST}`);
});
});
+7
View File
@@ -6,6 +6,13 @@
- Added `super` modifier support to native key parsing/matching and bound `super+alt+backspace` / `super+alt+delete` (and `super+alt+d`) into the word-delete defaults so Ghostty's default macOS Option+Backspace wire (`ESC [127;11u` — kitty modifier 11 = super|alt) deletes a word instead of falling through to single-char delete ([#2064](https://github.com/can1357/oh-my-pi/issues/2064)).
### Fixed
- Fixed the kitty keyboard progressive-enhancement probe to honor the `CSI ? <flags> u` reply even when the terminal answers the DA1 sentinel first. Previously the kitty reply was discarded once the DA1-driven `modifyOtherKeys` fallback engaged, so terminals like Superset/xterm-on-Electron stayed on the fallback and delivered Shift+Enter as a bare `\r` ([#2042](https://github.com/can1357/oh-my-pi/issues/2042)).
- Bounded TUI line fitting for oversized raw rows so ANSI-heavy subagent output cannot grow render buffers independently of the viewport ([#2045](https://github.com/can1357/oh-my-pi/issues/2045)).
- Fixed tmux offscreen-shrink frames to skip repainting when the visible tail is unchanged, avoiding intermittent blank/refresh flashes in pane terminals ([#2046](https://github.com/can1357/oh-my-pi/issues/2046)).
- Fixed Windows ConPTY hosts (Windows Terminal, Tabby, Hyper, VS Code) parking the viewport at the top of a full paint after a `/resume` or any long-session repaint. `ProcessTerminal#safeWrite` now splits oversized writes into ≤ 8 KiB pieces at line boundaries on `win32` and inside WSL (where stdout still crosses ConPTY at the `wslhost` boundary) so each underlying `WriteFile` stays below the ~32 KiB threshold where ConPTY stops tracking the cursor; the data was always delivered, but the host UI's scroll position would not follow until any focus event forced a re-query. ([#2034](https://github.com/can1357/oh-my-pi/issues/2034))
## [15.10.1] - 2026-06-07
### Breaking Changes
+79 -2
View File
@@ -9,6 +9,58 @@ const TERMINAL_PROGRESS_KEEPALIVE_MS = 1000;
const TERMINAL_PROGRESS_ACTIVE_SEQUENCE = "\x1b]9;4;3\x07";
const TERMINAL_PROGRESS_CLEAR_SEQUENCE = "\x1b]9;4;0;\x07";
/**
* Maximum bytes per `process.stdout.write` call on Windows.
*
* Windows ConPTY ties viewport tracking to per-`WriteFile` boundaries: when a
* single write exceeds ~32-64 KB, the pseudo-console stops following the
* cursor and the host UI's viewport stays parked at whatever scroll position
* the write started from. The visible symptom is that a full-paint of a long
* session (resume, history rebuild, large permission dialog) shows only the
* first ~30 lines until any focus event forces the host to re-query the
* cursor. The data is delivered correctly — it's purely a viewport-sync bug.
*
* 8 KiB is well below the 32 KiB threshold reported on Windows Terminal and
* leaves headroom for the other ConPTY hosts (Tabby, Hyper, VS Code) where
* the exact limit is undocumented. The cost is a handful of extra syscalls
* per full paint — invisible compared to the cost of the paint itself.
*/
const MAX_CONPTY_WRITE_CHUNK = 8 * 1024;
/**
* Split `data` into chunks no larger than `maxChunkSize`, preferring a line
* boundary (`\n`) as the cut point so escape sequences (which never contain
* `\n`) stay intact. The TUI's full-paint buffers are line-structured
* (`buffer += "\r\n"` between rows), so a newline almost always exists within
* the window. The fallback for a buffer with no newline in range is a hard
* cut at `maxChunkSize`: the ConPTY viewport bug from a single oversized
* write is strictly worse than a one-frame escape-sequence glitch on a buffer
* the renderer effectively never produces.
*
* Exported for unit testing of the chunking contract; `#safeWrite` is the
* sole production caller.
*/
export function chunkForConPTY(data: string, maxChunkSize: number = MAX_CONPTY_WRITE_CHUNK): string[] {
if (data.length <= maxChunkSize) return [data];
const chunks: string[] = [];
let pos = 0;
while (pos < data.length) {
const remaining = data.length - pos;
if (remaining <= maxChunkSize) {
chunks.push(data.slice(pos));
break;
}
const windowEnd = pos + maxChunkSize;
// Prefer the last newline inside the window so escape sequences stay
// intact within their chunk; hard-cut at `windowEnd` otherwise.
const nl = data.lastIndexOf("\n", windowEnd - 1);
const cut = nl >= pos ? nl + 1 : windowEnd;
chunks.push(data.slice(pos, cut));
pos = cut;
}
return chunks;
}
/**
* Minimal terminal interface for TUI
*/
@@ -504,11 +556,20 @@ export class ProcessTerminal implements Terminal {
}
const match = sequence.match(kittyResponsePattern);
if (match && !this.#modifyOtherKeysActive) {
if (match) {
if (this.#modifyOtherKeysTimeout) {
clearTimeout(this.#modifyOtherKeysTimeout);
this.#modifyOtherKeysTimeout = undefined;
}
// A DA1 sentinel that beat the kitty reply may have already
// engaged the modifyOtherKeys fallback (terminals such as
// Superset/xterm-on-Electron answer DA1 before `\x1b[?u`).
// Kitty is strictly preferred — undo the fallback so the two
// modes do not stack. See #2042.
if (this.#modifyOtherKeysActive) {
this.#safeWrite("\x1b[>4;0m");
this.#modifyOtherKeysActive = false;
}
// Any reply to `\x1b[?u` means the terminal speaks the kitty keyboard
// protocol. The reported flag value is the *current* stack-top — fresh
// terminals report 0 — so support is implied by the reply itself, not by
@@ -978,7 +1039,23 @@ export class ProcessTerminal implements Terminal {
// files). They serve no purpose there and would surface as visible noise.
if (!process.stdout.isTTY) return;
try {
process.stdout.write(data);
// Windows ConPTY drops viewport tracking when a single write exceeds
// ~32-64 KB: the host UI's scroll position stays parked at wherever
// the write began, even though every byte landed in scrollback. Split
// large paints into newline-aligned chunks so each underlying
// `WriteFile` stays well below the threshold. The gate also covers
// WSL — `process.platform === "linux"` there, but stdout still
// crosses into ConPTY at the `wslhost` boundary, so the same per-
// WriteFile cap applies. Non-ConPTY PTYs keep the single-write fast
// path. See #2034.
const conptyHosted = process.platform === "win32" || isWindowsSubsystemForLinux();
if (conptyHosted && data.length > MAX_CONPTY_WRITE_CHUNK) {
for (const chunk of chunkForConPTY(data, MAX_CONPTY_WRITE_CHUNK)) {
process.stdout.write(chunk);
}
} else {
process.stdout.write(data);
}
} catch (err) {
// Any write failure means terminal is dead - no recovery possible
this.#dead = true;
+101 -2
View File
@@ -47,6 +47,12 @@ const SEGMENT_RESET = "\x1b[0m";
const LINE_TERMINATOR = "\x1b[0m\x1b]8;;\x07";
const ERASE_LINE = "\x1b[2K";
const ERASE_TO_END_OF_LINE = "\x1b[K";
// Bound the raw code-unit span handed to native width/truncation. A terminal
// row can only display `width` cells, so oversized component rows should not
// force proportional JS/native copies while deciding what the viewport shows.
const LINE_FIT_MIN_SOURCE_CODE_UNITS = 4096;
const LINE_FIT_MAX_SOURCE_CODE_UNITS = 65536;
const LINE_FIT_SOURCE_WIDTH_MULTIPLIER = 64;
// Hide the hardware cursor before each paint/move write. Ghostty-style bar
// cursors can otherwise leave visual afterimages while the TUI repaints the
// row under a visible cursor. Paint writes also disable terminal autowrap:
@@ -2049,7 +2055,9 @@ export class TUI extends Container {
newLines.length < this.#previousLines.length &&
naturalViewportTop !== prevViewportTop
) {
return { kind: "viewportRepaint" };
return this.#bottomAnchoredViewportUnchanged(newLines, height)
? { kind: "deferredMutation" }
: { kind: "viewportRepaint" };
}
// Direct-input shrink can also move the natural viewport upward even when
@@ -2437,6 +2445,17 @@ export class TUI extends Container {
return { kind: "liveRegionPinned", appendFrom, appendTo, renderViewportTop };
}
#bottomAnchoredViewportUnchanged(newLines: string[], height: number): boolean {
const previousViewportTop = Math.max(0, this.#previousLines.length - height);
const newViewportTop = Math.max(0, newLines.length - height);
for (let row = 0; row < height; row++) {
if ((newLines[newViewportTop + row] ?? "") !== (this.#previousLines[previousViewportTop + row] ?? "")) {
return false;
}
}
return true;
}
#planDeferredTailRepaint(newLines: string[], prevViewportTop: number, height: number): RenderIntent {
const row = prevViewportTop + height - 1;
if (row < 0 || row >= this.#previousLines.length || newLines.length !== this.#previousLines.length) {
@@ -2478,7 +2497,8 @@ export class TUI extends Container {
if (TERMINAL.isImageLine(raw)) {
return { raw, width, line: raw };
}
const normalized = normalizeTerminalOutput(raw);
const source = this.#lineFitSource(raw, width);
const normalized = normalizeTerminalOutput(source);
const asciiWidth = this.#ansiAsciiLineWidth(normalized, width);
if ((asciiWidth ?? visibleWidth(normalized)) <= width) {
return { raw, width, line: normalized };
@@ -2487,6 +2507,85 @@ export class TUI extends Container {
return { raw, width, line };
}
#lineFitSource(raw: string, width: number): string {
const safeWidth = Number.isFinite(width) ? Math.max(1, Math.trunc(width)) : 1;
const maxSourceLength = Math.min(
LINE_FIT_MAX_SOURCE_CODE_UNITS,
Math.max(LINE_FIT_MIN_SOURCE_CODE_UNITS, safeWidth * LINE_FIT_SOURCE_WIDTH_MULTIPLIER),
);
if (raw.length <= maxSourceLength) return raw;
const chunks: string[] = [];
let emitted = 0;
for (let i = 0; i < raw.length && emitted < maxSourceLength; ) {
if (raw.charCodeAt(i) === 0x1b) {
const end = this.#ansiSequenceEnd(raw, i);
if (end === -1) break;
const sequenceLength = end - i;
if (this.#ansiSequenceHasVisiblePayload(raw, i)) {
// OSC 66 text-sizing spans carry their visible cells inside the
// OSC payload. Always include the whole sequence — splitting it
// would corrupt the terminator — and let the next loop iteration
// terminate on the budget overflow.
chunks.push(raw.slice(i, end));
emitted += sequenceLength;
i = end;
continue;
}
if (emitted > 0 && sequenceLength <= maxSourceLength - emitted) {
chunks.push(raw.slice(i, end));
emitted += sequenceLength;
}
i = end;
continue;
}
const start = i;
const end = Math.min(raw.length, start + maxSourceLength - emitted);
while (i < end && raw.charCodeAt(i) !== 0x1b) i++;
if (i === start) break;
chunks.push(raw.slice(start, i));
emitted += i - start;
}
return chunks.join("") + SEGMENT_RESET;
}
#ansiSequenceEnd(line: string, start: number): number {
const next = line.charCodeAt(start + 1);
if (next === 0x5b) {
let i = start + 2;
while (i < line.length) {
const final = line.charCodeAt(i);
if (final >= 0x40 && final <= 0x7e) return i + 1;
i++;
}
return -1;
}
if (next === 0x5d) {
let i = start + 2;
while (i < line.length) {
const osc = line.charCodeAt(i);
if (osc === 0x07) return i + 1;
if (osc === 0x1b && line.charCodeAt(i + 1) === 0x5c) return i + 2;
i++;
}
return -1;
}
return start + 2 <= line.length ? start + 2 : -1;
}
#ansiSequenceHasVisiblePayload(line: string, start: number): boolean {
// OSC 66 (`\x1b]66;META;TEXT\x1b\\`) carries its visible cells inside the
// payload, mirroring the special case in {@link #ansiAsciiLineWidth}.
return (
line.charCodeAt(start + 1) === 0x5d &&
line.charCodeAt(start + 2) === 0x36 &&
line.charCodeAt(start + 3) === 0x36 &&
line.charCodeAt(start + 4) === 0x3b
);
}
#ansiAsciiLineWidth(line: string, maxWidth: number): number | undefined {
let col = 0;
for (let i = 0; i < line.length; ) {
+203
View File
@@ -0,0 +1,203 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import { chunkForConPTY, ProcessTerminal } from "@oh-my-pi/pi-tui/terminal";
// Regression test for https://github.com/can1357/oh-my-pi/issues/2034
//
// Windows ConPTY ties viewport tracking to per-`WriteFile` boundaries: when
// a single `process.stdout.write` exceeds ~32-64 KB, the pseudo-console
// stops following the cursor and the host UI's scroll position stays parked
// at wherever the write began. The data lands in scrollback — Alt+Tab forces
// the host to re-query the cursor and the viewport jumps to the bottom —
// but until then the user sees only the first screenful of a long session
// or resume payload.
//
// Fix: `ProcessTerminal#safeWrite` chunks oversized writes into ≤ 8 KiB
// pieces on `process.platform === "win32"`. Non-win32 PTYs do not share the
// bug and keep the single-write fast path.
const ESC = "\x1b";
function buildFullPaint(lines: number, lineLength: number): string {
// Mirrors the shape of `TUI#emitFullPaint`'s buffer: a clear-screen prefix,
// rows terminated with `\r\n` and a per-line SGR reset, and a cursor/end
// sequence trailer. The exact bytes do not matter for the chunker — only
// that escapes are present and the buffer crosses the ConPTY threshold.
let buf = `${ESC}[2J${ESC}[H${ESC}[3J`;
for (let i = 0; i < lines; i++) {
if (i > 0) buf += "\r\n";
const content = `${ESC}[38;5;${i % 256}mrow-${i.toString().padStart(4, "0")}: ${"x".repeat(lineLength)}${ESC}[0m`;
buf += content;
}
buf += `${ESC}[H${ESC}[?25h`;
return buf;
}
describe("issue #2034: chunk large terminal writes on Windows ConPTY", () => {
describe("chunkForConPTY()", () => {
it("returns the original buffer untouched when under the chunk size", () => {
const data = "small payload";
expect(chunkForConPTY(data, 1024)).toEqual([data]);
});
it("splits a large multi-line buffer into pieces no larger than the chunk size", () => {
const data = buildFullPaint(2000, 60);
const max = 8 * 1024;
expect(data.length).toBeGreaterThan(max);
const chunks = chunkForConPTY(data, max);
expect(chunks.length).toBeGreaterThan(1);
for (const chunk of chunks) {
expect(chunk.length).toBeLessThanOrEqual(max);
}
});
it("preserves the full payload across chunks (no data loss or reordering)", () => {
const data = buildFullPaint(500, 120);
const chunks = chunkForConPTY(data, 4 * 1024);
expect(chunks.join("")).toBe(data);
});
it("splits at newline boundaries so escape sequences are never sliced apart", () => {
// Every row is bracketed by SGR escapes. If the chunker cut inside a
// chunk's escape sequence, the trailing chunk would not start with
// either an escape or the post-newline state — instead it would
// start with a stray CSI byte (`[`, digits, `m`).
const data = buildFullPaint(400, 80);
const chunks = chunkForConPTY(data, 4 * 1024);
// Exclude the head chunk (starts with the clear-screen prefix).
for (const chunk of chunks.slice(1)) {
// Every subsequent chunk begins on a fresh line: either the new
// line's first byte is the SGR escape, the row's plaintext
// prefix, or — for the trailing tail — the cursor sequence.
const firstByte = chunk.charCodeAt(0);
const startsWithEsc = chunk.startsWith(ESC);
const startsWithRowText = chunk.startsWith("row-");
expect(startsWithEsc || startsWithRowText).toBe(true);
if (!startsWithEsc) {
// Plain-text starts cannot be control characters that would
// indicate a sliced escape (CSI `[`, digits, or `m`).
expect(firstByte).toBeGreaterThanOrEqual(0x20);
}
}
});
it("makes forward progress on a single line longer than the chunk size", () => {
// Pathological case: one very long line with no embedded `\n`. The
// chunker must not loop, and the joined chunks must equal the input.
const giantLine = "a".repeat(20_000);
const data = `${giantLine}\nshort\n`;
const chunks = chunkForConPTY(data, 4 * 1024);
expect(chunks.length).toBeGreaterThanOrEqual(2);
expect(chunks.join("")).toBe(data);
});
it("falls back to a raw split when the buffer contains no newlines", () => {
const data = "x".repeat(20_000);
const chunks = chunkForConPTY(data, 4 * 1024);
expect(chunks.join("")).toBe(data);
expect(chunks.length).toBeGreaterThan(1);
// Every chunk except possibly the tail is exactly the chunk size.
for (const chunk of chunks.slice(0, -1)) {
expect(chunk.length).toBe(4 * 1024);
}
});
});
describe("ProcessTerminal#write platform gate", () => {
const stdinIsTtyDescriptor = Object.getOwnPropertyDescriptor(process.stdin, "isTTY");
const stdoutIsTtyDescriptor = Object.getOwnPropertyDescriptor(process.stdout, "isTTY");
const platformDescriptor = Object.getOwnPropertyDescriptor(process, "platform");
const originalWslDistro = Bun.env.WSL_DISTRO_NAME;
const originalWslInterop = Bun.env.WSL_INTEROP;
function setEnv(key: string, value: string | undefined): void {
if (value === undefined) delete Bun.env[key];
else Bun.env[key] = value;
}
beforeEach(() => {
Object.defineProperty(process.stdin, "isTTY", { value: true, configurable: true });
Object.defineProperty(process.stdout, "isTTY", { value: true, configurable: true });
// Clear WSL markers by default; tests opt in.
setEnv("WSL_DISTRO_NAME", undefined);
setEnv("WSL_INTEROP", undefined);
});
afterEach(() => {
vi.restoreAllMocks();
if (platformDescriptor) Object.defineProperty(process, "platform", platformDescriptor);
if (stdinIsTtyDescriptor) Object.defineProperty(process.stdin, "isTTY", stdinIsTtyDescriptor);
else Reflect.deleteProperty(process.stdin, "isTTY");
if (stdoutIsTtyDescriptor) Object.defineProperty(process.stdout, "isTTY", stdoutIsTtyDescriptor);
else Reflect.deleteProperty(process.stdout, "isTTY");
setEnv("WSL_DISTRO_NAME", originalWslDistro);
setEnv("WSL_INTEROP", originalWslInterop);
});
function captureStdoutWrites(): string[] {
const writes: string[] = [];
vi.spyOn(process.stdout, "write").mockImplementation(chunk => {
writes.push(typeof chunk === "string" ? chunk : chunk.toString());
return true;
});
return writes;
}
it("splits >8 KiB writes into chunks on win32 so ConPTY can track the viewport", () => {
Object.defineProperty(process, "platform", { value: "win32", configurable: true });
const writes = captureStdoutWrites();
const terminal = new ProcessTerminal();
const payload = buildFullPaint(2000, 60);
terminal.write(payload);
const conptyChunks = writes.filter(w => w.length > 0);
expect(conptyChunks.length).toBeGreaterThan(1);
for (const chunk of conptyChunks) {
expect(chunk.length).toBeLessThanOrEqual(8 * 1024);
}
expect(conptyChunks.join("")).toBe(payload);
});
it("splits >8 KiB writes inside WSL because stdout still crosses ConPTY at wslhost", () => {
Object.defineProperty(process, "platform", { value: "linux", configurable: true });
setEnv("WSL_DISTRO_NAME", "Ubuntu");
setEnv("WSL_INTEROP", "/run/WSL/123_interop");
const writes = captureStdoutWrites();
const terminal = new ProcessTerminal();
const payload = buildFullPaint(2000, 60);
terminal.write(payload);
const conptyChunks = writes.filter(w => w.length > 0);
expect(conptyChunks.length).toBeGreaterThan(1);
for (const chunk of conptyChunks) {
expect(chunk.length).toBeLessThanOrEqual(8 * 1024);
}
expect(conptyChunks.join("")).toBe(payload);
});
it("keeps the single-write fast path on non-ConPTY platforms (clean linux, darwin)", () => {
Object.defineProperty(process, "platform", { value: "linux", configurable: true });
const writes = captureStdoutWrites();
const terminal = new ProcessTerminal();
const payload = buildFullPaint(2000, 60);
terminal.write(payload);
expect(writes).toEqual([payload]);
});
it("does not chunk small writes on win32", () => {
Object.defineProperty(process, "platform", { value: "win32", configurable: true });
const writes = captureStdoutWrites();
const terminal = new ProcessTerminal();
const payload = `${ESC}[H${ESC}[K`;
terminal.write(payload);
expect(writes).toEqual([payload]);
});
});
});
+120
View File
@@ -0,0 +1,120 @@
import { describe, expect, it } from "bun:test";
import { type Component, TUI } from "@oh-my-pi/pi-tui";
import type { Terminal, TerminalAppearance } from "@oh-my-pi/pi-tui/terminal";
class CaptureTerminal implements Terminal {
writes: string[] = [];
#columns: number;
#rows: number;
constructor(columns = 80, rows = 4) {
this.#columns = columns;
this.#rows = rows;
}
get columns(): number {
return this.#columns;
}
get rows(): number {
return this.#rows;
}
get kittyProtocolActive(): boolean {
return false;
}
get appearance(): TerminalAppearance | undefined {
return undefined;
}
start(): void {}
stop(): void {}
async drainInput(): Promise<void> {}
write(data: string): void {
this.writes.push(data);
}
moveBy(): void {}
hideCursor(): void {}
showCursor(): void {}
clearLine(): void {}
clearFromCursor(): void {}
clearScreen(): void {}
setTitle(): void {}
setProgress(): void {}
onAppearanceChange(): void {}
}
class RawLinesComponent implements Component {
#lines: string[];
constructor(lines: string[]) {
this.#lines = lines;
}
invalidate(): void {}
render(): string[] {
return this.#lines;
}
}
async function settle(): Promise<void> {
await Bun.sleep(0);
}
describe("issue #2045: renderer bounds oversized rows", () => {
it("preserves visible text after pathological zero-width ANSI prefixes", async () => {
const term = new CaptureTerminal(80, 4);
const tui = new TUI(term);
const line = `${"\x1b[31m".repeat(20_000)}payload`;
tui.addChild(new RawLinesComponent([line]));
try {
tui.start();
await settle();
} finally {
tui.stop();
}
const rendered = term.writes.join("");
expect(rendered).toContain("payload");
expect(rendered.length).toBeLessThan(12_000);
});
it("preserves visible text after oversized OSC hyperlink prefixes", async () => {
const term = new CaptureTerminal(80, 4);
const tui = new TUI(term);
const line = `\x1b]8;;https://example.com/${"a".repeat(70_000)}\x07link-label\x1b]8;;\x07`;
tui.addChild(new RawLinesComponent([line]));
try {
tui.start();
await settle();
} finally {
tui.stop();
}
const rendered = term.writes.join("");
expect(rendered).toContain("link-label");
expect(rendered.length).toBeLessThan(12_000);
});
it("preserves OSC 66 text-sizing payloads at the start of long rows", async () => {
const term = new CaptureTerminal(80, 4);
const tui = new TUI(term);
const visibleText = "H".repeat(70);
const line = `\x1b]66;s=1;${visibleText}\x1b\\${"\x1b[31m".repeat(20_000)}`;
tui.addChild(new RawLinesComponent([line]));
try {
tui.start();
await settle();
} finally {
tui.stop();
}
const rendered = term.writes.join("");
expect(rendered).toContain(visibleText);
});
});
@@ -0,0 +1,69 @@
import { afterEach, describe, expect, it } from "bun:test";
import {
createProcessTerminalRenderHarness,
type ProcessTerminalRenderHarness,
} from "./process-terminal-render-harness";
// Progressive-enhancement probe ordering contract. omp sends `CSI ? u \\ CSI c`
// at startup: the kitty reply (`CSI ? <flags> u`) authoritatively says the
// terminal speaks the kitty keyboard protocol; the DA1 reply (`CSI ? ... c`)
// is only a sentinel that guarantees a reply even from terminals that ignore
// `CSI ? u`. Some terminals (Superset / xterm-on-Electron) answer DA1 first;
// the kitty reply must still be honored regardless of ordering.
describe("ProcessTerminal kitty keyboard progressive-enhancement ordering", () => {
let harness: ProcessTerminalRenderHarness | undefined;
afterEach(() => {
harness?.dispose();
harness = undefined;
});
it("enables kitty when the kitty reply arrives before the DA1 sentinel", async () => {
harness = createProcessTerminalRenderHarness(100, 30);
await harness.settle();
expect(harness.writes.join("")).toContain("\x1b[?u\x1b[c");
harness.writes.length = 0;
await harness.feed("\x1b[?0u", "\x1b[?1;2c");
const out = harness.writes.join("");
expect(harness.terminal.kittyProtocolActive).toBe(true);
expect(out).toContain("\x1b[>1u");
expect(out).not.toContain("\x1b[>4;2m");
});
it("enables kitty when the DA1 sentinel arrives before the kitty reply (#2042)", async () => {
harness = createProcessTerminalRenderHarness(100, 30);
await harness.settle();
harness.writes.length = 0;
// Superset/Electron-xterm answers DA1 before `CSI ? u`. The kitty reply
// must override the premature modifyOtherKeys fallback.
await harness.feed("\x1b[?1;2c", "\x1b[?0u");
const out = harness.writes.join("");
expect(harness.terminal.kittyProtocolActive).toBe(true);
expect(out).toContain("\x1b[>1u");
const enableIdx = out.indexOf("\x1b[>4;2m");
const disableIdx = out.indexOf("\x1b[>4;0m");
const kittyIdx = out.indexOf("\x1b[>1u");
expect(enableIdx).toBeGreaterThanOrEqual(0);
expect(disableIdx).toBeGreaterThan(enableIdx);
expect(kittyIdx).toBeGreaterThan(enableIdx);
});
it("keeps the modifyOtherKeys fallback when only DA1 ever replies", async () => {
harness = createProcessTerminalRenderHarness(100, 30);
await harness.settle();
harness.writes.length = 0;
// Terminals that ignore `CSI ? u` answer DA1 only — modifyOtherKeys is
// the right answer there.
await harness.feed("\x1b[?1;2c");
const out = harness.writes.join("");
expect(harness.terminal.kittyProtocolActive).toBe(false);
expect(out).toContain("\x1b[>4;2m");
expect(out).not.toContain("\x1b[>1u");
});
});
@@ -1469,6 +1469,40 @@ describe("TUI terminal-state regressions", () => {
});
});
it("tmux: offscreen shrink preserving the visible tail emits no repaint bytes", async () => {
await withEnvPatch({ TMUX: "1", STY: undefined, ZELLIJ: undefined }, async () => {
const term = new UnknownViewportTerminal(40, 4, 10_000);
const tui = new TUI(term);
const component = new MutableLinesComponent([
"old-0",
"remove-me",
"old-2",
"old-3",
"tail-0",
"tail-1",
"tail-2",
"tail-3",
]);
tui.addChild(component);
try {
tui.start();
await settle(term);
expect(visible(term)).toEqual(["tail-0", "tail-1", "tail-2", "tail-3"]);
const writes = captureWrites(term);
component.setLines(["old-0", "old-2", "old-3", "tail-0", "tail-1", "tail-2", "tail-3"]);
tui.requestRender();
await settle(term);
expect(visible(term)).toEqual(["tail-0", "tail-1", "tail-2", "tail-3"]);
expect(writes).toEqual([]);
} finally {
tui.stop();
}
});
});
// Root cause family: the dirty/replay machinery assumes native scrollback
// can be cleared and rebuilt, which is never true inside a multiplexer —
// tmux owns pane history, reflows it on resize itself, and a "replay" can