Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/pause-integration-continuations.md
Original file line number Diff line number Diff line change
@@ -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.
1 change: 1 addition & 0 deletions docs/environment-variables.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Expand Down
7 changes: 7 additions & 0 deletions packages/core/src/app-config/integrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
122 changes: 122 additions & 0 deletions packages/core/src/integrations/a2a-continuation-processor.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,10 @@ 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());
const dispatchPendingIntegrationTaskMock = vi.hoisted(() => vi.fn());
const getNextPendingTaskForThreadMock = vi.hoisted(() => vi.fn());
const getIntegrationCampaignForTaskMock = vi.hoisted(() => vi.fn());
Expand Down Expand Up @@ -40,6 +42,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(() =>
Expand Down Expand Up @@ -69,10 +72,12 @@ vi.mock("./a2a-continuations-store.js", () => ({
hasOnlyLegacyFailedA2AContinuationsForIntegrationTask:
hasOnlyLegacyFailedA2AContinuationsForIntegrationTaskMock,
listRecoverableA2AIntegrationTasks: listRecoverableA2ATasksMock,
deferA2AContinuationsForRuntime: deferA2AContinuationsForRuntimeMock,
recoverDueA2AContinuationIds: recoverDueA2AContinuationIdsMock,
recordA2ATerminalDeliveryReceipt: recordA2ATerminalDeliveryReceiptMock,
retainA2AUnconfirmedDeliveryClaim: retainA2AUnconfirmedDeliveryClaimMock,
rescheduleA2AContinuation: rescheduleA2AContinuationMock,
pauseA2AContinuationForRuntime: pauseA2AContinuationForRuntimeMock,
saveA2AVerifiedArtifactCheckpoint: saveA2AVerifiedArtifactCheckpointMock,
}));

Expand All @@ -83,6 +88,11 @@ vi.mock("./pending-tasks-store.js", () => ({

vi.mock("./integration-durable-dispatch.js", () => ({
isIntegrationDurableDispatchEnabledForTask: durableDispatchEnabledMock,
isIntegrationDurableDispatchExplicitlyDisabledForTask:
durableDispatchExplicitlyDisabledMock,
integrationDurableDispatchRuntimeUnavailableReasons: () => [
"background-route-unavailable",
],
dispatchPendingIntegrationTask: dispatchPendingIntegrationTaskMock,
}));

Expand Down Expand Up @@ -220,6 +230,7 @@ describe("A2A continuation processor", () => {
);
retainA2AUnconfirmedDeliveryClaimMock.mockResolvedValue(undefined);
rescheduleA2AContinuationMock.mockResolvedValue(undefined);
pauseA2AContinuationForRuntimeMock.mockResolvedValue(true);
saveA2AVerifiedArtifactCheckpointMock.mockImplementation(
async (_id: string, checkpoint: string) => checkpoint,
);
Expand Down Expand Up @@ -263,6 +274,7 @@ describe("A2A continuation processor", () => {
status: "processing",
});
durableDispatchEnabledMock.mockReturnValue(true);
durableDispatchExplicitlyDisabledMock.mockReturnValue(true);
getIntegrationCampaignForTaskMock.mockResolvedValue(null);
getNextPendingTaskForThreadMock.mockResolvedValue(null);
dispatchPendingIntegrationTaskMock.mockResolvedValue(
Expand Down Expand Up @@ -578,6 +590,73 @@ 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(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([
Expand Down Expand Up @@ -671,6 +750,49 @@ 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(pauseA2AContinuationForRuntimeMock).toHaveBeenCalledWith(
claimed.id,
claimed.attempts,
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);
Expand Down
51 changes: 50 additions & 1 deletion packages/core/src/integrations/a2a-continuation-processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,12 @@ import {
hasOnlyLegacyFailedA2AContinuationsForIntegrationTask,
hasPendingConfirmedA2ADeliveryForIntegrationTask,
listRecoverableA2AIntegrationTasks,
deferA2AContinuationsForRuntime,
recoverDueA2AContinuationIds,
recordA2ATerminalDeliveryReceipt,
retainA2AUnconfirmedDeliveryClaim,
rescheduleA2AContinuation,
pauseA2AContinuationForRuntime,
saveA2AVerifiedArtifactCheckpoint,
type A2AContinuation,
type A2ATerminalDeliveryKind,
Expand All @@ -47,7 +49,9 @@ import {
} from "./integration-campaigns-store.js";
import {
dispatchPendingIntegrationTask,
integrationDurableDispatchRuntimeUnavailableReasons,
isIntegrationDurableDispatchEnabledForTask,
isIntegrationDurableDispatchExplicitlyDisabledForTask,
} from "./integration-durable-dispatch.js";
import { signInternalToken } from "./internal-token.js";
import {
Expand Down Expand Up @@ -328,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,
Expand All @@ -340,11 +345,22 @@ 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);
} 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 =
Expand Down Expand Up @@ -576,6 +592,28 @@ 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 pauseA2AContinuationForRuntime(
Comment thread
builder-io-integration[bot] marked this conversation as resolved.
continuation.id,
continuation.attempts,
RESCHEDULE_DELAY_MS,
);
console.warn(
`[integrations] A2A continuation ${continuation.id} paused: durable dispatch runtime unavailable`,
integrationDurableDispatchRuntimeUnavailableReasons(),
);
return false;
}

if (
await hasPendingConfirmedA2ADeliveryForIntegrationTask(
continuation.integrationTaskId,
Expand Down Expand Up @@ -659,6 +697,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,
Expand Down
Loading
Loading