fix(google): recover interrupted Interactions streams (#159198)

Emit the existing incomplete-stream contract so partial Interactions replies remain eligible for transient recovery. Centralize reader cancellation, lock release, and pending-work tracking in finally while preserving the original failure.

Regression proof on Blacksmith: six expected failures before the repair; all 101 focused and owner tests pass afterward. The full changed gate passes, and independent Codex review found no actionable P0-P2 findings.

Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
Peter Steinberger
2026-09-26 15:18:11 -07:00
committed by GitHub
parent 1a2cda64d1
commit 727e4ed263
4 changed files with 152 additions and 66 deletions
+3
View File
@@ -511,6 +511,9 @@ roundtrip; pass `--openai-audio-cycles 3` for a short repeated lifecycle soak.
each request. Explicit Gemini `cachedContent` handles are not supported on
this route; use `google-generative-ai` for that feature.
Interrupted text replies use the normal transient-error retry and failover
policy. Malformed completed tool-call arguments remain rejected.
</Accordion>
<Accordion title="Direct Gemini cache reuse">
@@ -14,7 +14,7 @@ import { calculateCost } from "../model-utils.js";
import { buildGuardedModelFetch } from "../transports/host-policy.js";
import { parseJsonPreservingUnsafeIntegers } from "../transports/json-unsafe-integers.js";
import {
assignTransportErrorDetails,
failTransportStream,
notifyProviderHttpResponse,
notifyProviderStreamOpened,
parseTerminalToolCallArguments,
@@ -73,6 +73,7 @@ export async function runGoogleInteractionsLifecycle<T extends GoogleApiType>(pa
apiKey?: string;
}): Promise<void> {
const { stream, model, output, options, context, nextToolCallId } = params;
let reader: ReadableStreamDefaultReader<Uint8Array> | undefined;
try {
const host = getAiTransportHost();
@@ -149,12 +150,10 @@ export async function runGoogleInteractionsLifecycle<T extends GoogleApiType>(pa
throw new Error("Google Interactions API returned empty response body");
}
const reader = response.body.getReader();
reader = response.body.getReader();
await notifyProviderStreamOpened({
options,
cancelStream: async () => {
await reader.cancel();
},
cancelStream: () => reader?.cancel(),
});
stream.push({ type: "start", partial: output });
const decoder = new TextDecoder();
@@ -277,7 +276,6 @@ export async function runGoogleInteractionsLifecycle<T extends GoogleApiType>(pa
code: readStringField(providerError, "code"),
type: "google_interactions_stream_error",
});
await reader.cancel();
throw error;
} else if (eventType === "step.delta") {
const delta = asOptionalRecord(event.delta);
@@ -493,14 +491,13 @@ export async function runGoogleInteractionsLifecycle<T extends GoogleApiType>(pa
}
}
if (streamDone) {
void reader.cancel().catch(() => {});
}
endCurrentBlock();
if (!sawCompletion) {
throw new Error("Google Interactions stream ended before interaction.completed");
throw Object.assign(new Error("Google Interactions stream ended before a terminal event"), {
code: "STREAM_INCOMPLETE",
type: "google_incomplete_stream",
});
}
if (latestThoughtSignature) {
@@ -528,12 +525,13 @@ export async function runGoogleInteractionsLifecycle<T extends GoogleApiType>(pa
stream.end();
} catch (error) {
const failure = options?.signal?.aborted ? transportAbortError(options.signal) : error;
assignTransportErrorDetails(output, failure, options?.signal);
stream.push({
type: "error",
reason: output.stopReason === "aborted" ? "aborted" : "error",
error: output,
});
stream.end();
failTransportStream({ stream, output, signal: options?.signal, error: failure });
} finally {
if (reader) {
// Track cleanup without delaying terminal delivery or replacing the original failure.
const cancellation = reader.cancel().catch(() => undefined);
reader.releaseLock();
getAiTransportHost().observePendingProviderWork?.(cancellation);
}
}
}
@@ -1,4 +1,5 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../../../test/helpers/promise.js";
import { configureAiTransportHost } from "../host.js";
import type { AssistantMessage, Context, Model, ToolCall } from "../types.js";
import { streamGoogleInteractions, streamSimpleGoogleInteractions } from "./google-interactions.js";
@@ -96,6 +97,7 @@ describe("google-interactions provider", () => {
}
expect(cancelCalled).toBe(true);
expect(stream.locked).toBe(false);
const doneEvent = events.find(
(e): e is { type: "done"; message: { api: string; content: unknown[] } } =>
Boolean(e && typeof e === "object" && (e as { type: string }).type === "done"),
@@ -417,37 +419,67 @@ describe("google-interactions provider", () => {
]);
});
it("rejects malformed streamed tool call arguments", async () => {
const encoder = new TextEncoder();
const ssePayload = [
'data: {"event_type":"step.start","step":{"type":"function_call","id":"call_exec_1","name":"exec","arguments":{}}}\n\n',
'data: {"event_type":"step.delta","delta":{"type":"arguments_delta","arguments":"{\\"command\\":\\"ls"}}\n\n',
'data: {"event_type":"step.stop"}\n\n',
completedSse({ status: "requires_action" }),
"data: [DONE]\n\n",
].join("");
it.each(["resolved", "rejected", "pending"])(
"retires malformed tool streams when cancellation is %s",
async (cancellationState) => {
const encoder = new TextEncoder();
const ssePayload = [
'data: {"event_type":"step.start","step":{"type":"function_call","id":"call_exec_1","name":"exec","arguments":{}}}\n\n',
'data: {"event_type":"step.delta","delta":{"type":"arguments_delta","arguments":"{\\"command\\":\\"ls"}}\n\n',
'data: {"event_type":"step.stop"}\n\n',
].join("");
const cancellation = createDeferred();
const pendingWork: Promise<unknown>[] = [];
configureAiTransportHost({
observePendingProviderWork: (pending) => {
pendingWork.push(pending);
},
});
const cancel = vi.fn(() => {
if (cancellationState === "resolved") {
cancellation.resolve();
} else if (cancellationState === "rejected") {
cancellation.reject(new Error("cancel failed"));
}
return cancellation.promise;
});
const body = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(encoder.encode(ssePayload));
},
cancel,
});
vi.stubGlobal(
"fetch",
vi.fn(
async () =>
new Response(encoder.encode(ssePayload), {
status: 200,
headers: { "Content-Type": "text/event-stream" },
}),
),
);
vi.stubGlobal(
"fetch",
vi.fn(
async () =>
new Response(body, {
status: 200,
headers: { "Content-Type": "text/event-stream" },
}),
),
);
const result = await streamGoogleInteractions(makeInteractionsModel(), basicContext, {
apiKey: "test-key",
}).result();
try {
const result = await streamGoogleInteractions(makeInteractionsModel(), basicContext, {
apiKey: "test-key",
}).result();
expect(result).toMatchObject({
stopReason: "error",
errorCode: "malformed_tool_call_arguments",
errorMessage: "Provider completed tool call with malformed JSON arguments",
});
});
expect(result).toMatchObject({
stopReason: "error",
errorCode: "malformed_tool_call_arguments",
errorMessage: "Provider completed tool call with malformed JSON arguments",
});
expect(cancel).toHaveBeenCalledOnce();
expect(body.locked).toBe(false);
expect(pendingWork).toHaveLength(1);
} finally {
cancellation.resolve();
await Promise.all(pendingWork);
}
},
);
it("resolves API-key and custom-header sentinels before guarded egress", async () => {
const sentinel = "oc-sent-v2.AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA.end";
@@ -526,26 +558,6 @@ describe("google-interactions provider", () => {
expect(result.errorCode).toBe("gateway_timeout");
});
it("rejects a stream that ends before interaction.completed", async () => {
vi.stubGlobal(
"fetch",
vi.fn(
async () =>
new Response(new TextEncoder().encode("data: [DONE]\n\n"), {
status: 200,
headers: { "Content-Type": "text/event-stream" },
}),
),
);
const result = await streamGoogleInteractions(makeInteractionsModel(), basicContext, {
apiKey: "test-api-key",
}).result();
expect(result.stopReason).toBe("error");
expect(result.errorMessage).toContain("before interaction.completed");
});
it("maps cached, thought, and tool-use tokens into canonical usage", async () => {
vi.stubGlobal(
"fetch",
@@ -0,0 +1,73 @@
import type { Model } from "@openclaw/ai/types";
import { afterEach, describe, expect, it, vi } from "vitest";
import { configureAiTransportHost } from "../../packages/ai/src/host.js";
import { streamSimpleGoogleInteractions } from "../../packages/ai/src/providers/google-interactions.js";
import { classifyFailoverSignalCore } from "./failover/classify-core.js";
import { shouldRetryFailoverSignal } from "./failover/retry-evidence.js";
describe("Google Interactions stream recovery", () => {
afterEach(() => {
vi.unstubAllGlobals();
configureAiTransportHost({});
});
it.each(["EOF", "[DONE]"])(
"keeps partial replies retryable after %s without completion",
async (ending) => {
const model: Model<"google-interactions"> = {
id: "gemini-3-flash-preview",
name: "Gemini 3 Flash",
api: "google-interactions",
provider: "google",
baseUrl: "https://generativelanguage.googleapis.com/v1beta",
reasoning: false,
input: ["text"],
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
contextWindow: 128_000,
maxTokens: 8_192,
};
const payload =
'data: {"event_type":"step.delta","delta":{"type":"text","text":"Partial reply"}}\n\n' +
(ending === "[DONE]" ? "data: [DONE]\n\n" : "");
configureAiTransportHost({});
vi.stubGlobal(
"fetch",
vi.fn(
async () =>
new Response(payload, {
headers: { "Content-Type": "text/event-stream" },
}),
),
);
const result = await streamSimpleGoogleInteractions(
model,
{
messages: [{ role: "user", content: "Hello", timestamp: 0 }],
},
{ apiKey: "test-key" },
).result();
expect(result.stopReason).toBe("error");
expect(result.content).toEqual([{ type: "text", text: "Partial reply" }]);
const signal = {
provider: result.provider,
message: result.errorMessage,
code: result.errorCode,
errorType: result.errorType,
};
const classification = classifyFailoverSignalCore(signal);
expect({
classification,
retryable: shouldRetryFailoverSignal({ classification, signal }),
}).toEqual({
classification: { kind: "reason", reason: "timeout" },
retryable: true,
});
expect(result).toMatchObject({
errorCode: "STREAM_INCOMPLETE",
errorType: "google_incomplete_stream",
});
},
);
});