From 574a83f5f51254198e590764d5f33cb7a3fa6980 Mon Sep 17 00:00:00 2001 From: ranxianglei Date: Sun, 16 Aug 2026 19:22:49 +0800 Subject: [PATCH 1/2] fix(pi-ai): honor onPayload replacement payloads in openai-completions, bedrock and cursor MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The onPayload hook contract (README, docs/extensions.md) is that a non-undefined return replaces the provider request payload, and every provider except these three implements it (anthropic, openai-responses family, google, ollama — see the earlier fix for the responses providers). openai-completions, amazon-bedrock and cursor invoked the hook fire-and-forget and sent the original payload, so extensions hooking before_provider_request could never transform the wire body on these providers. - openai-completions: await the hook and apply a non-undefined replacement to the params used for the request body, raw request dump and error-path fallback state - amazon-bedrock: same for the ConverseStream command input - cursor: await the hook for the AgentRunRequest; buildGrpcRequest becomes async and is exported for direct testing (transport is HTTP/2) - devin-agent intentionally unchanged: it does not fire the hook at all (its payload is a protobuf object), which is a feature gap rather than a dropped replacement; documented in README/docs instead - regression tests: captured wire body reflects async/sync replacement, and an undefined return keeps the original payload (completions + bedrock over a mocked fetch; cursor by decoding the serialized run request) --- docs/extensions.md | 2 +- packages/ai/CHANGELOG.md | 4 + packages/ai/README.md | 2 +- packages/ai/src/providers/amazon-bedrock.ts | 5 +- packages/ai/src/providers/cursor.ts | 13 +-- .../ai/src/providers/openai-completions.ts | 5 +- packages/ai/test/bedrock-on-payload.test.ts | 74 +++++++++++++++ packages/ai/test/cursor-on-payload.test.ts | 67 +++++++++++++ .../openai-completions-on-payload.test.ts | 93 +++++++++++++++++++ 9 files changed, 253 insertions(+), 12 deletions(-) create mode 100644 packages/ai/test/bedrock-on-payload.test.ts create mode 100644 packages/ai/test/cursor-on-payload.test.ts create mode 100644 packages/ai/test/openai-completions-on-payload.test.ts diff --git a/docs/extensions.md b/docs/extensions.md index e589f9c88..28b56768c 100644 --- a/docs/extensions.md +++ b/docs/extensions.md @@ -245,7 +245,7 @@ Cancelable pre-events: - `input` - `before_agent_start` -- `before_provider_request` (may replace provider request payload) +- `before_provider_request` (may replace provider request payload — the replacement is applied by every provider that fires the hook, which is all of them except `devin-agent`, which does not fire it) - `after_provider_response` - `context` - `agent_start` / `agent_end` — agent loop lifecycle notification; `agent_end` remains notification-only diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 86d99d599..09ddc8abd 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed OpenAI Completions, Amazon Bedrock, and Cursor providers ignoring `onPayload` replacement payloads. The hook now transforms the actual request body sent upstream on these providers, matching the Anthropic/Gemini/OpenAI Responses replacement contract. `devin-agent` still does not fire the hook (its payload is a protobuf object). + ## [17.3.5] - 2026-08-16 ### Added diff --git a/packages/ai/README.md b/packages/ai/README.md index 8348bc291..0bdba9fc3 100644 --- a/packages/ai/README.md +++ b/packages/ai/README.md @@ -634,7 +634,7 @@ All providers accept the base `StreamOptions` (in addition to provider-specific - `headers`: Extra request headers merged on top of model-defined headers - `sessionId`: Provider-specific session identifier (prompt caching/routing) - `signal`: Abort in-flight requests -- `onPayload`: Callback invoked with the provider request payload just before sending +- `onPayload`: Callback invoked with the provider request payload just before sending. Return a replacement payload object (sync or async) to send it instead of the original; return `undefined` to keep the original. The replacement is applied by every provider that fires the hook — all of them except `devin-agent`, whose payload is a protobuf object and does not fire the hook yet. Example: diff --git a/packages/ai/src/providers/amazon-bedrock.ts b/packages/ai/src/providers/amazon-bedrock.ts index 4473f46d2..8aedcc4c2 100644 --- a/packages/ai/src/providers/amazon-bedrock.ts +++ b/packages/ai/src/providers/amazon-bedrock.ts @@ -342,7 +342,7 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( if (tc.any || tc.tool) additionalModelRequestFields = undefined; } - const commandInput: ConverseStreamRequest = { + let commandInput: ConverseStreamRequest = { messages: convertedMessages, system: buildSystemPrompt(context.systemPrompt, promptCachePolicy), inferenceConfig: { @@ -353,7 +353,8 @@ export const streamBedrock: StreamFunction<"bedrock-converse-stream"> = ( toolConfig, additionalModelRequestFields, }; - options?.onPayload?.(commandInput, model); + const replacementInput = await options?.onPayload?.(commandInput, model); + if (replacementInput !== undefined) commandInput = replacementInput as ConverseStreamRequest; const host = `bedrock-runtime.${region}.amazonaws.com`; const url = `https://${host}/model/${encodeURIComponent(model.id)}/converse-stream`; diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index de00062c8..2a9fd886f 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -615,7 +615,7 @@ export const streamCursor: StreamFunction<"cursor-agent"> = ( const blobStore = conversationBlobStores.get(conversationId) ?? new Map(); conversationBlobStores.set(conversationId, blobStore); const cachedState = conversationStateCache.get(conversationId); - const { requestBytes, conversationState } = buildGrpcRequest(model, context, options, { + const { requestBytes, conversationState } = await buildGrpcRequest(model, context, options, { conversationId, blobStore, conversationState: cachedState, @@ -4642,7 +4642,7 @@ function extractImages(content: (TextContent | ImageContent)[]) { ); } -function buildGrpcRequest( +export async function buildGrpcRequest( model: Model<"cursor-agent">, context: Context, options: CursorOptions | undefined, @@ -4651,11 +4651,11 @@ function buildGrpcRequest( blobStore: Map; conversationState?: ConversationStateStructure; }, -): { +): Promise<{ requestBytes: Uint8Array; blobStore: Map; conversationState: ConversationStateStructure; -} { +}> { const blobStore = state.blobStore; const systemPromptIds = buildCursorSystemPromptJsons(context.systemPrompt).map(json => @@ -4761,7 +4761,7 @@ function buildGrpcRequest( maxMode: cursorMaxMode, }); - const runRequest = create(AgentRunRequestSchema, { + let runRequest = create(AgentRunRequestSchema, { conversationState, action, modelDetails, @@ -4769,7 +4769,8 @@ function buildGrpcRequest( conversationId: state.conversationId, }); - options?.onPayload?.(runRequest, model); + const replacementRequest = await options?.onPayload?.(runRequest, model); + if (replacementRequest !== undefined) runRequest = replacementRequest as typeof runRequest; // Tools are sent later via requestContext (exec handshake) diff --git a/packages/ai/src/providers/openai-completions.ts b/packages/ai/src/providers/openai-completions.ts index cbbdfd226..c7d632b22 100644 --- a/packages/ai/src/providers/openai-completions.ts +++ b/packages/ai/src/providers/openai-completions.ts @@ -663,7 +663,7 @@ const streamOpenAICompletionsOnce = ( : `${trimmedBaseUrl}/chat/completions`; const createCompletionsStream = async (toolStrictModeOverride?: ToolStrictModeOverride) => { const effectiveToolStrictModeOverride = disableStrictTools ? "none" : toolStrictModeOverride; - const { params, strictToolsApplied } = buildParams( + let { params, strictToolsApplied } = buildParams( model, context, options, @@ -682,8 +682,9 @@ const streamOpenAICompletionsOnce = ( applyOpenAIReasoningEffortFallback(params, requestReasoningEffortFallback); } activeReasoningEffortFallbackKey = reasoningEffortFallbackKey; + const replacedParams = await options?.onPayload?.(params, model); + if (replacedParams !== undefined) params = replacedParams as typeof params; activeRequestParams = params; - options?.onPayload?.(params, model); rawRequestDump = { provider: model.provider, api: output.api, diff --git a/packages/ai/test/bedrock-on-payload.test.ts b/packages/ai/test/bedrock-on-payload.test.ts new file mode 100644 index 000000000..fb3298ddd --- /dev/null +++ b/packages/ai/test/bedrock-on-payload.test.ts @@ -0,0 +1,74 @@ +// Regression: amazon-bedrock ignored the onPayload replacement return value +// (fire-and-forget), so the hook could never change the body actually sent +// upstream. The replacement contract matches anthropic / openai-responses / +// google: await the hook and use its non-undefined return as the request body. +import { describe, expect, it, vi } from "bun:test"; +import { streamBedrock } from "@oh-my-pi/pi-ai/providers/amazon-bedrock"; +import type { Context, Model } from "@oh-my-pi/pi-ai/types"; +import { buildModel } from "@oh-my-pi/pi-catalog/build"; + +function model(): Model<"bedrock-converse-stream"> { + return buildModel({ + id: "us.anthropic.claude-haiku-4-5-20251001-v1:0", + name: "haiku", + api: "bedrock-converse-stream", + provider: "amazon-bedrock", + baseUrl: "https://bedrock-runtime.us-east-1.amazonaws.com", + reasoning: false, + input: ["text"], + cost: { input: 5, output: 25, cacheRead: 0.5, cacheWrite: 6.25 }, + contextWindow: 1_000_000, + maxTokens: 128_000, + }); +} + +const context: Context = { + messages: [{ role: "user", content: "hi", timestamp: 0 }], +}; + +// Capture the serialized body the provider sends. The response is an empty +// event stream: the fetch (and thus the body capture) happens before any +// response parsing, and the stream's outcome is irrelevant to the assertion. +async function captureSentBody(onPayload: (payload: unknown) => unknown | Promise): Promise> { + const { promise, resolve } = Promise.withResolvers>(); + const fetchMock = vi.fn(async (_input: string | URL | Request, init?: RequestInit) => { + const body = init?.body; + const text = body instanceof Uint8Array ? new TextDecoder().decode(body) : String(body); + resolve(JSON.parse(text)); + return new Response( + new ReadableStream({ start(controller) { controller.close(); } }), + { status: 200, headers: { "content-type": "application/vnd.amazon.eventstream" } }, + ); + }) as unknown as typeof fetch; + + const stream = streamBedrock(model(), context, { bearerToken: "test-token", fetch: fetchMock, onPayload }); + void (async () => { + try { + for await (const _ of stream) { + // ignore events + } + } catch { + // empty event stream: stream errors are expected and irrelevant + } + })(); + + return promise; +} + +describe("bedrock onPayload replacement", () => { + it("sends an async onPayload replacement body", async () => { + const body = await captureSentBody(async payload => ({ + ...(payload as Record), + messages: [{ role: "user", content: [{ text: "replacement" }] }], + })); + + expect(body.messages).toEqual([{ role: "user", content: [{ text: "replacement" }] }]); + expect(JSON.stringify(body.messages)).not.toContain("hi"); + }, 10_000); + + it("keeps the original body when onPayload returns undefined", async () => { + const body = await captureSentBody(async () => undefined); + + expect(body.messages[0].content[0].text).toBe("hi"); + }, 10_000); +}); diff --git a/packages/ai/test/cursor-on-payload.test.ts b/packages/ai/test/cursor-on-payload.test.ts new file mode 100644 index 000000000..766904d64 --- /dev/null +++ b/packages/ai/test/cursor-on-payload.test.ts @@ -0,0 +1,67 @@ +// Regression: cursor ignored the onPayload replacement return value +// (fire-and-forget), so the hook could never change the request actually sent +// upstream. The replacement contract matches anthropic / openai-responses / +// google: await the hook and use its non-undefined return as the request. +// buildGrpcRequest is exercised directly (the transport is HTTP/2), and the +// serialized run request is decoded back from the wire bytes. +import { describe, expect, it } from "bun:test"; +import { fromBinary } from "@bufbuild/protobuf"; +import { buildGrpcRequest } from "@oh-my-pi/pi-ai/providers/cursor"; +import { AgentClientMessageSchema } from "@oh-my-pi/pi-catalog/discovery/cursor-gen/agent_pb"; +import type { Context, Model } from "@oh-my-pi/pi-ai/types"; +import { buildModel } from "@oh-my-pi/pi-catalog/build"; + +const model: Model<"cursor-agent"> = buildModel({ + id: "cursor-composer-2.5", + name: "Cursor Composer 2.5", + api: "cursor-agent", + provider: "cursor", + baseUrl: "https://api2.cursor.sh", + reasoning: false, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 200_000, + maxTokens: 32_000, +}); + +const context: Context = { + messages: [{ role: "user", content: "Say hello", timestamp: 0 }], +}; + +function decodeRunRequest(requestBytes: Uint8Array): { case: string; value: Record } { + const decoded = fromBinary(AgentClientMessageSchema, requestBytes); + return decoded.message as unknown as { case: string; value: Record }; +} + +describe("cursor onPayload replacement", () => { + it("sends an async onPayload replacement body", async () => { + const { requestBytes } = await buildGrpcRequest( + model, + context, + { + onPayload: async payload => ({ + ...(payload as Record), + customSystemPrompt: "replacement", + }), + }, + { conversationId: "conv-1", blobStore: new Map() }, + ); + + const message = decodeRunRequest(requestBytes); + expect(message.case).toBe("runRequest"); + expect(message.value.customSystemPrompt).toBe("replacement"); + }); + + it("keeps the original body when onPayload returns undefined", async () => { + const { requestBytes } = await buildGrpcRequest( + model, + context, + { onPayload: async () => undefined }, + { conversationId: "conv-1", blobStore: new Map() }, + ); + + const message = decodeRunRequest(requestBytes); + expect(message.case).toBe("runRequest"); + expect(message.value.customSystemPrompt).toBeUndefined(); + }); +}); diff --git a/packages/ai/test/openai-completions-on-payload.test.ts b/packages/ai/test/openai-completions-on-payload.test.ts new file mode 100644 index 000000000..88de415c2 --- /dev/null +++ b/packages/ai/test/openai-completions-on-payload.test.ts @@ -0,0 +1,93 @@ +// Regression: openai-completions ignored the onPayload replacement return +// value (fire-and-forget), so extensions hooking before_provider_request +// could never transform the body actually sent upstream. The replacement +// contract matches anthropic / openai-responses / google: await the hook, +// and use its non-undefined return as the request body. +import { describe, expect, it } from "bun:test"; +import { streamOpenAICompletions } from "@oh-my-pi/pi-ai/providers/openai-completions"; +import type { Context, FetchImpl, Model } from "@oh-my-pi/pi-ai/types"; +import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; + +const completionsModel = { + ...(getBundledModel("openai", "gpt-4o-mini") as Model<"openai-completions">), + api: "openai-completions", +} satisfies Model<"openai-completions">; + +function baseContext(): Context { + return { + messages: [{ role: "user", content: "Say hello", timestamp: Date.now() }], + }; +} + +function createSseFetch(capture?: (body: unknown) => void): FetchImpl { + async function mockFetch(_input: string | URL | Request, init?: RequestInit): Promise { + capture?.(typeof init?.body === "string" ? JSON.parse(init.body) : undefined); + const encoder = new TextEncoder(); + const chunk = (extra: Record) => + `data: ${JSON.stringify({ id: "chatcmpl-payload", object: "chat.completion.chunk", created: 0, model: completionsModel.id, ...extra })}\n\n`; + const sse = + chunk({ choices: [{ index: 0, delta: { role: "assistant", content: "ok" } }] }) + + chunk({ choices: [{ index: 0, delta: {}, finish_reason: "stop" }] }) + + "data: [DONE]\n\n"; + const stream = new ReadableStream({ + start(controller) { + controller.enqueue(encoder.encode(sse)); + controller.close(); + }, + }); + return new Response(stream, { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } + return mockFetch as typeof fetch; +} + +type Body = Record; + +describe("openai-completions onPayload replacement", () => { + it("sends an async onPayload replacement body", async () => { + let captured: Body | undefined; + const result = await streamOpenAICompletions( + completionsModel, + baseContext(), + { + apiKey: "test-key", + fetch: createSseFetch(body => (captured = body as Body)), + onPayload: async payload => ({ + ...(payload as Record), + messages: [{ role: "user", content: "replacement" }], + }), + }, + ).result(); + + expect(result.stopReason).toBe("stop"); + expect(captured?.messages).toEqual([{ role: "user", content: "replacement" }]); + expect(JSON.stringify(captured)).not.toContain("Say hello"); + }, 10_000); + + it("sends a synchronous onPayload replacement body", async () => { + let captured: Body | undefined; + await streamOpenAICompletions(completionsModel, baseContext(), { + apiKey: "test-key", + fetch: createSseFetch(body => (captured = body as Body)), + onPayload: payload => ({ + ...(payload as Record), + messages: [{ role: "user", content: "sync-replacement" }], + }), + }).result(); + + expect(captured?.messages).toEqual([{ role: "user", content: "sync-replacement" }]); + }, 10_000); + + it("keeps the original body when onPayload returns undefined", async () => { + let captured: Body | undefined; + await streamOpenAICompletions(completionsModel, baseContext(), { + apiKey: "test-key", + fetch: createSseFetch(body => (captured = body as Body)), + onPayload: async () => undefined, + }).result(); + + expect(JSON.stringify(captured?.messages)).toContain("Say hello"); + }, 10_000); +}); From 91e27ebe9aa2f40d9f790dd28112f14154b651fc Mon Sep 17 00:00:00 2001 From: ranxianglei Date: Mon, 17 Aug 2026 09:33:12 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(pi-ai):=20cursor=20=E2=80=94=20apply=20?= =?UTF-8?q?customSystemPrompt=20before=20onPayload=20so=20the=20replacemen?= =?UTF-8?q?t=20is=20final?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Per review on #8717: the customSystemPrompt assignment ran after the hook, so when both options.customSystemPrompt and an extension payload replacement were set, the option silently won. Now the option is applied before the hook (the extension can inspect or drop it in its replacement), matching anthropic, where the hook runs right before serialization and is the last word on the wire body. Adds regression tests: replacement drops customSystemPrompt, replacement overrides it, and the option still applies when the hook returns undefined. --- packages/ai/src/providers/cursor.ts | 13 +++--- packages/ai/test/cursor-on-payload.test.ts | 50 ++++++++++++++++++++++ 2 files changed, 58 insertions(+), 5 deletions(-) diff --git a/packages/ai/src/providers/cursor.ts b/packages/ai/src/providers/cursor.ts index 2a9fd886f..dd77751c7 100644 --- a/packages/ai/src/providers/cursor.ts +++ b/packages/ai/src/providers/cursor.ts @@ -4769,15 +4769,18 @@ export async function buildGrpcRequest( conversationId: state.conversationId, }); - const replacementRequest = await options?.onPayload?.(runRequest, model); - if (replacementRequest !== undefined) runRequest = replacementRequest as typeof runRequest; - - // Tools are sent later via requestContext (exec handshake) - + // Apply customSystemPrompt BEFORE the hook so the onPayload replacement is the + // final word on the wire body — same contract as anthropic, where the hook runs + // right before serialization. An extension may inspect or drop it via the + // replacement it returns. if (options?.customSystemPrompt) { runRequest.customSystemPrompt = options.customSystemPrompt; } + // Tools are sent later via requestContext (exec handshake) + const replacementRequest = await options?.onPayload?.(runRequest, model); + if (replacementRequest !== undefined) runRequest = replacementRequest as typeof runRequest; + const clientMessage = create(AgentClientMessageSchema, { message: { case: "runRequest", value: runRequest }, }); diff --git a/packages/ai/test/cursor-on-payload.test.ts b/packages/ai/test/cursor-on-payload.test.ts index 766904d64..11b47fca7 100644 --- a/packages/ai/test/cursor-on-payload.test.ts +++ b/packages/ai/test/cursor-on-payload.test.ts @@ -64,4 +64,54 @@ describe("cursor onPayload replacement", () => { expect(message.case).toBe("runRequest"); expect(message.value.customSystemPrompt).toBeUndefined(); }); + + it("applies customSystemPrompt when onPayload returns undefined", async () => { + const { requestBytes } = await buildGrpcRequest( + model, + context, + { customSystemPrompt: "from-options", onPayload: async () => undefined }, + { conversationId: "conv-1", blobStore: new Map() }, + ); + + const message = decodeRunRequest(requestBytes); + expect(message.value.customSystemPrompt).toBe("from-options"); + }); + + it("lets the onPayload replacement drop customSystemPrompt (replacement is final)", async () => { + const { requestBytes } = await buildGrpcRequest( + model, + context, + { + customSystemPrompt: "from-options", + onPayload: async payload => { + const { customSystemPrompt: _dropped, ...rest } = payload as Record; + return rest; + }, + }, + { conversationId: "conv-1", blobStore: new Map() }, + ); + + const message = decodeRunRequest(requestBytes); + // The hook saw customSystemPrompt already applied (set before the hook) and + // returned a replacement that does not carry it — that replacement is final. + expect(message.value.customSystemPrompt).toBeUndefined(); + }); + + it("lets the onPayload replacement override customSystemPrompt", async () => { + const { requestBytes } = await buildGrpcRequest( + model, + context, + { + customSystemPrompt: "from-options", + onPayload: async payload => ({ + ...(payload as Record), + customSystemPrompt: "from-hook", + }), + }, + { conversationId: "conv-1", blobStore: new Map() }, + ); + + const message = decodeRunRequest(requestBytes); + expect(message.value.customSystemPrompt).toBe("from-hook"); + }); });