Files
openclaw/extensions/googlechat/src/api.ts
Ayaan Zaidi 73d4c07bd5 fix(delivery): record ambiguous final loss as durable notice debt (#121833)
A final reply whose platform send was accepted but whose response was lost
previously ended in silence. Custody that stays unknown after a claimed send
now records durable pendingDeliveryNotice debt; the next same-route turn
delivers one "could not confirm delivery" notice and acknowledges it into the
transcript. Restart recovery completes ambiguous sessions with the same debt
instead of a fire-and-forget notice; the debt survives reset and rollover, and
suppressed notice sends retain it instead of faking delivery. Permanent typed
no-send rejections settle as terminal suppression (no replay, no false
notice); retryable ones restore prepared custody for safe replay. Google Chat
media-only rejections use the typed no-send contract; Telegram native-command
replies join pending-final custody.

Fixes #80362

Co-authored-by: Ayaan Zaidi <hi@obviy.us>
2026-08-11 09:25:04 +00:00

337 lines
11 KiB
TypeScript

// Googlechat API module exposes the plugin public contract.
import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime";
import { redactToolPayloadText } from "openclaw/plugin-sdk/logging-core";
import {
MediaFetchError,
parseMediaContentLength,
readResponseTextSnippet,
} from "openclaw/plugin-sdk/media-runtime";
import { readProviderJsonResponse } from "openclaw/plugin-sdk/provider-http";
import { readResponseWithLimit } from "openclaw/plugin-sdk/response-limit-runtime";
import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime";
import type { ResolvedGoogleChatAccount } from "./accounts.js";
import { shouldSuppressGoogleChatManualExecApprovalFollowupText } from "./approval-card-actions.js";
import { getGoogleChatAccessToken } from "./auth.js";
import type { GoogleChatCardV2, GoogleChatSpace } from "./types.js";
const CHAT_API_BASE = "https://chat.googleapis.com/v1";
const GOOGLECHAT_API_TIMEOUT_MS = 30_000;
const GOOGLECHAT_MEDIA_TIMEOUT_GRACE_MS = 30_000;
const GOOGLECHAT_MEDIA_MIN_BYTES_PER_SECOND = 256 * 1024;
const GOOGLECHAT_MEDIA_MAX_TIMEOUT_MS = 15 * 60_000;
const GOOGLECHAT_RESPONSE_READ_IDLE_TIMEOUT_MS = 30_000;
const GOOGLECHAT_JSON_RESPONSE_MAX_BYTES = 16 * 1024 * 1024;
const GOOGLECHAT_ERROR_BODY_MAX_BYTES = 16 * 1024;
const GOOGLE_CHAT_DEFAULT_MEDIA_MAX_MB = 20;
const GOOGLE_CHAT_MEDIA_RESPONSE_MAX_BYTES = GOOGLE_CHAT_DEFAULT_MEDIA_MAX_MB * 1024 * 1024;
export class GoogleChatApiError extends Error {
constructor(
readonly status: number,
message: string,
) {
super(message);
this.name = "GoogleChatApiError";
}
}
function resolveGoogleChatMediaTimeoutMs(maxBytes?: number): number {
if (!maxBytes) {
return GOOGLECHAT_MEDIA_MAX_TIMEOUT_MS;
}
const transferMs = Math.ceil((maxBytes / GOOGLECHAT_MEDIA_MIN_BYTES_PER_SECOND) * 1000);
return Math.min(GOOGLECHAT_MEDIA_TIMEOUT_GRACE_MS + transferMs, GOOGLECHAT_MEDIA_MAX_TIMEOUT_MS);
}
async function readGoogleChatJsonResponse<T>(response: Response, label: string): Promise<T> {
return readProviderJsonResponse<T>(response, label, {
maxBytes: GOOGLECHAT_JSON_RESPONSE_MAX_BYTES,
chunkTimeoutMs: GOOGLECHAT_RESPONSE_READ_IDLE_TIMEOUT_MS,
onIdleTimeout: ({ chunkTimeoutMs }) =>
new Error(`${label}: response body stalled after ${chunkTimeoutMs}ms`),
});
}
async function readGoogleChatErrorResponse(response: Response, label: string): Promise<string> {
const text =
(await readResponseTextSnippet(response, {
maxBytes: GOOGLECHAT_ERROR_BODY_MAX_BYTES,
maxChars: GOOGLECHAT_ERROR_BODY_MAX_BYTES,
chunkTimeoutMs: GOOGLECHAT_RESPONSE_READ_IDLE_TIMEOUT_MS,
onIdleTimeout: ({ chunkTimeoutMs }) =>
new Error(`${label} error response stalled after ${chunkTimeoutMs}ms`),
})) ?? "";
// Remote API errors can reflect the request's Authorization header. Force
// tool-payload redaction before the text enters any surfaced error message.
return redactToolPayloadText(text);
}
const headersToObject = (headers?: HeadersInit): Record<string, string> =>
headers instanceof Headers
? Object.fromEntries(headers.entries())
: Array.isArray(headers)
? Object.fromEntries(headers)
: headers || {};
async function withGoogleChatResponse<T>(params: {
account: ResolvedGoogleChatAccount;
url: string;
init?: RequestInit;
auditContext: string;
errorPrefix?: string;
timeoutMs?: number;
handleResponse: (response: Response) => Promise<T>;
}): Promise<T> {
const {
account,
url,
init,
auditContext,
errorPrefix = "Google Chat API",
timeoutMs = GOOGLECHAT_API_TIMEOUT_MS,
handleResponse,
} = params;
const token = await getGoogleChatAccessToken(account);
const { response, release } = await fetchWithSsrFGuard({
url,
init: {
...init,
headers: {
...headersToObject(init?.headers),
Authorization: `Bearer ${token}`,
},
},
auditContext,
timeoutMs,
});
try {
if (!response.ok) {
const text = await readGoogleChatErrorResponse(response, errorPrefix);
throw new GoogleChatApiError(
response.status,
`${errorPrefix} ${response.status}: ${text || response.statusText}`,
);
}
return await handleResponse(response);
} finally {
// Status-only responses leave an unread body. Start cancellation before
// release; awaiting it can deadlock when debug capture tees the stream.
if (!response.bodyUsed) {
void response.body?.cancel().catch(() => undefined);
}
await release();
}
}
async function fetchJson<T>(
account: ResolvedGoogleChatAccount,
url: string,
init: RequestInit,
): Promise<T> {
return await withGoogleChatResponse({
account,
url,
init: {
...init,
headers: {
...headersToObject(init.headers),
"Content-Type": "application/json",
},
},
auditContext: "googlechat.api.json",
handleResponse: async (response) =>
await readGoogleChatJsonResponse<T>(response, "Google Chat API request failed"),
});
}
async function fetchOk(
account: ResolvedGoogleChatAccount,
url: string,
init: RequestInit,
): Promise<void> {
await withGoogleChatResponse({
account,
url,
init,
auditContext: "googlechat.api.ok",
handleResponse: async () => undefined,
});
}
async function fetchBuffer(
account: ResolvedGoogleChatAccount,
url: string,
init?: RequestInit,
options?: { maxBytes?: number },
): Promise<{ buffer: Buffer; contentType?: string }> {
return await withGoogleChatResponse({
account,
url,
init,
auditContext: "googlechat.api.buffer",
// Media gets transfer time proportional to its accepted size, while a silent
// response body is still bounded independently below.
timeoutMs: resolveGoogleChatMediaTimeoutMs(options?.maxBytes),
handleResponse: async (res) => {
const maxBytes = options?.maxBytes ?? GOOGLE_CHAT_MEDIA_RESPONSE_MAX_BYTES;
const lengthHeader = res.headers.get("content-length");
if (lengthHeader) {
const length = parseMediaContentLength(lengthHeader);
if (length !== null && length > maxBytes) {
throw new MediaFetchError(
"max_bytes",
`Google Chat media exceeds max bytes (${maxBytes})`,
);
}
}
const buffer = await readResponseWithLimit(res, maxBytes, {
chunkTimeoutMs: GOOGLECHAT_RESPONSE_READ_IDLE_TIMEOUT_MS,
onOverflow: () =>
new MediaFetchError("max_bytes", `Google Chat media exceeds max bytes (${maxBytes})`),
});
const contentType = res.headers.get("content-type") ?? undefined;
return { buffer, contentType };
},
});
}
/**
* A Google Chat `thread` must be a `spaces/{space}/threads/{thread}` resource
* name that belongs to the target space. Reply routing sometimes yields other
* shapes — a bare id, a `spaces/{space}/messages/{message}` name, or a thread
* from a different (or wrongly-cased) space — and passing any of those makes the
* Chat API reject the whole send with `400 INVALID_ARGUMENT`. Accept only a
* well-formed, same-space thread name; callers drop the rest so the message
* still delivers to the space (as a new thread) instead of failing outright.
*/
function isUsableGoogleChatThreadName(thread: string, space: string): boolean {
return /^spaces\/[^/]+\/threads\/[^/]+$/.test(thread) && thread.startsWith(`${space}/threads/`);
}
export async function sendGoogleChatMessage(params: {
account: ResolvedGoogleChatAccount;
space: string;
text?: string;
thread?: string;
cardsV2?: GoogleChatCardV2[];
}): Promise<{ messageName?: string; threadName?: string } | null> {
const { account, space, text, thread, cardsV2 } = params;
const usableThread = thread && isUsableGoogleChatThreadName(thread, space) ? thread : undefined;
if (
text &&
(!cardsV2 || cardsV2.length === 0) &&
shouldSuppressGoogleChatManualExecApprovalFollowupText(text)
) {
return null;
}
const body: Record<string, unknown> = {};
if (text) {
body.text = text;
}
if (cardsV2 && cardsV2.length > 0) {
body.cardsV2 = cardsV2;
}
if (usableThread) {
body.thread = { name: usableThread };
}
const urlObj = new URL(`${CHAT_API_BASE}/${space}/messages`);
if (usableThread) {
urlObj.searchParams.set("messageReplyOption", "REPLY_MESSAGE_FALLBACK_TO_NEW_THREAD");
}
const url = urlObj.toString();
const result = await fetchJson<{ name?: string; thread?: { name?: string } }>(account, url, {
method: "POST",
body: JSON.stringify(body),
});
return result ? { messageName: result.name, threadName: result.thread?.name } : null;
}
export async function updateGoogleChatMessage(params: {
account: ResolvedGoogleChatAccount;
messageName: string;
text?: string;
cardsV2?: GoogleChatCardV2[];
}): Promise<{ messageName?: string }> {
const { account, messageName, text, cardsV2 } = params;
const updateMask = [
...(text !== undefined ? ["text"] : []),
...(cardsV2 !== undefined ? ["cardsV2"] : []),
];
if (updateMask.length === 0) {
throw new Error("Google Chat message update requires text or cardsV2.");
}
const url = `${CHAT_API_BASE}/${messageName}?updateMask=${updateMask.join(",")}`;
const body: Record<string, unknown> = {};
if (text !== undefined) {
body.text = text;
}
if (cardsV2 !== undefined) {
body.cardsV2 = cardsV2;
}
const result = await fetchJson<{ name?: string }>(account, url, {
method: "PATCH",
body: JSON.stringify(body),
});
return { messageName: result.name };
}
export async function deleteGoogleChatMessage(params: {
account: ResolvedGoogleChatAccount;
messageName: string;
}): Promise<void> {
const { account, messageName } = params;
const url = `${CHAT_API_BASE}/${messageName}`;
await fetchOk(account, url, { method: "DELETE" });
}
export async function downloadGoogleChatMedia(params: {
account: ResolvedGoogleChatAccount;
resourceName: string;
maxBytes?: number;
}): Promise<{ buffer: Buffer; contentType?: string }> {
const { account, resourceName, maxBytes } = params;
const url = `${CHAT_API_BASE}/media/${resourceName}?alt=media`;
return await fetchBuffer(account, url, undefined, { maxBytes });
}
export async function findGoogleChatDirectMessage(params: {
account: ResolvedGoogleChatAccount;
userName: string;
}): Promise<GoogleChatSpace | null> {
const { account, userName } = params;
const url = new URL(`${CHAT_API_BASE}/spaces:findDirectMessage`);
url.searchParams.set("name", userName);
return await fetchJson<GoogleChatSpace>(account, url.toString(), {
method: "GET",
});
}
export async function getGoogleChatSpace(params: {
account: ResolvedGoogleChatAccount;
spaceName: string;
}): Promise<GoogleChatSpace> {
return await fetchJson<GoogleChatSpace>(params.account, `${CHAT_API_BASE}/${params.spaceName}`, {
method: "GET",
});
}
export async function probeGoogleChat(account: ResolvedGoogleChatAccount): Promise<{
ok: boolean;
status?: number;
error?: string;
}> {
try {
const url = new URL(`${CHAT_API_BASE}/spaces`);
url.searchParams.set("pageSize", "1");
await fetchJson<Record<string, unknown>>(account, url.toString(), {
method: "GET",
});
return { ok: true };
} catch (err) {
return {
ok: false,
error: formatErrorMessage(err),
};
}
}