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
This commit is contained in:
Josh Avant 2026-09-29 01:52:38 -05:00 • committed by GitHub
parent 2a34f3bb65
commit 59b5fd7759
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
12 changed files with 610 additions and 104 deletions

View file

@ -241,6 +241,9 @@ export function buildAssistantText(input: ResponsesInputItem[], body: Record<str
return `Protocol note: I checked memory and the project codename is ${orbitCode}.`;
}
if (isSnackRecallPrompt(prompt) && snackPreference) {
if (prompt.includes("Reply with only the snack preference, verbatim")) {
return snackPreference;
}
return `Protocol note: you usually want ${snackPreference} for QA movie night.`;
}
if (isSnackRecallPrompt(prompt)) {

View file

@ -0,0 +1,74 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createQaBusState } from "./bus-state.js";
import { createQaChannelTransport } from "./qa-channel-transport.js";
vi.mock("node:timers/promises", () => ({
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);
});
});

View file

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

View file

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

View file

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

View file

@ -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: [] });
}
},
);
});

View file

@ -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<ReturnType<typeof startQaBusServer>>;
let gateway: Promise<unknown>;
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<void>();
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<void>();
const releaseFinal = createDeferred<void>();
const acknowledged = createDeferred<void>();
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<string, unknown> = {
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<typeof transport.sendInbound>[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 <T>(check: () => T | Promise<T | undefined> | 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`),
});
},
);
});

View file

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

View file

@ -21,7 +21,6 @@ export type QaScenarioRuntimeEnv<
};
type QaScenarioRuntimeApiDeps = {
sleep: (ms?: number) => Promise<unknown>;
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 {

View file

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

View file

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

View file

@ -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('<active_memory_plugin>'))"
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('<active_memory_plugin>'))"
expr: "pausedRequests.length > 0 && pausedRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes('<active_memory_plugin>'))"
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('<active_memory_plugin>')) ? 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('<active_memory_plugin>'))"
expr: "disabledRequests.length > 0 && disabledRequests.every((request) => !String(request.allInputText ?? '').includes(config.privateFact) && !String(request.allInputText ?? '').includes('<active_memory_plugin>'))"
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')"