mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-24 19:35:28 -06:00
fix(agents): retire delivered requester finals (#123285)
* fix(agents): retire delivered requester finals * fix(agents): bind requester final receipts before yielded settlement --------- Co-authored-by: VACInc <3279061+VACInc@users.noreply.github.com> Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
@@ -1918,6 +1918,7 @@ describe("deliverSubagentAnnouncement completion delivery", () => {
|
||||
});
|
||||
|
||||
expectDeliveryPath(result, "direct");
|
||||
expect(result).toMatchObject({ requesterVisibleFinalDelivered: true });
|
||||
expect(callGateway).not.toHaveBeenCalled();
|
||||
expectInProcessAgentParams(dispatchGatewayMethodInProcess, {
|
||||
deliver: true,
|
||||
|
||||
@@ -519,35 +519,43 @@ export async function sendSubagentAnnounceDirectly(params: {
|
||||
error: "completion agent did not use the message tool for message-tool-only delivery",
|
||||
};
|
||||
}
|
||||
const hasVisibleCompletionReply = Boolean(
|
||||
const requesterVisibleFinalDelivered = Boolean(
|
||||
directAnnounceResult &&
|
||||
((params.requireVisibleReply
|
||||
? hasMessagingToolDeliveryToSource(directAnnounceResult, deliveryTarget, {
|
||||
requireFinalReply: true,
|
||||
})
|
||||
: hasMessagingToolDelivery) ||
|
||||
(hasVisibleAgentPayload(
|
||||
params.requireVisibleReply
|
||||
? {
|
||||
payloads: Array.isArray(directAnnounceResult.payloads)
|
||||
? directAnnounceResult.payloads.filter((payload) => {
|
||||
const flags = payload as Record<string, unknown>;
|
||||
return (
|
||||
flags?.isCommentary !== true &&
|
||||
flags?.isCompactionNotice !== true &&
|
||||
flags?.isFallbackNotice !== true &&
|
||||
flags?.isStatusNotice !== true &&
|
||||
flags?.visible !== false
|
||||
);
|
||||
})
|
||||
: [],
|
||||
}
|
||||
: directAnnounceResult,
|
||||
{ ...completionPayloadVisibility, includeSilentReplyPayloads: false },
|
||||
) &&
|
||||
(!params.requireVisibleReply ||
|
||||
directAnnounceResult.deliveryStatus?.status !== "suppressed"))),
|
||||
(hasMessagingToolDeliveryToSource(directAnnounceResult, deliveryTarget, {
|
||||
requireFinalReply: true,
|
||||
}) ||
|
||||
(shouldDeliverAgentFinal &&
|
||||
!requiresMessageToolDelivery &&
|
||||
hasVisibleAgentPayload(
|
||||
{
|
||||
payloads: Array.isArray(directAnnounceResult.payloads)
|
||||
? directAnnounceResult.payloads.filter((payload) => {
|
||||
const flags = payload as Record<string, unknown>;
|
||||
return (
|
||||
flags?.isCommentary !== true &&
|
||||
flags?.isCompactionNotice !== true &&
|
||||
flags?.isFallbackNotice !== true &&
|
||||
flags?.isStatusNotice !== true &&
|
||||
flags?.visible !== false
|
||||
);
|
||||
})
|
||||
: [],
|
||||
},
|
||||
{ ...completionPayloadVisibility, includeSilentReplyPayloads: false },
|
||||
) &&
|
||||
directAnnounceResult.deliveryStatus?.status !== "suppressed")),
|
||||
);
|
||||
const hasVisibleCompletionReply =
|
||||
requesterVisibleFinalDelivered ||
|
||||
(!params.requireVisibleReply &&
|
||||
Boolean(
|
||||
directAnnounceResult &&
|
||||
(hasMessagingToolDelivery ||
|
||||
hasVisibleAgentPayload(directAnnounceResult, {
|
||||
...completionPayloadVisibility,
|
||||
includeSilentReplyPayloads: false,
|
||||
})),
|
||||
));
|
||||
const acceptsIntentionalSilentCompletion =
|
||||
hasIntentionalSilentCompletionReply && !isSubagentCompletion;
|
||||
if (
|
||||
@@ -583,6 +591,11 @@ export async function sendSubagentAnnounceDirectly(params: {
|
||||
return {
|
||||
delivered: true,
|
||||
path: "direct",
|
||||
...(params.expectsCompletionMessage &&
|
||||
!params.requesterIsSubagent &&
|
||||
requesterVisibleFinalDelivered
|
||||
? { requesterVisibleFinalDelivered: true }
|
||||
: {}),
|
||||
};
|
||||
} catch (err) {
|
||||
const permanent = isPermanentAnnounceDeliveryError(err);
|
||||
|
||||
@@ -32,6 +32,8 @@ export type SubagentAnnounceDeliveryResult = {
|
||||
path: SubagentDeliveryPath;
|
||||
deliveredAt?: number;
|
||||
enqueuedAt?: number;
|
||||
/** Direct completion that already sent the yielded requester's visible final. */
|
||||
requesterVisibleFinalDelivered?: true;
|
||||
reason?: SubagentAnnounceDeliveryFailureReason;
|
||||
error?: string;
|
||||
// Stops fallback delivery when ownership changed or another terminal result
|
||||
|
||||
@@ -580,7 +580,7 @@ export const startSubagentAnnounceCleanupFlow = (
|
||||
retireSupersededCleanupInBackground(context, runId, entry, cleanupGeneration);
|
||||
return;
|
||||
}
|
||||
recordAnnounceDeliveryResult(entry, delivery);
|
||||
recordAnnounceDeliveryResult(entry, delivery, params.runs);
|
||||
if (delivery.delivered) {
|
||||
const deliveryState = ensureDeliveryState(entry);
|
||||
deliveryState.status = "delivered";
|
||||
|
||||
@@ -39,6 +39,7 @@ import type {
|
||||
} from "./subagent-registry-lifecycle-context.js";
|
||||
import type { PendingFinalDeliveryPayload, SubagentRunRecord } from "./subagent-registry.types.js";
|
||||
import { compareSubagentRunGeneration } from "./subagent-run-generation.js";
|
||||
import { hasSubagentRunEnded } from "./subagent-run-liveness.js";
|
||||
|
||||
const DELIVERY_MIRROR_HISTORY_MAX_CHARS = 128 * 1024;
|
||||
|
||||
@@ -78,6 +79,7 @@ export const formatAnnounceDeliveryError = (delivery: SubagentAnnounceDeliveryRe
|
||||
export const recordAnnounceDeliveryResult = (
|
||||
entry: SubagentRunRecord,
|
||||
delivery: SubagentAnnounceDeliveryResult,
|
||||
runs?: ReadonlyMap<string, SubagentRunRecord>,
|
||||
) => {
|
||||
const deliveryState = ensureDeliveryState(entry);
|
||||
if (typeof delivery.enqueuedAt === "number") {
|
||||
@@ -88,6 +90,34 @@ export const recordAnnounceDeliveryResult = (
|
||||
typeof delivery.deliveredAt === "number" ? delivery.deliveredAt : Date.now();
|
||||
deliveryState.deliveredAt = deliveredAt;
|
||||
deliveryState.lastDropReason = undefined;
|
||||
const requesterTurnRunId = entry.requesterTurnRunId?.trim();
|
||||
if (
|
||||
delivery.path === "direct" &&
|
||||
delivery.requesterVisibleFinalDelivered &&
|
||||
requesterTurnRunId
|
||||
) {
|
||||
const siblings = [...(runs?.values() ?? [])].filter(
|
||||
(sibling) =>
|
||||
sibling.requesterSessionKey === entry.requesterSessionKey &&
|
||||
sibling.requesterTurnRunId === requesterTurnRunId &&
|
||||
sibling.expectsCompletionMessage === true,
|
||||
);
|
||||
if (
|
||||
siblings.some((sibling) => sibling === entry) &&
|
||||
siblings.every(
|
||||
(sibling) =>
|
||||
sibling.execution.status === "terminal" &&
|
||||
hasSubagentRunEnded(sibling) &&
|
||||
(sibling === entry || sibling.delivery?.status === "delivered"),
|
||||
)
|
||||
) {
|
||||
// Bind final evidence before yielding; direct delivery is fenced once a yield is frozen.
|
||||
deliveryState.requesterVisibleFinal = {
|
||||
requesterTurnRunId,
|
||||
batchRunIds: siblings.map((sibling) => sibling.runId).toSorted(),
|
||||
};
|
||||
}
|
||||
}
|
||||
}
|
||||
deliveryState.disposition =
|
||||
delivery.disposition ?? (delivery.delivered ? "delivered" : "retryable");
|
||||
|
||||
@@ -79,6 +79,86 @@ describe("settleRequesterTurnAfterSessionSpawns", () => {
|
||||
expect(schedule).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("retires a completed yielded batch whose requester already produced its final", () => {
|
||||
const entry = makeRun("run-child");
|
||||
entry.cleanupCompletedAt = 2_100;
|
||||
entry.delivery = {
|
||||
status: "delivered",
|
||||
requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] },
|
||||
};
|
||||
const schedule = vi.fn();
|
||||
|
||||
expect(
|
||||
settleRequesterTurnAfterSessionSpawns({
|
||||
requesterSessionKey: REQUESTER,
|
||||
requesterTurnRunId: REQUESTER_TURN,
|
||||
requesterYielded: true,
|
||||
acceptedSessionSpawns: [accepted(entry)],
|
||||
runs: new Map([[entry.runId, entry]]),
|
||||
persistOrThrow: vi.fn(),
|
||||
schedule,
|
||||
}),
|
||||
).toBe(true);
|
||||
expect(entry.requesterSettleWake).toBeUndefined();
|
||||
expect(entry.requesterTurnRunId).toBeUndefined();
|
||||
expect(entry.delivery?.requesterVisibleFinal).toBeUndefined();
|
||||
expect(schedule).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it.each([
|
||||
[
|
||||
"another requester turn",
|
||||
(entry: SubagentRunRecord) => {
|
||||
entry.delivery!.requesterVisibleFinal!.requesterTurnRunId = "run-other";
|
||||
},
|
||||
],
|
||||
[
|
||||
"changed child membership",
|
||||
(entry: SubagentRunRecord) => {
|
||||
entry.delivery!.requesterVisibleFinal!.batchRunIds.push("run-later");
|
||||
},
|
||||
],
|
||||
[
|
||||
"unfinished cleanup",
|
||||
(entry: SubagentRunRecord) => {
|
||||
entry.cleanupCompletedAt = undefined;
|
||||
},
|
||||
],
|
||||
[
|
||||
"unfinished delivery",
|
||||
(entry: SubagentRunRecord) => {
|
||||
entry.delivery!.status = "in_progress";
|
||||
},
|
||||
],
|
||||
[
|
||||
"a replayed running child",
|
||||
(entry: SubagentRunRecord) => {
|
||||
entry.execution.status = "running";
|
||||
},
|
||||
],
|
||||
] as const)("keeps requester settlement when the final receipt has %s", (_, invalidate) => {
|
||||
const entry = makeRun("run-child");
|
||||
entry.cleanupCompletedAt = 2_100;
|
||||
entry.delivery = {
|
||||
status: "delivered",
|
||||
requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] },
|
||||
};
|
||||
invalidate(entry);
|
||||
|
||||
expect(
|
||||
settleRequesterTurnAfterSessionSpawns({
|
||||
requesterSessionKey: REQUESTER,
|
||||
requesterTurnRunId: REQUESTER_TURN,
|
||||
requesterYielded: true,
|
||||
acceptedSessionSpawns: [accepted(entry)],
|
||||
runs: new Map([[entry.runId, entry]]),
|
||||
persistOrThrow: vi.fn(),
|
||||
schedule: vi.fn(),
|
||||
}),
|
||||
).toBe(true);
|
||||
expect(entry.requesterSettleWake?.requesterYieldBatch).toBe(true);
|
||||
});
|
||||
|
||||
it.each([
|
||||
["matches", "agent:main:subagent:worker", true],
|
||||
["rejects", "agent:main:subagent:other", false],
|
||||
@@ -291,6 +371,31 @@ describe("settleRequesterTurnAfterSessionSpawns", () => {
|
||||
expect(entry.retireAfterRequesterTurn).toBeUndefined();
|
||||
});
|
||||
|
||||
it("retires a delete-mode row after its requester-owned final is already delivered", () => {
|
||||
const entry = makeRun("run-delete");
|
||||
entry.cleanup = "delete";
|
||||
entry.cleanupCompletedAt = 2_100;
|
||||
entry.retireAfterRequesterTurn = true;
|
||||
entry.delivery = {
|
||||
status: "delivered",
|
||||
requesterVisibleFinal: { requesterTurnRunId: REQUESTER_TURN, batchRunIds: [entry.runId] },
|
||||
};
|
||||
const runs = new Map([[entry.runId, entry]]);
|
||||
|
||||
expect(
|
||||
settleRequesterTurnAfterSessionSpawns({
|
||||
requesterSessionKey: REQUESTER,
|
||||
requesterTurnRunId: REQUESTER_TURN,
|
||||
requesterYielded: true,
|
||||
acceptedSessionSpawns: [accepted(entry)],
|
||||
runs,
|
||||
persistOrThrow: vi.fn(),
|
||||
schedule: vi.fn(),
|
||||
}),
|
||||
).toBe(true);
|
||||
expect(runs.has(entry.runId)).toBe(false);
|
||||
});
|
||||
|
||||
it("retires a completed delete-mode row after a normal requester answer", () => {
|
||||
const entry = makeRun("run-delete", false);
|
||||
entry.retireAfterRequesterTurn = true;
|
||||
|
||||
@@ -91,8 +91,25 @@ export function settleRequesterTurnAfterSessionSpawns(params: {
|
||||
requesterTurnYielded: entry.requesterTurnYielded,
|
||||
retireAfterRequesterTurn: entry.retireAfterRequesterTurn,
|
||||
}));
|
||||
const requesterAlreadyDeliveredFinal =
|
||||
params.requesterYielded &&
|
||||
entries.every(
|
||||
(entry) =>
|
||||
entry.execution.status === "terminal" &&
|
||||
typeof entry.execution.endedAt === "number" &&
|
||||
entry.delivery?.status === "delivered" &&
|
||||
typeof entry.cleanupCompletedAt === "number",
|
||||
) &&
|
||||
entries.some((entry) => {
|
||||
const receipt = entry.delivery?.requesterVisibleFinal;
|
||||
return (
|
||||
receipt?.requesterTurnRunId === requesterTurnRunId &&
|
||||
receipt.batchRunIds.length === batchRunIds.length &&
|
||||
receipt.batchRunIds.every((runId, index) => runId === batchRunIds[index])
|
||||
);
|
||||
});
|
||||
let rearmGeneration: number | undefined;
|
||||
if (params.requesterYielded) {
|
||||
if (params.requesterYielded && !requesterAlreadyDeliveredFinal) {
|
||||
rearmGeneration =
|
||||
Math.max(0, ...entries.map((entry) => entry.requesterSettleWake?.rearmGeneration ?? 0)) + 1;
|
||||
for (const entry of entries) {
|
||||
@@ -126,6 +143,9 @@ export function settleRequesterTurnAfterSessionSpawns(params: {
|
||||
}
|
||||
} else {
|
||||
for (const entry of entries) {
|
||||
if (entry.delivery) {
|
||||
delete entry.delivery.requesterVisibleFinal;
|
||||
}
|
||||
entry.requesterTurnRunId = undefined;
|
||||
entry.requesterTurnYielded = undefined;
|
||||
if (entry.retireAfterRequesterTurn === true) {
|
||||
|
||||
@@ -387,6 +387,46 @@ describe("subagent registry lifecycle error grace", () => {
|
||||
});
|
||||
}
|
||||
|
||||
it("does not replay a requester-owned final already delivered before its turn yields", async () => {
|
||||
const requesterTurnRunId = "run-requester-already-delivered";
|
||||
const runId = "run-completed-before-yield";
|
||||
const childSessionKey = "agent:main:subagent:completed-before-yield";
|
||||
registerCompletionRun(runId, "completed-before-yield", "finish once", requesterTurnRunId);
|
||||
setAssistantOutput(childSessionKey, "child complete");
|
||||
|
||||
emitLifecycleEvent(runId, { phase: "end", endedAt: Date.now() });
|
||||
await waitForDeliveredCleanup(runId);
|
||||
|
||||
const completed = mod
|
||||
.listSubagentRunsForRequester(MAIN_REQUESTER_SESSION_KEY)
|
||||
.find((run) => run.runId === runId);
|
||||
expect(completed?.delivery?.requesterVisibleFinal).toEqual({
|
||||
requesterTurnRunId,
|
||||
batchRunIds: [runId],
|
||||
});
|
||||
expect(getAgentCalls()).toHaveLength(1);
|
||||
expect(
|
||||
mod.markRequesterTurnYielded({
|
||||
requesterSessionKey: MAIN_REQUESTER_SESSION_KEY,
|
||||
requesterTurnRunId,
|
||||
}),
|
||||
).toBe(1);
|
||||
expect(
|
||||
mod.settleRequesterAfterSessionSpawns({
|
||||
requesterSessionKey: MAIN_REQUESTER_SESSION_KEY,
|
||||
requesterTurnRunId,
|
||||
requesterYielded: true,
|
||||
acceptedSessionSpawns: [{ runId, childSessionKey }],
|
||||
}),
|
||||
).toBe(true);
|
||||
|
||||
await vi.advanceTimersByTimeAsync(30_000);
|
||||
await flushAsync();
|
||||
expect(getAgentCalls()).toHaveLength(1);
|
||||
expect(getRequesterWakeCalls()).toHaveLength(0);
|
||||
expect(completed?.delivery?.requesterVisibleFinal).toBeUndefined();
|
||||
});
|
||||
|
||||
it("lets requester settlement own a yielded batch after sibling deliveries race", async () => {
|
||||
const requesterTurnRunId = "run-requester-yield-race";
|
||||
const alphaSessionKey = "agent:main:subagent:yield-alpha";
|
||||
|
||||
@@ -139,6 +139,24 @@ describe("subagent registry sqlite store", () => {
|
||||
});
|
||||
});
|
||||
|
||||
it("preserves requester-owned final receipts in the existing SQLite payload", async () => {
|
||||
await withTempStateEnv(async () => {
|
||||
const requesterVisibleFinal = {
|
||||
requesterTurnRunId: "run-requester",
|
||||
batchRunIds: ["run-one"],
|
||||
};
|
||||
const run = createRun({ delivery: { status: "delivered", requesterVisibleFinal } });
|
||||
|
||||
saveSubagentRegistryToSqlite(new Map([[run.runId, run]]));
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
|
||||
expect(loadSubagentRegistryFromSqlite().get(run.runId)?.delivery).toMatchObject({
|
||||
status: "delivered",
|
||||
requesterVisibleFinal,
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
it.each([
|
||||
{
|
||||
name: "visible",
|
||||
|
||||
@@ -139,6 +139,8 @@ export type SubagentCompletionDeliveryState = {
|
||||
enqueuedAt?: number;
|
||||
deliveredAt?: number;
|
||||
announcedAt?: number;
|
||||
/** Exact requester turn and completed child batch that already produced its visible final. */
|
||||
requesterVisibleFinal?: { requesterTurnRunId: string; batchRunIds: string[] };
|
||||
lastAttemptAt?: number;
|
||||
attemptCount?: number;
|
||||
lastError?: string | null;
|
||||
|
||||
Reference in New Issue
Block a user