Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions src/gateway/open-responses.schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof ResponseCreatedEventSchema>
| z.infer<typeof ResponseInProgressEventSchema>
Expand All @@ -358,4 +367,5 @@ export type StreamingEvent =
| z.infer<typeof ContentPartAddedEventSchema>
| z.infer<typeof ContentPartDoneEventSchema>
| z.infer<typeof OutputTextDeltaEventSchema>
| z.infer<typeof OutputTextReplacedEventSchema>
| z.infer<typeof OutputTextDoneEventSchema>;
50 changes: 48 additions & 2 deletions src/gateway/openresponses-http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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, {
Expand All @@ -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({
Expand Down Expand Up @@ -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<string, unknown>).toMatchObject({
content: finalAnswer,
reason: "assistant_reconciliation",
});

const done = events.find((event) => event.event === "response.output_text.done");
expect(JSON.parse(done?.data ?? "{}") as Record<string, unknown>).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();
Expand Down
66 changes: 49 additions & 17 deletions src/gateway/openresponses-http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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): {
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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();

Expand All @@ -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);

Expand Down Expand Up @@ -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;
Expand Down
Loading