diff --git a/docs/.generated/plugin-sdk-api-baseline.sha256 b/docs/.generated/plugin-sdk-api-baseline.sha256 index 5df6cba9dab0..435e79fbe025 100644 --- a/docs/.generated/plugin-sdk-api-baseline.sha256 +++ b/docs/.generated/plugin-sdk-api-baseline.sha256 @@ -1,2 +1,2 @@ -73531cdb2e35f3b37f7e1f1838792920457a5e950aed4a6232701ce47dcbf4b5 plugin-sdk-api-baseline.json -2f465b3acc86fe710a9716215034e812fc870cdd6046afce257f002b88ad914d plugin-sdk-api-baseline.jsonl +53c37adb223c5a6b497619fa688c04ffc024755a9c69e79bbc347c5d883ebe71 plugin-sdk-api-baseline.json +b537576389e0ee00ef83b6921e66cef9e1a89c0d1e14eaece14dccb9e116bdcc plugin-sdk-api-baseline.jsonl diff --git a/docs/plugins/sdk-channel-outbound.md b/docs/plugins/sdk-channel-outbound.md index 104b908b15dc..74048db450cb 100644 --- a/docs/plugins/sdk-channel-outbound.md +++ b/docs/plugins/sdk-channel-outbound.md @@ -75,6 +75,15 @@ Receipts without platform message identifiers are local receipt metadata only. Channels with read receipts or device-delivery state should track those facts through a separate channel-specific path. +If a channel adapter can prove that retrying a failure cannot duplicate a +recipient-visible send and no finalization-capable call began, throw +`new PlatformMessageNotDispatchedError("...", { cause: error })` from +`openclaw/plugin-sdk/error-runtime`. Core can then clear stale send-attempt +evidence and safely retry the queued intent. Only the adapter that owns the +final dispatch boundary may make this assertion. Never use the marker after a +finalization/send call begins or returns an ambiguous result; false marking can +duplicate messages. + ## Existing outbound adapters If the channel already has a compatible `outbound` adapter, derive the diff --git a/docs/plugins/sdk-migration.md b/docs/plugins/sdk-migration.md index 8c397f9f00b2..605077f490a6 100644 --- a/docs/plugins/sdk-migration.md +++ b/docs/plugins/sdk-migration.md @@ -444,7 +444,7 @@ SDK. | `plugin-sdk/exec-approvals-runtime` | Exec approval policy helpers | `loadExecApprovals`, `resolveExecApprovalsFromFile`, `ExecApprovalsFile` | | `plugin-sdk/collection-runtime` | Bounded cache helpers | `pruneMapToMaxSize` | | `plugin-sdk/diagnostic-runtime` | Diagnostic gating helpers | `isDiagnosticFlagEnabled`, `isDiagnosticsEnabled` | - | `plugin-sdk/error-runtime` | Error formatting helpers | `formatUncaughtError`, `isApprovalNotFoundError`, error graph helpers | + | `plugin-sdk/error-runtime` | Error helpers | `formatUncaughtError`, `isApprovalNotFoundError`, error graph helpers, `PlatformMessageNotDispatchedError` | | `plugin-sdk/fetch-runtime` | Wrapped fetch/proxy helpers | `resolveFetch`, proxy helpers, EnvHttpProxyAgent option helpers | | `plugin-sdk/host-runtime` | Host normalization helpers | `normalizeHostname`, `normalizeScpRemoteHost` | | `plugin-sdk/retry-runtime` | Retry helpers | `RetryConfig`, `retryAsync`, policy runners | diff --git a/docs/plugins/sdk-subpaths.md b/docs/plugins/sdk-subpaths.md index 44b7e152c58d..7e081ee3157f 100644 --- a/docs/plugins/sdk-subpaths.md +++ b/docs/plugins/sdk-subpaths.md @@ -326,7 +326,7 @@ usage endpoint failed or returned no usable usage data. | `plugin-sdk/infra-runtime` | Deprecated compatibility shim; use the focused runtime subpaths above | | `plugin-sdk/collection-runtime` | Small bounded cache helpers | | `plugin-sdk/diagnostic-runtime` | Diagnostic flag, event, and trace-context helpers | - | `plugin-sdk/error-runtime` | Error graph, formatting, shared error classification helpers, `isApprovalNotFoundError` | + | `plugin-sdk/error-runtime` | Error graph, formatting, shared error classification helpers, `PlatformMessageNotDispatchedError`, `isApprovalNotFoundError` | | `plugin-sdk/fetch-runtime` | Wrapped fetch, proxy, EnvHttpProxyAgent option, and pinned lookup helpers | | `plugin-sdk/runtime-fetch` | Dispatcher-aware runtime fetch without proxy/guarded-fetch imports | | `plugin-sdk/inline-image-data-url-runtime` | Inline image data URL sanitizer and signature sniffing helpers without the broad media runtime surface | diff --git a/extensions/slack/src/channel.message-adapter.test.ts b/extensions/slack/src/channel.message-adapter.test.ts index c3ef937ec92f..db54f1326d90 100644 --- a/extensions/slack/src/channel.message-adapter.test.ts +++ b/extensions/slack/src/channel.message-adapter.test.ts @@ -106,6 +106,7 @@ describe("slack channel message adapter", () => { const proveMedia = async () => { sendSlack.mockClear(); + const onPlatformSendDispatch = vi.fn(); const result = await sendMedia({ cfg, to: "C123", @@ -114,6 +115,7 @@ describe("slack channel message adapter", () => { mediaLocalRoots: ["/tmp/media"], accountId: "default", deliveryQueueId: "queue-1", + onPlatformSendDispatch, deps: { sendSlack }, }); const [to, text, options] = expectLastSendSlackCall(); @@ -123,6 +125,7 @@ describe("slack channel message adapter", () => { expect(options.mediaUrl).toBe("https://example.com/a.png"); expect(options.mediaLocalRoots).toEqual(["/tmp/media"]); expect(options.deliveryQueueId).toBeUndefined(); + expect(options.onPlatformSendDispatch).toBe(onPlatformSendDispatch); expect(result.receipt.parts[0]?.kind).toBe("media"); }; diff --git a/extensions/slack/src/channel.test.ts b/extensions/slack/src/channel.test.ts index 127be9ca715b..2cb17b9fe635 100644 --- a/extensions/slack/src/channel.test.ts +++ b/extensions/slack/src/channel.test.ts @@ -1162,8 +1162,12 @@ describe("slackPlugin outbound", () => { it("forwards mediaLocalRoots for sendMedia", async () => { const sendSlack = vi.fn().mockResolvedValue({ messageId: "m-media-local" }); - const sendMedia = requireSlackSendMedia(); + const sendMedia = slackOutbound.sendMedia; + if (!sendMedia) { + throw new Error("slack direct outbound.sendMedia unavailable"); + } const mediaLocalRoots = ["/tmp/workspace"]; + const onPlatformSendDispatch = vi.fn(); const result = await sendMedia({ cfg, @@ -1172,6 +1176,7 @@ describe("slackPlugin outbound", () => { mediaUrl: "/tmp/workspace/image.png", mediaLocalRoots, accountId: "default", + onPlatformSendDispatch, deps: { sendSlack }, }); @@ -1180,6 +1185,7 @@ describe("slackPlugin outbound", () => { expectRecordFields(requireMockCallArg(sendSlack, 0, 2), "send options", { mediaUrl: "/tmp/workspace/image.png", mediaLocalRoots, + onPlatformSendDispatch, }); expect(result).toEqual({ channel: "slack", messageId: "m-media-local" }); }); diff --git a/extensions/slack/src/channel.ts b/extensions/slack/src/channel.ts index 8016b1c32dc4..49986387e1cb 100644 --- a/extensions/slack/src/channel.ts +++ b/extensions/slack/src/channel.ts @@ -564,6 +564,7 @@ const slackChannelOutbound: ChannelOutboundAdapter = { deps: ctx.deps, send, tokenOverride, + onPlatformSendDispatch: ctx.onPlatformSendDispatch, }), }); }, diff --git a/extensions/slack/src/client-delivery.ts b/extensions/slack/src/client-delivery.ts index 9546490491b7..90343a00e660 100644 --- a/extensions/slack/src/client-delivery.ts +++ b/extensions/slack/src/client-delivery.ts @@ -1,19 +1,25 @@ // Slack plugin module owns WebClient-scoped message and file delivery primitives. import type { MessageMetadata } from "@slack/types"; import type { Block, KnownBlock, WebClient } from "@slack/web-api"; +import { + extractErrorCode, + PlatformMessageNotDispatchedError, + readErrorName, +} from "openclaw/plugin-sdk/error-runtime"; +import { buildTimeoutAbortSignal } from "openclaw/plugin-sdk/extension-shared"; import { withTrustedEnvProxyGuardedFetchMode } from "openclaw/plugin-sdk/fetch-runtime"; import { logVerbose } from "openclaw/plugin-sdk/runtime-env"; -import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; -import { - postSlackMessageWithIdentityFallback, - type SlackPostMessageIdentity, -} from "./post-message-identity.js"; +import { fetchWithSsrFGuard, type SsrFPolicy } from "openclaw/plugin-sdk/ssrf-runtime"; import { appendSlackDataVisualizationFallbackText, hasSlackDataVisualizationBlock, isSlackInvalidBlocksError, } from "./data-visualization.js"; import { SLACK_TEXT_LIMIT } from "./limits.js"; +import { + postSlackMessageWithIdentityFallback, + type SlackPostMessageIdentity, +} from "./post-message-identity.js"; import { buildSlackPostMessagePayload, type SlackPostMessagePayload, @@ -22,10 +28,19 @@ import { import { loadOutboundMediaFromUrl } from "./runtime-api.js"; import { truncateSlackText } from "./truncate.js"; -const SLACK_UPLOAD_SSRF_POLICY = { - allowedHostnames: ["*.slack.com", "*.slack-edge.com", "*.slack-files.com"], +const SLACK_COMMERCIAL_API_HOSTNAME = "slack.com"; +const SLACK_COMMERCIAL_UPLOAD_HOSTNAME = "files.slack.com"; +const SLACK_GOV_API_HOSTNAME = "slack-gov.com"; +const SLACK_GOV_UPLOAD_HOSTNAME = "files.slack-gov.com"; +const SLACK_COMMERCIAL_UPLOAD_SSRF_POLICY = { + hostnameAllowlist: [SLACK_COMMERCIAL_UPLOAD_HOSTNAME], allowRfc2544BenchmarkRange: true, -}; +} satisfies SsrFPolicy; +const SLACK_GOV_UPLOAD_SSRF_POLICY = { + hostnameAllowlist: [SLACK_GOV_UPLOAD_HOSTNAME], + allowRfc2544BenchmarkRange: true, +} satisfies SsrFPolicy; +const SLACK_UPLOAD_POST_TIMEOUT_MS = 120_000; const SLACK_DNS_RETRY_CODES = new Set(["EAI_AGAIN", "ENOTFOUND", "UND_ERR_DNS_RESOLVE_FAILED"]); const SLACK_DNS_RETRY_ATTEMPTS = 2; const SLACK_DNS_RETRY_BASE_DELAY_MS = 250; @@ -74,6 +89,115 @@ function delaySlackDnsRetry(attempt: number): Promise { }); } +function resolveSlackUploadTimeoutLogUrl(url: string): string | undefined { + // Slack puts the upload capability in the URL path. Timeout diagnostics may + // name the origin, but must not retain that capability-bearing path. + try { + return new URL(url).origin; + } catch { + return undefined; + } +} + +function buildSlackUploadFailureCause(error: unknown): Error { + const httpStatus = + error instanceof Error + ? /^Failed to upload file: HTTP (\d{3})$/u.exec(error.message)?.[1] + : undefined; + const cause = new Error( + httpStatus + ? `Slack external upload returned HTTP ${httpStatus}` + : "Slack external upload transfer failed", + ); + cause.name = readErrorName(error) || cause.name; + const code = extractErrorCode(error) ?? (httpStatus ? `HTTP_${httpStatus}` : undefined); + if (code) { + (cause as NodeJS.ErrnoException).code = code; + } + return cause; +} + +function parseSlackUploadHttpUrl(value: string, label: string): URL { + try { + const parsed = new URL(value); + if (parsed.protocol === "http:" || parsed.protocol === "https:") { + return parsed; + } + } catch { + // Fall through to the same capability-safe error below. + } + throw new Error(`${label} must use a valid HTTP or HTTPS URL`); +} + +function normalizeSlackHostname(hostname: string): string { + return hostname.trim().toLowerCase().replace(/\.$/, ""); +} + +function resolveSlackOwnedUploadPolicy(url: URL): SsrFPolicy | undefined { + if (url.protocol !== "https:") { + return undefined; + } + switch (normalizeSlackHostname(url.hostname)) { + case SLACK_COMMERCIAL_UPLOAD_HOSTNAME: + return SLACK_COMMERCIAL_UPLOAD_SSRF_POLICY; + case SLACK_GOV_UPLOAD_HOSTNAME: + return SLACK_GOV_UPLOAD_SSRF_POLICY; + default: + return undefined; + } +} + +function resolveOfficialSlackApiUploadPolicy(url: URL): SsrFPolicy | undefined { + if (url.protocol !== "https:" || url.port) { + return undefined; + } + switch (normalizeSlackHostname(url.hostname)) { + case SLACK_COMMERCIAL_API_HOSTNAME: + return SLACK_COMMERCIAL_UPLOAD_SSRF_POLICY; + case SLACK_GOV_API_HOSTNAME: + return SLACK_GOV_UPLOAD_SSRF_POLICY; + default: + return undefined; + } +} + +function normalizeSlackOrigin(url: URL): string { + const port = url.port ? `:${url.port}` : ""; + return `${url.protocol}//${normalizeSlackHostname(url.hostname)}${port}`; +} + +function resolveSlackUploadTransportPolicy(params: { uploadUrl: string; slackApiUrl?: string }): { + requireHttps: boolean; + policy: SsrFPolicy; +} { + if (!params.slackApiUrl) { + return { requireHttps: true, policy: SLACK_COMMERCIAL_UPLOAD_SSRF_POLICY }; + } + const apiUrl = parseSlackUploadHttpUrl(params.slackApiUrl, "Configured Slack API URL"); + const officialApiPolicy = resolveOfficialSlackApiUploadPolicy(apiUrl); + if (officialApiPolicy) { + return { requireHttps: true, policy: officialApiPolicy }; + } + const uploadUrl = parseSlackUploadHttpUrl(params.uploadUrl, "Slack external upload URL"); + const slackOwnedUploadPolicy = resolveSlackOwnedUploadPolicy(uploadUrl); + if (slackOwnedUploadPolicy) { + return { requireHttps: true, policy: slackOwnedUploadPolicy }; + } + // Default Slack capabilities stay Slack-hosted. An operator-selected API + // root may additionally return upload capabilities on its exact origin. + if (normalizeSlackOrigin(uploadUrl) !== normalizeSlackOrigin(apiUrl)) { + throw new Error("Slack external upload URL must match the configured Slack API origin"); + } + return { + requireHttps: apiUrl.protocol === "https:", + policy: { + hostnameAllowlist: [uploadUrl.hostname], + allowedOrigins: [uploadUrl.origin], + allowRfc2544BenchmarkRange: true, + }, + }; +} + export async function withSlackDnsRequestRetry( operation: string, fn: () => Promise, @@ -146,6 +270,7 @@ export async function postSlackMessageBestEffort(params: { export async function uploadSlackFile(params: { client: WebClient; + completionClient?: WebClient; channelId: string; mediaUrl: string; mediaAccess?: { @@ -180,29 +305,63 @@ export async function uploadSlackFile(params: { throw new Error(`Failed to get upload URL: ${uploadUrlResp.error ?? "unknown error"}`); } const uploadFileId = uploadUrlResp.file_id; - const { response: uploadResp, release } = await fetchWithSsrFGuard( - withTrustedEnvProxyGuardedFetchMode({ - url: uploadUrlResp.upload_url, - init: { - method: "POST", - ...(contentType ? { headers: { "Content-Type": contentType } } : {}), - body: new Uint8Array(buffer) as BodyInit, - }, - policy: SLACK_UPLOAD_SSRF_POLICY, - auditContext: params.auditContext ?? "slack-upload-file", - }), - ); + const uploadTransport = resolveSlackUploadTransportPolicy({ + uploadUrl: uploadUrlResp.upload_url, + slackApiUrl: params.client.slackApiUrl, + }); + // Bound only the byte transfer. Completion may commit server-side before its + // response arrives, so timing it out would create unsafe unknown-send retries. + const { signal: uploadTimeoutSignal, cleanup: cleanupUploadTimeout } = buildTimeoutAbortSignal({ + timeoutMs: SLACK_UPLOAD_POST_TIMEOUT_MS, + operation: "slack-upload-file", + url: resolveSlackUploadTimeoutLogUrl(uploadUrlResp.upload_url), + }); try { - if (!uploadResp.ok) { - throw new Error(`Failed to upload file: HTTP ${uploadResp.status}`); + const { response: uploadResp, release } = await fetchWithSsrFGuard( + withTrustedEnvProxyGuardedFetchMode({ + url: uploadUrlResp.upload_url, + init: { + method: "POST", + ...(contentType ? { headers: { "Content-Type": contentType } } : {}), + body: new Uint8Array(buffer) as BodyInit, + }, + signal: uploadTimeoutSignal, + requireHttps: uploadTransport.requireHttps, + policy: uploadTransport.policy, + capture: false, + auditContext: params.auditContext ?? "slack-upload-file", + }), + ); + try { + if (uploadResp.status !== 200) { + throw new Error(`Failed to upload file: HTTP ${uploadResp.status}`); + } + } finally { + // Slack's status is the upload result; discard any response body so its + // keep-alive or proxy socket cannot outlive this transfer. + await uploadResp.body?.cancel().catch(() => undefined); + await release(); } + } catch (error) { + // Slack discards raw uploads that never reach completion. Every failure in + // this transfer block is therefore safe to retry; finalization stays unmarked. + const outcome = uploadTimeoutSignal?.aborted ? "timed out" : "failed"; + throw new PlatformMessageNotDispatchedError( + `Slack external upload ${outcome} before completion dispatch`, + // Upload capabilities live in the URL path. Preserve only safe transport + // metadata so flattened cause logging cannot disclose that path. + { cause: buildSlackUploadFailureCause(error) }, + ); } finally { - await release(); + cleanupUploadTimeout(); } await params.onPlatformSendDispatch?.(); + // Slack allows this finalize call only once. Keep only the pre-connect DNS + // retry; a timeout or broader retry would create an unknown-send state. + const completionClient = params.completionClient ?? params.client; const completeResp = await withSlackDnsRequestRetry("files.completeUploadExternal", () => - params.client.files.completeUploadExternal({ + completionClient.files.completeUploadExternal({ files: [{ id: uploadFileId, title: uploadTitle }], channel_id: params.channelId, ...(params.caption ? { initial_comment: params.caption } : {}), diff --git a/extensions/slack/src/client.ts b/extensions/slack/src/client.ts index c09e2cae9a99..a4fbda539607 100644 --- a/extensions/slack/src/client.ts +++ b/extensions/slack/src/client.ts @@ -1,10 +1,18 @@ // Slack plugin module implements client behavior. import { createHash } from "node:crypto"; import { type WebClientOptions, WebClient } from "@slack/web-api"; -import { resolveSlackWebClientOptions, resolveSlackWriteClientOptions } from "./client-options.js"; +import { + resolveSlackWebClientOptions, + resolveSlackWriteClientOptions, + SLACK_WRITE_RETRY_OPTIONS, +} from "./client-options.js"; const SLACK_WRITE_CLIENT_CACHE_MAX = 32; const slackWriteClientCache = new Map(); +let slackListenerUploadCompletionClientCache = new WeakMap< + WebClient, + { teamId: string; client: WebClient } +>(); type SlackWriteClientCacheOptions = Pick; @@ -55,6 +63,45 @@ export function getSlackWriteClient( return client; } +export function getSlackListenerUploadCompletionClient(params: { + listenerClient: WebClient; + teamId: string; + clientOptions?: WebClientOptions; +}): WebClient | undefined { + const token = params.listenerClient.token?.trim(); + const teamId = params.teamId.trim().toUpperCase(); + if (!token || !teamId) { + return undefined; + } + const cached = slackListenerUploadCompletionClientCache.get(params.listenerClient); + if (cached) { + // Bolt pools listener clients by authorized team. Reusing one for a + // different team is invalid scope, not another completion-client key. + return cached.teamId === teamId ? cached.client : undefined; + } + const headers = Object.fromEntries( + Object.entries(params.clientOptions?.headers ?? {}).filter( + ([name]) => name.toLowerCase() !== "authorization", + ), + ); + // Completion is one-shot. Clone Bolt's public transport options and team + // scope, but never inherit its retry policy or request deadline. + const client = new WebClient( + token, + resolveSlackWriteClientOptions({ + ...params.clientOptions, + headers, + slackApiUrl: params.listenerClient.slackApiUrl, + teamId, + retryConfig: SLACK_WRITE_RETRY_OPTIONS, + timeout: 0, + }), + ); + slackListenerUploadCompletionClientCache.set(params.listenerClient, { teamId, client }); + return client; +} + export function clearSlackWriteClientCacheForTest(): void { slackWriteClientCache.clear(); + slackListenerUploadCompletionClientCache = new WeakMap(); } diff --git a/extensions/slack/src/client.web-api.test.ts b/extensions/slack/src/client.web-api.test.ts index 5674105c9608..70a023424ac4 100644 --- a/extensions/slack/src/client.web-api.test.ts +++ b/extensions/slack/src/client.web-api.test.ts @@ -1,8 +1,9 @@ // Slack tests cover real Web API routing behavior. import { createServer, type Server } from "node:http"; import type { AddressInfo } from "node:net"; +import { WebClient } from "@slack/web-api"; import { afterEach, describe, expect, it } from "vitest"; -import { createSlackWebClient } from "./client.js"; +import { createSlackWebClient, getSlackListenerUploadCompletionClient } from "./client.js"; const SLACK_API_URL_KEYS = ["SLACK_API_URL"] as const; const PROXY_KEYS = [ @@ -20,6 +21,7 @@ const originalEnv = { ...process.env }; type SlackApiRequest = { authorization?: string; + body?: string; method?: string; url?: string; }; @@ -46,7 +48,10 @@ async function closeServer(server: Server): Promise { }); } -async function startSlackApiServer(requests: SlackApiRequest[]): Promise<{ +async function startSlackApiServer( + requests: SlackApiRequest[], + responseDelayMs = 0, +): Promise<{ baseUrl: string; close(): Promise; }> { @@ -57,17 +62,51 @@ async function startSlackApiServer(requests: SlackApiRequest[]): Promise<{ url: request.url, }); request.resume(); - response.writeHead(200, { "content-type": "application/json" }); - response.end( - `${JSON.stringify({ - ok: true, - team: "Mock Slack", - team_id: "TMOCK", - url: "https://mock.slack.test/", - user: "mock-bot", - user_id: "UMOCK", - })}\n`, - ); + const sendResponse = () => { + response.writeHead(200, { "content-type": "application/json" }); + response.end( + `${JSON.stringify({ + ok: true, + team: "Mock Slack", + team_id: "TMOCK", + url: "https://mock.slack.test/", + user: "mock-bot", + user_id: "UMOCK", + })}\n`, + ); + }; + if (responseDelayMs > 0) { + setTimeout(sendResponse, responseDelayMs); + } else { + sendResponse(); + } + }); + await new Promise((resolve) => { + server.listen(0, "127.0.0.1", resolve); + }); + const address = server.address() as AddressInfo; + return { + baseUrl: `http://127.0.0.1:${address.port}`, + close: () => closeServer(server), + }; +} + +async function startDroppedResponseSlackApiServer(requests: SlackApiRequest[]): Promise<{ + baseUrl: string; + close(): Promise; +}> { + const server = createServer((request) => { + const chunks: Buffer[] = []; + request.on("data", (chunk) => chunks.push(Buffer.from(chunk))); + request.once("end", () => { + requests.push({ + authorization: request.headers.authorization, + body: Buffer.concat(chunks).toString("utf8"), + method: request.method, + url: request.url, + }); + request.socket.destroy(); + }); }); await new Promise((resolve) => { server.listen(0, "127.0.0.1", resolve); @@ -84,6 +123,101 @@ afterEach(() => { }); describe("Slack Web API routing", () => { + it("keeps dropped Enterprise upload completion responses to one team-scoped request", async () => { + for (const key of TEST_ENV_KEYS) { + delete process.env[key]; + } + const requests: SlackApiRequest[] = []; + const server = await startDroppedResponseSlackApiServer(requests); + try { + const clientOptions = { + headers: { + Authorization: "Bearer stale-fixture", + "X-Slack-Test": "preserved", + }, + slackApiUrl: `${server.baseUrl}/api/`, + retryConfig: { retries: 2 }, + }; + const listenerClient = new WebClient("listener-fixture", clientOptions); + const completionClient = getSlackListenerUploadCompletionClient({ + listenerClient, + teamId: "TENTERPRISE1", + clientOptions, + }); + expect(completionClient).toBeDefined(); + if (!completionClient) { + throw new Error("missing Enterprise upload completion client"); + } + expect( + getSlackListenerUploadCompletionClient({ + listenerClient, + teamId: "TENTERPRISE1", + clientOptions, + }), + ).toBe(completionClient); + expect( + getSlackListenerUploadCompletionClient({ + listenerClient, + teamId: "TENTERPRISE2", + clientOptions, + }), + ).toBeUndefined(); + + await expect( + completionClient.files.completeUploadExternal({ + files: [{ id: "F123", title: "proof.txt" }], + channel_id: "C123", + }), + ).rejects.toThrow(); + + expect(requests).toHaveLength(1); + expect(requests[0]).toMatchObject({ + authorization: "Bearer listener-fixture", + method: "POST", + url: "/api/files.completeUploadExternal", + }); + expect(new URLSearchParams(requests[0]?.body).get("team_id")).toBe("TENTERPRISE1"); + expect(requests[0]?.authorization).not.toContain("stale-fixture"); + } finally { + await server.close(); + } + }); + + it("does not inherit the listener request timeout for upload completion", async () => { + for (const key of TEST_ENV_KEYS) { + delete process.env[key]; + } + const requests: SlackApiRequest[] = []; + const server = await startSlackApiServer(requests, 80); + try { + const clientOptions = { + slackApiUrl: `${server.baseUrl}/api/`, + retryConfig: { retries: 2 }, + timeout: 20, + }; + const listenerClient = new WebClient("listener-fixture", clientOptions); + const completionClient = getSlackListenerUploadCompletionClient({ + listenerClient, + teamId: "TENTERPRISE1", + clientOptions, + }); + expect(completionClient).toBeDefined(); + if (!completionClient) { + throw new Error("missing Enterprise upload completion client"); + } + + const result = await completionClient.files.completeUploadExternal({ + files: [{ id: "F123", title: "proof.txt" }], + channel_id: "C123", + }); + + expect(result.ok).toBe(true); + expect(requests).toHaveLength(1); + } finally { + await server.close(); + } + }); + it("routes real WebClient requests to the SLACK_API_URL root", async () => { for (const key of TEST_ENV_KEYS) { delete process.env[key]; diff --git a/extensions/slack/src/monitor/event-scope.test.ts b/extensions/slack/src/monitor/event-scope.test.ts index 456f7cf1a630..36cfa81b31e0 100644 --- a/extensions/slack/src/monitor/event-scope.test.ts +++ b/extensions/slack/src/monitor/event-scope.test.ts @@ -3,27 +3,30 @@ import { describe, expect, it } from "vitest"; import { resolveSlackEventScope } from "./event-scope.js"; const identity = { kind: "enterprise", apiAppId: "A123", enterpriseId: "E123" } as const; -const client = {} as WebClient; +const client = new WebClient("listener-token"); describe("resolveSlackEventScope", () => { it.each(["T111", "T222"])("accepts authorized workspace %s in the same org", (teamId) => { + const listenerClient = new WebClient(`listener-token-${teamId.toLowerCase()}`); const result = resolveSlackEventScope({ identity, body: { api_app_id: "A123" }, context: { isEnterpriseInstall: true, enterpriseId: "E123", teamId }, - client, + client: listenerClient, }); - expect(result).toEqual({ + expect(result).toMatchObject({ ok: true, scope: { apiAppId: "A123", enterpriseId: "E123", teamId, isEnterpriseInstall: true, - client, + client: listenerClient, }, }); - expect(result.ok && result.scope?.client).toBe(client); + expect(result.ok && result.scope?.client).toBe(listenerClient); + expect(result.ok && result.scope?.uploadCompletionClient).toBeInstanceOf(WebClient); + expect(result.ok && result.scope?.uploadCompletionClient).not.toBe(listenerClient); }); it("accepts a signed enterprise event when startup auth.test omitted app_id", () => { @@ -33,7 +36,7 @@ describe("resolveSlackEventScope", () => { context: { isEnterpriseInstall: true, enterpriseId: "E123", teamId: "T111" }, client, }); - expect(result).toEqual({ + expect(result).toMatchObject({ ok: true, scope: { apiAppId: "A123", diff --git a/extensions/slack/src/monitor/event-scope.ts b/extensions/slack/src/monitor/event-scope.ts index 46b83e32eea0..a4945cb8d84e 100644 --- a/extensions/slack/src/monitor/event-scope.ts +++ b/extensions/slack/src/monitor/event-scope.ts @@ -1,6 +1,7 @@ // Slack plugin module validates non-serializable per-event Enterprise Grid scope. -import type { WebClient } from "@slack/web-api"; +import type { WebClient, WebClientOptions } from "@slack/web-api"; import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { getSlackListenerUploadCompletionClient } from "../client.js"; import type { SlackInstallationIdentity } from "./enterprise-install.js"; export type SlackEventScope = { @@ -8,9 +9,11 @@ export type SlackEventScope = { enterpriseId: string; teamId: string; isEnterpriseInstall: true; - // Keep Bolt's exact listener client: Bolt pools it by authorized team and WebClient injects - // that client's team_id into every API call. Do not recreate it or add team_id payloads here. + // Keep Bolt's exact listener client for ordinary reads and writes. client: WebClient; + // Completion is one-shot, so uploads finalize through a team-scoped client + // that cannot inherit Bolt's normal request retries. + uploadCompletionClient?: WebClient; }; export type SlackEventScopeResolution = @@ -38,6 +41,7 @@ export function resolveSlackEventScope(params: { teamId?: unknown; }; client?: WebClient; + clientOptions?: WebClientOptions; }): SlackEventScopeResolution { const context = params.context ?? {}; if (params.identity.kind !== "enterprise") { @@ -74,6 +78,11 @@ export function resolveSlackEventScope(params: { if (!params.client) { return { ok: false, reason: "missing_listener_client" }; } + const uploadCompletionClient = getSlackListenerUploadCompletionClient({ + listenerClient: params.client, + teamId, + clientOptions: params.clientOptions, + }); return { ok: true, scope: { @@ -82,6 +91,7 @@ export function resolveSlackEventScope(params: { teamId, isEnterpriseInstall: true, client: params.client, + ...(uploadCompletionClient ? { uploadCompletionClient } : {}), }, }; } diff --git a/extensions/slack/src/monitor/events/messages.ts b/extensions/slack/src/monitor/events/messages.ts index 89a426db1e09..c51eaec157fe 100644 --- a/extensions/slack/src/monitor/events/messages.ts +++ b/extensions/slack/src/monitor/events/messages.ts @@ -176,6 +176,7 @@ export function registerSlackMessageEvents(params: { body: args.body, context: args.context, client: args.client, + clientOptions: ctx.app.webClientOptions, }); if (!resolved.ok) { logVerbose(`slack: drop event (${resolved.reason})`); diff --git a/extensions/slack/src/outbound-adapter.ts b/extensions/slack/src/outbound-adapter.ts index 68d42644b29d..016fd466f209 100644 --- a/extensions/slack/src/outbound-adapter.ts +++ b/extensions/slack/src/outbound-adapter.ts @@ -465,6 +465,7 @@ export const slackOutbound: ChannelOutboundAdapter = { replyToId, threadId, identity, + onPlatformSendDispatch, onDeliveryResult, }) => await sendSlackOutboundMessage({ @@ -480,6 +481,7 @@ export const slackOutbound: ChannelOutboundAdapter = { replyToId, threadId, identity, + onPlatformSendDispatch, onDeliveryResult, }), }), diff --git a/extensions/slack/src/send.enterprise.test.ts b/extensions/slack/src/send.enterprise.test.ts index 3de5538f97aa..ef0ac33ab89b 100644 --- a/extensions/slack/src/send.enterprise.test.ts +++ b/extensions/slack/src/send.enterprise.test.ts @@ -63,21 +63,30 @@ function createEnterpriseClient(): EnterpriseTestClient { } as unknown as EnterpriseTestClient; } -function enterpriseEventScope(client: WebClient, teamId = "T1") { +function enterpriseEventScope( + client: WebClient, + teamId = "T1", + uploadCompletionClient: WebClient = client, +) { return { apiAppId: "A1", enterpriseId: "E1", isEnterpriseInstall: true as const, teamId, client, + uploadCompletionClient, }; } -function enterpriseOptions(client: WebClient, teamId = "T1") { +function enterpriseOptions( + client: WebClient, + teamId = "T1", + uploadCompletionClient: WebClient = client, +) { return { cfg: ENTERPRISE_CFG, client, - enterpriseEventScope: enterpriseEventScope(client, teamId), + enterpriseEventScope: enterpriseEventScope(client, teamId, uploadCompletionClient), }; } @@ -275,19 +284,38 @@ describe("sendMessageSlack Enterprise listener scope", () => { ); }); - it("uploads the first caption chunk and posts the remainder with the listener client", async () => { + it("rejects Enterprise media before upload without the one-shot completion client", async () => { + const client = createEnterpriseClient(); + const scope = enterpriseEventScope(client); + delete (scope as { uploadCompletionClient?: WebClient }).uploadCompletionClient; + + await expect( + sendMessageSlack("C123", "caption", { + cfg: ENTERPRISE_CFG, + client, + enterpriseEventScope: scope, + mediaUrl: "https://example.com/image.png", + }), + ).rejects.toThrow("missing_enterprise_slack_upload_completion_client"); + + expect(client.files.getUploadURLExternal).not.toHaveBeenCalled(); + expect(fetchWithSsrFGuard).not.toHaveBeenCalled(); + }); + + it("uses a completion-only client and keeps remaining posts on the listener client", async () => { const release = vi.fn(async () => {}); fetchWithSsrFGuard.mockResolvedValue({ response: { ok: true, status: 200 }, release, }); - const client = createEnterpriseClient(); - client.chat.postMessage + const listenerClient = createEnterpriseClient(); + const uploadCompletionClient = createEnterpriseClient(); + listenerClient.chat.postMessage .mockResolvedValueOnce({ ok: true, ts: "123.001", channel: "C123" }) .mockResolvedValueOnce({ ok: true, ts: "123.002", channel: "C123" }); const result = await sendMessageSlack("C123", "12345678abcdefghZ", { - ...enterpriseOptions(client), + ...enterpriseOptions(listenerClient, "T1", uploadCompletionClient), mediaUrl: "https://example.com/image.png", textLimit: 8, mediaMaxBytes: 5, @@ -300,15 +328,19 @@ describe("sendMessageSlack Enterprise listener scope", () => { expect(fetchWithSsrFGuard).toHaveBeenCalledWith( expect.objectContaining({ auditContext: "slack-enterprise-immediate-upload" }), ); - expect(client.files.completeUploadExternal).toHaveBeenCalledWith({ + expect(listenerClient.files.getUploadURLExternal).toHaveBeenCalledOnce(); + expect(listenerClient.files.completeUploadExternal).not.toHaveBeenCalled(); + expect(uploadCompletionClient.files.getUploadURLExternal).not.toHaveBeenCalled(); + expect(uploadCompletionClient.files.completeUploadExternal).toHaveBeenCalledWith({ files: [{ id: "F123", title: "image.png" }], channel_id: "C123", initial_comment: "12345678", }); - expect(client.chat.postMessage.mock.calls.map((call) => call[0]?.text)).toEqual([ + expect(listenerClient.chat.postMessage.mock.calls.map((call) => call[0]?.text)).toEqual([ "abcdefgh", "Z", ]); + expect(uploadCompletionClient.chat.postMessage).not.toHaveBeenCalled(); expect(release).toHaveBeenCalledOnce(); expect(result.receipt).toMatchObject({ primaryPlatformMessageId: "F123", diff --git a/extensions/slack/src/send.ts b/extensions/slack/src/send.ts index e3c11cfade57..948861cad6f7 100644 --- a/extensions/slack/src/send.ts +++ b/extensions/slack/src/send.ts @@ -81,11 +81,13 @@ type SlackEnterpriseEventScope = Readonly<{ teamId: string; isEnterpriseInstall: true; client: WebClient; + uploadCompletionClient?: WebClient; }>; type SlackEnterpriseDelivery = Readonly<{ client: WebClient; teamId: string; + uploadCompletionClient?: WebClient; }>; const slackDefaultSendIdentities = new Map(); @@ -116,7 +118,7 @@ type SlackSendOpts = { metadata?: MessageMetadata; /** Opaque durable intent id used to reconcile ambiguous platform outcomes. */ deliveryQueueId?: string; - /** Refresh durable timing after the per-target queue and before Slack API work. */ + /** Refresh durable timing before recipient-visible or finalizing platform I/O. */ onPlatformSendDispatch?: () => Promise; /** Persist each concrete platform send before any later chunk can fail. */ onDeliveryResult?: (result: SlackSendResult) => Promise | void; @@ -939,6 +941,9 @@ export async function sendMessageSlack( ? Object.freeze({ client: enterpriseEventScope.client, teamId: enterpriseEventScope.teamId, + ...(enterpriseEventScope.uploadCompletionClient + ? { uploadCompletionClient: enterpriseEventScope.uploadCompletionClient } + : {}), }) : undefined; if (isSilentReplyText(trimmedMessage) && !opts.mediaUrl && !opts.blocks) { @@ -1125,9 +1130,15 @@ async function sendMessageSlackQueuedInner(params: { let canonicalDeliveredThreadTs: string | undefined; let chunksToPost: string[]; if (opts.mediaUrl) { + if (enterpriseDelivery && !enterpriseDelivery.uploadCompletionClient) { + throw new Error("missing_enterprise_slack_upload_completion_client"); + } const [firstChunk, ...rest] = resolvedChunks; lastMessageId = await uploadSlackFile({ client, + ...(enterpriseDelivery?.uploadCompletionClient + ? { completionClient: enterpriseDelivery.uploadCompletionClient } + : {}), channelId, mediaUrl: opts.mediaUrl, mediaAccess: opts.mediaAccess, diff --git a/extensions/slack/src/send.upload.test.ts b/extensions/slack/src/send.upload.test.ts index d0ef73f875e2..6a7fba988c0e 100644 --- a/extensions/slack/src/send.upload.test.ts +++ b/extensions/slack/src/send.upload.test.ts @@ -1,5 +1,11 @@ // Slack tests cover send.upload plugin behavior. import type { WebClient } from "@slack/web-api"; +import { + formatErrorMessage, + PlatformMessageNotDispatchedError, +} from "openclaw/plugin-sdk/error-runtime"; +import type { LookupFn } from "openclaw/plugin-sdk/ssrf-runtime"; +import { withServer } from "openclaw/plugin-sdk/test-env"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import "./blocks.test-helpers.js"; import { @@ -16,19 +22,68 @@ const loadOutboundMediaFromUrlMock = vi.hoisted(() => fileName: "screenshot.png", })), ); +const cleanupUploadTimeout = vi.hoisted(() => vi.fn()); +const uploadTimeoutControllers = vi.hoisted(() => [] as AbortController[]); +const buildTimeoutAbortSignal = vi.hoisted(() => + vi.fn((params: { timeoutMs?: number }) => { + if (!Number.isFinite(params.timeoutMs) || (params.timeoutMs ?? 0) <= 0) { + throw new Error("Slack upload timeout requires a finite budget"); + } + const controller = new AbortController(); + uploadTimeoutControllers.push(controller); + return { + signal: controller.signal, + cleanup: () => { + const index = uploadTimeoutControllers.indexOf(controller); + if (index >= 0) { + uploadTimeoutControllers.splice(index, 1); + } + cleanupUploadTimeout(); + }, + refresh: () => {}, + }; + }), +); const fetchWithSsrFGuard = vi.fn( - async (params: { url: string; init?: RequestInit }) => - ({ - response: await fetch(params.url, params.init), + async ( + params: Parameters[0], + ) => { + const signal = params.signal; + if (!signal) { + throw new Error("guarded Slack upload fetch requires a finite timeout signal"); + } + return { + response: await fetch(params.url, { + ...params.init, + signal, + }), finalUrl: params.url, release: async () => {}, - }) as const, + } as const; + }, ); -vi.mock("openclaw/plugin-sdk/ssrf-runtime", () => ({ - fetchWithSsrFGuard: (...args: unknown[]) => - fetchWithSsrFGuard(...(args as [params: { url: string; init?: RequestInit }])), -})); +vi.mock("openclaw/plugin-sdk/ssrf-runtime", async () => { + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + return { + ...actual, + fetchWithSsrFGuard: (...args: unknown[]) => + fetchWithSsrFGuard(...(args as [params: Parameters[0]])), + }; +}); + +vi.mock("openclaw/plugin-sdk/extension-shared", async () => { + const actual = await vi.importActual( + "openclaw/plugin-sdk/extension-shared", + ); + return { + ...actual, + buildTimeoutAbortSignal: (...args: unknown[]) => + buildTimeoutAbortSignal(...(args as [params: { timeoutMs?: number }])), + }; +}); vi.mock("openclaw/plugin-sdk/fetch-runtime", async () => { const actual = await vi.importActual( @@ -131,8 +186,9 @@ function expectCompletedUpload(params: { return payload; } -function createUploadTestClient(): UploadTestClient { +function createUploadTestClient(slackApiUrl = "https://slack.com/api/"): UploadTestClient { return { + slackApiUrl, conversations: { open: vi.fn<(...args: unknown[]) => Promise>(async () => ({ channel: { id: "D99RESOLVED" }, @@ -146,7 +202,7 @@ function createUploadTestClient(): UploadTestClient { files: { getUploadURLExternal: vi.fn<(...args: unknown[]) => Promise>(async () => ({ ok: true, - upload_url: "https://uploads.slack.test/upload", + upload_url: "https://files.slack.com/upload", file_id: "F001", })), completeUploadExternal: vi.fn<(...args: unknown[]) => Promise>(async () => ({ @@ -164,6 +220,9 @@ describe("sendMessageSlack file upload with user IDs", () => { async () => new Response("ok", { status: 200 }), ) as unknown as typeof fetch; fetchWithSsrFGuard.mockClear(); + buildTimeoutAbortSignal.mockClear(); + cleanupUploadTimeout.mockClear(); + uploadTimeoutControllers.length = 0; loadOutboundMediaFromUrlMock.mockClear(); clearSlackDmChannelCache(); clearSlackThreadParticipationCache(); @@ -171,6 +230,7 @@ describe("sendMessageSlack file upload with user IDs", () => { afterEach(() => { globalThis.fetch = originalFetch; + vi.unstubAllEnvs(); vi.restoreAllMocks(); }); @@ -359,14 +419,37 @@ describe("sendMessageSlack file upload with user IDs", () => { it("uploads bytes to the presigned URL and completes with thread+caption", async () => { const client = createUploadTestClient(); + const events: string[] = []; + globalThis.fetch = vi.fn(async () => { + events.push("byte-upload"); + return new Response("ok", { status: 200 }); + }) as unknown as typeof fetch; + client.files.completeUploadExternal.mockImplementationOnce(async () => { + events.push("completion"); + return { ok: true }; + }); + let finishDispatch: () => void = () => {}; + const dispatchFinished = new Promise((resolve) => { + finishDispatch = resolve; + }); + const onPlatformSendDispatch = vi.fn(async () => { + events.push("dispatch-start"); + await dispatchFinished; + events.push("dispatch-end"); + }); - const result = await sendMessageSlack("channel:C123CHAN", "caption", { + const sendPromise = sendMessageSlack("channel:C123CHAN", "caption", { token: "xoxb-test", cfg: SLACK_TEST_CFG, client, mediaUrl: "/tmp/threaded.png", threadTs: "171.222", + onPlatformSendDispatch, }); + await vi.waitFor(() => expect(onPlatformSendDispatch).toHaveBeenCalledOnce()); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + finishDispatch(); + const result = await sendPromise; expect(client.files.getUploadURLExternal).toHaveBeenCalledWith({ filename: "screenshot.png", @@ -375,13 +458,29 @@ describe("sendMessageSlack file upload with user IDs", () => { const fetchCalls = (globalThis.fetch as unknown as MockCalls).mock.calls; expect(fetchCalls).toHaveLength(1); const [fetchUrl, fetchInit] = fetchCalls[0] ?? []; - expect(fetchUrl).toBe("https://uploads.slack.test/upload"); + expect(fetchUrl).toBe("https://files.slack.com/upload"); expectFields(requireRecord(fetchInit, "fetch init"), { method: "POST" }); + expectOnlyCallFirstArg(buildTimeoutAbortSignal, { + timeoutMs: 120_000, + operation: "slack-upload-file", + url: "https://files.slack.com", + }); expectOnlyCallFirstArg(fetchWithSsrFGuard, { - url: "https://uploads.slack.test/upload", + url: "https://files.slack.com/upload", mode: "trusted_env_proxy", + signal: expect.any(AbortSignal), + requireHttps: true, + policy: { + hostnameAllowlist: ["files.slack.com"], + allowRfc2544BenchmarkRange: true, + }, + capture: false, auditContext: "slack-upload-file", }); + expect(cleanupUploadTimeout).toHaveBeenCalledOnce(); + expect(uploadTimeoutControllers).toHaveLength(0); + expect(onPlatformSendDispatch).toHaveBeenCalledOnce(); + expect(events).toEqual(["byte-upload", "dispatch-start", "dispatch-end", "completion"]); expectCompletedUpload({ client, expected: { @@ -394,6 +493,497 @@ describe("sendMessageSlack file upload with user IDs", () => { expect(result.receipt.threadId).toBe("171.222"); }); + it("keeps the presigned upload capability out of timeout logging", async () => { + const client = createUploadTestClient(); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: "https://files.slack.com/upload/v1/secret-capability", + file_id: "F001", + }); + + await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/secret.png", + }); + + expectOnlyCallFirstArg(buildTimeoutAbortSignal, { + timeoutMs: 120_000, + operation: "slack-upload-file", + url: "https://files.slack.com", + }); + expectOnlyCallFirstArg(fetchWithSsrFGuard, { + url: "https://files.slack.com/upload/v1/secret-capability", + }); + }); + + it("preserves HTTP upload URLs on an alternate Slack API origin", async () => { + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + await withServer( + (req, res) => { + expect(req.method).toBe("POST"); + expect(req.url).toBe("/upload/v1/capability"); + req.resume(); + res.end("ok"); + }, + async (baseUrl) => { + vi.stubEnv("NO_PROXY", "127.0.0.1,localhost"); + vi.stubEnv("no_proxy", "127.0.0.1,localhost"); + const client = createUploadTestClient(`${baseUrl}/api/`); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: `${baseUrl}/upload/v1/capability`, + file_id: "F001", + }); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => { + const mockedFetch = globalThis.fetch; + globalThis.fetch = originalFetch; + try { + return await actual.fetchWithSsrFGuard(params); + } finally { + globalThis.fetch = mockedFetch; + } + }); + + await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/alternate-root.png", + }); + + expectCompletedUpload({ client, expected: { channel_id: "C123CHAN" } }); + expect(cleanupUploadTimeout).toHaveBeenCalledOnce(); + expect(uploadTimeoutControllers).toHaveLength(0); + }, + ); + }); + + it("allows an exact Slack upload host returned by a custom API root", async () => { + const client = createUploadTestClient("https://slack-relay.example/api/"); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: "https://files.slack.com/upload/v1/relayed-capability", + file_id: "F001", + }); + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + const networkFetch = vi.fn(async () => new Response("ok", { status: 200 })); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => + actual.fetchWithSsrFGuard({ + ...params, + fetchImpl: networkFetch, + lookupFn: (async () => [{ address: "93.184.216.34", family: 4 }]) as unknown as LookupFn, + }), + ); + + await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/relayed-upload.png", + }); + + expect(networkFetch).toHaveBeenCalledOnce(); + expectCompletedUpload({ client, expected: { channel_id: "C123CHAN" } }); + }); + + it("allows GovSlack upload destinations through the real hostname guard", async () => { + const client = createUploadTestClient("https://slack-gov.com/api/"); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: "https://files.slack-gov.com/upload/v1/gov-capability", + file_id: "F001", + }); + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + const networkFetch = vi.fn(async () => new Response("ok", { status: 200 })); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => + actual.fetchWithSsrFGuard({ ...params, fetchImpl: networkFetch }), + ); + + await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/gov-slack.png", + }); + + expect(networkFetch).toHaveBeenCalledOnce(); + expectCompletedUpload({ client, expected: { channel_id: "C123CHAN" } }); + }); + + it("retains the shipped RFC2544 fake-IP path for an exact Slack upload host", async () => { + const client = createUploadTestClient(); + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + const networkFetch = vi.fn(async () => new Response("ok", { status: 200 })); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => + actual.fetchWithSsrFGuard({ + ...params, + fetchImpl: networkFetch, + lookupFn: (async () => [{ address: "198.18.0.10", family: 4 }]) as unknown as LookupFn, + }), + ); + + await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/fake-ip.png", + }); + + expect(networkFetch).toHaveBeenCalledOnce(); + expectCompletedUpload({ client, expected: { channel_id: "C123CHAN" } }); + }); + + it.each(["10.0.0.1", "169.254.1.1"])( + "rejects exact Slack upload hosts resolving to blocked address %s", + async (address) => { + const client = createUploadTestClient(); + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + const networkFetch = vi.fn(async () => new Response("unexpected")); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => + actual.fetchWithSsrFGuard({ + ...params, + fetchImpl: networkFetch, + lookupFn: (async () => [{ address, family: 4 }]) as unknown as LookupFn, + }), + ); + + await expect( + sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/private-address.png", + }), + ).rejects.toThrow(); + + expect(networkFetch).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }, + ); + + it.each([ + [ + "public non-Slack", + "https://slack.com/api/", + "https://example.com/upload/v1/not-slack", + "SsrFBlockedError", + ], + [ + "plaintext Slack", + "https://slack.com/api/", + "http://files.slack.com/upload/v1/plaintext", + "Error", + ], + [ + "commercial Slack to GovSlack", + "https://slack.com/api/", + "https://files.slack-gov.com/upload/v1/cross-plane", + "SsrFBlockedError", + ], + [ + "GovSlack to commercial Slack", + "https://slack-gov.com/api/", + "https://files.slack.com/upload/v1/cross-plane", + "SsrFBlockedError", + ], + [ + "trailing-dot commercial Slack to GovSlack", + "https://slack.com./api/", + "https://files.slack-gov.com/upload/v1/cross-plane", + "SsrFBlockedError", + ], + [ + "trailing-dot GovSlack to commercial Slack", + "https://slack-gov.com./api/", + "https://files.slack.com/upload/v1/cross-plane", + "SsrFBlockedError", + ], + [ + "undocumented commercial subdomain", + "https://slack.com/api/", + "https://future-upload.slack.com/upload/v1/capability", + "SsrFBlockedError", + ], + [ + "undocumented GovSlack subdomain", + "https://slack-gov.com/api/", + "https://future-upload.slack-gov.com/upload/v1/capability", + "SsrFBlockedError", + ], + ])( + "rejects %s upload destinations before network access", + async (_label, apiUrl, uploadUrl, error) => { + const client = createUploadTestClient(apiUrl); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: uploadUrl, + file_id: "F001", + }); + const actual = await vi.importActual( + "openclaw/plugin-sdk/ssrf-runtime", + ); + const networkFetch = vi.fn(async () => new Response("unexpected")); + fetchWithSsrFGuard.mockImplementationOnce(async (params) => + actual.fetchWithSsrFGuard({ ...params, fetchImpl: networkFetch }), + ); + + const rejection = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/rejected-upload.png", + }).catch((cause: unknown) => cause); + + expect(rejection).toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(rejection).toMatchObject({ + cause: expect.objectContaining({ name: error }), + }); + + expect(networkFetch).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }, + ); + + it("rejects upload destinations outside an explicitly configured API origin", async () => { + const client = createUploadTestClient("http://slack-compatible.example/api/"); + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: "http://other-compatible.example/upload/v1/capability", + file_id: "F001", + }); + + await expect( + sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/wrong-origin.png", + }), + ).rejects.toThrow("must match the configured Slack API origin"); + + expect(fetchWithSsrFGuard).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }); + + it("times out a hanging presigned URL upload", async () => { + const client = createUploadTestClient(); + const closedResponses = vi.fn(); + + await withServer( + (req, res) => { + req.resume(); + const route = `${req.method ?? "GET"} ${req.url ?? "/"}`; + res.on("close", () => closedResponses(route)); + if (route === "POST /upload") { + // Fire the mocked deadline only after the request reaches the server, + // keeping the cancellation regression deterministic under CI load. + const controller = uploadTimeoutControllers.at(-1); + if (!controller) { + throw new Error("missing Slack upload timeout controller"); + } + const error = new Error("request timed out"); + error.name = "TimeoutError"; + controller.abort(error); + return; + } + res.statusCode = 500; + res.end(`unexpected ${route}`); + }, + async (baseUrl) => { + globalThis.fetch = originalFetch; + client.files.getUploadURLExternal.mockResolvedValueOnce({ + ok: true, + upload_url: `${baseUrl}/upload`, + file_id: "F001", + }); + + const onPlatformSendDispatch = vi.fn(); + const error = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/hanging.png", + onPlatformSendDispatch, + }).catch((cause: unknown) => cause); + + expect(error).toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(error).toMatchObject({ + name: "PlatformMessageNotDispatchedError", + code: "OPENCLAW_PLATFORM_MESSAGE_NOT_DISPATCHED", + cause: expect.objectContaining({ name: "TimeoutError" }), + }); + + await vi.waitFor(() => expect(closedResponses).toHaveBeenCalledWith("POST /upload")); + expectOnlyCallFirstArg(buildTimeoutAbortSignal, { + timeoutMs: 120_000, + operation: "slack-upload-file", + url: baseUrl, + }); + expectOnlyCallFirstArg(fetchWithSsrFGuard, { signal: expect.any(AbortSignal) }); + expect(cleanupUploadTimeout).toHaveBeenCalledOnce(); + expect(uploadTimeoutControllers).toHaveLength(0); + expect(onPlatformSendDispatch).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }, + ); + }); + + it.each([201, 204, 500])("rejects a non-200 byte-upload response (%s)", async (status) => { + const client = createUploadTestClient(); + const onPlatformSendDispatch = vi.fn(); + globalThis.fetch = vi.fn(async () => new Response(null, { status })) as unknown as typeof fetch; + + const error = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/non-200.png", + onPlatformSendDispatch, + }).catch((cause: unknown) => cause); + + expect(error).toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(error).toMatchObject({ + message: "Slack external upload failed before completion dispatch", + cause: expect.objectContaining({ + code: `HTTP_${status}`, + message: `Slack external upload returned HTTP ${status}`, + }), + }); + expect(onPlatformSendDispatch).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }); + + it("marks a non-timeout byte-upload transport failure as not dispatched", async () => { + const client = createUploadTestClient(); + const onPlatformSendDispatch = vi.fn(); + const transportError = Object.assign( + new Error( + "socket closed at https://files.slack.com/upload/v1/CAPABILITY_SENTINEL?token=QUERY_SENTINEL", + ), + { code: "ECONNRESET" }, + ); + globalThis.fetch = vi.fn(async () => { + throw transportError; + }) as unknown as typeof fetch; + + const error = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/transport-failure.png", + onPlatformSendDispatch, + }).catch((cause: unknown) => cause); + + expect(error).toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(error).toMatchObject({ + message: "Slack external upload failed before completion dispatch", + cause: expect.objectContaining({ + code: "ECONNRESET", + message: "Slack external upload transfer failed", + }), + }); + expect(formatErrorMessage(error)).not.toContain("CAPABILITY_SENTINEL"); + expect(formatErrorMessage(error)).not.toContain("QUERY_SENTINEL"); + expect(cleanupUploadTimeout).toHaveBeenCalledOnce(); + expect(onPlatformSendDispatch).not.toHaveBeenCalled(); + expect(client.files.completeUploadExternal).not.toHaveBeenCalled(); + }); + + it("disposes the byte-upload response and timeout before waiting for completion", async () => { + const client = createUploadTestClient(); + const uploadResponse = new Response("ok", { status: 200 }); + const cancelUploadBody = vi.spyOn(uploadResponse.body!, "cancel"); + cancelUploadBody.mockRejectedValueOnce(new Error("response body cleanup failed")); + globalThis.fetch = vi.fn(async () => uploadResponse) as unknown as typeof fetch; + let markCompletionStarted: () => void = () => {}; + const completionStarted = new Promise((resolve) => { + markCompletionStarted = resolve; + }); + let finishCompletion: (value: { ok: true }) => void = () => {}; + const completionResult = new Promise<{ ok: true }>((resolve) => { + finishCompletion = resolve; + }); + client.files.completeUploadExternal.mockImplementationOnce(() => { + markCompletionStarted(); + return completionResult; + }); + + const sendPromise = sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/completion-pending.png", + }); + + await completionStarted; + expect(cancelUploadBody).toHaveBeenCalledOnce(); + expect(cleanupUploadTimeout).toHaveBeenCalledOnce(); + expect(uploadTimeoutControllers).toHaveLength(0); + await expect( + Promise.race([ + sendPromise.then( + () => "settled", + () => "settled", + ), + Promise.resolve("pending"), + ]), + ).resolves.toBe("pending"); + + finishCompletion({ ok: true }); + await expect(sendPromise).resolves.toMatchObject({ messageId: "F001" }); + }); + + it("keeps completion failures unmarked because Slack may have finalized the upload", async () => { + const client = createUploadTestClient(); + const onPlatformSendDispatch = vi.fn(); + client.files.completeUploadExternal.mockRejectedValueOnce(new Error("completion unavailable")); + + const error = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/completion-failure.png", + onPlatformSendDispatch, + }).catch((cause: unknown) => cause); + + expect(error).toMatchObject({ message: "completion unavailable" }); + expect(error).not.toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(onPlatformSendDispatch).toHaveBeenCalledOnce(); + }); + + it("keeps completion error responses unmarked because Slack may have finalized the upload", async () => { + const client = createUploadTestClient(); + const onPlatformSendDispatch = vi.fn(); + client.files.completeUploadExternal.mockResolvedValueOnce({ + ok: false, + error: "completion_failed", + }); + + const error = await sendMessageSlack("channel:C123CHAN", "caption", { + token: "xoxb-test", + cfg: SLACK_TEST_CFG, + client, + mediaUrl: "/tmp/completion-error-response.png", + onPlatformSendDispatch, + }).catch((cause: unknown) => cause); + + expect(error).toMatchObject({ message: "Failed to complete upload: completion_failed" }); + expect(error).not.toBeInstanceOf(PlatformMessageNotDispatchedError); + expect(onPlatformSendDispatch).toHaveBeenCalledOnce(); + }); + it("uses explicit upload filename and title overrides when provided", async () => { const client = createUploadTestClient(); diff --git a/scripts/plugin-sdk-surface-report.mjs b/scripts/plugin-sdk-surface-report.mjs index 2cf7fa274c5a..3d4f89081472 100644 --- a/scripts/plugin-sdk-surface-report.mjs +++ b/scripts/plugin-sdk-surface-report.mjs @@ -195,7 +195,7 @@ export function readPluginSdkSurfaceBudgets(env = process.env) { ), publicExports: readPluginSdkSurfaceBudgetEnv( "OPENCLAW_PLUGIN_SDK_MAX_PUBLIC_EXPORTS", - 10492, + 10493, env, ), publicFunctionExports: readPluginSdkSurfaceBudgetEnv( diff --git a/src/channels/message/types.ts b/src/channels/message/types.ts index 9b34b414ae74..9d29ef8263d9 100644 --- a/src/channels/message/types.ts +++ b/src/channels/message/types.ts @@ -181,7 +181,7 @@ export type ChannelMessageSendTextContext = { gatewayClientScopes?: readonly string[]; /** @internal Opaque durable intent id for exact provider-side send reconciliation. */ deliveryQueueId?: string; - /** @internal Refresh durable timing after provider serialization and before I/O. */ + /** @internal Refresh durable timing before recipient-visible or finalizing platform I/O. */ onPlatformSendDispatch?: () => Promise; /** @internal Report each completed platform sub-send before another fallible step. */ onDeliveryResult?: (result: ChannelMessageSendResult) => Promise | void; diff --git a/src/channels/plugins/outbound.types.ts b/src/channels/plugins/outbound.types.ts index 1cf48ea203f1..e839b916f7f1 100644 --- a/src/channels/plugins/outbound.types.ts +++ b/src/channels/plugins/outbound.types.ts @@ -42,7 +42,7 @@ export type ChannelOutboundContext = { gatewayClientScopes?: readonly string[]; /** @internal Opaque durable intent id for exact provider-side send reconciliation. */ deliveryQueueId?: string; - /** @internal Refresh durable timing after provider serialization and before I/O. */ + /** @internal Refresh durable timing before recipient-visible or finalizing platform I/O. */ onPlatformSendDispatch?: () => Promise; /** @internal Report each completed platform sub-send before starting another fallible step. */ onDeliveryResult?: (result: OutboundDeliveryResult) => Promise | void; diff --git a/src/cli/send-runtime/channel-outbound-send.ts b/src/cli/send-runtime/channel-outbound-send.ts index c62a2ba8dc1e..f41f16f8d13e 100644 --- a/src/cli/send-runtime/channel-outbound-send.ts +++ b/src/cli/send-runtime/channel-outbound-send.ts @@ -27,7 +27,7 @@ type RuntimeSendOpts = { gatewayClientScopes?: readonly string[]; /** @internal Opaque durable intent id for provider-side reconciliation. */ deliveryQueueId?: string; - /** @internal Refresh durable timing after provider serialization and before I/O. */ + /** @internal Refresh durable timing before recipient-visible or finalizing platform I/O. */ onPlatformSendDispatch?: () => Promise; textMode?: "markdown" | "html"; }; diff --git a/src/infra/delivery-recovery.shared.ts b/src/infra/delivery-recovery.shared.ts index ff57d5cf728e..590f071d9b5b 100644 --- a/src/infra/delivery-recovery.shared.ts +++ b/src/infra/delivery-recovery.shared.ts @@ -1,5 +1,6 @@ import { sleep } from "../utils/sleep.js"; import { collectErrorGraphCandidates, extractErrorCode } from "./errors.js"; +import { isPlatformMessageNotDispatchedError } from "./outbound/deliver-types.js"; import { getRetryAttemptErrors } from "./retry-attempt-errors.js"; const RECOVERY_BACKOFF_MS: readonly number[] = [5_000, 25_000, 120_000, 600_000]; @@ -15,6 +16,11 @@ const PRE_CONNECT_ERROR_CODES = new Set([ ]); const TRANSPORT_ERROR_CODE_RE = /^(?:E(?:AI_|CONN|NET|HOST|ADDR|PIPE|TIMEDOUT|SOCKET)|UND_ERR_|ERR_(?:NETWORK|HTTP2|QUIC|TLS|SSL))/; +const UNPROVEN_ERROR_BRANCH = "unproven delivery error branch"; + +function preserveProofBranches(branches: readonly unknown[] | undefined): unknown[] { + return branches?.map((branch) => branch ?? UNPROVEN_ERROR_BRANCH) ?? []; +} function isProvenPreConnectCandidate(candidate: unknown): boolean { const code = extractErrorCode(candidate)?.trim().toUpperCase(); @@ -30,25 +36,28 @@ function isProvenPreConnectCandidate(candidate: unknown): boolean { function nestedErrorCandidates(current: Record): unknown[] { const retryAttempts = getRetryAttemptErrors(current); - if (isProvenPreConnectCandidate(current)) { - return retryAttempts ? [...retryAttempts] : []; + const retryBranches = preserveProofBranches(retryAttempts); + // The explicit marker covers its cause: the provider owns the final dispatch + // boundary and proved that no recipient-visible send could have completed. + if (isPlatformMessageNotDispatchedError(current) || isProvenPreConnectCandidate(current)) { + return retryBranches; } const nested = [current.cause, current.original, current.error, current.reason]; - if (Array.isArray(current.errors)) { - nested.push(...current.errors); - } const nestedObjects = nested.filter( (candidate) => candidate !== null && typeof candidate === "object", ); - return retryAttempts ? [...retryAttempts, ...nestedObjects] : nestedObjects; + const aggregateBranches = Array.isArray(current.errors) + ? preserveProofBranches(current.errors) + : []; + return [...retryBranches, ...aggregateBranches, ...nestedObjects]; } -export function isPreConnectNetworkError(err: unknown): boolean { - let foundPreConnectProof = false; +export function isProvenDeliveryNotSentError(err: unknown): boolean { + let foundNotSentProof = false; for (const candidate of collectErrorGraphCandidates(err, nestedErrorCandidates)) { const code = extractErrorCode(candidate)?.trim().toUpperCase(); - if (isProvenPreConnectCandidate(candidate)) { - foundPreConnectProof = true; + if (isPlatformMessageNotDispatchedError(candidate) || isProvenPreConnectCandidate(candidate)) { + foundNotSentProof = true; continue; } const nested = @@ -73,7 +82,7 @@ export function isPreConnectNetworkError(err: unknown): boolean { return false; } } - return foundPreConnectProof; + return foundNotSentProof; } export function computeBackoffMs(retryCount: number): number { diff --git a/src/infra/outbound/deliver-types.ts b/src/infra/outbound/deliver-types.ts index e84f9b94c5ab..ef89f68189c0 100644 --- a/src/infra/outbound/deliver-types.ts +++ b/src/infra/outbound/deliver-types.ts @@ -31,6 +31,28 @@ export type OutboundPayloadDeliverySuppressionReason = /** Delivery phase where a failure occurred. */ export type OutboundDeliveryFailureStage = "platform_send" | "queue" | "unknown"; +export const PLATFORM_MESSAGE_NOT_DISPATCHED_ERROR_CODE = + "OPENCLAW_PLATFORM_MESSAGE_NOT_DISPATCHED"; + +/** + * Provider assertion that retrying cannot duplicate a recipient-visible send. + * Never use this after a finalization/send call returned an ambiguous result. + */ +export class PlatformMessageNotDispatchedError extends Error { + readonly code = PLATFORM_MESSAGE_NOT_DISPATCHED_ERROR_CODE; + + constructor(message: string, options: { cause: unknown }) { + super(message, { cause: options.cause }); + this.name = "PlatformMessageNotDispatchedError"; + } +} + +export function isPlatformMessageNotDispatchedError( + error: unknown, +): error is PlatformMessageNotDispatchedError { + return error instanceof PlatformMessageNotDispatchedError; +} + /** Per-payload delivery status emitted to callers and channel send summaries. */ export type OutboundPayloadDeliveryOutcome = | { diff --git a/src/infra/outbound/deliver.queue-integration.test.ts b/src/infra/outbound/deliver.queue-integration.test.ts index 5c969bc4247c..a0bf067d57a3 100644 --- a/src/infra/outbound/deliver.queue-integration.test.ts +++ b/src/infra/outbound/deliver.queue-integration.test.ts @@ -7,6 +7,7 @@ import { setActivePluginRegistry, } from "../../plugins/runtime.js"; import { createOutboundTestPlugin, createTestRegistry } from "../../test-utils/channel-plugins.js"; +import { PlatformMessageNotDispatchedError } from "./deliver-types.js"; import { drainPendingDeliveries, type DeliverFn, loadPendingDeliveries } from "./delivery-queue.js"; import { createRecoveryLog, @@ -225,4 +226,42 @@ describe("deliverOutboundPayloads queue integration: mid-batch failure with send expect(recoverySendMatrix).toHaveBeenCalledTimes(2); expect(await loadPendingDeliveries(tmpDir)).toHaveLength(0); }); + + it("replays an entry after the provider proves no platform message was dispatched", async () => { + process.env.OPENCLAW_STATE_DIR = tmpDir; + const notDispatchedError = new PlatformMessageNotDispatchedError( + "upload timed out before completion dispatch", + { cause: new Error("request timed out") }, + ); + const sendMatrix = vi.fn().mockRejectedValueOnce(notDispatchedError); + + await expect( + deliverOutboundPayloads({ + cfg: {} as OpenClawConfig, + channel: "matrix", + to: "!room:example", + payloads: [{ text: "first" }], + deps: { matrix: sendMatrix }, + queuePolicy: "required", + }), + ).rejects.toThrow("upload timed out before completion dispatch"); + + const beforeDrain = await loadPendingDeliveries(tmpDir); + expect(beforeDrain).toHaveLength(1); + expect(beforeDrain[0]?.recoveryState).toBeUndefined(); + expect(beforeDrain[0]?.platformSendStartedAt).toBeUndefined(); + + const recoverySendMatrix = vi.fn().mockResolvedValueOnce({ messageId: "recovered" }); + const deliver = vi.fn(async (params) => + deliverOutboundPayloads({ + ...params, + deps: { matrix: recoverySendMatrix }, + }), + ); + await drainMatrixReconnect({ deliver, stateDir: tmpDir }); + + expect(deliver).toHaveBeenCalledOnce(); + expect(recoverySendMatrix).toHaveBeenCalledOnce(); + expect(await loadPendingDeliveries(tmpDir)).toHaveLength(0); + }); }); diff --git a/src/infra/outbound/deliver.test.ts b/src/infra/outbound/deliver.test.ts index bc96e5dd47b2..e14ffad995e1 100644 --- a/src/infra/outbound/deliver.test.ts +++ b/src/infra/outbound/deliver.test.ts @@ -31,6 +31,7 @@ import { } from "../diagnostic-events.js"; import { retryAsync } from "../retry.js"; import { resolvePreferredOpenClawTmpDir } from "../tmp-openclaw-dir.js"; +import { PlatformMessageNotDispatchedError } from "./deliver-types.js"; const mocks = vi.hoisted(() => ({ appendAssistantMessageToSessionTranscript: vi.fn<() => Promise>( @@ -1875,19 +1876,45 @@ describe("deliverOutboundPayloads", () => { expect(queueMocks.failDelivery).not.toHaveBeenCalled(); }); - it("preserves queued send evidence when any best-effort failure is ambiguous", async () => { + it("clears queued send evidence for an all-not-dispatched best-effort failure", async () => { + const sendMatrix = vi.fn().mockRejectedValueOnce( + new PlatformMessageNotDispatchedError("upload timed out before completion dispatch", { + cause: new Error("request timed out"), + }), + ); + + await expect( + deliverOutboundPayloads({ + cfg: {}, + channel: "matrix", + to: "!room:example", + payloads: [{ text: "hello" }], + deps: { matrix: sendMatrix }, + queuePolicy: "required", + bestEffort: true, + }), + ).resolves.toEqual([]); + + expect(queueMocks.failDeliveryBeforePlatformSend).toHaveBeenCalledWith( + "mock-queue-id", + "partial delivery failure (bestEffort)", + ); + expect(queueMocks.failDelivery).not.toHaveBeenCalled(); + }); + + it("preserves queued send evidence when a marked best-effort batch has an ambiguous failure", async () => { const ambiguousError = Object.assign(new Error("connect ECONNRESET"), { code: "ECONNRESET", syscall: "connect", }); - const preConnectError = Object.assign(new Error("connect ECONNREFUSED"), { - code: "ECONNREFUSED", - syscall: "connect", - }); + const notDispatchedError = new PlatformMessageNotDispatchedError( + "upload timed out before completion dispatch", + { cause: new Error("request timed out") }, + ); const sendMatrix = vi .fn() .mockRejectedValueOnce(ambiguousError) - .mockRejectedValueOnce(preConnectError); + .mockRejectedValueOnce(notDispatchedError); await expect( deliverOutboundPayloads({ diff --git a/src/infra/outbound/deliver.ts b/src/infra/outbound/deliver.ts index 351ebe3d80be..e67a64a3ee7c 100644 --- a/src/infra/outbound/deliver.ts +++ b/src/infra/outbound/deliver.ts @@ -43,7 +43,7 @@ import type { OutboundMediaAccess } from "../../media/load-options.js"; import { resolveAgentScopedOutboundMediaAccess } from "../../media/read-capability.js"; import { getGlobalHookRunner } from "../../plugins/hook-runner-global.js"; import { createLazyRuntimeModule } from "../../shared/lazy-runtime.js"; -import { isPreConnectNetworkError } from "../delivery-recovery.shared.js"; +import { isProvenDeliveryNotSentError } from "../delivery-recovery.shared.js"; import { diagnosticErrorCategory } from "../diagnostic-error-metadata.js"; import { emitInternalDiagnosticEvent as emitDiagnosticEvent, @@ -745,7 +745,7 @@ type DeliverOutboundPayloadsCoreParams = { requiredUnknownSendReconciliation?: boolean; /** @internal Caller preflight explicitly required provider unknown-send reconciliation. */ requireUnknownSendReconciliation?: boolean; - /** @internal Refresh durable timing after provider serialization and before I/O. */ + /** @internal Refresh durable timing before recipient-visible or finalizing platform I/O. */ onPlatformSendDispatch?: () => Promise; /** Session/agent context used for hooks and media local-root scoping. */ session?: OutboundSessionContext; @@ -1431,7 +1431,7 @@ async function deliverOutboundPayloadsWithQueueCleanup( // payload failed so we can call failDelivery instead of ackDelivery. let hadPartialFailure = false; let lastPayloadError: unknown; - let partialFailuresArePreConnect = true; + let partialFailuresAreProvenNotSent = true; const queuePolicy = params.queuePolicy ?? "best_effort"; const platformQueueId = queueId ?? params.deliveryQueueId; const platformQueuePolicy = queueId ? queuePolicy : (params.queuePolicy ?? "required"); @@ -1501,7 +1501,7 @@ async function deliverOutboundPayloadsWithQueueCleanup( onError: (err: unknown, payload: NormalizedOutboundPayload) => { hadPartialFailure = true; lastPayloadError = err; - partialFailuresArePreConnect &&= isPreConnectNetworkError(err); + partialFailuresAreProvenNotSent &&= isProvenDeliveryNotSentError(err); params.onError?.(err, payload); }, onDeliveryResult: async (result) => { @@ -1538,7 +1538,7 @@ async function deliverOutboundPayloadsWithQueueCleanup( const error = "partial delivery failure (bestEffort)"; if (postSendState === undefined || postSendState === "marked") { const recordFailure = - !partialSendEvidence && partialFailuresArePreConnect + !partialSendEvidence && partialFailuresAreProvenNotSent ? failDeliveryBeforePlatformSend : failDelivery; await recordFailure(queueId, error).catch((err: unknown) => { @@ -1634,7 +1634,7 @@ async function deliverOutboundPayloadsWithQueueCleanup( } await runCommitHooksAfterAck(); } else { - const recordFailure = isPreConnectNetworkError(err) + const recordFailure = isProvenDeliveryNotSentError(err) ? failDeliveryBeforePlatformSend : failDelivery; await recordFailure(queueId, formatErrorMessage(err)).catch((failErr: unknown) => { diff --git a/src/infra/outbound/delivery-queue-recovery.ts b/src/infra/outbound/delivery-queue-recovery.ts index 1b5e65fe1610..62a35b62b3da 100644 --- a/src/infra/outbound/delivery-queue-recovery.ts +++ b/src/infra/outbound/delivery-queue-recovery.ts @@ -14,7 +14,7 @@ import { computeBackoffMs, createRecoveryReplayPacer, getErrnoCode, - isPreConnectNetworkError, + isProvenDeliveryNotSentError, releaseRecoveryEntry as releaseSharedRecoveryEntry, } from "../delivery-recovery.shared.js"; import { formatErrorMessage } from "../errors.js"; @@ -569,7 +569,7 @@ async function drainQueuedEntry(opts: { } } else { const recordFailure = failedOutcomes.every((outcome) => - isPreConnectNetworkError(outcome.error), + isProvenDeliveryNotSentError(outcome.error), ) ? failDeliveryBeforePlatformSend : failDelivery; @@ -651,7 +651,7 @@ async function drainQueuedEntry(opts: { } } else { try { - const recordFailure = isPreConnectNetworkError(err) + const recordFailure = isProvenDeliveryNotSentError(err) ? failDeliveryBeforePlatformSend : failDelivery; await recordFailure(entry.id, errMsg, opts.stateDir); diff --git a/src/infra/outbound/delivery-queue-storage.ts b/src/infra/outbound/delivery-queue-storage.ts index ef1eae652ce2..d5cc4a715d4b 100644 --- a/src/infra/outbound/delivery-queue-storage.ts +++ b/src/infra/outbound/delivery-queue-storage.ts @@ -154,7 +154,7 @@ export async function failDelivery(id: string, error: string, stateDir?: string) })); } -/** Record a failed attempt that provably ended before any platform request was sent. */ +/** Record a failed attempt whose retry provably cannot duplicate a recipient-visible send. */ export async function failDeliveryBeforePlatformSend( id: string, error: string, @@ -208,7 +208,7 @@ export async function markDeliveryPlatformSendAttemptStarted( })); } -/** Refresh the attempt timestamp after provider serialization and immediately before I/O. */ +/** Refresh the attempt timestamp before recipient-visible or finalizing platform I/O. */ export async function markDeliveryPlatformSendDispatched( id: string, stateDir?: string, diff --git a/src/infra/outbound/delivery-queue.policy.test.ts b/src/infra/outbound/delivery-queue.policy.test.ts index ebb444d4d68f..c75457a2474d 100644 --- a/src/infra/outbound/delivery-queue.policy.test.ts +++ b/src/infra/outbound/delivery-queue.policy.test.ts @@ -1,6 +1,12 @@ // Covers delivery retry policy: permanent-error classification, backoff timing, // and first-replay eligibility after crashes. import { describe, expect, it } from "vitest"; +import { isProvenDeliveryNotSentError } from "../delivery-recovery.shared.js"; +import { recordRetryAttemptErrors } from "../retry-attempt-errors.js"; +import { + PlatformMessageNotDispatchedError, + PLATFORM_MESSAGE_NOT_DISPATCHED_ERROR_CODE, +} from "./deliver-types.js"; import { computeBackoffMs, isEntryEligibleForRecoveryRetry, @@ -8,6 +14,39 @@ import { } from "./delivery-queue.js"; describe("delivery-queue policy", () => { + describe("isProvenDeliveryNotSentError", () => { + const createMarker = () => + new PlatformMessageNotDispatchedError("upload stopped before finalization", { + cause: new Error("request timed out"), + }); + + it("accepts the channel-owned marker", () => { + expect(isProvenDeliveryNotSentError(createMarker())).toBe(true); + }); + + it("rejects a platform error that copies only the marker code", () => { + const forged = Object.assign(new Error("remote platform failure"), { + code: PLATFORM_MESSAGE_NOT_DISPATCHED_ERROR_CODE, + }); + expect(isProvenDeliveryNotSentError(forged)).toBe(false); + }); + + it.each(["connection reset after write", null])( + "rejects a marked aggregate with an unproven %s branch", + (unprovenBranch) => { + expect( + isProvenDeliveryNotSentError(new AggregateError([createMarker(), unprovenBranch])), + ).toBe(false); + }, + ); + + it("rejects a marker that follows an ambiguous retry attempt", () => { + const marker = createMarker(); + recordRetryAttemptErrors(marker, [new Error("connection reset after write"), marker]); + expect(isProvenDeliveryNotSentError(marker)).toBe(false); + }); + }); + describe("isPermanentDeliveryError", () => { it.each([ "No conversation reference found for user:abc", diff --git a/src/infra/outbound/delivery-queue.recovery.test.ts b/src/infra/outbound/delivery-queue.recovery.test.ts index cb30adb0db61..8519e23c880c 100644 --- a/src/infra/outbound/delivery-queue.recovery.test.ts +++ b/src/infra/outbound/delivery-queue.recovery.test.ts @@ -4,7 +4,11 @@ import { MAX_DATE_TIMESTAMP_MS } from "@openclaw/normalization-core/number-coerc import { beforeEach, describe, expect, it, vi } from "vitest"; import { openOpenClawStateDatabase } from "../../state/openclaw-state-db.js"; import { RECOVERY_REPLAY_SPACING_MS } from "../delivery-recovery.shared.js"; -import { OutboundDeliveryError, type OutboundPayloadDeliveryOutcome } from "./deliver-types.js"; +import { + OutboundDeliveryError, + PlatformMessageNotDispatchedError, + type OutboundPayloadDeliveryOutcome, +} from "./deliver-types.js"; import { attachOutboundDeliveryCommitHook } from "./delivery-commit-hooks.js"; import { enqueueDelivery, @@ -306,6 +310,28 @@ describe("delivery-queue recovery", () => { expect(entries[0]?.platformSendStartedAt).toBeUndefined(); }); + it("keeps a repeated provider-not-dispatched recovery failure replayable", async () => { + const id = await enqueueDelivery( + { channel: "demo-channel-c", to: "#ch", payloads: [{ text: "x" }] }, + tmpDir(), + ); + const deliver = vi.fn(async () => { + await markDeliveryPlatformSendAttemptStarted(id, tmpDir()); + throw new PlatformMessageNotDispatchedError("upload stopped before finalization", { + cause: new Error("request timed out"), + }); + }); + + const { result } = await runRecovery({ deliver }); + + expect(result).toMatchObject({ recovered: 0, failed: 1 }); + const entries = await loadPendingDeliveries(tmpDir()); + expect(entries).toHaveLength(1); + expect(entries[0]?.retryCount).toBe(1); + expect(entries[0]?.recoveryState).toBeUndefined(); + expect(entries[0]?.platformSendStartedAt).toBeUndefined(); + }); + it("does not replay a recovery batch that rejected after an earlier send succeeded", async () => { const id = await enqueueDelivery( { @@ -412,7 +438,46 @@ describe("delivery-queue recovery", () => { expect(entries[0]?.platformSendStartedAt).toBeUndefined(); }); - it("preserves send evidence when any best-effort recovery failure is ambiguous", async () => { + it("clears send evidence for an all-not-dispatched best-effort recovery failure", async () => { + const id = await enqueueDelivery( + { + channel: "demo-channel-c", + to: "#ch", + payloads: [{ text: "first" }], + bestEffort: true, + }, + tmpDir(), + ); + const deliver = vi.fn( + async (params: { + onPayloadDeliveryOutcome?: (outcome: OutboundPayloadDeliveryOutcome) => void; + }) => { + await markDeliveryPlatformSendAttemptStarted(id, tmpDir()); + params.onPayloadDeliveryOutcome?.({ + index: 0, + status: "failed", + error: new PlatformMessageNotDispatchedError( + "upload timed out before completion dispatch", + { cause: new Error("request timed out") }, + ), + sentBeforeError: false, + stage: "platform_send", + }); + return []; + }, + ); + + const { result } = await runRecovery({ deliver }); + + expect(result).toMatchObject({ recovered: 0, failed: 1 }); + const entries = await loadPendingDeliveries(tmpDir()); + expect(entries).toHaveLength(1); + expect(entries[0]?.retryCount).toBe(1); + expect(entries[0]?.recoveryState).toBeUndefined(); + expect(entries[0]?.platformSendStartedAt).toBeUndefined(); + }); + + it("preserves send evidence when a marked recovery batch has an ambiguous failure", async () => { const id = await enqueueDelivery( { channel: "demo-channel-c", @@ -430,10 +495,10 @@ describe("delivery-queue recovery", () => { params.onPayloadDeliveryOutcome?.({ index: 0, status: "failed", - error: Object.assign(new Error("connect ECONNREFUSED"), { - code: "ECONNREFUSED", - syscall: "connect", - }), + error: new PlatformMessageNotDispatchedError( + "upload timed out before completion dispatch", + { cause: new Error("request timed out") }, + ), sentBeforeError: false, stage: "platform_send", }); diff --git a/src/plugin-sdk/error-runtime.ts b/src/plugin-sdk/error-runtime.ts index d788335602d0..2e4cba29df0e 100644 --- a/src/plugin-sdk/error-runtime.ts +++ b/src/plugin-sdk/error-runtime.ts @@ -23,4 +23,5 @@ export { formatUncaughtError, readErrorName, } from "../infra/errors.js"; +export { PlatformMessageNotDispatchedError } from "../infra/outbound/deliver-types.js"; export { isApprovalNotFoundError } from "../infra/approval-errors.ts"; diff --git a/src/plugins/contracts/plugin-sdk-subpaths.test.ts b/src/plugins/contracts/plugin-sdk-subpaths.test.ts index 9ab27160a7ce..bf332f1d6a6a 100644 --- a/src/plugins/contracts/plugin-sdk-subpaths.test.ts +++ b/src/plugins/contracts/plugin-sdk-subpaths.test.ts @@ -1383,7 +1383,11 @@ describe("plugin-sdk subpath exports", () => { ); expectSourceMentions("delivery-queue-runtime", ["drainPendingDeliveries"]); expectSourceContains("delivery-queue-runtime", "../infra/outbound/deliver-runtime.js"); - expectSourceMentions("error-runtime", ["formatUncaughtError", "isApprovalNotFoundError"]); + expectSourceMentions("error-runtime", [ + "formatUncaughtError", + "isApprovalNotFoundError", + "PlatformMessageNotDispatchedError", + ]); expect(channelLifecycleSdk.createDraftStreamLoop).toBe( channelLifecycleDirectSdk.createDraftStreamLoop,