refactor(qa): keep Crabline thread identity structured (#124467)

* refactor(qa): preserve structured Crabline thread identity

* test(qa): use canonical Discord group thread target

* test(qa): preserve native Discord thread delivery

* fix(qa): preserve native Mattermost thread delivery
This commit is contained in:
Dallin Romney 2026-09-17 03:56:09 -07:00 • committed by GitHub
parent ce9a07f765
commit 35bfbb5377
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
31 changed files with 1076 additions and 293 deletions

View file

@ -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

View file

@ -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",
},

View file

@ -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,

View file

@ -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");
});
});

View file

@ -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",

View file

@ -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" });
});
});

View file

@ -190,6 +190,10 @@ export const qaChannelPlugin: ChannelPlugin<ResolvedQaChannelAccount> = 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<ResolvedQaChannelAccount> = 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 }) =>

View file

@ -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);

View file

@ -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) => {

View file

@ -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,

View file

@ -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:*",

View file

@ -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<string, unknown>) => ({
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?.();
}
});
});
});

View file

@ -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",

View file

@ -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<StartedOpenClawCrablineCorrelatedAdapter, "channel" | "createAgentDelivery">,
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<QaBusInboundMessageInput, "conversation" | "threadId">,
) {
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 } : {}),
}),
});
}

View file

@ -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<ReturnType<typeof createQaCrablineTransportAdapter>>,
) {
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<string, unknown>;
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",
});
});
});

View file

@ -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<Uint8Array>({ 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();
}
});
});
});

View file

@ -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<T>(params: {
url: string;
body: unknown;
headers?: Record<string, string>;
method?: string;
auditContext: string;
}): Promise<T> {
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<string, unknown>) =>
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<string, unknown>) => ({
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?.();
}
});
});
});

View file

@ -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<Uint8Array>({ 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 <T>(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: {

View file

@ -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<void>;
getOutboundEvents: () => Promise<readonly QaTransportOutboundEvent[]>;
observeEvent: (event: unknown) => void;
rememberProviderTarget: (providerTargetKey: string, qaTarget: string) => void;
rememberProviderTarget: (providerTargetKey: string, target: QaCrablineTarget) => void;
resetTransport: () => void;
};
type QaCrablineTarget = Pick<QaBusInboundMessageInput, "conversation" | "threadId">;
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<unknown>(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<string, string>();
const logicalRouteByTarget = new Map<string, { target: string; threadId?: string }>();
const targetByProviderTarget = new Map<string, QaCrablineTarget>();
const telegramMessageByProviderId = new Map<string, QaBusMessage>();
const pendingTelegramMessagesByChat = new Map<string, QaBusMessage[]>();
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 = () =>

View file

@ -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",

View file

@ -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);
}

View file

@ -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<void> {
const { command, ...message } = input;
await this.sendInbound({

View file

@ -376,11 +376,12 @@ export abstract class QaStateBackedTransportAdapter implements QaTransportAdapte
timeoutMs?: number;
pollIntervalMs?: number;
}) => Promise<void>;
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;

View file

@ -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<typeof extractQaToolPayload>[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 });

View file

@ -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");

View file

@ -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", () => {

View file

@ -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 } : {}),

3
pnpm-lock.yaml generated
View file

@ -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

View file

@ -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<NonNullable<QaRunnerCliRegistration["adapterFactory"]>["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");

View file

@ -163,11 +163,12 @@ type QaRunnerTransportAdapterDefinition = {
timeoutMs?: number;
pollIntervalMs?: number;
}) => Promise<void>;
buildAgentDelivery: (params: { target: string }) => {
buildAgentDelivery: (params: { target: string; threadId?: string }) => {
channel: string;
to?: string;
replyChannel: string;
replyTo: string;
threadId?: string;
};
createRuntimeEnvPatch?: () => NodeJS.ProcessEnv;
prepareFlow?: (

View file

@ -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,