diff --git a/extensions/slack/src/monitor/ingress.test.ts b/extensions/slack/src/monitor/ingress.test.ts index 9c6a57b434a1..65c83b0afa0c 100644 --- a/extensions/slack/src/monitor/ingress.test.ts +++ b/extensions/slack/src/monitor/ingress.test.ts @@ -5,7 +5,10 @@ import os from "node:os"; import path from "node:path"; import { App, type Receiver, type ReceiverEvent } from "@slack/bolt"; import type { WebClientOptions } from "@slack/web-api"; -import type { ChannelIngressQueue } from "openclaw/plugin-sdk/channel-outbound"; +import type { + ChannelIngressMonitorLifecycle, + ChannelIngressQueue, +} from "openclaw/plugin-sdk/channel-outbound"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import type { PluginJsonValue } from "openclaw/plugin-sdk/plugin-entry"; import { @@ -211,12 +214,13 @@ function createReceiverEventWithBody(body: Record): ReceiverEve function attachIngress( queue: ChannelIngressQueue, processEvent: (event: ReceiverEvent) => Promise, + options: { adoptionStallTimeoutMs?: number } = {}, ) { const ingress = createSlackDurableIngress({ accountId: "default", queue, pollIntervalMs: 60_000, - adoptionStallTimeoutMs: 5_000, + adoptionStallTimeoutMs: options.adoptionStallTimeoutMs ?? 5_000, }); const harness = createReceiverHarness(); ingress.wrapReceiver(harness.receiver).init({ processEvent } as App); @@ -302,6 +306,181 @@ describe("Slack durable ingress", () => { }); }); + it("dispatches independently routed threads concurrently after session ownership is established", async () => { + await withQueue(async (queue) => { + let releaseFirstDispatch: () => void = () => {}; + const firstDispatchGate = new Promise((resolve) => { + releaseFirstDispatch = resolve; + }); + const starts: string[] = []; + const processEvent = vi.fn(async (receiverEvent: ReceiverEvent) => { + const event = (receiverEvent.body as { event: { thread_ts: string } }).event; + const lifecycle = resolveSlackIngressTurnLifecycle(receiverEvent.customProperties); + await lifecycle?.onSessionRouted?.(`agent:main:slack:thread:${event.thread_ts}`); + starts.push(event.thread_ts); + if (event.thread_ts === "1700000000.000100") { + await firstDispatchGate; + } + await lifecycle?.onAdopted(); + }); + const { ingress, receive } = attachIngress(queue, processEvent); + ingress.start(); + + try { + for (const [eventId, threadTs, ts] of [ + ["Ev-thread-one", "1700000000.000100", "1700000000.000101"], + ["Ev-thread-two", "1700000000.000200", "1700000000.000201"], + ] as const) { + await receive( + createReceiverEvent(eventId, undefined, { + event: { + type: "message", + channel: "C_TEST", + channel_type: "channel", + user: "U_TEST", + thread_ts: threadTs, + ts, + text: "thread reply", + }, + }), + ); + } + + await vi.waitFor(() => expect(starts).toHaveLength(2), { timeout: 500 }); + expect(starts).toEqual(["1700000000.000100", "1700000000.000200"]); + } finally { + releaseFirstDispatch(); + await ingress.waitForIdle(); + await ingress.stop(); + } + }); + }); + + it.each<{ + name: string; + firstEvent: Record & { ts: string }; + secondEvent: Record & { ts: string }; + }>([ + { + name: "top-level channel messages", + firstEvent: { ts: "1700000000.000100" }, + secondEvent: { ts: "1700000000.000200" }, + }, + { + name: "threads bound to the same configured session", + firstEvent: { ts: "1700000000.000101", thread_ts: "1700000000.000100" }, + secondEvent: { ts: "1700000000.000201", thread_ts: "1700000000.000200" }, + }, + ])("serializes $name by their authoritative session", async ({ firstEvent, secondEvent }) => { + await withQueue(async (queue) => { + let releaseFirstDispatch: () => void = () => {}; + const firstDispatchGate = new Promise((resolve) => { + releaseFirstDispatch = resolve; + }); + const starts: string[] = []; + const processEvent = vi.fn(async (receiverEvent: ReceiverEvent) => { + const event = (receiverEvent.body as { event: { ts: string } }).event; + const lifecycle = resolveSlackIngressTurnLifecycle(receiverEvent.customProperties); + await lifecycle?.onSessionRouted?.("agent:main:slack:shared-session"); + starts.push(event.ts); + if (event.ts === firstEvent.ts) { + await firstDispatchGate; + } + await lifecycle?.onAdopted(); + }); + const { ingress, receive } = attachIngress(queue, processEvent); + ingress.start(); + + try { + for (const [eventId, event] of [ + ["Ev-shared-first", firstEvent], + ["Ev-shared-second", secondEvent], + ] as const) { + await receive( + createReceiverEvent(eventId, undefined, { + event: { + type: "message", + channel: "C_TEST", + channel_type: "channel", + user: "U_TEST", + text: "shared session", + ...event, + }, + }), + ); + } + + await vi.waitFor(() => expect(processEvent).toHaveBeenCalledTimes(2), { timeout: 500 }); + expect(starts).toEqual([firstEvent.ts]); + releaseFirstDispatch(); + await ingress.waitForIdle(); + expect(starts).toEqual([firstEvent.ts, secondEvent.ts]); + } finally { + releaseFirstDispatch(); + await ingress.waitForIdle(); + await ingress.stop(); + } + }); + }); + + it("keeps a queued same-session event alive past the adoption watchdog", async () => { + await withQueue(async (queue) => { + let releaseFirstSettlement: () => void = () => {}; + const firstSettlement = new Promise((resolve) => { + releaseFirstSettlement = resolve; + }); + const starts: string[] = []; + const processEvent = vi.fn(async (receiverEvent: ReceiverEvent) => { + const eventId = (receiverEvent.body as { event_id: string }).event_id; + const lifecycle = resolveSlackIngressTurnLifecycle(receiverEvent.customProperties); + await lifecycle?.onSessionRouted?.("agent:main:slack:shared-session"); + starts.push(eventId); + if (eventId === "Ev-session-watchdog-first") { + (lifecycle as ChannelIngressMonitorLifecycle).onAdoptionFinalizing(); + await firstSettlement; + } + await lifecycle?.onAdopted(); + }); + const { ingress, receive } = attachIngress(queue, processEvent, { + adoptionStallTimeoutMs: 80, + }); + ingress.start(); + + try { + await receive(createReceiverEvent("Ev-session-watchdog-first")); + await receive(createReceiverEvent("Ev-session-watchdog-second")); + await vi.waitFor(() => expect(processEvent).toHaveBeenCalledTimes(2)); + expect(starts).toEqual(["Ev-session-watchdog-first"]); + + await new Promise((resolve) => { + setTimeout(resolve, 120); + }); + await receive(createReceiverEvent("Ev-session-watchdog-third")); + await vi.waitFor(() => expect(processEvent).toHaveBeenCalledTimes(3)); + expect((await queue.listClaims()).map((claim) => claim.id)).toEqual([ + "Ev-session-watchdog-first", + "Ev-session-watchdog-second", + "Ev-session-watchdog-third", + ]); + expect(starts).toEqual(["Ev-session-watchdog-first"]); + + releaseFirstSettlement(); + await ingress.waitForIdle(); + expect(starts).toEqual([ + "Ev-session-watchdog-first", + "Ev-session-watchdog-second", + "Ev-session-watchdog-third", + ]); + expect(processEvent).toHaveBeenCalledTimes(3); + expect(await queue.listPending()).toEqual([]); + } finally { + releaseFirstSettlement(); + await ingress.waitForIdle(); + await ingress.stop(); + } + }); + }); + it("serializes new-channel messages behind channel-ID migration", async () => { await withQueue(async (queue) => { let markMigrationStarted: () => void = () => {}; @@ -337,8 +516,10 @@ describe("Slack durable ingress", () => { event: { type: "message", channel: "C_NEW", + channel_type: "channel", user: "U_TEST", ts: "1700000000.000200", + thread_ts: "1700000000.000100", text: "after migration", }, }), @@ -355,6 +536,103 @@ describe("Slack durable ingress", () => { }); }); + it.each([ + { name: "an already routed message", deferred: false }, + { name: "a deferred message", deferred: true }, + ])("serializes channel-ID migration behind $name through Bolt", async ({ deferred }) => { + await withQueue(async (queue) => { + let markMessageStarted: () => void = () => {}; + let releaseMessage: () => void = () => {}; + let releaseMigration: () => void = () => {}; + const messageStarted = new Promise((resolve) => { + markMessageStarted = resolve; + }); + const messageGate = new Promise((resolve) => { + releaseMessage = resolve; + }); + const migrationGate = new Promise((resolve) => { + releaseMigration = resolve; + }); + const starts: string[] = []; + const ingress = createSlackDurableIngress({ + accountId: "default", + queue, + pollIntervalMs: 60_000, + adoptionStallTimeoutMs: 5_000, + }); + const harness = createReceiverHarness(); + const app = new App({ + receiver: ingress.wrapReceiver(harness.receiver), + authorize: async () => ({ + botToken: "xoxb-test", + botId: "B_BOT", + botUserId: "U_BOT", + teamId: "T_TEST", + }), + convoStore: false, + ignoreSelf: false, + }); + app.event("message", async ({ context }) => { + const lifecycle = resolveSlackIngressTurnLifecycle(context); + await lifecycle?.onSessionRouted?.("agent:main:slack:thread:C_NEW"); + starts.push("message"); + if (deferred) { + lifecycle?.onDeferred(); + } + markMessageStarted(); + await messageGate; + await lifecycle?.onAdopted(); + }); + app.event("channel_id_changed", async ({ context }) => { + starts.push("channel_id_changed"); + await migrationGate; + await resolveSlackIngressTurnLifecycle(context)?.onAdopted(); + }); + ingress.start(); + + try { + await harness.receive( + createReceiverEventWithBody({ + ...createSlackEnvelope("Ev-routed-before-migration"), + event: { + type: "message", + channel: "C_NEW", + channel_type: "channel", + user: "U_TEST", + ts: "1700000000.000200", + thread_ts: "1700000000.000100", + text: "before migration", + }, + }), + ); + await messageStarted; + await harness.receive( + createReceiverEventWithBody( + createChannelIdChangedEnvelope("Ev-migration-after-route", "C_OLD", "C_NEW"), + ), + ); + await vi.waitFor(async () => { + expect((await queue.listClaims()).map((claim) => claim.id)).toEqual([ + "Ev-routed-before-migration", + "Ev-migration-after-route", + ]); + }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(starts).toEqual(["message"]); + + releaseMessage(); + await vi.waitFor(() => expect(starts).toEqual(["message", "channel_id_changed"])); + } finally { + releaseMessage(); + releaseMigration(); + await ingress.waitForIdle(); + await ingress.stop(); + } + }); + }); + it("drains a durable event when its acknowledgement fails", async () => { await withQueue(async (queue) => { const processEvent = vi.fn(async (event: ReceiverEvent) => { @@ -398,9 +676,20 @@ describe("Slack durable ingress", () => { }); }); - it("recovers a shipped row whose lane was derived only at drain time", async () => { + it.each([ + { name: "a lane derived only at drain time", laneKey: undefined }, + { name: "its persisted channel-only lane", laneKey: "team:T_TEST:conversation:C_TEST" }, + ])("recovers a shipped threaded row with $name", async ({ laneKey }) => { await withQueue(async (queue) => { - const body = createSlackEnvelope("Ev-legacy-lane"); + const body = createSlackEnvelope("Ev-legacy-lane", undefined, { + type: "message", + channel: "C_TEST", + channel_type: "channel", + user: "U_TEST", + ts: "1700000000.000101", + thread_ts: "1700000000.000100", + text: "persisted thread reply", + }); await queue.enqueue( "Ev-legacy-lane", { @@ -409,7 +698,7 @@ describe("Slack durable ingress", () => { kind: "events-api", body, }, - { receivedAt: 1_700_000_000_000 }, + { receivedAt: 1_700_000_000_000, ...(laneKey ? { laneKey } : {}) }, ); const dispatch = vi.fn(async (event: ReceiverEvent) => { await resolveSlackIngressTurnLifecycle(event.customProperties)?.onAdopted(); diff --git a/extensions/slack/src/monitor/ingress.ts b/extensions/slack/src/monitor/ingress.ts index be65f8814376..ea0ef9db6854 100644 --- a/extensions/slack/src/monitor/ingress.ts +++ b/extensions/slack/src/monitor/ingress.ts @@ -25,7 +25,9 @@ const SLACK_INGRESS_LIFECYCLE_CONTEXT_KEY = "openclawIngressLifecycle"; export type SlackIngressTurnLifecycle = Omit< ChannelIngressMonitorLifecycle, "onAdoptionFinalizing" ->; +> & { + onSessionRouted?: (sessionKey: string) => Promise; +}; type SlackIngressPayload = { version: number; @@ -223,6 +225,8 @@ export function createSlackDurableIngress( ): SlackDurableIngress { let app: App | undefined; let relayDispatch: SlackRelayIngressDispatch | undefined; + const activeSessionTurns = new Map>(); + const activeChannelTurns = new Map>>(); const monitor = createChannelIngressMonitor< SlackIngressRawEvent, SlackIngressBody, @@ -272,34 +276,159 @@ export function createSlackDurableIngress( await raw.afterDurableAdmission?.(); } }, - deliver: async (raw, lifecycle) => { - if (raw.kind === "relay") { - if (!relayDispatch) { - // Transient by design: a claim recovered before the relay source - // reattaches must retry, not dead-letter, or restart recovery loses it. - throw new Error("Slack relay ingress dispatcher is not attached."); - } - await relayDispatch(raw.message, lifecycle); - return; - } - if (!app) { - throw new Error("Slack ingress receiver is not attached to a Bolt app."); - } - await app.processEvent({ - body: raw.body as ReceiverEvent["body"], - ack: async () => {}, - ...(raw.retryNum === undefined ? {} : { retryNum: raw.retryNum }), - ...(raw.retryReason === undefined ? {} : { retryReason: raw.retryReason }), - customProperties: { - [SLACK_INGRESS_LIFECYCLE_CONTEXT_KEY]: lifecycle, + deliver: async (raw, lifecycle, claim) => { + const laneKey = claim.laneKey ?? inspectSlackIngress(raw).laneKey; + let releaseSession: (() => void) | undefined; + let releaseChannel: (() => void) | undefined; + let routedSession: string | undefined; + let migrationAwaitedChannelTurns = false; + let downstreamDeferred = false; + let settled = false; + const settleSession = () => { + settled = true; + releaseSession?.(); + }; + const settleTurn = () => { + settleSession(); + releaseChannel?.(); + }; + const routedLifecycle: SlackIngressTurnLifecycle = { + ...lifecycle, + onSessionRouted: async (sessionKey) => { + if (routedSession !== undefined) { + if (routedSession !== sessionKey) { + throw new Error("Slack ingress session ownership changed after routing."); + } + return; + } + lifecycle.abortSignal.throwIfAborted(); + routedSession = sessionKey; + const previousTurn = activeSessionTurns.get(sessionKey); + let resolveCurrentTurn: () => void = () => {}; + const releasedCurrentTurn = new Promise((resolve) => { + resolveCurrentTurn = resolve; + }); + const currentTurn = previousTurn + ? previousTurn.then(() => releasedCurrentTurn) + : releasedCurrentTurn; + activeSessionTurns.set(sessionKey, currentTurn); + let resolveChannelTurn: () => void = () => {}; + const channelTurn = new Promise((resolve) => { + resolveChannelTurn = resolve; + }); + const channelTurns = activeChannelTurns.get(laneKey) ?? new Set>(); + channelTurns.add(channelTurn); + activeChannelTurns.set(laneKey, channelTurns); + const releaseCurrentSession = () => { + lifecycle.abortSignal.removeEventListener("abort", releaseCurrentSession); + resolveCurrentTurn(); + }; + const releaseCurrentChannel = () => { + lifecycle.abortSignal.removeEventListener("abort", releaseCurrentChannel); + resolveChannelTurn(); + }; + void currentTurn.then(() => { + if (activeSessionTurns.get(sessionKey) === currentTurn) { + activeSessionTurns.delete(sessionKey); + } + }); + void channelTurn.then(() => { + channelTurns.delete(channelTurn); + if (channelTurns.size === 0 && activeChannelTurns.get(laneKey) === channelTurns) { + activeChannelTurns.delete(laneKey); + } + }); + releaseSession = releaseCurrentSession; + releaseChannel = releaseCurrentChannel; + lifecycle.abortSignal.addEventListener("abort", releaseSession, { once: true }); + lifecycle.abortSignal.addEventListener("abort", releaseChannel, { once: true }); + // Preserve shipped channel lanes until the prepared route proves its + // session; channel-ID migration therefore still fences all traffic. + lifecycle.onDeferred(); + if (previousTurn) { + // A queued session turn owns its durable claim; its predecessor may + // legitimately outlive the pre-adoption watchdog. + lifecycle.onAdoptionFinalizing(); + } + monitor.requestDrain(); + await previousTurn; + lifecycle.abortSignal.throwIfAborted(); }, - }); + onAdopted: async () => { + try { + await lifecycle.onAdopted(); + } finally { + settleTurn(); + } + }, + onDeferred: () => { + downstreamDeferred = true; + lifecycle.onDeferred(); + // Reply handoff releases session order; migration remains fenced + // until the durable turn is actually adopted or abandoned. + settleSession(); + monitor.requestDrain(); + }, + onAbandoned: async () => { + try { + await lifecycle.onAbandoned(); + } finally { + settleTurn(); + } + }, + }; + try { + const event = raw.kind === "events-api" ? asOptionalRecord(raw.body)?.event : undefined; + if (asOptionalRecord(event)?.type === "channel_id_changed") { + const channelTurns = activeChannelTurns.get(laneKey); + if (channelTurns && channelTurns.size > 0) { + // A migration owns the channel lane while earlier routed sessions + // settle; later channel traffic cannot overtake the config change. + migrationAwaitedChannelTurns = true; + lifecycle.onAdoptionFinalizing(); + await Promise.all(channelTurns); + lifecycle.abortSignal.throwIfAborted(); + } + } + if (raw.kind === "relay") { + if (!relayDispatch) { + // Transient by design: a claim recovered before the relay source + // reattaches must retry, not dead-letter, or restart recovery loses it. + throw new Error("Slack relay ingress dispatcher is not attached."); + } + await relayDispatch(raw.message, routedLifecycle); + } else { + if (!app) { + throw new Error("Slack ingress receiver is not attached to a Bolt app."); + } + await app.processEvent({ + body: raw.body as ReceiverEvent["body"], + ack: async () => {}, + ...(raw.retryNum === undefined ? {} : { retryNum: raw.retryNum }), + ...(raw.retryReason === undefined ? {} : { retryReason: raw.retryReason }), + customProperties: { + [SLACK_INGRESS_LIFECYCLE_CONTEXT_KEY]: routedLifecycle, + }, + }); + } + if ( + (routedSession !== undefined || migrationAwaitedChannelTurns) && + !settled && + !downstreamDeferred + ) { + await routedLifecycle.onAdopted(); + } + } catch (error) { + settleTurn(); + throw error; + } }, pollIntervalMs: options.pollIntervalMs ?? SLACK_INGRESS_POLL_INTERVAL_MS, retention: "standard", appendRetryDelaysMs: [0], drain: { resolveNonRetryableFailure: resolveSlackIngressNonRetryableFailure, + deferredLaneOccupancy: "release", // Shipped Slack rows did not store lanes, so replay still derives them from payloads. deriveLaneKey: (record) => record.payload.kind === "relay" diff --git a/extensions/slack/src/monitor/message-handler.test.ts b/extensions/slack/src/monitor/message-handler.test.ts index 8bc7b3449380..19605b1c2418 100644 --- a/extensions/slack/src/monitor/message-handler.test.ts +++ b/extensions/slack/src/monitor/message-handler.test.ts @@ -21,7 +21,10 @@ const prepareSlackMessageMock = vi.fn( async (_params?: { ctx: Parameters[0]["ctx"]; opts: { onVisibleDrop?: () => void }; - }): Promise<{ ctxPayload: Record } | null> => ({ ctxPayload: {} }), + }): Promise<{ + ctxPayload: Record; + route?: { sessionKey: string }; + } | null> => ({ ctxPayload: {} }), ); const dispatchPreparedSlackMessageMock = vi.fn(async (_prepared: unknown) => {}); const resolveThreadTsMock = vi.fn(async ({ message }: { message: Record }) => ({ @@ -632,12 +635,17 @@ describe("createSlackMessageHandler", () => { }); it("carries durable ingress ownership into prepared dispatch", async () => { + prepareSlackMessageMock.mockResolvedValueOnce({ + ctxPayload: {}, + route: { sessionKey: "agent:main:slack:channel:C111" }, + }); const turnAdoptionLifecycle = { admission: "exclusive" as const, abortSignal: new AbortController().signal, onAdopted: vi.fn(), onDeferred: vi.fn(), onAbandoned: vi.fn(), + onSessionRouted: vi.fn(async () => {}), }; const { handler } = createHandlerWithTracker(); const handled = handler( @@ -664,6 +672,12 @@ describe("createSlackMessageHandler", () => { // The flush wraps the lifecycle to settle dispatch-dedupe claims, so assert // ownership forwarding rather than function identity. expect(dispatchPreparedSlackMessageMock).toHaveBeenCalledTimes(1); + expect(turnAdoptionLifecycle.onSessionRouted).toHaveBeenCalledExactlyOnceWith( + "agent:main:slack:channel:C111", + ); + expect(turnAdoptionLifecycle.onSessionRouted.mock.invocationCallOrder[0]).toBeLessThan( + dispatchPreparedSlackMessageMock.mock.invocationCallOrder[0] ?? 0, + ); const prepared = dispatchPreparedSlackMessageMock.mock.calls[0]?.[0] as { turnAdoptionLifecycle?: typeof turnAdoptionLifecycle; }; diff --git a/extensions/slack/src/monitor/message-handler.ts b/extensions/slack/src/monitor/message-handler.ts index 6f913eee7823..03eda4bc9ca5 100644 --- a/extensions/slack/src/monitor/message-handler.ts +++ b/extensions/slack/src/monitor/message-handler.ts @@ -343,6 +343,7 @@ export function createSlackMessageHandler(params: { releaseClaims(); return; } + await turnAdoptionLifecycle?.onSessionRouted?.(prepared.route.sessionKey); // Commit at adoption (durable turn ownership), release on abandonment; // deferred turns hand settlement to the reply lane with the claim held. prepared.turnAdoptionLifecycle = {