From b6f518d69eb0c4b492e38438e51a749d2d49c0ed Mon Sep 17 00:00:00 2001 From: Alice Alexandra Moore <86723305+3mdistal@users.noreply.github.com> Date: Wed, 23 Sep 2026 13:16:41 -0400 Subject: [PATCH 1/4] fix(integrations): preserve A2A work during runtime gaps --- .changeset/pause-integration-continuations.md | 5 ++ docs/environment-variables.md | 1 + packages/core/src/app-config/integrations.ts | 7 ++ .../a2a-continuation-processor.spec.ts | 65 +++++++++++++++++++ .../a2a-continuation-processor.ts | 41 +++++++++++- .../integration-campaign-recovery.spec.ts | 20 ++++++ .../integration-campaign-recovery.ts | 13 ++++ .../integration-durable-dispatch-config.ts | 12 +--- .../integration-durable-dispatch.spec.ts | 26 ++++++++ .../integration-durable-dispatch.ts | 39 ++++++++--- packages/core/src/integrations/plugin.spec.ts | 30 ++++++++- packages/core/src/integrations/plugin.ts | 13 ++++ 12 files changed, 251 insertions(+), 21 deletions(-) create mode 100644 .changeset/pause-integration-continuations.md diff --git a/.changeset/pause-integration-continuations.md b/.changeset/pause-integration-continuations.md new file mode 100644 index 00000000000..32e96f6929f --- /dev/null +++ b/.changeset/pause-integration-continuations.md @@ -0,0 +1,5 @@ +--- +"@agent-native/core": patch +--- + +Preserve in-flight integration campaigns and A2A continuation identities when durable dispatch is temporarily unavailable, while retaining explicit rollout cancellation. diff --git a/docs/environment-variables.md b/docs/environment-variables.md index c45d83f4250..c38d87fda17 100644 --- a/docs/environment-variables.md +++ b/docs/environment-variables.md @@ -465,6 +465,7 @@ only in code. | `app.description` | — | string | — | One-line description of the app, from first-party template metadata. | | `auth.disableDesktopSsoFallbackInDevelopment` | `AGENT_NATIVE_DISABLE_DESKTOP_SSO_FALLBACK` | boolean | `false` | Disable the loopback Desktop SSO fallback in development so isolated acceptance runs can use their configured local identity. Ignored in production. | | `auth.requireEmailVerification` | `AUTH_REQUIRE_EMAIL_VERIFICATION` | boolean | — | Whether password signup must verify email before a session. Unset: hosted deployments verify with a provider and skip verification without one; local development skips it. Setting false accepts unverified email in production. Email remains optional. | +| `integrations.durableDispatch` | `AGENT_INTEGRATION_DURABLE_DISPATCH` | boolean | — | Use a durable background function for integration work. Explicit false stops an enabled rollout; unset leaves it unavailable. | | `integrations.allowUnverifiedWebhooks` | `AGENT_NATIVE_ALLOW_UNVERIFIED_WEBHOOKS` | boolean | `false` | Skip inbound webhook signature verification. Development only — every adapter that reads this treats it as a bypass of sender authentication. | | `integrations.webhookBaseUrl` | `WEBHOOK_BASE_URL` | string | — | Optional public base URL for self-callback and webhook targets. | | `integrations.platforms` | `AGENT_NATIVE_INTEGRATION_PLATFORMS` | array | — | Integration platforms to mount, comma-separated, each matched against an adapter's `platform` id (slack, telegram, whatsapp, microsoft-teams, discord, google-docs, email). Unset mounts every adapter; a name no adapter provides throws at plugin init. | diff --git a/packages/core/src/app-config/integrations.ts b/packages/core/src/app-config/integrations.ts index f683fb1e9d7..1dd79007184 100644 --- a/packages/core/src/app-config/integrations.ts +++ b/packages/core/src/app-config/integrations.ts @@ -2,6 +2,13 @@ import { z } from "zod"; /** Inbound integration webhook policy. */ export const integrationsConfig = z.object({ + durableDispatch: z + .boolean() + .optional() + .meta({ + env: ["AGENT_INTEGRATION_DURABLE_DISPATCH"], + doc: "Use a durable background function for integration work. Explicit false stops an enabled rollout; unset leaves it unavailable.", + }), allowUnverifiedWebhooks: z .boolean() .default(false) diff --git a/packages/core/src/integrations/a2a-continuation-processor.spec.ts b/packages/core/src/integrations/a2a-continuation-processor.spec.ts index 79ef18959d0..fba0f56ed13 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.spec.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.spec.ts @@ -13,6 +13,7 @@ const recoverDueA2AContinuationIdsMock = vi.hoisted(() => vi.fn()); const listRecoverableA2ATasksMock = vi.hoisted(() => vi.fn()); const getPendingTaskMock = vi.hoisted(() => vi.fn()); const durableDispatchEnabledMock = vi.hoisted(() => vi.fn()); +const durableDispatchExplicitlyDisabledMock = vi.hoisted(() => vi.fn()); const dispatchPendingIntegrationTaskMock = vi.hoisted(() => vi.fn()); const getNextPendingTaskForThreadMock = vi.hoisted(() => vi.fn()); const getIntegrationCampaignForTaskMock = vi.hoisted(() => vi.fn()); @@ -83,6 +84,11 @@ vi.mock("./pending-tasks-store.js", () => ({ vi.mock("./integration-durable-dispatch.js", () => ({ isIntegrationDurableDispatchEnabledForTask: durableDispatchEnabledMock, + isIntegrationDurableDispatchExplicitlyDisabledForTask: + durableDispatchExplicitlyDisabledMock, + integrationDurableDispatchRuntimeUnavailableReasons: () => [ + "background-route-unavailable", + ], dispatchPendingIntegrationTask: dispatchPendingIntegrationTaskMock, })); @@ -263,6 +269,7 @@ describe("A2A continuation processor", () => { status: "processing", }); durableDispatchEnabledMock.mockReturnValue(true); + durableDispatchExplicitlyDisabledMock.mockReturnValue(true); getIntegrationCampaignForTaskMock.mockResolvedValue(null); getNextPendingTaskForThreadMock.mockResolvedValue(null); dispatchPendingIntegrationTaskMock.mockResolvedValue( @@ -578,6 +585,22 @@ describe("A2A continuation processor", () => { expect(vi.mocked(fetch)).not.toHaveBeenCalled(); }); + it("keeps due continuations intact while the runtime prerequisite is unavailable", async () => { + durableDispatchEnabledMock.mockReturnValue(false); + durableDispatchExplicitlyDisabledMock.mockReturnValue(false); + recoverDueA2AContinuationIdsMock.mockResolvedValue([]); + const { recoverDueA2AContinuations } = + await import("./a2a-continuation-processor.js"); + + await expect(recoverDueA2AContinuations()).resolves.toEqual({ + dispatched: 0, + failed: 0, + }); + expect(failA2AContinuationsForIntegrationTaskMock).not.toHaveBeenCalled(); + expect(failDisabledIntegrationCampaignTaskMock).not.toHaveBeenCalled(); + expect(vi.mocked(fetch)).not.toHaveBeenCalled(); + }); + it("still wakes receipt-confirmed history when the rollout scope is disabled", async () => { durableDispatchEnabledMock.mockReturnValue(false); listRecoverableA2ATasksMock.mockResolvedValueOnce([ @@ -671,6 +694,48 @@ describe("A2A continuation processor", () => { }); }); + it("reschedules an already-started Content continuation without repeating its edit", async () => { + const claimed = continuation(); + const sendResponse = vi.fn(async () => ({ status: "delivered" as const })); + claimA2AContinuationMock.mockResolvedValue(claimed); + getIntegrationCampaignForTaskMock.mockResolvedValue({ + id: "campaign-1", + status: "waiting", + }); + getPendingTaskMock.mockResolvedValue({ + id: claimed.integrationTaskId, + platform: "slack", + externalThreadId: claimed.externalThreadId, + dispatchScope: "C123", + status: "processing", + }); + durableDispatchEnabledMock.mockReturnValueOnce(false); + durableDispatchExplicitlyDisabledMock.mockReturnValueOnce(false); + const { processA2AContinuationById } = + await import("./a2a-continuation-processor.js"); + const adapters = new Map([["slack", adapter(sendResponse)]]); + + await processA2AContinuationById(claimed.id, { + adapters, + }); + + expect(rescheduleA2AContinuationMock).toHaveBeenCalledWith( + claimed.id, + 20_000, + ); + expect(failA2AContinuationsForIntegrationTaskMock).not.toHaveBeenCalled(); + expect(failDisabledIntegrationCampaignTaskMock).not.toHaveBeenCalled(); + expect(getTaskMock).not.toHaveBeenCalled(); + expect(dispatchPendingIntegrationTaskMock).not.toHaveBeenCalled(); + + await processA2AContinuationById(claimed.id, { adapters }); + + expect(getTaskMock).toHaveBeenCalledWith(claimed.a2aTaskId); + expect(A2AClientMock).toHaveBeenCalledOnce(); + expect(sendResponse).toHaveBeenCalledOnce(); + expect(dispatchPendingIntegrationTaskMock).not.toHaveBeenCalled(); + }); + it("cancels only an unconfirmed sibling while confirmed history still owns custody", async () => { const claimed = continuation({ id: "cont-unconfirmed" }); claimA2AContinuationMock.mockResolvedValueOnce(claimed); diff --git a/packages/core/src/integrations/a2a-continuation-processor.ts b/packages/core/src/integrations/a2a-continuation-processor.ts index 7e57c683a14..ddb51e3ff5c 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.ts @@ -47,7 +47,9 @@ import { } from "./integration-campaigns-store.js"; import { dispatchPendingIntegrationTask, + integrationDurableDispatchRuntimeUnavailableReasons, isIntegrationDurableDispatchEnabledForTask, + isIntegrationDurableDispatchExplicitlyDisabledForTask, } from "./integration-durable-dispatch.js"; import { signInternalToken } from "./internal-token.js"; import { @@ -340,7 +342,15 @@ export async function recoverDueA2AContinuations(options?: { eligibleTaskIds.push(task.id); } else if (task.hasPendingConfirmedDelivery) { confirmedHistoryTaskIds.push(task.id); - } else { + } else if ( + isIntegrationDurableDispatchExplicitlyDisabledForTask({ + platform: task.platform, + externalThreadId: task.externalThreadId, + platformContext: task.dispatchScope + ? { channelId: task.dispatchScope } + : undefined, + }) + ) { await failDisabledDurableA2ATask(task); } if (eligibleTaskIds.length + confirmedHistoryTaskIds.length >= limit) break; @@ -576,6 +586,24 @@ async function durableContinuationScopeStillEnabled( }); if (enabled) return true; + if ( + task?.status === "processing" && + !isIntegrationDurableDispatchExplicitlyDisabledForTask({ + platform: task.platform, + externalThreadId: task.externalThreadId, + platformContext: task.dispatchScope + ? { channelId: task.dispatchScope } + : undefined, + }) + ) { + await rescheduleA2AContinuation(continuation.id, RESCHEDULE_DELAY_MS); + console.warn( + `[integrations] A2A continuation ${continuation.id} paused: durable dispatch runtime unavailable`, + integrationDurableDispatchRuntimeUnavailableReasons(), + ); + return false; + } + if ( await hasPendingConfirmedA2ADeliveryForIntegrationTask( continuation.integrationTaskId, @@ -659,6 +687,17 @@ export async function reconcileTerminalA2AParentIfDisabled( ) { return false; } + if ( + !isIntegrationDurableDispatchExplicitlyDisabledForTask({ + platform: task.platform, + externalThreadId: task.externalThreadId, + platformContext: task.dispatchScope + ? { channelId: task.dispatchScope } + : undefined, + }) + ) { + return false; + } await failDisabledDurableA2ATask({ id: task.id, platform: task.platform, diff --git a/packages/core/src/integrations/integration-campaign-recovery.spec.ts b/packages/core/src/integrations/integration-campaign-recovery.spec.ts index 382a37a3fa9..7b60aa0890f 100644 --- a/packages/core/src/integrations/integration-campaign-recovery.spec.ts +++ b/packages/core/src/integrations/integration-campaign-recovery.spec.ts @@ -5,6 +5,7 @@ const getCampaignMock = vi.hoisted(() => vi.fn()); const getTaskMock = vi.hoisted(() => vi.fn()); const dispatchMock = vi.hoisted(() => vi.fn()); const durableEnabledMock = vi.hoisted(() => vi.fn()); +const durableExplicitlyDisabledMock = vi.hoisted(() => vi.fn()); const failDisabledMock = vi.hoisted(() => vi.fn()); const getNextTaskMock = vi.hoisted(() => vi.fn()); const getA2AContinuationTaskOutcomeMock = vi.hoisted(() => vi.fn()); @@ -23,6 +24,8 @@ vi.mock("./pending-tasks-store.js", () => ({ vi.mock("./integration-durable-dispatch.js", () => ({ dispatchPendingIntegrationTask: dispatchMock, isIntegrationDurableDispatchEnabledForTask: durableEnabledMock, + isIntegrationDurableDispatchExplicitlyDisabledForTask: + durableExplicitlyDisabledMock, })); vi.mock("./a2a-continuations-store.js", () => ({ @@ -47,6 +50,7 @@ describe("integration campaign recovery", () => { }); dispatchMock.mockResolvedValue("background-acknowledged"); durableEnabledMock.mockReturnValue(true); + durableExplicitlyDisabledMock.mockReturnValue(true); getNextTaskMock.mockResolvedValue(null); getA2AContinuationTaskOutcomeMock.mockResolvedValue("missing"); }); @@ -131,6 +135,22 @@ describe("integration campaign recovery", () => { expect(failDisabledMock).toHaveBeenCalledWith("task-1"); }); + it("leaves a due campaign recoverable when runtime prerequisites are unavailable", async () => { + durableEnabledMock.mockReturnValueOnce(false); + durableExplicitlyDisabledMock.mockReturnValueOnce(false); + const { recoverDueIntegrationCampaigns } = + await import("./integration-campaign-recovery.js"); + + await expect(recoverDueIntegrationCampaigns({})).resolves.toEqual({ + selected: 1, + dispatched: 0, + skipped: 1, + failed: 0, + }); + expect(failDisabledMock).not.toHaveBeenCalled(); + expect(dispatchMock).not.toHaveBeenCalled(); + }); + it("still wakes confirmed receipt reconciliation after scope is disabled", async () => { durableEnabledMock.mockReturnValueOnce(false); getTaskMock.mockResolvedValueOnce({ diff --git a/packages/core/src/integrations/integration-campaign-recovery.ts b/packages/core/src/integrations/integration-campaign-recovery.ts index 65230691456..9a55c84c847 100644 --- a/packages/core/src/integrations/integration-campaign-recovery.ts +++ b/packages/core/src/integrations/integration-campaign-recovery.ts @@ -7,6 +7,7 @@ import { import { dispatchPendingIntegrationTask, isIntegrationDurableDispatchEnabledForTask, + isIntegrationDurableDispatchExplicitlyDisabledForTask, } from "./integration-durable-dispatch.js"; import { getNextPendingTaskForThread, @@ -77,6 +78,18 @@ export async function recoverDueIntegrationCampaigns(options: { hasConfirmedDeliveryReceipt(task.payload) || (await getA2AContinuationTaskOutcome(task.id)) === "terminal-delivered"; if (!durableDispatchEnabled && !confirmedReceipt) { + if ( + !isIntegrationDurableDispatchExplicitlyDisabledForTask({ + platform: task.platform, + externalThreadId: task.externalThreadId, + platformContext: task.dispatchScope + ? { channelId: task.dispatchScope } + : undefined, + }) + ) { + result.skipped += 1; + continue; + } await failDisabledIntegrationCampaignTask(task.id); const nextTask = await getNextPendingTaskForThread( task.platform, diff --git a/packages/core/src/integrations/integration-durable-dispatch-config.ts b/packages/core/src/integrations/integration-durable-dispatch-config.ts index 8562af0dd07..b7424f7f3a9 100644 --- a/packages/core/src/integrations/integration-durable-dispatch-config.ts +++ b/packages/core/src/integrations/integration-durable-dispatch-config.ts @@ -1,3 +1,5 @@ +import { getAppConfig } from "../app-config/index.js"; + export const INTEGRATION_DURABLE_DISPATCH_ENV = "AGENT_INTEGRATION_DURABLE_DISPATCH"; export const INTEGRATION_DURABLE_DISPATCH_SCOPES_ENV = @@ -20,13 +22,5 @@ export function isInIntegrationRecoveryRuntime(): boolean { } export function isIntegrationDurableDispatchConfigured(): boolean { - const value = process.env.AGENT_INTEGRATION_DURABLE_DISPATCH; - if (!value) return false; - const normalized = value.trim().toLowerCase(); - return ( - normalized === "1" || - normalized === "true" || - normalized === "yes" || - normalized === "on" - ); + return getAppConfig().integrations.durableDispatch === true; } diff --git a/packages/core/src/integrations/integration-durable-dispatch.spec.ts b/packages/core/src/integrations/integration-durable-dispatch.spec.ts index e9ff0424246..3a5e545fe4e 100644 --- a/packages/core/src/integrations/integration-durable-dispatch.spec.ts +++ b/packages/core/src/integrations/integration-durable-dispatch.spec.ts @@ -40,6 +40,32 @@ describe("durable integration dispatch", () => { ); }); + it("distinguishes a missing runtime prerequisite from an explicit rollout stop", async () => { + const { isIntegrationDurableDispatchExplicitlyDisabledForTask } = + await import("./integration-durable-dispatch.js"); + const task = { + platform: "slack", + externalThreadId: "slack:team:C123:1", + }; + + vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH", "true"); + vi.stubEnv("A2A_SECRET", ""); + expect(isIntegrationDurableDispatchExplicitlyDisabledForTask(task)).toBe( + false, + ); + + vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH", "false"); + expect(isIntegrationDurableDispatchExplicitlyDisabledForTask(task)).toBe( + true, + ); + + vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH", "true"); + vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH_SCOPES", "slack:C999"); + expect(isIntegrationDurableDispatchExplicitlyDisabledForTask(task)).toBe( + true, + ); + }); + it("uses an acknowledged Netlify handoff for an enabled Slack scope", async () => { vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH", "true"); vi.stubEnv("AGENT_INTEGRATION_DURABLE_DISPATCH_SCOPES", "slack:C123"); diff --git a/packages/core/src/integrations/integration-durable-dispatch.ts b/packages/core/src/integrations/integration-durable-dispatch.ts index 7d20d9f059f..3397036c344 100644 --- a/packages/core/src/integrations/integration-durable-dispatch.ts +++ b/packages/core/src/integrations/integration-durable-dispatch.ts @@ -5,6 +5,7 @@ import { dispatchPathTargetsNetlifyBackgroundFunction, resolveDurableBackgroundDispatchPath, } from "../agent/durable-background.js"; +import { getAppConfig } from "../app-config/index.js"; import { fireInternalDispatch } from "../server/self-dispatch.js"; import { INTEGRATION_DURABLE_DISPATCH_ENV, @@ -93,15 +94,7 @@ export function isIntegrationDurableDispatchEnabledForTask( task: IntegrationDispatchTaskScope, ): boolean { if (!isIntegrationDurableDispatchConfigured()) return false; - const scopes = configuredIntegrationDurableDispatchScopes(); - if ( - scopes && - !taskScopeCandidates(task).some((candidate) => - scopes.some((scope) => `${scope.platform}:${scope.value}` === candidate), - ) - ) { - return false; - } + if (isIntegrationDurableDispatchExplicitlyDisabledForTask(task)) return false; const path = resolveDurableBackgroundDispatchPath( INTEGRATION_PROCESS_TASK_PATH, ); @@ -111,6 +104,34 @@ export function isIntegrationDurableDispatchEnabledForTask( ); } +export function isIntegrationDurableDispatchExplicitlyDisabledForTask( + task: IntegrationDispatchTaskScope, +): boolean { + if (getAppConfig().integrations.durableDispatch === false) return true; + const scopes = configuredIntegrationDurableDispatchScopes(); + return Boolean( + scopes && + !taskScopeCandidates(task).some((candidate) => + scopes.some((scope) => `${scope.platform}:${scope.value}` === candidate), + ), + ); +} + +export function integrationDurableDispatchRuntimeUnavailableReasons(): string[] { + const reasons: string[] = []; + if (!isIntegrationDurableDispatchConfigured()) + reasons.push("flag-unavailable"); + if ( + !dispatchPathTargetsNetlifyBackgroundFunction( + resolveDurableBackgroundDispatchPath(INTEGRATION_PROCESS_TASK_PATH), + ) + ) { + reasons.push("background-route-unavailable"); + } + if (!hasConfiguredA2ASecret()) reasons.push("a2a-secret-unavailable"); + return reasons; +} + async function recordDispatch( taskId: string, outcome: IntegrationDispatchOutcome, diff --git a/packages/core/src/integrations/plugin.spec.ts b/packages/core/src/integrations/plugin.spec.ts index 03591a848b3..a59742f6b95 100644 --- a/packages/core/src/integrations/plugin.spec.ts +++ b/packages/core/src/integrations/plugin.spec.ts @@ -1232,7 +1232,7 @@ describe("integrations plugin routes", () => { it("fails a queued continuation closed after the durable scope is disabled", async () => { process.env.NODE_ENV = "development"; - delete process.env.AGENT_INTEGRATION_DURABLE_DISPATCH; + process.env.AGENT_INTEGRATION_DURABLE_DISPATCH = "false"; delete process.env.A2A_SECRET; const task = claimedTask(1); getPendingTaskMock.mockResolvedValueOnce(task); @@ -1256,6 +1256,32 @@ describe("integrations plugin routes", () => { expect(markTaskCompletedMock).not.toHaveBeenCalled(); }); + it("pauses a queued continuation when the runtime flag is unavailable", async () => { + process.env.NODE_ENV = "development"; + delete process.env.AGENT_INTEGRATION_DURABLE_DISPATCH; + delete process.env.A2A_SECRET; + const task = claimedTask(1); + getPendingTaskMock.mockResolvedValueOnce(task); + const nitroApp = createNitroApp(); + await createIntegrationsPlugin({ adapters: [adapter] })(nitroApp); + + const result = await dispatch( + nitroApp, + "/_agent-native/integrations/process-task", + "POST", + { taskId: task.id, __integrationCampaignContinuation: true }, + ); + + expect(result.status).toBe(202); + expect(result.body).toEqual({ + ok: true, + paused: "durable-runtime-unavailable", + }); + expect(failDisabledIntegrationCampaignTaskMock).not.toHaveBeenCalled(); + expect(processIntegrationTaskMock).not.toHaveBeenCalled(); + expect(markTaskCompletedMock).not.toHaveBeenCalled(); + }); + it("finishes an A2A parent after its partial receipt is reconciled last", async () => { process.env.NODE_ENV = "development"; process.env.NETLIFY = "true"; @@ -1337,7 +1363,7 @@ describe("integrations plugin routes", () => { it("does not send an unreceipted campaign delivery after scope is disabled", async () => { process.env.NODE_ENV = "development"; - delete process.env.AGENT_INTEGRATION_DURABLE_DISPATCH; + process.env.AGENT_INTEGRATION_DURABLE_DISPATCH = "false"; const baseTask = claimedTask(1); const task = { ...baseTask, diff --git a/packages/core/src/integrations/plugin.ts b/packages/core/src/integrations/plugin.ts index c18037dce2b..24eb92f4a67 100644 --- a/packages/core/src/integrations/plugin.ts +++ b/packages/core/src/integrations/plugin.ts @@ -108,6 +108,7 @@ import { integrationDispatchScopeValue, isInIntegrationRecoveryRuntime, isIntegrationDurableDispatchEnabledForTask, + isIntegrationDurableDispatchExplicitlyDisabledForTask, } from "./integration-durable-dispatch.js"; import { forgetIntegrationMemory, @@ -2060,6 +2061,18 @@ export function createIntegrationsPlugin( !durableCampaignEnabled && !confirmedDeliveryProof ) { + if ( + !isIntegrationDurableDispatchExplicitlyDisabledForTask({ + platform: task.platform, + externalThreadId: task.externalThreadId, + platformContext: task.dispatchScope + ? { channelId: task.dispatchScope } + : undefined, + }) + ) { + setResponseStatus(event, 202); + return { ok: true, paused: "durable-runtime-unavailable" }; + } await failDisabledIntegrationCampaignTask(task.id); const nextTask = await getNextPendingTaskForThread( task.platform, From a561ad7a6eddfef82675a265e93c5847126b2d43 Mon Sep 17 00:00:00 2001 From: Alice Alexandra Moore <86723305+3mdistal@users.noreply.github.com> Date: Wed, 23 Sep 2026 13:22:08 -0400 Subject: [PATCH 2/4] fix(integrations): keep paused recovery fair and retryable --- .../a2a-continuation-processor.spec.ts | 6 ++- .../a2a-continuation-processor.ts | 7 ++- .../a2a-continuations-store.spec.ts | 25 ++++++++++ .../integrations/a2a-continuations-store.ts | 20 +++++++- .../integration-campaign-recovery.spec.ts | 50 +++++++++++++++++++ .../integration-campaign-recovery.ts | 2 + .../integration-campaigns-store.spec.ts | 22 ++++++++ .../integration-campaigns-store.ts | 22 ++++++++ 8 files changed, 151 insertions(+), 3 deletions(-) diff --git a/packages/core/src/integrations/a2a-continuation-processor.spec.ts b/packages/core/src/integrations/a2a-continuation-processor.spec.ts index fba0f56ed13..568ee1cba53 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.spec.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.spec.ts @@ -41,6 +41,7 @@ const failA2AContinuationMock = vi.hoisted(() => vi.fn()); const failA2AContinuationsForIntegrationTaskMock = vi.hoisted(() => vi.fn()); const getA2AContinuationMock = vi.hoisted(() => vi.fn()); const rescheduleA2AContinuationMock = vi.hoisted(() => vi.fn()); +const pauseA2AContinuationForRuntimeMock = vi.hoisted(() => vi.fn()); const saveA2AVerifiedArtifactCheckpointMock = vi.hoisted(() => vi.fn()); const getTaskMock = vi.hoisted(() => vi.fn()); const signA2ATokenMock = vi.hoisted(() => @@ -74,6 +75,7 @@ vi.mock("./a2a-continuations-store.js", () => ({ recordA2ATerminalDeliveryReceipt: recordA2ATerminalDeliveryReceiptMock, retainA2AUnconfirmedDeliveryClaim: retainA2AUnconfirmedDeliveryClaimMock, rescheduleA2AContinuation: rescheduleA2AContinuationMock, + pauseA2AContinuationForRuntime: pauseA2AContinuationForRuntimeMock, saveA2AVerifiedArtifactCheckpoint: saveA2AVerifiedArtifactCheckpointMock, })); @@ -226,6 +228,7 @@ describe("A2A continuation processor", () => { ); retainA2AUnconfirmedDeliveryClaimMock.mockResolvedValue(undefined); rescheduleA2AContinuationMock.mockResolvedValue(undefined); + pauseA2AContinuationForRuntimeMock.mockResolvedValue(true); saveA2AVerifiedArtifactCheckpointMock.mockImplementation( async (_id: string, checkpoint: string) => checkpoint, ); @@ -719,8 +722,9 @@ describe("A2A continuation processor", () => { adapters, }); - expect(rescheduleA2AContinuationMock).toHaveBeenCalledWith( + expect(pauseA2AContinuationForRuntimeMock).toHaveBeenCalledWith( claimed.id, + claimed.attempts, 20_000, ); expect(failA2AContinuationsForIntegrationTaskMock).not.toHaveBeenCalled(); diff --git a/packages/core/src/integrations/a2a-continuation-processor.ts b/packages/core/src/integrations/a2a-continuation-processor.ts index ddb51e3ff5c..cd7a8b3cc7b 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.ts @@ -33,6 +33,7 @@ import { recordA2ATerminalDeliveryReceipt, retainA2AUnconfirmedDeliveryClaim, rescheduleA2AContinuation, + pauseA2AContinuationForRuntime, saveA2AVerifiedArtifactCheckpoint, type A2AContinuation, type A2ATerminalDeliveryKind, @@ -596,7 +597,11 @@ async function durableContinuationScopeStillEnabled( : undefined, }) ) { - await rescheduleA2AContinuation(continuation.id, RESCHEDULE_DELAY_MS); + await pauseA2AContinuationForRuntime( + continuation.id, + continuation.attempts, + RESCHEDULE_DELAY_MS, + ); console.warn( `[integrations] A2A continuation ${continuation.id} paused: durable dispatch runtime unavailable`, integrationDurableDispatchRuntimeUnavailableReasons(), diff --git a/packages/core/src/integrations/a2a-continuations-store.spec.ts b/packages/core/src/integrations/a2a-continuations-store.spec.ts index 2b0f96f4037..631f5bd134b 100644 --- a/packages/core/src/integrations/a2a-continuations-store.spec.ts +++ b/packages/core/src/integrations/a2a-continuations-store.spec.ts @@ -606,9 +606,34 @@ describe("A2A continuations store", () => { querySql(query).includes("INNER JOIN integration_pending_tasks"), ); expect(joinedReads).toHaveLength(1); + expect(querySql(joinedReads[0]![0])).toContain( + "ORDER BY has_pending_confirmed_delivery DESC", + ); expect(queryArgs(joinedReads[0]![0]).at(-1)).toBe(10); }); + it("returns the claim attempt when runtime pauses before remote polling", async () => { + const { pauseA2AContinuationForRuntime } = await loadStore(); + executeMock.mockResolvedValue({ rows: [{ id: "cont-1" }] }); + + await expect( + pauseA2AContinuationForRuntime("cont-1", 31, 20_000), + ).resolves.toBe(true); + + const update = executeMock.mock.calls.find(([query]) => + querySql(query).includes("attempts = attempts - 1"), + )?.[0]; + expect(querySql(update!)).toContain( + "status = 'processing' AND attempts = ?", + ); + expect(queryArgs(update!)).toEqual([ + expect.any(Number), + expect.any(Number), + "cont-1", + 31, + ]); + }); + it("terminalizes all active A2A rows for a disabled durable task", async () => { const { failA2AContinuationsForIntegrationTask } = await loadStore(); executeMock.mockResolvedValue({ rows: [], rowsAffected: 2 }); diff --git a/packages/core/src/integrations/a2a-continuations-store.ts b/packages/core/src/integrations/a2a-continuations-store.ts index 28605bc974f..2ebfef803d7 100644 --- a/packages/core/src/integrations/a2a-continuations-store.ts +++ b/packages/core/src/integrations/a2a-continuations-store.ts @@ -782,7 +782,7 @@ export async function listRecoverableA2AIntegrationTasks( OR (c.status = 'delivering' AND ((c.terminal_delivery_confirmed_at IS NOT NULL AND c.next_check_at <= ?) OR c.updated_at <= ?))) - ORDER BY c.integration_task_id ASC + ORDER BY has_pending_confirmed_delivery DESC, c.integration_task_id ASC LIMIT ?`, args: [ now, @@ -837,6 +837,24 @@ export async function rescheduleA2AContinuation( }); } +export async function pauseA2AContinuationForRuntime( + id: string, + claimedAttempts: number, + delayMs: number, +): Promise { + await ensureTable(); + const now = Date.now(); + const result = await getDbExec().execute({ + sql: `UPDATE integration_a2a_continuations + SET status = 'pending', attempts = attempts - 1, + next_check_at = ?, updated_at = ? + WHERE id = ? AND status = 'processing' AND attempts = ? + RETURNING id`, + args: [now + delayMs, now, id, claimedAttempts], + }); + return (result.rows?.length ?? 0) > 0; +} + export async function retainA2AUnconfirmedDeliveryClaim( id: string, ): Promise { diff --git a/packages/core/src/integrations/integration-campaign-recovery.spec.ts b/packages/core/src/integrations/integration-campaign-recovery.spec.ts index 7b60aa0890f..846423a855c 100644 --- a/packages/core/src/integrations/integration-campaign-recovery.spec.ts +++ b/packages/core/src/integrations/integration-campaign-recovery.spec.ts @@ -1,6 +1,7 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; const listDueMock = vi.hoisted(() => vi.fn()); +const deferRuntimeMock = vi.hoisted(() => vi.fn()); const getCampaignMock = vi.hoisted(() => vi.fn()); const getTaskMock = vi.hoisted(() => vi.fn()); const dispatchMock = vi.hoisted(() => vi.fn()); @@ -12,6 +13,7 @@ const getA2AContinuationTaskOutcomeMock = vi.hoisted(() => vi.fn()); vi.mock("./integration-campaigns-store.js", () => ({ listDueIntegrationCampaignIds: listDueMock, + deferIntegrationCampaignForRuntime: deferRuntimeMock, getIntegrationCampaign: getCampaignMock, failDisabledIntegrationCampaignTask: failDisabledMock, })); @@ -36,6 +38,7 @@ describe("integration campaign recovery", () => { beforeEach(() => { vi.clearAllMocks(); listDueMock.mockResolvedValue(["campaign-1"]); + deferRuntimeMock.mockResolvedValue(true); getCampaignMock.mockResolvedValue({ id: "campaign-1", integrationTaskId: "task-1", @@ -149,6 +152,53 @@ describe("integration campaign recovery", () => { }); expect(failDisabledMock).not.toHaveBeenCalled(); expect(dispatchMock).not.toHaveBeenCalled(); + expect(deferRuntimeMock).toHaveBeenCalledWith("campaign-1", 60_000); + }); + + it("clears a full paused scan window before a confirmed delivery", async () => { + const pausedIds = Array.from( + { length: 20 }, + (_, index) => `paused-${index}`, + ); + listDueMock + .mockResolvedValueOnce(pausedIds) + .mockResolvedValueOnce(["receipt"]); + getCampaignMock.mockImplementation(async (id: string) => ({ + id, + integrationTaskId: id, + })); + getTaskMock.mockImplementation(async (id: string) => ({ + id, + platform: "slack", + externalThreadId: "slack:team:C123:1", + dispatchScope: "C123", + status: "processing", + payload: JSON.stringify( + id === "receipt" + ? { + kind: "response-delivery", + deliveryReceipt: { status: "delivered" }, + } + : { incoming: { text: "continue" } }, + ), + })); + durableEnabledMock.mockReturnValue(false); + durableExplicitlyDisabledMock.mockReturnValue(false); + const { recoverDueIntegrationCampaigns } = + await import("./integration-campaign-recovery.js"); + + await expect(recoverDueIntegrationCampaigns({})).resolves.toMatchObject({ + selected: 20, + skipped: 20, + }); + expect(deferRuntimeMock).toHaveBeenCalledTimes(20); + await expect(recoverDueIntegrationCampaigns({})).resolves.toMatchObject({ + selected: 1, + dispatched: 1, + }); + expect(dispatchMock).toHaveBeenCalledWith( + expect.objectContaining({ taskId: "receipt" }), + ); }); it("still wakes confirmed receipt reconciliation after scope is disabled", async () => { diff --git a/packages/core/src/integrations/integration-campaign-recovery.ts b/packages/core/src/integrations/integration-campaign-recovery.ts index 9a55c84c847..3b1c0d0bfe8 100644 --- a/packages/core/src/integrations/integration-campaign-recovery.ts +++ b/packages/core/src/integrations/integration-campaign-recovery.ts @@ -3,6 +3,7 @@ import { getIntegrationCampaign, failDisabledIntegrationCampaignTask, listDueIntegrationCampaignIds, + deferIntegrationCampaignForRuntime, } from "./integration-campaigns-store.js"; import { dispatchPendingIntegrationTask, @@ -87,6 +88,7 @@ export async function recoverDueIntegrationCampaigns(options: { : undefined, }) ) { + await deferIntegrationCampaignForRuntime(campaign.id, 60_000); result.skipped += 1; continue; } diff --git a/packages/core/src/integrations/integration-campaigns-store.spec.ts b/packages/core/src/integrations/integration-campaigns-store.spec.ts index 52d0726b4e4..62540f44c20 100644 --- a/packages/core/src/integrations/integration-campaigns-store.spec.ts +++ b/packages/core/src/integrations/integration-campaigns-store.spec.ts @@ -276,6 +276,28 @@ describe("integration campaigns store", () => { ]); }); + it("defers only due runtime-paused campaigns so later receipts can be swept", async () => { + const { deferIntegrationCampaignForRuntime } = await loadStore(); + executeMock.mockResolvedValue({ rows: [{ id: "campaign-1" }] }); + + await expect( + deferIntegrationCampaignForRuntime("campaign-1", 60_000), + ).resolves.toBe(true); + + const update = executeMock.mock.calls.find(([query]) => + sqlOf(query).includes("SET next_run_at = ?"), + )?.[0]; + expect(sqlOf(update!)).toContain("lease_expires_at <= ?"); + expect(argsOf(update!)).toEqual([ + expect.any(Number), + expect.any(Number), + expect.any(Number), + "campaign-1", + expect.any(Number), + expect.any(Number), + ]); + }); + it("reports the chunk ceiling instead of looping a campaign forever", async () => { const { claimIntegrationCampaign } = await loadStore(); executeMock.mockImplementation(async (query: string | { sql: string }) => { diff --git a/packages/core/src/integrations/integration-campaigns-store.ts b/packages/core/src/integrations/integration-campaigns-store.ts index 675a777b39c..afa197d5425 100644 --- a/packages/core/src/integrations/integration-campaigns-store.ts +++ b/packages/core/src/integrations/integration-campaigns-store.ts @@ -967,6 +967,28 @@ export async function completeIntegrationCampaignTaskAfterA2A( throw new Error("Database does not support atomic A2A parent completion"); } +export async function deferIntegrationCampaignForRuntime( + id: string, + delayMs: number, +): Promise { + await ensureTable(); + const now = Date.now(); + const nextRunAt = now + delayMs; + const result = await getDbExec().execute({ + sql: `UPDATE integration_campaigns + SET next_run_at = ?, + lease_expires_at = CASE WHEN status = 'processing' THEN ? ELSE lease_expires_at END, + updated_at = ? + WHERE id = ? AND ( + (status IN ('pending', 'waiting') AND next_run_at <= ?) + OR (status = 'processing' AND lease_expires_at IS NOT NULL AND lease_expires_at <= ?) + ) + RETURNING id`, + args: [nextRunAt, nextRunAt, now, id, now, now], + }); + return (result.rows?.length ?? 0) > 0; +} + export async function listDueIntegrationCampaignIds( limit = 25, ): Promise { From 021349b5bc6e300e5d84287b66328a84749679e3 Mon Sep 17 00:00:00 2001 From: Alice Alexandra Moore <86723305+3mdistal@users.noreply.github.com> Date: Wed, 23 Sep 2026 13:25:22 -0400 Subject: [PATCH 3/4] test(integrations): prove pauses do not burn polling attempts --- .../a2a-continuations-store.spec.ts | 41 +++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/packages/core/src/integrations/a2a-continuations-store.spec.ts b/packages/core/src/integrations/a2a-continuations-store.spec.ts index 631f5bd134b..9e6dde096e0 100644 --- a/packages/core/src/integrations/a2a-continuations-store.spec.ts +++ b/packages/core/src/integrations/a2a-continuations-store.spec.ts @@ -634,6 +634,47 @@ describe("A2A continuations store", () => { ]); }); + it("does not exhaust the remote polling budget through repeated runtime pauses", async () => { + const state = continuationRow({ status: "pending", attempts: 0 }); + executeMock.mockImplementation( + async (query: string | { sql: string; args?: unknown[] }) => { + const sql = querySql(query); + if (sql.includes("attempts = attempts + 1")) { + if (state.status !== "pending") return { rows: [] }; + state.status = "processing"; + state.attempts = Number(state.attempts) + 1; + return { rows: [{ ...state }] }; + } + if (sql.includes("attempts = attempts - 1")) { + const expected = queryArgs(query).at(-1); + if (state.status !== "processing" || state.attempts !== expected) { + return { rows: [] }; + } + state.status = "pending"; + state.attempts = Number(state.attempts) - 1; + return { rows: [{ id: state.id }] }; + } + return { rows: [] }; + }, + ); + const { claimA2AContinuation, pauseA2AContinuationForRuntime } = + await loadStore(); + + for (let pause = 0; pause < 35; pause += 1) { + const claimed = await claimA2AContinuation("cont-1"); + expect(claimed?.attempts).toBe(1); + await expect( + pauseA2AContinuationForRuntime("cont-1", claimed!.attempts, 20_000), + ).resolves.toBe(true); + } + + expect(state.attempts).toBe(0); + await expect(claimA2AContinuation("cont-1")).resolves.toMatchObject({ + attempts: 1, + a2aTaskId: "a2a-task-1", + }); + }); + it("terminalizes all active A2A rows for a disabled durable task", async () => { const { failA2AContinuationsForIntegrationTask } = await loadStore(); executeMock.mockResolvedValue({ rows: [], rowsAffected: 2 }); From d334f298c0bc303b15cf66e252b615fdf8476082 Mon Sep 17 00:00:00 2001 From: Alice Alexandra Moore <86723305+3mdistal@users.noreply.github.com> Date: Wed, 23 Sep 2026 15:29:09 -0400 Subject: [PATCH 4/4] fix(integrations): keep recovery scans fair during runtime gaps --- .../a2a-continuation-processor.spec.ts | 53 +++++++++++++++++++ .../a2a-continuation-processor.ts | 5 ++ .../a2a-continuations-store.spec.ts | 28 ++++++++++ .../integrations/a2a-continuations-store.ts | 36 ++++++++++++- 4 files changed, 120 insertions(+), 2 deletions(-) diff --git a/packages/core/src/integrations/a2a-continuation-processor.spec.ts b/packages/core/src/integrations/a2a-continuation-processor.spec.ts index 568ee1cba53..825969d0c81 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.spec.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.spec.ts @@ -11,6 +11,7 @@ const claimA2AContinuationMock = vi.hoisted(() => vi.fn()); const claimDueA2AContinuationsMock = vi.hoisted(() => vi.fn(async () => [])); const recoverDueA2AContinuationIdsMock = vi.hoisted(() => vi.fn()); const listRecoverableA2ATasksMock = vi.hoisted(() => vi.fn()); +const deferA2AContinuationsForRuntimeMock = vi.hoisted(() => vi.fn()); const getPendingTaskMock = vi.hoisted(() => vi.fn()); const durableDispatchEnabledMock = vi.hoisted(() => vi.fn()); const durableDispatchExplicitlyDisabledMock = vi.hoisted(() => vi.fn()); @@ -71,6 +72,7 @@ vi.mock("./a2a-continuations-store.js", () => ({ hasOnlyLegacyFailedA2AContinuationsForIntegrationTask: hasOnlyLegacyFailedA2AContinuationsForIntegrationTaskMock, listRecoverableA2AIntegrationTasks: listRecoverableA2ATasksMock, + deferA2AContinuationsForRuntime: deferA2AContinuationsForRuntimeMock, recoverDueA2AContinuationIds: recoverDueA2AContinuationIdsMock, recordA2ATerminalDeliveryReceipt: recordA2ATerminalDeliveryReceiptMock, retainA2AUnconfirmedDeliveryClaim: retainA2AUnconfirmedDeliveryClaimMock, @@ -601,9 +603,60 @@ describe("A2A continuation processor", () => { }); expect(failA2AContinuationsForIntegrationTaskMock).not.toHaveBeenCalled(); expect(failDisabledIntegrationCampaignTaskMock).not.toHaveBeenCalled(); + expect(deferA2AContinuationsForRuntimeMock).toHaveBeenCalledWith( + ["task-1", "task-2"], + 120_000, + ); expect(vi.mocked(fetch)).not.toHaveBeenCalled(); }); + it("moves a full unavailable scan window aside so a later task can recover", async () => { + const unavailable = Array.from({ length: 200 }, (_, index) => ({ + id: `paused-${index}`, + platform: "slack", + externalThreadId: `slack:team:C123:${index}`, + dispatchScope: "C123", + status: "processing", + hasPendingConfirmedDelivery: false, + })); + listRecoverableA2ATasksMock + .mockResolvedValueOnce(unavailable) + .mockResolvedValueOnce([ + { + id: "eligible", + platform: "slack", + externalThreadId: "slack:team:C123:eligible", + dispatchScope: "C123", + status: "processing", + hasPendingConfirmedDelivery: false, + }, + ]); + durableDispatchEnabledMock.mockReturnValueOnce(false); + durableDispatchEnabledMock.mockImplementation((task) => + task.externalThreadId.endsWith(":eligible"), + ); + durableDispatchExplicitlyDisabledMock.mockReturnValue(false); + recoverDueA2AContinuationIdsMock + .mockResolvedValueOnce([]) + .mockResolvedValueOnce(["cont-eligible"]); + const { recoverDueA2AContinuations } = + await import("./a2a-continuation-processor.js"); + + await expect(recoverDueA2AContinuations()).resolves.toEqual({ + dispatched: 0, + failed: 0, + }); + expect(deferA2AContinuationsForRuntimeMock).toHaveBeenCalledWith( + unavailable.map((task) => task.id), + 120_000, + ); + await expect(recoverDueA2AContinuations()).resolves.toEqual({ + dispatched: 1, + failed: 0, + }); + expect(fetch).toHaveBeenCalledOnce(); + }); + it("still wakes receipt-confirmed history when the rollout scope is disabled", async () => { durableDispatchEnabledMock.mockReturnValue(false); listRecoverableA2ATasksMock.mockResolvedValueOnce([ diff --git a/packages/core/src/integrations/a2a-continuation-processor.ts b/packages/core/src/integrations/a2a-continuation-processor.ts index cd7a8b3cc7b..f8afcb150d4 100644 --- a/packages/core/src/integrations/a2a-continuation-processor.ts +++ b/packages/core/src/integrations/a2a-continuation-processor.ts @@ -29,6 +29,7 @@ import { hasOnlyLegacyFailedA2AContinuationsForIntegrationTask, hasPendingConfirmedA2ADeliveryForIntegrationTask, listRecoverableA2AIntegrationTasks, + deferA2AContinuationsForRuntime, recoverDueA2AContinuationIds, recordA2ATerminalDeliveryReceipt, retainA2AUnconfirmedDeliveryClaim, @@ -331,6 +332,7 @@ export async function recoverDueA2AContinuations(options?: { const candidateTasks = await listRecoverableA2AIntegrationTasks(200); const eligibleTaskIds: string[] = []; const confirmedHistoryTaskIds: string[] = []; + const unavailableTaskIds: string[] = []; for (const task of candidateTasks) { const enabled = isIntegrationDurableDispatchEnabledForTask({ platform: task.platform, @@ -353,9 +355,12 @@ export async function recoverDueA2AContinuations(options?: { }) ) { await failDisabledDurableA2ATask(task); + } else { + unavailableTaskIds.push(task.id); } if (eligibleTaskIds.length + confirmedHistoryTaskIds.length >= limit) break; } + await deferA2AContinuationsForRuntime(unavailableTaskIds, 2 * 60_000); const ids = await recoverDueA2AContinuationIds(limit, eligibleTaskIds); const remaining = Math.max(0, limit - ids.length); const confirmedHistoryIds = diff --git a/packages/core/src/integrations/a2a-continuations-store.spec.ts b/packages/core/src/integrations/a2a-continuations-store.spec.ts index 9e6dde096e0..aa2179df31c 100644 --- a/packages/core/src/integrations/a2a-continuations-store.spec.ts +++ b/packages/core/src/integrations/a2a-continuations-store.spec.ts @@ -609,9 +609,37 @@ describe("A2A continuations store", () => { expect(querySql(joinedReads[0]![0])).toContain( "ORDER BY has_pending_confirmed_delivery DESC", ); + expect(querySql(joinedReads[0]![0])).toContain("MIN(c.next_check_at) ASC"); expect(queryArgs(joinedReads[0]![0]).at(-1)).toBe(10); }); + it("defers only due unconfirmed continuations without spending a claim", async () => { + const { deferA2AContinuationsForRuntime } = await loadStore(); + executeMock.mockResolvedValue({ rows: [], rowsAffected: 2 }); + + await deferA2AContinuationsForRuntime(["task-1", "task-2"], 120_000); + + const update = executeMock.mock.calls.find(([query]) => + querySql(query).includes("SET next_check_at = ?"), + )?.[0]; + expect(querySql(update!)).toContain( + "terminal_delivery_confirmed_at IS NULL", + ); + expect(querySql(update!)).toContain("status = 'processing'"); + expect(querySql(update!)).toContain("status = 'delivering'"); + expect(querySql(update!)).not.toContain("attempts = attempts + 1"); + expect(queryArgs(update!)).toEqual([ + expect.any(Number), + expect.any(Number), + "task-1", + "task-2", + expect.any(Number), + expect.any(Number), + expect.any(Number), + expect.any(Number), + ]); + }); + it("returns the claim attempt when runtime pauses before remote polling", async () => { const { pauseA2AContinuationForRuntime } = await loadStore(); executeMock.mockResolvedValue({ rows: [{ id: "cont-1" }] }); diff --git a/packages/core/src/integrations/a2a-continuations-store.ts b/packages/core/src/integrations/a2a-continuations-store.ts index 2ebfef803d7..d18cae2402a 100644 --- a/packages/core/src/integrations/a2a-continuations-store.ts +++ b/packages/core/src/integrations/a2a-continuations-store.ts @@ -764,7 +764,7 @@ export async function listRecoverableA2AIntegrationTasks( await ensureTable(); const now = Date.now(); const { rows } = await getDbExec().execute({ - sql: `SELECT DISTINCT c.integration_task_id, t.platform, + sql: `SELECT c.integration_task_id, t.platform, t.external_thread_id, t.dispatch_scope, t.status, EXISTS ( SELECT 1 FROM integration_a2a_continuations receipt @@ -782,7 +782,10 @@ export async function listRecoverableA2AIntegrationTasks( OR (c.status = 'delivering' AND ((c.terminal_delivery_confirmed_at IS NOT NULL AND c.next_check_at <= ?) OR c.updated_at <= ?))) - ORDER BY has_pending_confirmed_delivery DESC, c.integration_task_id ASC + GROUP BY c.integration_task_id, t.platform, t.external_thread_id, + t.dispatch_scope, t.status + ORDER BY has_pending_confirmed_delivery DESC, + MIN(c.next_check_at) ASC, c.integration_task_id ASC LIMIT ?`, args: [ now, @@ -805,6 +808,35 @@ export async function listRecoverableA2AIntegrationTasks( })); } +export async function deferA2AContinuationsForRuntime( + integrationTaskIds: string[], + delayMs: number, +): Promise { + if (integrationTaskIds.length === 0) return; + await ensureTable(); + const now = Date.now(); + const taskFilter = integrationTaskIds.map(() => "?").join(", "); + await getDbExec().execute({ + sql: `UPDATE integration_a2a_continuations + SET next_check_at = ?, updated_at = ? + WHERE integration_task_id IN (${taskFilter}) + AND terminal_delivery_confirmed_at IS NULL + AND ((status = 'pending' AND next_check_at <= ?) + OR (status = 'processing' AND + (updated_at <= ? OR next_check_at <= ?)) + OR (status = 'delivering' AND updated_at <= ?))`, + args: [ + now + delayMs, + now, + ...integrationTaskIds, + now, + now - PROCESSING_STUCK_AFTER_MS, + now - PROCESSING_NEXT_CHECK_STALE_AFTER_MS, + now - PROCESSING_STUCK_AFTER_MS, + ], + }); +} + export async function claimA2AContinuationDelivery( id: string, ): Promise {