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 <dallinromney@gmail.com>
This commit is contained in:
RoboClaw 2026-10-02 18:12:47 -07:00 • committed by GitHub
parent 8c6ace2d7d
commit fa56f5e2e9
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
15 changed files with 418 additions and 86 deletions

View file

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

View file

@ -42,6 +42,7 @@ export type QaMockProviderDispatchResult = {
failure?: QaMockProviderFailure;
onResponseSent?: () => void;
previewPauseMs?: number;
previewPause?: () => Promise<void>;
responsePauseMs?: number;
};
@ -464,13 +465,14 @@ export async function writeSse(
events: Array<StreamEvent | AnthropicStreamEvent>,
protocol: "responses" | "anthropic",
pauseMs?: number,
pause?: () => Promise<void>,
) {
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);
}

View file

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

View file

@ -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<void> {
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?.();
}

View file

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

View file

@ -2,6 +2,7 @@ export type QaMockOpenAiServerOptions = {
host?: string;
port?: number;
finalOnlyMarkerPauseMs?: number;
telegramChannelStreamingPause?: () => Promise<void>;
modelRefs?: readonly string[];
repeatedRequestResponsePauseMs?: number;
repeatedRequestStalledResponsePauseMs?: number;

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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<T>(
): { value: T; receipt: SessionPendingInputWorkerReceipt } {
const hydrate = (current: SessionPendingInputWorkerFacts): SessionPendingInputOwner => ({
...current,
workerDatabasePath: current.databasePath,
sources: current.sources?.map(hydrate),
assertCurrent,
finish: () => {

View file

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