mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix(qa): judge the completed Discord final in the reply-shape scenario (#163046)
* test(qa): distinguish Discord readiness from runtime evidence The readiness snapshot has been startup-only sincea2c698f3d3, while the Discord roundtrip test still expected later runtime traffic there. Assert that the startup recorder is distinct, contains accepted probe activity, and excludes subsequent roundtrip delivery; retain the independent full runtime REST, WebSocket, thread, quote, and cleanup assertions. The old expectation failed twice on Blacksmith Testbox despite the actual roundtrip succeeding. The repaired opt-in Discord E2E file passes both tests in 46.21s (roundtrip case 9.52s). Full typechecks and targeted changed-file checks pass. Host restrictions require remote formatting and tests; local staged content guards were run separately. * fix(qa): judge completed native Discord replies Preserve native inbound message IDs and message_reference metadata, apply the existing topLevelReplies transport policy, and select the retained Discord final only after the channel run queue drains. Read native messages after completion so edits and deletions cannot leave a matching preview as false proof. Validate text, destination, and reply shape after final selection. QA-channel reuses its existing processing-acknowledged reply owner; unsupported completion adapters fail explicitly. Repairs metadata omission from35bfbb5377and scenario/policy omission from253b784468. Before the fix, the real QA CLI reported pass while the native recorder showed a matching preview followed by a quoted final. The new regression failed on the original code for ignored policy and premature preview completion. Afterward, held previews cannot finish the scenario: quoted final, wrong edited final, and deleted-only delivery fail; a quoted preview followed by a top-level final passes. Blacksmith Testbox proof on rebasedd53b32bb96: private QA build; 177 tests across 20 related files; real Discord plugin plus Crabline and mock provider controls in both directions; thread roundtrip and cleanup. Final opt-in E2E file: 2 passed, 46.21s, including 20.22s for the dual final-shape control. Measured pnpm test --maxWorkers=1 wall: crabline-discord-transport.test.ts 19.51s, crabline-transport.test.ts 2.90s. No live credentials or external channel were used. Full check-changed passed all types and guards, then found a test-only filter/at lint violation. Replaced it with findLast; final targeted check-changed, formatting, types, lint, dead-export scan, and runtime import-cycle gate passed. Remote-only execution respected host limits; local staged content guards passed separately. No changelog, config surface, SQLite schema, Gateway protocol, or plugin SDK changes. * fix(qa): pin Discord final replies to delivery receipts Record the native Discord message ID at successful send completion, then read that exact retained message after the channel run drains. A deleted or wrongly edited final cannot fall back to an earlier matching preview. Keep the QA relay's REST and WebSocket origins aligned. Regression: preview -> final -> delete-final passed before the fix and now fails with the missing final delivery ID. Extend the existing native transport fixture with deleted-final and edited-final cases. Blacksmith Testbox proof on the exact candidate source: 188 tests across 21 QA suites; both Discord Crabline E2E tests (32.153s test time); check-changed including 6 doctor guard tests; private-QA build. No local builds or tests, and no real provider/channel credentials. * fix(qa): preserve Discord relay message content Replace recursive string rewriting with Discord payload-specific endpoint handling: Get Gateway/Get Gateway Bot url, Ready resume_gateway_url, and message attachment url/proxy_url including referenced messages. Leave content, embeds, attachment descriptions, and unrelated payload fields unchanged. The native HTTP/WebSocket regression fails on the reviewed head by changing content and embed text to the relay origin. It passes after the fix for both API and Gateway origin prefixes, while endpoint rewriting and attachment downloads still work. Blacksmith Testbox exact-source proof: 188 tests across 22 QA suites (48.96s wall), both Discord Crabline E2E tests (35.265s test time), private-QA build, and check-changed including 6 doctor guard tests. Dependency installation preceded remote execution; no local build/test or real provider/channel credentials. Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
parent
b3cb94506b
commit
268fdc552c
10 changed files with 790 additions and 32 deletions
118
extensions/qa-lab/src/crabline-discord-replies.test.ts
Normal file
118
extensions/qa-lab/src/crabline-discord-replies.test.ts
Normal file
|
|
@ -0,0 +1,118 @@
|
|||
import path from "node:path";
|
||||
import { startOpenClawCrablineAdapter } from "@openclaw/crabline";
|
||||
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
|
||||
import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
|
||||
import { withTempDir } from "openclaw/plugin-sdk/test-env";
|
||||
import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress";
|
||||
import { expect, it } from "vitest";
|
||||
import WebSocket from "ws";
|
||||
import { createQaBusState } from "./bus-state.js";
|
||||
import { startCrablineDiscordReplies } from "./crabline-discord-replies.js";
|
||||
|
||||
it("relays Discord message bytes unchanged while rewriting transport endpoints", async () => {
|
||||
await withTempDir("qa-discord-relay-payload-", async (outputDir) => {
|
||||
const adapter = await startOpenClawCrablineAdapter({
|
||||
channel: "discord",
|
||||
openclawConfig: {},
|
||||
recorderPath: path.join(outputDir, "discord.jsonl"),
|
||||
});
|
||||
const relay = await startCrablineDiscordReplies({
|
||||
adapter,
|
||||
state: createQaBusState(),
|
||||
targets: new Map(),
|
||||
});
|
||||
let socket: WebSocket | undefined;
|
||||
try {
|
||||
if (!relay || adapter.manifest.provider !== "discord") {
|
||||
throw new Error("Discord relay did not start");
|
||||
}
|
||||
const manifest = adapter.manifest;
|
||||
const headers = { authorization: `Bot ${manifest.botToken}` };
|
||||
const relayOrigin = new URL(relay.apiBaseUrl).origin;
|
||||
const gatewayUrl = `${relayOrigin.replace(/^http/u, "ws")}/gateway`;
|
||||
for (const route of ["gateway", "gateway/bot"]) {
|
||||
const response = await fetch(`${relay.apiBaseUrl}/${route}`, { headers });
|
||||
expect(response.ok).toBe(true);
|
||||
await expect(response.json()).resolves.toMatchObject({ url: gatewayUrl });
|
||||
}
|
||||
socket = new WebSocket(gatewayUrl);
|
||||
const ready = createDeferred<unknown>();
|
||||
let delivered = createDeferred<unknown>();
|
||||
socket.on("error", (error) => {
|
||||
ready.reject(error);
|
||||
delivered.reject(error);
|
||||
});
|
||||
socket.on("message", (data) => {
|
||||
const event: unknown = JSON.parse(rawDataToString(data));
|
||||
if (isRecord(event)) {
|
||||
if (event.t === "READY") {
|
||||
ready.resolve(event.d);
|
||||
} else if (event.t === "MESSAGE_CREATE" || event.t === "MESSAGE_UPDATE") {
|
||||
delivered.resolve(event.d);
|
||||
}
|
||||
}
|
||||
});
|
||||
socket.on("open", () => {
|
||||
socket?.send(JSON.stringify({ op: 2, d: { token: manifest.botToken, intents: 0 } }));
|
||||
});
|
||||
await expect(ready.promise).resolves.toMatchObject({ resume_gateway_url: gatewayUrl });
|
||||
|
||||
const messagesUrl = `${relay.apiBaseUrl}/channels/${manifest.fixture.channelId}/messages`;
|
||||
for (const origin of [
|
||||
new URL(manifest.endpoints.apiRoot).origin,
|
||||
new URL(manifest.endpoints.gatewayUrl).origin,
|
||||
]) {
|
||||
const text = `${origin}/user-visible?keep=%2F#fragment\r\n\tUnicode 🦞`;
|
||||
const payload = {
|
||||
content: text,
|
||||
embeds: [{ description: text, title: text, url: `${origin}/semantic-link` }],
|
||||
};
|
||||
const form = new FormData();
|
||||
form.set(
|
||||
"payload_json",
|
||||
JSON.stringify({
|
||||
...payload,
|
||||
attachments: [{ id: "0", description: text }],
|
||||
}),
|
||||
);
|
||||
form.set("files[0]", new Blob(["attachment bytes"]), "proof.txt");
|
||||
delivered = createDeferred<unknown>();
|
||||
const response = await fetch(messagesUrl, { method: "POST", headers, body: form });
|
||||
expect(response.ok).toBe(true);
|
||||
const message: unknown = await response.json();
|
||||
if (!isRecord(message) || typeof message.id !== "string") {
|
||||
throw new Error("Discord reply omitted its message ID");
|
||||
}
|
||||
const event = await delivered.promise;
|
||||
const retainedResponse = await fetch(`${messagesUrl}/${message.id}`, { headers });
|
||||
expect(retainedResponse.ok).toBe(true);
|
||||
for (const received of [message, event, await retainedResponse.json()]) {
|
||||
expect(received).toMatchObject(payload);
|
||||
if (!isRecord(received) || !Array.isArray(received.attachments)) {
|
||||
throw new Error("Discord reply omitted its attachments");
|
||||
}
|
||||
expect(received.attachments).toHaveLength(1);
|
||||
const attachment: unknown = received.attachments[0];
|
||||
expect(attachment).toMatchObject({ description: text });
|
||||
if (!isRecord(attachment)) {
|
||||
throw new Error("Discord attachment is invalid");
|
||||
}
|
||||
for (const field of ["url", "proxy_url"]) {
|
||||
const url = attachment[field];
|
||||
if (typeof url !== "string") {
|
||||
throw new Error(`Discord attachment omitted ${field}`);
|
||||
}
|
||||
expect(new URL(url).origin).toBe(relayOrigin);
|
||||
const download = await fetch(url);
|
||||
expect(download.ok).toBe(true);
|
||||
await expect(download.text()).resolves.toBe("attachment bytes");
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
socket?.terminate();
|
||||
await relay?.cleanup();
|
||||
await adapter.close();
|
||||
}
|
||||
});
|
||||
});
|
||||
314
extensions/qa-lab/src/crabline-discord-replies.ts
Normal file
314
extensions/qa-lab/src/crabline-discord-replies.ts
Normal file
|
|
@ -0,0 +1,314 @@
|
|||
import { once } from "node:events";
|
||||
import { createServer, request, type IncomingMessage } from "node:http";
|
||||
import { pipeline } from "node:stream/promises";
|
||||
import type { StartedOpenClawCrablineCorrelatedAdapter } from "@openclaw/crabline";
|
||||
import { fetchWithSsrFGuard } from "openclaw/plugin-sdk/ssrf-runtime";
|
||||
import { isRecord, readStringValue } from "openclaw/plugin-sdk/string-coerce-runtime";
|
||||
import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress";
|
||||
import WebSocket, { WebSocketServer } from "ws";
|
||||
import { closeQaHttpServer, dispatchQaHttpRequest } from "./bus-server.js";
|
||||
import { readQaJsonResponse } from "./ignored-response-body.js";
|
||||
import { readLiveQaChannelAccounts } from "./live-transports/shared/live-channel-status.js";
|
||||
import {
|
||||
waitForQaTransportCondition,
|
||||
type QaTransportAdapter,
|
||||
type QaTransportState,
|
||||
} from "./qa-transport.js";
|
||||
import { extractQaFailureReplyText } from "./reply-failure.js";
|
||||
import type { QaBusInboundMessageInput, QaBusMessage } from "./runtime-api.js";
|
||||
|
||||
export async function startCrablineDiscordReplies(params: {
|
||||
adapter: StartedOpenClawCrablineCorrelatedAdapter;
|
||||
state: QaTransportState;
|
||||
targets: ReadonlyMap<string, Pick<QaBusInboundMessageInput, "conversation" | "threadId">>;
|
||||
}) {
|
||||
const { state } = params;
|
||||
const manifest = params.adapter.manifest;
|
||||
if (manifest.provider !== "discord") {
|
||||
return undefined;
|
||||
}
|
||||
const upstream = new URL(manifest.endpoints.apiRoot);
|
||||
const lifecycle = new AbortController();
|
||||
const sockets = new Set<WebSocket>();
|
||||
const gatewayServer = new WebSocketServer({ noServer: true });
|
||||
let generation = 0;
|
||||
let lastDelivery: { messageId: string; channelId: string } | undefined;
|
||||
const server = createServer((req, res) => {
|
||||
dispatchQaHttpRequest(res, async () => {
|
||||
const channelId =
|
||||
req.method === "POST" && req.headers.authorization === `Bot ${manifest.botToken}`
|
||||
? req.url?.match(/^\/api\/v10\/channels\/(\d+)\/messages(?:\?.*)?$/u)?.[1]
|
||||
: undefined;
|
||||
const requestGeneration = generation;
|
||||
if (channelId) {
|
||||
lastDelivery = undefined;
|
||||
}
|
||||
const forwarded = request(upstream, {
|
||||
method: req.method,
|
||||
path: req.url,
|
||||
headers: { ...req.headers, host: upstream.host },
|
||||
signal: lifecycle.signal,
|
||||
});
|
||||
const received = new Promise<IncomingMessage>((resolve, reject) => {
|
||||
forwarded.once("response", resolve);
|
||||
forwarded.once("error", reject);
|
||||
});
|
||||
const [response] = await Promise.all([received, pipeline(req, forwarded)]);
|
||||
if (!channelId && !response.headers["content-type"]?.includes("application/json")) {
|
||||
res.writeHead(response.statusCode ?? 502, response.headers);
|
||||
await pipeline(response, res);
|
||||
return;
|
||||
}
|
||||
const chunks: Buffer[] = [];
|
||||
let size = 0;
|
||||
for await (const chunk of response) {
|
||||
const bytes = Buffer.from(chunk);
|
||||
size += bytes.length;
|
||||
if (size > 8 << 20) {
|
||||
throw new Error("Discord delivery response exceeded the reply receipt limit");
|
||||
}
|
||||
chunks.push(bytes);
|
||||
}
|
||||
const message: unknown = JSON.parse(Buffer.concat(chunks).toString("utf8"));
|
||||
if (
|
||||
channelId &&
|
||||
response.statusCode &&
|
||||
response.statusCode >= 200 &&
|
||||
response.statusCode < 300
|
||||
) {
|
||||
const messageId = isRecord(message) ? readStringValue(message.id) : undefined;
|
||||
if (!messageId || !/^\d+$/u.test(messageId)) {
|
||||
throw new Error("Discord delivery response omitted its message id");
|
||||
}
|
||||
if (requestGeneration === generation) {
|
||||
lastDelivery = { messageId, channelId };
|
||||
}
|
||||
}
|
||||
const pathname = new URL(req.url ?? "/", upstream).pathname;
|
||||
if (req.method === "GET" && /^\/api\/v10\/gateway(?:\/bot)?$/u.test(pathname)) {
|
||||
rewriteEndpointFields(message, ["url"], origins);
|
||||
} else if (/^\/api\/v10\/channels\/\d+\/messages(?:\/\d+)?$/u.test(pathname)) {
|
||||
rewriteMessageAttachmentUrls(message, origins);
|
||||
}
|
||||
const body = Buffer.from(JSON.stringify(message));
|
||||
const headers = { ...response.headers, "content-length": body.length };
|
||||
delete headers["transfer-encoding"];
|
||||
res.writeHead(response.statusCode ?? 502, headers);
|
||||
res.end(body);
|
||||
});
|
||||
});
|
||||
await once(server.listen(0, "127.0.0.1"), "listening");
|
||||
const address = server.address();
|
||||
if (!address || typeof address === "string") {
|
||||
throw new Error("Discord delivery receipt server failed to bind");
|
||||
}
|
||||
const origin = `http://127.0.0.1:${address.port}`;
|
||||
const origins = new Map([
|
||||
[upstream.origin, origin],
|
||||
[new URL(manifest.endpoints.gatewayUrl).origin, origin.replace(/^http/u, "ws")],
|
||||
]);
|
||||
server.on("upgrade", (req, socket, head) => {
|
||||
const target = new URL(manifest.endpoints.gatewayUrl);
|
||||
const requested = new URL(req.url ?? "/", upstream);
|
||||
if (requested.pathname !== target.pathname) {
|
||||
socket.destroy();
|
||||
return;
|
||||
}
|
||||
target.search = requested.search;
|
||||
const remote = new WebSocket(target);
|
||||
sockets.add(remote);
|
||||
let client: WebSocket | undefined;
|
||||
socket.once("close", () => remote.terminate());
|
||||
remote.once("error", () => {
|
||||
socket.destroy();
|
||||
client?.terminate();
|
||||
});
|
||||
remote.once("close", (code, reason) => {
|
||||
sockets.delete(remote);
|
||||
if (client) {
|
||||
closeGatewayPeer(client, code, reason);
|
||||
} else {
|
||||
socket.destroy();
|
||||
}
|
||||
});
|
||||
remote.once("open", () => {
|
||||
gatewayServer.handleUpgrade(req, socket, head, (connection) => {
|
||||
client = connection;
|
||||
connection.on("error", () => remote.terminate());
|
||||
connection.on("close", (code, reason) => closeGatewayPeer(remote, code, reason));
|
||||
connection.on("message", (data, binary) => remote.send(data, { binary }));
|
||||
remote.on("message", (data, binary) => {
|
||||
if (!binary) {
|
||||
const event: unknown = JSON.parse(rawDataToString(data));
|
||||
if (isRecord(event)) {
|
||||
if (event.t === "READY") {
|
||||
rewriteEndpointFields(event.d, ["resume_gateway_url"], origins);
|
||||
} else if (event.t === "MESSAGE_CREATE" || event.t === "MESSAGE_UPDATE") {
|
||||
rewriteMessageAttachmentUrls(event.d, origins);
|
||||
}
|
||||
}
|
||||
connection.send(JSON.stringify(event));
|
||||
return;
|
||||
}
|
||||
connection.send(data, { binary });
|
||||
});
|
||||
});
|
||||
});
|
||||
});
|
||||
const waitForCompletedReply: NonNullable<QaTransportAdapter["waitForCompletedReply"]> = async ({
|
||||
inbound,
|
||||
gateway,
|
||||
timeoutMs = 60_000,
|
||||
}) => {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
await waitForQaTransportCondition(() => {
|
||||
const messages = state.getSnapshot().messages;
|
||||
const inboundIndex = messages.findIndex(
|
||||
(message) => message.id === inbound.id && message.direction === "inbound",
|
||||
);
|
||||
return inboundIndex >= 0 &&
|
||||
messages
|
||||
.slice(inboundIndex + 1)
|
||||
.some(
|
||||
(message) =>
|
||||
message.direction === "outbound" && message.accountId === inbound.accountId,
|
||||
)
|
||||
? true
|
||||
: undefined;
|
||||
}, timeoutMs);
|
||||
await waitForQaTransportCondition(
|
||||
async () => {
|
||||
const accounts = await readLiveQaChannelAccounts(gateway, "discord", {
|
||||
timeoutMs: Math.max(1, deadline - Date.now()),
|
||||
});
|
||||
const account = accounts.find((entry) => entry.accountId === inbound.accountId);
|
||||
return account?.running === true &&
|
||||
account.connected === true &&
|
||||
account.restartPending !== true &&
|
||||
account.busy === false &&
|
||||
account.activeRuns === 0
|
||||
? true
|
||||
: undefined;
|
||||
},
|
||||
Math.max(1, deadline - Date.now()),
|
||||
);
|
||||
const delivery = lastDelivery;
|
||||
if (!delivery || BigInt(delivery.messageId) <= BigInt(inbound.id)) {
|
||||
throw new Error(`Discord inbound ${inbound.id} completed without a final delivery receipt`);
|
||||
}
|
||||
const target = params.targets.get(delivery.channelId);
|
||||
if (!target) {
|
||||
throw new Error(`Discord final delivery ${delivery.messageId} has no observed destination`);
|
||||
}
|
||||
const { response, release } = await fetchWithSsrFGuard({
|
||||
url: `${manifest.endpoints.apiRoot}/v10/channels/${delivery.channelId}/messages/${delivery.messageId}`,
|
||||
init: { headers: { authorization: `Bot ${manifest.botToken}` } },
|
||||
policy: { allowPrivateNetwork: true },
|
||||
timeoutMs: Math.max(1, deadline - Date.now()),
|
||||
auditContext: "qa-lab-crabline-discord-retained-reply",
|
||||
});
|
||||
const message = await readQaJsonResponse<unknown>(
|
||||
response,
|
||||
release,
|
||||
`Discord inbound ${inbound.id} completed without a retained reply: final delivery ${delivery.messageId} unavailable`,
|
||||
);
|
||||
if (
|
||||
!isRecord(message) ||
|
||||
message.id !== delivery.messageId ||
|
||||
!isRecord(message.author) ||
|
||||
message.author.id !== manifest.botUserId
|
||||
) {
|
||||
throw new Error(`Discord final delivery ${delivery.messageId} returned an invalid message`);
|
||||
}
|
||||
const replyToId = isRecord(message.message_reference)
|
||||
? readStringValue(message.message_reference.message_id)
|
||||
: undefined;
|
||||
const reply: QaBusMessage = {
|
||||
id: delivery.messageId,
|
||||
accountId: params.adapter.accountId,
|
||||
direction: "outbound",
|
||||
...target,
|
||||
senderId: manifest.botUserId,
|
||||
text: readStringValue(message.content) ?? "",
|
||||
timestamp: Date.parse(String(message.timestamp)),
|
||||
...(replyToId ? { replyToId } : {}),
|
||||
reactions: [],
|
||||
};
|
||||
const failure = extractQaFailureReplyText(reply);
|
||||
if (failure) {
|
||||
throw new Error(failure);
|
||||
}
|
||||
return reply;
|
||||
};
|
||||
let closing: Promise<void> | undefined;
|
||||
return {
|
||||
apiBaseUrl: `${origin}/api/v10`,
|
||||
waitForCompletedReply,
|
||||
reset() {
|
||||
generation += 1;
|
||||
lastDelivery = undefined;
|
||||
},
|
||||
cleanup() {
|
||||
lifecycle.abort();
|
||||
return (closing ??= (async () => {
|
||||
for (const socket of [...sockets, ...gatewayServer.clients]) {
|
||||
socket.terminate();
|
||||
}
|
||||
await Promise.all([
|
||||
closeQaHttpServer(server),
|
||||
new Promise<void>((resolve, reject) => {
|
||||
gatewayServer.close((error) => (error ? reject(error) : resolve()));
|
||||
}),
|
||||
]);
|
||||
})());
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function rewriteEndpointFields(
|
||||
value: unknown,
|
||||
fields: readonly string[],
|
||||
origins: ReadonlyMap<string, string>,
|
||||
): void {
|
||||
if (!isRecord(value)) {
|
||||
return;
|
||||
}
|
||||
for (const field of fields) {
|
||||
const url = value[field];
|
||||
if (typeof url !== "string") {
|
||||
continue;
|
||||
}
|
||||
for (const [source, target] of origins) {
|
||||
if (url === source || url.startsWith(`${source}/`) || url.startsWith(`${source}?`)) {
|
||||
value[field] = `${target}${url.slice(source.length)}`;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function rewriteMessageAttachmentUrls(value: unknown, origins: ReadonlyMap<string, string>): void {
|
||||
if (Array.isArray(value)) {
|
||||
for (const message of value) {
|
||||
rewriteMessageAttachmentUrls(message, origins);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (!isRecord(value)) {
|
||||
return;
|
||||
}
|
||||
if (Array.isArray(value.attachments)) {
|
||||
for (const attachment of value.attachments) {
|
||||
rewriteEndpointFields(attachment, ["url", "proxy_url"], origins);
|
||||
}
|
||||
}
|
||||
rewriteMessageAttachmentUrls(value.referenced_message, origins);
|
||||
}
|
||||
|
||||
function closeGatewayPeer(socket: WebSocket, code: number, reason: Buffer): void {
|
||||
if (code === 1005 || code === 1006 || code === 1015) {
|
||||
socket.terminate();
|
||||
} else {
|
||||
socket.close(code, reason);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,7 +1,13 @@
|
|||
import { createChannelRunQueue } from "openclaw/plugin-sdk/channel-outbound";
|
||||
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
|
||||
import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
|
||||
import { withTempDir } from "openclaw/plugin-sdk/test-env";
|
||||
import { rawDataToString } from "openclaw/plugin-sdk/webhook-ingress";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import WebSocket from "ws";
|
||||
import { createQaBusState } from "./bus-state.js";
|
||||
import { createQaCrablineTransportAdapter } from "./crabline-transport.js";
|
||||
import { runLoadedScenarioFlow } from "./scenario-flow-runner.test-support.js";
|
||||
|
||||
function requireString(value: unknown, label: string): string {
|
||||
if (typeof value !== "string" || value.length === 0) {
|
||||
|
|
@ -22,7 +28,11 @@ describe("Crabline Discord transport", () => {
|
|||
providerReadinessArtifactPath: "crabline-provider-readiness.json",
|
||||
},
|
||||
state: createQaBusState(),
|
||||
transportPolicy: { requireGroupMention: true, senderAllowlist: ["driver"] },
|
||||
transportPolicy: {
|
||||
requireGroupMention: true,
|
||||
senderAllowlist: ["driver"],
|
||||
topLevelReplies: true,
|
||||
},
|
||||
});
|
||||
const config = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" });
|
||||
const discord = config.channels?.discord;
|
||||
|
|
@ -34,6 +44,7 @@ describe("Crabline Discord transport", () => {
|
|||
DISCORD_API_URL: expect.stringMatching(/^http:\/\/127\.0\.0\.1:\d+\/api\/v10$/u),
|
||||
});
|
||||
expect(discord).toMatchObject({
|
||||
replyToMode: "off",
|
||||
allowFrom: [expect.stringMatching(/^\d{17,20}$/u)],
|
||||
dmPolicy: "allowlist",
|
||||
groupPolicy: "allowlist",
|
||||
|
|
@ -45,6 +56,39 @@ describe("Crabline Discord transport", () => {
|
|||
},
|
||||
});
|
||||
|
||||
const gatewayResponse = await fetch(`${runtimeEnv.DISCORD_API_URL}/gateway/bot`, {
|
||||
headers: { authorization: `Bot ${requireString(discord?.token, "Discord bot token")}` },
|
||||
});
|
||||
expect(gatewayResponse.ok).toBe(true);
|
||||
const gateway: unknown = await gatewayResponse.json();
|
||||
const gatewayUrl = requireString(
|
||||
isRecord(gateway) ? gateway.url : undefined,
|
||||
"Gateway URL",
|
||||
);
|
||||
expect(new URL(gatewayUrl).origin).toBe(
|
||||
new URL(requireString(runtimeEnv.DISCORD_API_URL, "Discord API URL")).origin.replace(
|
||||
/^http/u,
|
||||
"ws",
|
||||
),
|
||||
);
|
||||
const socket = new WebSocket(gatewayUrl);
|
||||
const ready = createDeferred<unknown>();
|
||||
socket.on("error", ready.reject);
|
||||
socket.on("message", (data) => {
|
||||
const event: unknown = JSON.parse(rawDataToString(data));
|
||||
if (isRecord(event) && event.t === "READY") {
|
||||
ready.resolve(event.d);
|
||||
}
|
||||
});
|
||||
socket.on("open", () => {
|
||||
socket.send(JSON.stringify({ op: 2, d: { token: discord?.token, intents: 0 } }));
|
||||
});
|
||||
try {
|
||||
await expect(ready.promise).resolves.toMatchObject({ resume_gateway_url: gatewayUrl });
|
||||
} finally {
|
||||
socket.terminate();
|
||||
}
|
||||
|
||||
const inbound = await transport.sendInbound({
|
||||
conversation: { id: "discord-crabline-primary", kind: "group" },
|
||||
senderId: "driver",
|
||||
|
|
@ -53,6 +97,7 @@ describe("Crabline Discord transport", () => {
|
|||
threadId: "discord-crabline-thread",
|
||||
});
|
||||
expect(inbound).toMatchObject({
|
||||
id: expect.stringMatching(/^\d{17,20}$/u),
|
||||
conversation: { id: "discord-crabline-primary", kind: "group" },
|
||||
threadId: "discord-crabline-thread",
|
||||
});
|
||||
|
|
@ -74,7 +119,10 @@ describe("Crabline Discord transport", () => {
|
|||
const response = await fetch(
|
||||
`${runtimeEnv.DISCORD_API_URL}/channels/${channelId}/messages`,
|
||||
{
|
||||
body: JSON.stringify({ content: "Discord provider outbound marker." }),
|
||||
body: JSON.stringify({
|
||||
content: "Discord provider outbound marker.",
|
||||
message_reference: { message_id: inbound.id },
|
||||
}),
|
||||
headers: {
|
||||
authorization: `Bot ${requireString(discord?.token, "Discord bot token")}`,
|
||||
"content-type": "application/json",
|
||||
|
|
@ -92,6 +140,7 @@ describe("Crabline Discord transport", () => {
|
|||
timeoutMs: 1_000,
|
||||
}),
|
||||
).resolves.toMatchObject({
|
||||
replyToId: inbound.id,
|
||||
conversation: { id: "discord-crabline-primary", kind: "group" },
|
||||
threadId: "discord-crabline-thread",
|
||||
});
|
||||
|
|
@ -109,4 +158,155 @@ describe("Crabline Discord transport", () => {
|
|||
}
|
||||
});
|
||||
});
|
||||
|
||||
it("judges retained final deliveries after the channel run drains, not matching previews", async () => {
|
||||
await withTempDir("qa-crabline-discord-finals-", async (outputDir) => {
|
||||
const transport = await createQaCrablineTransportAdapter({
|
||||
outputDir,
|
||||
selection: {
|
||||
capabilityMatrixPath: "crabline-channel-driver-capabilities.json",
|
||||
channel: "discord",
|
||||
channelDriver: "crabline",
|
||||
providerReadinessArtifactPath: "crabline-provider-readiness.json",
|
||||
},
|
||||
});
|
||||
const config = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" });
|
||||
const apiUrl = transport.createRuntimeEnvPatch?.().DISCORD_API_URL;
|
||||
const headers = {
|
||||
authorization: `Bot ${requireString(config.channels?.discord?.token, "Discord bot token")}`,
|
||||
"content-type": "application/json",
|
||||
};
|
||||
const marker = "QA-TOP-LEVEL-REPLY-OK";
|
||||
let accountStatus = { accountId: transport.accountId, running: true, connected: true };
|
||||
const runQueue = createChannelRunQueue({
|
||||
setStatus: (patch) => {
|
||||
accountStatus = { ...accountStatus, ...patch };
|
||||
},
|
||||
});
|
||||
try {
|
||||
for (const finalKind of [
|
||||
"top-level",
|
||||
"reply",
|
||||
"edited",
|
||||
"deleted",
|
||||
"deleted-final",
|
||||
"edited-final",
|
||||
] as const) {
|
||||
const releaseFinal = createDeferred<void>();
|
||||
const waitingForDelivery = createDeferred<"waiting">();
|
||||
const deliveryDone = createDeferred<void>();
|
||||
const gateway = {
|
||||
call: async () => {
|
||||
waitingForDelivery.resolve("waiting");
|
||||
return { channelAccounts: { discord: [accountStatus] } };
|
||||
},
|
||||
};
|
||||
const flowTransport = {
|
||||
...transport,
|
||||
sendInbound: async (input: Parameters<typeof transport.sendInbound>[0]) => {
|
||||
const inbound = await transport.sendInbound(input);
|
||||
const delivery = transport.buildAgentDelivery({
|
||||
target: `group:${input.conversation.id}`,
|
||||
});
|
||||
const channelId = requireString(delivery.to, "Discord delivery target").replace(
|
||||
/^channel:/u,
|
||||
"",
|
||||
);
|
||||
const messagesUrl = `${apiUrl}/channels/${channelId}/messages`;
|
||||
runQueue.enqueue(inbound.id, async () => {
|
||||
try {
|
||||
const previewResponse = await fetch(messagesUrl, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify({
|
||||
content: marker,
|
||||
...(finalKind === "top-level"
|
||||
? { message_reference: { message_id: inbound.id } }
|
||||
: {}),
|
||||
}),
|
||||
});
|
||||
expect(previewResponse.ok).toBe(true);
|
||||
const preview = (await previewResponse.json()) as { id: string };
|
||||
await releaseFinal.promise;
|
||||
const response = await fetch(
|
||||
finalKind === "edited" || finalKind === "deleted"
|
||||
? `${messagesUrl}/${preview.id}`
|
||||
: messagesUrl,
|
||||
{
|
||||
method:
|
||||
finalKind === "edited"
|
||||
? "PATCH"
|
||||
: finalKind === "deleted"
|
||||
? "DELETE"
|
||||
: "POST",
|
||||
headers,
|
||||
...(finalKind === "deleted"
|
||||
? {}
|
||||
: {
|
||||
body: JSON.stringify({
|
||||
content: finalKind === "edited" ? "wrong final" : marker,
|
||||
...(finalKind === "reply"
|
||||
? { message_reference: { message_id: inbound.id } }
|
||||
: {}),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
);
|
||||
expect(response.ok).toBe(true);
|
||||
if (finalKind === "deleted-final" || finalKind === "edited-final") {
|
||||
const final = (await response.json()) as { id: string };
|
||||
expect(final.id).not.toBe(preview.id);
|
||||
const mutation = await fetch(`${messagesUrl}/${final.id}`, {
|
||||
method: finalKind === "deleted-final" ? "DELETE" : "PATCH",
|
||||
headers,
|
||||
...(finalKind === "edited-final"
|
||||
? { body: JSON.stringify({ content: "wrong final" }) }
|
||||
: {}),
|
||||
});
|
||||
expect(mutation.ok).toBe(true);
|
||||
await mutation.body?.cancel();
|
||||
} else {
|
||||
await response.body?.cancel();
|
||||
}
|
||||
deliveryDone.resolve();
|
||||
} catch (error) {
|
||||
deliveryDone.reject(error);
|
||||
}
|
||||
});
|
||||
return inbound;
|
||||
},
|
||||
};
|
||||
const result = runLoadedScenarioFlow("channel-top-level-reply-shape", {
|
||||
api: { transport: flowTransport, env: { gateway } },
|
||||
});
|
||||
const observed = result.then(
|
||||
() => "completed",
|
||||
() => "failed",
|
||||
);
|
||||
try {
|
||||
expect(await Promise.race([waitingForDelivery.promise, observed])).toBe("waiting");
|
||||
} finally {
|
||||
releaseFinal.resolve();
|
||||
await deliveryDone.promise;
|
||||
}
|
||||
if (finalKind === "top-level") {
|
||||
await expect(result).resolves.toMatchObject({ status: "pass" });
|
||||
} else {
|
||||
await expect(result).rejects.toThrow(
|
||||
finalKind === "reply"
|
||||
? "expected top-level reply"
|
||||
: finalKind === "edited" || finalKind === "edited-final"
|
||||
? "completed delivery did not contain"
|
||||
: finalKind === "deleted-final"
|
||||
? "final delivery"
|
||||
: "completed without a retained reply",
|
||||
);
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
runQueue.deactivate();
|
||||
await transport.cleanupAfterGatewayStop?.();
|
||||
}
|
||||
});
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -93,6 +93,7 @@ describe("crabline transport", () => {
|
|||
transportPolicy: {
|
||||
requireGroupMention: true,
|
||||
senderAllowlist: ["driver"],
|
||||
topLevelReplies: true,
|
||||
},
|
||||
selection: createSelection(),
|
||||
state: createQaBusState(),
|
||||
|
|
@ -102,6 +103,7 @@ describe("crabline transport", () => {
|
|||
const gatewayConfig = transport.createGatewayConfig({ baseUrl: "http://127.0.0.1:1" });
|
||||
const telegramConfig = gatewayConfig.channels?.telegram;
|
||||
expect(telegramConfig).toMatchObject({
|
||||
replyToMode: "off",
|
||||
allowFrom: [expect.stringMatching(/^[1-9]\d+$/u)],
|
||||
groupAllowFrom: [expect.stringMatching(/^[1-9]\d+$/u)],
|
||||
groupPolicy: "allowlist",
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ import {
|
|||
readStringValue,
|
||||
} from "openclaw/plugin-sdk/string-coerce-runtime";
|
||||
import { createQaBusState, type QaBusState } from "./bus-state.js";
|
||||
import { startCrablineDiscordReplies } from "./crabline-discord-replies.js";
|
||||
import {
|
||||
createCrablineProviderCorrelation,
|
||||
createCrablineProviderDelivery,
|
||||
|
|
@ -41,6 +42,7 @@ import type { QaBusInboundMessageInput, QaBusMessage } from "./runtime-api.js";
|
|||
|
||||
type QaCrablineTransportState = QaTransportState & {
|
||||
slackIngress?: ReturnType<typeof createCrablineSlackIngress>;
|
||||
discordReplies?: Awaited<ReturnType<typeof startCrablineDiscordReplies>>;
|
||||
cleanup: () => Promise<void>;
|
||||
getOutboundEvents: () => Promise<readonly QaTransportOutboundEvent[]>;
|
||||
observeEvent: (event: unknown) => void;
|
||||
|
|
@ -191,6 +193,8 @@ async function postCrablineInbound(params: {
|
|||
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 === "discord" && isRecord(result) && isRecord(result.message)) {
|
||||
providerMessageId = readStringValue(result.message.id);
|
||||
} else if (
|
||||
params.adapter.channel === "telegram" &&
|
||||
isRecord(result) &&
|
||||
|
|
@ -202,10 +206,10 @@ async function postCrablineInbound(params: {
|
|||
return { providerMessageId, response: result };
|
||||
}
|
||||
|
||||
function createCrablineState(params: {
|
||||
async function createCrablineState(params: {
|
||||
adapter: StartedOpenClawCrablineCorrelatedAdapter;
|
||||
state: QaBusState;
|
||||
}): QaCrablineTransportState {
|
||||
}): Promise<QaCrablineTransportState> {
|
||||
const baseState = params.state;
|
||||
const slackIngress =
|
||||
params.adapter.manifest.provider === "slack"
|
||||
|
|
@ -215,15 +219,20 @@ function createCrablineState(params: {
|
|||
const telegramMessageByProviderId = new Map<string, QaBusMessage>();
|
||||
const pendingTelegramMessagesByChat = new Map<string, QaBusMessage[]>();
|
||||
const outboundEvents: QaTransportOutboundEvent[] = [];
|
||||
const discordTargets = new Map<string, QaCrablineTarget>();
|
||||
const discordReplies = await startCrablineDiscordReplies({ ...params, targets: discordTargets });
|
||||
const resetTransport = () => {
|
||||
targetByProviderTarget.clear();
|
||||
telegramMessageByProviderId.clear();
|
||||
pendingTelegramMessagesByChat.clear();
|
||||
outboundEvents.length = 0;
|
||||
discordTargets.clear();
|
||||
discordReplies?.reset();
|
||||
};
|
||||
|
||||
return {
|
||||
...(slackIngress ? { slackIngress } : {}),
|
||||
...(discordReplies ? { discordReplies } : {}),
|
||||
reset() {
|
||||
resetTransport();
|
||||
baseState.reset();
|
||||
|
|
@ -277,11 +286,34 @@ function createCrablineState(params: {
|
|||
? { to: observation.fallbackTarget }
|
||||
: undefined;
|
||||
if (destination) {
|
||||
const replyToId =
|
||||
params.adapter.channel === "discord" &&
|
||||
isRecord(event) &&
|
||||
isRecord(event.body) &&
|
||||
isRecord(event.body.message_reference)
|
||||
? readStringValue(event.body.message_reference.message_id)
|
||||
: undefined;
|
||||
if (params.adapter.channel === "discord" && isRecord(event)) {
|
||||
const channelId = readStringValue(event.path)?.match(
|
||||
/^\/api\/v10\/channels\/(\d+)\/messages$/u,
|
||||
)?.[1];
|
||||
if (channelId) {
|
||||
const parsed = parseQaTarget(destination.to);
|
||||
discordTargets.set(
|
||||
channelId,
|
||||
target ?? {
|
||||
conversation: { id: parsed.conversationId, kind: parsed.chatType },
|
||||
...(parsed.threadId ? { threadId: parsed.threadId } : {}),
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
baseState.addOutboundMessage({
|
||||
accountId: observation.accountId,
|
||||
senderId: observation.senderId,
|
||||
senderName: observation.senderName,
|
||||
text: observation.text,
|
||||
...(replyToId ? { replyToId } : {}),
|
||||
...destination,
|
||||
});
|
||||
}
|
||||
|
|
@ -333,6 +365,7 @@ function createCrablineState(params: {
|
|||
waitFor: baseState.waitFor.bind(baseState),
|
||||
async cleanup() {
|
||||
await slackIngress?.cleanup();
|
||||
await discordReplies?.cleanup();
|
||||
await params.adapter.close();
|
||||
},
|
||||
};
|
||||
|
|
@ -349,14 +382,20 @@ function createQaCrablineTransport(params: {
|
|||
const stateMethods = createQaTransportStateMethods({ state, accountId: adapter.accountId });
|
||||
const hooks: Pick<
|
||||
QaTransportAdapter,
|
||||
"prepareFlow" | "cleanup" | "sendNativeCommand" | "waitForOutboundSequence"
|
||||
| "prepareFlow"
|
||||
| "cleanup"
|
||||
| "sendNativeCommand"
|
||||
| "waitForOutboundSequence"
|
||||
| "waitForCompletedReply"
|
||||
> = {};
|
||||
let releaseDiscordQaApiBase: (() => void) | undefined;
|
||||
if (params.state.slackIngress) {
|
||||
hooks.prepareFlow = params.state.slackIngress.prepareFlow;
|
||||
hooks.cleanup = params.state.slackIngress.cleanup;
|
||||
}
|
||||
if (params.selection.channel === "discord" && params.adapter.manifest.provider === "discord") {
|
||||
if (state.discordReplies && params.adapter.manifest.provider === "discord") {
|
||||
const { apiBaseUrl, waitForCompletedReply } = state.discordReplies;
|
||||
hooks.waitForCompletedReply = waitForCompletedReply;
|
||||
const manifest = params.adapter.manifest;
|
||||
let prepared:
|
||||
| Promise<
|
||||
|
|
@ -369,7 +408,7 @@ function createQaCrablineTransport(params: {
|
|||
prepared ??= (async () => {
|
||||
const scenarioRuntime = await import("./live-transports/discord/discord-live.runtime.js");
|
||||
releaseDiscordQaApiBase = scenarioRuntime.registerDiscordQaApiBase({
|
||||
apiBaseUrl: `${manifest.endpoints.apiRoot}/v10`,
|
||||
apiBaseUrl,
|
||||
tokens: [manifest.botToken, manifest.driverBotToken],
|
||||
});
|
||||
const [sutIdentity, driverIdentity] = await Promise.all([
|
||||
|
|
@ -433,6 +472,7 @@ function createQaCrablineTransport(params: {
|
|||
...config.channels,
|
||||
discord: {
|
||||
...discord,
|
||||
...(transportPolicy?.topLevelReplies ? { replyToMode: "off" as const } : {}),
|
||||
allowFrom: [...dmAllowlist],
|
||||
...(dmAllowlist.includes("*") ? {} : { dmPolicy: "allowlist" as const }),
|
||||
...(senderAllowlist ? { groupPolicy: "allowlist" as const } : {}),
|
||||
|
|
@ -460,7 +500,11 @@ function createQaCrablineTransport(params: {
|
|||
const senderAllowlist = transportPolicy?.senderAllowlist?.map(
|
||||
(senderId) => adapter.createAgentDelivery({ target: `dm:${senderId}` }).providerTargetKey,
|
||||
);
|
||||
if (!transportPolicy?.requireGroupMention && !senderAllowlist) {
|
||||
if (
|
||||
!transportPolicy?.requireGroupMention &&
|
||||
!senderAllowlist &&
|
||||
!transportPolicy?.topLevelReplies
|
||||
) {
|
||||
return config as QaTransportGatewayConfig;
|
||||
}
|
||||
return {
|
||||
|
|
@ -469,6 +513,7 @@ function createQaCrablineTransport(params: {
|
|||
...config.channels,
|
||||
telegram: {
|
||||
...config.channels?.telegram,
|
||||
...(transportPolicy?.topLevelReplies ? { replyToMode: "off" as const } : {}),
|
||||
...(senderAllowlist
|
||||
? {
|
||||
allowFrom: [...senderAllowlist],
|
||||
|
|
@ -547,9 +592,9 @@ function createQaCrablineTransport(params: {
|
|||
},
|
||||
|
||||
createRuntimeEnvPatch: () =>
|
||||
adapter.manifest.provider === "discord"
|
||||
state.discordReplies
|
||||
? {
|
||||
DISCORD_API_URL: `${adapter.manifest.endpoints.apiRoot}/v10`,
|
||||
DISCORD_API_URL: state.discordReplies.apiBaseUrl,
|
||||
}
|
||||
: adapter.createProviderReadinessEnv({}),
|
||||
|
||||
|
|
@ -623,13 +668,24 @@ export async function createQaCrablineTransportAdapter(params: {
|
|||
recorderPath,
|
||||
});
|
||||
// Readiness owns the startup probe; runtime transcripts may contain provider-specific API records.
|
||||
let readiness: Awaited<ReturnType<typeof runOpenClawCrablineProviderReadiness>>;
|
||||
try {
|
||||
readiness = await runOpenClawCrablineProviderReadiness({
|
||||
const readiness = await runOpenClawCrablineProviderReadiness({
|
||||
adapter,
|
||||
outputDir: params.outputDir,
|
||||
selection: params.selection,
|
||||
});
|
||||
const state = await createCrablineState({
|
||||
adapter,
|
||||
state: params.state ?? createQaBusState(),
|
||||
});
|
||||
observeEvent = state.observeEvent;
|
||||
return createQaCrablineTransport({
|
||||
adapter,
|
||||
readiness,
|
||||
transportPolicy: params.transportPolicy,
|
||||
selection: params.selection,
|
||||
state,
|
||||
});
|
||||
} catch (error) {
|
||||
try {
|
||||
await adapter.close();
|
||||
|
|
@ -640,18 +696,6 @@ export async function createQaCrablineTransportAdapter(params: {
|
|||
}
|
||||
throw error;
|
||||
}
|
||||
const state = createCrablineState({
|
||||
adapter,
|
||||
state: params.state ?? createQaBusState(),
|
||||
});
|
||||
observeEvent = state.observeEvent;
|
||||
return createQaCrablineTransport({
|
||||
adapter,
|
||||
readiness,
|
||||
transportPolicy: params.transportPolicy,
|
||||
selection: params.selection,
|
||||
state,
|
||||
});
|
||||
}
|
||||
|
||||
export async function createQaCrablineTransportDefinition(
|
||||
|
|
|
|||
|
|
@ -34,6 +34,59 @@ async function readRecorderEvents(recorderPath: string): Promise<RecorderEvent[]
|
|||
}
|
||||
|
||||
describe("Discord Crabline real-plugin roundtrip", () => {
|
||||
it.runIf(RUN_DISCORD_CRABLINE_E2E)(
|
||||
"accepts top-level finals and rejects native reply finals through the shipped scenario",
|
||||
async () => {
|
||||
for (const forceReply of [false, true]) {
|
||||
const outputDir = path.join(
|
||||
process.cwd(),
|
||||
".artifacts",
|
||||
"qa-e2e",
|
||||
`discord-final-shape-${process.pid}-${Date.now()}`,
|
||||
);
|
||||
const suite = await runQaSuite({
|
||||
channelDriver: "crabline",
|
||||
channelId: "discord",
|
||||
controlUiEnabled: false,
|
||||
outputDir,
|
||||
providerMode: "mock-openai",
|
||||
repoRoot: process.cwd(),
|
||||
scenarioIds: ["channel-top-level-reply-shape"],
|
||||
...(forceReply
|
||||
? {
|
||||
mutateConfig: (cfg) => ({
|
||||
...cfg,
|
||||
channels: {
|
||||
...cfg.channels,
|
||||
discord: { ...cfg.channels?.discord, replyToMode: "all" },
|
||||
},
|
||||
}),
|
||||
}
|
||||
: {}),
|
||||
});
|
||||
expect(suite.result.scenarios).toEqual([
|
||||
expect.objectContaining({ status: forceReply ? "fail" : "pass" }),
|
||||
]);
|
||||
if (forceReply) {
|
||||
expect(JSON.stringify(suite.result.scenarios)).toContain("expected top-level reply");
|
||||
}
|
||||
const events = await readRecorderEvents(
|
||||
path.join(outputDir, "artifacts", "crabline", "discord-provider-server.jsonl"),
|
||||
);
|
||||
const final = events.findLast(
|
||||
(event) =>
|
||||
event.type === "api" &&
|
||||
event.method === "POST" &&
|
||||
event.accepted === true &&
|
||||
event.body?.content === "QA-TOP-LEVEL-REPLY-OK",
|
||||
);
|
||||
expect(final).toBeDefined();
|
||||
expect(Boolean(readObject(final?.body?.message_reference)?.message_id)).toBe(forceReply);
|
||||
}
|
||||
},
|
||||
360_000,
|
||||
);
|
||||
|
||||
it.runIf(RUN_DISCORD_CRABLINE_E2E)(
|
||||
"crosses the real Discord REST and Gateway boundaries and closes every owned resource",
|
||||
async () => {
|
||||
|
|
@ -144,6 +197,8 @@ describe("Discord Crabline real-plugin roundtrip", () => {
|
|||
const snapshotEvents = await readRecorderEvents(
|
||||
path.resolve(suite.result.outputDir, snapshotRecorderPath),
|
||||
);
|
||||
expect(path.resolve(suite.result.outputDir, snapshotRecorderPath)).not.toBe(recorderPath);
|
||||
expect(snapshotEvents.some((event) => event.accepted === true)).toBe(true);
|
||||
expect(
|
||||
snapshotEvents.some(
|
||||
(event) =>
|
||||
|
|
@ -153,7 +208,7 @@ describe("Discord Crabline real-plugin roundtrip", () => {
|
|||
(readStringValue(readObject(event.body)?.content) ?? "").includes(EXPECTED_MARKER) &&
|
||||
event.accepted === true,
|
||||
),
|
||||
).toBe(true);
|
||||
).toBe(false);
|
||||
|
||||
// The suite returns only after Gateway, WebSocket, HTTP, recorder, and temporary runtime
|
||||
// owners have all completed their ordered cleanup.
|
||||
|
|
|
|||
|
|
@ -14,6 +14,7 @@ import type {
|
|||
QaTransportPolicy,
|
||||
QaTransportReportParams,
|
||||
} from "./qa-transport.js";
|
||||
import { waitForCompletedQaReply } from "./suite-runtime-transport.js";
|
||||
|
||||
const QA_CHANNEL_ID = "qa-channel";
|
||||
const QA_CHANNEL_ACCOUNT_ID = "default";
|
||||
|
|
@ -94,6 +95,8 @@ export function createQaChannelTransport(state: QaBusState, transportPolicy?: Qa
|
|||
accountId: QA_CHANNEL_ACCOUNT_ID,
|
||||
requiredPluginIds: QA_CHANNEL_REQUIRED_PLUGIN_IDS,
|
||||
supportedActions: ["delete", "edit", "react", "thread-create"],
|
||||
waitForCompletedReply: ({ inbound, timeoutMs }) =>
|
||||
waitForCompletedQaReply(state, inbound, timeoutMs),
|
||||
async reset() {
|
||||
await waitForQaTransportCondition(() => {
|
||||
if (
|
||||
|
|
|
|||
|
|
@ -242,6 +242,11 @@ export type QaTransportAdapter = Omit<
|
|||
reset: () => Promise<void>;
|
||||
waitForNoOutbound: (input?: QaTransportWaitForNoOutboundInput) => Promise<void>;
|
||||
waitForOutbound: (input: QaTransportOutboundMatch) => Promise<QaBusMessage>;
|
||||
waitForCompletedReply?: (input: {
|
||||
inbound: QaBusMessage;
|
||||
gateway: Parameters<QaTransportAdapterDefinition["waitReady"]>[0]["gateway"];
|
||||
timeoutMs?: number;
|
||||
}) => Promise<QaBusMessage>;
|
||||
waitForCondition: <T>(
|
||||
check: () => T | Promise<T | null | undefined> | null | undefined,
|
||||
timeoutMs?: number,
|
||||
|
|
|
|||
14
qa/README.md
14
qa/README.md
|
|
@ -27,6 +27,20 @@ message returned by `sendInbound`. It waits for the channel's processing
|
|||
acknowledgment before reading the retained reply. QA-channel reset also waits for
|
||||
pending inbound turns before clearing observations.
|
||||
|
||||
Reply-shape scenarios use `transport.waitForCompletedReply` with the inbound
|
||||
message and Gateway client. For Crabline Discord, the adapter waits for an
|
||||
observed delivery and the channel run queue to become idle. Its local API relay
|
||||
keeps REST, Gateway discovery, and resumed WebSockets on the same origin. It
|
||||
rewrites only Gateway discovery `url`, READY `resume_gateway_url`, and message
|
||||
attachment `url`/`proxy_url` (including referenced messages). Message content,
|
||||
embeds, and attachment descriptions pass through unchanged. It records the
|
||||
message ID from the final send response, then reads that exact native
|
||||
message before checking its text, destination, or quote relation. A deleted final
|
||||
fails explicitly; a retained preview cannot replace it, and edits are judged from
|
||||
the retained final. QA-channel uses its processing acknowledgment instead.
|
||||
Adapters without a completion boundary fail explicitly rather than falling back
|
||||
to the first matching outbound observation.
|
||||
|
||||
Generated-media scenarios count attachment deliveries separately from text
|
||||
progress and check the saved bytes plus the persisted completion reply.
|
||||
|
||||
|
|
|
|||
|
|
@ -45,14 +45,17 @@ flow:
|
|||
senderName: QA Driver
|
||||
text:
|
||||
expr: "`@openclaw reply exactly: ${config.expectedMarker}`"
|
||||
- waitForOutbound:
|
||||
conversation:
|
||||
id: { ref: config.conversationId }
|
||||
kind: group
|
||||
textIncludes: { ref: config.expectedMarker }
|
||||
timeoutMs:
|
||||
expr: liveTurnTimeoutMs(env, 60000)
|
||||
saveAs: inbound
|
||||
- call: transport.waitForCompletedReply
|
||||
args:
|
||||
- inbound: { ref: inbound }
|
||||
gateway: { expr: env.gateway }
|
||||
timeoutMs:
|
||||
expr: liveTurnTimeoutMs(env, 60000)
|
||||
saveAs: outbound
|
||||
- assert:
|
||||
expr: "outbound.text.includes(config.expectedMarker) && outbound.conversation.id === inbound.conversation.id && outbound.conversation.kind === inbound.conversation.kind"
|
||||
message: "completed delivery did not contain the exact marker in the originating conversation"
|
||||
- assert:
|
||||
expr: "!outbound.threadId && (transport.id === 'qa-channel' || !outbound.replyToId)"
|
||||
message:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue