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..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({ @@ -581,6 +582,51 @@ 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" } }); + // 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, { + 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..a41fec76d8d57 100644 --- a/src/gateway/openresponses-http.ts +++ b/src/gateway/openresponses-http.ts @@ -184,6 +184,32 @@ function extractUsageFromResult(result: unknown): Usage { ); } +function extractFinalPayloadText(result: unknown): string | undefined { + const payloads = (result as { payloads?: Array<{ text?: string; isError?: boolean }> } | null) + ?.payloads; + if (!Array.isArray(payloads)) { + return 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 }; function resolveStopReasonAndPendingToolCalls(meta: unknown): { @@ -483,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); @@ -511,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, @@ -621,6 +640,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 +735,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(); @@ -719,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); @@ -789,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;