From 59b5fd7759a380e1c7b83ae5be188909b625d4f5 Mon Sep 17 00:00:00 2001 From: Josh Avant <830519+joshavant@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:52:38 -0500 Subject: [PATCH] fix(qa): await completed memory replies before assertions and reset (#160929) * fix(qa): await completed memory replies before assertions and reset * fix(qa): restore memory config after completion timeout --- .../mock-openai/mock-openai-assistant-text.ts | 3 + .../qa-lab/src/qa-channel-reset.test.ts | 74 ++++++ .../qa-lab/src/qa-channel-transport.test.ts | 5 +- extensions/qa-lab/src/qa-channel-transport.ts | 23 ++ .../qa-lab/src/scenario-catalog.test.ts | 6 - .../src/scenario-memory-cleanup.test.ts | 111 +++++++++ .../scenario-memory-completed-reply.test.ts | 231 ++++++++++++++++++ .../qa-lab/src/scenario-runtime-api.test.ts | 4 - extensions/qa-lab/src/scenario-runtime-api.ts | 2 - .../qa-lab/src/suite-runtime-transport.ts | 58 +++++ qa/README.md | 5 + .../memory/remember-across-conversations.yaml | 192 ++++++++------- 12 files changed, 610 insertions(+), 104 deletions(-) create mode 100644 extensions/qa-lab/src/qa-channel-reset.test.ts create mode 100644 extensions/qa-lab/src/scenario-memory-cleanup.test.ts create mode 100644 extensions/qa-lab/src/scenario-memory-completed-reply.test.ts diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-assistant-text.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-assistant-text.ts index 87ebd3842936..5feadf9909a4 100644 --- a/extensions/qa-lab/src/providers/mock-openai/mock-openai-assistant-text.ts +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-assistant-text.ts @@ -241,6 +241,9 @@ export function buildAssistantText(input: ResponsesInputItem[], body: Record ({ + setTimeout: (ms: number) => + new Promise((resolve) => { + setTimeout(resolve, ms); + }), +})); + +beforeEach(() => vi.useFakeTimers()); +afterEach(() => vi.useRealTimers()); + +describe("QA channel reset processing boundary", () => { + it.each(["default", "other"])( + "retains a pending %s turn through its final edit", + async (accountId) => { + const state = createQaBusState(); + const transport = createQaChannelTransport(state); + state.addInboundMessage({ + accountId, + conversation: { id: "alice", kind: "direct" }, + senderId: "alice", + text: "Recall my preference.", + }); + const inboundCursor = state.getSnapshot().cursor; + const preview = state.addOutboundMessage({ accountId, to: "dm:alice", text: "You" }); + // Fetch progress and another account's completion cannot release this turn. + state.resolvePollCursor({ accountId, cursor: state.getSnapshot().cursor }); + state.resolvePollCursor({ + accountId: "unrelated", + acknowledgedCursor: state.getSnapshot().cursor, + }); + const reset = transport.reset(); + try { + await vi.advanceTimersByTimeAsync(100); + expect(state.getSnapshot().messages.map(({ text }) => text)).toEqual([ + "Recall my preference.", + "You", + ]); + state.editMessage({ + accountId, + messageId: preview.id, + text: "lemon pepper wings with blue cheese", + }); + } finally { + state.resolvePollCursor({ accountId, acknowledgedCursor: inboundCursor }); + await vi.runAllTimersAsync(); + await reset; + } + expect(state.getSnapshot().messages).toEqual([]); + expect(state.getSnapshot().cursor).toBeGreaterThanOrEqual(inboundCursor); + }, + ); + + it("leaves unsettled messages intact when completion times out", async () => { + const state = createQaBusState(); + state.addInboundMessage({ + conversation: { id: "alice", kind: "direct" }, + senderId: "alice", + text: "Still running.", + }); + const before = state.getSnapshot(); + const reset = createQaChannelTransport(state).reset(); + const outcome = reset.then( + () => "cleared", + (error: unknown) => String(error), + ); + await vi.runAllTimersAsync(); + expect(await outcome).toBe("Error: timed out after 15000ms"); + expect(state.getSnapshot()).toEqual(before); + }); +}); diff --git a/extensions/qa-lab/src/qa-channel-transport.test.ts b/extensions/qa-lab/src/qa-channel-transport.test.ts index 91557be191f8..120625e6f29d 100644 --- a/extensions/qa-lab/src/qa-channel-transport.test.ts +++ b/extensions/qa-lab/src/qa-channel-transport.test.ts @@ -162,7 +162,8 @@ describe("qa channel transport", () => { }); it("implements the portable scenario transport actions", async () => { - const transport = createQaChannelTransport(createQaBusState()); + const state = createQaBusState(); + const transport = createQaChannelTransport(state); const conversation = { id: "alice", kind: "direct" as const }; await transport.sendInbound({ @@ -178,6 +179,8 @@ describe("qa channel transport", () => { await expect( transport.waitForOutbound({ conversation, textIncludes: "QA-PORTABLE-OK" }), ).resolves.toMatchObject({ text: "QA-PORTABLE-OK" }); + // The synthetic fixture has no channel poller; record its completed turn. + state.resolvePollCursor({ acknowledgedCursor: state.getSnapshot().cursor }); await transport.reset(); expect(transport.state.getSnapshot().messages).toEqual([]); }); diff --git a/extensions/qa-lab/src/qa-channel-transport.ts b/extensions/qa-lab/src/qa-channel-transport.ts index f11705f93aeb..2bf56b48eb59 100644 --- a/extensions/qa-lab/src/qa-channel-transport.ts +++ b/extensions/qa-lab/src/qa-channel-transport.ts @@ -3,6 +3,7 @@ import { getQaProvider } from "./providers/index.js"; import { QaStateBackedTransportAdapter, waitForQaTransportAccountReady, + waitForQaTransportCondition, waitForQaTransportOutboundSequence, } from "./qa-transport.js"; import type { @@ -87,6 +88,7 @@ async function handleQaChannelAction( class QaChannelTransport extends QaStateBackedTransportAdapter { readonly #transportPolicy?: QaTransportPolicy; + readonly #busState: QaBusState; constructor(state: QaBusState, transportPolicy?: QaTransportPolicy) { super({ @@ -98,6 +100,27 @@ class QaChannelTransport extends QaStateBackedTransportAdapter { state, }); this.#transportPolicy = transportPolicy; + this.#busState = state; + } + + override async reset() { + await waitForQaTransportCondition(() => { + if ( + this.#busState + .getSnapshot() + .events.some( + (event) => + event.kind === "inbound-message" && + this.#busState.getAcknowledgedPollCursor(event.accountId) < event.cursor, + ) + ) { + return undefined; + } + // Reset clears every account. Check and clear together so a newly admitted + // turn cannot lose its message while an earlier turn is being drained. + this.#busState.reset(); + return true; + }); } createGatewayConfig = ({ baseUrl }: { baseUrl: string }) => diff --git a/extensions/qa-lab/src/scenario-catalog.test.ts b/extensions/qa-lab/src/scenario-catalog.test.ts index 711a495fe12c..28b25fa9736d 100644 --- a/extensions/qa-lab/src/scenario-catalog.test.ts +++ b/extensions/qa-lab/src/scenario-catalog.test.ts @@ -700,12 +700,6 @@ describe("qa scenario catalog", () => { expect(flow).toContain("[sourceSessionKey, targetSessionKey, groupSessionKey]"); expect(flow).toContain("readSessionTranscriptSummary"); expect(flow).toContain("transcript.eventCursor > 0"); - expect(flow).toContain( - "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length", - ); - expect(flow).toContain('"saveAs":"pauseCommandOutbound"'); - expect(flow).toContain("candidate.conversation.id === config.pausedConversationId"); - expect(flow).toContain('"sinceIndex":{"ref":"pauseCommandStartIndex"}'); expect(flow).not.toContain('"call":"sleep"'); expect(flow).not.toContain(".sessionFile"); }); diff --git a/extensions/qa-lab/src/scenario-memory-cleanup.test.ts b/extensions/qa-lab/src/scenario-memory-cleanup.test.ts new file mode 100644 index 000000000000..95242c82c0c7 --- /dev/null +++ b/extensions/qa-lab/src/scenario-memory-cleanup.test.ts @@ -0,0 +1,111 @@ +import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { createQaBusState } from "./bus-state.js"; +import { readQaScenarioById } from "./scenario-catalog.js"; +import { runScenarioFlow } from "./scenario-flow-runner.js"; +import { waitForQaInboundCompletion } from "./suite-runtime-transport.js"; + +vi.mock("node:timers/promises", () => ({ + setTimeout: (ms: number) => + new Promise((resolve) => { + setTimeout(resolve, ms); + }), +})); + +beforeEach(() => vi.useFakeTimers()); +afterEach(() => vi.useRealTimers()); + +const scenario = readQaScenarioById("remember-across-conversations"); +const guarded = scenario.execution.flow!.steps[0]!.actions.find( + (action) => isRecord(action) && "try" in action, +); +if (!isRecord(guarded) || !isRecord(guarded.try)) { + throw new Error("expected memory scenario cleanup boundary"); +} +const memoryTry = guarded.try; + +describe("memory scenario cleanup", () => { + it.each([ + { scenarioFails: true, completionTimesOut: true, restorationFails: false }, + { scenarioFails: false, completionTimesOut: true, restorationFails: false }, + { scenarioFails: true, completionTimesOut: true, restorationFails: true }, + { scenarioFails: false, completionTimesOut: true, restorationFails: true }, + { scenarioFails: true, completionTimesOut: false, restorationFails: true }, + { scenarioFails: false, completionTimesOut: false, restorationFails: true }, + { scenarioFails: false, completionTimesOut: false, restorationFails: false }, + ])( + "restores config and preserves the primary failure: %j", + async ({ scenarioFails, completionTimesOut, restorationFails }) => { + const state = createQaBusState(); + const lastInbound = state.addInboundMessage({ + conversation: { id: "remember-disabled", kind: "direct" }, + senderId: "remember-disabled", + text: "Recall my preference.", + }); + if (!completionTimesOut) { + state.resolvePollCursor({ acknowledgedCursor: state.getSnapshot().cursor }); + } + const before = state.getSnapshot(); + const originalMemorySearch = { rememberAcrossConversations: true, sources: ["sessions"] }; + let memorySearch = { ...originalMemorySearch, rememberAcrossConversations: false }; + const scenarioError = new Error("wrong recalled preference"); + const restorationError = new Error("restored config failed its health check"); + const outcome = runScenarioFlow({ + scenarioTitle: scenario.title, + // Substitute only the scenario body; execute its authored error and cleanup handlers. + flow: { + steps: [ + { + name: "cleanup", + actions: [{ try: { ...memoryTry, actions: [{ call: "recall" }] } }], + }, + ], + }, + vars: { lastInbound, originalMemorySearch }, + api: { + scenario, + config: scenario.execution.config ?? {}, + state, + env: {}, + liveTurnTimeoutMs: (_env: unknown, timeoutMs: number) => timeoutMs, + waitForQaInboundCompletion, + recall: () => { + if (scenarioFails) { + throw scenarioError; + } + }, + patchConfig: ({ patch }: { patch: { memory: { search: typeof memorySearch } } }) => { + memorySearch = patch.memory.search; + }, + waitForGatewayHealthy: () => { + if (restorationFails) { + throw restorationError; + } + }, + waitForQaChannelReady: () => {}, + runScenario: async (name, steps) => { + for (const step of steps) { + await step.run(); + } + return { name, status: "pass", steps: [] }; + }, + }, + }).catch((error: unknown) => error); + await vi.runAllTimersAsync(); + const error = await outcome; + expect + .soft(memorySearch) + .toEqual({ rememberAcrossConversations: true, sources: ["sessions"] }); + expect(state.getSnapshot()).toEqual(before); + if (scenarioFails) { + expect(error).toBe(scenarioError); + } else if (completionTimesOut) { + expect(error).toEqual(new Error("timed out after 60000ms")); + } else if (restorationFails) { + expect(error).toBe(restorationError); + } else { + expect(error).toEqual({ name: scenario.title, status: "pass", steps: [] }); + } + }, + ); +}); diff --git a/extensions/qa-lab/src/scenario-memory-completed-reply.test.ts b/extensions/qa-lab/src/scenario-memory-completed-reply.test.ts new file mode 100644 index 000000000000..bd2669660709 --- /dev/null +++ b/extensions/qa-lab/src/scenario-memory-completed-reply.test.ts @@ -0,0 +1,231 @@ +import path from "node:path"; +import { buildChannelInboundEventContext } from "openclaw/plugin-sdk/channel-inbound"; +import { + createPluginRuntimeMock, + createStartAccountContext, +} from "openclaw/plugin-sdk/channel-test-helpers"; +import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; +import { createReplyDispatcher, settleReplyDispatcher } from "openclaw/plugin-sdk/reply-runtime"; +import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; +import { injectQaBusInboundMessage, qaChannelPlugin } from "../../qa-channel/api.js"; +import { startQaBusServer } from "./bus-server.js"; +import { createQaBusState } from "./bus-state.js"; +import { createQaChannelTransport } from "./qa-channel-transport.js"; +import { readQaScenarioById } from "./scenario-catalog.js"; +import { runScenarioFlow } from "./scenario-flow-runner.js"; +import { runQaSuiteScenarioSteps } from "./suite-runtime-flow.js"; +import { waitForCompletedQaReply } from "./suite-runtime-transport.js"; + +const scenario = readQaScenarioById("remember-across-conversations"); +const config = scenario.execution.config ?? {}; +const step = scenario.execution.flow!.steps[0]!; +const guarded = step.actions.find((action) => isRecord(action) && "try" in action); +if (!isRecord(guarded) || !isRecord(guarded.try) || !Array.isArray(guarded.try.actions)) { + throw new Error("expected memory recall actions"); +} +const actions: unknown[] = guarded.try.actions; +const recallStart = actions.findIndex( + (action) => isRecord(action) && action.set === "requestCursorBeforeRecall", +); +const recallEnd = actions.findIndex((action) => isRecord(action) && "if" in action); +if (recallStart < 0 || recallEnd <= recallStart) { + throw new Error("expected memory recall and evidence boundary"); +} + +describe("memory scenario completed reply", () => { + const state = createQaBusState(); + const transport = createQaChannelTransport(state); + const runtime = createPluginRuntimeMock({ + channel: { inbound: { buildContext: buildChannelInboundEventContext } }, + }); + const controller = new AbortController(); + let bus: Awaited>; + let gateway: Promise; + let observeAcknowledgment = (_cursor: number) => {}; + let observeCompletionRead = (_accountId: string | undefined, _cursor: number) => {}; + + beforeAll(async () => { + bus = await startQaBusServer({ state }); + const resolvePollCursor = state.resolvePollCursor.bind(state); + vi.spyOn(state, "resolvePollCursor").mockImplementation((input) => { + const cursor = resolvePollCursor(input); + observeAcknowledgment(input?.acknowledgedCursor ?? 0); + return cursor; + }); + const getAcknowledgedPollCursor = state.getAcknowledgedPollCursor.bind(state); + vi.spyOn(state, "getAcknowledgedPollCursor").mockImplementation((accountId) => { + const cursor = getAcknowledgedPollCursor(accountId); + observeCompletionRead(accountId, cursor); + return cursor; + }); + const cfg = transport.createGatewayConfig({ baseUrl: bus.baseUrl }); + const ready = createDeferred(); + const context = createStartAccountContext({ + account: qaChannelPlugin.config.resolveAccount(cfg, transport.accountId), + cfg, + abortSignal: controller.signal, + statusPatchSink: (snapshot) => { + if (snapshot.lifecycle === "ready") { + ready.resolve(); + } + }, + }); + const startAccount = qaChannelPlugin.gateway?.startAccount; + if (!startAccount) { + throw new Error("expected QA channel gateway entry point"); + } + gateway = Promise.resolve(startAccount({ ...context, channelRuntime: runtime.channel })); + await Promise.race([ + ready.promise, + gateway.then(() => { + throw new Error("QA channel stopped before ready"); + }), + ]); + }); + + afterAll(async () => { + controller.abort(); + await gateway; + await bus.stop(); + }); + + async function runRecall(finalText?: string, helperLeak?: "group" | "anchor") { + state.reset(); + const previewSent = createDeferred(); + const releaseFinal = createDeferred(); + const acknowledged = createDeferred(); + const waitingForCompletion = createDeferred<"waiting">(); + let inboundCursor = 0; + observeAcknowledgment = (cursor) => { + if (inboundCursor > 0 && cursor >= inboundCursor) { + acknowledged.resolve(); + } + }; + observeCompletionRead = (accountId, cursor) => { + if (accountId === "default" && inboundCursor > 0 && cursor < inboundCursor) { + waitingForCompletion.resolve("waiting"); + } + }; + // Model execution is held; the real channel preview, dispatcher, HTTP bus, + // and poller's processing acknowledgment establish completion independently. + vi.mocked( + runtime.channel.reply.dispatchReplyWithBufferedBlockDispatcher, + ).mockImplementationOnce(async ({ dispatcherOptions, replyOptions }) => { + await replyOptions?.onPartialReply?.({ text: "You usually want" }); + previewSent.resolve(); + await releaseFinal.promise; + const dispatcher = createReplyDispatcher(dispatcherOptions); + try { + if (finalText !== undefined) { + dispatcher.sendFinalReply({ text: finalText }); + } + } finally { + await settleReplyDispatcher({ dispatcher }); + } + return { queuedFinal: finalText !== undefined, counts: dispatcher.getQueuedCounts() }; + }); + const vars: Record = { + transcriptRoot: "/synthetic-memory-transcripts", + sourceSession: { sessionId: "private-source-transcript" }, + targetSession: { sessionId: "anchor-transcript" }, + groupSession: { sessionId: "group-transcript" }, + initialSessionsVisibility: "tree", + }; + const result = runScenarioFlow({ + scenarioTitle: scenario.title, + // Execute the shipped recall selector and all its fact/source assertions. + // Independent session seeding and config-toggle scenarios are live-suite proof. + flow: { steps: [{ name: step.name, actions: actions.slice(recallStart, recallEnd) }] }, + vars, + api: { + scenario, + config, + state, + path, + env: { providerMode: "live-frontier" }, + waitForCompletedQaReply, + transport: { + accountId: transport.accountId, + sendInbound: async (input: Parameters[0]) => { + const { message } = await injectQaBusInboundMessage({ baseUrl: bus.baseUrl, input }); + inboundCursor = state + .getSnapshot() + .events.findLast( + (event) => event.kind === "inbound-message" && event.message.id === message.id, + )!.cursor; + return message; + }, + }, + fs: { + readdir: async () => ["recall.jsonl"], + readFile: async () => + ["memory_search", "private-source-transcript", helperLeak && `${helperLeak}-transcript`] + .filter(Boolean) + .join("\n"), + }, + readConfigSnapshot: async () => ({ + config: { tools: { sessions: { visibility: "tree" } } }, + }), + liveTurnTimeoutMs: (_env: unknown, timeoutMs: number) => timeoutMs, + waitForCondition: async (check: () => T | Promise | undefined) => { + const value = await check(); + if (value === undefined) { + throw new Error("expected helper transcript after completed recall"); + } + return value; + }, + runScenario: runQaSuiteScenarioSteps, + }, + }); + try { + await previewSent.promise; + expect(await Promise.race([waitingForCompletion.promise, result])).toBe("waiting"); + expect(state.getAcknowledgedPollCursor("default")).toBeLessThan(inboundCursor); + } finally { + releaseFinal.resolve(); + await acknowledged.promise; + await result; + } + if (finalText === undefined) { + expect(vars.targetOutbound).toBeUndefined(); + } else { + expect(vars.targetOutbound).toMatchObject({ text: finalText }); + } + return await result; + } + + it("reads the edited final only after the originating turn completes", async () => { + expect(await runRecall("lemon pepper wings with blue cheese")).toMatchObject({ + status: "pass", + }); + }); + + it("rejects an acknowledged turn whose preview was removed without a retained reply", async () => { + expect(await runRecall()).toMatchObject({ + status: "fail", + details: expect.stringContaining("completed without a retained reply"), + }); + }); + + it.each([ + "lemon pepper wings with ranch", + "lemon pepper wings with blue cheese; GROUP-ONLY loaded nachos with black olives", + "lemon pepper wings with blue cheese; ANCHOR-ONLY pretzel bites test marker", + ])("rejects the completed wrong or leaking preference: %s", async (text) => { + expect(await runRecall(text)).toMatchObject({ + status: "fail", + details: expect.stringContaining("private target missed recalled preference"), + }); + }); + + it.each(["group", "anchor"] as const)( + "rejects completed recall with %s helper leakage", + async (leak) => { + expect(await runRecall("lemon pepper wings with blue cheese", leak)).toMatchObject({ + status: "fail", + details: expect.stringContaining(`${leak} transcript`), + }); + }, + ); +}); diff --git a/extensions/qa-lab/src/scenario-runtime-api.test.ts b/extensions/qa-lab/src/scenario-runtime-api.test.ts index 96a0a6362a80..f6af88d2f912 100644 --- a/extensions/qa-lab/src/scenario-runtime-api.test.ts +++ b/extensions/qa-lab/src/scenario-runtime-api.test.ts @@ -30,7 +30,6 @@ describe("createQaScenarioRuntimeApi", () => { } return value; }; - const sleep = vi.fn(async () => undefined); const env = { lab: { baseUrl: "http://127.0.0.1:1234" }, transport: { @@ -78,7 +77,6 @@ describe("createQaScenarioRuntimeApi", () => { }, }; const deps = { - sleep, waitForTransportReady: vi.fn(), waitForAgentHistoryReply: vi.fn(), browserRequest: vi.fn(), @@ -126,7 +124,6 @@ describe("createQaScenarioRuntimeApi", () => { expect(outboundSpy).toHaveBeenCalledTimes(1); expect(readSpy).toHaveBeenCalledTimes(1); expect(resetSpy).toHaveBeenCalledTimes(3); - expect(sleep).toHaveBeenCalledTimes(3); }); it("routes scenario injection through a factory-created transport", async () => { @@ -181,7 +178,6 @@ describe("createQaScenarioRuntimeApi", () => { execution: { kind: "flow", flow: { steps: [] } }, }, deps: { - sleep: vi.fn(async () => undefined), waitForTransportReady: vi.fn(), }, constants, diff --git a/extensions/qa-lab/src/scenario-runtime-api.ts b/extensions/qa-lab/src/scenario-runtime-api.ts index d1bcf0fb191d..bccc6a8ef472 100644 --- a/extensions/qa-lab/src/scenario-runtime-api.ts +++ b/extensions/qa-lab/src/scenario-runtime-api.ts @@ -21,7 +21,6 @@ export type QaScenarioRuntimeEnv< }; type QaScenarioRuntimeApiDeps = { - sleep: (ms?: number) => Promise; waitForTransportReady: (...args: never[]) => unknown; }; @@ -69,7 +68,6 @@ export function createQaScenarioRuntimeApi< const transportState = transport.state; const resetTransportState = async () => { await transport.reset(); - await params.deps.sleep(100); }; return { diff --git a/extensions/qa-lab/src/suite-runtime-transport.ts b/extensions/qa-lab/src/suite-runtime-transport.ts index 620ab720d8c3..73b3f9dc2804 100644 --- a/extensions/qa-lab/src/suite-runtime-transport.ts +++ b/extensions/qa-lab/src/suite-runtime-transport.ts @@ -1,4 +1,5 @@ import { setTimeout as sleep } from "node:timers/promises"; +import type { QaBusState } from "./bus-state.js"; import { findFailureOutboundMessage, waitForQaTransportCondition, @@ -11,6 +12,61 @@ type WaitForNoOutboundOptions = { sinceIndex?: number; }; +async function waitForQaInboundCompletion( + state: QaBusState, + inbound: QaBusMessage, + timeoutMs = 15_000, +) { + const event = state + .getSnapshot() + .events.find( + (candidate) => + candidate.kind === "inbound-message" && + candidate.accountId === inbound.accountId && + candidate.message.id === inbound.id, + ); + if (!event) { + throw new Error(`QA inbound event missing for ${inbound.id}`); + } + await waitForQaTransportCondition( + () => state.getAcknowledgedPollCursor(inbound.accountId) >= event.cursor || undefined, + timeoutMs, + ); +} + +async function waitForCompletedQaReply( + state: QaBusState, + inbound: QaBusMessage, + timeoutMs = 15_000, +) { + await waitForQaInboundCompletion(state, inbound, timeoutMs); + // Snapshots are clones. Read again after the channel has drained preview edits + // and final delivery, rather than retaining the first streamed fragment. + const replies = state + .getSnapshot() + .messages.filter( + (message) => + message.direction === "outbound" && + !message.deleted && + message.accountId === inbound.accountId && + message.conversation.id === inbound.conversation.id && + message.conversation.kind === inbound.conversation.kind && + message.threadId === inbound.threadId && + message.replyToId === inbound.id, + ); + for (const reply of replies) { + const failure = extractQaFailureReplyText(reply); + if (failure) { + throw new Error(failure); + } + } + const reply = replies.at(-1); + if (!reply) { + throw new Error(`QA inbound ${inbound.id} completed without a retained reply`); + } + return reply; +} + async function waitForOutboundMessage( state: QaTransportState, predicate: (message: QaBusMessage) => boolean, @@ -127,6 +183,8 @@ export { formatTransportTranscript, readTransportTranscript, recentOutboundSummary, + waitForCompletedQaReply, + waitForQaInboundCompletion, waitForNoOutbound, waitForNoOutbound as waitForNoTransportOutbound, waitForOutboundMessage, diff --git a/qa/README.md b/qa/README.md index d5a0aaef18ea..8f7a2824fa3e 100644 --- a/qa/README.md +++ b/qa/README.md @@ -22,6 +22,11 @@ Operator workflows: Keep this folder in git. Add new scenarios here before wiring them into automation. +QA-channel flows that assert final replies use `waitForCompletedQaReply` with the +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. + Generated-media scenarios count attachment deliveries separately from text progress and check the saved bytes plus the persisted completion reply. diff --git a/qa/scenarios/memory/remember-across-conversations.yaml b/qa/scenarios/memory/remember-across-conversations.yaml index a30a747b4e8a..cfb6ae3a1f0b 100644 --- a/qa/scenarios/memory/remember-across-conversations.yaml +++ b/qa/scenarios/memory/remember-across-conversations.yaml @@ -35,7 +35,7 @@ scenario: maxSummaryChars: 220 successCriteria: - Two private conversations keep distinct session keys and persisted transcript identities. - - A private reply recalls the relevant fact from the other private conversation. + - A completed private reply recalls the exact relevant fact from the other private conversation. - Accepted memory-search evidence excludes the group transcript and anchor transcript. - Group destinations, the disabled setting, and a session-scoped pause do not start product recall. - The feature does not widen tools.sessions.visibility. @@ -66,8 +66,7 @@ scenario: groupFact: GROUP-ONLY loaded nachos with black olives anchorFact: ANCHOR-ONLY pretzel bites test marker sourceMarker: QA-REMEMBER-SOURCE-SEEDED - recallPrompt: "Remember across conversations QA check: what snack do I usually want for QA movie night? Reply in one short sentence." - expectedNeedle: lemon pepper wings with blue cheese + recallPrompt: "Remember across conversations QA check: what snack do I usually want for QA movie night? Reply with only the snack preference, verbatim, without punctuation or any other text." transcriptDir: qa-remember-across-conversations flow: @@ -144,17 +143,18 @@ flow: expr: config.sourceConversationId senderName: Remember Source text: - expr: "`Stable QA movie night usual favorite snack preference: ${config.privateFact}. Reply exactly: ${config.sourceMarker}.`" - - waitForOutbound: - conversation: - id: - expr: config.sourceConversationId - kind: direct - textIncludes: - expr: config.sourceMarker - timeoutMs: - expr: liveTurnTimeoutMs(env, 60000) + expr: "`Stable QA movie night usual favorite snack preference: ${config.privateFact}. Reply exactly: ${config.sourceMarker}`" + saveAs: lastInbound + - call: waitForCompletedQaReply saveAs: sourceOutbound + args: + - ref: state + - ref: lastInbound + - expr: liveTurnTimeoutMs(env, 60000) + - assert: + expr: sourceOutbound.text === config.sourceMarker + message: + expr: "`private source did not acknowledge the seeded fact: ${sourceOutbound.text}`" - sendInbound: conversation: id: @@ -165,14 +165,12 @@ flow: senderName: Remember Group Member text: expr: "`@openclaw Stable QA movie night usual favorite snack preference: ${config.groupFact}. This applies only inside this group. Acknowledge briefly.`" - - waitForOutbound: - conversation: - id: - expr: config.groupConversationId - kind: channel - timeoutMs: - expr: liveTurnTimeoutMs(env, 60000) - saveAs: groupOutbound + saveAs: lastInbound + - call: waitForCompletedQaReply + args: + - ref: state + - ref: lastInbound + - expr: liveTurnTimeoutMs(env, 60000) - sendInbound: conversation: id: @@ -183,14 +181,12 @@ flow: senderName: Remember Target text: expr: "`QA movie night snack recall anchor fixture: ${config.anchorFact}. This is a test marker, not a preference. Acknowledge briefly.`" - - waitForOutbound: - conversation: - id: - expr: config.targetConversationId - kind: direct - timeoutMs: - expr: liveTurnTimeoutMs(env, 60000) - saveAs: anchorOutbound + saveAs: lastInbound + - call: waitForCompletedQaReply + args: + - ref: state + - ref: lastInbound + - expr: liveTurnTimeoutMs(env, 60000) - call: readRawQaSessionStore saveAs: seededStore args: @@ -248,9 +244,6 @@ flow: - set: requestCursorBeforeRecall value: expr: "env.mock ? (await fetchJson(`${env.mock.baseUrl}/debug/request-cursor`)).cursor : 0" - - set: targetStartIndex - value: - expr: "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length" - sendInbound: conversation: id: @@ -261,16 +254,13 @@ flow: senderName: Remember Target text: expr: config.recallPrompt - - call: waitForOutboundMessage + saveAs: lastInbound + - call: waitForCompletedQaReply saveAs: targetOutbound args: - ref: state - - lambda: - params: [candidate] - expr: "candidate.conversation.id === config.targetConversationId && candidate.direction === 'outbound'" + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - sinceIndex: - ref: targetStartIndex - call: waitForCondition saveAs: helperTranscriptPath args: @@ -291,7 +281,7 @@ flow: value: expr: "recallRequests.map((request) => ({ plannedToolName: request.plannedToolName ?? null, plannedToolArgs: request.plannedToolArgs ?? null, toolOutput: String(request.toolOutput ?? '').slice(0, 1200), finalText: String(request.finalText ?? '').slice(0, 300), allInputText: String(request.allInputText ?? '').slice(-500) }))" - assert: - expr: "normalizeLowercaseStringOrEmpty(targetOutbound.text).includes(normalizeLowercaseStringOrEmpty(config.expectedNeedle))" + expr: targetOutbound.text === config.privateFact message: expr: "`private target missed recalled preference: reply=${targetOutbound.text}; sessions=${JSON.stringify({ sourceSession, targetSession, groupSession })}; helper=${helperTranscriptText}; requests=${JSON.stringify(recallRequestDebug)}`" - assert: @@ -342,14 +332,15 @@ flow: senderName: Remember Group Member text: expr: "`@openclaw ${config.recallPrompt}`" - - call: waitForCondition - saveAs: groupDestinationRequests + saveAs: lastInbound + - call: waitForCompletedQaReply args: - - lambda: - async: true - expr: "await (async () => { const requests = await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${groupRequestCursorBefore}`); return requests.length > 0 ? requests : undefined; })()" + - ref: state + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - 100 + - set: groupDestinationRequests + value: + expr: "await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${groupRequestCursorBefore}`)" - set: helperCountAfterGroupDestination value: expr: "(await fs.readdir(transcriptRoot).catch(() => [])).filter((entry) => entry.endsWith('.jsonl')).length" @@ -363,9 +354,6 @@ flow: expr: "groupDestinationRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes(''))" message: expr: "`group destination received private recall context: ${JSON.stringify(groupDestinationRequests)}`" - - set: pauseCommandStartIndex - value: - expr: "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length" - sendInbound: conversation: id: @@ -374,18 +362,15 @@ flow: senderId: qa-operator senderName: QA Operator text: /active-memory off - - call: waitForOutboundMessage + saveAs: lastInbound + - call: waitForCompletedQaReply saveAs: pauseCommandOutbound args: - ref: state - - lambda: - params: [candidate] - expr: "candidate.direction === 'outbound' && candidate.conversation.id === config.pausedConversationId && candidate.text.includes('Active Memory: off for this session.')" + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - sinceIndex: - ref: pauseCommandStartIndex - assert: - expr: "pauseCommandOutbound.text.includes('Active Memory: off for this session.')" + expr: "pauseCommandOutbound.text === 'Active Memory: off for this session.'" message: expr: "`unexpected Active Memory command response: ${JSON.stringify(pauseCommandOutbound)}`" - set: helperCountBeforePaused @@ -404,14 +389,15 @@ flow: senderName: Remember Paused text: expr: config.recallPrompt - - call: waitForCondition - saveAs: pausedRequests + saveAs: lastInbound + - call: waitForCompletedQaReply args: - - lambda: - async: true - expr: "await (async () => { const requests = await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${pausedRequestCursorBefore}`); return requests.length > 0 ? requests : undefined; })()" + - ref: state + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - 100 + - set: pausedRequests + value: + expr: "await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${pausedRequestCursorBefore}`)" - set: helperCountAfterPaused value: expr: "(await fs.readdir(transcriptRoot).catch(() => [])).filter((entry) => entry.endsWith('.jsonl')).length" @@ -419,7 +405,7 @@ flow: expr: helperCountAfterPaused === helperCountBeforePaused message: session-scoped pause unexpectedly started recall - assert: - expr: "pausedRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes(''))" + expr: "pausedRequests.length > 0 && pausedRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes(''))" message: expr: "`paused conversation received private recall context: ${JSON.stringify(pausedRequests)}`" - set: freshRequestCursorBefore @@ -435,14 +421,15 @@ flow: senderName: Remember Fresh text: expr: config.recallPrompt - - call: waitForCondition - saveAs: freshRequests + saveAs: lastInbound + - call: waitForCompletedQaReply args: - - lambda: - async: true - expr: "await (async () => { const requests = await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${freshRequestCursorBefore}`); return requests.some((request) => String(request.allInputText ?? '').includes('')) ? requests : undefined; })()" + - ref: state + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - 100 + - set: freshRequests + value: + expr: "await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${freshRequestCursorBefore}`)" - set: helperCountAfterFresh value: expr: "(await fs.readdir(transcriptRoot).catch(() => [])).filter((entry) => entry.endsWith('.jsonl')).length" @@ -485,14 +472,15 @@ flow: senderName: Remember Disabled text: expr: config.recallPrompt - - call: waitForCondition - saveAs: disabledRequests + saveAs: lastInbound + - call: waitForCompletedQaReply args: - - lambda: - async: true - expr: "await (async () => { const requests = await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${disabledRequestCursorBefore}`); return requests.length > 0 ? requests : undefined; })()" + - ref: state + - ref: lastInbound - expr: liveTurnTimeoutMs(env, 60000) - - 100 + - set: disabledRequests + value: + expr: "await fetchJson(`${env.mock.baseUrl}/debug/requests?after=${disabledRequestCursorBefore}`)" - set: helperCountAfterDisabled value: expr: "(await fs.readdir(transcriptRoot).catch(() => [])).filter((entry) => entry.endsWith('.jsonl')).length" @@ -500,24 +488,46 @@ flow: expr: helperCountAfterDisabled === helperCountBeforeDisabled message: disabled setting unexpectedly started Remember across conversations recall - assert: - expr: "disabledRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes(''))" + expr: "disabledRequests.length > 0 && disabledRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes(''))" message: expr: "`disabled setting received private recall context: ${JSON.stringify(disabledRequests)}`" + catchAs: scenarioError finally: - - call: patchConfig - args: - - env: - ref: env - patch: - memory: - search: - expr: "originalMemorySearch === undefined ? null : structuredClone(originalMemorySearch)" - - call: waitForGatewayHealthy - args: - - ref: env - - 60000 - - call: waitForQaChannelReady - args: - - ref: env - - 60000 + - try: + actions: + - if: + expr: "typeof lastInbound !== 'undefined'" + then: + - call: waitForQaInboundCompletion + args: + - ref: state + - ref: lastInbound + - expr: liveTurnTimeoutMs(env, 60000) + catchAs: completionError + finally: + - try: + actions: + - call: patchConfig + args: + - env: + ref: env + patch: + memory: + search: + expr: "originalMemorySearch === undefined ? null : structuredClone(originalMemorySearch)" + - call: waitForGatewayHealthy + args: + - ref: env + - 60000 + - call: waitForQaChannelReady + args: + - ref: env + - 60000 + finally: + # Cleanup must restore config without replacing the first failure. + - if: + expr: "typeof scenarioError !== 'undefined' || typeof completionError !== 'undefined'" + then: + - throw: + expr: "typeof scenarioError !== 'undefined' ? scenarioError : completionError" detailsExpr: "[`sourceSession=${sourceSessionKey}`, `targetSession=${targetSessionKey}`, `groupSession=${groupSessionKey}`, `reply=${targetOutbound.text}`, `helperTranscript=${helperTranscriptPath}`, `visibility=${String(afterSessionsVisibility)}`].join('\\n')"