From 46cf4194daa7d0557b3b05cfa974db9e47cada0d Mon Sep 17 00:00:00 2001 From: brianChen-ttc Date: Wed, 22 Jul 2026 16:57:33 +0800 Subject: [PATCH 1/2] Gateway: reconcile streamed output to final payload --- src/gateway/open-responses.schema.ts | 10 ++++++ src/gateway/openresponses-http.test.ts | 43 ++++++++++++++++++++++++++ src/gateway/openresponses-http.ts | 32 +++++++++++++++++++ 3 files changed, 85 insertions(+) diff --git a/src/gateway/open-responses.schema.ts b/src/gateway/open-responses.schema.ts index ca23f8de23599..f586a949f5444 100644 --- a/src/gateway/open-responses.schema.ts +++ b/src/gateway/open-responses.schema.ts @@ -348,6 +348,15 @@ export const OutputTextDoneEventSchema = z.object({ text: z.string(), }); +export const OutputTextReplacedEventSchema = z.object({ + type: z.literal("response.output_text.replaced"), + item_id: z.string(), + output_index: z.number().int().nonnegative(), + content_index: z.number().int().nonnegative(), + content: z.string(), + reason: z.string(), +}); + export type StreamingEvent = | z.infer | z.infer @@ -358,4 +367,5 @@ export type StreamingEvent = | z.infer | z.infer | z.infer + | z.infer | z.infer; diff --git a/src/gateway/openresponses-http.test.ts b/src/gateway/openresponses-http.test.ts index 3f6cb43917d37..001d191d8ddfb 100644 --- a/src/gateway/openresponses-http.test.ts +++ b/src/gateway/openresponses-http.test.ts @@ -581,6 +581,49 @@ describe("OpenResponses HTTP API (e2e)", () => { } }); + it("reconciles streamed tool narration to the final assistant payload", async () => { + const port = enabledPort; + const narration = "I'll help you create this document. Let me check the skill first."; + const finalAnswer = "文档已创建成功\n\nhttps://example.test/docx/final"; + + agentCommand.mockClear(); + agentCommand.mockImplementationOnce((async (opts: unknown) => { + const runId = (opts as { runId?: string }).runId ?? ""; + emitAgentEvent({ runId, stream: "assistant", data: { delta: narration } }); + emitAgentEvent({ runId, stream: "assistant", data: { delta: finalAnswer } }); + // The embedded runner can finish its lifecycle before the command promise + // resolves with the authoritative final payload. + emitAgentEvent({ runId, stream: "lifecycle", data: { phase: "end" } }); + return { payloads: [{ text: finalAnswer }] }; + }) as never); + + const response = await postResponses(port, { + stream: true, + model: "openclaw", + input: "create a document", + }); + expect(response.status).toBe(200); + + const events = parseSseEvents(await response.text()); + const replacement = events.find((event) => event.event === "response.output_text.replaced"); + expect(replacement).toBeDefined(); + expect(JSON.parse(replacement?.data ?? "{}") as Record).toMatchObject({ + content: finalAnswer, + reason: "assistant_reconciliation", + }); + + const done = events.find((event) => event.event === "response.output_text.done"); + expect(JSON.parse(done?.data ?? "{}") as Record).toMatchObject({ + text: finalAnswer, + }); + + const completed = events.find((event) => event.event === "response.completed"); + const completedPayload = JSON.parse(completed?.data ?? "{}") as { + response?: { output?: Array<{ content?: Array<{ text?: string }> }> }; + }; + expect(completedPayload.response?.output?.[0]?.content?.[0]?.text).toBe(finalAnswer); + }); + it("blocks unsafe URL-based file/image inputs", async () => { const port = enabledPort; agentCommand.mockClear(); diff --git a/src/gateway/openresponses-http.ts b/src/gateway/openresponses-http.ts index 97a5fee3c6632..336f26021d5b7 100644 --- a/src/gateway/openresponses-http.ts +++ b/src/gateway/openresponses-http.ts @@ -184,6 +184,18 @@ function extractUsageFromResult(result: unknown): Usage { ); } +function extractFinalPayloadText(result: unknown): string | undefined { + const payloads = (result as { payloads?: Array<{ text?: string }> } | null)?.payloads; + if (!Array.isArray(payloads)) { + return undefined; + } + const text = payloads + .map((payload) => (typeof payload.text === "string" ? payload.text : "")) + .filter(Boolean) + .join("\n\n"); + return text || undefined; +} + type PendingToolCall = { id: string; name: string; arguments: string }; function resolveStopReasonAndPendingToolCalls(meta: unknown): { @@ -621,6 +633,12 @@ export async function handleOpenResponsesHttpRequest( maybeFinalize(); }; + const replaceRequestedFinalText = (text: string) => { + if (finalizeRequested) { + finalizeRequested.text = text; + } + }; + // Send initial events const initialResponse = createResponseResource({ id: responseId, @@ -710,6 +728,20 @@ export async function handleOpenResponsesHttpRequest( deps, }); + const finalPayloadText = extractFinalPayloadText(result); + if (sawAssistantDelta && finalPayloadText && finalPayloadText !== accumulatedText) { + accumulatedText = finalPayloadText; + replaceRequestedFinalText(finalPayloadText); + writeSseEvent(res, { + type: "response.output_text.replaced", + item_id: outputItemId, + output_index: 0, + content_index: 0, + content: finalPayloadText, + reason: "assistant_reconciliation", + }); + } + finalUsage = extractUsageFromResult(result); maybeFinalize(); From 53d31962567d74d895bbfc8d24e25c4ba7c9558f Mon Sep 17 00:00:00 2001 From: brianChen-ttc Date: Wed, 22 Jul 2026 17:11:53 +0800 Subject: [PATCH 2/2] Gateway: select final assistant payload --- src/gateway/openresponses-http.test.ts | 9 +++-- src/gateway/openresponses-http.ts | 46 +++++++++++++------------- 2 files changed, 29 insertions(+), 26 deletions(-) diff --git a/src/gateway/openresponses-http.test.ts b/src/gateway/openresponses-http.test.ts index 001d191d8ddfb..b259e60e1b1bf 100644 --- a/src/gateway/openresponses-http.test.ts +++ b/src/gateway/openresponses-http.test.ts @@ -458,7 +458,7 @@ describe("OpenResponses HTTP API (e2e)", () => { expect(usageJson.usage).toEqual({ input_tokens: 3, output_tokens: 5, total_tokens: 10 }); await ensureResponseConsumed(resUsage); - mockAgentOnce([{ text: "hello" }]); + mockAgentOnce([{ text: "Let me check the skill." }, { text: "hello" }]); const resShape = await postResponses(port, { stream: false, model: "openclaw", @@ -542,7 +542,7 @@ describe("OpenResponses HTTP API (e2e)", () => { agentCommand.mockClear(); agentCommand.mockResolvedValueOnce({ - payloads: [{ text: "hello" }], + payloads: [{ text: "Let me check the skill." }, { text: "hello" }], } as never); const resFallback = await postResponses(port, { @@ -554,6 +554,7 @@ describe("OpenResponses HTTP API (e2e)", () => { const fallbackText = await resFallback.text(); expect(fallbackText).toContain("[DONE]"); expect(fallbackText).toContain("hello"); + expect(fallbackText).not.toContain("Let me check the skill."); agentCommand.mockClear(); agentCommand.mockResolvedValueOnce({ @@ -594,7 +595,9 @@ describe("OpenResponses HTTP API (e2e)", () => { // The embedded runner can finish its lifecycle before the command promise // resolves with the authoritative final payload. emitAgentEvent({ runId, stream: "lifecycle", data: { phase: "end" } }); - return { payloads: [{ text: finalAnswer }] }; + // The embedded runner returns one payload per assistant message, including + // tool-call preambles. The Responses adapter must select the final reply. + return { payloads: [{ text: narration }, { text: finalAnswer }] }; }) as never); const response = await postResponses(port, { diff --git a/src/gateway/openresponses-http.ts b/src/gateway/openresponses-http.ts index 336f26021d5b7..a41fec76d8d57 100644 --- a/src/gateway/openresponses-http.ts +++ b/src/gateway/openresponses-http.ts @@ -185,15 +185,29 @@ function extractUsageFromResult(result: unknown): Usage { } function extractFinalPayloadText(result: unknown): string | undefined { - const payloads = (result as { payloads?: Array<{ text?: string }> } | null)?.payloads; + const payloads = (result as { payloads?: Array<{ text?: string; isError?: boolean }> } | null) + ?.payloads; if (!Array.isArray(payloads)) { return undefined; } - const text = payloads - .map((payload) => (typeof payload.text === "string" ? payload.text : "")) - .filter(Boolean) - .join("\n\n"); - return text || undefined; + + // The embedded runner returns one payload per assistant message in the tool + // loop. Earlier payloads are tool-call preambles; the last non-error text is + // the user-facing final assistant reply. + for (let index = payloads.length - 1; index >= 0; index -= 1) { + const payload = payloads[index]; + if (payload?.isError !== true && typeof payload?.text === "string" && payload.text) { + return payload.text; + } + } + + for (let index = payloads.length - 1; index >= 0; index -= 1) { + const text = payloads[index]?.text; + if (typeof text === "string" && text) { + return text; + } + } + return undefined; } type PendingToolCall = { id: string; name: string; arguments: string }; @@ -495,7 +509,6 @@ export async function handleOpenResponsesHttpRequest( deps, }); - const payloads = (result as { payloads?: Array<{ text?: string }> } | null)?.payloads; const usage = extractUsageFromResult(result); const meta = (result as { meta?: unknown } | null)?.meta; const { stopReason, pendingToolCalls } = resolveStopReasonAndPendingToolCalls(meta); @@ -523,13 +536,7 @@ export async function handleOpenResponsesHttpRequest( return true; } - const content = - Array.isArray(payloads) && payloads.length > 0 - ? payloads - .map((p) => (typeof p.text === "string" ? p.text : "")) - .filter(Boolean) - .join("\n\n") - : "No response from OpenClaw."; + const content = extractFinalPayloadText(result) ?? "No response from OpenClaw."; const response = createResponseResource({ id: responseId, @@ -751,8 +758,7 @@ export async function handleOpenResponsesHttpRequest( // Fallback: if no streaming deltas were received, send the full response if (!sawAssistantDelta) { - const resultAny = result as { payloads?: Array<{ text?: string }>; meta?: unknown }; - const payloads = resultAny.payloads; + const resultAny = result as { meta?: unknown }; const meta = resultAny.meta; const { stopReason, pendingToolCalls } = resolveStopReasonAndPendingToolCalls(meta); @@ -821,13 +827,7 @@ export async function handleOpenResponsesHttpRequest( return; } - const content = - Array.isArray(payloads) && payloads.length > 0 - ? payloads - .map((p) => (typeof p.text === "string" ? p.text : "")) - .filter(Boolean) - .join("\n\n") - : "No response from OpenClaw."; + const content = extractFinalPayloadText(result) ?? "No response from OpenClaw."; accumulatedText = content; sawAssistantDelta = true;