Skip to content

Commit 3181112

Browse files
committed
fix: preserve strict final-answer delivery phases (openclaw#59643) (thanks @ringlochid)
1 parent 98ce1c2 commit 3181112

2 files changed

Lines changed: 46 additions & 53 deletions

File tree

src/agents/openai-ws-stream.test.ts

Lines changed: 43 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import { createAssistantMessageEventStream } from "@mariozechner/pi-ai";
1212
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
1313
import type { ResponseObject } from "./openai-ws-connection.js";
14+
import { buildOpenAIWebSocketResponseCreatePayload } from "./openai-ws-request.js";
1415
import {
1516
__testing as openAIWsStreamTesting,
1617
buildAssistantMessageFromResponse,
@@ -123,6 +124,12 @@ const { MockManager } = vi.hoisted(() => {
123124
close(): void {
124125
this.closeCallCount++;
125126
this._connected = false;
127+
this._lastCloseInfo = {
128+
code: 1000,
129+
reason: "closed",
130+
retryable: false,
131+
};
132+
this.emit("close", 1000, "closed");
126133
}
127134

128135
// Test helper: simulate WebSocket connection drop mid-request
@@ -1470,19 +1477,35 @@ describe("createOpenAIWebSocketStreamFn", () => {
14701477
});
14711478
releaseWsSession("sess-1");
14721479
releaseWsSession("sess-2");
1480+
releaseWsSession("sess-boundary");
14731481
releaseWsSession("sess-fallback");
14741482
releaseWsSession("sess-boundary-http-fallback");
1483+
releaseWsSession("sess-full-context-replay");
14751484
releaseWsSession("sess-incremental");
14761485
releaseWsSession("sess-full");
1486+
releaseWsSession("sess-onpayload");
1487+
releaseWsSession("sess-onpayload-async");
14771488
releaseWsSession("sess-phase");
14781489
releaseWsSession("sess-phase-stream");
14791490
releaseWsSession("sess-phase-late-map");
1491+
releaseWsSession("sess-reason");
1492+
releaseWsSession("sess-reason-none");
14801493
releaseWsSession("sess-tools");
14811494
releaseWsSession("sess-store-default");
14821495
releaseWsSession("sess-store-compat");
1496+
releaseWsSession("sess-store-proxy");
14831497
releaseWsSession("sess-max-tokens-zero");
1498+
releaseWsSession("sess-runtime-fallback-nested");
14841499
releaseWsSession("sess-runtime-fallback");
1500+
releaseWsSession("sess-runtime-retry");
1501+
releaseWsSession("sess-send-fail-reset");
1502+
releaseWsSession("sess-temp");
1503+
releaseWsSession("sess-text-verbosity");
1504+
releaseWsSession("sess-text-verbosity-invalid");
1505+
releaseWsSession("sess-topp");
14851506
releaseWsSession("sess-turn-metadata-retry");
1507+
releaseWsSession("sess-warmup-disabled");
1508+
releaseWsSession("sess-warmup-enabled");
14861509
releaseWsSession("sess-degraded-cooldown");
14871510
releaseWsSession("sess-drop");
14881511
openAIWsStreamTesting.setWsDegradeCooldownMsForTest();
@@ -1501,8 +1524,10 @@ describe("createOpenAIWebSocketStreamFn", () => {
15011524

15021525
const manager = MockManager.lastInstance;
15031526
expect(manager?.connectCallCount).toBe(1);
1504-
// Consume stream to avoid dangling promise
1505-
void resolveStream(stream);
1527+
releaseWsSession("sess-1");
1528+
for await (const _ of await resolveStream(stream)) {
1529+
// consume
1530+
}
15061531
});
15071532

15081533
it("sends a response.create event on first turn (full context)", async () => {
@@ -1611,39 +1636,27 @@ describe("createOpenAIWebSocketStreamFn", () => {
16111636
expect(sent).not.toHaveProperty("store");
16121637
});
16131638

1614-
it("keeps store=false for proxied openai-responses routes when store is still supported", async () => {
1615-
releaseWsSession("sess-store-proxy");
1639+
it("keeps store=false for proxied openai-responses routes when store is still supported", () => {
16161640
const proxiedModel = {
16171641
...modelStub,
16181642
baseUrl: "https://proxy.example.com/v1",
16191643
};
1620-
const streamFn = createOpenAIWebSocketStreamFn("sk-test", "sess-store-proxy");
1621-
const stream = streamFn(
1622-
proxiedModel as Parameters<typeof streamFn>[0],
1623-
contextStub as Parameters<typeof streamFn>[1],
1624-
);
1625-
1626-
const completed = new Promise<void>((res, rej) => {
1627-
queueMicrotask(async () => {
1628-
try {
1629-
await new Promise((r) => setImmediate(r));
1630-
const manager = MockManager.lastInstance!;
1631-
manager.simulateEvent({
1632-
type: "response.completed",
1633-
response: makeResponseObject("resp_store_proxy", "ok"),
1634-
});
1635-
for await (const _ of await resolveStream(stream)) {
1636-
// consume
1637-
}
1638-
res();
1639-
} catch (e) {
1640-
rej(e);
1641-
}
1642-
});
1644+
const turnInput = planTurnInput({
1645+
context: contextStub as Parameters<typeof planTurnInput>[0]["context"],
1646+
model: proxiedModel as Parameters<typeof planTurnInput>[0]["model"],
1647+
previousResponseId: null,
1648+
lastContextLength: 0,
16431649
});
1644-
await completed;
1645-
1646-
const sent = MockManager.lastInstance!.sentEvents[0] as Record<string, unknown>;
1650+
const sent = buildOpenAIWebSocketResponseCreatePayload({
1651+
model: proxiedModel as Parameters<
1652+
typeof buildOpenAIWebSocketResponseCreatePayload
1653+
>[0]["model"],
1654+
context: contextStub as Parameters<
1655+
typeof buildOpenAIWebSocketResponseCreatePayload
1656+
>[0]["context"],
1657+
turnInput,
1658+
tools: [],
1659+
}) as Record<string, unknown>;
16471660
expect(sent.store).toBe(false);
16481661
});
16491662

src/agents/pi-embedded-subscribe.handlers.messages.ts

Lines changed: 3 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -43,27 +43,6 @@ const stripTrailingDirective = (text: string): string => {
4343
return text.slice(0, openIndex);
4444
};
4545

46-
type AssistantPhase = "commentary" | "final_answer";
47-
48-
const normalizeAssistantPhase = (value: unknown): AssistantPhase | undefined => {
49-
return value === "commentary" || value === "final_answer" ? value : undefined;
50-
};
51-
52-
const getAssistantTextSignaturePhase = (value: unknown): AssistantPhase | undefined => {
53-
if (typeof value !== "string" || value.trim().length === 0) {
54-
return undefined;
55-
}
56-
if (!value.startsWith("{")) {
57-
return undefined;
58-
}
59-
try {
60-
const parsed = JSON.parse(value) as { phase?: unknown; v?: unknown };
61-
return parsed.v === 1 ? normalizeAssistantPhase(parsed.phase) : undefined;
62-
} catch {
63-
return undefined;
64-
}
65-
};
66-
6746
const coerceText = (value: unknown): string => {
6847
if (typeof value === "string") {
6948
return value;
@@ -106,6 +85,7 @@ function resolveAssistantDeliveryPhase(
10685
if (!Array.isArray(message.content)) {
10786
return undefined;
10887
}
88+
const explicitStructuredPhases = new Set<AssistantDeliveryPhase>();
10989
for (const part of message.content) {
11090
if (!part || typeof part !== "object") {
11191
continue;
@@ -118,13 +98,13 @@ function resolveAssistantDeliveryPhase(
11898
const parsed = JSON.parse(block.textSignature) as { phase?: unknown };
11999
const phase = normalizeAssistantDeliveryPhase(parsed.phase);
120100
if (phase) {
121-
return phase;
101+
explicitStructuredPhases.add(phase);
122102
}
123103
} catch {
124104
continue;
125105
}
126106
}
127-
return undefined;
107+
return explicitStructuredPhases.size === 1 ? [...explicitStructuredPhases][0] : undefined;
128108
}
129109

130110
function shouldSuppressAssistantVisibleOutput(message: AgentMessage | undefined): boolean {

0 commit comments

Comments
 (0)