mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-17 16:12:21 -06:00
73d4c07bd5
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>
337 lines
11 KiB
TypeScript
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),
|
|
};
|
|
}
|
|
}
|