mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
* fix(subagents): preserve completion during concurrent registry writes Replace speculative row staging with per-row FIFO planning, immutable post-ACK publication, digest CAS, and bounded conflict replanning. Route registry and native completion writes through the same owner, and commit eligible terminal signals with their terminal rows. Retain runtime custody across metadata publication and acknowledged collector rekeys. Fence retired Gateway incarnations and unknown writes until canonical reconciliation. This local checkpoint has independent review clean through P2, core and selected test type checks, typed caller lint, and focused worker/lifecycle proof. Main integration and complete validation remain in progress; descriptorless recovery fixture setup and warm restoration failures remain recorded in the task notes for the next checkpoint. * fix(subagents): preserve reset and stop ownership after publication Prepare reset revocation after session generation checks and revalidate current child captures before commit. Reacquire queued Stop rows through their runtime owner after terminal publication. Adapt remaining fixtures to immutable registry observations and awaited acknowledgements. * fix(subagents): retain cancellation and completion custody across acknowledgements Preserve semantic kill ownership during replanning, allow confirmation of the same cancellation, and retain grace callbacks through their own terminal acknowledgement. Keep compact projections noncustodial and accept current runtime owners in Stop. Complete immutable-row fixture cutovers and join acknowledged terminal work during teardown. * test(subagents): assert committed cancellation failure boundaries * fix(agents): retain retirement receipts and lifecycle attempt custody Keep acknowledged quiet deletion receipts through publication recovery and bind delayed terminal grace to its admitted execution attempt. Adapt immutable publication fixtures and requester reply proof to the integrated yield policy. * refactor(agents): centralize canonical subagent restoration Keep restore admission, immutable row installation, and uncertain write recovery in the registry mutation owner. Retain cache and notification ownership in state without the reverse import, and migrate all callers and fixtures to the owning contracts. * fix(agents): preserve child ownership during collector settlement Resolve collector usage through the existing child-session owner before terminal publication. Cover recorded child agents and configured defaults for unscoped session keys through real queued registration failure settlement. * fix(subagents): retain Gateway custody after lifecycle commits Return the acknowledged published row from lifecycle mutations instead of their private planning draft. The draft has no Gateway binding, so cleanup could leave an owned managed-worktree session behind after falling back to raw RPC. The real production-boundary deletion case failed before this fix and now deletes the session/worktree while preserving its result snapshot. Separate raw row hashing and mutation contracts from persisted record types to remove static import cycles. Finish the cutover by removing obsolete exports, eliminate a redundant assertion, and give sweep deletion a unique export name. Keep fixture checks on actual published rows and provide the scheduler logger's trace method. No test guard, budget, baseline, or schema change is needed. Final focused proof: five affected files / 311 cases, including all 25 spawn production-boundary cases. Additional unchanged-contract CI repair tests retain their passing replays. Core and all 27 affected test graphs, typed lint, architecture/dependency and assertion guards, formatting, diff checks, and isolated review through P2 pass. * test(agents): align registry proof with owner lifetimes Keep requester replay guards until isolated database teardown, avoiding schema writes while worker cleanup is active. Update benchmark imports to the extracted registry contract and row owners.
843 lines
32 KiB
TypeScript
843 lines
32 KiB
TypeScript
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
|
|
import type { AddressInfo } from "node:net";
|
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
|
import type { ChatEvent } from "../packages/gateway-protocol/src/schema/logs-chat.js";
|
|
import { saveSubagentRegistryToSqlite } from "../src/agents/subagents/registry/subagent-registry-state.fixture.test-support.js";
|
|
import {
|
|
settleSubagentRegistryPersistenceWork,
|
|
writeSubagentSessionEntry,
|
|
} from "../src/agents/subagents/registry/subagent-registry.persistence.test-support.js";
|
|
import { loadSubagentRegistryFromSqlite } from "../src/agents/subagents/registry/subagent-registry.store.sqlite.js";
|
|
import type { SubagentRunRecord } from "../src/agents/subagents/registry/subagent-registry.types.js";
|
|
import { getSessionKysely } from "../src/config/sessions/session-accessor.sqlite-scope.js";
|
|
import type { OpenClawConfig } from "../src/config/types.openclaw.js";
|
|
import { connectGatewayClient, disconnectGatewayClient } from "../src/gateway/test-helpers.e2e.js";
|
|
import { executeSqliteQuerySync } from "../src/infra/kysely-sync.js";
|
|
import { extractFirstTextBlock } from "../src/shared/chat-message-content.js";
|
|
import { withOpenClawAgentDatabaseReadOnly } from "../src/state/openclaw-agent-db-readonly.js";
|
|
import { closeOpenClawStateDatabaseForTest } from "../src/state/openclaw-state-db.js";
|
|
import {
|
|
writeOpenAiResponsesSse,
|
|
writeOpenAiResponsesText,
|
|
} from "./helpers/openai-responses-sse.js";
|
|
import {
|
|
createOpenClawTestInstance,
|
|
type OpenClawTestInstance,
|
|
} from "./helpers/openclaw-test-instance.js";
|
|
import { createDeferred, withTestTimeout } from "./helpers/promise.js";
|
|
import { runQaGatewayFixture } from "./helpers/qa-gateway-cleanup.js";
|
|
|
|
const TEST_TIMEOUT_MS = 180_000;
|
|
const MODEL_REF = "requester-owner/synthetic";
|
|
const REQUESTER_KEY = "requester-owner-requester";
|
|
const REQUESTER_AGENT_ID = "beta";
|
|
const OTHER_AGENT_ID = "alpha";
|
|
const PARENT_PROMPT = "REQUESTER-OWNER parent: spawn one worker and finish without waiting.";
|
|
const CHILD_TASK = "REQUESTER-OWNER child task: reply with the agreed child token.";
|
|
const CHILD_MARKER = "REQUESTER-OWNER-CHILD-OK";
|
|
const ANNOUNCE_FAILURE_MARKER = "Subagent announce failed";
|
|
const RESTORED_RUN_ID = "run-requester-owner-legacy";
|
|
const RESTORED_REQUESTER_KEY = "requester-owner-legacy-requester";
|
|
const RESTORED_CHILD_RESULT = "REQUESTER-OWNER-LEGACY-CHILD-RESULT";
|
|
|
|
type SseEvent = Record<string, unknown>;
|
|
|
|
type ProofModelServer = {
|
|
bodies: () => readonly string[];
|
|
close: () => Promise<void>;
|
|
countRequestsContaining: (marker: string) => number;
|
|
completionResponseCount: () => number;
|
|
requestCount: () => number;
|
|
url: string;
|
|
};
|
|
|
|
const instances: OpenClawTestInstance[] = [];
|
|
const modelServers: ProofModelServer[] = [];
|
|
|
|
afterEach(async () => {
|
|
const results = await Promise.allSettled([
|
|
...instances.splice(0).map((instance) => instance.cleanup()),
|
|
...modelServers.splice(0).map((server) => server.close()),
|
|
]);
|
|
const errors = results.flatMap((result) => (result.status === "rejected" ? [result.reason] : []));
|
|
if (errors.length > 0) {
|
|
throw new AggregateError(errors, "Requester fixture cleanup failed");
|
|
}
|
|
});
|
|
|
|
describe("REQUESTER-OWNER requester agent id survives completion dispatch", () => {
|
|
it(
|
|
"answers status through WebSocket while parent and child provider requests are held",
|
|
{ timeout: TEST_TIMEOUT_MS },
|
|
async () => {
|
|
const parentGate = createDeferred();
|
|
const childGate = createDeferred();
|
|
const modelServer = await startProofModelServer({
|
|
yieldAfterSpawn: parentGate.promise,
|
|
childReply: childGate.promise,
|
|
});
|
|
modelServers.push(modelServer);
|
|
const instance = await createOpenClawTestInstance({
|
|
name: "requester-owner-busy-status",
|
|
config: createTestConfig(modelServer.url),
|
|
env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined },
|
|
});
|
|
instances.push(instance);
|
|
instance.state.applyEnv();
|
|
await instance.startGateway();
|
|
const sessionKey = `agent:${REQUESTER_AGENT_ID}:${REQUESTER_KEY}`;
|
|
const statusRunId = "busy-parent-status";
|
|
const statusReply = createDeferred<ChatEvent>();
|
|
const client = await connectGatewayClient({
|
|
url: instance.url,
|
|
token: instance.gatewayToken,
|
|
onEvent: (event) => {
|
|
if (event.event !== "chat") {
|
|
return;
|
|
}
|
|
const chat = event.payload as ChatEvent;
|
|
if (
|
|
chat.runId === statusRunId &&
|
|
(chat.state === "final" || chat.state === "error" || chat.state === "aborted")
|
|
) {
|
|
statusReply.resolve(chat);
|
|
}
|
|
},
|
|
});
|
|
let parent: Promise<unknown> | undefined;
|
|
let parentSettled = false;
|
|
try {
|
|
parent = client.request(
|
|
"agent",
|
|
{
|
|
sessionKey: REQUESTER_KEY,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
idempotencyKey: "busy-status-parent-turn",
|
|
message: PARENT_PROMPT,
|
|
deliver: false,
|
|
},
|
|
{ expectFinal: true },
|
|
);
|
|
void parent.then(
|
|
() => {
|
|
parentSettled = true;
|
|
},
|
|
() => {
|
|
parentSettled = true;
|
|
},
|
|
);
|
|
await vi.waitFor(
|
|
() => {
|
|
expect(modelServer.requestCount(), instance.logs()).toBe(3);
|
|
expect(
|
|
modelServer
|
|
.bodies()
|
|
.filter(
|
|
(body) => body.includes(PARENT_PROMPT) && body.includes("function_call_output"),
|
|
),
|
|
).toHaveLength(1);
|
|
expect(
|
|
modelServer
|
|
.bodies()
|
|
.filter(
|
|
(body) => body.includes(CHILD_TASK) && !body.includes("function_call_output"),
|
|
),
|
|
).toHaveLength(1);
|
|
},
|
|
{ interval: 50, timeout: 60_000 },
|
|
);
|
|
const runs = [...loadSubagentRegistryFromSqlite().values()];
|
|
expect(runs, instance.logs()).toHaveLength(1);
|
|
const run = runs[0]!;
|
|
expect(run).toMatchObject({
|
|
requesterAgentId: REQUESTER_AGENT_ID,
|
|
requesterSessionKey: sessionKey,
|
|
requesterTurnRunId: "busy-status-parent-turn",
|
|
completionTarget: "parent",
|
|
});
|
|
expect(run.childSessionKey).toMatch(/^agent:beta:subagent:/);
|
|
expect(run.execution.status).toBe("running");
|
|
expect(run.task).toBe(CHILD_TASK);
|
|
expect(runs.filter((entry) => entry.requesterAgentId === OTHER_AGENT_ID)).toEqual([]);
|
|
|
|
const requestsBeforeStatus = modelServer.requestCount();
|
|
expect(
|
|
await client.request("chat.send", {
|
|
sessionKey,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
message: "/status",
|
|
idempotencyKey: statusRunId,
|
|
}),
|
|
).toMatchObject({ runId: statusRunId });
|
|
const reply = await withTestTimeout(
|
|
statusReply.promise,
|
|
30_000,
|
|
"status did not finish while the parent provider was held",
|
|
);
|
|
expect(reply, instance.logs()).toMatchObject({
|
|
state: "final",
|
|
runId: statusRunId,
|
|
sessionKey,
|
|
});
|
|
const statusText = "message" in reply ? extractFirstTextBlock(reply.message) : undefined;
|
|
expect(statusText).toBeTruthy();
|
|
expect(statusText).not.toContain(CHILD_MARKER);
|
|
expect(parentSettled).toBe(false);
|
|
expect(modelServer.requestCount()).toBe(requestsBeforeStatus);
|
|
|
|
childGate.resolve();
|
|
await vi.waitFor(
|
|
() => {
|
|
expect(loadSubagentRegistryFromSqlite().get(run.runId), instance.logs()).toMatchObject({
|
|
execution: { status: "terminal", outcome: { status: "ok" } },
|
|
completion: { resultText: CHILD_MARKER },
|
|
delivery: { status: "pending" },
|
|
});
|
|
},
|
|
{ interval: 50, timeout: 30_000 },
|
|
);
|
|
expect(parentSettled).toBe(false);
|
|
expect(modelServer.requestCount()).toBe(requestsBeforeStatus);
|
|
expect(modelServer.countRequestsContaining(CHILD_MARKER)).toBe(0);
|
|
|
|
parentGate.resolve();
|
|
expect(await parent, instance.logs()).toMatchObject({ status: "ok" });
|
|
await vi.waitFor(
|
|
async () => {
|
|
expect(loadSubagentRegistryFromSqlite().get(run.runId), instance.logs()).toMatchObject({
|
|
execution: { status: "terminal", outcome: { status: "ok" } },
|
|
delivery: { status: "delivered" },
|
|
});
|
|
},
|
|
{ interval: 50, timeout: 30_000 },
|
|
);
|
|
const history = await client.request<{
|
|
messages: Array<{ role?: string; content?: unknown }>;
|
|
}>("chat.history", { sessionKey, agentId: REQUESTER_AGENT_ID, limit: 30 });
|
|
const visible = history.messages.map((message) => ({
|
|
role: message.role,
|
|
text: extractFirstTextBlock(message),
|
|
}));
|
|
expect(
|
|
visible.filter(({ role, text }) => role === "user" && text === "/status"),
|
|
).toHaveLength(1);
|
|
expect(
|
|
visible.filter(({ role, text }) => role === "assistant" && text === statusText),
|
|
).toHaveLength(1);
|
|
const statusIndex = visible.findIndex(
|
|
({ role, text }) => role === "user" && text === "/status",
|
|
);
|
|
expect(visible[statusIndex + 1]).toEqual({ role: "assistant", text: statusText });
|
|
expect(
|
|
visible.filter(({ role, text }) => role === "assistant" && text === CHILD_MARKER),
|
|
).toHaveLength(1);
|
|
expect(JSON.stringify(history.messages)).not.toContain("This turn ended before a reply");
|
|
const laterModelText = modelServer
|
|
.bodies()
|
|
.slice(requestsBeforeStatus)
|
|
.flatMap((body) => {
|
|
const request = JSON.parse(body) as {
|
|
input: Array<{ role?: string; content?: unknown }>;
|
|
};
|
|
return request.input
|
|
.filter(({ role }) => role === "user" || role === "assistant")
|
|
.map(extractFirstTextBlock);
|
|
});
|
|
expect(laterModelText.some((text) => text?.includes(PARENT_PROMPT))).toBe(true);
|
|
expect(laterModelText).not.toContain("/status");
|
|
expect(laterModelText).not.toContain(statusText);
|
|
} finally {
|
|
childGate.resolve();
|
|
parentGate.resolve();
|
|
await runQaGatewayFixture(
|
|
async () => {
|
|
await withTestTimeout(
|
|
Promise.allSettled(parent ? [parent] : []),
|
|
30_000,
|
|
"parent request did not settle after releasing provider gates",
|
|
);
|
|
},
|
|
() => disconnectGatewayClient(client),
|
|
() => instance.stopGateway(),
|
|
() => closeOpenClawStateDatabaseForTest(),
|
|
);
|
|
}
|
|
},
|
|
);
|
|
|
|
it.each([undefined, { kind: "local" } as const])(
|
|
"resumes a yielded requester once after its private child finishes (placement: %j)",
|
|
{ timeout: TEST_TIMEOUT_MS },
|
|
async (placement) => {
|
|
const yieldGate = createDeferred();
|
|
const modelServer = await startProofModelServer({
|
|
yieldAfterSpawn: yieldGate.promise,
|
|
placement,
|
|
});
|
|
modelServers.push(modelServer);
|
|
const instance = await createOpenClawTestInstance({
|
|
name: "private-completion-before-yield",
|
|
config: createTestConfig(modelServer.url),
|
|
env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined },
|
|
});
|
|
instances.push(instance);
|
|
instance.state.applyEnv();
|
|
const sessionId = "private-yield-requester-session";
|
|
const sessionKey = `agent:${REQUESTER_AGENT_ID}:${REQUESTER_KEY}`;
|
|
await writeSubagentSessionEntry({
|
|
stateDir: instance.stateDir,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
sessionKey,
|
|
sessionId,
|
|
defaultSessionId: sessionId,
|
|
});
|
|
closeOpenClawStateDatabaseForTest();
|
|
await instance.startGateway();
|
|
const chatErrors: unknown[] = [];
|
|
const client = await connectGatewayClient({
|
|
url: instance.url,
|
|
token: instance.gatewayToken,
|
|
onEvent: (event) => {
|
|
if (event.event === "chat" && (event.payload as { state?: string })?.state === "error") {
|
|
chatErrors.push(event.payload);
|
|
}
|
|
},
|
|
});
|
|
const readInputs = () =>
|
|
withOpenClawAgentDatabaseReadOnly(
|
|
({ db }) =>
|
|
executeSqliteQuerySync(
|
|
db,
|
|
getSessionKysely(db)
|
|
.selectFrom("session_pending_inputs")
|
|
.selectAll()
|
|
.where("session_id", "=", sessionId),
|
|
).rows,
|
|
{ agentId: REQUESTER_AGENT_ID },
|
|
);
|
|
try {
|
|
const parent = client.request(
|
|
"agent",
|
|
{
|
|
sessionKey,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
idempotencyKey: "private-yield-parent-turn",
|
|
message: PARENT_PROMPT,
|
|
deliver: false,
|
|
},
|
|
{ expectFinal: true },
|
|
);
|
|
void parent.catch(() => {});
|
|
instance.state.applyEnv();
|
|
await vi.waitFor(
|
|
() => {
|
|
const runs = [...loadSubagentRegistryFromSqlite().values()];
|
|
expect(runs, instance.logs()).toHaveLength(1);
|
|
expect(runs[0]).toMatchObject({
|
|
completionTarget: "parent",
|
|
requesterTurnRunId: "private-yield-parent-turn",
|
|
execution: { status: "terminal", outcome: { status: "ok" } },
|
|
completion: { resultText: CHILD_MARKER },
|
|
});
|
|
expect(modelServer.countRequestsContaining(CHILD_MARKER)).toBe(0);
|
|
},
|
|
{ interval: 50, timeout: 60_000 },
|
|
);
|
|
// The parent checks the finished child through its real tool before yielding.
|
|
yieldGate.resolve();
|
|
expect(await parent, instance.logs()).toMatchObject({ status: "ok" });
|
|
await vi.waitFor(
|
|
() => {
|
|
const runs = [...loadSubagentRegistryFromSqlite().values()];
|
|
expect(runs, instance.logs()).toHaveLength(1);
|
|
expect(runs[0]?.delivery?.status, instance.logs()).toBe("delivered");
|
|
expect(runs[0]?.requesterSettleWake).toBeUndefined();
|
|
},
|
|
{ interval: 50, timeout: 30_000 },
|
|
);
|
|
const receipts = withOpenClawAgentDatabaseReadOnly(
|
|
({ db }) =>
|
|
executeSqliteQuerySync(
|
|
db,
|
|
getSessionKysely(db)
|
|
.selectFrom("session_input_completions")
|
|
.selectAll()
|
|
.where("session_id", "=", sessionId),
|
|
).rows,
|
|
{ agentId: REQUESTER_AGENT_ID },
|
|
);
|
|
expect(receipts.found).toBe(true);
|
|
if (!receipts.found) {
|
|
throw new Error("Expected requester database for completion verification");
|
|
}
|
|
expect(receipts.value).toEqual([]);
|
|
expect(modelServer.completionResponseCount()).toBe(1);
|
|
expect(chatErrors).toEqual([]);
|
|
const history = await client.request<{
|
|
messages: Array<{ role?: string; content?: unknown }>;
|
|
}>("chat.history", { sessionKey, agentId: REQUESTER_AGENT_ID, limit: 30 });
|
|
expect(
|
|
history.messages.filter(
|
|
(message) =>
|
|
message.role === "assistant" && extractFirstTextBlock(message) === CHILD_MARKER,
|
|
),
|
|
).toHaveLength(1);
|
|
expect(JSON.stringify(history.messages)).not.toContain("This turn ended before a reply");
|
|
const inputs = readInputs();
|
|
expect(inputs.found && inputs.value.filter((input) => input.state !== "cancelled")).toEqual(
|
|
[],
|
|
);
|
|
} finally {
|
|
yieldGate.resolve();
|
|
await disconnectGatewayClient(client);
|
|
await instance.stopGateway();
|
|
}
|
|
expect(instance.logs()).not.toContain(
|
|
"subagent source lifecycle changed before completion delivery",
|
|
);
|
|
},
|
|
);
|
|
|
|
it(
|
|
"preserves the requester owner through a fresh normalized spawn",
|
|
{ timeout: TEST_TIMEOUT_MS },
|
|
async () => {
|
|
const modelServer = await startProofModelServer();
|
|
modelServers.push(modelServer);
|
|
const instance = await createOpenClawTestInstance({
|
|
name: "requester-owner-requester-agent-id",
|
|
config: createTestConfig(modelServer.url),
|
|
env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined },
|
|
});
|
|
instances.push(instance);
|
|
|
|
instance.state.applyEnv();
|
|
try {
|
|
await writeSubagentSessionEntry({
|
|
stateDir: instance.stateDir,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
sessionKey: REQUESTER_KEY,
|
|
sessionId: "requester-owner-requester-session",
|
|
defaultSessionId: "requester-owner-requester-session",
|
|
});
|
|
} finally {
|
|
closeOpenClawStateDatabaseForTest();
|
|
}
|
|
|
|
await instance.startGateway();
|
|
const client = await connectGatewayClient({
|
|
url: instance.url,
|
|
token: instance.gatewayToken,
|
|
});
|
|
try {
|
|
const parent = client.request(
|
|
"agent",
|
|
{
|
|
sessionKey: REQUESTER_KEY,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
idempotencyKey: "requester-owner-parent-turn",
|
|
message: PARENT_PROMPT,
|
|
deliver: false,
|
|
},
|
|
{ expectFinal: true },
|
|
);
|
|
void parent.catch(() => {});
|
|
await vi.waitFor(
|
|
() => expect(modelServer.countRequestsContaining(CHILD_MARKER)).toBeGreaterThan(0),
|
|
{ interval: 50, timeout: 90_000 },
|
|
);
|
|
expect(await parent).toMatchObject({ status: "ok" });
|
|
instance.state.applyEnv();
|
|
await vi.waitFor(
|
|
() => {
|
|
const runs = [...loadSubagentRegistryFromSqlite().values()];
|
|
expect(runs).toHaveLength(1);
|
|
expect(runs[0]?.delivery?.status).toBe("delivered");
|
|
},
|
|
{ interval: 50, timeout: 25_000 },
|
|
);
|
|
const requester = await client.request<{ messages: unknown[] }>("chat.history", {
|
|
sessionKey: REQUESTER_KEY,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
limit: 30,
|
|
});
|
|
const other = await client.request<{ messages: unknown[] }>("chat.history", {
|
|
sessionKey: REQUESTER_KEY,
|
|
agentId: OTHER_AGENT_ID,
|
|
limit: 30,
|
|
});
|
|
expect(
|
|
requester.messages.filter((message) => JSON.stringify(message).includes(CHILD_MARKER)),
|
|
).toHaveLength(1);
|
|
expect(other.messages).toEqual([]);
|
|
} finally {
|
|
await disconnectGatewayClient(client);
|
|
await instance.stopGateway();
|
|
}
|
|
|
|
const logs = instance.logs();
|
|
instance.state.applyEnv();
|
|
try {
|
|
const runs = [...loadSubagentRegistryFromSqlite().values()];
|
|
expect(runs, logs).toHaveLength(1);
|
|
const run = runs[0]!;
|
|
expect(run.requesterAgentId, logs).toBe(REQUESTER_AGENT_ID);
|
|
expect(run.requesterSessionKey, logs).toBe(`agent:${REQUESTER_AGENT_ID}:${REQUESTER_KEY}`);
|
|
expect(run.execution.status, logs).toBe("terminal");
|
|
expect(run.execution.outcome, logs).toMatchObject({ status: "ok" });
|
|
expect(run.delivery?.status, logs).toBe("delivered");
|
|
expect(run.requesterSettleWake, logs).toBeUndefined();
|
|
} finally {
|
|
closeOpenClawStateDatabaseForTest();
|
|
}
|
|
expect(logs).not.toContain(ANNOUNCE_FAILURE_MARKER);
|
|
},
|
|
);
|
|
|
|
it(
|
|
"delivers a restored unscoped completion once across Gateway restarts",
|
|
{ timeout: TEST_TIMEOUT_MS },
|
|
async () => {
|
|
const modelServer = await startProofModelServer();
|
|
modelServers.push(modelServer);
|
|
const instance = await createOpenClawTestInstance({
|
|
name: "requester-owner-legacy-unscoped-requester",
|
|
config: createTestConfig(modelServer.url),
|
|
env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined },
|
|
});
|
|
instances.push(instance);
|
|
|
|
instance.state.applyEnv();
|
|
const registry =
|
|
await import("../src/agents/subagents/registry/subagent-registry.test-helpers.js");
|
|
await runQaGatewayFixture(
|
|
async () => {
|
|
await registry.resetSubagentRegistryForTests({ persist: false });
|
|
const childSessionKey = `agent:${REQUESTER_AGENT_ID}:subagent:requester-owner-legacy`;
|
|
await writeSubagentSessionEntry({
|
|
stateDir: instance.stateDir,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
sessionKey: RESTORED_REQUESTER_KEY,
|
|
sessionId: "requester-owner-legacy-session",
|
|
defaultSessionId: "requester-owner-legacy-session",
|
|
});
|
|
await writeSubagentSessionEntry({
|
|
stateDir: instance.stateDir,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
sessionKey: childSessionKey,
|
|
sessionId: "requester-owner-legacy-child-session",
|
|
defaultSessionId: "requester-owner-legacy-child-session",
|
|
});
|
|
// Registration owns native child execution and physical-store provenance even for a bare key.
|
|
// Queue only while preparing stored completion state; no child RPC is dispatched.
|
|
await registry.registerSubagentRun({
|
|
runId: RESTORED_RUN_ID,
|
|
childSessionKey,
|
|
requesterSessionKey: RESTORED_REQUESTER_KEY,
|
|
requesterDisplayKey: RESTORED_REQUESTER_KEY,
|
|
requesterAgentId: REQUESTER_AGENT_ID,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
task: "REQUESTER-OWNER legacy restored completion",
|
|
cleanup: "keep",
|
|
expectsCompletionMessage: true,
|
|
queued: true,
|
|
});
|
|
const registered = loadSubagentRegistryFromSqlite().get(RESTORED_RUN_ID);
|
|
if (!registered) {
|
|
throw new Error("Restored requester fixture registration was not persisted");
|
|
}
|
|
expect(registered.requesterStorePath).toBeTruthy();
|
|
expect(registered.controllerStorePath).toBe(registered.requesterStorePath);
|
|
// Retire setup callbacks, not durable ownership, before freezing the completed fixture.
|
|
await registry.resetSubagentRegistryForTests({ persist: false });
|
|
const endedAt = Date.now();
|
|
const restored: SubagentRunRecord = {
|
|
...registered,
|
|
endedReason: "subagent-complete",
|
|
execution: {
|
|
...registered.execution,
|
|
status: "terminal",
|
|
startedAt: registered.createdAt,
|
|
endedAt,
|
|
outcome: { status: "ok" },
|
|
},
|
|
completion: { required: true, resultText: RESTORED_CHILD_RESULT, capturedAt: endedAt },
|
|
};
|
|
saveSubagentRegistryToSqlite(new Map([[restored.runId, restored]]));
|
|
const seeded = loadSubagentRegistryFromSqlite().get(RESTORED_RUN_ID);
|
|
expect(seeded?.requesterAgentId).toBe(REQUESTER_AGENT_ID);
|
|
expect(seeded?.requesterSessionKey).toBe(RESTORED_REQUESTER_KEY);
|
|
expect(seeded?.delivery?.status).toBe("pending");
|
|
},
|
|
() => settleSubagentRegistryPersistenceWork(),
|
|
() => registry.resetSubagentRegistryForTests({ persist: false }),
|
|
() => closeOpenClawStateDatabaseForTest(),
|
|
);
|
|
|
|
let settledRequests: number | undefined;
|
|
for (let boot = 0; boot < 2; boot += 1) {
|
|
await instance.startGateway();
|
|
const client = await connectGatewayClient({
|
|
url: instance.url,
|
|
token: instance.gatewayToken,
|
|
});
|
|
try {
|
|
instance.state.applyEnv();
|
|
await vi.waitFor(
|
|
() => {
|
|
const run = loadSubagentRegistryFromSqlite().get(RESTORED_RUN_ID);
|
|
expect(run?.delivery?.status, instance.logs()).toBe("delivered");
|
|
expect(run?.execution.outcome).toMatchObject({ status: "ok" });
|
|
expect(run?.requesterSettleWake).toBeUndefined();
|
|
},
|
|
{ interval: 50, timeout: 25_000 },
|
|
);
|
|
const requester = await client.request<{ messages: unknown[] }>("chat.history", {
|
|
sessionKey: RESTORED_REQUESTER_KEY,
|
|
agentId: REQUESTER_AGENT_ID,
|
|
limit: 30,
|
|
});
|
|
const other = await client.request<{ messages: unknown[] }>("chat.history", {
|
|
sessionKey: RESTORED_REQUESTER_KEY,
|
|
agentId: OTHER_AGENT_ID,
|
|
limit: 30,
|
|
});
|
|
expect(
|
|
requester.messages.filter((message) =>
|
|
JSON.stringify(message).includes(RESTORED_CHILD_RESULT),
|
|
),
|
|
).toHaveLength(1);
|
|
expect(other.messages).toEqual([]);
|
|
expect(modelServer.countRequestsContaining(RESTORED_CHILD_RESULT)).toBe(1);
|
|
if (settledRequests !== undefined) {
|
|
expect(modelServer.requestCount()).toBe(settledRequests);
|
|
}
|
|
settledRequests = modelServer.requestCount();
|
|
} finally {
|
|
await disconnectGatewayClient(client);
|
|
await instance.stopGateway();
|
|
}
|
|
}
|
|
|
|
const logs = instance.logs();
|
|
expect(logs).not.toContain(ANNOUNCE_FAILURE_MARKER);
|
|
instance.state.applyEnv();
|
|
try {
|
|
const run = loadSubagentRegistryFromSqlite().get(RESTORED_RUN_ID);
|
|
expect(run?.requesterAgentId, logs).toBe(REQUESTER_AGENT_ID);
|
|
expect(run?.requesterSessionKey, logs).toBe(RESTORED_REQUESTER_KEY);
|
|
expect(run?.delivery?.status, logs).toBe("delivered");
|
|
} finally {
|
|
closeOpenClawStateDatabaseForTest();
|
|
}
|
|
},
|
|
);
|
|
});
|
|
|
|
function createTestConfig(baseUrl: string): OpenClawConfig {
|
|
return {
|
|
logging: { file: "${OPENCLAW_STATE_DIR}/logs/requester-owner-e2e.log" },
|
|
plugins: { enabled: false },
|
|
agents: {
|
|
ownership: "explicit",
|
|
entries: { [OTHER_AGENT_ID]: {}, [REQUESTER_AGENT_ID]: {} },
|
|
defaults: {
|
|
heartbeat: { every: "0m" },
|
|
maxConcurrent: 8,
|
|
model: { primary: MODEL_REF },
|
|
models: { [MODEL_REF]: { agentRuntime: { id: "openclaw" } } },
|
|
skipBootstrap: true,
|
|
skills: [],
|
|
},
|
|
},
|
|
tools: { profile: "coding", codeMode: false, toolSearch: false },
|
|
models: {
|
|
mode: "replace",
|
|
providers: {
|
|
"requester-owner": {
|
|
baseUrl: `${baseUrl}/v1`,
|
|
apiKey: "test-token-placeholder",
|
|
api: "openai-responses",
|
|
request: { allowPrivateNetwork: true },
|
|
models: [
|
|
{
|
|
id: "synthetic",
|
|
name: "requester-owner",
|
|
api: "openai-responses",
|
|
reasoning: false,
|
|
input: ["text"],
|
|
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
|
|
contextWindow: 128_000,
|
|
maxTokens: 4_096,
|
|
},
|
|
],
|
|
},
|
|
},
|
|
},
|
|
};
|
|
}
|
|
|
|
let responseSequence = 0;
|
|
|
|
function buildToolCallEvents(name: string, args: Record<string, unknown>): SseEvent[] {
|
|
const sequence = ++responseSequence;
|
|
const responseId = `resp_requester-owner_tool_${sequence}`;
|
|
const itemId = `fc_requester-owner_${sequence}`;
|
|
const callId = `call_requester-owner_${sequence}`;
|
|
const argumentsText = JSON.stringify(args);
|
|
const item = {
|
|
type: "function_call",
|
|
id: itemId,
|
|
call_id: callId,
|
|
name,
|
|
arguments: argumentsText,
|
|
};
|
|
return [
|
|
{
|
|
type: "response.created",
|
|
response: { id: responseId, object: "response", status: "in_progress", output: [] },
|
|
},
|
|
{ type: "response.output_item.added", output_index: 0, item: { ...item, arguments: "" } },
|
|
{
|
|
type: "response.function_call_arguments.delta",
|
|
item_id: itemId,
|
|
output_index: 0,
|
|
delta: argumentsText,
|
|
},
|
|
{
|
|
type: "response.function_call_arguments.done",
|
|
item_id: itemId,
|
|
output_index: 0,
|
|
arguments: argumentsText,
|
|
},
|
|
{ type: "response.output_item.done", output_index: 0, item },
|
|
{
|
|
type: "response.completed",
|
|
response: {
|
|
id: responseId,
|
|
object: "response",
|
|
status: "completed",
|
|
output: [item],
|
|
usage: { input_tokens: 8, output_tokens: 4, total_tokens: 12 },
|
|
},
|
|
},
|
|
];
|
|
}
|
|
|
|
async function startProofModelServer(options?: {
|
|
yieldAfterSpawn: Promise<void>;
|
|
childReply?: Promise<void>;
|
|
placement?: { kind: "local" };
|
|
}): Promise<ProofModelServer> {
|
|
const requestBodies: string[] = [];
|
|
let parentCheckedChildren = false;
|
|
let parentYielded = false;
|
|
let completionResponses = 0;
|
|
const server = createServer((request, response) => {
|
|
void handleModelRequest(request, response).catch((error: unknown) => {
|
|
if (!response.headersSent) {
|
|
response.writeHead(500, { "content-type": "application/json" });
|
|
}
|
|
response.end(JSON.stringify({ error: { message: String(error) } }));
|
|
});
|
|
});
|
|
|
|
async function handleModelRequest(
|
|
request: IncomingMessage,
|
|
response: ServerResponse,
|
|
): Promise<void> {
|
|
const url = new URL(request.url ?? "/", "http://127.0.0.1");
|
|
if (request.method === "GET" && url.pathname === "/v1/models") {
|
|
response.writeHead(200, { "content-type": "application/json" });
|
|
response.end(JSON.stringify({ data: [{ id: "requester-owner", object: "model" }] }));
|
|
return;
|
|
}
|
|
if (request.method !== "POST" || url.pathname !== "/v1/responses") {
|
|
response.writeHead(404).end();
|
|
return;
|
|
}
|
|
let body = "";
|
|
for await (const chunk of request) {
|
|
body += typeof chunk === "string" ? chunk : Buffer.from(chunk).toString("utf8");
|
|
}
|
|
requestBodies.push(body);
|
|
const requestBody = JSON.parse(body) as { tools?: Array<{ type?: string; name?: string }> };
|
|
const respondWithTool = (name: string, args: Record<string, unknown>) => {
|
|
expect(requestBody.tools).toEqual(
|
|
expect.arrayContaining([expect.objectContaining({ type: "function", name })]),
|
|
);
|
|
writeOpenAiResponsesSse(response, buildToolCallEvents(name, args));
|
|
};
|
|
if (options?.yieldAfterSpawn && parentCheckedChildren && !parentYielded) {
|
|
parentYielded = true;
|
|
respondWithTool("sessions_yield", {});
|
|
return;
|
|
}
|
|
const completion = [RESTORED_CHILD_RESULT, CHILD_MARKER].find((marker) =>
|
|
body.includes(marker),
|
|
);
|
|
if (completion) {
|
|
completionResponses += 1;
|
|
writeOpenAiResponsesText(response, {
|
|
text: completion,
|
|
responseId: `response-${++responseSequence}`,
|
|
messageId: `message-${responseSequence}`,
|
|
});
|
|
return;
|
|
}
|
|
|
|
if (body.includes(CHILD_TASK) && !body.includes("function_call_output")) {
|
|
if (options?.childReply) {
|
|
await options.childReply;
|
|
}
|
|
writeOpenAiResponsesText(response, {
|
|
text: CHILD_MARKER,
|
|
responseId: `response-${++responseSequence}`,
|
|
messageId: `message-${responseSequence}`,
|
|
});
|
|
return;
|
|
}
|
|
if (body.includes(PARENT_PROMPT) && !body.includes("function_call_output")) {
|
|
respondWithTool("sessions_spawn", {
|
|
task: CHILD_TASK,
|
|
label: "requester-owner-child",
|
|
...(options?.placement ? { placement: options.placement } : {}),
|
|
thread: false,
|
|
mode: "run",
|
|
...(options?.yieldAfterSpawn ? { completionTarget: "parent" } : {}),
|
|
});
|
|
return;
|
|
}
|
|
if (options?.yieldAfterSpawn) {
|
|
await options.yieldAfterSpawn;
|
|
parentCheckedChildren = true;
|
|
respondWithTool("subagents", { action: "list" });
|
|
return;
|
|
}
|
|
writeOpenAiResponsesText(response, {
|
|
text: "REQUESTER-OWNER-PARENT-OK",
|
|
responseId: `response-${++responseSequence}`,
|
|
messageId: `message-${responseSequence}`,
|
|
});
|
|
}
|
|
|
|
await new Promise<void>((resolve, reject) => {
|
|
server.once("error", reject);
|
|
server.listen(0, "127.0.0.1", resolve);
|
|
});
|
|
const address = server.address() as AddressInfo;
|
|
return {
|
|
bodies: () => requestBodies,
|
|
countRequestsContaining: (marker) =>
|
|
requestBodies.filter((entry) => entry.includes(marker)).length,
|
|
completionResponseCount: () => completionResponses,
|
|
requestCount: () => requestBodies.length,
|
|
url: `http://127.0.0.1:${address.port}`,
|
|
close: async () => {
|
|
server.closeAllConnections();
|
|
await new Promise<void>((resolve) => {
|
|
server.close(() => resolve());
|
|
});
|
|
},
|
|
};
|
|
}
|