From fa56f5e2e9e6f2ec91b0eeae6dec7d9717c6e8ca Mon Sep 17 00:00:00 2001 From: RoboClaw Date: Fri, 2 Oct 2026 18:12:47 -0700 Subject: [PATCH] fix(qa): restore Telegram release validation (#163488) * fix(qa): restore Telegram release checks Co-authored-by: RomneyDa <6581799+RomneyDa@users.noreply.github.com> Co-authored-by: vincentkoc <25068+vincentkoc@users.noreply.github.com> * test(qa): gate Telegram stream completion Co-authored-by: RomneyDa <6581799+RomneyDa@users.noreply.github.com> Co-authored-by: vincentkoc <25068+vincentkoc@users.noreply.github.com> * fix(sessions): canonicalize pending input worker custody * fix(sessions): preserve native pending input custody --------- Co-authored-by: roboclaw-bot <309084314+roboclaw-bot@users.noreply.github.com> Co-authored-by: vincentkoc <25068+vincentkoc@users.noreply.github.com> Co-authored-by: Dallin Romney --- ...n-subagent-self-yield-followup.e2e.test.ts | 18 ++- .../mock-openai/mock-openai-contracts.ts | 10 +- .../mock-openai/mock-openai-events.ts | 47 ++++++- .../mock-openai/mock-openai-responses-http.ts | 45 +++++++ .../mock-openai-responses-websocket.ts | 11 +- .../providers/mock-openai/server-options.ts | 1 + .../server.telegram-policy.test.ts | 101 ++++++++++++++ .../src/providers/mock-openai/server.ts | 47 ++----- .../src/scenario-flow-runner.test-support.ts | 52 +++++--- .../index.js | 13 +- .../telegram-rich-inline-composition.yaml | 25 ++-- ...ssor.pending-inputs-worker-custody.test.ts | 125 ++++++++++++++++++ .../session-accessor.pending-inputs.ts | 2 + .../session-accessor.sqlite-pending-inputs.ts | 6 +- .../session-pending-input-history.test.ts | 1 + 15 files changed, 418 insertions(+), 86 deletions(-) create mode 100644 extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-http.ts create mode 100644 extensions/qa-lab/src/providers/mock-openai/server.telegram-policy.test.ts create mode 100644 src/config/sessions/session-accessor.pending-inputs-worker-custody.test.ts diff --git a/extensions/qa-lab/src/plugin-subagent-self-yield-followup.e2e.test.ts b/extensions/qa-lab/src/plugin-subagent-self-yield-followup.e2e.test.ts index 5ba9211455b0..9d6e480ef704 100644 --- a/extensions/qa-lab/src/plugin-subagent-self-yield-followup.e2e.test.ts +++ b/extensions/qa-lab/src/plugin-subagent-self-yield-followup.e2e.test.ts @@ -162,11 +162,15 @@ describe("plugin subagent sessions_yield follow-up", () => { throw failureContext(error); } - const outbound = state - .getSnapshot() - .messages.filter((message) => message.direction === "outbound"); - // Exactly one announce for the whole continued run: the paused kickoff must - // not announce separately, and the follow-up must not announce twice. + const outbound = await transport.waitForCondition(() => { + const messages = state + .getSnapshot() + .messages.filter((message) => message.direction === "outbound"); + return messages.length >= outboundStartIndex + 2 ? messages : undefined; + }); + // The pause notice and the final completion are distinct requester outcomes. + // Each must arrive once; the continued run must not announce its final twice. + expect(outbound).toHaveLength(outboundStartIndex + 2); expect( outbound.filter((message) => message.text.includes(QA_SUBAGENT_SELF_YIELD_MARKER)), ).toHaveLength(1); @@ -203,7 +207,7 @@ describe("plugin subagent sessions_yield follow-up", () => { request.prompt?.includes("Subagent self yield qa worker") || request.prompt?.includes("Subagent self yield qa remote job finished"), ); - expect(requests).toHaveLength(2); + expect(requests).toHaveLength(3); const verdict = { schemaVersion: 1, scenario: "channel-handoff-adoption", @@ -216,6 +220,7 @@ describe("plugin subagent sessions_yield follow-up", () => { (request) => request.plannedToolName === "sessions_yield", ).length, childModelRequests: handoffRequests.length, + pauseNoticeRequests: requests.length - handoffRequests.length, visibleReplies: outbound.filter((message) => message.text.includes(QA_SUBAGENT_SELF_YIELD_MARKER), ).length, @@ -227,6 +232,7 @@ describe("plugin subagent sessions_yield follow-up", () => { expect(verdict.facts).toEqual({ sessionsYieldCalls: 1, childModelRequests: 2, + pauseNoticeRequests: 1, visibleReplies: 1, duplicateRepliesAfterQuietWindow: 0, duplicateRepliesAfterGatewayRestart: 0, diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts index a131ab61a8fd..85314bc1ecd9 100644 --- a/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-contracts.ts @@ -42,6 +42,7 @@ export type QaMockProviderDispatchResult = { failure?: QaMockProviderFailure; onResponseSent?: () => void; previewPauseMs?: number; + previewPause?: () => Promise; responsePauseMs?: number; }; @@ -464,13 +465,14 @@ export async function writeSse( events: Array, protocol: "responses" | "anthropic", pauseMs?: number, + pause?: () => Promise, ) { const frames = events.map( (event) => `${protocol === "anthropic" ? `event: ${event.type}\n` : ""}data: ${JSON.stringify(event)}\n\n`, ); const completionIndex = - pauseMs === undefined + pauseMs === undefined && pause === undefined ? -1 : events.findIndex((event, index) => isPreviewCompletion(event, events[index - 1])); const body = @@ -485,7 +487,11 @@ export async function writeSse( if (completionIndex >= 0) { // Flush preview deltas before delaying the final text and completion frames. res.write(frames.slice(0, completionIndex).join("")); - await sleep(pauseMs); + if (pause) { + await pause(); + } else { + await sleep(pauseMs); + } } res.end(body); } diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-events.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-events.ts index 0c8912f10354..65a366afca19 100644 --- a/extensions/qa-lab/src/providers/mock-openai/mock-openai-events.ts +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-events.ts @@ -203,13 +203,58 @@ function buildQaLongFinalText({ return `${startMarker}\n${body}\n${endMarker}`; } -export const QA_TELEGRAM_PREPARED_DELIVERY_RE = /Telegram prepared delivery QA: (\{[^\n]+\})/u; +const QA_TELEGRAM_PREPARED_DELIVERY_RE = /Telegram prepared delivery QA: (\{[^\n]+\})/u; +const QA_TELEGRAM_POLICY_HOT_RELOAD_RE = + /^Write (40|12) numbered plain-text lines\. Every line must contain (TG-RELOAD-(?:root|account)-[0-9a-f]{8}(?:-NEXT)?) and the words ((?:hot reload|new policy) keeps this conversation connected)\. Finish with a separate final line containing \2-END\. Do not use tools, Markdown, or explicit reply tags\.$/u; + +function readTelegramPolicyHotReloadPrompt(prompt: string) { + const match = QA_TELEGRAM_POLICY_HOT_RELOAD_RE.exec(prompt); + const lineCount = Number(match?.[1]); + const marker = match?.[2]; + const phrase = match?.[3]; + if (!Number.isSafeInteger(lineCount) || !marker || !phrase) { + return undefined; + } + const isHeldTurn = + lineCount === 40 && !marker.endsWith("-NEXT") && phrase.startsWith("hot reload"); + const isNextTurn = + lineCount === 12 && marker.endsWith("-NEXT") && phrase.startsWith("new policy"); + return isHeldTurn || isNextTurn ? { lineCount, marker, phrase } : undefined; +} + +function buildTelegramPolicyHotReloadText(prompt: string): string | undefined { + const fixture = readTelegramPolicyHotReloadPrompt(prompt); + if (!fixture) { + return undefined; + } + const { lineCount, marker, phrase } = fixture; + return [ + ...Array.from({ length: lineCount }, (_, index) => `${index + 1}. ${marker} ${phrase}`), + `${marker}-END`, + ].join("\n"); +} + +export function resolveTelegramChannelStreamingPause( + prompt: string, +): { previewPauseMs: number } | undefined { + return QA_TELEGRAM_PREPARED_DELIVERY_RE.test(prompt) || + readTelegramPolicyHotReloadPrompt(prompt)?.lineCount === 40 + ? { previewPauseMs: 3_000 } + : undefined; +} export function buildChannelStreamingFixtureEvents(params: { currentPrompt: string; allInputText: string; hasCompletedToolOutput: boolean; }): StreamEvent[] | undefined { + const policyHotReloadText = buildTelegramPolicyHotReloadText(params.currentPrompt); + if (policyHotReloadText) { + return buildStreamingFinalAnswerEvents( + "msg_mock_telegram_policy_hot_reload", + policyHotReloadText, + ); + } if (QA_TELEGRAM_LONG_FINAL_THREE_CHUNK_PROMPT_RE.test(params.allInputText)) { const text = buildQaLongFinalText({ endMarker: "TELEGRAM-LONG-FINAL-3CHUNK-END", diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-http.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-http.ts new file mode 100644 index 000000000000..9e60435811d5 --- /dev/null +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-http.ts @@ -0,0 +1,45 @@ +import type { ServerResponse } from "node:http"; +import { setTimeout as sleep } from "node:timers/promises"; +import { writeJson } from "../shared/http-json.js"; +import { type QaMockProviderDispatchResult, writeSse } from "./mock-openai-contracts.js"; + +export async function writeMockOpenAiResponsesHttp( + res: ServerResponse, + stream: boolean, + dispatched: QaMockProviderDispatchResult, +): Promise { + if (dispatched.failure) { + if (dispatched.failure.retryAfterSeconds !== undefined) { + res.setHeader("retry-after", String(dispatched.failure.retryAfterSeconds)); + } + writeJson(res, dispatched.failure.status, { + error: { + type: dispatched.failure.type, + ...(dispatched.failure.code ? { code: dispatched.failure.code } : {}), + message: dispatched.failure.message, + }, + }); + return; + } + if (dispatched.responsePauseMs !== undefined) { + await sleep(dispatched.responsePauseMs); + } + if (!stream) { + const completion = dispatched.events.at(-1); + if (!completion || completion.type !== "response.completed") { + writeJson(res, 500, { error: "mock completion failed" }); + return; + } + writeJson(res, 200, completion.response); + dispatched.onResponseSent?.(); + return; + } + await writeSse( + res, + dispatched.events, + "responses", + dispatched.previewPauseMs, + dispatched.previewPause, + ); + dispatched.onResponseSent?.(); +} diff --git a/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-websocket.ts b/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-websocket.ts index ce45320bfedf..c8dc96e8bce1 100644 --- a/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-websocket.ts +++ b/extensions/qa-lab/src/providers/mock-openai/mock-openai-responses-websocket.ts @@ -202,8 +202,15 @@ export function attachQaMockResponsesWebSocketServer(params: { }; } for (const [index, event] of events.entries()) { - if (dispatched.previewPauseMs && isPreviewCompletion(event, events[index - 1])) { - await sleep(dispatched.previewPauseMs); + if ( + (dispatched.previewPauseMs || dispatched.previewPause) && + isPreviewCompletion(event, events[index - 1]) + ) { + if (dispatched.previewPause) { + await dispatched.previewPause(); + } else if (dispatched.previewPauseMs) { + await sleep(dispatched.previewPauseMs); + } } sendEvent(event); } diff --git a/extensions/qa-lab/src/providers/mock-openai/server-options.ts b/extensions/qa-lab/src/providers/mock-openai/server-options.ts index 41c28419e3b2..2e42d4b8da4d 100644 --- a/extensions/qa-lab/src/providers/mock-openai/server-options.ts +++ b/extensions/qa-lab/src/providers/mock-openai/server-options.ts @@ -2,6 +2,7 @@ export type QaMockOpenAiServerOptions = { host?: string; port?: number; finalOnlyMarkerPauseMs?: number; + telegramChannelStreamingPause?: () => Promise; modelRefs?: readonly string[]; repeatedRequestResponsePauseMs?: number; repeatedRequestStalledResponsePauseMs?: number; diff --git a/extensions/qa-lab/src/providers/mock-openai/server.telegram-policy.test.ts b/extensions/qa-lab/src/providers/mock-openai/server.telegram-policy.test.ts new file mode 100644 index 000000000000..1b8c5d9a6b7e --- /dev/null +++ b/extensions/qa-lab/src/providers/mock-openai/server.telegram-policy.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it } from "vitest"; +import { + buildChannelStreamingFixtureEvents, + resolveTelegramChannelStreamingPause, +} from "./mock-openai-events.js"; +import { + createMockServerTestHarness, + makeUserInput, + postResponses, +} from "./server.test-harness.js"; + +const { startMockServer } = createMockServerTestHarness(); + +describe("Telegram policy hot-reload mock provider", () => { + it("keeps the held turn active with the requested long marked response", async () => { + let releaseCompletion: (() => void) | undefined; + const completionGate = new Promise((resolve) => { + releaseCompletion = resolve; + }); + const server = await startMockServer({ + telegramChannelStreamingPause: () => completionGate, + }); + const marker = "TG-RELOAD-root-a1b2c3d4"; + const prompt = `Write 40 numbered plain-text lines. Every line must contain ${marker} and the words hot reload keeps this conversation connected. Finish with a separate final line containing ${marker}-END. Do not use tools, Markdown, or explicit reply tags.`; + const expected = [ + ...Array.from( + { length: 40 }, + (_, index) => `${index + 1}. ${marker} hot reload keeps this conversation connected`, + ), + `${marker}-END`, + ].join("\n"); + const response = await postResponses(server, { + model: "gpt-5.6-luna", + stream: true, + input: [makeUserInput(prompt)], + }); + expect(response.status).toBe(200); + expect(response.headers.get("content-type")).toContain("text/event-stream"); + const reader = response.body?.getReader(); + expect(reader).toBeDefined(); + const decoder = new TextDecoder(); + let streamed = ""; + while (!streamed.includes('"type":"response.output_text.delta"')) { + const part = await reader?.read(); + expect(part?.done).toBe(false); + streamed += decoder.decode(part?.value, { stream: true }); + } + expect(streamed).not.toContain('"type":"response.output_text.done"'); + releaseCompletion?.(); + while (!streamed.includes('"type":"response.completed"')) { + const part = await reader?.read(); + if (part?.done) { + break; + } + streamed += decoder.decode(part?.value, { stream: true }); + } + const streamedDeltas = streamed + .split("\n") + .filter((line) => line.startsWith("data: {") && line.endsWith("}")) + .map((line) => JSON.parse(line.slice("data: ".length)) as { type?: string; delta?: string }) + .flatMap((event) => (event.type === "response.output_text.delta" ? [event.delta ?? ""] : [])); + expect(streamedDeltas.join("")).toBe(expected); + + const events = buildChannelStreamingFixtureEvents({ + currentPrompt: prompt, + allInputText: prompt, + hasCompletedToolOutput: false, + }); + expect(events).toBeDefined(); + const deltas = events?.flatMap((event) => + event.type === "response.output_text.delta" ? [event.delta] : [], + ); + const doneIndex = + events?.findIndex((event) => event.type === "response.output_text.done") ?? -1; + const lastDeltaIndex = events?.findLastIndex( + (event) => event.type === "response.output_text.delta", + ); + expect(deltas?.join("")).toBe(expected); + expect(lastDeltaIndex).toBeGreaterThanOrEqual(0); + expect(doneIndex).toBeGreaterThan(lastDeltaIndex ?? -1); + expect(resolveTelegramChannelStreamingPause(prompt)).toEqual({ previewPauseMs: 3_000 }); + }); + + it("does not capture unrelated numbered-line prompts", () => { + const prompts = [ + "Write 40 numbered plain-text lines. Every line must contain OTHER-MARKER and the words hot reload keeps this conversation connected. Finish with a separate final line containing OTHER-MARKER-END. Do not use tools, Markdown, or explicit reply tags.", + "Write 40 numbered plain-text lines. Every line must contain TG-RELOAD-account-a1b2c3d4 and the words new policy keeps this conversation connected. Finish with a separate final line containing TG-RELOAD-account-a1b2c3d4-END. Do not use tools, Markdown, or explicit reply tags.", + "Write 12 numbered plain-text lines. Every line must contain TG-RELOAD-root-a1b2c3d4-NEXT and the words hot reload keeps this conversation connected. Finish with a separate final line containing TG-RELOAD-root-a1b2c3d4-NEXT-END. Do not use tools, Markdown, or explicit reply tags.", + ]; + for (const prompt of prompts) { + expect( + buildChannelStreamingFixtureEvents({ + currentPrompt: prompt, + allInputText: prompt, + hasCompletedToolOutput: false, + }), + ).toBeUndefined(); + expect(resolveTelegramChannelStreamingPause(prompt)).toBeUndefined(); + } + }); +}); diff --git a/extensions/qa-lab/src/providers/mock-openai/server.ts b/extensions/qa-lab/src/providers/mock-openai/server.ts index b331862ccdaa..7ec1171df1f8 100644 --- a/extensions/qa-lab/src/providers/mock-openai/server.ts +++ b/extensions/qa-lab/src/providers/mock-openai/server.ts @@ -143,7 +143,7 @@ import { extractPlannedToolIdentity, splitMockStreamingText, buildChannelStreamingFixtureEvents, - QA_TELEGRAM_PREPARED_DELIVERY_RE, + resolveTelegramChannelStreamingPause, buildAssistantThenToolCallEvents, buildAssistantEvents, buildStreamingFinalAnswerEvents, @@ -180,6 +180,7 @@ import { parseToolOutputJson, } from "./mock-openai-input.js"; import { createMockOpenAiRequestLog } from "./mock-openai-request-log.js"; +import { writeMockOpenAiResponsesHttp } from "./mock-openai-responses-http.js"; import { attachQaMockResponsesWebSocketServer } from "./mock-openai-responses-websocket.js"; import { buildSlackOwnedRequesterEvents, @@ -2221,11 +2222,17 @@ export async function startQaMockOpenAiServer(params?: QaMockOpenAiServerOptions } : {}), ...(failure ? { failure } : {}), - ...(QA_TELEGRAM_PREPARED_DELIVERY_RE.test(splitMockConversationContext(prompt).current) - ? { previewPauseMs: 3_000 } - : QA_FINAL_ONLY_MARKER_STREAMING_PROMPT_RE.test(allInputText) + ...((() => { + const telegramPause = resolveTelegramChannelStreamingPause( + splitMockConversationContext(prompt).current, + ); + return telegramPause && params?.telegramChannelStreamingPause + ? { previewPause: params.telegramChannelStreamingPause } + : telegramPause; + })() ?? + (QA_FINAL_ONLY_MARKER_STREAMING_PROMPT_RE.test(allInputText) ? { previewPauseMs: finalOnlyMarkerPauseMs } - : {}), + : {})), // Stall one request; later failures let the normal retry budget settle the turn. ...(repeatedRequestRecovery && scenarioState.repeatedRequestRecoveryAttempts <= QA_REPEATED_REQUEST_STALL_ATTEMPT @@ -2361,35 +2368,7 @@ export async function startQaMockOpenAiServer(params?: QaMockOpenAiServerOptions } if (url.pathname === "/v1/responses") { const dispatched = await dispatchResponses({ body, raw, headers: req.headers }); - if (dispatched.failure) { - if (dispatched.failure.retryAfterSeconds !== undefined) { - res.setHeader("retry-after", String(dispatched.failure.retryAfterSeconds)); - } - writeJson(res, dispatched.failure.status, { - error: { - type: dispatched.failure.type, - ...(dispatched.failure.code ? { code: dispatched.failure.code } : {}), - message: dispatched.failure.message, - }, - }); - return; - } - const { events } = dispatched; - if (dispatched.responsePauseMs !== undefined) { - await sleep(dispatched.responsePauseMs); - } - if (body.stream !== true) { - const completion = events.at(-1); - if (!completion || completion.type !== "response.completed") { - writeJson(res, 500, { error: "mock completion failed" }); - return; - } - writeJson(res, 200, completion.response); - dispatched.onResponseSent?.(); - return; - } - await writeSse(res, events, "responses", dispatched.previewPauseMs); - dispatched.onResponseSent?.(); + await writeMockOpenAiResponsesHttp(res, body.stream === true, dispatched); return; } const dispatched = await dispatchProvider({ diff --git a/extensions/qa-lab/src/scenario-flow-runner.test-support.ts b/extensions/qa-lab/src/scenario-flow-runner.test-support.ts index 33d5b5103603..fc6bd55cefd9 100644 --- a/extensions/qa-lab/src/scenario-flow-runner.test-support.ts +++ b/extensions/qa-lab/src/scenario-flow-runner.test-support.ts @@ -373,8 +373,37 @@ export async function assertTelegramRichObservationFlow( providerMode: "mock-openai", cfg: { channels: { telegram: { accounts: { sut: { richMessages: true } } } } }, gateway: { - call: async (method: string, args: { message: string }) => { + call: async ( + method: string, + args: { + message?: string; + action?: string; + params?: { messageId?: number; content?: string }; + }, + ) => { + if (method === "message.action") { + assert.equal(args.action, "edit"); + assert.equal(args.params?.messageId, 1); + edits += 1; + const firstMarker = [...markers][0]; + assert.ok(firstMarker !== undefined); + const marker = readMarker(args.params?.content); + deliver("edit", () => { + observed[0] = observation( + testCase === "wrong-edit-id" ? Number(args.params?.messageId) + 40 : 1, + list, + testCase === "wrong-edit-kind" ? "message" : "edit", + testCase === "wrong-edit-marker" + ? marker + "-wrong" + : testCase === "stale-edit" + ? firstMarker + : marker, + ); + }); + return { ok: true }; + } assert.equal(method, "send"); + assert.ok(args.message); // Telegram extracts Markdown images before dispatching to its renderer. const [planned] = createOutboundPayloadPlan([{ text: args.message }], { extractMarkdownImages: true, @@ -449,27 +478,6 @@ export async function assertTelegramRichObservationFlow( await joined; } }, - runQaCli: async (_env: unknown, args: string[]) => { - assert.deepEqual(args.slice(0, 2), ["message", "edit"]); - const id = Number(args[args.indexOf("--message-id") + 1]); - assert.equal(id, 1); - edits += 1; - const firstMarker = [...markers][0]; - assert.ok(firstMarker !== undefined); - const marker = readMarker(args[args.indexOf("--message") + 1]); - deliver("edit", () => { - observed[0] = observation( - testCase === "wrong-edit-id" ? id + 40 : id, - list, - testCase === "wrong-edit-kind" ? "message" : "edit", - testCase === "wrong-edit-marker" - ? marker + "-wrong" - : testCase === "stale-edit" - ? firstMarker - : marker, - ); - }); - }, runAgentPrompt: () => { throw new Error("direct delivery must not invoke a model"); }, diff --git a/extensions/qa-lab/test-fixtures/self-yield-followup-subagent-plugin/index.js b/extensions/qa-lab/test-fixtures/self-yield-followup-subagent-plugin/index.js index 71fb0c1fd6d4..c52b9928aed9 100644 --- a/extensions/qa-lab/test-fixtures/self-yield-followup-subagent-plugin/index.js +++ b/extensions/qa-lab/test-fixtures/self-yield-followup-subagent-plugin/index.js @@ -20,7 +20,6 @@ function getState() { kickoffRunId: undefined, resolveFinalReply, yieldEntered: false, - releaseYield: undefined, }); } @@ -49,15 +48,11 @@ function writeJson(res, statusCode, body) { export default { id: "qa-self-yield-followup-subagent", register(api) { - api.on("before_tool_call", async (event) => { + api.on("after_tool_call", (event) => { if (event.toolName !== "sessions_yield") { return; } - const state = getState(); - state.yieldEntered = true; - await new Promise((resolve) => { - state.releaseYield = resolve; - }); + getState().yieldEntered = true; }); api.on("before_dispatch", async (event) => { @@ -150,12 +145,10 @@ export default { gatewayRuntimeScopeSurface: "trusted-operator", async handler(_req, res) { const state = getState(); - if (!state.followUpRunId || !state.releaseYield) { + if (!state.followUpRunId) { writeJson(res, 409, { ok: false, error: "follow-up was not queued" }); return true; } - state.releaseYield(); - state.releaseYield = undefined; const terminal = await api.runtime.subagent.waitForRun({ runId: state.followUpRunId, timeoutMs: 90_000, diff --git a/qa/scenarios/channels/telegram-rich-inline-composition.yaml b/qa/scenarios/channels/telegram-rich-inline-composition.yaml index f8766e847813..ecb9f33a861a 100644 --- a/qa/scenarios/channels/telegram-rich-inline-composition.yaml +++ b/qa/scenarios/channels/telegram-rich-inline-composition.yaml @@ -37,7 +37,7 @@ scenario: isolationReason: Rich-message account configuration and native message revisions belong to this leased DM. transportPolicy: directMessageOnly: true - summary: Send eleven synthetic formatting cases through the Gateway and edit one through the normal CLI; inspect full TDLib trees from the leased Test Server DM. This proves direct operator delivery, not model generation. The CLI harness and package candidate must contain the same reviewed renderer. + summary: Send eleven synthetic formatting cases and edit one through the Gateway action boundary; inspect full TDLib trees from the leased Test Server DM. This proves direct operator delivery, not model generation. config: requiredProviderMode: mock-openai requiredChannelDriver: live @@ -129,7 +129,7 @@ flow: value: [] - set: failures value: [] - - set: editCliCompleted + - set: editCompleted value: false - set: editVerified value: false @@ -176,7 +176,7 @@ flow: receipts: vars.receipts.slice(0, 11).map((receipt, index) => ({ case: receipt.case, label: `message-${index + 1}` })), cases: vars.proofs.slice(0, 11).map(proof => ({ case: proof.case, label: proof.label, accountMessageIdsDiffer: proof.accountMessageIdsDiffer, richMessage: tree(proof.richMessage) })), failures: vars.failures.map(failure => ({ case: failure.case, label: failure.label, contentType: failure.contentType, full: failure.full, matchesPattern: failure.matchesPattern, richMessage: tree(failure.richMessage) })), - edit: { completedWithBotReceipt: vars.editCliCompleted, verified: vars.editVerified }, + edit: { completedWithBotReceipt: vars.editCompleted, verified: vars.editVerified }, candidates, unmatched: { count: unmatchedCount, contentTypes: [...unmatchedTypes] } }; const content = redactQaGatewayDebugText(JSON.stringify(artifact).replaceAll(runId, 'run')); @@ -243,13 +243,22 @@ flow: - set: editMarker value: expr: "`QA-RICH-${runId}-edit-END`" - - call: runQaCli + - call: env.gateway.call args: - - ref: env - - expr: "['message', 'edit', '--channel', 'telegram', '--account', transport.accountId, '--target', delivery.to, '--message-id', String(editTarget.botMessageId), '--message', `${fixtures[1].markdown}\n\n${editMarker}`, '--json']" + - message.action + - channel: telegram + accountId: { ref: transport.accountId } + action: edit + params: + chatId: { ref: delivery.to } + messageId: + expr: "Number(editTarget.botMessageId)" + content: + expr: "`${fixtures[1].markdown}\n\n${editMarker}`" + idempotencyKey: + expr: "`${runId}:edit`" - timeoutMs: 60000 - json: true - - set: editCliCompleted + - set: editCompleted value: true - call: waitForCondition saveAs: edited diff --git a/src/config/sessions/session-accessor.pending-inputs-worker-custody.test.ts b/src/config/sessions/session-accessor.pending-inputs-worker-custody.test.ts new file mode 100644 index 000000000000..356cf9661f01 --- /dev/null +++ b/src/config/sessions/session-accessor.pending-inputs-worker-custody.test.ts @@ -0,0 +1,125 @@ +import fs from "node:fs"; +import path from "node:path"; +import { afterEach, describe, expect, it } from "vitest"; +import type { PersistedUserTurnMessage } from "../../sessions/user-turn-transcript.types.js"; +import { + closeOpenClawAgentDatabasesForTest, + openOpenClawAgentDatabase, +} from "../../state/openclaw-agent-db.js"; +import { + appendTranscriptMessageSync, + loadTranscriptEvents, + upsertSessionEntryCore, +} from "./session-accessor.js"; +import { + listSessionPendingInputs, + stageSessionPendingInput, + type SessionPendingInputReceipt, +} from "./session-accessor.pending-inputs.js"; +import { + captureSessionPendingInputWorkerCustody, + runWithSessionPendingInputWorkerCustody, +} from "./session-accessor.sqlite-pending-inputs.js"; +import { resolveSqliteScope, toDatabaseOptions } from "./session-accessor.sqlite-scope.js"; +import { useTempSessionsFixture } from "./test-helpers.js"; + +describe("accepted input worker custody", () => { + const fixture = useTempSessionsFixture("openclaw-pending-worker-custody-"); + let receipt: SessionPendingInputReceipt | undefined; + + afterEach(() => { + receipt?.finish("interrupted"); + receipt = undefined; + closeOpenClawAgentDatabasesForTest(); + }); + + it("appends, consumes, and finishes worker custody across a state-directory alias", async () => { + const fixtureRoot = path.resolve(fixture.sessionsDir(), "../../.."); + const aliasRoot = path.join(fixtureRoot, "state-alias"); + fs.symlinkSync(fixtureRoot, aliasRoot, process.platform === "win32" ? "junction" : "dir"); + const scope = { + agentId: "alias-agent", + env: { OPENCLAW_STATE_DIR: aliasRoot }, + sessionId: "alias-session", + sessionKey: "agent:alias-agent:pending-inputs", + }; + await upsertSessionEntryCore(scope, { sessionId: scope.sessionId, updatedAt: 1 }); + const message: PersistedUserTurnMessage = { + role: "user", + content: "Continue through worker custody", + timestamp: 100, + idempotencyKey: "worker-alias:user", + }; + receipt = await stageSessionPendingInput(scope, { + runId: "worker-alias", + message, + assertCurrent: () => {}, + }); + if (!receipt) { + throw new Error("Expected aliased pending input custody"); + } + + const custody = receipt.run(() => captureSessionPendingInputWorkerCustody()); + if (!custody) { + throw new Error("Expected captured worker custody"); + } + const database = openOpenClawAgentDatabase(toDatabaseOptions(resolveSqliteScope(scope))); + expect(custody.facts.databasePath).toBe(fs.realpathSync(database.path)); + + const workerScope = { ...scope, storePath: custody.facts.databasePath }; + const result = runWithSessionPendingInputWorkerCustody( + custody.facts, + custody.relocation, + custody.assertCurrent, + () => appendTranscriptMessageSync(workerScope, { message: receipt!.message }), + ); + expect(result.value).toMatchObject({ ok: true, value: { appended: true } }); + custody.publish(result.receipt); + receipt.finish("cancelled"); + receipt = undefined; + + expect(await loadTranscriptEvents(scope)).toContainEqual( + expect.objectContaining({ message: expect.objectContaining({ content: message.content }) }), + ); + expect(await listSessionPendingInputs(scope)).toMatchObject({ items: [], total: 0 }); + }); + + it("appends and finishes pending input through the native incognito owner", async () => { + const scope = { + agentId: "incognito-agent", + env: { OPENCLAW_STATE_DIR: path.resolve(fixture.sessionsDir(), "../../..") }, + sessionId: "incognito-session", + sessionKey: "agent:incognito-agent:dashboard:incognito-pending-input", + }; + await upsertSessionEntryCore(scope, { + incognito: true, + sessionId: scope.sessionId, + updatedAt: 1, + }); + const message: PersistedUserTurnMessage = { + role: "user", + content: "Continue in memory", + timestamp: 100, + idempotencyKey: "incognito-native:user", + }; + receipt = await stageSessionPendingInput(scope, { + runId: "incognito-native", + message, + assertCurrent: () => {}, + }); + if (!receipt) { + throw new Error("Expected incognito pending input custody"); + } + + expect( + receipt.run(() => appendTranscriptMessageSync(scope, { message: receipt!.message })), + ).toMatchObject({ ok: true, value: { appended: true } }); + receipt.finish("cancelled"); + receipt = undefined; + + expect(await loadTranscriptEvents(scope)).toContainEqual( + expect.objectContaining({ message: expect.objectContaining({ content: message.content }) }), + ); + expect(await listSessionPendingInputs(scope)).toMatchObject({ items: [], total: 0 }); + }); +}); diff --git a/src/config/sessions/session-accessor.pending-inputs.ts b/src/config/sessions/session-accessor.pending-inputs.ts index 09443abde33b..3dececaf2fe7 100644 --- a/src/config/sessions/session-accessor.pending-inputs.ts +++ b/src/config/sessions/session-accessor.pending-inputs.ts @@ -406,6 +406,7 @@ export async function stageSessionPendingInput( path: current.path, databaseIdentity: physical.identity, databaseBirthtime: physical.birthtime, + workerDatabasePath: physical.canonicalPath || current.path, }; }, databaseOptions, @@ -420,6 +421,7 @@ export async function stageSessionPendingInput( sessionId: scope.sessionId, sessionKey: resolved.sessionKey, databasePath: source.path, + workerDatabasePath: source.workerDatabasePath, idempotencyKey, lifecycleGeneration, messageJson, diff --git a/src/config/sessions/session-accessor.sqlite-pending-inputs.ts b/src/config/sessions/session-accessor.sqlite-pending-inputs.ts index b31e1352011a..487e063ac7f7 100644 --- a/src/config/sessions/session-accessor.sqlite-pending-inputs.ts +++ b/src/config/sessions/session-accessor.sqlite-pending-inputs.ts @@ -50,7 +50,10 @@ export type SessionPendingInputOwner = { transcriptInputId: string; sessionId: string; sessionKey: string; + /** Native cache locator; may be the process-held incognito sentinel. */ databasePath: string; + /** Prepared physical locator serialized only to a database worker. */ + workerDatabasePath: string; idempotencyKey: string; lifecycleGeneration: string; messageJson: string; @@ -101,7 +104,7 @@ export function captureSessionPendingInputWorkerCustody() { transcriptInputId: current.transcriptInputId, sessionId: current.sessionId, sessionKey: current.sessionKey, - databasePath: current.databasePath, + databasePath: current.workerDatabasePath, idempotencyKey: current.idempotencyKey, lifecycleGeneration: current.lifecycleGeneration, messageJson: current.messageJson, @@ -133,6 +136,7 @@ export function runWithSessionPendingInputWorkerCustody( ): { value: T; receipt: SessionPendingInputWorkerReceipt } { const hydrate = (current: SessionPendingInputWorkerFacts): SessionPendingInputOwner => ({ ...current, + workerDatabasePath: current.databasePath, sources: current.sources?.map(hydrate), assertCurrent, finish: () => { diff --git a/src/config/sessions/session-pending-input-history.test.ts b/src/config/sessions/session-pending-input-history.test.ts index c2cf28491bf6..68f35feef002 100644 --- a/src/config/sessions/session-pending-input-history.test.ts +++ b/src/config/sessions/session-pending-input-history.test.ts @@ -138,6 +138,7 @@ it.each(["transaction", "commit"] as const)( sessionId: scope.sessionId, sessionKey: scope.sessionKey, databasePath: database.path, + workerDatabasePath: database.path, idempotencyKey: "late:user", lifecycleGeneration: getAgentEventLifecycleGeneration(), messageJson: JSON.stringify(receipt.message),