diff --git a/extensions/feishu/src/tool-account.test.ts b/extensions/feishu/src/tool-account.test.ts index f346dab83cd4..aca974186e5f 100644 --- a/extensions/feishu/src/tool-account.test.ts +++ b/extensions/feishu/src/tool-account.test.ts @@ -57,6 +57,26 @@ describe("resolveFeishuToolAccount", () => { expect(resolved.accountId).toBe("work"); }); + it("allows the explicit unlisted default backed by top-level credentials", () => { + const resolved = resolveFeishuToolAccount({ + api: { + config: { + channels: { + feishu: { + defaultAccount: "ops", + appId: "base-app-id", + appSecret: "base-app-secret", // pragma: allowlist secret + }, + }, + }, + }, + executeParams: { accountId: "OPS" }, + }); + + expect(resolved.accountId).toBe("ops"); + expect(resolved.configured).toBe(true); + }); + it.each([ { name: "malformed", accountId: "!!!", error: "Invalid Feishu account ID" }, { name: "unknown", accountId: "missing", error: "Unknown Feishu account" }, diff --git a/extensions/feishu/src/tool-account.ts b/extensions/feishu/src/tool-account.ts index 14de9f0329ab..b7b3dbfbefd6 100644 --- a/extensions/feishu/src/tool-account.ts +++ b/extensions/feishu/src/tool-account.ts @@ -5,6 +5,7 @@ import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runti import type { OpenClawPluginApi } from "../runtime-api.js"; import { listFeishuAccountIds, + resolveDefaultFeishuAccountId, resolveFeishuAccount, resolveFeishuRuntimeAccount, } from "./accounts.js"; @@ -31,9 +32,16 @@ function resolveImplicitToolAccountId(params: { if (!normalizedAccountId) { throw new Error(`Invalid Feishu account ID "${explicitAccountId}"`); } - const listedAccountId = listFeishuAccountIds(params.api.config).find( - (accountId) => normalizeOptionalAccountId(accountId) === normalizedAccountId, - ); + const listedAccountId = + listFeishuAccountIds(params.api.config).find( + (accountId) => normalizeOptionalAccountId(accountId) === normalizedAccountId, + ) ?? + (() => { + const defaultAccountId = resolveDefaultFeishuAccountId(params.api.config); + return normalizeOptionalAccountId(defaultAccountId) === normalizedAccountId + ? defaultAccountId + : undefined; + })(); if (!listedAccountId) { throw new Error(`Unknown Feishu account "${explicitAccountId}"`); } diff --git a/src/gateway/server-methods/send.test.ts b/src/gateway/server-methods/send.test.ts index 226323120bbd..cce29b53a3c3 100644 --- a/src/gateway/server-methods/send.test.ts +++ b/src/gateway/server-methods/send.test.ts @@ -1034,8 +1034,9 @@ describe("gateway send mirroring", () => { isWebchatConnect: () => false, }); - await Promise.resolve(); - expect(mocks.dispatchChannelMessageAction).toHaveBeenCalledTimes(2); + await vi.waitFor(() => { + expect(mocks.dispatchChannelMessageAction).toHaveBeenCalledTimes(2); + }); expect(mocks.dispatchChannelMessageAction.mock.calls[0]?.[0]).toMatchObject({ conversationReadOrigin: "direct-operator", }); @@ -1441,6 +1442,85 @@ describe("gateway send mirroring", () => { expect(firstRespondCall(retryRespond)?.[3]?.cached).toBe(true); }); + it.each([ + { + name: "send", + method: "send" as const, + request: { + to: "channel:C1", + message: "hi", + idempotencyKey: "idem-send-deferred-route-race", + }, + providerCall: mocks.deliverOutboundPayloads, + }, + { + name: "poll", + method: "poll" as const, + request: { + to: "channel:C1", + question: "Q?", + options: ["A", "B"], + idempotencyKey: "idem-poll-deferred-route-race", + }, + providerCall: mocks.sendPoll, + }, + ])( + "keeps the first deferred $name route when a concurrent retry sees newer defaults", + async (testCase) => { + const firstSelection = createDeferred<{ channel: string; configured: string[] }>(); + mocks.resolveMessageChannelSelection + .mockImplementationOnce(async () => await firstSelection.promise) + .mockResolvedValue({ channel: "discord", configured: ["discord"] }); + mockMutableMessageRouteAccounts(() => "primary"); + + const providerDeferred = createDeferred(); + if (testCase.method === "send") { + mocks.deliverOutboundPayloads.mockReturnValueOnce(providerDeferred.promise as never); + } else { + mocks.sendPoll.mockReturnValueOnce(providerDeferred.promise as never); + } + + const context = makeContext(); + const firstRespond = vi.fn(); + const retryRespond = vi.fn(); + const firstRequest = invokeGatewayMessageMethod({ + method: testCase.method, + request: testCase.request, + respond: firstRespond, + context, + }); + await vi.waitFor(() => { + expect(mocks.resolveMessageChannelSelection).toHaveBeenCalledTimes(1); + }); + + const retryRequest = invokeGatewayMessageMethod({ + method: testCase.method, + request: testCase.request, + respond: retryRespond, + context, + }); + await Promise.resolve(); + expect(mocks.resolveMessageChannelSelection).toHaveBeenCalledTimes(1); + + firstSelection.resolve({ channel: "slack", configured: ["slack"] }); + await vi.waitFor(() => { + expect(testCase.providerCall).toHaveBeenCalledTimes(1); + }); + if (testCase.method === "send") { + providerDeferred.resolve([{ messageId: "m-race", channel: "slack" }]); + } else { + providerDeferred.resolve({ messageId: "poll-race", pollId: "poll-race" }); + } + await Promise.all([firstRequest, retryRequest]); + + expect(mocks.resolveMessageChannelSelection).toHaveBeenCalledTimes(1); + expect(testCase.providerCall).toHaveBeenCalledTimes(1); + expect(firstRespondCall(firstRespond)?.[0]).toBe(true); + expect(firstRespondCall(retryRespond)?.[0]).toBe(true); + expect(firstRespondCall(retryRespond)?.[3]?.cached).toBe(true); + }, + ); + it("dedupes omitted and explicit default poll routes", async () => { const context = makeContext(); const firstRespond = vi.fn(); diff --git a/src/gateway/server-methods/send.ts b/src/gateway/server-methods/send.ts index db90f532cae4..2df88a5aec54 100644 --- a/src/gateway/server-methods/send.ts +++ b/src/gateway/server-methods/send.ts @@ -100,6 +100,54 @@ type MessageOperationRouteBinding = { reservedRoute?: MessageOperationRoute; }; +const messageOperationRouteBindingQueues = new WeakMap(); + +function getMessageOperationRouteBindingQueue(context: GatewayRequestContext): KeyedAsyncQueue { + let queue = messageOperationRouteBindingQueues.get(context); + if (!queue) { + queue = new KeyedAsyncQueue(); + messageOperationRouteBindingQueues.set(context, queue); + } + return queue; +} + +async function acquireMessageOperationRouteBindingLock(params: { + context: GatewayRequestContext; + binding: MessageOperationRouteBinding | undefined; +}): Promise<() => void> { + if (!params.binding) { + return () => undefined; + } + + let signalAcquired: (() => void) | undefined; + let signalRelease: (() => void) | undefined; + const acquired = new Promise((resolve) => { + signalAcquired = resolve; + }); + const held = new Promise((resolve) => { + signalRelease = resolve; + }); + // The lock covers mutable route selection through canonical in-flight registration. + // Otherwise a later retry can bind newer defaults while the first request is resolving. + void getMessageOperationRouteBindingQueue(params.context).enqueue( + params.binding.key, + async () => { + signalAcquired?.(); + await held; + }, + ); + await acquired; + + let released = false; + return () => { + if (released) { + return; + } + released = true; + signalRelease?.(); + }; +} + function resolveTrustedMessageActionToolContext(params: { client: Parameters[0]["client"]; request: { @@ -718,7 +766,7 @@ export const sendHandlers: GatewayRequestHandlers = { client, requestedOrigin: request.conversationReadOrigin, }); - const routeBinding = resolveMessageOperationRouteBinding({ + let routeBinding = resolveMessageOperationRouteBinding({ context, prefix: "message.action", idempotencyKey: request.idempotencyKey, @@ -726,201 +774,221 @@ export const sendHandlers: GatewayRequestHandlers = { requestChannel: request.channel, accountIds: [request.accountId, request.params.accountId], }); - const reservedReplay = replayReservedMessageOperationRoute({ + const releaseRouteBindingLock = await acquireMessageOperationRouteBindingLock({ context, binding: routeBinding, - prefix: "message.action", - idempotencyKey: request.idempotencyKey, - respond, - conversationReadOrigin, }); - if (reservedReplay) { - await reservedReplay; - return; - } - const resolvedChannel = await resolveRequestedChannel({ - requestChannel: routeBinding?.reservedRoute?.channel ?? request.channel, - unsupportedMessage: (input) => `unsupported channel: ${input}`, - context, - rejectWebchatAsInternalOnly: true, - }); - if ("error" in resolvedChannel) { - respond(false, undefined, resolvedChannel.error); - return; - } - const { cfg: selectedCfg, sourceCfg, channel } = resolvedChannel; - const cfg = resolveMessageActionRuntimeConfig({ cfg: selectedCfg, sourceCfg }); - const plugin = resolveOutboundChannelPlugin({ channel, cfg }); - if (!plugin?.actions?.handleAction) { - respond( - false, - undefined, - errorShape( - ErrorCodes.INVALID_REQUEST, - `Channel ${channel} does not support action ${request.action}.`, - ), - ); - return; - } - let accountRoute: ReturnType; try { - accountRoute = resolveMessageOperationAccountRoute({ - cfg, - channel, - plugin, - accountIds: [ - request.accountId, - request.params.accountId, - routeBinding?.reservedRoute?.accountId, - ], - conflictMessage: "message.action accountId does not match params.accountId", + routeBinding = resolveMessageOperationRouteBinding({ + context, + prefix: "message.action", + idempotencyKey: request.idempotencyKey, + conversationReadOrigin, + requestChannel: request.channel, + accountIds: [request.accountId, request.params.accountId], }); - } catch (error) { - respondGatewayInvalidRequest({ respond, channel, error }); - return; - } - if ( - !bindMessageOperationRoute({ + const reservedReplay = replayReservedMessageOperationRoute({ context, binding: routeBinding, - requestScope: accountRoute.requestScope, - }) - ) { - respondGatewayInvalidRequest({ + prefix: "message.action", + idempotencyKey: request.idempotencyKey, respond, - channel, - error: "idempotency key is already bound to a different message route", + conversationReadOrigin, }); - return; - } - const inflight = resolveGatewayInflightRequest({ - context, - prefix: "message.action", - idempotencyKey: request.idempotencyKey, - respond, - conversationReadOrigin, - requestScope: accountRoute.requestScope, - }); - if (inflight.kind === "handled") { - await inflight.done; - return; - } - const { dedupeKey, inflightMap } = inflight; - retainMessageOperationRouteBinding({ - context, - binding: routeBinding, - requestScope: accountRoute.requestScope, - }); - const work = (async (): Promise => { - try { - const sessionKey = normalizeOptionalString(request.sessionKey) ?? undefined; - const agentId = - normalizeOptionalString(request.agentId) ?? - (sessionKey ? resolveSessionAgentId({ sessionKey, config: cfg }) : undefined); - const accountId = accountRoute.accountId; - if (accountId) { - request.params.accountId = accountId; - } - if (request.action === "send") { - await hydrateAttachmentParamsForAction({ - cfg, - channel, - accountId, - args: request.params, - action: "send", - mediaPolicy: resolveAttachmentMediaPolicy({ - mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, agentId), - }), - }); - } - const sourceReplyMirror = { - action: request.action, - channel, - actionParams: request.params, - cfg, - accountId, - currentAccountId: trustedContext.requesterAccountId, - sessionKey, - sessionId: trustedContext.sessionId, - agentId, - toolContext: trustedContext.toolContext, - idempotencyKey: request.idempotencyKey, - toolCallId: trustedContext.sourceReplyToolCallId, - ...(trustedContext.sourceReplyFinal !== undefined - ? { sourceReplyFinal: trustedContext.sourceReplyFinal } - : {}), - }; - const terminalDeliveryReceipt = - trustedContext.sourceReplyFinal === true - ? await beginTerminalSourceReplyDelivery(sourceReplyMirror) - : undefined; - const gatewayClientScopes = client?.connect?.scopes ?? []; - const handled = await dispatchChannelMessageAction({ - channel, - action: request.action as never, - cfg, - params: request.params, - accountId, - requesterAccountId: trustedContext.requesterAccountId, - requesterSenderId: trustedContext.requesterSenderId, - senderIsOwner: gatewayClientScopes.includes(ADMIN_SCOPE) - ? request.senderIsOwner === true - : false, - conversationReadOrigin, - sessionKey, - sessionId: normalizeOptionalString(request.sessionId) ?? undefined, - inboundEventKind: request.inboundTurnKind, - agentId, - mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, agentId), - toolContext: trustedContext.toolContext, - dryRun: false, - gatewayClientScopes, - }); - if (!handled) { - await cancelTerminalSourceReplyDelivery(terminalDeliveryReceipt); - const error = errorShape( - ErrorCodes.INVALID_REQUEST, - `Message action ${request.action} not supported for channel ${channel}.`, - ); - cacheGatewayDedupeFailure({ context, dedupeKey, error }); - return { ok: false, error, meta: { channel } }; - } - const payload = extractToolPayload(handled); - try { - await reconcileTerminalSourceReplyDelivery({ - deliveredPayload: payload, - mirror: sourceReplyMirror, - receipt: terminalDeliveryReceipt, - }); - } catch (err) { - // The pre-send intent remains durable. Return the provider result so - // the model does not retry an external effect with an unknown outcome. - context.logGateway?.warn?.("Terminal source reply receipt reconciliation failed.", { - error: formatForLog(err), - channel, - sessionKey, - }); - } - await scheduleDeliveredSourceReplyTranscriptMirror({ - context, - mirror: { - ...sourceReplyMirror, - deliveredPayload: payload, - }, - }); - return createGatewayInflightSuccess({ context, dedupeKey, payload, channel }); - } catch (err) { - return createGatewayInflightUnavailableFailure({ context, dedupeKey, channel, err }); + if (reservedReplay) { + releaseRouteBindingLock(); + await reservedReplay; + return; } - })().finally(() => { - refreshMessageOperationRouteBinding({ + const resolvedChannel = await resolveRequestedChannel({ + requestChannel: routeBinding?.reservedRoute?.channel ?? request.channel, + unsupportedMessage: (input) => `unsupported channel: ${input}`, + context, + rejectWebchatAsInternalOnly: true, + }); + if ("error" in resolvedChannel) { + respond(false, undefined, resolvedChannel.error); + return; + } + const { cfg: selectedCfg, sourceCfg, channel } = resolvedChannel; + const cfg = resolveMessageActionRuntimeConfig({ cfg: selectedCfg, sourceCfg }); + const plugin = resolveOutboundChannelPlugin({ channel, cfg }); + if (!plugin?.actions?.handleAction) { + respond( + false, + undefined, + errorShape( + ErrorCodes.INVALID_REQUEST, + `Channel ${channel} does not support action ${request.action}.`, + ), + ); + return; + } + let accountRoute: ReturnType; + try { + accountRoute = resolveMessageOperationAccountRoute({ + cfg, + channel, + plugin, + accountIds: [ + request.accountId, + request.params.accountId, + routeBinding?.reservedRoute?.accountId, + ], + conflictMessage: "message.action accountId does not match params.accountId", + }); + } catch (error) { + respondGatewayInvalidRequest({ respond, channel, error }); + return; + } + if ( + !bindMessageOperationRoute({ + context, + binding: routeBinding, + requestScope: accountRoute.requestScope, + }) + ) { + respondGatewayInvalidRequest({ + respond, + channel, + error: "idempotency key is already bound to a different message route", + }); + return; + } + const inflight = resolveGatewayInflightRequest({ + context, + prefix: "message.action", + idempotencyKey: request.idempotencyKey, + respond, + conversationReadOrigin, + requestScope: accountRoute.requestScope, + }); + if (inflight.kind === "handled") { + releaseRouteBindingLock(); + await inflight.done; + return; + } + const { dedupeKey, inflightMap } = inflight; + retainMessageOperationRouteBinding({ context, binding: routeBinding, requestScope: accountRoute.requestScope, }); - }); + const work = (async (): Promise => { + try { + const sessionKey = normalizeOptionalString(request.sessionKey) ?? undefined; + const agentId = + normalizeOptionalString(request.agentId) ?? + (sessionKey ? resolveSessionAgentId({ sessionKey, config: cfg }) : undefined); + const accountId = accountRoute.accountId; + if (accountId) { + request.params.accountId = accountId; + } + if (request.action === "send") { + await hydrateAttachmentParamsForAction({ + cfg, + channel, + accountId, + args: request.params, + action: "send", + mediaPolicy: resolveAttachmentMediaPolicy({ + mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, agentId), + }), + }); + } + const sourceReplyMirror = { + action: request.action, + channel, + actionParams: request.params, + cfg, + accountId, + currentAccountId: trustedContext.requesterAccountId, + sessionKey, + sessionId: trustedContext.sessionId, + agentId, + toolContext: trustedContext.toolContext, + idempotencyKey: request.idempotencyKey, + toolCallId: trustedContext.sourceReplyToolCallId, + ...(trustedContext.sourceReplyFinal !== undefined + ? { sourceReplyFinal: trustedContext.sourceReplyFinal } + : {}), + }; + const terminalDeliveryReceipt = + trustedContext.sourceReplyFinal === true + ? await beginTerminalSourceReplyDelivery(sourceReplyMirror) + : undefined; + const gatewayClientScopes = client?.connect?.scopes ?? []; + const handled = await dispatchChannelMessageAction({ + channel, + action: request.action as never, + cfg, + params: request.params, + accountId, + requesterAccountId: trustedContext.requesterAccountId, + requesterSenderId: trustedContext.requesterSenderId, + senderIsOwner: gatewayClientScopes.includes(ADMIN_SCOPE) + ? request.senderIsOwner === true + : false, + conversationReadOrigin, + sessionKey, + sessionId: normalizeOptionalString(request.sessionId) ?? undefined, + inboundEventKind: request.inboundTurnKind, + agentId, + mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, agentId), + toolContext: trustedContext.toolContext, + dryRun: false, + gatewayClientScopes, + }); + if (!handled) { + await cancelTerminalSourceReplyDelivery(terminalDeliveryReceipt); + const error = errorShape( + ErrorCodes.INVALID_REQUEST, + `Message action ${request.action} not supported for channel ${channel}.`, + ); + cacheGatewayDedupeFailure({ context, dedupeKey, error }); + return { ok: false, error, meta: { channel } }; + } + const payload = extractToolPayload(handled); + try { + await reconcileTerminalSourceReplyDelivery({ + deliveredPayload: payload, + mirror: sourceReplyMirror, + receipt: terminalDeliveryReceipt, + }); + } catch (err) { + // The pre-send intent remains durable. Return the provider result so + // the model does not retry an external effect with an unknown outcome. + context.logGateway?.warn?.("Terminal source reply receipt reconciliation failed.", { + error: formatForLog(err), + channel, + sessionKey, + }); + } + await scheduleDeliveredSourceReplyTranscriptMirror({ + context, + mirror: { + ...sourceReplyMirror, + deliveredPayload: payload, + }, + }); + return createGatewayInflightSuccess({ context, dedupeKey, payload, channel }); + } catch (err) { + return createGatewayInflightUnavailableFailure({ context, dedupeKey, channel, err }); + } + })().finally(() => { + refreshMessageOperationRouteBinding({ + context, + binding: routeBinding, + requestScope: accountRoute.requestScope, + }); + }); - await runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + const inflightWork = runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + releaseRouteBindingLock(); + await inflightWork; + } finally { + releaseRouteBindingLock(); + } }, send: async ({ params, respond, context, client }) => { const p = params; @@ -968,276 +1036,295 @@ export const sendHandlers: GatewayRequestHandlers = { const requestedAccountId = normalizeOptionalString(request.accountId); const replyToId = normalizeOptionalString(request.replyToId); const threadId = normalizeOptionalString(request.threadId); - const routeBinding = resolveMessageOperationRouteBinding({ + let routeBinding = resolveMessageOperationRouteBinding({ context, prefix: "send", idempotencyKey: request.idempotencyKey, requestChannel: request.channel, accountIds: [request.accountId], }); - const reservedReplay = replayReservedMessageOperationRoute({ + const releaseRouteBindingLock = await acquireMessageOperationRouteBindingLock({ context, binding: routeBinding, - prefix: "send", - idempotencyKey: request.idempotencyKey, - respond, }); - if (reservedReplay) { - await reservedReplay; - return; - } - const resolvedChannel = await resolveInternalDeliveryChannel( - routeBinding?.reservedRoute?.channel ?? request.channel, - context, - ); - if (resolvedChannel.kind !== "ready") { - const result = resolvedChannel.result; - respond(result.ok, result.payload, result.error, result.meta); - return; - } - const { cfg, channel } = resolvedChannel; - const outboundChannel = channel; - const plugin = resolveOutboundChannelPlugin({ channel, cfg }); - if (!plugin) { - respond( - false, - undefined, - errorShape(ErrorCodes.INVALID_REQUEST, `unsupported channel: ${channel}`), - ); - return; - } - let accountRoute: ReturnType; try { - accountRoute = resolveMessageOperationAccountRoute({ - cfg, - channel, - plugin, - accountIds: [requestedAccountId, routeBinding?.reservedRoute?.accountId], - conflictMessage: "send account selections do not match", + routeBinding = resolveMessageOperationRouteBinding({ + context, + prefix: "send", + idempotencyKey: request.idempotencyKey, + requestChannel: request.channel, + accountIds: [request.accountId], }); - } catch (error) { - respondGatewayInvalidRequest({ respond, channel, error }); - return; - } - if ( - !bindMessageOperationRoute({ + const reservedReplay = replayReservedMessageOperationRoute({ + context, + binding: routeBinding, + prefix: "send", + idempotencyKey: request.idempotencyKey, + respond, + }); + if (reservedReplay) { + releaseRouteBindingLock(); + await reservedReplay; + return; + } + const resolvedChannel = await resolveInternalDeliveryChannel( + routeBinding?.reservedRoute?.channel ?? request.channel, + context, + ); + if (resolvedChannel.kind !== "ready") { + const result = resolvedChannel.result; + respond(result.ok, result.payload, result.error, result.meta); + return; + } + const { cfg, channel } = resolvedChannel; + const outboundChannel = channel; + const plugin = resolveOutboundChannelPlugin({ channel, cfg }); + if (!plugin) { + respond( + false, + undefined, + errorShape(ErrorCodes.INVALID_REQUEST, `unsupported channel: ${channel}`), + ); + return; + } + let accountRoute: ReturnType; + try { + accountRoute = resolveMessageOperationAccountRoute({ + cfg, + channel, + plugin, + accountIds: [requestedAccountId, routeBinding?.reservedRoute?.accountId], + conflictMessage: "send account selections do not match", + }); + } catch (error) { + respondGatewayInvalidRequest({ respond, channel, error }); + return; + } + if ( + !bindMessageOperationRoute({ + context, + binding: routeBinding, + requestScope: accountRoute.requestScope, + }) + ) { + respondGatewayInvalidRequest({ + respond, + channel, + error: "idempotency key is already bound to a different message route", + }); + return; + } + const accountId = accountRoute.accountId; + const inflight = resolveGatewayInflightRequest({ + context, + prefix: "send", + idempotencyKey: request.idempotencyKey, + respond, + requestScope: accountRoute.requestScope, + }); + if (inflight.kind === "handled") { + releaseRouteBindingLock(); + await inflight.done; + return; + } + const { idem, dedupeKey, inflightMap } = inflight; + retainMessageOperationRouteBinding({ context, binding: routeBinding, requestScope: accountRoute.requestScope, - }) - ) { - respondGatewayInvalidRequest({ - respond, - channel, - error: "idempotency key is already bound to a different message route", }); - return; - } - const accountId = accountRoute.accountId; - const inflight = resolveGatewayInflightRequest({ - context, - prefix: "send", - idempotencyKey: request.idempotencyKey, - respond, - requestScope: accountRoute.requestScope, - }); - if (inflight.kind === "handled") { - await inflight.done; - return; - } - const { idem, dedupeKey, inflightMap } = inflight; - retainMessageOperationRouteBinding({ - context, - binding: routeBinding, - requestScope: accountRoute.requestScope, - }); - const work = (async (): Promise => { - try { - const resolvedTarget = resolveGatewayOutboundTarget({ - channel: outboundChannel, - to, - cfg, - accountId, - }); - if (!resolvedTarget.ok) { - return { - ok: false, - error: resolvedTarget.error, - meta: { channel }, - }; - } - const idLikeTarget = await maybeResolveIdLikeTarget({ - cfg, - channel, - input: resolvedTarget.to, - accountId, - }); - const deliveryTarget = idLikeTarget?.to ?? resolvedTarget.to; - // Preserve opaque, case-sensitive peer IDs (e.g. Matrix room ids) on an - // explicit session key instead of raw-lowercasing it (openclaw#75670). - // Non-enrolled channels still canonicalize to lowercase via the registry. - const providedSessionKey = - normalizeSessionKeyPreservingOpaquePeerIds(request.sessionKey) || undefined; - const explicitAgentId = normalizeOptionalString(request.agentId); - const sessionAgentId = providedSessionKey - ? resolveSessionAgentId({ sessionKey: providedSessionKey, config: cfg }) - : undefined; - const defaultAgentId = resolveSessionAgentId({ config: cfg }); - const effectiveAgentId = explicitAgentId ?? sessionAgentId ?? defaultAgentId; - const sendArgs: Record = { - mediaUrl, - mediaUrls, - buffer, - filename: normalizeOptionalString(request.filename) ?? undefined, - contentType: normalizeOptionalString(request.contentType) ?? undefined, - }; - await hydrateAttachmentParamsForAction({ - cfg, - channel, - accountId, - args: sendArgs, - action: "send", - mediaPolicy: resolveAttachmentMediaPolicy({ - mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, effectiveAgentId), - }), - }); - const hydratedMediaUrl = normalizeOptionalString(sendArgs.mediaUrl); - const hydratedMediaUrls = Array.isArray(sendArgs.mediaUrls) - ? sendArgs.mediaUrls - .map((entry) => normalizeOptionalString(entry)) - .filter((entry): entry is string => Boolean(entry)) - : undefined; - const outboundDeps = context.deps ? createOutboundSendDeps(context.deps) : undefined; - const outboundPayloads = [ - { - text: message, - mediaUrl: hydratedMediaUrl, - mediaUrls: hydratedMediaUrls, - ...(request.asVoice === true ? { audioAsVoice: true } : {}), - }, - ]; - const outboundPayloadPlan = createOutboundPayloadPlan(outboundPayloads); - const mirrorProjection = projectOutboundPayloadPlanForMirror(outboundPayloadPlan); - const mirrorText = mirrorProjection.text; - const mirrorMediaUrls = mirrorProjection.mediaUrls; - const derivedRoute = await resolveOutboundSessionRoute({ - cfg, - channel, - agentId: effectiveAgentId, - accountId, - target: deliveryTarget, - currentSessionKey: providedSessionKey, - resolvedTarget: idLikeTarget, - replyToId, - threadId, - }); - const providedSessionBaseKey = - parseThreadSessionSuffix(providedSessionKey).baseSessionKey ?? providedSessionKey; - const shouldUseDerivedThreadSessionKey = - resolveChannelThreadAddressing(channel) === "message" && - Boolean(providedSessionKey) && - Boolean(normalizeOptionalString(derivedRoute?.threadId)) && - normalizeOptionalLowercaseString(derivedRoute?.baseSessionKey) === - normalizeOptionalLowercaseString(providedSessionBaseKey) && - normalizeOptionalLowercaseString(derivedRoute?.sessionKey) !== providedSessionKey; - // Message-scoped threads can refine an existing base session only after target lookup. - const outboundRoute = derivedRoute - ? providedSessionKey - ? shouldUseDerivedThreadSessionKey - ? { - ...derivedRoute, - baseSessionKey: derivedRoute.baseSessionKey ?? providedSessionKey, - } - : { - ...derivedRoute, - sessionKey: providedSessionKey, - baseSessionKey: providedSessionKey, - } - : derivedRoute - : null; - const outboundSessionKey = outboundRoute?.sessionKey ?? providedSessionKey; - if (outboundSessionKey && isAgentHarnessSessionKey(outboundSessionKey)) { - const { canonicalKey, entry } = loadSessionEntry(outboundSessionKey); - const missingHarnessSessionError = resolveMissingAgentHarnessSessionError( - canonicalKey, - entry, - ); - if (missingHarnessSessionError) { + const work = (async (): Promise => { + try { + const resolvedTarget = resolveGatewayOutboundTarget({ + channel: outboundChannel, + to, + cfg, + accountId, + }); + if (!resolvedTarget.ok) { return { ok: false, - error: errorShape(ErrorCodes.INVALID_REQUEST, missingHarnessSessionError), + error: resolvedTarget.error, meta: { channel }, }; } - } - if (outboundRoute) { - await ensureOutboundSessionEntry({ + const idLikeTarget = await maybeResolveIdLikeTarget({ + cfg, + channel, + input: resolvedTarget.to, + accountId, + }); + const deliveryTarget = idLikeTarget?.to ?? resolvedTarget.to; + // Preserve opaque, case-sensitive peer IDs (e.g. Matrix room ids) on an + // explicit session key instead of raw-lowercasing it (openclaw#75670). + // Non-enrolled channels still canonicalize to lowercase via the registry. + const providedSessionKey = + normalizeSessionKeyPreservingOpaquePeerIds(request.sessionKey) || undefined; + const explicitAgentId = normalizeOptionalString(request.agentId); + const sessionAgentId = providedSessionKey + ? resolveSessionAgentId({ sessionKey: providedSessionKey, config: cfg }) + : undefined; + const defaultAgentId = resolveSessionAgentId({ config: cfg }); + const effectiveAgentId = explicitAgentId ?? sessionAgentId ?? defaultAgentId; + const sendArgs: Record = { + mediaUrl, + mediaUrls, + buffer, + filename: normalizeOptionalString(request.filename) ?? undefined, + contentType: normalizeOptionalString(request.contentType) ?? undefined, + }; + await hydrateAttachmentParamsForAction({ cfg, channel, accountId, - route: outboundRoute, + args: sendArgs, + action: "send", + mediaPolicy: resolveAttachmentMediaPolicy({ + mediaLocalRoots: getAgentScopedMediaLocalRoots(cfg, effectiveAgentId), + }), }); - } - const outboundSession = buildOutboundSessionContext({ - cfg, - agentId: effectiveAgentId, - sessionKey: outboundSessionKey, - conversationType: outboundRoute?.chatType, - }); - const send = await sendDurableMessageBatch({ - cfg, - channel: outboundChannel, - to: deliveryTarget, - accountId, - payloads: outboundPayloads, - replyToId: replyToId ?? null, - session: outboundSession, - gifPlayback: request.gifPlayback, - forceDocument: request.forceDocument, - threadId: outboundRoute?.threadId ?? threadId ?? null, - deps: outboundDeps, - gatewayClientScopes: client?.connect?.scopes ?? [], - silent: request.silent, - formatting: request.parseMode ? { parseMode: request.parseMode } : undefined, - mirror: outboundSessionKey - ? { - sessionKey: outboundSessionKey, - agentId: effectiveAgentId, - text: mirrorText || message, - mediaUrls: mirrorMediaUrls.length > 0 ? mirrorMediaUrls : undefined, - idempotencyKey: idem, - } - : undefined, - }); - if (send.status === "failed" || send.status === "partial_failed") { - throw send.error; - } - const results = send.status === "sent" ? send.results : []; + const hydratedMediaUrl = normalizeOptionalString(sendArgs.mediaUrl); + const hydratedMediaUrls = Array.isArray(sendArgs.mediaUrls) + ? sendArgs.mediaUrls + .map((entry) => normalizeOptionalString(entry)) + .filter((entry): entry is string => Boolean(entry)) + : undefined; + const outboundDeps = context.deps ? createOutboundSendDeps(context.deps) : undefined; + const outboundPayloads = [ + { + text: message, + mediaUrl: hydratedMediaUrl, + mediaUrls: hydratedMediaUrls, + ...(request.asVoice === true ? { audioAsVoice: true } : {}), + }, + ]; + const outboundPayloadPlan = createOutboundPayloadPlan(outboundPayloads); + const mirrorProjection = projectOutboundPayloadPlanForMirror(outboundPayloadPlan); + const mirrorText = mirrorProjection.text; + const mirrorMediaUrls = mirrorProjection.mediaUrls; + const derivedRoute = await resolveOutboundSessionRoute({ + cfg, + channel, + agentId: effectiveAgentId, + accountId, + target: deliveryTarget, + currentSessionKey: providedSessionKey, + resolvedTarget: idLikeTarget, + replyToId, + threadId, + }); + const providedSessionBaseKey = + parseThreadSessionSuffix(providedSessionKey).baseSessionKey ?? providedSessionKey; + const shouldUseDerivedThreadSessionKey = + resolveChannelThreadAddressing(channel) === "message" && + Boolean(providedSessionKey) && + Boolean(normalizeOptionalString(derivedRoute?.threadId)) && + normalizeOptionalLowercaseString(derivedRoute?.baseSessionKey) === + normalizeOptionalLowercaseString(providedSessionBaseKey) && + normalizeOptionalLowercaseString(derivedRoute?.sessionKey) !== providedSessionKey; + // Message-scoped threads can refine an existing base session only after target lookup. + const outboundRoute = derivedRoute + ? providedSessionKey + ? shouldUseDerivedThreadSessionKey + ? { + ...derivedRoute, + baseSessionKey: derivedRoute.baseSessionKey ?? providedSessionKey, + } + : { + ...derivedRoute, + sessionKey: providedSessionKey, + baseSessionKey: providedSessionKey, + } + : derivedRoute + : null; + const outboundSessionKey = outboundRoute?.sessionKey ?? providedSessionKey; + if (outboundSessionKey && isAgentHarnessSessionKey(outboundSessionKey)) { + const { canonicalKey, entry } = loadSessionEntry(outboundSessionKey); + const missingHarnessSessionError = resolveMissingAgentHarnessSessionError( + canonicalKey, + entry, + ); + if (missingHarnessSessionError) { + return { + ok: false, + error: errorShape(ErrorCodes.INVALID_REQUEST, missingHarnessSessionError), + meta: { channel }, + }; + } + } + if (outboundRoute) { + await ensureOutboundSessionEntry({ + cfg, + channel, + accountId, + route: outboundRoute, + }); + } + const outboundSession = buildOutboundSessionContext({ + cfg, + agentId: effectiveAgentId, + sessionKey: outboundSessionKey, + conversationType: outboundRoute?.chatType, + }); + const send = await sendDurableMessageBatch({ + cfg, + channel: outboundChannel, + to: deliveryTarget, + accountId, + payloads: outboundPayloads, + replyToId: replyToId ?? null, + session: outboundSession, + gifPlayback: request.gifPlayback, + forceDocument: request.forceDocument, + threadId: outboundRoute?.threadId ?? threadId ?? null, + deps: outboundDeps, + gatewayClientScopes: client?.connect?.scopes ?? [], + silent: request.silent, + formatting: request.parseMode ? { parseMode: request.parseMode } : undefined, + mirror: outboundSessionKey + ? { + sessionKey: outboundSessionKey, + agentId: effectiveAgentId, + text: mirrorText || message, + mediaUrls: mirrorMediaUrls.length > 0 ? mirrorMediaUrls : undefined, + idempotencyKey: idem, + } + : undefined, + }); + if (send.status === "failed" || send.status === "partial_failed") { + throw send.error; + } + const results = send.status === "sent" ? send.results : []; - const result = results.at(-1); - if (!result) { - throw new Error("No delivery result"); + const result = results.at(-1); + if (!result) { + throw new Error("No delivery result"); + } + return createGatewayDeliveryInflightSuccess({ + context, + dedupeKey, + runId: idem, + channel, + result, + }); + } catch (err) { + return createGatewayInflightUnavailableFailure({ context, dedupeKey, channel, err }); } - return createGatewayDeliveryInflightSuccess({ + })().finally(() => { + refreshMessageOperationRouteBinding({ context, - dedupeKey, - runId: idem, - channel, - result, + binding: routeBinding, + requestScope: accountRoute.requestScope, }); - } catch (err) { - return createGatewayInflightUnavailableFailure({ context, dedupeKey, channel, err }); - } - })().finally(() => { - refreshMessageOperationRouteBinding({ - context, - binding: routeBinding, - requestScope: accountRoute.requestScope, }); - }); - await runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + const inflightWork = runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + releaseRouteBindingLock(); + await inflightWork; + } finally { + releaseRouteBindingLock(); + } }, poll: async ({ params, respond, context, client }) => { const p = params; @@ -1258,170 +1345,192 @@ export const sendHandlers: GatewayRequestHandlers = { accountId?: string; idempotencyKey: string; }; - const routeBinding = resolveMessageOperationRouteBinding({ + let routeBinding = resolveMessageOperationRouteBinding({ context, prefix: "poll", idempotencyKey: request.idempotencyKey, requestChannel: request.channel, accountIds: [request.accountId], }); - const reservedReplay = replayReservedMessageOperationRoute({ + const releaseRouteBindingLock = await acquireMessageOperationRouteBindingLock({ context, binding: routeBinding, - prefix: "poll", - idempotencyKey: request.idempotencyKey, - respond, }); - if (reservedReplay) { - await reservedReplay; - return; - } - const resolvedChannel = await resolveRequestedChannel({ - requestChannel: routeBinding?.reservedRoute?.channel ?? request.channel, - unsupportedMessage: (input) => `unsupported poll channel: ${input}`, - context, - }); - if ("error" in resolvedChannel) { - respond(false, undefined, resolvedChannel.error); - return; - } - const { cfg, channel } = resolvedChannel; - const plugin = resolveOutboundChannelPlugin({ channel, cfg }); - const outbound = plugin?.outbound; - if ( - typeof request.durationSeconds === "number" && - outbound?.supportsPollDurationSeconds !== true - ) { - // Duration support is channel-specific; reject before normalizing to avoid silent truncation. - respond( - false, - undefined, - errorShape( - ErrorCodes.INVALID_REQUEST, - `durationSeconds is not supported for ${channel} polls`, - ), - ); - return; - } - if (typeof request.isAnonymous === "boolean" && outbound?.supportsAnonymousPolls !== true) { - respond( - false, - undefined, - errorShape(ErrorCodes.INVALID_REQUEST, `isAnonymous is not supported for ${channel} polls`), - ); - return; - } - if (!plugin || !outbound?.sendPoll) { - respond( - false, - undefined, - errorShape(ErrorCodes.INVALID_REQUEST, `unsupported poll channel: ${channel}`), - ); - return; - } - const sendPoll = outbound.sendPoll; - let accountRoute: ReturnType; try { - accountRoute = resolveMessageOperationAccountRoute({ - cfg, - channel, - plugin, - accountIds: [request.accountId, routeBinding?.reservedRoute?.accountId], - conflictMessage: "poll account selections do not match", + routeBinding = resolveMessageOperationRouteBinding({ + context, + prefix: "poll", + idempotencyKey: request.idempotencyKey, + requestChannel: request.channel, + accountIds: [request.accountId], }); - } catch (error) { - respondGatewayInvalidRequest({ respond, channel, error }); - return; - } - if ( - !bindMessageOperationRoute({ + const reservedReplay = replayReservedMessageOperationRoute({ context, binding: routeBinding, - requestScope: accountRoute.requestScope, - }) - ) { - respondGatewayInvalidRequest({ + prefix: "poll", + idempotencyKey: request.idempotencyKey, respond, - channel, - error: "idempotency key is already bound to a different message route", }); - return; - } - const accountId = accountRoute.accountId; - const inflight = resolveGatewayInflightRequest({ - context, - prefix: "poll", - idempotencyKey: request.idempotencyKey, - respond, - requestScope: accountRoute.requestScope, - }); - if (inflight.kind === "handled") { - await inflight.done; - return; - } - const { idem, dedupeKey, inflightMap } = inflight; - retainMessageOperationRouteBinding({ - context, - binding: routeBinding, - requestScope: accountRoute.requestScope, - }); - const work = (async (): Promise => { - const poll = { - question: request.question, - options: request.options, - maxSelections: request.maxSelections, - durationSeconds: request.durationSeconds, - durationHours: request.durationHours, - }; - const threadId = normalizeOptionalString(request.threadId); - try { - const resolvedTarget = resolveGatewayOutboundTarget({ - channel, - to: request.to.trim(), - cfg, - accountId, - }); - if (!resolvedTarget.ok) { - return { ok: false, error: resolvedTarget.error }; - } - const normalized = outbound.pollMaxOptions - ? normalizePollInput(poll, { maxOptions: outbound.pollMaxOptions }) - : normalizePollInput(poll); - const result = await sendPoll({ - cfg, - to: resolvedTarget.to, - poll: normalized, - accountId, - threadId, - silent: request.silent, - isAnonymous: request.isAnonymous, - gatewayClientScopes: client?.connect?.scopes ?? [], - }); - const payload = buildGatewayDeliveryPayload({ runId: idem, channel, result }); - cacheGatewayDedupeSuccess({ - context, - dedupeKey, - payload, - }); - return { ok: true, payload, meta: { channel } }; - } catch (err) { - const error = errorShape(ErrorCodes.UNAVAILABLE, String(err)); - cacheGatewayDedupeFailure({ - context, - dedupeKey, - error, - }); - return { ok: false, error, meta: { channel, error: formatForLog(err) } }; + if (reservedReplay) { + releaseRouteBindingLock(); + await reservedReplay; + return; } - })().finally(() => { - refreshMessageOperationRouteBinding({ + const resolvedChannel = await resolveRequestedChannel({ + requestChannel: routeBinding?.reservedRoute?.channel ?? request.channel, + unsupportedMessage: (input) => `unsupported poll channel: ${input}`, + context, + }); + if ("error" in resolvedChannel) { + respond(false, undefined, resolvedChannel.error); + return; + } + const { cfg, channel } = resolvedChannel; + const plugin = resolveOutboundChannelPlugin({ channel, cfg }); + const outbound = plugin?.outbound; + if ( + typeof request.durationSeconds === "number" && + outbound?.supportsPollDurationSeconds !== true + ) { + // Duration support is channel-specific; reject before normalizing to avoid silent truncation. + respond( + false, + undefined, + errorShape( + ErrorCodes.INVALID_REQUEST, + `durationSeconds is not supported for ${channel} polls`, + ), + ); + return; + } + if (typeof request.isAnonymous === "boolean" && outbound?.supportsAnonymousPolls !== true) { + respond( + false, + undefined, + errorShape( + ErrorCodes.INVALID_REQUEST, + `isAnonymous is not supported for ${channel} polls`, + ), + ); + return; + } + if (!plugin || !outbound?.sendPoll) { + respond( + false, + undefined, + errorShape(ErrorCodes.INVALID_REQUEST, `unsupported poll channel: ${channel}`), + ); + return; + } + const sendPoll = outbound.sendPoll; + let accountRoute: ReturnType; + try { + accountRoute = resolveMessageOperationAccountRoute({ + cfg, + channel, + plugin, + accountIds: [request.accountId, routeBinding?.reservedRoute?.accountId], + conflictMessage: "poll account selections do not match", + }); + } catch (error) { + respondGatewayInvalidRequest({ respond, channel, error }); + return; + } + if ( + !bindMessageOperationRoute({ + context, + binding: routeBinding, + requestScope: accountRoute.requestScope, + }) + ) { + respondGatewayInvalidRequest({ + respond, + channel, + error: "idempotency key is already bound to a different message route", + }); + return; + } + const accountId = accountRoute.accountId; + const inflight = resolveGatewayInflightRequest({ + context, + prefix: "poll", + idempotencyKey: request.idempotencyKey, + respond, + requestScope: accountRoute.requestScope, + }); + if (inflight.kind === "handled") { + releaseRouteBindingLock(); + await inflight.done; + return; + } + const { idem, dedupeKey, inflightMap } = inflight; + retainMessageOperationRouteBinding({ context, binding: routeBinding, requestScope: accountRoute.requestScope, }); - }); + const work = (async (): Promise => { + const poll = { + question: request.question, + options: request.options, + maxSelections: request.maxSelections, + durationSeconds: request.durationSeconds, + durationHours: request.durationHours, + }; + const threadId = normalizeOptionalString(request.threadId); + try { + const resolvedTarget = resolveGatewayOutboundTarget({ + channel, + to: request.to.trim(), + cfg, + accountId, + }); + if (!resolvedTarget.ok) { + return { ok: false, error: resolvedTarget.error }; + } + const normalized = outbound.pollMaxOptions + ? normalizePollInput(poll, { maxOptions: outbound.pollMaxOptions }) + : normalizePollInput(poll); + const result = await sendPoll({ + cfg, + to: resolvedTarget.to, + poll: normalized, + accountId, + threadId, + silent: request.silent, + isAnonymous: request.isAnonymous, + gatewayClientScopes: client?.connect?.scopes ?? [], + }); + const payload = buildGatewayDeliveryPayload({ runId: idem, channel, result }); + cacheGatewayDedupeSuccess({ + context, + dedupeKey, + payload, + }); + return { ok: true, payload, meta: { channel } }; + } catch (err) { + const error = errorShape(ErrorCodes.UNAVAILABLE, String(err)); + cacheGatewayDedupeFailure({ + context, + dedupeKey, + error, + }); + return { ok: false, error, meta: { channel, error: formatForLog(err) } }; + } + })().finally(() => { + refreshMessageOperationRouteBinding({ + context, + binding: routeBinding, + requestScope: accountRoute.requestScope, + }); + }); - await runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + const inflightWork = runGatewayInflightWork({ inflightMap, dedupeKey, work, respond }); + releaseRouteBindingLock(); + await inflightWork; + } finally { + releaseRouteBindingLock(); + } }, }; /* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */ diff --git a/src/infra/outbound/message-account-selection.test.ts b/src/infra/outbound/message-account-selection.test.ts new file mode 100644 index 000000000000..ca6cf02b69eb --- /dev/null +++ b/src/infra/outbound/message-account-selection.test.ts @@ -0,0 +1,41 @@ +import { describe, expect, it } from "vitest"; +import type { ChannelPlugin } from "../../channels/plugins/types.plugin.js"; +import type { OpenClawConfig } from "../../config/types.openclaw.js"; +import { validateExplicitMessageAccountSelection } from "./message-account-selection.js"; + +describe("validateExplicitMessageAccountSelection", () => { + const cfg = {} as OpenClawConfig; + const plugin = { + id: "feishu", + config: { + listAccountIds: () => ["default"], + defaultAccountId: () => "ops", + resolveAccount: (_cfg: OpenClawConfig, accountId?: string | null) => ({ + accountId, + enabled: true, + }), + }, + } as unknown as ChannelPlugin; + + it("accepts the plugin-resolved default when it is intentionally unlisted", () => { + expect( + validateExplicitMessageAccountSelection({ + cfg, + channel: "feishu", + accountId: "OPS", + plugin, + }), + ).toBe("ops"); + }); + + it("still rejects a non-default unlisted account", () => { + expect(() => + validateExplicitMessageAccountSelection({ + cfg, + channel: "feishu", + accountId: "missing", + plugin, + }), + ).toThrow('Unknown account "missing"'); + }); +}); diff --git a/src/infra/outbound/message-account-selection.ts b/src/infra/outbound/message-account-selection.ts index 39d271bc93b9..d54639d94a54 100644 --- a/src/infra/outbound/message-account-selection.ts +++ b/src/infra/outbound/message-account-selection.ts @@ -1,5 +1,6 @@ import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { resolveChannelAccountEnabled } from "../../channels/account-summary.js"; +import { resolveChannelDefaultAccountId } from "../../channels/plugins/helpers.js"; import { getChannelPlugin, listChannelPlugins } from "../../channels/plugins/index.js"; import type { ChannelPlugin } from "../../channels/plugins/types.plugin.js"; import type { ChannelId } from "../../channels/plugins/types.public.js"; @@ -21,9 +22,19 @@ function resolveListedAccountId(params: { cfg: OpenClawConfig; accountId: string; }): string | undefined { - return params.plugin.config + const listedAccountId = params.plugin.config .listAccountIds(params.cfg) .find((candidate) => normalizeOptionalAccountId(candidate) === params.accountId); + if (listedAccountId) { + return listedAccountId; + } + const defaultAccountId = resolveChannelDefaultAccountId({ + plugin: params.plugin, + cfg: params.cfg, + }); + return normalizeOptionalAccountId(defaultAccountId) === params.accountId + ? defaultAccountId + : undefined; } function isExplicitAccountDisabled(params: {