Merge remote-tracking branch 'upstream/main' into feat/error-notify

This commit is contained in:
Mathews-Tom
2026-07-02 06:00:36 +05:30
62 changed files with 1483 additions and 163 deletions
+2 -2
View File
@@ -43,7 +43,7 @@ Patch language inside `input`:
- Every body row is `+TEXT`; `+` alone adds a blank line.
- `DEL` never has body rows.
- There is no repeat row kind. To keep a line, leave it out of every range; split edits into multiple hunks when needed.
- `-` rows are invalid. Literal text beginning with `-` or `+` must be written as `+-text` / `++text`.
- `-` rows are invalid. Literal Markdown bullets or text beginning with `-` / `+` must be written as `+- item` / `++ item`.
Anchors come from `read`/`grep` output. `read` emits a `[PATH#TAG]` header from the session snapshot store and lines as `LINE:TEXT`; copy the header into the edit section and copy only the line number into hunk headers.
@@ -165,7 +165,7 @@ DEL 20
- Stray payload line:
- `line N: payload line has no preceding hunk header. Use \`SWAP N.=M:\`, \`DEL N.=M\`, or \`INS.PRE|POST|HEAD|TAIL:\` above the body. Got "...".`
- Minus row:
- ``line N: `-` rows are not valid; the range already names the lines being changed. For a literal `-` line, write `+-…`.``
- ``line N: `-` rows are not valid; the range already names the lines being changed. For Markdown bullets or other literal `-` lines, prefix the literal row with `+`: `+- item`.``
- Empty body-bearing hunk:
- `line N: \`INS\` needs at least one \`+TEXT\` body row.`
- `line N: \`SWAP.BLK N:\` needs at least one \`+TEXT\` body row. To delete a block, use \`DEL.BLK N\`.`
+1 -2
View File
@@ -13,7 +13,7 @@ import {
resolveWireModelId,
} from "@oh-my-pi/pi-catalog/model-thinking";
import { CATALOG_PROVIDERS, type ProviderCatalogEntry } from "@oh-my-pi/pi-catalog/provider-models";
import { $env, $pickenv, getConfigRootDir, isEnoent, logger } from "@oh-my-pi/pi-utils";
import { $env, $pickenv, getConfigRootDir, isEnoent, logger, withExtraCaFetch } from "@oh-my-pi/pi-utils";
import { getCustomApi } from "./api-registry";
import { AUTH_RETRY_STEPS, isApiKeyResolver, resolveRetryKey } from "./auth-retry";
import * as AIError from "./error";
@@ -75,7 +75,6 @@ import { wrapLeakedThinkingStream } from "./utils/leaked-thinking-stream";
import { wrapFetchForProxy } from "./utils/proxy";
import { withRequestDebugFetch } from "./utils/request-debug";
import { withGeminiThinkingLoopGuard } from "./utils/thinking-loop";
import { withExtraCaFetch } from "./utils/tls-fetch";
function isGoogleVertexAuthenticatedModel(model: Model<Api>): boolean {
return (
+1
View File
@@ -5,6 +5,7 @@
### Fixed
- Fixed the Xiaomi provider's default model to use the supported mimo-v2.5 model.
- Fixed model discovery (`/models` probes) failing behind private-CA gateways: the discovery fetch fallback now honors `NODE_EXTRA_CA_CERTS`, matching provider chat requests.
- Fixed CoreWeave Serverless Inference project-header detection to ensure blank OpenAI-Project overrides do not block the COREWEAVE_PROJECT fallback.
- Fixed LiteLLM MiniMax M3 discovery to remove reseller-only (3x usage) display suffixes.
- Fixed users with a warm LiteLLM model cache keeping stale reseller display-name suffixes for up to 24 hours by bumping the cache namespace to rich-v2.
@@ -1,3 +1,4 @@
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { type } from "arktype";
import type { ModelSpec } from "../types";
import { toPositiveNumber } from "../utils";
@@ -161,7 +162,7 @@ export interface FetchAntigravityDiscoveryModelsOptions {
export async function fetchAntigravityDiscoveryModels(
options: FetchAntigravityDiscoveryModelsOptions,
): Promise<ModelSpec<"google-gemini-cli">[] | null> {
const fetcher = options.fetcher ?? fetch;
const fetcher = options.fetcher ?? wrapFetchForExtraCa(fetch);
const endpoints = options.endpoint
? [trimTrailingSlashes(options.endpoint)]
: DEFAULT_ANTIGRAVITY_DISCOVERY_ENDPOINTS.map(trimTrailingSlashes);
+3 -2
View File
@@ -1,3 +1,4 @@
import { type FetchImpl, wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { type } from "arktype";
import type { ModelSpec } from "../types";
import { isRecord } from "../utils";
@@ -81,7 +82,7 @@ export interface CodexModelDiscoveryResult {
* Returns `{ models: [] }` when a route succeeds but yields no usable models.
*/
export async function fetchCodexModels(options: CodexModelDiscoveryOptions): Promise<CodexModelDiscoveryResult | null> {
const fetchFn = options.fetchFn ?? fetch;
const fetchFn = options.fetchFn ?? wrapFetchForExtraCa(fetch);
const baseUrl = normalizeBaseUrl(options.baseUrl);
const paths = normalizePaths(options.paths);
const headers = buildCodexHeaders(options);
@@ -168,7 +169,7 @@ function buildCodexHeaders(options: CodexModelDiscoveryOptions): Headers {
async function resolveCodexClientVersion(
clientVersion: string | undefined,
fetchFn: typeof fetch,
fetchFn: FetchImpl,
signal: AbortSignal | undefined,
): Promise<string> {
const normalizedClientVersion = normalizeClientVersion(clientVersion);
+2 -1
View File
@@ -1,5 +1,6 @@
import { gunzipSync } from "node:zlib";
import { create, fromBinary, toBinary } from "@bufbuild/protobuf";
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import type { FetchImpl, ModelSpec } from "../types";
import {
GetCliModelConfigsRequestSchema,
@@ -76,7 +77,7 @@ export async function fetchDevinModels(
accept: "*/*",
};
const fetchImpl = options.fetch ?? fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(fetch);
const response = await fetchImpl(requestUrl, { method: "POST", headers, body, signal });
if (!response.ok) {
return null;
+2 -1
View File
@@ -1,3 +1,4 @@
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { type } from "arktype";
import { getBundledModels } from "../models";
import { toModelSpec } from "../provider-models/bundled-references";
@@ -77,7 +78,7 @@ export async function fetchGeminiModels(
return null;
}
const fetchImpl = options.fetch ?? fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(fetch);
const baseUrl = normalizeBaseUrl(options.baseUrl);
const pageSize = normalizePositiveInt(options.pageSize, DEFAULT_PAGE_SIZE);
const maxPages = normalizePositiveInt(options.maxPages, DEFAULT_MAX_PAGES);
@@ -1,5 +1,6 @@
import * as fs from "node:fs/promises";
import * as path from "node:path";
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { z } from "zod/v4";
import type { FetchImpl, ModelSpec } from "../types";
import { isRecord } from "../utils";
@@ -340,7 +341,7 @@ async function fetchNamespaceOverrideCandidate(
if (!restNamespaceId) {
return null;
}
const fetchImpl = config.fetch ?? fetch;
const fetchImpl = config.fetch ?? wrapFetchForExtraCa(fetch);
let response: Response;
try {
response = await fetchImpl(`${baseUrl}${GROUPS_PATH}/${encodeURIComponent(restNamespaceId)}`, {
@@ -406,7 +407,7 @@ async function fetchProjectRootNamespaceViaRest(
baseUrl: string,
projectIdOrPath: string,
): Promise<GitLabDuoWorkflowRestProject | null> {
const fetchImpl = config.fetch ?? fetch;
const fetchImpl = config.fetch ?? wrapFetchForExtraCa(fetch);
let response: Response;
try {
response = await fetchImpl(`${baseUrl}${PROJECTS_PATH}/${encodeURIComponent(projectIdOrPath)}`, {
@@ -449,7 +450,7 @@ async function fetchTopLevelGroupNamespaceCandidates(
config: GitLabDuoWorkflowDiscoveryConfig,
baseUrl: string,
): Promise<GitLabDuoWorkflowCandidate[]> {
const fetchImpl = config.fetch ?? fetch;
const fetchImpl = config.fetch ?? wrapFetchForExtraCa(fetch);
const candidates: (GitLabDuoWorkflowCandidate & { preferred: boolean })[] = [];
// GitLab paginates `/groups`; a token can belong to more than one page of top-level
// groups, and a usable Duo namespace may live on a later page. Follow the keyset/
@@ -518,7 +519,7 @@ async function postGraphQL(
query: string,
variables: Record<string, string>,
): Promise<unknown | null> {
const fetchImpl = config.fetch ?? fetch;
const fetchImpl = config.fetch ?? wrapFetchForExtraCa(fetch);
let response: Response;
try {
response = await fetchImpl(`${baseUrl}${GRAPHQL_PATH}`, {
@@ -1,3 +1,4 @@
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { type } from "arktype";
import type { Api, FetchImpl, ModelSpec, Provider } from "../types";
@@ -137,7 +138,7 @@ export async function fetchOpenAICompatibleModels<TApi extends Api>(
requestHeaders.Authorization = `Bearer ${options.apiKey}`;
}
const fetchImpl = options.fetch ?? globalThis.fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(globalThis.fetch);
const fetchPayload = async (signal?: AbortSignal): Promise<unknown | null> => {
let response: Response;
try {
@@ -1,4 +1,4 @@
import { fetchWithRetry } from "@oh-my-pi/pi-utils";
import { fetchWithRetry, wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { Effort } from "../effort";
import { isGlm52ReasoningEffortModelId } from "../identity/family";
import type { ModelManagerOptions } from "../model-manager";
@@ -75,7 +75,7 @@ async function fetchShowMetadata(
baseUrl: string,
apiKey: string,
model: string,
fetchImpl: FetchImpl = fetch,
fetchImpl: FetchImpl = wrapFetchForExtraCa(fetch),
): Promise<OllamaShowResponse | undefined> {
const response = await fetchImpl(`${baseUrl}/api/show`, {
method: "POST",
@@ -107,7 +107,7 @@ export function ollamaCloudModelManagerOptions(
const response = await fetchWithRetry(`${baseUrl}/api/tags`, {
method: "GET",
headers: createCloudHeaders(apiKey),
fetch: config?.fetch,
fetch: config?.fetch ?? wrapFetchForExtraCa(fetch),
defaultDelayMs: OLLAMA_RETRY_DELAYS_MS,
});
if (!response.ok) {
@@ -1,3 +1,4 @@
import { wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import {
fetchOpenAICompatibleModels,
type OpenAICompatibleModelMapperContext,
@@ -743,7 +744,7 @@ async function fetchUmansModelsInfo(options: {
if (options.apiKey) {
requestHeaders["x-api-key"] = options.apiKey;
}
const fetchImpl = options.fetch ?? fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(fetch);
let payload: unknown;
try {
const response = await fetchImpl(`${discoveryBaseUrl}${UMANS_MODELS_INFO_PATH}`, {
@@ -1541,7 +1542,7 @@ async function fetchFireworksServerlessModels(options: {
}): Promise<ModelSpec<"openai-completions">[] | null> {
const listUrl = toFireworksControlPlaneModelsUrl(options.baseUrl, FIREWORKS_CONTROL_PLANE_ACCOUNT);
if (!listUrl) return null;
const fetchImpl = options.fetch ?? fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(fetch);
const collected = new Map<string, ModelSpec<"openai-completions">>();
let pageToken = "";
for (let page = 0; page < FIREWORKS_CONTROL_PLANE_MAX_PAGES; page++) {
@@ -3079,7 +3080,7 @@ async function fetchLiteLLMRichEndpoint<TApi extends Api>(
runtimeBaseUrl: string,
signal?: AbortSignal,
): Promise<ModelSpec<TApi>[] | null> {
const fetchImpl = options.fetch ?? globalThis.fetch;
const fetchImpl = options.fetch ?? wrapFetchForExtraCa(globalThis.fetch);
const requestHeaders: Record<string, string> = {
Accept: "application/json",
...options.headers,
+3 -9
View File
@@ -1,5 +1,8 @@
import type { Effort } from "./effort";
// Re-exported from @oh-my-pi/pi-utils so the whole workspace shares one
// `fetch`-compatible signature (tls-fetch's wrappers produce/accept it).
export type { FetchImpl } from "@oh-my-pi/pi-utils";
export type { KnownProvider } from "./provider-models/descriptors";
export type KnownApi =
@@ -89,15 +92,6 @@ export type Provider = string;
/** Token budgets for each thinking level (token-based providers only) */
export type ThinkingBudgets = { [key in Effort]?: number };
/**
* `fetch`-compatible function. Accepts any callable matching the standard
* fetch signature; `preconnect` is optional because non-Bun runtimes (browsers,
* test mocks) won't expose it.
*/
export type FetchImpl = ((input: string | URL | Request, init?: RequestInit) => Promise<Response>) & {
preconnect?: typeof globalThis.fetch.preconnect;
};
export interface Usage {
/** Non-cached input tokens (matches the bucket the provider bills as new input). */
input: number;
@@ -1,4 +1,8 @@
import { describe, expect, it } from "bun:test";
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { __resetExtraCaCache } from "@oh-my-pi/pi-utils";
import { fetchOpenAICompatibleModels } from "../src/discovery/openai-compatible";
describe("discovery null limits", () => {
@@ -30,3 +34,75 @@ describe("discovery null limits", () => {
expect(models![0].maxTokens).toBeNull();
});
});
describe("discovery extra-CA fallback fetch", () => {
const SAMPLE_PEM =
"-----BEGIN CERTIFICATE-----\nMIIBkTCCATegAwIBAgIUF/sample/extra/ca/for/discovery/123=\n-----END CERTIFICATE-----\n";
const MODELS_BODY = JSON.stringify({ data: [{ id: "some-model" }] });
let tmpDir: string;
let originalEnv: string | undefined;
let originalFetch: typeof globalThis.fetch;
beforeEach(async () => {
__resetExtraCaCache();
originalEnv = Bun.env.NODE_EXTRA_CA_CERTS;
originalFetch = globalThis.fetch;
tmpDir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-discovery-ca-"));
});
afterEach(async () => {
__resetExtraCaCache();
globalThis.fetch = originalFetch;
if (originalEnv === undefined) delete Bun.env.NODE_EXTRA_CA_CERTS;
else Bun.env.NODE_EXTRA_CA_CERTS = originalEnv;
await fs.rm(tmpDir, { recursive: true, force: true });
});
it("carries the NODE_EXTRA_CA_CERTS bundle on the fallback fetch", async () => {
const caPath = path.join(tmpDir, "corp.pem");
await Bun.write(caPath, SAMPLE_PEM);
Bun.env.NODE_EXTRA_CA_CERTS = caPath;
const inits: (RequestInit & { tls?: { ca?: string | string[] } })[] = [];
globalThis.fetch = (async (_input: string | URL | Request, init?: RequestInit) => {
inits.push((init ?? {}) as (typeof inits)[number]);
return new Response(MODELS_BODY, { status: 200, headers: { "Content-Type": "application/json" } });
}) as typeof globalThis.fetch;
// No `fetch` injected — the wrapped globalThis.fetch fallback must be used.
const models = await fetchOpenAICompatibleModels({
provider: "custom",
api: "openai-completions",
baseUrl: "https://gateway.corp.example/v1",
});
expect(models).not.toBeNull();
expect(models!.length).toBe(1);
expect(inits).toHaveLength(1);
expect(inits[0].tls?.ca).toContain(SAMPLE_PEM);
});
it("leaves a caller-injected fetch untouched", async () => {
const caPath = path.join(tmpDir, "corp.pem");
await Bun.write(caPath, SAMPLE_PEM);
Bun.env.NODE_EXTRA_CA_CERTS = caPath;
const inits: (RequestInit & { tls?: { ca?: string | string[] } })[] = [];
const injectedFetch = async (_input: string | URL | Request, init?: RequestInit) => {
inits.push((init ?? {}) as (typeof inits)[number]);
return new Response(MODELS_BODY, { status: 200, headers: { "Content-Type": "application/json" } });
};
const models = await fetchOpenAICompatibleModels({
provider: "custom",
api: "openai-completions",
baseUrl: "https://gateway.corp.example/v1",
fetch: injectedFetch,
});
expect(models).not.toBeNull();
expect(inits).toHaveLength(1);
expect(inits[0].tls).toBeUndefined();
});
});
+10
View File
@@ -15,6 +15,11 @@
### Fixed
- Fixed task.maxConcurrency being breachable when a queued spawn was cancelled: the spawn path could release a semaphore permit it never acquired, letting a later task start while the cap was saturated.
- Fixed session exit diagnostics recording signal and crash exits (SIGTERM, SIGHUP, uncaught exceptions) as a normal "dispose": the postmortem teardown now threads the real reason into session disposal.
- Fixed the subagent yield-label guard ignoring JTD discriminator (oneOf) output schemas, which let stale incremental labels pass into successful results when final validation was skipped after retries.
- Fixed grep/ast_grep search scopes rejecting `www.` and collapsed-scheme (`https:/host`) URL spellings that the read tool accepts; unresolvable URL-shaped scopes now fail with an explicit external-URL error instead of "Path not found".
- Fixed model discovery ignoring `NODE_EXTRA_CA_CERTS`: the model registry's default fetch now applies the extra-CA wrapper, so `/models` probes work behind private-CA gateways like provider chat requests.
- Fixed ctrl+p role-model cycling getting stuck on one transition and skipping every other role: a session-branch traversal regression returned entries leaf-to-root, so the cycle (and session model restore) read the oldest recorded model change instead of the newest.
- Fixed ctrl+p cycling from a stale slot after the model was switched through another surface (alt+m, /model, retry fallback): the recorded role is now trusted only while its resolved model is still the active model, falling back to matching by model.
- Fixed the apply_patch tool to prevent silently overwriting pre-existing files during creation or renaming, rejecting upfront with an error instead.
@@ -25,6 +30,7 @@
- Fixed RPC mode abort_bash being blocked by running bash commands by dispatching bash in the background.
- Fixed task.maxConcurrency and task.maxRecursionDepth limits being bypassed by sub-spawn paths, ensuring limits are dynamically resized and respected.
- Fixed the edit tool inflating session files by pruning extremely large file snapshots from tool-result details.
- Fixed edit-tool Markdown list guidance so hashline parser errors and the model-facing prompt teach `+- item` escaping instead of steering agents toward full-file `write` fallbacks. ([#4179](https://github.com/can1357/oh-my-pi/issues/4179))
- Fixed workstation OS detection rendering "Kernel: unknown" on macOS 15+.
- Fixed /copy code and /copy cmd commands being treated as normal prompts instead of copying the requested blocks.
- Fixed interactive bash status line not updating after directory changes (cd).
@@ -64,6 +70,10 @@
- Fixed transcript rebuilds (theme change, /shake, focus replay) showing stale streamed write/edit/eval content by sharing the partial-JSON decode between the live streaming path and every rebuild path.
- Fixed an explicitly configured compaction.reserveTokens equal to the built-in default being silently replaced by the proportional small-window fallback; the setting now defaults to unset and explicit values are always honored.
- Fixed user-configured LiteLLM discovery providers keeping stale reseller display-name suffixes for up to 24 hours after upgrade by invalidating the warm model cache.
### Fixed
- Fixed `mergeTaskBranches` and `applyNestedPatches` leaving stage 1/2/3 unmerged entries in `.git/index` when the post-merge stash pop conflicted with the cherry-picked HEAD. The corrupted index survived indefinitely and every subsequent overlay-isolated task inherited it through the lower layer, causing `captureRepoDeltaPatch` to emit `diff --cc` output that `git apply` rejects with "No valid patches in input". The stash restore now runs behind a 3-way preflight check (`git apply --3way --check`) and a `reset --hard HEAD` cleanup fallback; the stash entry is preserved for manual recovery on conflict, and the merged commits still land on HEAD. ([#4175](https://github.com/can1357/oh-my-pi/issues/4175))
- Fixed macOS `Command+V` image pastes in Ghostty by binding the Kitty `super+v` key event to the image-paste action alongside `Ctrl+V`. ([#4178](https://github.com/can1357/oh-my-pi/issues/4178))
## [16.2.13] - 2026-07-01
+10
View File
@@ -275,6 +275,16 @@ export class CollabGuestLink {
}
return;
}
if (frame.t === "error" && !this.#welcomed && !this.#left) {
// Pre-welcome errors are the host's targeted reply to our
// hello (e.g. protocol mismatch): no welcome will follow.
// Fail the join with the host's message instead of hanging
// until the welcome timeout.
this.#clearWelcomeTimer();
if (joined) this.#ctx.showError(`Collab host: ${frame.message}`);
else firstWelcome.reject(new Error(frame.message));
return;
}
if (!this.#welcomed || this.#left) return;
this.#applyFrame(frame);
})
@@ -65,7 +65,9 @@ declare module "@oh-my-pi/pi-tui" {
* Resolve default image-paste shortcuts for the current terminal platform.
*/
export function getDefaultPasteImageKeys(platform: NodeJS.Platform = process.platform): KeyId[] {
return platform === "win32" ? ["ctrl+v", "alt+v"] : ["ctrl+v"];
if (platform === "win32") return ["ctrl+v", "alt+v"];
if (platform === "darwin") return ["ctrl+v", "super+v"];
return ["ctrl+v"];
}
/**
@@ -52,7 +52,7 @@ import type { ApiKeyResolver, FetchImpl } from "@oh-my-pi/pi-ai";
import { registerOAuthProvider, unregisterOAuthProviders } from "@oh-my-pi/pi-ai/oauth";
import type { OAuthCredentials, OAuthLoginCallbacks } from "@oh-my-pi/pi-ai/oauth/types";
import { getBundledModelReferenceIndex, resolveModelReference } from "@oh-my-pi/pi-catalog/identity";
import { isBunTestRuntime, isRecord, logger } from "@oh-my-pi/pi-utils";
import { isBunTestRuntime, isRecord, logger, wrapFetchForExtraCa } from "@oh-my-pi/pi-utils";
import { parseModelString, resolveProviderModelReference } from "../config/model-resolver";
import type { AuthStorage, OAuthCredential } from "../session/auth-storage";
import { type ApiKeyResolverModel, type ApiKeyResolverOptions, createApiKeyResolver } from "./api-key-resolver";
@@ -763,7 +763,7 @@ export class ModelRegistry {
options?.fetch ??
(isBunTestRuntime()
? () => Promise.reject(new Error("network disabled in model-registry runtime test"))
: fetch);
: wrapFetchForExtraCa(fetch));
this.#modelsConfigFile = ModelsConfigFile.relocate(modelsPath);
this.#cacheDbPath = modelsPath ? path.join(path.dirname(modelsPath), "models.db") : undefined;
// Set up fallback resolver for custom provider API keys
@@ -180,7 +180,8 @@ function resolvePreviewEdits(args: {
}): readonly Edit[] {
const { section, absolutePath, normalized, snapshots, expected, liveMatches, edits } = args;
if (!hasBlockEdit(edits)) return edits;
const baseText = expected === undefined || liveMatches ? normalized : snapshots.byHash(absolutePath, expected)?.text;
const baseText =
expected === undefined || liveMatches ? normalized : snapshots.byHashExact(absolutePath, expected)?.text;
if (baseText === undefined) {
throw createMismatchError(section, absolutePath, normalized, snapshots, expected ?? "");
}
@@ -199,7 +200,17 @@ function applyPreviewEdits(args: {
if (!options.skipHashValidation && expected === undefined) {
throw new Error(missingSnapshotTagMessage(section.path));
}
const liveMatches = expected !== undefined && computeFileHash(normalized) === expected;
// A 16-bit tag can collide across two different file states, so hash
// equality alone does not prove the live text IS the snapshot the model's
// anchors were minted against (mirrors Patcher's apply-time guard). When
// the store retains text for the tag, require it to be unambiguous and
// byte-identical to live; otherwise fall through to recovery/reject below
// exactly as if the hash had not matched.
const liveMatches =
expected !== undefined &&
computeFileHash(normalized) === expected &&
(snapshots.byHash(absolutePath, expected) === null ||
snapshots.byHashExact(absolutePath, expected)?.text === normalized);
const edits = parsePreviewEdits(section, options.streaming);
const resolved = resolvePreviewEdits({ section, absolutePath, normalized, snapshots, expected, liveMatches, edits });
if (options.skipHashValidation || expected === undefined || liveMatches) return applyEdits(normalized, resolved);
@@ -1,5 +1,7 @@
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "bun:test";
import { setKittyProtocolActive } from "@oh-my-pi/pi-tui/keys";
import { $ } from "bun";
import { getDefaultPasteImageKeys } from "../../config/keybindings";
import { getEditorTheme, initTheme } from "../theme/theme";
import {
CustomEditor,
@@ -124,6 +126,24 @@ describe("CustomEditor bracketed path paste", () => {
expect(imagePathCalls).toBe(0);
});
});
describe("CustomEditor configured paste image keys", () => {
it("routes Ghostty Cmd+V kitty key events through the macOS image-paste default", () => {
const { editor } = makeEditor();
const onPasteImage = vi.fn();
editor.onPasteImage = onPasteImage;
editor.setActionKeys("app.clipboard.pasteImage", getDefaultPasteImageKeys("darwin"));
setKittyProtocolActive(true);
try {
editor.handleInput("\x1b[118;9u");
} finally {
setKittyProtocolActive(false);
}
expect(onPasteImage).toHaveBeenCalledTimes(1);
expect(editor.getText()).toBe("");
});
});
describe("extractImagePathFromText (issue #3506)", () => {
it("returns the path when the text is a single image file path", () => {
@@ -794,9 +794,15 @@ export class InteractiveMode implements InteractiveModeContext {
getDraftText: () => this.editor.getText(),
beginDispose: () => this.session.beginDispose(),
saveDraft: text => this.sessionManager.saveDraft(text),
disposeSession: () => this.session.dispose({ mnemopiConsolidateTimeoutMs: SHUTDOWN_CONSOLIDATE_BUDGET_MS }),
disposeSession: reason =>
this.session.dispose({ mnemopiConsolidateTimeoutMs: SHUTDOWN_CONSOLIDATE_BUDGET_MS, reason }),
});
this.#cleanupUnsubscribe = postmortem.register("session-teardown", () => this.#signalTeardown!());
// Forward the postmortem reason (SIGTERM/SIGHUP/uncaughtException/…) so the
// persisted `session_exit` diagnostic carries the real trigger. Postmortem
// runs callbacks in REVERSE registration order — this callback (registered
// after the AgentSession constructor's `agent-session:<id>` recorder) runs
// FIRST and its dispose() would otherwise persist the generic "dispose".
this.#cleanupUnsubscribe = postmortem.register("session-teardown", reason => this.#signalTeardown!(reason));
// Wire the report_tool_issue consent gate to the Yes/No dialog popup.
// The handler is process-global — subagent tools (which can't reach
@@ -1,4 +1,5 @@
import { describe, expect, it } from "bun:test";
import { postmortem } from "@oh-my-pi/pi-utils";
import { createSessionTeardown } from "./session-teardown";
/**
@@ -156,4 +157,63 @@ describe("createSessionTeardown", () => {
expect(captured).toBe("before");
});
it("forwards the postmortem reason into disposeSession so signal exits record the real trigger", async () => {
const received: Array<postmortem.Reason | undefined> = [];
const teardown = createSessionTeardown({
getDraftText: () => "",
beginDispose: () => {},
saveDraft: async () => {},
disposeSession: async reason => {
received.push(reason);
},
});
await teardown(postmortem.Reason.SIGTERM);
expect(received).toEqual([postmortem.Reason.SIGTERM]);
});
it("keypress path passes no reason — dispose falls back to a normal exit record", async () => {
const received: Array<postmortem.Reason | undefined> = [];
const teardown = createSessionTeardown({
getDraftText: () => "",
beginDispose: () => {},
saveDraft: async () => {},
disposeSession: async reason => {
received.push(reason);
},
});
await teardown();
expect(received).toEqual([undefined]);
});
it("first call's reason wins: a later caller with a different reason awaits the same promise", async () => {
const received: Array<postmortem.Reason | undefined> = [];
const release = Promise.withResolvers<void>();
const teardown = createSessionTeardown({
getDraftText: () => "",
beginDispose: () => {},
saveDraft: async () => {},
disposeSession: async reason => {
received.push(reason);
await release.promise;
},
});
// SIGTERM lands first; postmortem.quit(0)'s MANUAL pass arrives while the
// teardown is still draining — it must not restart the run or mutate the
// recorded reason.
const first = teardown(postmortem.Reason.SIGTERM);
const second = teardown(postmortem.Reason.MANUAL);
release.resolve();
await Promise.all([first, second]);
expect(received).toEqual([postmortem.Reason.SIGTERM]);
});
});
@@ -9,7 +9,7 @@
* Extracted (rather than inlined into `InteractiveMode`) so the callback body
* is directly unit-testable without instantiating the full TUI stack.
*/
import { logger } from "@oh-my-pi/pi-utils";
import { logger, type postmortem } from "@oh-my-pi/pi-utils";
/** Dependencies the teardown captures at construction time. */
export interface SessionTeardownDeps {
@@ -26,12 +26,22 @@ export interface SessionTeardownDeps {
* previously-persisted draft sidecar is cleared on a clean exit.
*/
saveDraft: (text: string) => Promise<void>;
/** Dispose the session — emits `session_shutdown`, drains async jobs, closes the manager. */
disposeSession: () => Promise<void>;
/**
* Dispose the session — emits `session_shutdown`, drains async jobs, closes
* the manager. Receives the postmortem reason that triggered the teardown
* (undefined on the keypress/`/exit` path) so `AgentSession.dispose()` can
* persist the real exit reason instead of the generic `"dispose"`.
*/
disposeSession: (reason?: postmortem.Reason) => Promise<void>;
}
/** Idempotent teardown: concurrent/repeat invocations share one settled promise. */
export type SessionTeardown = () => Promise<void>;
/**
* Idempotent teardown: concurrent/repeat invocations share one settled
* promise. The optional `reason` is the postmortem reason that triggered the
* teardown (`sigterm`, `sighup`, `uncaught_exception`, …); only the FIRST
* call's reason is used — later callers await the same settled promise.
*/
export type SessionTeardown = (reason?: postmortem.Reason) => Promise<void>;
/**
* Build a promise-memoized teardown function. The first call snapshots the
@@ -42,13 +52,20 @@ export type SessionTeardown = () => Promise<void>;
* double-emit `session_shutdown`, double-dispose the session's async-job
* manager, or race each other.
*
* The postmortem callback forwards its `Reason` so the persisted
* `session_exit` diagnostic carries the real trigger (`sigterm`, `sighup`,
* `uncaught_exception`, …) instead of the generic `"dispose"` that plain
* programmatic disposal records. First call wins: a signal arriving after a
* keypress-initiated teardown awaits the in-flight promise and its reason is
* dropped — by then the exit entry is already being written as a normal exit.
*
* `saveDraft` failures are logged but never abort the disposal chain — a
* draft-write error must not leak background bash/task jobs or skip the
* extension `session_shutdown` event.
*/
export function createSessionTeardown(deps: SessionTeardownDeps): SessionTeardown {
let pending: Promise<void> | undefined;
const run = async (): Promise<void> => {
const run = async (reason?: postmortem.Reason): Promise<void> => {
const draftText = deps.getDraftText();
deps.beginDispose();
try {
@@ -56,10 +73,10 @@ export function createSessionTeardown(deps: SessionTeardownDeps): SessionTeardow
} catch (err) {
logger.warn("Failed to save session draft during teardown", { error: String(err) });
}
await deps.disposeSession();
await deps.disposeSession(reason);
};
return () => {
if (!pending) pending = run();
return (reason?: postmortem.Reason) => {
if (!pending) pending = run(reason);
return pending;
};
}
@@ -458,6 +458,13 @@ export const SHUTDOWN_CONSOLIDATE_BUDGET_MS = 1_500;
export interface AgentSessionDisposeOptions {
mnemopiConsolidateTimeoutMs?: number;
/**
* Postmortem reason that triggered this dispose (signal/fatal teardown
* paths). When set, the persisted `session_exit` diagnostic records it
* instead of the generic `"dispose"` used for normal programmatic disposal
* (`/quit`, test teardown, subagent completion).
*/
reason?: postmortem.Reason;
}
type CompactionCheckResult = Readonly<{
@@ -5343,7 +5350,7 @@ export class AgentSession {
async #doDispose(options: AgentSessionDisposeOptions = {}): Promise<void> {
this.beginDispose();
this.#recordSessionExit("dispose");
this.#recordSessionExit(options.reason ?? "dispose");
this.#cancelExitRecorder?.();
this.#cancelExitRecorder = undefined;
try {
+13 -2
View File
@@ -814,6 +814,17 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
const startedAt = Date.now();
const semaphore = this.#getSpawnSemaphore();
let semaphoreHeld = false;
// Every release funnels through here: the flag flips before the
// release so no path — acquire-time abort, executor failure, or a
// future refactor that reorders the branches — can return a permit
// twice. Releasing a permit this job never acquired would steal one
// from a running job and let a later spawn start past
// task.maxConcurrency.
const releasePermit = () => {
if (!semaphoreHeld) return;
semaphoreHeld = false;
this.#releaseSpawnSemaphore();
};
try {
await semaphore.acquire(runSignal);
semaphoreHeld = true;
@@ -825,7 +836,7 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
}
const acquiredAt = Date.now();
if (!semaphoreHeld || runSignal.aborted) {
if (semaphoreHeld) this.#releaseSpawnSemaphore();
releasePermit();
progress.status = "aborted";
onSettled?.(true);
throw new Error("Aborted before execution");
@@ -897,7 +908,7 @@ export class TaskTool implements AgentTool<TaskToolSchemaInstance, TaskToolDetai
const hint = AgentRegistry.global().get(agentId) ? buildFollowUpHint(false) : "";
throw new TaskJobError(`${message}${hint}`);
} finally {
this.#releaseSpawnSemaphore();
releasePermit();
}
},
{
+12 -13
View File
@@ -305,16 +305,13 @@ export async function applyNestedPatches(
}
} finally {
if (stashed) {
try {
await git.stash.pop(nestedDir, { index: true });
} catch (popErr) {
const message = popErr instanceof Error ? popErr.message : String(popErr);
const restored = await git.stash.tryPop(nestedDir, { index: true });
if (!restored) {
logger.warn("Pre-existing nested-repo dirty state could not be auto-restored", {
nestedDir,
error: message,
});
warnings.push(
`Pre-existing dirty state in nested repo \`${relativePath}\` could not be auto-restored after the agent commit; stash entry preserved (${message}).`,
`Pre-existing dirty state in nested repo \`${relativePath}\` could not be auto-restored after the agent commit; stash entry preserved.`,
);
}
}
@@ -755,13 +752,15 @@ export async function mergeTaskBranches(
}
} finally {
if (didStash) {
try {
await git.stash.pop(repoRoot, { index: true });
} catch {
// Stash-pop conflicts mean the replayed changes clash with the user's
// uncommitted edits. The cherry-picked commits are already on HEAD, so
// the merged branches DID land — report them as merged and surface the
// stash conflict separately instead of claiming they are unmerged.
const restored = await git.stash.tryPop(repoRoot, { index: true });
if (!restored) {
// Stash pop would leave stage 1/2/3 unmerged entries in `.git/index`
// that overlay-isolated subsequent tasks inherit through the lower
// layer, corrupting every downstream `captureRepoDeltaPatch`. `tryPop`
// short-circuits the pop when the WIP would conflict with the
// cherry-picked HEAD (and reset-cleans up if a rarer conflict slips
// past). The merged branches DID land — surface a stash-restore
// warning without claiming the merge failed.
logger.warn("Failed to restore stashed changes after task merge; stash entry preserved");
const stashConflict =
"stash pop: cherry-picked changes conflict with uncommitted edits. The merged commits are on HEAD; run `git stash pop` and resolve manually.";
+1 -5
View File
@@ -27,7 +27,7 @@ import { finalizeOutput, loadPage, looksLikeHtml, MAX_BYTES, MAX_OUTPUT_CHARS }
import { convertWithMarkit, fetchBinary } from "../web/scrapers/utils";
import { applyListLimit } from "./list-limit";
import { formatStyledArtifactReference, type OutputMeta } from "./output-meta";
import { type LineRange, parseLineRanges } from "./path-utils";
import { isReadableUrlPath, type LineRange, parseLineRanges } from "./path-utils";
import { formatBytes, formatExpandHint, getDomain, replaceTabs } from "./render-utils";
import { listTables, looksLikeSqlite, renderTableList } from "./sqlite-reader";
import { ToolAbortError, ToolError } from "./tool-errors";
@@ -145,10 +145,6 @@ function normalizeUrl(url: string): string {
return url;
}
export function isReadableUrlPath(value: string): boolean {
return /^https?:\/\/?/i.test(value) || /^www\./i.test(value);
}
// URL line selectors mirror the file form: `:50`, `:50-100`, `:50+150`, `:5-10,20-30`, `:raw`,
// or `:raw:N-M` / `:N-M:raw` to combine raw mode with a range. If a URL would otherwise look
// like `host:port`, add a trailing slash before the selector (e.g. `https://example.com/:80`
@@ -139,14 +139,38 @@ interface SectionLabelMetadata {
isKnown(label: string): boolean;
}
/**
* Derive incremental-label metadata from top-level schema closure.
*
* The unknown-label gate (`rejectUnknownSections`) engages when the schema constrains top-level
* property names anywhere: a closed conjunct (root or recursive `allOf` child with
* `additionalProperties: false`) or a `oneOf`/`anyOf` union whose EVERY variant is closed. A label
* is known iff every closed conjunct accepts it AND, per closed union, at least one variant
* accepts it (union semantics are disjunctive — the assembled output only has to match one
* variant). Unions containing any open variant never gate: the open variant accepts arbitrary
* labels, so rejection would be a false positive.
*/
function buildSectionLabelMetadata(jsonSchema: Record<string, unknown>): SectionLabelMetadata {
const closedSchemas = collectClosedTopLevelSchemas(jsonSchema);
const labels = [...new Set(closedSchemas.flatMap(schema => declaredPropertyLabels(schema)))];
const closedConjuncts = collectClosedTopLevelSchemas(jsonSchema);
const closedUnions = collectClosedTopLevelUnions(jsonSchema);
const closed = closedConjuncts.length > 0 || closedUnions.length > 0;
const acceptedByAll = (conjuncts: readonly Record<string, unknown>[], label: string): boolean =>
conjuncts.every(schema => schemaAcceptsSectionLabel(schema, label));
const labels = [
...new Set([
...closedConjuncts.flatMap(schema => declaredPropertyLabels(schema)),
...closedUnions.flatMap(variants =>
variants.flatMap(conjuncts => conjuncts.flatMap(schema => declaredPropertyLabels(schema))),
),
]),
];
return {
labels,
rejectUnknownSections: closedSchemas.length > 0,
rejectUnknownSections: closed,
isKnown: label =>
closedSchemas.length === 0 || closedSchemas.every(schema => schemaAcceptsSectionLabel(schema, label)),
!closed ||
(acceptedByAll(closedConjuncts, label) &&
closedUnions.every(variants => variants.some(conjuncts => acceptedByAll(conjuncts, label)))),
};
}
@@ -164,6 +188,43 @@ function collectClosedTopLevelSchemas(jsonSchema: Record<string, unknown>): Reco
return schemas;
}
/** One fully-closed `oneOf`/`anyOf` union: per variant, that variant's closed conjunct schemas. */
type ClosedUnionVariants = Record<string, unknown>[][];
/**
* Collect top-level `oneOf`/`anyOf` unions in which EVERY variant is closed — i.e. each variant
* (or one of its `allOf` conjuncts, resolved via `collectClosedTopLevelSchemas`) carries
* `additionalProperties: false`. JTD discriminator output schemas compile to exactly this shape:
* a root `oneOf` of closed object variants. Unions with any open (or non-object) variant are
* skipped entirely so the unknown-label gate cannot fire false rejections. Unions nested under
* `allOf` conjuncts gate identically (intersection semantics).
*/
function collectClosedTopLevelUnions(jsonSchema: Record<string, unknown>): ClosedUnionVariants[] {
const unions: ClosedUnionVariants[] = [];
for (const key of ["oneOf", "anyOf"] as const) {
const rawVariants = jsonSchema[key];
if (!Array.isArray(rawVariants) || rawVariants.length === 0) continue;
const variants: ClosedUnionVariants = [];
let allClosed = true;
for (const raw of rawVariants) {
const conjuncts = isRecord(raw) ? collectClosedTopLevelSchemas(raw) : [];
if (conjuncts.length === 0) {
allClosed = false;
break;
}
variants.push(conjuncts);
}
if (allClosed) unions.push(variants);
}
const allOf = jsonSchema.allOf;
if (Array.isArray(allOf)) {
for (const raw of allOf) {
if (isRecord(raw)) unions.push(...collectClosedTopLevelUnions(raw));
}
}
return unions;
}
function declaredPropertyLabels(jsonSchema: Record<string, unknown>): string[] {
const properties = jsonSchema.properties;
if (properties === null || typeof properties !== "object" || Array.isArray(properties)) return [];
+37 -5
View File
@@ -2,7 +2,7 @@ import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import * as url from "node:url";
import { isEnoent, stripWindowsExtendedLengthPathPrefix } from "@oh-my-pi/pi-utils";
import { isEnoent, isEnotdir, stripWindowsExtendedLengthPathPrefix } from "@oh-my-pi/pi-utils";
import type { Skill } from "../extensibility/skills";
import { InternalUrlRouter, type LocalProtocolOptions } from "../internal-urls";
import { ToolError } from "./tool-errors";
@@ -435,6 +435,16 @@ export function isSshUrl(path: string): boolean {
return /^ssh:\/\//i.test(path.trim());
}
/**
* True when the read tool's URL parser (`parseReadUrlTarget` in fetch.ts) would
* recognize this path as a readable external URL: a strict `http(s)://`, a
* collapsed `http(s):/host` (Node path normalization folds `//` → `/`), or a
* scheme-less `www.` spelling. Keep in sync with `parseReadUrlTarget`.
*/
export function isReadableUrlPath(value: string): boolean {
return /^https?:\/\/?/i.test(value) || /^www\./i.test(value);
}
/**
* Resolve a path relative to the given cwd.
* Handles ~ expansion and absolute paths.
@@ -1090,14 +1100,30 @@ export async function resolveToolSearchScope(opts: ToolScopeOptions): Promise<To
if (rawPaths.some(rawPath => rawPath.length === 0)) {
throw new ToolError("`paths` must contain non-empty paths or globs");
}
const externalUrlRe = /^(?:https?|ftp|file|ws|wss):\/\//i;
// Strict external-URL schemes. `file://` is intentionally absent: it has
// local-path semantics (expandPath strips it downstream), so it flows through
// the ordinary filesystem pipeline instead of the external-URL resolver.
const strictExternalUrlRe = /^(?:https?|ftp|ws|wss):\/\//i;
const internalRouter = InternalUrlRouter.instance();
const resolvedPathInputs: string[] = [];
const immutableSourcePaths = new Set<string>();
for (const rawPath of rawPaths) {
const externalUrl = externalUrlRe.test(rawPath);
if (externalUrl && opts.resolveExternalUrl) {
const resolved = await opts.resolveExternalUrl(rawPath);
let externalUrl = strictExternalUrlRe.test(rawPath);
if (!externalUrl && isReadableUrlPath(rawPath) && !hasGlobPathChars(rawPath)) {
// Fuzzy spelling the read parser accepts (`www.host/…`, collapsed
// `https:/host/…`). An existing local path wins over URL
// interpretation so a directory literally named `www.foo` stays
// searchable; only a definitive ENOENT/ENOTDIR flips to URL handling
// (any other stat error means the path exists — let the local
// pipeline surface it).
try {
await fs.promises.stat(resolveToCwd(rawPath, cwd));
} catch (err) {
externalUrl = isEnoent(err) || isEnotdir(err);
}
}
if (externalUrl) {
const resolved = opts.resolveExternalUrl ? await opts.resolveExternalUrl(rawPath) : undefined;
if (resolved) {
resolvedPathInputs.push(resolved.sourcePath);
if (opts.trackImmutableSources && resolved.immutable) {
@@ -1105,6 +1131,12 @@ export async function resolveToolSearchScope(opts: ToolScopeOptions): Promise<To
}
continue;
}
// Resolver missing or declined (e.g. ftp/ws/wss): fail explicitly
// instead of letting the local-path fallthrough surface a confusing
// "Path not found" for a URL-shaped input.
throw new ToolError(
`Cannot ${internalUrlAction} external URL: ${rawPath}. Use \`read\` to fetch web content, then search the returned text.`,
);
}
if (!internalRouter.canHandle(rawPath)) {
resolvedPathInputs.push(rawPath);
+1 -1
View File
@@ -77,7 +77,6 @@ import {
} from "./conflict-detect";
import {
executeReadUrl,
isReadableUrlPath,
loadReadUrlCacheEntry,
parseReadUrlTarget,
type ReadUrlToolDetails,
@@ -95,6 +94,7 @@ import {
import {
expandPath,
formatPathRelativeToCwd,
isReadableUrlPath,
type LineRange,
parseLineRanges,
pathTargetsSsh,
+3 -2
View File
@@ -390,8 +390,9 @@ export class YieldTool implements AgentTool<TSchema, YieldDetails> {
}
/**
* Return incremental yield labels that are not declared as top-level properties of a closed
* caller schema. Open schemas (no `additionalProperties: false`) accept any label.
* Return incremental yield labels the closed caller schema does not accept. Closure covers the
* root, `allOf` conjuncts, and `oneOf`/`anyOf` unions whose every variant is closed (e.g. JTD
* discriminators). Open schemas accept any label.
*/
#unknownIncrementalLabels(labels: string[]): string[] {
if (!this.#rejectUnknownSections) return [];
+74 -2
View File
@@ -1750,6 +1750,69 @@ export const stash = {
if (options?.index) args.push("--index");
await runEffect(cwd, args);
},
/**
* Return the working-tree patch that `stash@{0}` would apply, in a form
* that `git apply --check` can consume. Empty string when no stash entry
* exists or the stash contains no diffable working-tree changes.
*/
async showPatch(cwd: string): Promise<string> {
return (await tryText(cwd, ["stash", "show", "-p", "--binary", "stash@{0}"], { readOnly: true })) ?? "";
},
/** Return untracked paths stored in the top stash entry. */
async untrackedFiles(cwd: string): Promise<string[]> {
const output = await tryText(cwd, ["ls-tree", "-r", "-z", "--name-only", "stash@{0}^3"], { readOnly: true });
return output?.split("\0").filter(Boolean) ?? [];
},
/**
* Attempt to restore the top stash entry. On success returns `true` and
* git drops the stash entry. On conflict returns `false`, leaves the stash
* entry preserved for manual resolution, and guarantees the failed restore
* leaves no unmerged index entries or partially-restored untracked files.
*
* The historical raw `pop` catches the failure in a `finally` block and
* only logs — it leaves `.git/index` with stage 1/2/3 unmerged entries
* that survive indefinitely, corrupting every subsequent overlay-isolated
* task that reads through this repo's `.git/`. See issue #4175.
*/
async tryPop(cwd: string, options?: { index?: boolean }): Promise<boolean> {
// Preflight: `git stash pop` internally does a 3-way merge, so a plain
// `git apply --check` is too strict — it rejects hunks whose context
// drifted from HEAD even when 3-way merge would resolve them cleanly.
// Match pop's semantics with `--3way --check`, which succeeds iff the
// patch either applies directly or merges without conflict against
// the patch's `index abc..def` base blobs.
const workingPatch = await stash.showPatch(cwd);
if (workingPatch.trim() && !(await patch.canApplyText(cwd, workingPatch, { threeWay: true }))) {
return false;
}
const restoredUntracked = await stash.untrackedFiles(cwd);
try {
await stash.pop(cwd, options);
return true;
} catch {
// Preflight can still miss mode-only or delete/modify conflicts. If
// the pop left unmerged entries, wipe them: HEAD holds the merged
// state so `reset --hard HEAD` restores a clean index and working
// tree without losing the cherry-picked commits. A failed pop can
// still restore unrelated untracked files before exiting while
// preserving the stash entry, so clean only the untracked paths
// recorded in that stash. The user's WIP remains recoverable via
// `git stash pop`.
try {
await reset(cwd, { hard: true });
} catch {
/* best-effort cleanup — do not mask the primary conflict */
}
if (restoredUntracked.length > 0) {
try {
await clean(cwd, { includeIgnored: true, literalPathspecs: true, paths: restoredUntracked });
} catch {
/* best-effort cleanup — do not mask the primary conflict */
}
}
return false;
}
},
};
// ════════════════════════════════════════════════════════════════════════════
@@ -1819,9 +1882,18 @@ export async function reset(
export async function clean(
cwd: string,
options: { ignoredOnly?: boolean; paths?: readonly string[]; signal?: AbortSignal } = {},
options: {
ignoredOnly?: boolean;
includeIgnored?: boolean;
literalPathspecs?: boolean;
paths?: readonly string[];
signal?: AbortSignal;
} = {},
): Promise<void> {
const args = ["clean", options.ignoredOnly ? "-fdX" : "-fd"];
const args = [options.literalPathspecs ? "--literal-pathspecs" : undefined, "clean"].filter(
(arg): arg is string => arg !== undefined,
);
args.push(options.ignoredOnly ? "-fdX" : options.includeIgnored ? "-fdx" : "-fd");
if (options.paths?.length) args.push("--", ...options.paths);
await runEffect(cwd, args, { signal: options.signal });
}
@@ -14,10 +14,13 @@
import { afterEach, beforeEach, describe, expect, it, spyOn } from "bun:test";
import { generateRoomKey, importRoomKey } from "@oh-my-pi/pi-coding-agent/collab/crypto";
import { CollabGuestLink } from "@oh-my-pi/pi-coding-agent/collab/guest";
import { CollabHost } from "@oh-my-pi/pi-coding-agent/collab/host";
import {
COLLAB_PROTO,
type CollabFrame,
type CollabSessionState,
formatCollabLink,
parseCollabLink,
rewriteEnvelopePeer,
unpackEnvelope,
} from "@oh-my-pi/pi-coding-agent/collab/protocol";
@@ -523,3 +526,167 @@ describe("collab TUI guest ui-request handling (#4049)", () => {
expect(h.uiResponses).toEqual([]);
});
});
// ── Proto handshake (#4049: ui-request frames require COLLAB_PROTO >= 3) ───
//
// The ui-request/ui-response grammar shipped without a proto bump, so v2
// guests joined fine and silently dropped host asks. These tests pin the
// enforcement: a real CollabHost must reject stale-proto hellos with an
// observable error frame (never a welcome), current-proto guests must still
// complete a full ui-request round trip, and a rejected CollabGuestLink
// join must fail fast with the host's reason instead of hanging until the
// welcome timeout.
/** Minimal InteractiveModeContext double: only the members CollabHost touches. */
function makeHostContext(): InteractiveModeContext {
return {
settings: { get: () => "" },
sessionManager: {
getSessionId: () => "sess-proto",
getCwd: () => "/tmp",
snapshotForReplication: () => ({
header: { type: "session", id: "sess-proto", timestamp: new Date().toISOString(), cwd: "/tmp" },
entries: [],
}),
onEntryAppended: undefined,
},
session: {
isStreaming: false,
queuedMessageCount: 0,
sessionName: "proto test",
model: undefined,
thinkingLevel: undefined,
subscribe: () => () => {},
emitNotice: () => {},
promptCustomMessage: () => Promise.resolve(),
abort: () => Promise.resolve(),
},
eventBus: undefined,
statusLine: {
setCollabStatus: () => {},
invalidate: () => {},
getCachedContextBreakdown: () => ({ usedTokens: 0, contextWindow: 0 }),
},
ui: { requestRender: () => {} },
showStatus: () => {},
collabHost: undefined,
} as unknown as InteractiveModeContext;
}
/** Raw wire-speaking guest with a configurable hello proto. */
async function joinRawGuest(
link: string,
proto: number,
): Promise<{ socket: CollabSocket; nextFrame(): Promise<CollabFrame> }> {
const parsed = parseCollabLink(link);
if ("error" in parsed) throw new Error(parsed.error);
const writeToken = parsed.writeToken ? Buffer.from(parsed.writeToken).toString("base64url") : undefined;
const key = await importRoomKey(parsed.key);
const socket = new CollabSocket({ wsUrl: parsed.wsUrl, role: "guest", key });
const queue: CollabFrame[] = [];
const waiters: ((frame: CollabFrame) => void)[] = [];
// Directed welcome/error/ui frames only: the host's debounced broadcasts
// (state/agents/entry/event/bus) and the snapshot-chunk train interleave
// nondeterministically with the frames these tests assert on.
const filtered: Record<string, true> = {
state: true,
agents: true,
entry: true,
event: true,
bus: true,
"snapshot-chunk": true,
};
socket.onFrame = frame => {
if (filtered[frame.t]) return;
const waiter = waiters.shift();
if (waiter) waiter(frame);
else queue.push(frame);
};
socket.onOpen = () => socket.send({ t: "hello", proto, name: `guest-v${proto}`, writeToken });
socket.connect();
const nextFrame = (): Promise<CollabFrame> => {
const queued = queue.shift();
if (queued) return Promise.resolve(queued);
const { promise, resolve } = Promise.withResolvers<CollabFrame>();
waiters.push(resolve);
return promise;
};
return { socket, nextFrame };
}
describe("collab proto handshake (#4049)", () => {
it("host rejects a stale-proto hello with a protocol-mismatch error and never welcomes or admits the guest", async () => {
const host = new CollabHost(makeHostContext());
await host.start("ws://localhost:8787");
const guest = await joinRawGuest(host.link, COLLAB_PROTO - 1);
try {
const reply = await guest.nextFrame();
if (reply.t !== "error") throw new Error(`expected error, got ${reply.t}`);
expect(reply.message).toContain("protocol mismatch");
expect(reply.message).toContain(`host speaks v${COLLAB_PROTO}`);
expect(reply.message).toContain(`guest sent v${COLLAB_PROTO - 1}`);
// The rejected guest was never admitted: no participant entry, and a
// host ask finds no writable peer to route to.
expect(host.participants.filter(p => p.role !== "host")).toEqual([]);
expect(host.requestGuestUi({ kind: "select", title: "anyone?", options: ["Yes"] })).toBeNull();
} finally {
guest.socket.close();
await host.stop("test done");
}
});
it("welcomes a current-proto guest at v3 and round-trips a ui-request", async () => {
const host = new CollabHost(makeHostContext());
await host.start("ws://localhost:8787");
const guest = await joinRawGuest(host.link, COLLAB_PROTO);
try {
const welcome = await guest.nextFrame();
if (welcome.t !== "welcome") throw new Error(`expected welcome, got ${welcome.t}`);
expect(welcome.proto).toBe(COLLAB_PROTO);
expect(welcome.proto).toBe(3);
const pending = host.requestGuestUi({ kind: "select", title: "Continue?", options: ["Yes"] });
if (!pending) throw new Error("expected writable guest UI request");
const request = await guest.nextFrame();
if (request.t !== "ui-request") throw new Error(`expected ui-request, got ${request.t}`);
guest.socket.send({ t: "ui-response", reqId: request.request.reqId, value: "Yes" });
expect(await pending).toBe("Yes");
} finally {
guest.socket.close();
await host.stop("test done");
}
});
it("CollabGuestLink.join fails fast with the host's rejection message instead of hanging for the welcome", async () => {
// Scripted host that rejects every hello the way CollabHost does for a
// proto mismatch. The real guest must surface that message from join().
const roomId = "proto-reject-room";
const roomKey = generateRoomKey();
const cryptoKey = await importRoomKey(roomKey);
const link = formatCollabLink("ws://localhost:8788", roomId, roomKey);
const hostSocket = new CollabSocket({ wsUrl: `ws://localhost:8788/r/${roomId}`, role: "host", key: cryptoKey });
const hostOpen = Promise.withResolvers<void>();
hostSocket.onOpen = () => hostOpen.resolve();
hostSocket.onFrame = frame => {
if (frame.t === "hello") {
hostSocket.send({
t: "error",
message: `protocol mismatch: host speaks v${COLLAB_PROTO + 1}, guest sent v${frame.proto}`,
});
}
};
hostSocket.connect();
await hostOpen.promise;
const ctx = {
settings: { get: () => "" },
sessionManager: { getSessionFile: () => null },
} as unknown as InteractiveModeContext;
const guest = new CollabGuestLink(ctx);
try {
await expect(guest.join(link)).rejects.toThrow(/protocol mismatch/);
} finally {
hostSocket.close();
}
});
});
@@ -298,6 +298,63 @@ describe("computeHashlineDiff", () => {
expect(result.error).toContain('internal scheme "local://"');
}
});
// A 16-bit snapshot tag can collide across two different file states. The
// preview must mirror Patcher's apply-time guard: hash equality alone never
// proves the live text IS the snapshot the anchors were minted against.
test("rejects the no-drift path when the live text is a colliding ambiguous tag", async () => {
// Both texts hash to `1D84` (pinned in hashline's collision tests).
const SNAPSHOT_TEXT = "line one 263\nline two 4471\n";
const LIVE_TEXT = "line one 410\nline two 6970\n";
const sourcePath = path.join(tempDir, "source.txt");
await Bun.write(sourcePath, LIVE_TEXT);
const snapshotStore = new InMemorySnapshotStore();
// Anchors were minted against SNAPSHOT_TEXT; the live file is the
// colliding LIVE_TEXT, also retained (e.g. read after an external
// write). The tag is ambiguous — previewing the SWAP against live
// would show the model's payload landing on unrelated content.
const tag = snapshotStore.record(sourcePath, SNAPSHOT_TEXT);
snapshotStore.record(sourcePath, LIVE_TEXT);
const result = await computeHashlineDiff(
{ input: `${formatHashlineHeader(sourcePath, tag)}\nSWAP 2.=2:\n+edited from snapshot` },
tempDir,
snapshotStore,
);
expect("error" in result).toBe(true);
if ("error" in result) {
expect(result.error).toContain("file changed between read and edit");
}
});
test("rejects the no-drift path when the live text collides with the single retained snapshot", async () => {
const SNAPSHOT_TEXT = "line one 263\nline two 4471\n";
const LIVE_TEXT = "line one 410\nline two 6970\n";
const sourcePath = path.join(tempDir, "source.txt");
await Bun.write(sourcePath, LIVE_TEXT);
const snapshotStore = new InMemorySnapshotStore();
// Only SNAPSHOT_TEXT is retained; live drifted to a colliding text the
// store never saw. computeFileHash(live) === tag, but the retained
// text differs — recovery (3-way merge from SNAPSHOT_TEXT) must run
// instead of anchoring directly onto the collider. Here the merge
// cannot apply (every line differs), so the preview surfaces the
// drift error rather than a bogus diff.
const tag = snapshotStore.record(sourcePath, SNAPSHOT_TEXT);
const result = await computeHashlineDiff(
{ input: `${formatHashlineHeader(sourcePath, tag)}\nSWAP 2.=2:\n+edited from snapshot` },
tempDir,
snapshotStore,
);
expect("error" in result).toBe(true);
if ("error" in result) {
expect(result.error).toContain("file changed between read and edit");
}
});
});
describe("computeEditDiff", () => {
@@ -616,7 +616,7 @@ describe("InputController escape behavior", () => {
const controller = new InputController(ctx);
controller.setupKeyHandlers();
expect(() => editor.onEscape?.()).not.toThrow();
editor.onEscape?.();
expect(viewSession.abortCompaction).toHaveBeenCalledTimes(1);
expect(viewSession.abortHandoff).toHaveBeenCalledTimes(1);
@@ -38,8 +38,8 @@ describe("getDefaultPasteImageKeys", () => {
expect(getDefaultPasteImageKeys("win32")).toEqual(["ctrl+v", "alt+v"]);
});
it("uses Ctrl+V as the image-paste shortcut on non-Windows platforms", () => {
it("adds the macOS Command key event to Ctrl+V for image paste", () => {
expect(getDefaultPasteImageKeys("linux")).toEqual(["ctrl+v"]);
expect(getDefaultPasteImageKeys("darwin")).toEqual(["ctrl+v"]);
expect(getDefaultPasteImageKeys("darwin")).toEqual(["ctrl+v", "super+v"]);
});
});
@@ -6,6 +6,7 @@ import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { createSessionTeardown } from "@oh-my-pi/pi-coding-agent/modes/session-teardown";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import {
@@ -17,7 +18,7 @@ import {
} from "@oh-my-pi/pi-coding-agent/session/exit-diagnostics";
import { convertToLlm } from "@oh-my-pi/pi-coding-agent/session/messages";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { TempDir } from "@oh-my-pi/pi-utils";
import { postmortem, TempDir } from "@oh-my-pi/pi-utils";
const pendingAssistant: AssistantMessage = {
role: "assistant",
@@ -131,6 +132,70 @@ describe("session exit diagnostics", () => {
});
});
it("signal teardown persists the postmortem reason, not the generic dispose", async () => {
tempDir = TempDir.createSync("@pi-session-exit-signal-");
authStorage = await AuthStorage.create(path.join(tempDir.path(), "auth.db"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
const modelRegistry = new ModelRegistry(authStorage);
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!model) throw new Error("Expected built-in anthropic model to exist");
const sessionManager = SessionManager.inMemory(tempDir.path());
const agent = new Agent({
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
messages: [],
},
convertToLlm,
});
session = new AgentSession({
agent,
sessionManager,
settings: Settings.isolated({ "compaction.enabled": false }),
modelRegistry,
});
const activeSession = session;
// The assistant message persists through an async queue; the tool start
// marker is appended synchronously and is what makes the session durable
// enough for #recordSessionExit to write the exit entry (same setup as
// the plain-dispose test above).
agent.emitExternalEvent({ type: "message_end", message: pendingAssistant });
await Promise.resolve();
agent.emitExternalEvent({
type: "tool_execution_start",
toolCallId: "toolu_repro",
toolName: "bash",
args: { command: "bun run check:ts" },
});
await Promise.resolve();
// Mirror InteractiveMode.init(): the postmortem "session-teardown"
// callback runs FIRST on SIGTERM/SIGHUP/uncaughtException (reverse
// registration order) and calls dispose(). Without reason threading,
// #doDispose would persist the generic "dispose"/"normal" and cancel the
// reason-specific agent-session recorder — losing the real trigger.
const teardown = createSessionTeardown({
getDraftText: () => "",
beginDispose: () => activeSession.beginDispose(),
saveDraft: async () => {},
disposeSession: reason => activeSession.dispose({ reason }),
});
await teardown(postmortem.Reason.SIGTERM);
session = undefined;
const exitEntry = sessionManager
.getEntries()
.find(entry => entry.type === "custom" && entry.customType === SESSION_EXIT_CUSTOM_TYPE);
if (exitEntry?.type !== "custom") throw new Error("Expected session exit marker");
expect(exitEntry.data).toMatchObject({
reason: "sigterm",
kind: "signal",
});
});
it("does not materialize an empty session just to write an exit marker", async () => {
tempDir = TempDir.createSync("@pi-empty-session-exit-");
authStorage = await AuthStorage.create(path.join(tempDir.path(), "auth.db"));
@@ -231,6 +231,73 @@ describe("task spawn routing", () => {
expect(secondJob.status).toBe("cancelled");
});
it("keeps the concurrency cap intact when a queued spawn is cancelled (no permit leak)", async () => {
vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({
agents: [taskAgent],
projectAgentsDir: null,
});
const started: string[] = [];
const gates = new Map<string, Deferred>();
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
const id = options.id ?? "?";
started.push(id);
const gate = deferred();
gates.set(id, gate);
await gate.promise;
return makeResult(id);
});
const manager = createManager();
const tool = await TaskTool.create(createSession({ manager, settings: { "task.maxConcurrency": 1 } }));
// A holds the only permit, gated inside the executor.
const first = await tool.execute("tc-1", { agent: "task", id: "First", assignment: "Work A." } as TaskParams);
const firstJob = manager.getJob(first.details!.async!.jobId)!;
await pollUntil(() => started.length === 1);
// B parks at the semaphore, then is cancelled while queued. Its
// teardown must NOT release a permit it never acquired.
const second = await tool.execute("tc-2", { agent: "task", id: "Second", assignment: "Work B." } as TaskParams);
const secondJob = manager.getJob(second.details!.async!.jobId)!;
expect(secondJob.queued).toBe(true);
expect(manager.cancel(secondJob.id)).toBe(true);
await secondJob.promise;
expect(secondJob.status).toBe("cancelled");
// C must stay parked while A still holds the cap. A phantom release
// from B's cancellation would admit C here, running 2 bodies at cap 1.
const third = await tool.execute("tc-3", { agent: "task", id: "Third", assignment: "Work C." } as TaskParams);
const thirdJob = manager.getJob(third.details!.async!.jobId)!;
await Bun.sleep(50);
expect(started).toEqual(["First"]);
expect(thirdJob.queued).toBe(true);
// A finishing admits C — the cap still cycles normally.
gates.get("First")!.resolve();
await firstJob.promise;
await pollUntil(() => started.length === 2);
expect(started).toEqual(["First", "Third"]);
// D queued behind running C stays serialized: if B's teardown had
// double-released, two permits would be free and D would start now.
const fourth = await tool.execute("tc-4", { agent: "task", id: "Fourth", assignment: "Work D." } as TaskParams);
const fourthJob = manager.getJob(fourth.details!.async!.jobId)!;
await Bun.sleep(50);
expect(started).toEqual(["First", "Third"]);
expect(fourthJob.queued).toBe(true);
gates.get("Third")!.resolve();
await thirdJob.promise;
await pollUntil(() => started.length === 3);
gates.get("Fourth")!.resolve();
await fourthJob.promise;
expect(started).toEqual(["First", "Third", "Fourth"]);
expect(firstJob.status).toBe("completed");
expect(thirdJob.status).toBe("completed");
expect(fourthJob.status).toBe("completed");
});
for (const maxConcurrency of [0, 0.5]) {
it(`runs spawn job bodies unbounded when task.maxConcurrency is ${maxConcurrency}`, async () => {
vi.spyOn(discoveryModule, "discoverAgents").mockResolvedValue({
@@ -15,6 +15,7 @@ import {
parseIsolationMode,
} from "@oh-my-pi/pi-coding-agent/task/worktree";
import * as jj from "@oh-my-pi/pi-coding-agent/utils/jj";
import * as git from "@oh-my-pi/pi-coding-agent/utils/git";
import * as natives from "@oh-my-pi/pi-natives";
import { removeWithRetries, setWorktreesDir } from "@oh-my-pi/pi-utils";
@@ -235,6 +236,108 @@ describe("worktree isolation helpers", () => {
expect(stashList).toBe("");
});
// Regression for #4175: a stash-pop conflict used to leave stage 1/2/3
// unmerged entries in `.git/index` (no `MERGE_HEAD`, no way to abort).
// The corrupted index survived indefinitely and every subsequent
// overlay-isolated task read it through the lower layer, so
// `captureRepoDeltaPatch` produced `diff --cc` output that `git apply`
// rejects. mergeTaskBranches MUST leave the index clean regardless of
// whether the stash could be popped.
it("keeps the index clean when stash pop would conflict with a cherry-picked change", async () => {
// User's WIP touches the same file the task branch modifies, so a
// naive stash push → cherry-pick → stash pop conflicts on pop.
await fs.writeFile(path.join(repo, "merged.txt"), "user wip\n");
const result = await mergeTaskBranches(repo, [{ branchName: TASK_BRANCH, taskId: "task-1" }]);
const [status, unmerged, stashList, headContent] = await Promise.all([
runGit(repo, ["status", "--porcelain=v1"]),
runGit(repo, ["ls-files", "--unmerged"]),
runGit(repo, ["stash", "list"]),
fs.readFile(path.join(repo, "merged.txt"), "utf8"),
]);
// Cherry-pick landed on HEAD; only the WIP restore was declined.
expect(result.merged).toEqual([TASK_BRANCH]);
expect(result.failed).toEqual([]);
expect(result.stashConflict).toBeDefined();
// The invariant that was previously broken: no unmerged entries.
expect(unmerged).toBe("");
// Working tree matches the merged HEAD, and the WIP is preserved
// as a stash entry for the user to reconcile manually.
expect(status).toBe("");
expect(headContent).toBe("task branch change\n");
expect(stashList).toContain("omp-task-merge");
// Downstream contract: with a clean index, captureDeltaPatch
// produces a valid unified diff (not `diff --cc`) that a
// subsequent isolated task's `git apply --cached` accepts.
// Editing a tracked file keeps the shared fixture clean —
// `reset --hard` on the next test restores it.
const baseline = await captureBaseline(repo);
await fs.writeFile(path.join(repo, "staged.txt"), "downstream edit\n");
const delta = await captureDeltaPatch(repo, baseline);
expect(delta.rootPatch).not.toContain("diff --cc");
expect(delta.rootPatch).toContain("+downstream edit");
});
it("cleans restored stash files with literal pathspecs", async () => {
// Force the fallback branch: preflight would normally refuse this
// pop before Git can restore anything, but mode/delete edge cases can
// still pass preflight and fail during the actual stash pop. Git can
// restore unrelated untracked files before reporting the tracked
// conflict. If the task branch also adds an ignore rule for that
// restored path, the fallback must clean the restored ignored path
// without interpreting stash-derived filenames as pathspec magic.
const magicName = ":(glob)*";
const buildLog = path.join(repo, "build.log");
const ignoredBranch = "task/ignored-restored-untracked";
await fs.writeFile(path.join(repo, ".gitignore"), "*.log\n");
await runGit(repo, ["add", ".gitignore"]);
await runGit(repo, ["commit", "-q", "-m", "ignore-build-artifacts"]);
await runGit(repo, ["checkout", "-q", "-b", ignoredBranch]);
await Promise.all([
fs.writeFile(path.join(repo, "merged.txt"), "task branch change\n"),
fs.writeFile(path.join(repo, ".gitignore"), `*.log\n${magicName}\n`),
]);
await runGit(repo, ["add", ".gitignore", "merged.txt"]);
await runGit(repo, ["commit", "-q", "-m", "task-change-ignored-note"]);
await runGit(repo, ["checkout", "-q", BASE_BRANCH]);
try {
vi.spyOn(git.patch, "canApplyText").mockResolvedValue(true);
await fs.writeFile(path.join(repo, "merged.txt"), "user wip\n");
await fs.writeFile(path.join(repo, magicName), "untracked wip\n");
await fs.writeFile(buildLog, "ignored build artifact\n");
const result = await mergeTaskBranches(repo, [{ branchName: ignoredBranch, taskId: "task-1" }]);
const [status, unmerged, stashList, headContent, magicExists, buildLogExists] = await Promise.all([
runGit(repo, ["status", "--porcelain=v1"]),
runGit(repo, ["ls-files", "--unmerged"]),
runGit(repo, ["stash", "list"]),
fs.readFile(path.join(repo, "merged.txt"), "utf8"),
Bun.file(path.join(repo, magicName)).exists(),
Bun.file(buildLog).exists(),
]);
expect(result.merged).toEqual([ignoredBranch]);
expect(result.failed).toEqual([]);
expect(result.stashConflict).toBeDefined();
expect(unmerged).toBe("");
expect(status).toBe("");
expect(magicExists).toBe(false);
expect(buildLogExists).toBe(true);
expect(headContent).toBe("task branch change\n");
expect(stashList).toContain("omp-task-merge");
} finally {
await cleanupTaskBranches(repo, [ignoredBranch]);
await Promise.all([
fs.rm(path.join(repo, magicName), { force: true }),
fs.rm(buildLog, { force: true }),
]);
}
});
it("commits isolated edits when parent dirt only changes nearby context", async () => {
const fixtureName = "EXP_DIRTY_TEST.txt";
const fixturePath = path.join(repo, fixtureName);
@@ -32,8 +32,8 @@ function resultText(result: { content: Array<{ type: string; text?: string }> })
.join("\n");
}
function stubLoadPage(body: string, contentType: string): void {
vi.spyOn(scrapers, "loadPage").mockImplementation(async requestedUrl => ({
function stubLoadPage(body: string, contentType: string) {
return vi.spyOn(scrapers, "loadPage").mockImplementation(async requestedUrl => ({
ok: true,
status: 200,
finalUrl: requestedUrl,
@@ -116,4 +116,82 @@ describe("search tools with external URL paths", () => {
expect(text).toContain("remoteNeedle");
expect(text).not.toContain("Parse issues");
});
it("search materializes a scheme-less www. scope like its canonical spelling", async () => {
const loadPage = stubLoadPage("alpha\nremote needle\nomega\n", "text/plain");
const tools = await createTools(createSession(testDir));
const tool = tools.find(entry => entry.name === "grep");
expect(tool).toBeDefined();
const result = await tool!.execute("search-url-www", {
pattern: "remote needle",
paths: ["www.example.com/notes.txt"],
});
expect(resultText(result)).toContain("remote needle");
expect(loadPage).toHaveBeenCalledWith("https://www.example.com/notes.txt", expect.anything());
});
it("search repairs a collapsed https:/ scheme before materializing", async () => {
const loadPage = stubLoadPage("alpha\nremote needle\nomega\n", "text/plain");
const tools = await createTools(createSession(testDir));
const tool = tools.find(entry => entry.name === "grep");
expect(tool).toBeDefined();
const result = await tool!.execute("search-url-collapsed", {
pattern: "remote needle",
paths: ["https:/example.com/notes.txt"],
});
expect(resultText(result)).toContain("remote needle");
expect(loadPage).toHaveBeenCalledWith("https://example.com/notes.txt", expect.anything());
});
it("search prefers an existing local directory named like a www. host", async () => {
const loadPage = stubLoadPage("remote body\n", "text/plain");
await fs.mkdir(path.join(testDir, "www.example.com"), { recursive: true });
await fs.writeFile(path.join(testDir, "www.example.com", "notes.txt"), "local needle\n");
const tools = await createTools(createSession(testDir));
const tool = tools.find(entry => entry.name === "grep");
expect(tool).toBeDefined();
const result = await tool!.execute("search-local-dir", {
pattern: "local needle",
paths: ["www.example.com"],
});
expect(resultText(result)).toContain("local needle");
expect(loadPage).not.toHaveBeenCalled();
});
it("search leaves plain relative paths untouched by URL materialization", async () => {
const loadPage = stubLoadPage("remote body\n", "text/plain");
await fs.mkdir(path.join(testDir, "src"), { recursive: true });
await fs.writeFile(path.join(testDir, "src", "notes.txt"), "local needle\n");
const tools = await createTools(createSession(testDir));
const tool = tools.find(entry => entry.name === "grep");
expect(tool).toBeDefined();
const result = await tool!.execute("search-local-rel", {
pattern: "local needle",
paths: ["src/notes.txt"],
});
expect(resultText(result)).toContain("local needle");
expect(loadPage).not.toHaveBeenCalled();
});
it("search rejects unsupported URL schemes explicitly", async () => {
stubLoadPage("remote body\n", "text/plain");
const tools = await createTools(createSession(testDir));
const tool = tools.find(entry => entry.name === "grep");
expect(tool).toBeDefined();
await expect(
tool!.execute("search-url-ftp", {
pattern: "needle",
paths: ["ftp://example.com/notes.txt"],
}),
).rejects.toThrow("Cannot search external URL");
});
});
@@ -410,6 +410,89 @@ describe("YieldTool", () => {
);
});
it("rejects unknown incremental labels for JTD discriminator (oneOf-of-closed) caller schemas", async () => {
// JTD discriminator output schemas compile to a top-level `oneOf` of closed object
// variants (no root `additionalProperties: false`, no `allOf`). Closure MUST be derived
// across the union — otherwise a stale label is accepted here and, once a sibling section
// exhausts MAX_SCHEMA_RETRIES, finalizeSubprocessOutput honors schemaOverridden and the
// stale label lands in a "successful" result (PR #3927 review).
const tool = new YieldTool(
createSession({
outputSchema: {
discriminator: "verdict",
mapping: {
clean: {
properties: { issue_key: { type: "string" } },
},
blockers: {
properties: { issue_key: { type: "string" } },
optionalProperties: {
blockers: { elements: { properties: { title: { type: "string" } } } },
},
},
},
},
}),
);
await expect(
tool.execute("call-jtd-discriminator-stale-label", {
type: ["findings"],
result: { data: { title: "native reviewer finding" } },
} as never),
).rejects.toThrow(
/Section "findings" uses unknown incremental yield label\(s\): "findings"\. Resubmit with one of the schema's labels: "issue_key", "verdict", "blockers"\./,
);
// A label declared by only ONE variant is known: union semantics are disjunctive, the
// assembled output only has to match one variant.
const singleVariantLabel = await tool.execute("call-jtd-discriminator-variant-label", {
type: ["blockers"],
result: { data: { title: "blocker from the blockers variant" } },
} as never);
expect(singleVariantLabel.details?.data).toEqual({ title: "blocker from the blockers variant" });
// The discriminator property itself is declared by every variant.
const discriminatorLabel = await tool.execute("call-jtd-discriminator-tag-label", {
type: ["verdict"],
result: { data: "blockers" },
} as never);
expect(discriminatorLabel.details?.data).toBe("blockers");
});
it("does not gate incremental labels when a oneOf variant is open", async () => {
// One open variant (`additionalProperties` not false) accepts arbitrary top-level
// properties, so the union places no constraint on labels — engaging the gate would
// reject labels the schema actually allows.
const tool = new YieldTool(
createSession({
outputSchema: {
oneOf: [
{
type: "object",
properties: {
issue_key: { type: "string" },
verdict: { enum: ["clean", "blockers"] },
},
required: ["issue_key", "verdict"],
additionalProperties: false,
},
{
type: "object",
properties: { notes: { type: "string" } },
},
],
},
}),
);
const result = await tool.execute("call-open-variant-label", {
type: ["findings"],
result: { data: { title: "accepted because one variant is open" } },
} as never);
expect(result.details?.data).toEqual({ title: "accepted because one variant is open" });
});
it("detects array-valued labels when the closed caller schema is a root $ref", () => {
const labels = arrayValuedLabels({
$ref: "#/$defs/Closed",
+1
View File
@@ -6,6 +6,7 @@
- Fixed an issue in the mobile collaboration web UI where 'ask' questions were displayed without response controls.
- Fixed the agent transcript drawer hot-retrying forever when the host reports a terminal transcript error (such as an oversized row); the error now stops polling and is shown below any rows already loaded.
- Fixed pre-welcome host `error` frames (such as a protocol-version rejection) being invisible until the welcome timeout; the session now ends immediately with the host's reason.
## [16.2.0] - 2026-06-27
+8
View File
@@ -379,6 +379,14 @@ export class GuestClient {
this.#end(frame.reason);
return; // #end already committed
case "error":
if (!this.#welcomed) {
// Pre-welcome errors are the host's targeted reply to our
// hello (e.g. protocol mismatch): no welcome will follow.
// End with the host's reason instead of waiting out the
// welcome timeout.
this.#end(frame.message);
return; // #end already committed
}
this.#pushNotice("error", frame.message);
break;
default:
+14 -2
View File
@@ -11,7 +11,7 @@ import type {
WireMessage,
} from "@oh-my-pi/pi-wire";
import { GuestClient } from "../src/lib/client";
import { encodeBase64Url } from "../src/lib/link";
import { COLLAB_PROTO, encodeBase64Url } from "../src/lib/link";
import { CollabSocket } from "../src/lib/socket";
const LINK = `roomroomroom1234#${encodeBase64Url(new Uint8Array(32))}`;
@@ -53,7 +53,7 @@ function messageEntry(id: string, message: WireMessage): SessionEntry {
}
function welcomeFrame(entryCount = 0, readOnly?: boolean): HostFrame {
return { t: "welcome", proto: 2, header: HEADER, state: STATE, agents: AGENTS, entryCount, readOnly };
return { t: "welcome", proto: COLLAB_PROTO, header: HEADER, state: STATE, agents: AGENTS, entryCount, readOnly };
}
function snapshotChunk(entries: SessionEntry[], final = true): HostFrame {
@@ -234,6 +234,18 @@ describe("GuestClient frame apply", () => {
expect(notices[0]).toMatchObject({ level: "error", message: "boom" });
});
it("a pre-welcome error (hello rejection, e.g. protocol mismatch) ends the session with the host's reason", () => {
const client = new GuestClient(LINK, "tester");
client.applyFrameForTest({
t: "error",
message: `protocol mismatch: host speaks v${COLLAB_PROTO}, guest sent v${COLLAB_PROTO - 1}`,
});
const snap = client.getSnapshot();
expect(snap.phase).toBe("ended");
expect(snap.endedReason).toContain("protocol mismatch");
expect(snap.endedReason).toContain(`v${COLLAB_PROTO}`);
});
it("tracks host UI requests and sends responses", () => {
const sent: GuestFrame[] = [];
const sendSpy = vi.spyOn(CollabSocket.prototype, "send").mockImplementation((frame: GuestFrame) => {
+3
View File
@@ -6,6 +6,9 @@
- Fixed an issue where snapshot tag collisions could cause line-anchored edits to be incorrectly applied to unrelated content.
- Fixed tracking of edit anchors when earlier in-session insertions or deletions shift unchanged target lines.
- Fixed recovery and edit-preview paths still treating a 16-bit snapshot tag as exact identity after collision-aware retention landed: ambiguous colliding tags now fall through to recovery/rejection (`byHashExact`) instead of applying anchors against the most-recent collider.
- Reduced stale-anchor remap validation from quadratic to linear: duplicate-line detection and anchor-neighbor context are now precomputed once per pass instead of scanning the whole file per anchor.
- Fixed hashline edit guidance for Markdown list rows by teaching `+- item` escaping in the model prompt and minus-row parser error. ([#4179](https://github.com/can1357/oh-my-pi/issues/4179))
## [16.2.8] - 2026-06-30
+1 -1
View File
@@ -54,7 +54,7 @@ export const BARE_BODY_AUTO_PIPED_WARNING =
/** Unified-diff-style `-` row in a hunk body. */
export const MINUS_ROW_REJECTED =
"`-` rows are not valid; the range already names the lines being changed. For a literal `-` line, write `+-…`.";
"`-` rows are not valid; the range already names the lines being changed. For Markdown bullets or other literal `-` lines, prefix the literal row with `+`: `+- item`.";
/** Replace hunk with no body. */
export const EMPTY_REPLACE = `\`SWAP N${HL_RANGE_SEP}M:\` needs at least one \`+TEXT\` body row. To delete lines, use \`DEL N${HL_RANGE_SEP}M\`.`;
+10 -2
View File
@@ -19,7 +19,7 @@ Single line: `SWAP N.=N:` / `DEL N`. The range is the ORIGINAL lines you touch;
</ops>
<body-rows>
Body rows appear only under a `:` header. Every body row is `+TEXT` — add a literal line `TEXT`, verbatim (leading whitespace kept); `+` alone adds a blank line. No other row kind. NEVER write `-old` or a bare/context line. To keep a line, leave it out of every range. To insert a literal line starting with `-` or `+`, prefix it: `+-x`, `++x`.
Body rows appear only under a `:` header. Every body row is `+TEXT` — add a literal line `TEXT`, verbatim (leading whitespace kept); `+` alone adds a blank line. No other row kind. NEVER write `-old` or a bare/context line. To keep a line, leave it out of every range. Literal lines starting with `-`/`+` still need the body prefix: Markdown `- item` → `+- item`, `+ item` → `++ item`.
</body-rows>
<rules>
@@ -103,6 +103,14 @@ INS.TAIL:
+greet("everyone")
```
Insert Markdown bullets — the leading `+` is the body-row marker; the file receives `- task`:
```
[PLAN.md#A1B2]
INS.POST 2:
+- task
+ - nested task
```
Replace the whole `greet` function block — `SWAP.BLK 1:` resolves lines 1–3 (the `def` header through `print(msg)`); line 4 is a separate statement and stays:
```
[greet.py#A1B2]
@@ -160,5 +168,5 @@ INS.POST 3:
If you remember nothing else:
1. RE-GROUND AFTER EVERY EDIT. Every apply mints a fresh `#TAG` and renumbers — take the next edit's numbers from the edit response or a fresh `read`. Stale tag or surprise? STOP, re-`read`.
2. RANGES ARE TIGHT. Cover only lines that change; a stale wide range shreds everything it spans. Whole construct → `SWAP.BLK N`.
3. THE BODY IS THE FINAL CONTENT. Only `+TEXT` rows; never `-old`/context lines. The range does the deleting.
3. THE BODY IS THE FINAL CONTENT. Every body row starts with `+`; Markdown bullets use `+- item`, not `- item`.
</critical>
+57 -25
View File
@@ -130,37 +130,59 @@ function buildLineMap(previousText: string, currentText: string): Map<number, nu
return map;
}
function lineIsDuplicated(lines: readonly string[], line: number): boolean {
const value = lines[line - 1];
return lines.indexOf(value) !== lines.lastIndexOf(value);
/** Values appearing two or more times in `lines`, for O(1) duplicate checks. */
function collectDuplicatedValues(lines: readonly string[]): Set<string> {
const seen = new Set<string>();
const duplicated = new Set<string>();
for (const value of lines) {
if (seen.has(value)) duplicated.add(value);
else seen.add(value);
}
return duplicated;
}
function nearestContextLine(
line: number,
direction: -1 | 1,
anchorLines: ReadonlySet<number>,
lineCount: number,
): number | undefined {
for (let candidate = line + direction; candidate >= 1 && candidate <= lineCount; candidate += direction) {
if (!anchorLines.has(candidate)) return candidate;
interface AnchorNeighbors {
/** Nearest non-anchor line below the anchor's run, or `undefined` at the file edge. */
before: number | undefined;
/** Nearest non-anchor line above the anchor's run, or `undefined` at the file edge. */
after: number | undefined;
}
/**
* Nearest non-anchor context line on each side of every anchor, computed in
* one sweep over the sorted anchor set. Anchors in one contiguous run share
* both neighbors (the lines just outside the run), so this replaces the
* per-anchor directional walk across anchored ranges — O(anchors²) on a
* large block replacement — with one O(anchors log anchors) pass.
*/
function computeAnchorNeighbors(anchorLines: ReadonlySet<number>, lineCount: number): Map<number, AnchorNeighbors> {
const sorted = [...anchorLines].sort((a, b) => a - b);
const neighbors = new Map<number, AnchorNeighbors>();
for (let i = 0; i < sorted.length; ) {
let j = i;
while (j + 1 < sorted.length && sorted[j + 1] === sorted[j] + 1) j++;
const start = sorted[i];
const end = sorted[j];
const before = start - 1 >= 1 && start - 1 <= lineCount ? start - 1 : undefined;
const after = end + 1 <= lineCount ? end + 1 : undefined;
for (let k = i; k <= j; k++) neighbors.set(sorted[k], { before, after });
i = j + 1;
}
return undefined;
return neighbors;
}
function validateDuplicateAnchorContext(
line: number,
mapped: number,
previousLines: readonly string[],
neighbors: AnchorNeighbors,
lineMap: ReadonlyMap<number, number>,
anchorLines: ReadonlySet<number>,
): boolean {
let checked = false;
const before = nearestContextLine(line, -1, anchorLines, previousLines.length);
const { before, after } = neighbors;
if (before !== undefined) {
checked = true;
if (lineMap.get(before) !== mapped - (line - before)) return false;
}
const after = nearestContextLine(line, 1, anchorLines, previousLines.length);
if (after !== undefined) {
checked = true;
if (lineMap.get(after) !== mapped + (after - line)) return false;
@@ -171,14 +193,12 @@ function validateDuplicateAnchorContext(
function validateUniqueAnchorContext(
line: number,
mapped: number,
previousLines: readonly string[],
neighbors: AnchorNeighbors,
lineMap: ReadonlyMap<number, number>,
anchorLines: ReadonlySet<number>,
): boolean {
const offset = mapped - line;
const after = nearestContextLine(line, 1, anchorLines, previousLines.length);
const { before, after } = neighbors;
if (after !== undefined) return lineMap.get(after) === after + offset;
const before = nearestContextLine(line, -1, anchorLines, previousLines.length);
return before !== undefined && lineMap.get(before) === before + offset;
}
@@ -191,17 +211,25 @@ function validateRemappedAnchorContext(
const previousLines = previousText.split("\n");
const currentLines = currentText.split("\n");
const anchorLines = new Set(collectAnchorLines(edits));
// Precompute once per validation pass: which line values are duplicated,
// and each anchor's nearest non-anchor context. The per-anchor forms —
// indexOf/lastIndexOf full-file scans plus directional walks across
// anchored ranges — are O(anchors×lines) + O(anchors²) and blow up on
// large block replacements.
const duplicatedPrevious = collectDuplicatedValues(previousLines);
const duplicatedCurrent = collectDuplicatedValues(currentLines);
const anchorNeighbors = computeAnchorNeighbors(anchorLines, previousLines.length);
for (const line of anchorLines) {
for (const [line, neighbors] of anchorNeighbors) {
const mapped = lineMap.get(line);
if (mapped === undefined) return false;
if (!lineIsDuplicated(previousLines, line) && !lineIsDuplicated(currentLines, mapped)) {
if (!validateUniqueAnchorContext(line, mapped, previousLines, lineMap, anchorLines)) {
if (!duplicatedPrevious.has(previousLines[line - 1]) && !duplicatedCurrent.has(currentLines[mapped - 1])) {
if (!validateUniqueAnchorContext(line, mapped, neighbors, lineMap)) {
return false;
}
continue;
}
if (!validateDuplicateAnchorContext(line, mapped, previousLines, lineMap, anchorLines)) {
if (!validateDuplicateAnchorContext(line, mapped, neighbors, lineMap)) {
return false;
}
}
@@ -364,7 +392,11 @@ export class Recovery {
*/
tryRecover(args: RecoveryArgs): RecoveryResult | null {
const { path, currentText, fileHash, edits } = args;
const snapshot = this.store.byHash(path, fileHash);
// Collision-safe lookup: when two retained texts share the 16-bit tag
// there is no way to know which one the model's anchors were minted
// against — replaying against the wrong collider would land the edit
// on unrelated content. Refuse and let the caller reject (re-read).
const snapshot = this.store.byHashExact(path, fileHash);
if (!snapshot) return null;
const isHead = isHeadSnapshot(this.store.head(path), snapshot);
const recoveryWarning = isHead ? RECOVERY_EXTERNAL_WARNING : RECOVERY_SESSION_CHAIN_WARNING;
+30 -6
View File
@@ -10,9 +10,9 @@
* Producers (typically `read` / `search` / `write` tools) call
* {@link SnapshotStore.record} with the full normalized text they observed.
* The store hashes it, dedups against the per-path history, and returns the
* tag. Consumers (the patcher) resolve a stale tag back to the recorded full
* text via {@link SnapshotStore.byHash} and 3-way-merge the would-be edit onto
* the live content.
* tag. Consumers (recovery, the patcher) resolve a stale tag back to the
* recorded full text via {@link SnapshotStore.byHashExact} and 3-way-merge the
* would-be edit onto the live content.
*
* The abstract base class lets callers plug in whatever storage they like
* (LRU, persistent SQLite, etc.). {@link InMemorySnapshotStore} ships as a
@@ -49,7 +49,7 @@ export interface Snapshot {
/**
* Storage seam for full-file version snapshots. The patcher calls {@link head}
* for the latest version of a path and {@link byHash} when it needs the
* for the latest version of a path and {@link byHashExact} when it needs the
* specific historical version a section's stale tag names.
*/
export abstract class SnapshotStore {
@@ -59,11 +59,21 @@ export abstract class SnapshotStore {
/**
* Recorded version for `path` whose tag equals `hash`, or `null`. When two
* distinct texts collide on the 16-bit tag, returns the most-recently
* recorded one; callers that need the exact content must verify
* {@link Snapshot.text} against the live text (or use {@link byContent}).
* recorded one; callers that treat the tag as content identity must use
* {@link byHashExact} (or verify {@link Snapshot.text} via {@link byContent}).
*/
abstract byHash(path: string, hash: string): Snapshot | null;
/**
* Collision-safe {@link byHash}: the single retained version for `path`
* whose tag equals `hash`, or `null` when none is retained OR when two or
* more distinct texts collide on the tag. In the collision case there is
* no way to know which retained text the model's line anchors were minted
* against, so consumers that replay anchors (recovery, previews) must
* refuse rather than pick one.
*/
abstract byHashExact(path: string, hash: string): Snapshot | null;
/**
* Recorded version for `path` whose {@link Snapshot.text} equals `fullText`,
* or `null`. Disambiguates hash collisions where two distinct file states
@@ -179,6 +189,20 @@ export class InMemorySnapshotStore extends SnapshotStore {
return history?.find(version => version.hash === hash) ?? null;
}
byHashExact(path: string, hash: string): Snapshot | null {
const history = this.#versions.get(path);
if (history === undefined) return null;
let match: Snapshot | null = null;
for (const version of history) {
if (version.hash !== hash) continue;
// Two retained versions with one tag are distinct texts by
// construction (record() dedups on full-text equality) — ambiguous.
if (match !== null) return null;
match = version;
}
return match;
}
byContent(path: string, fullText: string): Snapshot | null {
const history = this.#versions.get(path);
return history?.find(version => version.text === fullText) ?? null;
+8 -4
View File
@@ -166,12 +166,16 @@ describe("hashline body contracts", () => {
expect(applyEdits(FILE, result.edits).text).toBe('a\n1: "one",\n2: "two",\nd\ne');
});
it("rejects `-` body rows with a teaching error", () => {
expect(() => parsePatch("SWAP 2.=2:\n-old\n+new")).toThrow(/`-` rows are not valid/);
it("rejects `-` body rows with Markdown bullet escape guidance", () => {
expect(() => parsePatch("SWAP 2.=2:\n-old\n+new")).toThrow(
/Markdown bullets or other literal `-` lines.*`\+- item`/,
);
});
it("allows literal text that begins with `-` or `+` when prefixed with `+`", () => {
expect(applyPatch(FILE, "SWAP 2.=2:\n+-literal\n++plus")).toBe("a\n-literal\n+plus\nc\nd\ne");
it("allows literal Markdown bullets and plus-prefixed text when prefixed with `+`", () => {
expect(applyPatch(FILE, "SWAP 2.=2:\n+- item\n+ - nested\n++plus")).toBe(
"a\n- item\n - nested\n+plus\nc\nd\ne",
);
});
it("treats empty replace as delete and still rejects empty insert", () => {
@@ -11,6 +11,7 @@
*/
import { describe, expect, it } from "bun:test";
import {
computeFileHash,
InMemorySnapshotStore,
parsePatch,
RECOVERY_LINE_REMAP_WARNING,
@@ -157,4 +158,95 @@ describe("Recovery — session-chain replay anchor-content gate", () => {
expect(recovered).toBeNull();
});
it("recovers duplicate-line anchors shifted by a prior insertion when context still matches", () => {
// Remap-parity pin for the linearized validator: an anchor RANGE
// covering a duplicated line ("DUP" appears twice) plus a unique line
// must still remap through a prior insertion — the duplicate-context
// and unique-context branches both accept exactly as before.
const store = new InMemorySnapshotStore();
const v0Text = lines("alpha", "DUP", "beta", "DUP", "omega");
const h0 = store.record(PATH, v0Text);
const v1Text = lines("alpha", "INSERTED", "DUP", "beta", "DUP", "omega");
store.record(PATH, v1Text);
const { edits } = parsePatch("SWAP 3.=4:\n+B-MODEL\n+MODEL");
const recovered = new Recovery(store).tryRecover({
path: PATH,
currentText: v1Text,
fileHash: h0,
edits,
});
expect(recovered).not.toBeNull();
expect(recovered?.text).toBe(lines("alpha", "INSERTED", "DUP", "B-MODEL", "MODEL", "omega"));
expect(recovered?.warnings).toContain(RECOVERY_LINE_REMAP_WARNING);
});
});
/**
* Brute-force two distinct texts sharing one 4-hex tag. 16-bit tags collide
* within a few hundred candidates (birthday bound), so this stays cheap.
* Texts share `template` around a varying middle line so line-anchored edits
* against one collider are plausible-but-wrong against the other.
*/
function findCollidingTexts(): { older: string; newer: string } {
const textFor = (n: number): string => lines("shared head", `unique payload ${n}`, "shared tail");
const byTag = new Map<string, number>();
for (let n = 0; ; n++) {
const text = textFor(n);
const tag = computeFileHash(text);
const prior = byTag.get(tag);
if (prior !== undefined) return { older: textFor(prior), newer: text };
byTag.set(tag, n);
}
}
describe("Recovery — colliding snapshot tags", () => {
it("refuses recovery when two retained texts share the section's tag", () => {
const { older, newer } = findCollidingTexts();
const tag = computeFileHash(older);
expect(computeFileHash(newer)).toBe(tag);
expect(newer).not.toBe(older);
const store = new InMemorySnapshotStore();
store.record(PATH, older);
store.record(PATH, newer);
// Live file drifted away from both colliders, so recovery cannot
// shortcut via live==snapshot; it must pick a base text for the tag.
// The model's edit was authored against `older` (line 2 = its unique
// payload). Resolving the tag to the most-recent collider would 3-way
// merge the stale payload onto `newer`'s unrelated line 2 — silent
// corruption. The ambiguous tag must refuse instead.
const currentText = `${newer}drifted trailer\n`;
const recovered = new Recovery(store).tryRecover({
path: PATH,
currentText,
fileHash: tag,
edits: parsePatch("SWAP 2.=2:\n+model payload").edits,
});
expect(recovered).toBeNull();
});
it("still recovers when exactly one retained text carries the tag", () => {
// Same drift scenario minus the collision: the unambiguous tag keeps
// recovering via 3-way merge, proving the collision gate above does
// not overreach.
const { older } = findCollidingTexts();
const store = new InMemorySnapshotStore();
const tag = store.record(PATH, older);
const currentText = `${older}drifted trailer\n`;
const recovered = new Recovery(store).tryRecover({
path: PATH,
currentText,
fileHash: tag,
edits: parsePatch("SWAP 2.=2:\n+model payload").edits,
});
expect(recovered).not.toBeNull();
expect(recovered?.text).toBe(lines("shared head", "model payload", "shared tail", "drifted trailer"));
});
});
+20
View File
@@ -151,5 +151,25 @@ describe("InMemorySnapshotStore", () => {
expect(store.byContent(PATH, COLLIDE_A)?.seenLines).toEqual(new Set([1, 2]));
expect(store.byContent(PATH, COLLIDE_B)).toBeNull();
});
it("byHashExact returns the single retained collider, or null when the tag is ambiguous", () => {
const store = new InMemorySnapshotStore();
expect(store.byHashExact(PATH, computeFileHash(COLLIDE_A))).toBeNull();
const tag = store.record(PATH, COLLIDE_A);
// Exactly one retained text for the tag → safe to resolve.
expect(store.byHashExact(PATH, tag)?.text).toBe(COLLIDE_A);
store.record(PATH, COLLIDE_B);
// Two distinct texts now share the tag: there is no way to know
// which one a section's anchors were minted against, so the
// collision-safe lookup must refuse while byHash still surfaces
// the most-recent collider.
expect(store.byHashExact(PATH, tag)).toBeNull();
expect(store.byHash(PATH, tag)?.text).toBe(COLLIDE_B);
// Other paths and unknown tags stay null.
expect(store.byHashExact(OTHER, tag)).toBeNull();
expect(store.byHashExact(PATH, tag === "0000" ? "FFFF" : "0000")).toBeNull();
});
});
});
+1
View File
@@ -4,6 +4,7 @@
### Fixed
- Fixed the bounded stdin escape parser scanning past its own cap: OSC/DCS/APC terminator lookups used unbounded `indexOf`, so one oversized unterminated sequence could still block the event loop; scans are now strictly bounded to the per-sequence byte cap.
- Fixed an issue where large Windows terminal session restores could get truncated mid-frame during ConPTY full-paint resume.
## [16.2.13] - 2026-07-01
+28 -10
View File
@@ -132,25 +132,43 @@ function resolveEscapeEnd(buffer: string, pos: number, length: number, resumeSea
}
case 0x5d /* ] */:
{
// OSC: ESC ] ... BEL or ST (ESC \).
// OSC: ESC ] ... BEL or ST (ESC \). Scan is bounded to
// [searchFrom, scanLimit): `String#indexOf` has no end bound, so
// an unterminated payload delivered as one huge chunk would
// otherwise be scanned to the end of the buffer — past the cap
// this function exists to enforce. `resumeSearchFrom - 1` keeps
// the one-byte overlap so an `ESC \` split across chunks is
// still found (the prior call's trailing ESC is re-inspected).
const searchFrom = Math.max(pos + 2, resumeSearchFrom - 1);
const scanLimit = Math.min(length, pos + MAX_STRING_SEQ_BYTES);
const belIndex = buffer.indexOf("\x07", searchFrom);
const stIndex = buffer.indexOf("\x1b\\", searchFrom);
let end = -1;
if (belIndex !== -1 && belIndex + 1 <= scanLimit) end = belIndex + 1;
if (stIndex !== -1 && stIndex + 2 <= scanLimit && (end === -1 || stIndex + 2 < end)) end = stIndex + 2;
if (end !== -1) return end;
for (let i = searchFrom; i < scanLimit; i++) {
const code = buffer.charCodeAt(i);
if (code === 0x07 /* BEL */) return i + 1;
if (code === 0x1b /* ESC */) {
// `ESC \` (ST) must end within the cap; a lone trailing
// ESC at the buffer edge stays incomplete and is
// re-examined next call via the resume overlap.
if (i + 1 < scanLimit && buffer.charCodeAt(i + 1) === 0x5c /* \ */) return i + 2;
}
}
return length - pos >= MAX_STRING_SEQ_BYTES ? -2 : -1;
}
case 0x50 /* P */:
case 0x5f /* _ */:
{
// DCS / APC: ESC P/_ ... ST (ESC \).
// DCS / APC: ESC P/_ ... ST (ESC \). Same bounded scan and
// split-ST overlap as the OSC branch, minus BEL.
const searchFrom = Math.max(pos + 2, resumeSearchFrom - 1);
const scanLimit = Math.min(length, pos + MAX_STRING_SEQ_BYTES);
const stIndex = buffer.indexOf("\x1b\\", searchFrom);
if (stIndex !== -1 && stIndex + 2 <= scanLimit) return stIndex + 2;
for (let i = searchFrom; i < scanLimit; i++) {
if (
buffer.charCodeAt(i) === 0x1b /* ESC */ &&
i + 1 < scanLimit &&
buffer.charCodeAt(i + 1) === 0x5c /* \ */
) {
return i + 2;
}
}
return length - pos >= MAX_STRING_SEQ_BYTES ? -2 : -1;
}
case 0x4f /* O */:
+49
View File
@@ -103,6 +103,29 @@ describe("StdinBuffer", () => {
expect(emittedSequences).toEqual(["\x1b[<35;20;5m"]);
});
it("reassembles an OSC whose ST is split exactly at the chunk boundary", () => {
// Chunk 1 ends on the ESC of `ESC \`; chunk 2 opens with the `\`.
// The resume overlap (`resumeSearchFrom - 1`) must re-inspect the
// trailing ESC, or the terminator is never seen and the payload
// leaks via timeout flush as raw bytes.
processInput("\x1b]52;c;aGVsbG8=\x1b");
expect(emittedSequences).toEqual([]);
expect(buffer.getBuffer()).toBe("\x1b]52;c;aGVsbG8=\x1b");
processInput("\\");
expect(emittedSequences).toEqual(["\x1b]52;c;aGVsbG8=\x1b\\"]);
expect(buffer.getBuffer()).toBe("");
});
it("reassembles a DCS whose ST is split exactly at the chunk boundary", () => {
processInput("\x1bPq#0;2;0;0;0\x1b");
expect(emittedSequences).toEqual([]);
processInput("\\");
expect(emittedSequences).toEqual(["\x1bPq#0;2;0;0;0\x1b\\"]);
expect(buffer.getBuffer()).toBe("");
});
it("should flush incomplete sequence after timeout", async () => {
// Non-mouse CSI partial: ambiguous, so it flushes after the timeout.
processInput("\x1b[1;5");
@@ -623,6 +646,32 @@ describe("StdinBuffer", () => {
expect(emittedSequences).toEqual(["\x1b]z\x07", "a", "b", "c"]);
expect(buffer.getBuffer()).toBe("");
});
it("caps an unterminated OSC delivered as one oversized chunk and keeps parsing", () => {
// MAX_STRING_SEQ_BYTES = 16 MiB. A single chunk whose OSC payload
// exceeds the cap with no BEL/ST must cap-flush the capped prefix
// as ONE raw sequence (progress guaranteed, scan bounded to the
// cap — not the whole chunk), deliver the tail per scalar, and
// leave the buffer clean so later input still parses.
const cap = 16 * 1024 * 1024;
const head = "\x1b]5522;";
const tail = "xy";
// Total pre-tail length is exactly `cap`, so the cap-flush consumes
// the whole unterminated sequence and only `tail` remains.
processInput(`${head}${"a".repeat(cap - head.length)}${tail}`);
expect(emittedSequences.length).toBe(1 + tail.length);
expect(emittedSequences[0]!.length).toBe(cap);
expect(emittedSequences[0]!.startsWith("\x1b]5522;")).toBe(true);
expect(emittedSequences.slice(1)).toEqual(["x", "y"]);
expect(buffer.getBuffer()).toBe("");
// Parser state is clean: a normal OSC afterwards completes.
emittedSequences.length = 0;
processInput("\x1b]z\x07");
expect(emittedSequences).toEqual(["\x1b]z\x07"]);
expect(buffer.getBuffer()).toBe("");
});
});
describe("Destroy", () => {
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased]
### Added
- Added `wrapFetchForExtraCa` / `withExtraCaFetch` (moved from `@oh-my-pi/pi-ai` internals): a fetch wrapper that applies `NODE_EXTRA_CA_CERTS` to Bun's `RequestInit.tls.ca`, shared by provider streaming and catalog model discovery.
## [16.2.9] - 2026-06-30
### Added
+1
View File
@@ -29,6 +29,7 @@ export * from "./snowflake";
export * from "./stream";
export * from "./tab-spacing";
export * from "./temp";
export * from "./tls-fetch";
export * from "./type-guards";
export * from "./which";
@@ -1,15 +1,13 @@
/**
* `NODE_EXTRA_CA_CERTS` shim for the Bun fetch path used by every provider
* stream.
* `NODE_EXTRA_CA_CERTS` shim for Bun's `fetch`.
*
* Node's TLS layer honours `NODE_EXTRA_CA_CERTS` natively, but Bun's
* `fetch` does not, and the OpenAI-compatible providers (`openai-responses`,
* `openai-completions`, `openai-codex-responses`, `ollama-chat`, ...) route
* every request through Bun's runtime. Without this wrapper, corporate
* relays and private gateways behind a custom CA bundle fail with
* `fetch` does not, and both the provider streams (`openai-responses`,
* `openai-completions`, `openai-codex-responses`, `ollama-chat`, ...) and
* catalog model discovery (`/models` probes) route every request through
* Bun's runtime. Without this wrapper, corporate relays and private
* gateways behind a custom CA bundle fail with
* `unknown certificate verification error` even when the env var is set.
* (Foundry mTLS in {@link resolveFoundryTlsOptions} consumed the env var for
* the Anthropic path only — issue #3731.)
*
* The wrapper merges the resolved CA bundle into Bun's `RequestInit.tls.ca`.
* Bun's `tls.ca` REPLACES the default trust store when set, so the wrapper
@@ -18,9 +16,29 @@
*/
import * as fs from "node:fs";
import * as tls from "node:tls";
import { $env, isEnoent } from "@oh-my-pi/pi-utils";
import * as AIError from "../error";
import type { FetchImpl } from "../types";
import { $env } from "./env";
import { isEnoent } from "./fs-error";
/**
* `fetch`-compatible function. Accepts any callable matching the standard
* fetch signature; `preconnect` is optional because non-Bun runtimes
* (browsers, test mocks) won't expose it.
*/
export type FetchImpl = ((input: string | URL | Request, init?: RequestInit) => Promise<Response>) & {
preconnect?: typeof globalThis.fetch.preconnect;
};
/**
* `NODE_EXTRA_CA_CERTS` was set but unusable (path does not exist). This is
* a config/contract error, not a transient transport fault — it is never
* retried.
*/
export class ExtraCaError extends Error {
constructor(message: string, options?: { cause?: unknown }) {
super(message, options?.cause === undefined ? undefined : { cause: options.cause });
this.name = "ExtraCaError";
}
}
/** Bun extension to `RequestInit` for the TLS options we touch. */
type BunTlsOptions = {
@@ -41,7 +59,7 @@ type ExtraCaFetch = FetchImpl & { [EXTRA_CA_FETCH_MARKER]?: true };
* Cached resolution of `NODE_EXTRA_CA_CERTS`. Keyed on the env value plus
* the file mtime for path values so on-disk cert rotation (short-lived
* corporate bundles) invalidates the cache instead of pinning the first
* read forever. Mirrors {@link foundryTlsOptionsCacheKey}.
* read forever.
*/
let cacheKey: string | undefined;
let cacheValue: string | undefined;
@@ -56,7 +74,7 @@ let cacheValue: string | undefined;
* shell exports.
* - File path. Anything that does not contain a PEM header is treated as a
* path, matching Node's "extensionless filename is still a path" contract.
* `ENOENT` becomes {@link AIError.ValidationError}; other I/O errors bubble.
* `ENOENT` becomes {@link ExtraCaError}; other I/O errors bubble.
*/
function resolveExtraCa(): string | undefined {
const raw = $env.NODE_EXTRA_CA_CERTS?.trim();
@@ -81,7 +99,7 @@ function resolveExtraCa(): string | undefined {
cacheValue = fs.readFileSync(raw, "utf8");
} catch (error) {
if (isEnoent(error)) {
throw new AIError.ValidationError(`NODE_EXTRA_CA_CERTS path does not exist: ${raw}`);
throw new ExtraCaError(`NODE_EXTRA_CA_CERTS path does not exist: ${raw}`);
}
throw error;
}
@@ -101,7 +119,7 @@ export function __resetExtraCaCache(): void {
* list, the system root store is included alongside the extra bundle —
* Bun's `tls.ca` replaces the default trust store, so omitting roots would
* break every public host. When the caller already curated a list (e.g.
* Anthropic Foundry's {@link resolveFoundryTlsOptions}, which already seeds
* Anthropic Foundry's mTLS options, which already seed
* `tls.rootCertificates`), only the extra CA is appended.
*/
function withExtraCaInit(init: RequestInit | undefined, extraCa: string): RequestInit {
@@ -146,9 +164,10 @@ export function wrapFetchForExtraCa(fetchImpl: FetchImpl): FetchImpl {
}
/**
* Convenience for the stream-entry composition in `stream.ts`. Mirrors
* {@link withRequestDebugFetch} so the proxy/debug/extra-CA wrappers compose
* uniformly. No-op when the env var is unset.
* Convenience for options-bag composition (e.g. the stream-entry path in
* `@oh-my-pi/pi-ai`'s `stream.ts`, which mirrors `withRequestDebugFetch` so
* the proxy/debug/extra-CA wrappers compose uniformly). No-op when the env
* var is unset.
*/
export function withExtraCaFetch<T extends { fetch?: FetchImpl } | undefined>(options: T): T {
if (!$env.NODE_EXTRA_CA_CERTS?.trim()) return options;
@@ -3,9 +3,13 @@ import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import * as tls from "node:tls";
import * as AIError from "../../error";
import type { FetchImpl } from "../../types";
import { __resetExtraCaCache, withExtraCaFetch, wrapFetchForExtraCa } from "../tls-fetch";
import {
__resetExtraCaCache,
ExtraCaError,
type FetchImpl,
withExtraCaFetch,
wrapFetchForExtraCa,
} from "@oh-my-pi/pi-utils/tls-fetch";
const SAMPLE_PEM =
"-----BEGIN CERTIFICATE-----\nMIIBkTCCATegAwIBAgIUF/sample/extra/ca/for/tests/1234567=\n-----END CERTIFICATE-----\n";
@@ -128,12 +132,12 @@ describe("wrapFetchForExtraCa", () => {
expect(ca).not.toContain(SAMPLE_PEM);
});
it("throws ValidationError when the configured path does not exist", async () => {
it("throws ExtraCaError when the configured path does not exist", async () => {
Bun.env.NODE_EXTRA_CA_CERTS = path.join(tmpDir, "missing.pem");
const { fetchImpl } = makeRecordingFetch();
const wrapped = wrapFetchForExtraCa(fetchImpl);
await expect(wrapped("https://corp.example/v1")).rejects.toBeInstanceOf(AIError.ValidationError);
await expect(wrapped("https://corp.example/v1")).rejects.toBeInstanceOf(ExtraCaError);
});
it("is idempotent — wrapping a wrapped fetch returns the same reference", async () => {
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased]
### Breaking Changes
- Bumped `COLLAB_PROTO` to `3`: the `ui-request`/`ui-request-end` host frames and `ui-response` guest frame are part of the handshake contract. Proto-2 guests, which silently dropped host ask requests, are now rejected at hello with the protocol-mismatch error.
### Added
- Added collaboration UI request/response wireframes, enabling browser guests to respond to host-side interactive prompts.
+5 -1
View File
@@ -389,8 +389,12 @@ export type WireFrame = GuestFrame | HostFrame;
* transcript entries follow in `snapshot-chunk` frames, so multi-MB
* sessions are not gated on a single welcome frame fitting under the
* guest's first-welcome timeout.
* - `3`: host asks guests through `ui-request`/`ui-request-end` host frames
* answered by the `ui-response` guest frame. Guests that predate the
* grammar would silently drop `ui-request` (asks hang forever on the
* host), so they must be rejected at hello.
*/
export const COLLAB_PROTO = 2;
export const COLLAB_PROTO = 3;
/** Parameter key used for intent tracing (e.g. prompt explanation/reasoning) */
export const INTENT_FIELD = "i";
+1 -1
View File
@@ -9,7 +9,7 @@ import {
describe("collab wire constants", () => {
it("exports the protocol constants consumed by host, guest, and relay links", () => {
expect(COLLAB_PROTO).toBe(2);
expect(COLLAB_PROTO).toBe(3);
expect(COLLAB_PROMPT_MESSAGE_TYPE).toBe("collab-prompt");
expect(ENVELOPE_HEADER_LENGTH).toBe(4);
expect(ROOM_ID_BYTES).toBe(16);
+1 -1
View File
@@ -170,7 +170,7 @@ mkdir -p "$TARBALL_APP_DIR"
exit 1
}
wire_proto="$(bun -e 'import { COLLAB_PROTO } from "@oh-my-pi/pi-wire"; process.stdout.write(String(COLLAB_PROTO));')"
[ "$wire_proto" = "2" ] || {
[ "$wire_proto" = "3" ] || {
echo "Unexpected @oh-my-pi/pi-wire COLLAB_PROTO: $wire_proto"
exit 1
}