diff --git a/config/assertion-safety-baseline.txt b/config/assertion-safety-baseline.txt index c9b1f8c35248..471ec4af1e0d 100644 --- a/config/assertion-safety-baseline.txt +++ b/config/assertion-safety-baseline.txt @@ -921,7 +921,7 @@ extensions/qa-lab/src/ci-smoke-plan.ts 1 extensions/qa-lab/src/cli.runtime.ts 9 extensions/qa-lab/src/codex-plugin.fixture.ts 2 extensions/qa-lab/src/confidence-report.ts 2 -extensions/qa-lab/src/crabline-transport.ts 6 +extensions/qa-lab/src/crabline-transport.ts 5 extensions/qa-lab/src/cron-run-wait.ts 1 extensions/qa-lab/src/discovery-eval.ts 1 extensions/qa-lab/src/docker-up.runtime.ts 1 diff --git a/extensions/qa-channel/src/channel-actions.test.ts b/extensions/qa-channel/src/channel-actions.test.ts index 615bdbd0590e..1506cb74c2dd 100644 --- a/extensions/qa-channel/src/channel-actions.test.ts +++ b/extensions/qa-channel/src/channel-actions.test.ts @@ -74,10 +74,12 @@ describe("qa-channel direct message actions", () => { const threadPayload = extractToolPayload(threadResult) as { thread: { id: string; title: string }; target: string; + threadId: string; }; expect(threadPayload.thread.id).toMatch(/^thread-/); expect(threadPayload.thread.title).toBe("QA thread"); - expect(threadPayload.target).toContain(threadPayload.thread.id); + expect(threadPayload.target).toBe("channel:qa-room"); + expect(threadPayload.threadId).toBe(threadPayload.thread.id); const replyResult = await handleAction({ channel: "qa-channel", @@ -86,6 +88,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { target: threadPayload.target, + threadId: threadPayload.threadId, message: "thread reply", text: "ignored legacy reply", }, @@ -111,6 +114,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { to: threadPayload.target, + threadId: threadPayload.threadId, messageId: outbound.id, emoji: "white_check_mark", }, @@ -123,6 +127,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { target: threadPayload.target, + threadId: threadPayload.threadId, messageId: outbound.id, message: "message (edited)", text: "ignored legacy edit", @@ -136,6 +141,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { to: threadPayload.target, + threadId: threadPayload.threadId, messageId: outbound.id, }, }); @@ -165,6 +171,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { to: threadPayload.target, + threadId: threadPayload.threadId, messageId: outbound.id, }, }); @@ -188,7 +195,7 @@ describe("qa-channel direct message actions", () => { cfg, accountId: "default", params: { - target: "channel:canonical-room", + target: "channel:canonical/room", to: "channel:legacy-to-room", channelId: "legacy-channel-id-room", threadName: "Canonical target thread", @@ -197,10 +204,12 @@ describe("qa-channel direct message actions", () => { const threadPayload = extractToolPayload(threadResult) as { thread: { id: string; conversationId: string }; target: string; + threadId: string; }; - expect(threadPayload.thread.conversationId).toBe("canonical-room"); - expect(threadPayload.target).toBe(`thread:canonical-room/${threadPayload.thread.id}`); + expect(threadPayload.thread.conversationId).toBe("canonical/room"); + expect(threadPayload.target).toBe("channel:canonical/room"); + expect(threadPayload.threadId).toBe(threadPayload.thread.id); const replyResult = await handleAction({ channel: "qa-channel", @@ -209,13 +218,14 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { target: threadPayload.target, + threadId: threadPayload.threadId, channelId: "legacy-reply-room", message: "canonical target reply", }, }); expect(extractToolPayload(replyResult)).toMatchObject({ message: { - conversation: { id: "canonical-room", kind: "channel" }, + conversation: { id: "canonical/room", kind: "channel" }, text: "canonical target reply", threadId: threadPayload.thread.id, }, @@ -291,6 +301,7 @@ describe("qa-channel direct message actions", () => { const threadPayload = extractToolPayload(threadResult) as { thread: { id: string; title: string }; target: string; + threadId: string; }; expect(threadPayload.thread.title).toBe("Legacy thread"); @@ -324,6 +335,7 @@ describe("qa-channel direct message actions", () => { accountId: "default", params: { to: threadPayload.target, + threadId: threadPayload.threadId, messageId: outbound.id, text: "legacy edit", }, diff --git a/extensions/qa-channel/src/channel-actions.ts b/extensions/qa-channel/src/channel-actions.ts index 07c451ddc40b..f62897723a66 100644 --- a/extensions/qa-channel/src/channel-actions.ts +++ b/extensions/qa-channel/src/channel-actions.ts @@ -147,7 +147,12 @@ export const qaChannelMessageActions: ChannelMessageActionAdapter = { if (action === "thread-reply") { const channelId = typeof args.channelId === "string" ? args.channelId.trim() : ""; const threadId = typeof args.threadId === "string" ? args.threadId.trim() : ""; - return channelId && threadId ? { to: `thread:${channelId}/${threadId}` } : null; + return channelId && threadId + ? { + to: buildQaTarget({ chatType: "channel", conversationId: channelId }), + threadId, + } + : null; } return null; }, @@ -191,7 +196,6 @@ export const qaChannelMessageActions: ChannelMessageActionAdapter = { to: buildQaTarget({ chatType: parsed.chatType, conversationId: parsed.conversationId, - threadId: resolved.threadId, }), text, senderId: account.botUserId, @@ -217,7 +221,11 @@ export const qaChannelMessageActions: ChannelMessageActionAdapter = { }); return jsonResult({ thread, - target: `thread:${target.conversationId}/${thread.id}`, + target: buildQaTarget({ + chatType: target.conversationKind, + conversationId: target.conversationId, + }), + threadId: thread.id, }); } case "thread-reply": { @@ -234,7 +242,6 @@ export const qaChannelMessageActions: ChannelMessageActionAdapter = { to: buildQaTarget({ chatType: target.conversationKind, conversationId: target.conversationId, - threadId: target.threadId, }), text, senderId: account.botUserId, diff --git a/extensions/qa-channel/src/channel-thread-routing.test.ts b/extensions/qa-channel/src/channel-thread-routing.test.ts new file mode 100644 index 000000000000..a04f9d00d47f --- /dev/null +++ b/extensions/qa-channel/src/channel-thread-routing.test.ts @@ -0,0 +1,48 @@ +import { describe, expect, it } from "vitest"; +import { qaChannelPlugin } from "../api.js"; + +describe("qa-channel structured thread routing", () => { + it("derives thread-aware outbound session routes from explicit thread targets", async () => { + const route = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ + cfg: {}, + agentId: "main", + accountId: "default", + target: "thread:qa-room/thread-1", + }); + + expect(route?.sessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room:thread:thread-1"); + expect(route?.baseSessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room"); + expect(route?.threadId).toBe("thread-1"); + }); + + it("does not duplicate routing metadata on explicit thread targets", async () => { + const route = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ + cfg: {}, + agentId: "main", + accountId: "default", + target: "thread:qa-room/thread-1", + replyToId: "reply-1", + threadId: "thread-1", + currentSessionKey: "agent:main:qa-channel:channel:thread:qa-room/thread-1:thread:stale", + }); + + expect(route?.sessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room:thread:thread-1"); + expect(route?.baseSessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room"); + expect(route?.threadId).toBe("thread-1"); + }); + + it("keeps structured thread identity authoritative over reply metadata", async () => { + const route = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ + cfg: {}, + agentId: "main", + accountId: "default", + target: "channel:qa-room", + replyToId: "reply-1", + threadId: "thread-1", + }); + + expect(route?.sessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room:thread:thread-1"); + expect(route?.baseSessionKey).toBe("agent:main:qa-channel:channel:channel:qa-room"); + expect(route?.threadId).toBe("thread-1"); + }); +}); diff --git a/extensions/qa-channel/src/channel.test.ts b/extensions/qa-channel/src/channel.test.ts index 584f5b05e954..3dc700a1f430 100644 --- a/extensions/qa-channel/src/channel.test.ts +++ b/extensions/qa-channel/src/channel.test.ts @@ -185,35 +185,6 @@ async function startQaChannelTestHarness(params?: { } describe("qa-channel plugin", () => { - it("derives thread-aware outbound session routes from explicit thread targets", async () => { - const route = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ - cfg: {}, - agentId: "main", - accountId: "default", - target: "thread:qa-room/thread-1", - }); - - expect(route?.sessionKey).toBe("agent:main:qa-channel:channel:thread:qa-room/thread-1"); - expect(route?.baseSessionKey).toBe("agent:main:qa-channel:channel:thread:qa-room/thread-1"); - expect(route?.threadId).toBeUndefined(); - }); - - it("does not append routing metadata to explicit thread targets", async () => { - const route = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ - cfg: {}, - agentId: "main", - accountId: "default", - target: "thread:qa-room/thread-1", - replyToId: "reply-1", - threadId: "thread-1", - currentSessionKey: "agent:main:qa-channel:channel:thread:qa-room/thread-1:thread:stale", - }); - - expect(route?.sessionKey).toBe("agent:main:qa-channel:channel:thread:qa-room/thread-1"); - expect(route?.baseSessionKey).toBe("agent:main:qa-channel:channel:thread:qa-room/thread-1"); - expect(route?.threadId).toBeUndefined(); - }); - it("rejects conflicting explicit thread routing metadata", () => { expect(() => qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ @@ -987,6 +958,18 @@ describe("qa-channel plugin", () => { }); expect(sendTarget).toEqual({ to: "channel:qa-room", threadId: undefined }); + const legacyThreadTarget = qaChannelPlugin.actions?.extractToolSend?.({ + args: { + action: "thread-reply", + channelId: "canonical/room", + threadId: "thread-1", + }, + }); + expect(legacyThreadTarget).toEqual({ + to: "channel:canonical/room", + threadId: "thread-1", + }); + const result = await qaChannelPlugin.actions?.handleAction?.({ channel: "qa-channel", action: "send", diff --git a/extensions/qa-channel/src/channel.threading.test.ts b/extensions/qa-channel/src/channel.threading.test.ts index 3a9978845933..526402f2ebd0 100644 --- a/extensions/qa-channel/src/channel.threading.test.ts +++ b/extensions/qa-channel/src/channel.threading.test.ts @@ -211,7 +211,7 @@ describe("qa-channel thread delivery contracts", () => { }); }); - it("extracts thread replies as canonical QA thread targets", () => { + it("extracts thread replies with structured thread identity", () => { expect( qaChannelPlugin.actions?.extractToolSend?.({ args: { @@ -221,6 +221,6 @@ describe("qa-channel thread delivery contracts", () => { message: "hello thread", }, }), - ).toEqual({ to: "thread:qa-room/thread-1" }); + ).toEqual({ to: "channel:qa-room", threadId: "thread-1" }); }); }); diff --git a/extensions/qa-channel/src/channel.ts b/extensions/qa-channel/src/channel.ts index 32303ea6ac26..249b73a437a6 100644 --- a/extensions/qa-channel/src/channel.ts +++ b/extensions/qa-channel/src/channel.ts @@ -190,6 +190,10 @@ export const qaChannelPlugin: ChannelPlugin = createCh }) => { const resolved = resolveQaTargetThread({ target, threadId }); const parsed = resolved.target; + const baseTarget = buildQaTarget({ + chatType: parsed.chatType, + conversationId: parsed.conversationId, + }); const baseRoute = buildChannelOutboundSessionRoute({ cfg, agentId, @@ -203,20 +207,16 @@ export const qaChannelPlugin: ChannelPlugin = createCh : parsed.chatType === "group" ? "group" : "channel", - id: buildQaTarget(parsed), + id: baseTarget, }, chatType: parsed.chatType, from: `${QA_CHANNEL_ID}:${accountId ?? DEFAULT_ACCOUNT_ID}`, - to: buildQaTarget(parsed), + to: baseTarget, }); - // An explicit thread target already owns the complete session identity; - // applying reply or current-thread metadata would append a second thread. - if (parsed.threadId !== undefined) { - return baseRoute; - } return buildThreadAwareOutboundSessionRoute({ route: baseRoute, - replyToId, + // Structured thread identity is authoritative; reply metadata must not replace its thread. + replyToId: resolved.threadId === undefined ? replyToId : undefined, threadId: resolved.threadId, currentSessionKey, canRecoverCurrentThread: ({ route }) => diff --git a/extensions/qa-channel/src/inbound.test.ts b/extensions/qa-channel/src/inbound.test.ts index 39a88050139b..a22894f9d8f9 100644 --- a/extensions/qa-channel/src/inbound.test.ts +++ b/extensions/qa-channel/src/inbound.test.ts @@ -6,6 +6,7 @@ import { loadOutboundMediaFromUrl } from "openclaw/plugin-sdk/outbound-media"; import { beforeEach, describe, expect, it, vi } from "vitest"; import { setQaChannelRuntime } from "../api.js"; import { deleteQaBusMessage, editQaBusMessage, sendQaBusMessage } from "./bus-client.js"; +import { qaChannelPlugin } from "./channel.js"; import { handleQaInbound } from "./inbound.js"; import { createQaInboundParams, @@ -110,9 +111,18 @@ describe("handleQaInbound", () => { replyToId: "msg-1", text: "preview", threadId: "42", - to: "thread:/v1/group/qa-room/42", + to: "group:qa-room", }), ); + expect(assembled.ctxPayload).toMatchObject({ + To: "group:qa-room", + OriginatingTo: "group:qa-room", + SessionKey: expect.stringMatching(/:thread:42$/u), + }); + expect(assembled.ctxPayload.ParentSessionKey).toBe( + assembled.ctxPayload.SessionKey?.replace(/:thread:42$/u, ""), + ); + expect(assembled.route.sessionKey).toBe(assembled.ctxPayload.SessionKey); expect(editQaBusMessage).toHaveBeenNthCalledWith( 1, expect.objectContaining({ messageId: "preview-1", text: "preview expanded" }), @@ -123,6 +133,30 @@ describe("handleQaInbound", () => { ); }); + it("uses one session for inbound channel threads and explicit thread replies", async () => { + const runtime = createPluginRuntimeMock(); + setQaChannelRuntime(runtime); + + await handleQaInbound( + createQaInboundParams({ + message: { + conversation: { id: "qa-room", kind: "channel" }, + threadId: "42", + }, + }), + ); + + const assembled = firstRunAssembledParams(runtime); + const outboundRoute = await qaChannelPlugin.messaging?.resolveOutboundSessionRoute?.({ + cfg: {}, + agentId: "main", + accountId: "default", + target: "thread:qa-room/42", + }); + + expect(outboundRoute?.sessionKey).toBe(assembled.route.sessionKey); + }); + it("treats deliveries without dispatcher metadata as final replies", async () => { const runtime = createPluginRuntimeMock(); setQaChannelRuntime(runtime); diff --git a/extensions/qa-channel/src/inbound.ts b/extensions/qa-channel/src/inbound.ts index 9cb2d03eaa62..1fb55c227eaa 100644 --- a/extensions/qa-channel/src/inbound.ts +++ b/extensions/qa-channel/src/inbound.ts @@ -1,6 +1,7 @@ import { createAsyncLock } from "openclaw/plugin-sdk/async-lock-runtime"; import { buildChannelInboundEventContext, + createChannelInboundEnvelopeBuilder, formatInboundMediaUnavailableText, resolveChannelInboundRouteEnvelope, toInboundMediaFactsWithMetadata, @@ -15,6 +16,7 @@ import { sanitizeQaBusToolCallArguments, type QaBusToolCall, } from "openclaw/plugin-sdk/qa-channel-protocol"; +import { resolveThreadSessionKeys } from "openclaw/plugin-sdk/routing"; import { buildQaTarget, deleteQaBusMessage, @@ -319,10 +321,9 @@ export async function handleQaInbound(params: { const target = buildQaTarget({ chatType: inbound.conversation.kind, conversationId: inbound.conversation.id, - threadId: inbound.threadId, }); const toolCalls: QaBusToolCall[] = []; - const { route, buildEnvelope } = resolveChannelInboundRouteEnvelope({ + const { route } = resolveChannelInboundRouteEnvelope({ cfg: params.config, channel: params.channelId, accountId: params.account.accountId, @@ -345,6 +346,15 @@ export async function handleQaInbound(params: { mediaLocalRoots: getAgentScopedMediaLocalRoots(params.config, route.agentId), }); const isGroup = inbound.conversation.kind !== "direct"; + const threadKeys = resolveThreadSessionKeys({ + baseSessionKey: route.sessionKey, + threadId: inbound.threadId, + parentSessionKey: isGroup ? route.sessionKey : undefined, + }); + const buildEnvelope = createChannelInboundEnvelopeBuilder({ + cfg: params.config, + route: { agentId: route.agentId, sessionKey: threadKeys.sessionKey }, + }); const wasMentioned = isGroup ? channelRuntime.mentions.matchesMentionPatterns( inbound.text, @@ -364,10 +374,10 @@ export async function handleQaInbound(params: { agentId: route.agentId, sessionPrefix: "qa-channel:slash", userId: inbound.senderId, - targetSessionKey: route.sessionKey, + targetSessionKey: threadKeys.sessionKey, }) : undefined; - const sessionKey = commandTargets?.sessionKey ?? route.sessionKey; + const sessionKey = commandTargets?.sessionKey ?? threadKeys.sessionKey; const access = await resolveStableChannelMessageIngress({ cfg: params.config, channelId: params.channelId, @@ -450,6 +460,7 @@ export async function handleQaInbound(params: { accountId: route.accountId, routeSessionKey: sessionKey, dispatchSessionKey: sessionKey, + parentSessionKey: threadKeys.parentSessionKey, }, reply: { to: target, @@ -487,7 +498,7 @@ export async function handleQaInbound(params: { cfg: params.config, channel: params.channelId, accountId: params.account.accountId, - route: { agentId: route.agentId, dmScope: route.dmScope, sessionKey: route.sessionKey }, + route: { agentId: route.agentId, dmScope: route.dmScope, sessionKey: threadKeys.sessionKey }, ctxPayload, delivery: { deliver: async (payload, info) => { diff --git a/extensions/qa-channel/src/outbound.ts b/extensions/qa-channel/src/outbound.ts index d530611c3918..d665e75cc5cf 100644 --- a/extensions/qa-channel/src/outbound.ts +++ b/extensions/qa-channel/src/outbound.ts @@ -38,7 +38,6 @@ export async function sendQaChannelText(params: QaChannelTextSendParams) { to: buildQaTarget({ chatType: parsed.chatType, conversationId: parsed.conversationId, - threadId: resolved.threadId, }), text: params.text, isError: params.isError, diff --git a/extensions/qa-lab/package.json b/extensions/qa-lab/package.json index a3ca517e046b..276182419fe4 100644 --- a/extensions/qa-lab/package.json +++ b/extensions/qa-lab/package.json @@ -17,6 +17,7 @@ "devDependencies": { "@openclaw/crabbox-provider": "workspace:*", "@openclaw/discord": "workspace:*", + "@openclaw/mattermost": "workspace:*", "@openclaw/matrix": "workspace:*", "@openclaw/plugin-sdk": "workspace:*", "@openclaw/slack": "workspace:*", diff --git a/extensions/qa-lab/src/crabline-discord-thread-delivery.test.ts b/extensions/qa-lab/src/crabline-discord-thread-delivery.test.ts new file mode 100644 index 000000000000..9129cc8b8c69 --- /dev/null +++ b/extensions/qa-lab/src/crabline-discord-thread-delivery.test.ts @@ -0,0 +1,73 @@ +// Integration tests preserve provider-native Discord thread ids through QA Lab and final send routing. +import { discordPlugin } from "@openclaw/discord/channel-plugin-api.js"; +import { withTempDir } from "openclaw/plugin-sdk/test-env"; +import { describe, expect, it, vi } from "vitest"; +import { createQaBusState } from "./bus-state.js"; +import { createQaCrablineTransportAdapter } from "./crabline-transport.js"; +import { startAgentRun } from "./suite-runtime-agent-process.js"; + +describe("QA Crabline Discord thread delivery", () => { + it.each([ + { taskTracking: true, method: "agent", toField: "to", threadField: "threadId" }, + { + taskTracking: false, + method: "chat.send", + toField: "originatingTo", + threadField: "originatingThreadId", + }, + ] as const)("routes symbolic threads natively through $method", async (testCase) => { + await withTempDir("qa-crabline-discord-thread-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: { + capabilityMatrixPath: "crabline-channel-driver-capabilities.json", + channel: "discord", + channelDriver: "crabline", + providerReadinessArtifactPath: "crabline-provider-readiness.json", + }, + state: createQaBusState(), + }); + const gatewayCall = vi.fn(async (_method: string, _payload: Record) => ({ + runId: `run-${testCase.method}`, + })); + + try { + await startAgentRun({ gateway: { call: gatewayCall }, transport } as never, { + sessionKey: `agent:qa:${testCase.method}`, + message: "Discord thread delivery proof", + to: "group:discord-crabline-primary", + threadId: "discord-crabline-thread", + taskTracking: testCase.taskTracking, + }); + + const call = gatewayCall.mock.calls[0]; + if (!call) { + throw new Error("Gateway call was not recorded"); + } + const [method, payload] = call; + expect(method).toBe(testCase.method); + const to = String(payload[testCase.toField]); + const threadId = String(payload[testCase.threadField]); + expect(to).toMatch(/^channel:\d{17,20}$/u); + expect(threadId).toMatch(/^\d{17,20}$/u); + + const sendDiscord = vi.fn(async () => ({ messageId: `sent-${testCase.method}` })); + await discordPlugin.outbound!.sendText!({ + cfg: {}, + to, + threadId, + text: "Discord thread delivery proof", + silent: true, + deps: { discord: sendDiscord }, + }); + expect(sendDiscord).toHaveBeenCalledWith( + `channel:${threadId}`, + "Discord thread delivery proof", + expect.any(Object), + ); + } finally { + await transport.cleanupAfterGatewayStop?.(); + } + }); + }); +}); diff --git a/extensions/qa-lab/src/crabline-discord-transport.test.ts b/extensions/qa-lab/src/crabline-discord-transport.test.ts index 2819dff78114..b4981bacc18d 100644 --- a/extensions/qa-lab/src/crabline-discord-transport.test.ts +++ b/extensions/qa-lab/src/crabline-discord-transport.test.ts @@ -57,7 +57,8 @@ describe("Crabline Discord transport", () => { threadId: "discord-crabline-thread", }); const delivery = transport.buildAgentDelivery({ - target: "thread:discord-crabline-primary/discord-crabline-thread", + target: "group:discord-crabline-primary", + threadId: "discord-crabline-thread", }); expect(delivery).toMatchObject({ channel: "discord", diff --git a/extensions/qa-lab/src/crabline-provider-targets.ts b/extensions/qa-lab/src/crabline-provider-targets.ts index d7c26d7f85eb..64c5f7e007cc 100644 --- a/extensions/qa-lab/src/crabline-provider-targets.ts +++ b/extensions/qa-lab/src/crabline-provider-targets.ts @@ -2,14 +2,11 @@ import { createHash } from "node:crypto"; import type { OpenClawCrablineInbound, OpenClawCrablineInboundInput, - StartedOpenClawCrablineAdapter, StartedOpenClawCrablineCorrelatedAdapter, } from "@openclaw/crabline"; import { parseQaTarget } from "./qa-bus-protocol.js"; import type { QaBusInboundMessageInput } from "./runtime-api.js"; -const TELEGRAM_QA_DRIVER_ID = "100001"; -const TELEGRAM_QA_OBSERVER_ID = "100002"; const MATRIX_QA_SERVER_NAME = "matrix-qa.test"; const MATRIX_QA_DRIVER_ID = `@driver:${MATRIX_QA_SERVER_NAME}`; const DISCORD_ID_PATTERN = /^\d{17,20}$/u; @@ -24,14 +21,6 @@ export function resolveDiscordQaId(value: string) { return String(DISCORD_ID_FLOOR + (digest % DISCORD_ID_FLOOR)); } -export function resolveTelegramQaSenderId(senderId: string) { - return senderId === "driver" - ? TELEGRAM_QA_DRIVER_ID - : senderId === "observer" - ? TELEGRAM_QA_OBSERVER_ID - : senderId; -} - function resolveMatrixQaSenderId(senderId: string) { return senderId === "driver" ? MATRIX_QA_DRIVER_ID @@ -45,8 +34,9 @@ function resolveMatrixQaConversationId(conversationId: string) { if (!trimmed) { throw new Error("Matrix QA conversation id must be non-empty"); } - if (trimmed.startsWith("!") && trimmed.includes(":")) { - return trimmed; + const explicitTarget = normalizeExplicitMatrixTarget(trimmed); + if (explicitTarget) { + return explicitTarget; } const digest = createHash("sha256").update(trimmed).digest("hex").slice(0, 16); return `!${digest}:${MATRIX_QA_SERVER_NAME}`; @@ -139,7 +129,7 @@ function resolveDiscordQaTarget(target: string) { } export function createCrablineProviderInboundInput( - adapter: StartedOpenClawCrablineAdapter, + adapter: StartedOpenClawCrablineCorrelatedAdapter, input: QaBusInboundMessageInput, ): OpenClawCrablineInboundInput { const kind = input.conversation.kind === "direct" ? "direct" : "group"; @@ -156,13 +146,11 @@ export function createCrablineProviderInboundInput( kind, }, senderId: - adapter.channel === "telegram" - ? resolveTelegramQaSenderId(input.senderId) - : adapter.channel === "matrix" - ? resolveMatrixQaSenderId(input.senderId) - : adapter.channel === "discord" - ? resolveDiscordQaId(input.senderId) - : input.senderId, + adapter.channel === "matrix" + ? resolveMatrixQaSenderId(input.senderId) + : adapter.channel === "discord" + ? resolveDiscordQaId(input.senderId) + : input.senderId, text: adapter.channel === "matrix" && adapter.manifest.provider === "matrix" ? resolveMatrixQaText(input.text, adapter.manifest.botUserId) @@ -176,7 +164,7 @@ export function createCrablineProviderInboundInput( } export function resolveCrablineStateConversation(params: { - adapter: StartedOpenClawCrablineAdapter; + adapter: StartedOpenClawCrablineCorrelatedAdapter; input: QaBusInboundMessageInput; providerInbound: OpenClawCrablineInbound; }) { @@ -188,6 +176,7 @@ export function resolveCrablineStateConversation(params: { export function createCrablineProviderDelivery( adapter: Pick, target: string, + threadId?: string, ) { const { providerTargetKey, ...delivery } = adapter.createAgentDelivery({ target: @@ -196,6 +185,21 @@ export function createCrablineProviderDelivery( : adapter.channel === "discord" ? resolveDiscordQaTarget(target) : target, + threadId: adapter.channel === "discord" && threadId ? resolveDiscordQaId(threadId) : threadId, }); return { delivery, providerTargetKey }; } + +export function createCrablineProviderCorrelation( + adapter: StartedOpenClawCrablineCorrelatedAdapter, + target: Pick, +) { + return adapter.createInbound({ + input: createCrablineProviderInboundInput(adapter, { + conversation: target.conversation, + senderId: target.conversation.kind === "direct" ? target.conversation.id : "driver", + text: "QA provider correlation", + ...(target.threadId ? { threadId: target.threadId } : {}), + }), + }); +} diff --git a/extensions/qa-lab/src/crabline-transport-provider-identity.test.ts b/extensions/qa-lab/src/crabline-transport-provider-identity.test.ts new file mode 100644 index 000000000000..86402ae901a3 --- /dev/null +++ b/extensions/qa-lab/src/crabline-transport-provider-identity.test.ts @@ -0,0 +1,137 @@ +// Qa Lab tests cover provider-authoritative Telegram target correlation. +import type { OpenClawCrablineChannelDriverSelection } from "@openclaw/crabline"; +import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; +import { withTempDir } from "openclaw/plugin-sdk/test-env"; +import { describe, expect, it } from "vitest"; +import { createQaBusState } from "./bus-state.js"; +import { createQaCrablineTransportAdapter } from "./crabline-transport.js"; + +const selection = { + capabilityMatrixPath: "crabline-channel-driver-capabilities.json", + channel: "telegram", + channelDriver: "crabline", + providerReadinessArtifactPath: "crabline-provider-readiness.json", +} as const satisfies OpenClawCrablineChannelDriverSelection; + +function requireString(value: unknown, label: string): string { + if (typeof value !== "string" || value.length === 0) { + throw new Error(`${label} is required`); + } + return value; +} + +async function readTelegramInbound( + transport: Awaited>, +) { + const config = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" }); + const telegram = config.channels?.telegram as { apiRoot?: string; botToken?: string } | undefined; + const apiRoot = requireString(telegram?.apiRoot, "Telegram API root"); + const botToken = requireString(telegram?.botToken, "Telegram bot token"); + const response = await fetch(`${apiRoot}/bot${botToken}/getUpdates`); + const updates = (await response.json()) as { + result?: Array<{ message?: { chat?: { id?: number }; message_thread_id?: number } }>; + }; + return { apiRoot, botToken, message: updates.result?.at(-1)?.message }; +} + +async function postTelegramMessage(params: { + apiRoot: string; + body: Record; + botToken: string; +}) { + const { response, release } = await fetchWithSsrFGuard({ + url: `${params.apiRoot}/bot${params.botToken}/sendMessage`, + init: { + body: JSON.stringify(params.body), + headers: { "content-type": "application/json" }, + method: "POST", + }, + policy: { allowPrivateNetwork: true }, + auditContext: "qa-lab-crabline-telegram-provider-correlation-test", + }); + await release(); + expect(response.ok).toBe(true); +} + +async function expectTelegramProviderCorrelation(params: { + conversation: { id: string; kind: "channel" | "direct" }; + expectedInbound?: { conversation: { id: string; kind: "direct" }; threadId: string }; + inboundText: string; + outboundText: string; + senderId: string; + threadId?: string; +}) { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection, + state: createQaBusState(), + }); + try { + const inbound = await transport.sendInbound({ + conversation: params.conversation, + senderId: params.senderId, + text: params.inboundText, + ...(params.threadId ? { threadId: params.threadId } : {}), + }); + if (params.expectedInbound) { + expect(inbound).toMatchObject(params.expectedInbound); + } + + const { apiRoot, botToken, message } = await readTelegramInbound(transport); + expect(message?.chat?.id).toEqual(expect.any(Number)); + if (params.threadId) { + expect(message?.message_thread_id).toEqual(expect.any(Number)); + } + await postTelegramMessage({ + apiRoot, + botToken, + body: { + chat_id: message?.chat?.id, + ...(message?.message_thread_id ? { message_thread_id: message.message_thread_id } : {}), + text: params.outboundText, + }, + }); + + await expect( + transport.waitForOutbound({ + conversation: params.conversation, + textIncludes: params.outboundText, + ...(params.threadId ? { threadId: params.threadId } : {}), + timeoutMs: 1_000, + }), + ).resolves.toMatchObject({ + conversation: params.conversation, + text: params.outboundText, + ...(params.threadId ? { threadId: params.threadId } : {}), + }); + } finally { + await transport.cleanupAfterGatewayStop(); + } + }); +} + +describe("Crabline Telegram provider identity", () => { + it("correlates private-topic sends with the canonical direct-thread target", async () => { + await expectTelegramProviderCorrelation({ + conversation: { id: "alice/team", kind: "direct" }, + expectedInbound: { + conversation: { id: "alice/team", kind: "direct" }, + threadId: "42", + }, + inboundText: "Private topic baseline marker.", + outboundText: "assistant via private topic", + senderId: "alice/team", + threadId: "42", + }); + }); + + it("preserves a channel target without outbound delivery registration", async () => { + await expectTelegramProviderCorrelation({ + conversation: { id: "telegram-announcements", kind: "channel" }, + inboundText: "Channel identity baseline.", + outboundText: "assistant via channel identity", + senderId: "alice", + }); + }); +}); diff --git a/extensions/qa-lab/src/crabline-transport-response.test.ts b/extensions/qa-lab/src/crabline-transport-response.test.ts new file mode 100644 index 000000000000..a65eeea4155c --- /dev/null +++ b/extensions/qa-lab/src/crabline-transport-response.test.ts @@ -0,0 +1,91 @@ +// Qa Lab tests cover bounded Crabline provider responses and failed-body cleanup. +import type { OpenClawCrablineChannelDriverSelection } from "@openclaw/crabline"; +import { withTempDir } from "openclaw/plugin-sdk/test-env"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { createQaBusState } from "./bus-state.js"; +import { createQaCrablineTransportAdapter } from "./crabline-transport.js"; + +afterEach(() => { + vi.unstubAllGlobals(); +}); + +function createSelection() { + return { + capabilityMatrixPath: "crabline-channel-driver-capabilities.json", + channel: "telegram", + channelDriver: "crabline", + providerReadinessArtifactPath: "crabline-provider-readiness.json", + } as const satisfies OpenClawCrablineChannelDriverSelection; +} + +describe("crabline transport responses", () => { + it("rejects oversized successful inbound responses before parsing provider metadata", async () => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection(), + state: createQaBusState(), + }); + vi.stubGlobal( + "fetch", + vi.fn( + async () => + new Response( + JSON.stringify({ + update: { message: { message_id: 42, padding: "x".repeat(1024 * 1024) } }, + }), + ), + ), + ); + + try { + await expect( + transport.sendInbound({ + conversation: { id: "-1001234567890", kind: "group" }, + senderId: "100001", + senderName: "Alice", + text: "Oversized response marker.", + }), + ).rejects.toThrow("JSON response exceeds 1048576 bytes"); + } finally { + await transport.cleanupAfterGatewayStop(); + } + }); + }); + + it("cancels a failed inbound response before surfacing the provider error", async () => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection(), + state: createQaBusState(), + }); + const cancel = vi.fn(() => { + throw new Error("cancel failed"); + }); + vi.stubGlobal( + "fetch", + vi.fn( + async () => + new Response(new ReadableStream({ cancel }), { + status: 503, + }), + ), + ); + + try { + await expect( + transport.sendInbound({ + conversation: { id: "-1001234567890", kind: "group" }, + senderId: "100001", + senderName: "Alice", + text: "Telegram failure marker.", + }), + ).rejects.toThrow("Crabline telegram inbound injection failed with HTTP 503"); + expect(cancel).toHaveBeenCalledOnce(); + } finally { + await transport.cleanupAfterGatewayStop(); + } + }); + }); +}); diff --git a/extensions/qa-lab/src/crabline-transport-thread-routing.test.ts b/extensions/qa-lab/src/crabline-transport-thread-routing.test.ts new file mode 100644 index 000000000000..e5b54d51cc12 --- /dev/null +++ b/extensions/qa-lab/src/crabline-transport-thread-routing.test.ts @@ -0,0 +1,292 @@ +// Crabline runner tests preserve provider-native thread ownership at the Gateway boundary. +import type { OpenClawCrablineChannelDriverSelection } from "@openclaw/crabline"; +import { mattermostPlugin } from "@openclaw/mattermost/channel-plugin-api.js"; +import { setMattermostRuntime } from "@openclaw/mattermost/runtime-api.js"; +import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime"; +import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; +import { withTempDir } from "openclaw/plugin-sdk/test-env"; +import { describe, expect, it, vi } from "vitest"; +import { createQaBusState } from "./bus-state.js"; +import { createQaCrablineTransportAdapter } from "./crabline-transport.js"; +import { startAgentRun } from "./suite-runtime-agent-process.js"; + +function createSelection(channel: OpenClawCrablineChannelDriverSelection["channel"]) { + return { + capabilityMatrixPath: "crabline-channel-driver-capabilities.json", + channel, + channelDriver: "crabline", + providerReadinessArtifactPath: "crabline-provider-readiness.json", + } as const; +} + +function requireString(value: unknown, label: string): string { + if (typeof value !== "string" || !value) { + throw new Error(`${label} is required`); + } + return value; +} + +async function postJson(params: { + url: string; + body: unknown; + headers?: Record; + method?: string; + auditContext: string; +}): Promise { + const { response, release } = await fetchWithSsrFGuard({ + url: params.url, + init: { + body: JSON.stringify(params.body), + headers: { "content-type": "application/json", ...params.headers }, + method: params.method ?? "POST", + }, + policy: { allowPrivateNetwork: true }, + auditContext: params.auditContext, + }); + try { + expect(response.ok).toBe(true); + return (await response.json()) as T; + } finally { + await release(); + } +} + +describe("Crabline provider thread routing", () => { + it.each([ + { channel: "matrix", threadId: "$native-thread:matrix.test" }, + { channel: "mattermost", threadId: "threadroot0000000000000000" }, + ] as const)( + "forwards $channel threads at the Gateway boundary", + async ({ channel, threadId }) => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection(channel), + state: createQaBusState(), + }); + const gatewayCall = vi.fn(async () => ({ runId: `run-${channel}` })); + + try { + await expect( + startAgentRun({ gateway: { call: gatewayCall }, transport } as never, { + sessionKey: `agent:qa:${channel}`, + message: "thread routing proof", + to: "group:qa-channel", + threadId, + }), + ).resolves.toEqual({ runId: `run-${channel}` }); + expect(gatewayCall).toHaveBeenCalledWith( + "agent", + expect.objectContaining({ channel, threadId }), + { timeoutMs: 30_000 }, + ); + } finally { + await transport.cleanupAfterGatewayStop?.(); + } + }); + }, + ); + + it("keeps Matrix root and thread correlation distinct", async () => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection("matrix"), + state: createQaBusState(), + }); + const conversationId = "matrix:room:!qa:matrix.test"; + const threadId = "$native-thread:matrix.test"; + const secondThreadId = "$second-thread:matrix.test"; + try { + await transport.state.addInboundMessage({ + conversation: { id: conversationId, kind: "group" }, + senderId: "driver", + text: "provision Matrix thread room", + }); + await transport.state.reset(); + const delivery = transport.buildAgentDelivery({ + target: conversationId, + threadId, + }); + transport.buildAgentDelivery({ target: conversationId, threadId: secondThreadId }); + const roomId = delivery.to.replace(/^room:/u, ""); + const env = transport.createRuntimeEnvPatch?.() ?? {}; + const matrixBaseUrl = requireString(env.MATRIX_BASE_URL, "Matrix base URL"); + const accessToken = requireString(env.MATRIX_ACCESS_TOKEN, "Matrix access token"); + const send = async (transactionId: string, body: Record) => + await postJson({ + url: `${matrixBaseUrl}/_matrix/client/v3/rooms/${encodeURIComponent(roomId)}/send/m.room.message/${transactionId}`, + body, + headers: { authorization: `Bearer ${accessToken}` }, + method: "PUT", + auditContext: "qa-lab-crabline-matrix-thread-correlation-test", + }); + + await send("qa-root-send", { body: "matrix root reply", msgtype: "m.text" }); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "direct" }, + threadId, + textIncludes: "matrix root reply", + timeoutMs: 50, + }), + ).rejects.toThrow(); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "direct" }, + threadId: secondThreadId, + textIncludes: "matrix root reply", + timeoutMs: 50, + }), + ).rejects.toThrow(); + await send("qa-thread-send", { + body: "matrix threaded reply", + msgtype: "m.text", + "m.relates_to": { rel_type: "m.thread", event_id: threadId }, + }); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "direct" }, + threadId, + textIncludes: "matrix threaded reply", + timeoutMs: 1_000, + }), + ).resolves.toMatchObject({ threadId, text: "matrix threaded reply" }); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "direct" }, + threadId: secondThreadId, + textIncludes: "matrix threaded reply", + timeoutMs: 50, + }), + ).rejects.toThrow(); + await send("qa-second-thread-send", { + body: "matrix second threaded reply", + msgtype: "m.text", + "m.relates_to": { rel_type: "m.thread", event_id: secondThreadId }, + }); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "direct" }, + threadId: secondThreadId, + textIncludes: "matrix second threaded reply", + timeoutMs: 1_000, + }), + ).resolves.toMatchObject({ + threadId: secondThreadId, + text: "matrix second threaded reply", + }); + } finally { + await transport.cleanupAfterGatewayStop?.(); + } + }); + }); + + it("keeps Mattermost root and thread correlation distinct", async () => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection("mattermost"), + state: createQaBusState(), + }); + const conversationId = "thread-channel"; + try { + await transport.state.addInboundMessage({ + conversation: { id: conversationId, kind: "group" }, + senderId: "alice", + text: "provision Mattermost thread channel", + }); + await transport.state.reset(); + const rootDelivery = transport.buildAgentDelivery({ target: `group:${conversationId}` }); + const channelId = rootDelivery.to.replace(/^channel:/u, ""); + const env = transport.createRuntimeEnvPatch?.() ?? {}; + const mattermostUrl = requireString(env.MATTERMOST_URL, "Mattermost URL"); + const botToken = requireString(env.MATTERMOST_BOT_TOKEN, "Mattermost bot token"); + const send = async (message: string, rootId?: string) => + await postJson<{ id: string }>({ + url: `${mattermostUrl}/api/v4/posts`, + body: { channel_id: channelId, message, ...(rootId ? { root_id: rootId } : {}) }, + headers: { authorization: `Bearer ${botToken}` }, + auditContext: "qa-lab-crabline-mattermost-thread-correlation-test", + }); + + const root = await send("mattermost seed root"); + transport.buildAgentDelivery({ target: `group:${conversationId}`, threadId: root.id }); + await send("mattermost unrelated root"); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "group" }, + threadId: root.id, + textIncludes: "mattermost unrelated root", + timeoutMs: 50, + }), + ).rejects.toThrow(); + await send("mattermost threaded reply", root.id); + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "group" }, + threadId: root.id, + textIncludes: "mattermost threaded reply", + timeoutMs: 1_000, + }), + ).resolves.toMatchObject({ threadId: root.id, text: "mattermost threaded reply" }); + } finally { + await transport.cleanupAfterGatewayStop?.(); + } + }); + }); + + it("delivers symbolic Mattermost threads through their native provider root", async () => { + await withTempDir("qa-crabline-transport-", async (outputDir) => { + const transport = await createQaCrablineTransportAdapter({ + outputDir, + selection: createSelection("mattermost"), + state: createQaBusState(), + }); + const conversationId = "symbolic-thread-channel"; + const threadId = "post-root"; + const gatewayCall = vi.fn(async (_method: string, _payload: Record) => ({ + runId: "run-mattermost-symbolic-thread", + })); + try { + setMattermostRuntime(createPluginRuntimeMock()); + await transport.state.addInboundMessage({ + conversation: { id: conversationId, kind: "group" }, + senderId: "alice", + text: "mattermost symbolic thread seed", + threadId, + }); + await startAgentRun({ gateway: { call: gatewayCall }, transport } as never, { + sessionKey: "agent:qa:mattermost-symbolic-thread", + message: "mattermost symbolic threaded reply", + to: `group:${conversationId}`, + threadId, + }); + const gatewayPayload = gatewayCall.mock.calls[0]?.[1]; + expect(gatewayPayload).toMatchObject({ + channel: "mattermost", + threadId: expect.stringMatching(/^[a-z0-9]{26}$/u), + to: expect.stringMatching(/^channel:[a-z0-9]{26}$/u), + }); + + await mattermostPlugin.outbound?.sendText?.({ + cfg: transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" }), + to: String(gatewayPayload?.to), + threadId: String(gatewayPayload?.threadId), + text: "mattermost symbolic threaded reply", + }); + + await expect( + transport.waitForOutbound({ + conversation: { id: conversationId, kind: "group" }, + threadId, + textIncludes: "mattermost symbolic threaded reply", + timeoutMs: 1_000, + }), + ).resolves.toMatchObject({ threadId, text: "mattermost symbolic threaded reply" }); + } finally { + await transport.cleanupAfterGatewayStop?.(); + } + }); + }); +}); diff --git a/extensions/qa-lab/src/crabline-transport.test.ts b/extensions/qa-lab/src/crabline-transport.test.ts index f50c210ac965..0589f6d52e08 100644 --- a/extensions/qa-lab/src/crabline-transport.test.ts +++ b/extensions/qa-lab/src/crabline-transport.test.ts @@ -4,14 +4,10 @@ import path from "node:path"; import type { OpenClawCrablineChannelDriverSelection } from "@openclaw/crabline"; import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime"; import { withTempDir } from "openclaw/plugin-sdk/test-env"; -import { afterEach, describe, expect, it, vi } from "vitest"; +import { describe, expect, it } from "vitest"; import { createQaBusState } from "./bus-state.js"; import { createQaCrablineTransportAdapter } from "./crabline-transport.js"; -afterEach(() => { - vi.unstubAllGlobals(); -}); - function createSelection(channel: OpenClawCrablineChannelDriverSelection["channel"] = "telegram") { return { capabilityMatrixPath: "crabline-channel-driver-capabilities.json", @@ -29,76 +25,6 @@ function requireString(value: unknown, label: string): string { } describe("crabline transport", () => { - it("rejects oversized successful inbound responses before parsing provider metadata", async () => { - await withTempDir("qa-crabline-transport-", async (outputDir) => { - const transport = await createQaCrablineTransportAdapter({ - outputDir, - selection: createSelection(), - state: createQaBusState(), - }); - vi.stubGlobal( - "fetch", - vi.fn( - async () => - new Response( - JSON.stringify({ - update: { message: { message_id: 42, padding: "x".repeat(1024 * 1024) } }, - }), - ), - ), - ); - - try { - await expect( - transport.sendInbound({ - conversation: { id: "-1001234567890", kind: "group" }, - senderId: "100001", - senderName: "Alice", - text: "Oversized response marker.", - }), - ).rejects.toThrow("JSON response exceeds 1048576 bytes"); - } finally { - await transport.cleanupAfterGatewayStop?.(); - } - }); - }); - - it("cancels a failed inbound response before surfacing the provider error", async () => { - await withTempDir("qa-crabline-transport-", async (outputDir) => { - const transport = await createQaCrablineTransportAdapter({ - outputDir, - selection: createSelection(), - state: createQaBusState(), - }); - const cancel = vi.fn(() => { - throw new Error("cancel failed"); - }); - vi.stubGlobal( - "fetch", - vi.fn( - async () => - new Response(new ReadableStream({ cancel }), { - status: 503, - }), - ), - ); - - try { - await expect( - transport.sendInbound({ - conversation: { id: "-1001234567890", kind: "group" }, - senderId: "100001", - senderName: "Alice", - text: "Telegram failure marker.", - }), - ).rejects.toThrow("Crabline telegram inbound injection failed with HTTP 503"); - expect(cancel).toHaveBeenCalledOnce(); - } finally { - await transport.cleanupAfterGatewayStop?.(); - } - }); - }); - it("configures OpenClaw's Telegram plugin against a Crabline local provider server", async () => { await withTempDir("qa-crabline-transport-", async (outputDir) => { const transport = await createQaCrablineTransportAdapter({ @@ -131,6 +57,14 @@ describe("crabline transport", () => { to: expect.stringMatching(/^[1-9]\d+$/u), }); expect(delivery.replyTo).toBe(delivery.to); + expect(transport.buildAgentDelivery({ target: "dm:alice", threadId: "42" })).toMatchObject({ + threadId: "42", + }); + expect(transport.buildAgentDelivery({ target: "-1001234567890" })).toMatchObject({ + channel: "telegram", + replyTo: "-1001234567890", + to: "-1001234567890", + }); await expect( fs.access(path.join(outputDir, "crabline-provider-server.json")), @@ -165,15 +99,14 @@ describe("crabline transport", () => { }); try { - expect(transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" })).toMatchObject({ - channels: { - telegram: { - allowFrom: ["100001"], - groupAllowFrom: ["100001"], - groupPolicy: "allowlist", - }, - }, + const gatewayConfig = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" }); + const telegramConfig = gatewayConfig.channels?.telegram; + expect(telegramConfig).toMatchObject({ + allowFrom: [expect.stringMatching(/^[1-9]\d+$/u)], + groupAllowFrom: [expect.stringMatching(/^[1-9]\d+$/u)], + groupPolicy: "allowlist", }); + const allowedDriverId = Number(telegramConfig?.allowFrom?.[0]); await transport.state.addInboundMessage({ conversation: { id: "qa-routing-ordering", kind: "group" }, senderId: "observer", @@ -185,8 +118,7 @@ describe("crabline transport", () => { text: "driver", }); - const config = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" }); - const telegram = config.channels?.telegram as + const telegram = gatewayConfig.channels?.telegram as | { apiRoot?: string; botToken?: string } | undefined; const apiRoot = requireString(telegram?.apiRoot, "Telegram API root"); @@ -195,7 +127,9 @@ describe("crabline transport", () => { const payload = (await response.json()) as { result?: Array<{ message?: { from?: { id?: number }; text?: string } }>; }; - expect(payload.result?.map((update) => update.message?.from?.id)).toEqual([100002, 100001]); + const senderIds = payload.result?.map((update) => update.message?.from?.id); + expect(senderIds?.[0]).not.toBe(allowedDriverId); + expect(senderIds?.[1]).toBe(allowedDriverId); } finally { await transport.cleanupAfterGatewayStop?.(); } @@ -360,6 +294,11 @@ describe("crabline transport", () => { SLACK_BOT_TOKEN: "xoxb-crabline-slack-token", SLACK_SIGNING_SECRET: "crabline-slack-signing-secret", }); + expect(transport.buildAgentDelivery({ target: "C1234567890" })).toMatchObject({ + channel: "slack", + replyTo: "C1234567890", + to: "C1234567890", + }); } finally { await transport.cleanupAfterGatewayStop?.(); } @@ -662,6 +601,12 @@ describe("crabline transport", () => { replyTo: expect.stringMatching(/^channel:[a-z0-9]{26}$/u), to: expect.stringMatching(/^channel:[a-z0-9]{26}$/u), }); + expect( + transport.buildAgentDelivery({ target: "group:qa-channel", threadId: "post-root" }), + ).toMatchObject({ + channel: "mattermost", + threadId: expect.stringMatching(/^[a-z0-9]{26}$/u), + }); expect(mattermostGatewayConfig.channels?.mattermost?.streaming).toEqual({ mode: "off" }); await expect( @@ -683,7 +628,7 @@ describe("crabline transport", () => { }); }); - it("normalizes native Mattermost post creation into outbound state", async () => { + it("correlates Mattermost's authoritative inbound channel with symbolic QA state", async () => { await withTempDir("qa-crabline-transport-", async (outputDir) => { const transport = await createQaCrablineTransportAdapter({ outputDir, @@ -693,43 +638,57 @@ describe("crabline transport", () => { try { await transport.state.addInboundMessage({ - conversation: { id: "qa-channel", kind: "group" }, + conversation: { id: "alice", kind: "direct" }, senderId: "alice", senderName: "Alice", text: "Mattermost baseline marker check.", }); - await transport.state.reset(); - const delivery = transport.buildAgentDelivery({ target: "group:qa-channel" }); const env = transport.createRuntimeEnvPatch?.() ?? {}; const mattermostUrl = requireString(env.MATTERMOST_URL, "Mattermost URL"); const botToken = requireString(env.MATTERMOST_BOT_TOKEN, "Mattermost bot token"); - const { response, release } = await fetchWithSsrFGuard({ - url: `${mattermostUrl}/api/v4/posts`, - init: { - body: JSON.stringify({ - channel_id: delivery.to.replace(/^channel:/u, ""), - message: "assistant via fake mattermost", - }), - headers: { - authorization: `Bearer ${botToken}`, - "content-type": "application/json", + const mattermostRequest = async (apiPath: string, init?: RequestInit) => { + const headers = new Headers(init?.headers); + headers.set("authorization", `Bearer ${botToken}`); + headers.set("content-type", "application/json"); + const { response, release } = await fetchWithSsrFGuard({ + url: `${mattermostUrl}/api/v4${apiPath}`, + init: { + ...init, + headers, }, - method: "POST", - }, - policy: { allowPrivateNetwork: true }, - auditContext: "qa-lab-crabline-mattermost-transport-test", + policy: { allowPrivateNetwork: true }, + auditContext: "qa-lab-crabline-mattermost-transport-test", + }); + try { + expect(response.ok).toBe(true); + return (await response.json()) as T; + } finally { + await release(); + } + }; + const bot = await mattermostRequest<{ id: string }>("/users/me"); + const alice = await mattermostRequest<{ id: string }>("/users/username/alice"); + const directChannel = await mattermostRequest<{ id: string }>("/channels/direct", { + body: JSON.stringify([bot.id, alice.id]), + method: "POST", }); - await release(); - expect(response.ok).toBe(true); + const outboundPost = await mattermostRequest<{ channel_id: string }>("/posts", { + body: JSON.stringify({ + channel_id: directChannel.id, + message: "assistant via fake mattermost", + }), + method: "POST", + }); + expect(outboundPost.channel_id).toBe(directChannel.id); await expect( transport.waitForOutbound({ - conversation: { id: "qa-channel", kind: "group" }, + conversation: { id: "alice", kind: "direct" }, textIncludes: "assistant via fake mattermost", timeoutMs: 1_000, }), ).resolves.toMatchObject({ - conversation: { id: "qa-channel", kind: "group" }, + conversation: { id: "alice", kind: "direct" }, text: "assistant via fake mattermost", }); } finally { @@ -794,12 +753,18 @@ describe("crabline transport", () => { replyTo: "room:!qa:matrix.test", to: "room:!qa:matrix.test", }); + expect( + transport.buildAgentDelivery({ target: "group:main", threadId: "$event:matrix.test" }), + ).toMatchObject({ + channel: "matrix", + threadId: "$event:matrix.test", + }); expect(() => transport.buildAgentDelivery({ target: "group:" })).toThrow( - "Matrix QA conversation id must be non-empty", - ); - expect(() => transport.buildAgentDelivery({ target: "thread:/v1/main/%24event" })).toThrow( - "Matrix thread targets require OpenClaw QA thread forwarding", + "invalid qa-channel group target", ); + expect(() => + transport.buildAgentDelivery({ target: "thread:main/$event:matrix.test" }), + ).toThrow("Matrix thread targets require OpenClaw QA thread forwarding"); await expect( transport.state.addInboundMessage({ conversation: { id: " ", kind: "group" }, @@ -983,6 +948,18 @@ describe("crabline transport", () => { | undefined; expect(telegram?.apiRoot).toBeTruthy(); expect(telegram?.botToken).toBeTruthy(); + const updatesResponse = await fetch( + `${telegram?.apiRoot}/bot${telegram?.botToken}/getUpdates`, + ); + const updates = (await updatesResponse.json()) as { + result?: Array<{ message?: { chat?: { id?: number } } }>; + }; + const authoritativeChatId = updates.result?.at(-1)?.message?.chat?.id; + expect(authoritativeChatId).toEqual(expect.any(Number)); + const groupDelivery = transport.buildAgentDelivery({ + target: "channel:telegram-command-room", + }); + expect(groupDelivery.to).toBe(String(authoritativeChatId)); const { response, release } = await fetchWithSsrFGuard({ url: `${telegram?.apiRoot}/bot${telegram?.botToken}/sendMessage`, init: { diff --git a/extensions/qa-lab/src/crabline-transport.ts b/extensions/qa-lab/src/crabline-transport.ts index b1480716a3a1..7f1998aad03d 100644 --- a/extensions/qa-lab/src/crabline-transport.ts +++ b/extensions/qa-lab/src/crabline-transport.ts @@ -7,7 +7,6 @@ import { startOpenClawCrablineAdapter, type OpenClawCrablineChannelDriverSelection, type OpenClawCrablineInbound, - type StartedOpenClawCrablineAdapter, type StartedOpenClawCrablineCorrelatedAdapter, } from "@openclaw/crabline"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; @@ -19,13 +18,14 @@ import { } from "openclaw/plugin-sdk/string-coerce-runtime"; import { createQaBusState, type QaBusState } from "./bus-state.js"; import { + createCrablineProviderCorrelation, createCrablineProviderDelivery, createCrablineProviderInboundInput, resolveCrablineStateConversation, resolveDiscordQaId, - resolveTelegramQaSenderId, } from "./crabline-provider-targets.js"; import { readQaJsonResponse } from "./ignored-response-body.js"; +import { buildQaConversationTarget, parseQaTarget } from "./qa-bus-protocol.js"; import { QaStateBackedTransportAdapter, type QaTransportActionName, @@ -50,10 +50,19 @@ type QaCrablineTransportState = QaTransportState & { cleanup: () => Promise; getOutboundEvents: () => Promise; observeEvent: (event: unknown) => void; - rememberProviderTarget: (providerTargetKey: string, qaTarget: string) => void; + rememberProviderTarget: (providerTargetKey: string, target: QaCrablineTarget) => void; resetTransport: () => void; }; +type QaCrablineTarget = Pick; + +function qaTargetForInput(input: QaBusInboundMessageInput): QaCrablineTarget { + return { + conversation: { ...input.conversation }, + ...(input.threadId ? { threadId: input.threadId } : {}), + }; +} + function normalizeCrablineSignalGatewayConfig(config: OpenClawConfig): OpenClawConfig { const signal = config.channels?.signal as unknown; if (!isRecord(signal)) { @@ -85,23 +94,6 @@ function normalizeCrablineSignalGatewayConfig(config: OpenClawConfig): OpenClawC } as OpenClawConfig; } -function resolveLogicalQaTarget( - { conversation, threadId }: QaBusInboundMessageInput, - providerQaTarget: string, - providerPreservesConversationId: boolean, -) { - if (conversation.kind !== "channel" && providerPreservesConversationId) { - return providerQaTarget; - } - const prefix = conversation.kind === "direct" ? "dm" : conversation.kind; - return threadId ? `thread:${conversation.id}/${threadId}` : `${prefix}:${conversation.id}`; -} - -function formatLogicalQaConversationTarget({ conversation }: QaBusInboundMessageInput) { - const prefix = conversation.kind === "direct" ? "dm" : conversation.kind; - return `${prefix}:${conversation.id}`; -} - const TELEGRAM_LIFECYCLE_METHOD_RE = /\/(sendMessage|editMessageText|deleteMessage)$/u; function readTelegramLifecycleEvent(params: { @@ -180,7 +172,7 @@ function readTelegramLifecycleEvent(params: { } async function postCrablineInbound(params: { - adapter: StartedOpenClawCrablineAdapter; + adapter: StartedOpenClawCrablineCorrelatedAdapter; providerInbound: OpenClawCrablineInbound; }) { const { response, release } = await fetchWithSsrFGuard({ @@ -198,36 +190,33 @@ async function postCrablineInbound(params: { }); const label = `Crabline ${params.adapter.channel} inbound injection failed`; const result = await readQaJsonResponse(response, release, label); + let providerMessageId: string | undefined; if (params.adapter.channel === "matrix" && isRecord(result) && isRecord(result.event)) { - return readStringValue(result.event.event_id); - } - if (params.adapter.channel === "slack" && isRecord(result) && isRecord(result.message)) { - return readStringValue(result.message.ts); - } - if ( + providerMessageId = readStringValue(result.event.event_id); + } else if (params.adapter.channel === "slack" && isRecord(result) && isRecord(result.message)) { + providerMessageId = readStringValue(result.message.ts); + } else if ( params.adapter.channel === "telegram" && isRecord(result) && isRecord(result.update) && isRecord(result.update.message) ) { - return normalizeStringifiedOptionalString(result.update.message.message_id); + providerMessageId = normalizeStringifiedOptionalString(result.update.message.message_id); } - return undefined; + return { providerMessageId, response: result }; } function createCrablineState(params: { - adapter: StartedOpenClawCrablineAdapter; + adapter: StartedOpenClawCrablineCorrelatedAdapter; state: QaBusState; }): QaCrablineTransportState { const baseState = params.state; - const targetByProviderTarget = new Map(); - const logicalRouteByTarget = new Map(); + const targetByProviderTarget = new Map(); const telegramMessageByProviderId = new Map(); const pendingTelegramMessagesByChat = new Map(); const outboundEvents: QaTransportOutboundEvent[] = []; const resetTransport = () => { targetByProviderTarget.clear(); - logicalRouteByTarget.clear(); telegramMessageByProviderId.clear(); pendingTelegramMessagesByChat.clear(); outboundEvents.length = 0; @@ -268,45 +257,54 @@ function createCrablineState(params: { }, } : event; - const outbound = params.adapter.createOutboundFromRecorderEvent({ - event: normalizedEvent, - targetByProviderTarget, - }) as QaBusOutboundMessageInput | null; - if (outbound) { - const logicalRoute = logicalRouteByTarget.get(outbound.to); - baseState.addOutboundMessage( - logicalRoute + const observation = params.adapter.createOutboundObservation({ event: normalizedEvent }); + const target = observation?.providerTargetKeys + .map((key) => targetByProviderTarget.get(key)) + .find((candidate) => candidate !== undefined); + const outbound: QaBusOutboundMessageInput | null = observation + ? target + ? { + accountId: observation.accountId, + senderId: observation.senderId, + senderName: observation.senderName, + text: observation.text, + to: buildQaConversationTarget({ + chatType: target.conversation.kind, + conversationId: target.conversation.id, + }), + threadId: target.threadId, + } + : observation.fallbackTarget ? { - ...outbound, - to: logicalRoute.target, - ...(logicalRoute.threadId ? { threadId: logicalRoute.threadId } : {}), + accountId: observation.accountId, + senderId: observation.senderId, + senderName: observation.senderName, + text: observation.text, + to: observation.fallbackTarget, } - : outbound, - ); + : null + : null; + if (outbound) { + baseState.addOutboundMessage(outbound); } }, async addInboundMessage(input: QaBusInboundMessageInput) { const providerInbound = params.adapter.createInbound({ input: createCrablineProviderInboundInput(params.adapter, input), }); - // Provider targets carry typed thread identity. Synthetic channels and Matrix's - // provider-native room ids still need their scenario-owned logical target restored. - const logicalTarget = resolveLogicalQaTarget( - input, - providerInbound.qaTarget, - params.adapter.channel !== "matrix" && params.adapter.channel !== "discord", - ); - targetByProviderTarget.set(providerInbound.providerTargetKey, logicalTarget); - if (params.adapter.channel === "discord") { - logicalRouteByTarget.set(logicalTarget, { - target: formatLogicalQaConversationTarget(input), - ...(input.threadId ? { threadId: input.threadId } : {}), - }); - } - const providerMessageId = await postCrablineInbound({ + const ingress = await postCrablineInbound({ adapter: params.adapter, providerInbound, }); + // Register only the provider identity confirmed by successful ingress. Provider servers may + // realize a symbolic target as a different native conversation than the provisional input. + targetByProviderTarget.set( + params.adapter.resolveInboundProviderTargetKey({ + inbound: providerInbound, + response: ingress.response, + }), + qaTargetForInput(input), + ); return baseState.addInboundMessage( { ...input, @@ -315,17 +313,12 @@ function createCrablineState(params: { input, providerInbound, }), - ...(input.threadId && params.adapter.channel === "discord" - ? { threadId: input.threadId } - : providerInbound.threadId - ? { threadId: providerInbound.threadId } - : {}), }, - providerMessageId, + ingress.providerMessageId, ); }, - rememberProviderTarget(providerTargetKey, qaTarget) { - targetByProviderTarget.set(providerTargetKey, qaTarget); + rememberProviderTarget(providerTargetKey, target) { + targetByProviderTarget.set(providerTargetKey, target); }, addOutboundMessage: baseState.addOutboundMessage.bind(baseState), readMessage: baseState.readMessage.bind(baseState), @@ -473,7 +466,10 @@ class QaCrablineTransport extends QaStateBackedTransportAdapter { if (this.#selection.channel !== "telegram") { return config as QaTransportGatewayConfig; } - const senderAllowlist = this.#transportPolicy?.senderAllowlist?.map(resolveTelegramQaSenderId); + const senderAllowlist = this.#transportPolicy?.senderAllowlist?.map( + (senderId) => + this.#adapter.createAgentDelivery({ target: `dm:${senderId}` }).providerTargetKey, + ); if (!this.#transportPolicy?.requireGroupMention && !senderAllowlist) { return config as QaTransportGatewayConfig; } @@ -509,10 +505,55 @@ class QaCrablineTransport extends QaStateBackedTransportAdapter { channel: this.#adapter.channel, }); - buildAgentDelivery = ({ target }: { target: string }) => { - const { delivery, providerTargetKey } = createCrablineProviderDelivery(this.#adapter, target); - this.#state.rememberProviderTarget(providerTargetKey, target); - return delivery; + buildAgentDelivery = ({ target, threadId }: { target: string; threadId?: string }) => { + const parsed = parseQaTarget(target); + if (parsed.threadId && threadId && parsed.threadId !== threadId) { + throw new Error("Crabline delivery received conflicting thread targets"); + } + const logicalTarget = { + conversation: { id: parsed.conversationId, kind: parsed.chatType }, + threadId: threadId ?? parsed.threadId, + }; + // Provider-native targets must retain their own classification (for example, + // Telegram negative group ids and Slack C/G conversation ids). Matrix and + // Mattermost also require OpenClaw to forward threads separately at the + // Gateway request boundary instead of passing them into Crabline delivery setup. + const providerThreadId = + this.#selection.channel === "matrix" || this.#selection.channel === "mattermost" + ? undefined + : logicalTarget.threadId; + const { delivery, providerTargetKey } = createCrablineProviderDelivery( + this.#adapter, + target, + providerThreadId, + ); + let deliveryThreadId = logicalTarget.threadId; + if (providerThreadId === undefined && logicalTarget.threadId) { + const providerCorrelation = createCrablineProviderCorrelation(this.#adapter, logicalTarget); + this.#state.rememberProviderTarget(providerTargetKey, { + conversation: logicalTarget.conversation, + }); + this.#state.rememberProviderTarget(providerCorrelation.providerTargetKey, logicalTarget); + if (this.#selection.channel === "mattermost") { + if (!providerCorrelation.threadId) { + throw new Error("Crabline Mattermost correlation did not resolve a native thread root"); + } + deliveryThreadId = providerCorrelation.threadId; + } + } else { + this.#state.rememberProviderTarget(providerTargetKey, logicalTarget); + } + return { + ...delivery, + ...(deliveryThreadId + ? { + threadId: + this.#selection.channel === "discord" + ? resolveDiscordQaId(deliveryThreadId) + : deliveryThreadId, + } + : {}), + }; }; createRuntimeEnvPatch = () => diff --git a/extensions/qa-lab/src/qa-bus-protocol.test.ts b/extensions/qa-lab/src/qa-bus-protocol.test.ts index c3e0b9ca9937..253deb3e4ad3 100644 --- a/extensions/qa-lab/src/qa-bus-protocol.test.ts +++ b/extensions/qa-lab/src/qa-bus-protocol.test.ts @@ -3,9 +3,23 @@ import { sanitizeQaBusToolCalls as sanitizeCanonicalQaBusToolCalls, } from "openclaw/plugin-sdk/qa-channel-protocol"; import { describe, expect, it } from "vitest"; -import { parseQaTarget, sanitizeQaBusToolCalls } from "./qa-bus-protocol.js"; +import { + buildQaConversationTarget, + parseQaTarget, + sanitizeQaBusToolCalls, +} from "./qa-bus-protocol.js"; describe("QA Lab package bus protocol", () => { + it.each(["direct", "group", "channel"] as const)( + "builds the canonical base target for %s conversations", + (chatType) => { + const conversationId = "Case/Room"; + expect(buildQaConversationTarget({ chatType, conversationId })).toBe( + `${chatType === "direct" ? "dm" : chatType}:${conversationId}`, + ); + }, + ); + it.each([ "bare-id", "channel:CaseSensitive", diff --git a/extensions/qa-lab/src/qa-bus-protocol.ts b/extensions/qa-lab/src/qa-bus-protocol.ts index 70f4175ed237..313be68a3a4a 100644 --- a/extensions/qa-lab/src/qa-bus-protocol.ts +++ b/extensions/qa-lab/src/qa-bus-protocol.ts @@ -1,2 +1,16 @@ +import { + buildQaTarget, + parseQaTarget, + sanitizeQaBusToolCalls, +} from "openclaw/plugin-sdk/qa-channel-protocol"; +import type { QaBusConversation } from "./runtime-api.js"; + // This subpath is part of the published package used by package-acceptance mounts. -export { parseQaTarget, sanitizeQaBusToolCalls } from "openclaw/plugin-sdk/qa-channel-protocol"; +export { parseQaTarget, sanitizeQaBusToolCalls }; + +export function buildQaConversationTarget(params: { + chatType: QaBusConversation["kind"]; + conversationId: string; +}): string { + return buildQaTarget(params); +} diff --git a/extensions/qa-lab/src/qa-channel-transport.ts b/extensions/qa-lab/src/qa-channel-transport.ts index 2d6bb041a7cc..8138189a633f 100644 --- a/extensions/qa-lab/src/qa-channel-transport.ts +++ b/extensions/qa-lab/src/qa-channel-transport.ts @@ -114,11 +114,14 @@ class QaChannelTransport extends QaStateBackedTransportAdapter { accountId: QA_CHANNEL_ACCOUNT_ID, channel: QA_CHANNEL_ID, }); - buildAgentDelivery = ({ target }: { target: string }) => ({ - channel: QA_CHANNEL_ID, - replyChannel: QA_CHANNEL_ID, - replyTo: target, - }); + buildAgentDelivery = ({ target, threadId }: { target: string; threadId?: string }) => { + return { + channel: QA_CHANNEL_ID, + replyChannel: QA_CHANNEL_ID, + replyTo: target, + ...(threadId ? { threadId } : {}), + }; + }; async sendNativeCommand(input: QaTransportNativeCommandInput): Promise { const { command, ...message } = input; await this.sendInbound({ diff --git a/extensions/qa-lab/src/qa-transport.ts b/extensions/qa-lab/src/qa-transport.ts index 12d4c975903b..ddd5b30c6061 100644 --- a/extensions/qa-lab/src/qa-transport.ts +++ b/extensions/qa-lab/src/qa-transport.ts @@ -376,11 +376,12 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte timeoutMs?: number; pollIntervalMs?: number; }) => Promise; - abstract buildAgentDelivery: (params: { target: string }) => { + abstract buildAgentDelivery: (params: { target: string; threadId?: string }) => { channel: string; to?: string; replyChannel: string; replyTo: string; + threadId?: string; }; abstract handleAction: (params: { action: QaTransportActionName; diff --git a/extensions/qa-lab/src/self-check-scenario.ts b/extensions/qa-lab/src/self-check-scenario.ts index 6b6f389c0178..24a01f6a1d75 100644 --- a/extensions/qa-lab/src/self-check-scenario.ts +++ b/extensions/qa-lab/src/self-check-scenario.ts @@ -9,7 +9,7 @@ export function createQaSelfCheckScenario(options?: { waitTimeoutMs?: number; }): QaScenarioDefinition { const waitTimeoutMs = options?.waitTimeoutMs ?? 5_000; - let lifecycle: { target: string; message: QaBusMessage } | undefined; + let lifecycle: { target: string; threadId: string; message: QaBusMessage } | undefined; const waitForReply = (state: QaTransportState, inbound: QaBusMessage) => waitForOutboundMessage( state, @@ -48,12 +48,11 @@ export function createQaSelfCheckScenario(options?: { }); const threadPayload = extractQaToolPayload( threadResult as Parameters[0], - ) as { target?: string; thread?: { id?: string } } | undefined; - const threadId = threadPayload?.thread?.id; - if (!threadId || !threadPayload?.target) { + ) as { target?: string; threadId?: string; thread?: { id?: string } } | undefined; + const threadId = threadPayload?.threadId; + if (!threadId || threadId !== threadPayload?.thread?.id || !threadPayload.target) { throw new Error("thread-create did not return thread id and target"); } - const inbound = await state.addInboundMessage({ conversation: { id: "qa-room", kind: "channel", title: "QA Room" }, senderId: "alice", @@ -64,6 +63,7 @@ export function createQaSelfCheckScenario(options?: { }); lifecycle = { target: threadPayload.target, + threadId, message: await waitForReply(state, inbound), }; return threadId; @@ -78,10 +78,11 @@ export function createQaSelfCheckScenario(options?: { if (!lifecycle) { throw new Error("threaded outbound message and target not found"); } - const { target, message: outboundMessage } = lifecycle; + const { target, threadId, message: outboundMessage } = lifecycle; await performAction("react", { to: target, + threadId, messageId: outboundMessage.id, emoji: "white_check_mark", }); @@ -95,6 +96,7 @@ export function createQaSelfCheckScenario(options?: { await performAction("edit", { to: target, + threadId, messageId: outboundMessage.id, text: "qa-echo: inside thread (edited)", }); @@ -108,6 +110,7 @@ export function createQaSelfCheckScenario(options?: { await performAction("delete", { to: target, + threadId, messageId: outboundMessage.id, }); const deleted = await state.readMessage({ messageId: outboundMessage.id }); diff --git a/extensions/qa-lab/src/self-check.test.ts b/extensions/qa-lab/src/self-check.test.ts index 71a59df161fb..644c2adae3cc 100644 --- a/extensions/qa-lab/src/self-check.test.ts +++ b/extensions/qa-lab/src/self-check.test.ts @@ -105,13 +105,14 @@ describe("createQaSelfCheckScenario", () => { }); return { details: { - target: `thread:${thread.conversationId}/${thread.id}`, + target: `channel:${thread.conversationId}`, + threadId: thread.id, thread, }, }; } const message = state.readMessage({ messageId: String(args.messageId) }); - if (args.to !== `thread:${message.conversation.id}/${String(message.threadId)}`) { + if (args.to !== `channel:${message.conversation.id}` || args.threadId !== message.threadId) { throw new Error("qa-channel message is not in the selected conversation"); } targets.push(args.to); @@ -157,11 +158,7 @@ describe("createQaSelfCheckScenario", () => { const thread = state.getSnapshot().threads[0]; expect(thread).toBeDefined(); - expect(targets).toEqual([ - `thread:qa-room/${String(thread?.id)}`, - `thread:qa-room/${String(thread?.id)}`, - `thread:qa-room/${String(thread?.id)}`, - ]); + expect(targets).toEqual(["channel:qa-room", "channel:qa-room", "channel:qa-room"]); const deletedMessage = state.getSnapshot().messages.find((message) => message.deleted); if (!deletedMessage) { throw new Error("self-check did not preserve its deleted message tombstone"); diff --git a/extensions/qa-lab/src/suite-runtime-agent-process.test.ts b/extensions/qa-lab/src/suite-runtime-agent-process.test.ts index 9bbeca1cf9dc..6884928e6b49 100644 --- a/extensions/qa-lab/src/suite-runtime-agent-process.test.ts +++ b/extensions/qa-lab/src/suite-runtime-agent-process.test.ts @@ -497,6 +497,7 @@ describe("qa suite runtime agent process helpers", () => { to: "transport-target", replyChannel: "reply-channel", replyTo: "reply-target", + threadId: "adapter-thread", })), }, } as never; @@ -516,6 +517,7 @@ describe("qa suite runtime agent process helpers", () => { replyChannel?: string; replyTo?: string; sessionKey?: string; + threadId?: string; to?: string; } | undefined; @@ -525,15 +527,17 @@ describe("qa suite runtime agent process helpers", () => { expect(agentPayload?.to).toBe("transport-target"); expect(agentPayload?.replyChannel).toBe("reply-channel"); expect(agentPayload?.replyTo).toBe("reply-target"); + expect(agentPayload?.threadId).toBe("adapter-thread"); expect(gatewayArgs?.[2]).toBeTypeOf("object"); }); - it("starts an interactive run without CLI task tracking", async () => { + it("preserves thread routing for an interactive run without CLI task tracking", async () => { const gatewayCall = vi.fn(async () => ({ runId: "run-chat", status: "started" })); const buildAgentDelivery = vi.fn(() => ({ channel: "qa-channel", replyChannel: "qa-channel", replyTo: "dm:qa-operator", + threadId: "provider-topic-42", })); const env = { gateway: { call: gatewayCall }, @@ -546,6 +550,7 @@ describe("qa suite runtime agent process helpers", () => { startAgentRun(env, { sessionKey: "agent:qa:main", message: "hello", + threadId: "topic-42", taskTracking: false, }), ).resolves.toEqual({ runId: "run-chat", status: "started" }); @@ -558,10 +563,14 @@ describe("qa suite runtime agent process helpers", () => { deliver: true, originatingChannel: "qa-channel", originatingTo: "dm:qa-operator", + originatingThreadId: "provider-topic-42", }, { timeoutMs: 30_000 }, ); - expect(buildAgentDelivery).toHaveBeenCalledWith({ target: "dm:qa-operator" }); + expect(buildAgentDelivery).toHaveBeenCalledWith({ + target: "dm:qa-operator", + threadId: "topic-42", + }); }); it("finds managed dreaming cron jobs across legacy and current payload contracts", () => { diff --git a/extensions/qa-lab/src/suite-runtime-agent-process.ts b/extensions/qa-lab/src/suite-runtime-agent-process.ts index 76848f511db8..e8c0409a6754 100644 --- a/extensions/qa-lab/src/suite-runtime-agent-process.ts +++ b/extensions/qa-lab/src/suite-runtime-agent-process.ts @@ -77,7 +77,10 @@ async function startAgentRun( ) { if (params.taskTracking === false) { const target = params.to ?? "dm:qa-operator"; - const delivery = env.transport.buildAgentDelivery({ target }); + const delivery = env.transport.buildAgentDelivery({ + target, + ...(params.threadId ? { threadId: params.threadId } : {}), + }); const started = (await env.gateway.call( "chat.send", { @@ -87,6 +90,8 @@ async function startAgentRun( deliver: true, originatingChannel: delivery.replyChannel, originatingTo: delivery.replyTo, + // chat.send routes threads separately; omitting this replies at the conversation root. + originatingThreadId: delivery.threadId ?? params.threadId, }, { timeoutMs: params.timeoutMs ?? 30_000, @@ -98,7 +103,10 @@ async function startAgentRun( return started; } const target = params.to ?? "dm:qa-operator"; - const delivery = env.transport.buildAgentDelivery({ target }); + const delivery = env.transport.buildAgentDelivery({ + target, + ...(params.threadId ? { threadId: params.threadId } : {}), + }); const started = (await env.gateway.call( "agent", { @@ -111,7 +119,9 @@ async function startAgentRun( to: delivery.to ?? target, replyChannel: delivery.replyChannel, replyTo: delivery.replyTo, - ...(params.threadId ? { threadId: params.threadId } : {}), + ...((delivery.threadId ?? params.threadId) + ? { threadId: delivery.threadId ?? params.threadId } + : {}), ...(params.provider ? { provider: params.provider } : {}), ...(params.model ? { model: params.model } : {}), ...(params.attachments ? { attachments: params.attachments } : {}), diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 7203d602cf49..ae5499ff3ad4 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1914,6 +1914,9 @@ importers: '@openclaw/matrix': specifier: workspace:* version: link:../matrix + '@openclaw/mattermost': + specifier: workspace:* + version: link:../mattermost '@openclaw/plugin-sdk': specifier: workspace:* version: link:../../packages/plugin-sdk diff --git a/src/plugin-sdk/qa-runner-runtime.test.ts b/src/plugin-sdk/qa-runner-runtime.test.ts index ac9d9eb3a2b9..b7205f75c0c8 100644 --- a/src/plugin-sdk/qa-runner-runtime.test.ts +++ b/src/plugin-sdk/qa-runner-runtime.test.ts @@ -4,6 +4,7 @@ import path from "node:path"; import type { Command } from "commander"; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import type { QaRunnerCliRegistration } from "./qa-runner-runtime.js"; import { cleanupTempDirs, expectPrivateQaLabRuntimeSurfaceLoad, @@ -80,6 +81,21 @@ describe("plugin-sdk qa-runner-runtime", () => { } }); + it("exposes structured thread identity to transport delivery adapters", () => { + type Adapter = Awaited< + ReturnType["create"]> + >; + const buildAgentDelivery: Adapter["buildAgentDelivery"] = ({ target, threadId }) => ({ + channel: "linked", + replyChannel: "linked", + replyTo: threadId ? `${target}:thread:${threadId}` : target, + }); + + expect(buildAgentDelivery({ target: "channel:room", threadId: "topic-1" }).replyTo).toBe( + "channel:room:thread:topic-1", + ); + }); + it("stays cold until runner discovery is requested", async () => { vi.resetModules(); await import("./qa-runner-runtime.js"); diff --git a/src/plugin-sdk/qa-runner-runtime.ts b/src/plugin-sdk/qa-runner-runtime.ts index ee19080ebb8a..9e8ec51a8fb3 100644 --- a/src/plugin-sdk/qa-runner-runtime.ts +++ b/src/plugin-sdk/qa-runner-runtime.ts @@ -163,11 +163,12 @@ type QaRunnerTransportAdapterDefinition = { timeoutMs?: number; pollIntervalMs?: number; }) => Promise; - buildAgentDelivery: (params: { target: string }) => { + buildAgentDelivery: (params: { target: string; threadId?: string }) => { channel: string; to?: string; replyChannel: string; replyTo: string; + threadId?: string; }; createRuntimeEnvPatch?: () => NodeJS.ProcessEnv; prepareFlow?: ( diff --git a/test/qa-channel-message-tool-delivery.test.ts b/test/qa-channel-message-tool-delivery.test.ts index 8a8bed3876d2..67af4859f1ab 100644 --- a/test/qa-channel-message-tool-delivery.test.ts +++ b/test/qa-channel-message-tool-delivery.test.ts @@ -302,7 +302,8 @@ describe("QA message-tool current conversation delivery", () => { { target: root }, { to: root }, { channelId: root }, - ...(threadId ? [{ target: `${thread}/${threadId}` }, { threadId }] : []), + // Read-capable actions use the host-owned root plus explicit thread identity. + ...(threadId ? [{ threadId }] : []), ]) { const result = await tool.execute(`own-${args.action}`, { ...args,