mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-25 11:55:47 -06:00
fix(agents): prevent requester settle while child is still running (#120601)
* fix(agents): keep requester settle attached to live children * fix(agents): gate requester settle on terminal children --------- Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com>
This commit is contained in:
@@ -437,6 +437,63 @@ describe("maybeWakeRequesterAfterAllChildrenSettled", () => {
|
||||
expect(completeBatchSpy).toHaveBeenCalledWith(["run-b"], 1);
|
||||
});
|
||||
|
||||
it.each([
|
||||
["is running without an end timestamp", { status: "running", startedAt: 2_000 }],
|
||||
[
|
||||
"is still marked running with an end timestamp",
|
||||
{ status: "running", startedAt: 2_000, endedAt: 3_000 },
|
||||
],
|
||||
["has no end timestamp", { status: "terminal", startedAt: 2_000 }],
|
||||
] as const)(
|
||||
"does not wake a yielded requester while its only frozen child %s",
|
||||
async (_description, execution) => {
|
||||
const activeChild = makeSettledChild({
|
||||
runId: "run-b",
|
||||
execution,
|
||||
delivery: { status: "pending" },
|
||||
requesterSettleWake: {
|
||||
status: "pending",
|
||||
attemptCount: 0,
|
||||
batchRunIds: ["run-b"],
|
||||
requesterYieldBatch: true,
|
||||
rearmGeneration: 1,
|
||||
},
|
||||
});
|
||||
registryRuntimeMock.listSubagentRunsForRequester.mockReturnValue([activeChild]);
|
||||
|
||||
const woke = await maybeWakeRequesterAfterAllChildrenSettled(
|
||||
wakeParams({ settledEntry: activeChild }),
|
||||
);
|
||||
|
||||
expect(woke).toBe(false);
|
||||
expect(deliverSpy).not.toHaveBeenCalled();
|
||||
expect(completeBatchSpy).not.toHaveBeenCalled();
|
||||
},
|
||||
);
|
||||
|
||||
it("wakes after a retired frozen member disappears from the registry", async () => {
|
||||
const remainingChild = makeSettledChild({
|
||||
runId: "run-a",
|
||||
delivery: { status: "delivered" },
|
||||
requesterSettleWake: {
|
||||
status: "pending",
|
||||
attemptCount: 0,
|
||||
batchRunIds: ["run-a", "run-b"],
|
||||
requesterYieldBatch: true,
|
||||
rearmGeneration: 1,
|
||||
},
|
||||
});
|
||||
registryRuntimeMock.listSubagentRunsForRequester.mockReturnValue([remainingChild]);
|
||||
|
||||
const woke = await maybeWakeRequesterAfterAllChildrenSettled(
|
||||
wakeParams({ settledEntry: remainingChild }),
|
||||
);
|
||||
|
||||
expect(woke).toBe(true);
|
||||
expect(deliverSpy).toHaveBeenCalledOnce();
|
||||
expect(completeBatchSpy).toHaveBeenCalledWith(["run-a"], 1);
|
||||
});
|
||||
|
||||
it("wakes after a requester yields with one already-delivered completion", async () => {
|
||||
const child = makeSettledChild({
|
||||
runId: "run-b",
|
||||
|
||||
@@ -257,6 +257,8 @@ export async function maybeWakeRequesterAfterAllChildrenSettled(params: {
|
||||
let settledBatch: SubagentRunRecord[];
|
||||
if (frozenBatchRunIds && frozenBatchRunIds.length > 0) {
|
||||
const runsById = new Map(requesterRuns.map((entry) => [entry.runId, entry]));
|
||||
// Retired rows no longer own completion, but every surviving frozen member
|
||||
// must be terminal before this batch can wake its requester.
|
||||
settledBatch = frozenBatchRunIds
|
||||
.map((runId) => runsById.get(runId))
|
||||
.filter(
|
||||
@@ -264,9 +266,21 @@ export async function maybeWakeRequesterAfterAllChildrenSettled(params: {
|
||||
Boolean(entry?.requesterSettleWake) &&
|
||||
entry?.requesterSettleWake?.rearmGeneration === currentRearmGeneration,
|
||||
);
|
||||
if (
|
||||
settledBatch.some(
|
||||
(entry) => entry.execution.status === "running" || !hasSubagentRunEnded(entry),
|
||||
)
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
} else {
|
||||
settledBatch = buildConnectedSettledWave(
|
||||
requesterRuns.filter((entry) => entry.requesterSettleWake && hasSubagentRunEnded(entry)),
|
||||
requesterRuns.filter(
|
||||
(entry) =>
|
||||
entry.requesterSettleWake &&
|
||||
entry.execution.status !== "running" &&
|
||||
hasSubagentRunEnded(entry),
|
||||
),
|
||||
currentSettledEntry,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -6,6 +6,7 @@ import type {
|
||||
SubagentRegistryLifecycleState,
|
||||
} from "./subagent-registry-lifecycle-contracts.js";
|
||||
import type { RequesterSettleWakeState, SubagentRunRecord } from "./subagent-registry.types.js";
|
||||
import { hasSubagentRunEnded } from "./subagent-run-liveness.js";
|
||||
|
||||
type RequesterSettleWakeBatchState =
|
||||
import("./subagent-announce.requester-settle-wake.js").RequesterSettleWakeBatchState;
|
||||
@@ -207,8 +208,12 @@ export function createSubagentRegistryLifecycleRequesterWake(
|
||||
|
||||
function scheduleRequesterSettleWake(runId: string, entry: SubagentRunRecord): void {
|
||||
const requesterSessionKey = entry.requesterSessionKey?.trim();
|
||||
// A replayed lifecycle start can retain an older endedAt; require both
|
||||
// terminal status and end evidence so a live child never wakes its requester.
|
||||
if (
|
||||
entry.collect ||
|
||||
entry.execution.status === "running" ||
|
||||
!hasSubagentRunEnded(entry) ||
|
||||
!requesterSessionKey ||
|
||||
(entry.requesterTurnRunId && entry.requesterTurnYielded === true) ||
|
||||
scheduledRequesterSettleWakeRuns.has(runId) ||
|
||||
|
||||
@@ -4121,6 +4121,26 @@ describe("requester settle wake trigger", () => {
|
||||
expect(later.requesterSettleWake).toEqual({ status: "pending", attemptCount: 0 });
|
||||
});
|
||||
|
||||
it("does not resume a persisted settle wake until its registry row is terminal", async () => {
|
||||
const entry = createRunEntry({
|
||||
requesterSettleWake: { status: "pending", attemptCount: 0 },
|
||||
});
|
||||
const settleWake = vi.fn(async () => false);
|
||||
const controller = createLifecycleController({
|
||||
entry,
|
||||
maybeWakeRequesterAfterAllChildrenSettled: settleWake,
|
||||
});
|
||||
|
||||
controller.resumeRequesterSettleWake(entry.runId, entry);
|
||||
await Promise.resolve();
|
||||
expect(settleWake).not.toHaveBeenCalled();
|
||||
|
||||
entry.execution = { ...entry.execution, status: "terminal", endedAt: 4_000 };
|
||||
controller.resumeRequesterSettleWake(entry.runId, entry);
|
||||
|
||||
await waitForLifecycleState(() => expect(settleWake).toHaveBeenCalledOnce());
|
||||
});
|
||||
|
||||
it("keeps a yielded completion parked until its requester turn settles", async () => {
|
||||
const entry = createRunEntry({
|
||||
requesterTurnRunId: "run-requester",
|
||||
|
||||
@@ -4,7 +4,10 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { testing as subagentAnnounceDeliveryTesting } from "./subagent-announce-delivery.test-support.js";
|
||||
import { testing as subagentAnnounceOutputTesting } from "./subagent-announce-output.test-support.js";
|
||||
import { testing as subagentAnnounceTesting } from "./subagent-announce.js";
|
||||
import { testing as settleWakeTesting } from "./subagent-announce.requester-settle-wake.js";
|
||||
import {
|
||||
maybeWakeRequesterAfterAllChildrenSettled,
|
||||
testing as settleWakeTesting,
|
||||
} from "./subagent-announce.requester-settle-wake.js";
|
||||
import * as announceRead from "./subagent-registry-announce-read.js";
|
||||
import * as mod from "./subagent-registry.test-helpers.js";
|
||||
|
||||
@@ -363,6 +366,17 @@ describe("subagent registry lifecycle error grace", () => {
|
||||
.filter((request): request is GatewayRequest => request.method === "agent");
|
||||
}
|
||||
|
||||
function getRequesterWakeCalls() {
|
||||
return getAgentCalls().filter((request) => {
|
||||
const idempotencyKey = (request.params as Record<string, unknown> | undefined)
|
||||
?.idempotencyKey;
|
||||
return (
|
||||
typeof idempotencyKey === "string" &&
|
||||
idempotencyKey.startsWith("announce:requester-settle:")
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
function getAgentResultsForChildSession(childSessionKey: string): string[] {
|
||||
return getAgentCalls()
|
||||
.filter((request) => {
|
||||
@@ -479,25 +493,16 @@ describe("subagent registry lifecycle error grace", () => {
|
||||
rearmGeneration: undefined,
|
||||
},
|
||||
]);
|
||||
const requesterWakeCalls = () =>
|
||||
getAgentCalls().filter((request) => {
|
||||
const idempotencyKey = (request.params as Record<string, unknown> | undefined)
|
||||
?.idempotencyKey;
|
||||
return (
|
||||
typeof idempotencyKey === "string" &&
|
||||
idempotencyKey.startsWith("announce:requester-settle:")
|
||||
);
|
||||
});
|
||||
await waitForAgentCallCount(3);
|
||||
expect(requesterWakeCalls()).toHaveLength(1);
|
||||
expect(getRequesterWakeCalls()).toHaveLength(1);
|
||||
|
||||
agentCallGates.delete(betaSessionKey);
|
||||
releaseBetaDelivery?.();
|
||||
await waitForDeliveredCleanup("run-yield-alpha");
|
||||
await waitForDeliveredCleanup("run-yield-beta");
|
||||
|
||||
expect(requesterWakeCalls()).toHaveLength(1);
|
||||
const requesterWakeParams = requesterWakeCalls()[0]?.params as
|
||||
expect(getRequesterWakeCalls()).toHaveLength(1);
|
||||
const requesterWakeParams = getRequesterWakeCalls()[0]?.params as
|
||||
| Record<string, unknown>
|
||||
| undefined;
|
||||
expect(requesterWakeParams?.idempotencyKey).toContain(":yield-1");
|
||||
@@ -511,7 +516,66 @@ describe("subagent registry lifecycle error grace", () => {
|
||||
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
await flushAsync();
|
||||
expect(requesterWakeCalls()).toHaveLength(1);
|
||||
expect(getRequesterWakeCalls()).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("keeps a frozen live child asleep until its real registry row becomes terminal", async () => {
|
||||
const requesterTurnRunId = "run-requester-live-child";
|
||||
const liveChildSessionKey = "agent:main:subagent:frozen-live-child";
|
||||
registerCompletionRun(
|
||||
"run-frozen-live-child",
|
||||
"frozen-live-child",
|
||||
"live child",
|
||||
requesterTurnRunId,
|
||||
);
|
||||
setAssistantOutput(liveChildSessionKey, "live child complete");
|
||||
|
||||
expect(
|
||||
mod.markRequesterTurnYielded({
|
||||
requesterSessionKey: MAIN_REQUESTER_SESSION_KEY,
|
||||
requesterTurnRunId,
|
||||
}),
|
||||
).toBe(1);
|
||||
expect(
|
||||
mod.settleRequesterAfterSessionSpawns({
|
||||
requesterSessionKey: MAIN_REQUESTER_SESSION_KEY,
|
||||
requesterTurnRunId,
|
||||
requesterYielded: true,
|
||||
acceptedSessionSpawns: [
|
||||
{ runId: "run-frozen-live-child", childSessionKey: liveChildSessionKey },
|
||||
],
|
||||
}),
|
||||
).toBe(true);
|
||||
|
||||
const liveChild = mod
|
||||
.listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY)
|
||||
.find((run) => run.runId === "run-frozen-live-child");
|
||||
if (!liveChild) {
|
||||
throw new Error("expected frozen live child");
|
||||
}
|
||||
expect(
|
||||
await maybeWakeRequesterAfterAllChildrenSettled({
|
||||
requesterSessionKey: MAIN_REQUESTER_SESSION_KEY,
|
||||
settledEntry: liveChild,
|
||||
transitionBatch: noop,
|
||||
completeBatch: noop,
|
||||
}),
|
||||
).toBe(false);
|
||||
await flushAsync();
|
||||
|
||||
expect(getRequesterWakeCalls()).toHaveLength(0);
|
||||
expect(liveChild.execution.status).toBe("running");
|
||||
|
||||
emitLifecycleEvent("run-frozen-live-child", { phase: "end", endedAt: Date.now() + 1 });
|
||||
await waitForAgentCallCount(1);
|
||||
await waitForDeliveredCleanup("run-frozen-live-child");
|
||||
|
||||
expect(
|
||||
mod
|
||||
.listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY)
|
||||
.find((run) => run.runId === "run-frozen-live-child")?.execution,
|
||||
).toMatchObject({ status: "terminal", endedAt: expect.any(Number) });
|
||||
expect(getRequesterWakeCalls()).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("ignores transient lifecycle errors when run retries and then ends successfully", async () => {
|
||||
|
||||
Reference in New Issue
Block a user