mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix(agents): a subagent's previous result is lost when its next turn starts quickly (#163352)
* fix(agents): preserve completed turn delivery during follow-up Keep completed subagent results and their parent-delivery custody across ordinary next-turn admission. Fence predecessor session effects while preserving exact execution guards, independent receipts, and fresh successor authority. Paused continuation and steer retain replacement semantics. Regression proof uses registered agent admission with paused transcript and delivery boundaries, distinct turn outputs, durable receipts, restored pending delivery, and replay deduplication. Focused announce, reactivation, registry, type, lint, and formatting checks passed; independent review found no actionable P0-P2 findings. * test(agents): align follow-up delivery CI coverage Update sibling expectations for retained completed runs, keep requester completion task-scoped, and wait for the real Gateway dispatch signal in the admission regression. Migrate the benchmark to the cleanup ownership callback. * test(agents): prove independent completion delivery authority Exercise real child Gateway admission and requester input writes with independently revoked predecessor and successor sources. Assert persisted receipts, transcript attribution, exact source identity, and replay deduplication. All three cases reject a borrowed-authority negative control.
This commit is contained in:
parent
e86a14698f
commit
a9442b7a6d
29 changed files with 979 additions and 116 deletions
|
|
@ -472,6 +472,14 @@ outcomes fence affected rows until canonical restore; they never authorize repla
|
|||
The digest requires no new column, schema migration, or updater behavior, and also
|
||||
detects writes from older processes that cannot maintain a new revision column.
|
||||
|
||||
An ordinary follow-up to a completed child creates a new task and execution while
|
||||
retaining the completed row's pending parent delivery and receipt. The same row
|
||||
transaction fences the predecessor's session effects; its exact execution and
|
||||
requester custody still own announcement retries. Each turn reads its own transcript
|
||||
result and uses its own announcement idempotency key. Paused continuation and steer
|
||||
replacement retain their same-task adoption semantics. No schema migration or
|
||||
updater change is required.
|
||||
|
||||
Provisional cancellation claims keep registration pending until their owner releases
|
||||
or confirms them. Existing persistence notifications wake the wait; work cancellation,
|
||||
Gateway drain, and database retirement dispose its subscriptions. A released claim
|
||||
|
|
|
|||
|
|
@ -524,7 +524,7 @@ async function runSweepSample(childCount: number): Promise<Sample> {
|
|||
resumeRequesterSettleWake: () => {},
|
||||
startSubagentAnnounceCleanupFlow: () => true,
|
||||
completeCleanupBookkeeping: async () => {},
|
||||
isEndedHookOwnerCurrent: (runId, entry) =>
|
||||
isCleanupOwnerCurrent: (runId, entry) =>
|
||||
isSameSubagentRunOwner(runs.get(runId), entry) || !runs.has(runId),
|
||||
sessionEffectsHostCurrent: (entry) => entry.execution.suppressSessionEffects !== true,
|
||||
shouldSuppressSessionEffects: async (entry) => entry.execution.suppressSessionEffects === true,
|
||||
|
|
|
|||
|
|
@ -184,9 +184,11 @@ it.each([false, true])(
|
|||
task: "Check the remaining finding.",
|
||||
}),
|
||||
).toBe(true);
|
||||
expect(subagentRuns.has(childRunId)).toBe(false);
|
||||
expect(subagentRuns.get(childRunId)).toMatchObject({
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
});
|
||||
expect(subagentRuns.get(nextRunId)?.execution.status).toBe("running");
|
||||
expect(leasedSteering.isCurrent()).toBe(false);
|
||||
expect(leasedSteering.isCurrent()).toBe(true);
|
||||
return { content: [{ type: "text" as const, text: "Follow-up accepted." }], details: {} };
|
||||
});
|
||||
const { session } = await createTestSession({
|
||||
|
|
@ -226,8 +228,12 @@ it.each([false, true])(
|
|||
|
||||
expect(followUp).toHaveBeenCalledOnce();
|
||||
expect(JSON.stringify(requests[0])).toContain(escapedAnswer);
|
||||
expect(subagentRuns.has(childRunId)).toBe(false);
|
||||
expect(subagentRuns.get(childRunId)).toMatchObject({
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
delivery: { status: "delivered" },
|
||||
});
|
||||
expect(subagentRuns.get(nextRunId)).toMatchObject({
|
||||
taskRunId: nextRunId,
|
||||
task: "Check the remaining finding.",
|
||||
execution: { status: "running" },
|
||||
});
|
||||
|
|
|
|||
|
|
@ -215,11 +215,7 @@ export async function maybeWakeRequesterAfterAllChildrenSettled(
|
|||
const currentCompletionRows = (rows: SubagentRunRecord[]) =>
|
||||
frozenBatchRunIds?.length
|
||||
? rows.filter((entry) =>
|
||||
isRequesterCompletionCohortCurrent(
|
||||
entry,
|
||||
settledBatch,
|
||||
getLatestLiveSubagentRunByChildSessionKey,
|
||||
),
|
||||
isRequesterCompletionCohortCurrent(entry, getLatestLiveSubagentRunByChildSessionKey),
|
||||
)
|
||||
: dedupeLatestChildCompletionRows(
|
||||
filterCurrentDirectChildCompletionRows(rows, {
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ import {
|
|||
SUBAGENT_COMPLETION_OUTCOME_INSTRUCTION,
|
||||
SUBAGENT_PRIVATE_COMPLETION_INSTRUCTION,
|
||||
} from "../completion/subagent-completion-instructions.js";
|
||||
import { subagentRuns } from "../registry/subagent-registry-memory.js";
|
||||
import {
|
||||
countPendingDescendantRuns,
|
||||
getLatestSubagentRunByChildSessionKey,
|
||||
|
|
@ -195,11 +196,11 @@ async function runSubagentAnnounceFlowBound(
|
|||
childSessionEffectsAllowed() &&
|
||||
(await params.prepareChildSessionEffects?.()) !== false &&
|
||||
childSessionEffectsAllowed();
|
||||
let isOwnResultCurrent = () => true;
|
||||
let isOwnResultCurrent: (() => boolean) | undefined;
|
||||
let isChildResultsCurrent = () => true;
|
||||
const completionDeliveryAllowed = () =>
|
||||
params.isCompletionDeliveryAllowed?.() !== false &&
|
||||
isOwnResultCurrent() &&
|
||||
(isOwnResultCurrent?.() ?? true) &&
|
||||
isChildResultsCurrent();
|
||||
let childSessionId: string | undefined;
|
||||
let childSessionLifecycleRevision: string | undefined;
|
||||
|
|
@ -361,12 +362,8 @@ async function runSubagentAnnounceFlowBound(
|
|||
? (stripAndClassifyReply(fallbackReply ?? "") ?? undefined)
|
||||
: undefined;
|
||||
|
||||
const childRun = getLatestSubagentRunByChildSessionKey(params.childSessionKey);
|
||||
if (
|
||||
childRun?.runId === params.childRunId &&
|
||||
(await prepareChildSessionEffects()) &&
|
||||
childSessionEffectsAllowed()
|
||||
) {
|
||||
const childRun = subagentRuns.get(params.childRunId);
|
||||
if (childRun?.childSessionKey === params.childSessionKey && completionDeliveryAllowed()) {
|
||||
const prepared = await readSubagentRunAnnounceResult(childRun);
|
||||
reply = prepared.text;
|
||||
isOwnResultCurrent = prepared.isCurrent;
|
||||
|
|
@ -440,7 +437,7 @@ async function runSubagentAnnounceFlowBound(
|
|||
}
|
||||
|
||||
const childSessionCurrent = await prepareChildSessionEffects();
|
||||
if (!childSessionCurrent || !childSessionEffectsAllowed()) {
|
||||
if (!isOwnResultCurrent && (!childSessionCurrent || !childSessionEffectsAllowed())) {
|
||||
reply = params.roundOneReply ?? params.fallbackReply;
|
||||
if (
|
||||
expectsCompletionMessage &&
|
||||
|
|
|
|||
|
|
@ -277,8 +277,10 @@ it.each(["not-committed", "unknown", "successor"] as const)(
|
|||
release.resolve();
|
||||
await registration;
|
||||
await fixture.settle();
|
||||
expect(fixture.wake).not.toHaveBeenCalled();
|
||||
expect(loadSubagentRegistryFromSqlite().get(run.runId)?.cleanupCompletedAt).toBeUndefined();
|
||||
if (change !== "successor") {
|
||||
expect(fixture.wake).not.toHaveBeenCalled();
|
||||
expect(loadSubagentRegistryFromSqlite().get(run.runId)?.cleanupCompletedAt).toBeUndefined();
|
||||
}
|
||||
if (change === "not-committed") {
|
||||
expect(subagentRuns.get(run.runId)?.cleanupHandled).not.toBe(true);
|
||||
expect(loadSubagentRegistryFromSqlite().get(run.runId)).toEqual(before);
|
||||
|
|
@ -301,6 +303,15 @@ it.each(["not-committed", "unknown", "successor"] as const)(
|
|||
await fixture.settle();
|
||||
expect(fixture.wake).not.toHaveBeenCalled();
|
||||
} else {
|
||||
expect(fixture.wake).toHaveBeenCalledExactlyOnceWith(
|
||||
expect.objectContaining({
|
||||
settledEntry: expect.objectContaining({ runId: run.runId }),
|
||||
}),
|
||||
);
|
||||
expect(loadSubagentRegistryFromSqlite().get(run.runId)).toMatchObject({
|
||||
cleanupCompletedAt: expect.any(Number),
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
});
|
||||
expect(subagentRuns.get(`${run.runId}-successor`)?.execution.status).toBe("running");
|
||||
}
|
||||
} finally {
|
||||
|
|
|
|||
|
|
@ -355,8 +355,11 @@ it.each([
|
|||
),
|
||||
).toBe(true);
|
||||
const b1 = subagentRuns.get("publication-b1")!;
|
||||
expect(subagentRuns.has(b0.runId)).toBe(false);
|
||||
expect(b1).toMatchObject({ taskRunId: b0.taskRunId, execution: { status: "running" } });
|
||||
expect(subagentRuns.get(b0.runId)).toMatchObject({
|
||||
taskRunId: b0.taskRunId,
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
});
|
||||
expect(b1).toMatchObject({ taskRunId: b1.runId, execution: { status: "running" } });
|
||||
if (typeof b0.generation !== "number") {
|
||||
throw new Error("Registration did not mint a run generation");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -76,7 +76,7 @@ export interface SubagentLifecycleCleanupContext extends SubagentLifecycleCommon
|
|||
isCleanupAttemptCurrent(runId: string, entry: SubagentRunRecord, generation: number): boolean;
|
||||
isCleanupGeneration(entry: SubagentRunRecord, generation: number): boolean;
|
||||
isCleanupGenerationCurrent(runId: string, entry: SubagentRunRecord, generation: number): boolean;
|
||||
isEndedHookOwnerCurrent(runId: string, entry: SubagentRunRecord): boolean;
|
||||
isCleanupOwnerCurrent(runId: string, entry: SubagentRunRecord): boolean;
|
||||
startSubagentAnnounceCleanupFlow(runId: string, entry: SubagentRunRecord): boolean;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
import { expect, it, vi } from "vitest";
|
||||
import { getActiveGatewayRootWorkCount } from "../../../process/gateway-work-admission.js";
|
||||
import { createDeferredCore } from "../../../shared/deferred.js";
|
||||
import {
|
||||
runSubagentAnnounceDispatch,
|
||||
|
|
@ -316,7 +315,7 @@ export function registerLifecycleDeliveryReceiptCases({
|
|||
expect(readLifecycleRun(entry).delivery?.nextAttemptAt).toBeUndefined();
|
||||
});
|
||||
|
||||
it("keeps a late superseded-delivery retirement root-admitted", async () => {
|
||||
it("keeps a delivered receipt when a late failure arrives after the next turn starts", async () => {
|
||||
const entry = createRunEntry({ expectsCompletionMessage: true, generation: 1 });
|
||||
const runs = new Map([[entry.runId, entry]]);
|
||||
let onDeliveryResult: Parameters<
|
||||
|
|
@ -328,13 +327,7 @@ export function registerLifecycleDeliveryReceiptCases({
|
|||
return "delivered" as const;
|
||||
},
|
||||
);
|
||||
let releaseRetirement = () => {};
|
||||
const retirementPending = new Promise<void>((resolve) => {
|
||||
releaseRetirement = resolve;
|
||||
});
|
||||
const retireSupersededRun = vi.fn(async () => {
|
||||
await retirementPending;
|
||||
});
|
||||
const retireSupersededRun = vi.fn(async () => {});
|
||||
const controller = createLifecycleController({
|
||||
entry,
|
||||
runs,
|
||||
|
|
@ -342,8 +335,9 @@ export function registerLifecycleDeliveryReceiptCases({
|
|||
runSubagentAnnounceFlow,
|
||||
});
|
||||
|
||||
await completeRun(controller, entry, { triggerCleanup: true });
|
||||
await waitForLifecycleState(() => expect(getActiveGatewayRootWorkCount()).toBe(0));
|
||||
await completeAndJoinCleanup(controller, entry, { triggerCleanup: true });
|
||||
const delivered = readLifecycleRun(entry).delivery;
|
||||
expect(delivered?.status).toBe("delivered");
|
||||
const newer = createRunEntry({
|
||||
runId: "run-2",
|
||||
childSessionKey: entry.childSessionKey,
|
||||
|
|
@ -353,15 +347,9 @@ export function registerLifecycleDeliveryReceiptCases({
|
|||
|
||||
await onDeliveryResult?.({ delivered: false, path: "none" });
|
||||
|
||||
await waitForLifecycleState(() =>
|
||||
expect(retireSupersededRun).toHaveBeenCalledWith(
|
||||
entry.runId,
|
||||
expect.objectContaining({ runId: entry.runId, childSessionKey: entry.childSessionKey }),
|
||||
),
|
||||
);
|
||||
expect(getActiveGatewayRootWorkCount()).toBe(1);
|
||||
releaseRetirement();
|
||||
await waitForLifecycleState(() => expect(getActiveGatewayRootWorkCount()).toBe(0));
|
||||
expect(readLifecycleRun(entry).delivery).toEqual(delivered);
|
||||
expect(retireSupersededRun).not.toHaveBeenCalled();
|
||||
expect(runs.get(newer.runId)).toBe(newer);
|
||||
});
|
||||
|
||||
it("finalizes terminal visible-send failures without scheduling completion retry", async () => {
|
||||
|
|
|
|||
|
|
@ -132,8 +132,9 @@ export async function finishSubagentCleanup(
|
|||
const { cleanup, cleanupGeneration, stateContext, isCurrent } = args;
|
||||
let entry = args.entry;
|
||||
let runId = entry.runId;
|
||||
if (cleanup === "delete" || !entry.retainAttachmentsOnKeep) {
|
||||
await safeRemoveAttachmentsDir(entry, isCurrent);
|
||||
const sessionEffectsCurrent = () => isCurrent() && context.sessionEffectsHostCurrent(entry);
|
||||
if ((cleanup === "delete" || !entry.retainAttachmentsOnKeep) && sessionEffectsCurrent()) {
|
||||
await safeRemoveAttachmentsDir(entry, sessionEffectsCurrent);
|
||||
}
|
||||
if (!isCurrent()) {
|
||||
if (cleanupGeneration !== undefined) {
|
||||
|
|
@ -151,7 +152,7 @@ export async function finishSubagentCleanup(
|
|||
}
|
||||
const cleanupOwnerCurrent = () =>
|
||||
(cleanupGeneration === undefined || context.isCleanupGeneration(entry, cleanupGeneration)) &&
|
||||
context.isEndedHookOwnerCurrent(runId, entry);
|
||||
context.isCleanupOwnerCurrent(runId, entry);
|
||||
// Hook loading is best-effort; durable delivery and cleanup must already
|
||||
// be terminal before plugin code can fail or stall.
|
||||
await context.completeCleanupBookkeeping({
|
||||
|
|
|
|||
|
|
@ -3346,7 +3346,7 @@ describe("subagent registry lifecycle hardening", () => {
|
|||
expect(finalPostimage?.delivery?.announcedAt).toBeUndefined();
|
||||
});
|
||||
|
||||
it("retires a stale cleanup before deleting a newer session generation", async () => {
|
||||
it("finishes old cleanup without deleting a newer session generation", async () => {
|
||||
const entry = createRunEntry({
|
||||
cleanup: "delete",
|
||||
expectsCompletionMessage: false,
|
||||
|
|
@ -3354,11 +3354,10 @@ describe("subagent registry lifecycle hardening", () => {
|
|||
generation: 1,
|
||||
});
|
||||
const runs = new Map([[entry.runId, entry]]);
|
||||
const retireSupersededRun = vi.fn(async (runId: string) => {
|
||||
runs.delete(runId);
|
||||
});
|
||||
const retireSupersededRun = vi.fn(async () => {});
|
||||
const controller = createLifecycleController({ entry, runs, retireSupersededRun });
|
||||
|
||||
const join = observeRootWork();
|
||||
expect(controller.startSubagentAnnounceCleanupFlow(entry.runId, entry)).toBe(true);
|
||||
const newer = createRunEntry({
|
||||
runId: "run-2",
|
||||
|
|
@ -3368,12 +3367,12 @@ describe("subagent registry lifecycle hardening", () => {
|
|||
});
|
||||
runs.set(newer.runId, newer);
|
||||
|
||||
await waitForLifecycleState(() =>
|
||||
expect(retireSupersededRun).toHaveBeenCalledWith(
|
||||
entry.runId,
|
||||
expect.objectContaining({ runId: entry.runId, childSessionKey: entry.childSessionKey }),
|
||||
),
|
||||
);
|
||||
await join();
|
||||
expect(retireSupersededRun).not.toHaveBeenCalled();
|
||||
expect(readLifecycleRun(entry)).toMatchObject({
|
||||
cleanupCompletedAt: expect.any(Number),
|
||||
execution: { suppressSessionEffects: true },
|
||||
});
|
||||
expect(runs.get(newer.runId)).toBe(newer);
|
||||
expect(gatewayMocks.callGateway).not.toHaveBeenCalledWith(
|
||||
expect.objectContaining({ method: "sessions.delete" }),
|
||||
|
|
|
|||
|
|
@ -214,6 +214,7 @@ export class SubagentLifecycleController {
|
|||
this.options.runs.get(entry.runId) ?? getCurrentSubagentRunOwner(this.options.runs, entry);
|
||||
return (
|
||||
(current !== undefined && !isSameSubagentRunOwner(current, entry)) ||
|
||||
this.newerGenerationOwnsSession(entry) ||
|
||||
shouldSuppressSubagentRecoverySessionEffects(current ?? entry)
|
||||
);
|
||||
};
|
||||
|
|
@ -230,6 +231,7 @@ export class SubagentLifecycleController {
|
|||
this.options.runs.get(entry.runId) ?? getCurrentSubagentRunOwner(this.options.runs, entry);
|
||||
if (
|
||||
(current !== undefined && !isSameSubagentRunOwner(current, entry)) ||
|
||||
this.newerGenerationOwnsSession(entry) ||
|
||||
shouldSuppressSubagentRecoverySessionEffects(current ?? entry)
|
||||
) {
|
||||
return false;
|
||||
|
|
@ -366,8 +368,7 @@ export class SubagentLifecycleController {
|
|||
return (
|
||||
current !== undefined &&
|
||||
current.pauseReason !== "sessions_yield" &&
|
||||
this.isCleanupGeneration(entry, generation) &&
|
||||
!this.newerGenerationOwnsSession(current)
|
||||
this.isCleanupGeneration(entry, generation)
|
||||
);
|
||||
};
|
||||
isCleanupAttemptCurrent = (
|
||||
|
|
@ -377,15 +378,16 @@ export class SubagentLifecycleController {
|
|||
): boolean =>
|
||||
getCurrentSubagentRunOwner(this.options.runs, entry)?.cleanupHandled === true &&
|
||||
this.isCleanupGenerationCurrent(runId, entry, generation);
|
||||
isEndedHookOwnerCurrent = (_runId: string, entry: SubagentRunRecord): boolean => {
|
||||
isCleanupOwnerCurrent = (_runId: string, entry: SubagentRunRecord): boolean => {
|
||||
const current =
|
||||
this.options.runs.get(entry.runId) ?? getCurrentSubagentRunOwner(this.options.runs, entry);
|
||||
return (
|
||||
(current === undefined || isSameSubagentRunOwner(current, entry)) &&
|
||||
(current ?? entry).pauseReason !== "sessions_yield" &&
|
||||
!this.newerGenerationOwnsSession(entry)
|
||||
(current ?? entry).pauseReason !== "sessions_yield"
|
||||
);
|
||||
};
|
||||
isEndedHookOwnerCurrent = (runId: string, entry: SubagentRunRecord): boolean =>
|
||||
this.isCleanupOwnerCurrent(runId, entry) && !this.newerGenerationOwnsSession(entry);
|
||||
|
||||
bumpTerminalGeneration(entry: SubagentRunRecord, bindingChanged = false): number {
|
||||
const identity = this.trackRun(entry);
|
||||
|
|
|
|||
|
|
@ -298,7 +298,6 @@ class SubagentRunMap extends Map<string, SubagentRunRecord> {
|
|||
!record ||
|
||||
!isSameSubagentRunOwner(this.get(entry.runId), entry) ||
|
||||
record.generation !== entry.generation ||
|
||||
record.execution.suppressSessionEffects === true ||
|
||||
(!record.requesterTurnRunId &&
|
||||
!record.requesterSettleWake &&
|
||||
record.pauseReason !== "sessions_yield" &&
|
||||
|
|
|
|||
|
|
@ -94,7 +94,7 @@ describe("requester settle wake commit retry", () => {
|
|||
});
|
||||
|
||||
it.each(["another requester", "the same task"])(
|
||||
"rejects a frozen completion whose newer run belongs to %s",
|
||||
"keeps frozen completion custody task-scoped when a newer run belongs to %s",
|
||||
async (replacement) => {
|
||||
const entry = makeRetainedChild();
|
||||
entry.generation = 1;
|
||||
|
|
@ -108,8 +108,12 @@ describe("requester settle wake commit retry", () => {
|
|||
};
|
||||
const { context } = makeContext([entry, successor]);
|
||||
const commit = vi.fn(() => true);
|
||||
await commitRequesterWake(context, [entry, successor], undefined, commit, false);
|
||||
expect(commit).not.toHaveBeenCalled();
|
||||
await commitRequesterWake(context, [entry], undefined, commit, false);
|
||||
if (replacement === "another requester") {
|
||||
expect(commit).toHaveBeenCalledExactlyOnceWith([entry], expect.any(Object));
|
||||
} else {
|
||||
expect(commit).not.toHaveBeenCalled();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
|
|
|
|||
|
|
@ -249,7 +249,7 @@ export function commitRequesterInitialTransfer(
|
|||
(entry) =>
|
||||
context.pendingRequesterSettleWakeCommits.get(getSubagentRunRuntimeKey(entry)) !==
|
||||
pending ||
|
||||
!isRequesterCompletionCohortCurrent(entry, entries, (key, matches) =>
|
||||
!isRequesterCompletionCohortCurrent(entry, (key, matches) =>
|
||||
context.options.getLatestRunForChildSession(key, matches),
|
||||
),
|
||||
)
|
||||
|
|
@ -594,7 +594,7 @@ export function commitRequesterWake(
|
|||
!isDeepStrictEqual(captureRequesterSettleRunIdentity(live), owner.identity) ||
|
||||
!isDeepStrictEqual(live.killIntent, owner.killIntent) ||
|
||||
!isDeepStrictEqual(live.killReconciliation, owner.killReconciliation) ||
|
||||
!isRequesterCompletionCohortCurrent(live, pending.entries, (key, matches) =>
|
||||
!isRequesterCompletionCohortCurrent(live, (key, matches) =>
|
||||
context.options.getLatestRunForChildSession(key, matches),
|
||||
)
|
||||
) {
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import type { GatewayContextResolver } from "../../../gateway/server-methods/types.js";
|
||||
import { captureOperatorToolGatewayContinuationContext } from "../../../gateway/server-plugin-in-process-dispatch.js";
|
||||
import {
|
||||
getAgentEventLifecycleGeneration,
|
||||
isAgentEventLifecycleGenerationCurrent,
|
||||
|
|
@ -136,6 +137,8 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
expected?: SubagentRunRecord;
|
||||
runTimeoutSeconds?: number;
|
||||
allowEndedSource?: boolean;
|
||||
/** Ordinary next turns retain the completed execution's independent delivery. */
|
||||
preserveCompletedRun?: boolean;
|
||||
preserveFrozenResultFallback?: boolean;
|
||||
// A follow-up that continues a paused run inherits the original requester's
|
||||
// wake credential. An operator steer intentionally drops it: the operator is
|
||||
|
|
@ -172,6 +175,11 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
if (!selected) {
|
||||
return false;
|
||||
}
|
||||
const preserveCompletedRun = replaceParams.preserveCompletedRun === true;
|
||||
const authority = preserveCompletedRun
|
||||
? await captureOperatorToolGatewayContinuationContext()
|
||||
: undefined;
|
||||
let custodyTransferred = false;
|
||||
const runIds = new Set([
|
||||
previousRunId,
|
||||
nextRunId,
|
||||
|
|
@ -191,6 +199,10 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
if (
|
||||
!source ||
|
||||
!isSameSubagentRunOwner(source, selected) ||
|
||||
(preserveCompletedRun &&
|
||||
(source.execution.status !== "terminal" ||
|
||||
source.pauseReason === "sessions_yield" ||
|
||||
previousRunId === nextRunId)) ||
|
||||
(replaceParams.expected && !isSameSubagentRunOwner(source, replaceParams.expected)) ||
|
||||
(replaceParams.expected &&
|
||||
((typeof source.execution.endedAt === "number" && !replaceParams.allowEndedSource) ||
|
||||
|
|
@ -256,9 +268,13 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
const next: SubagentRunRecord = normalizeSubagentRunState({
|
||||
...source,
|
||||
runId: nextRunId,
|
||||
// Materialize the legacy run-id fallback so later replacements keep the
|
||||
// same canonical task owner after this source row is retired.
|
||||
taskRunId: source.taskRunId ?? source.runId,
|
||||
// Completed follow-ups start a new task; steer retains its task's lineage.
|
||||
taskRunId: preserveCompletedRun ? nextRunId : (source.taskRunId ?? source.runId),
|
||||
requesterTurnRunId: preserveCompletedRun ? undefined : source.requesterTurnRunId,
|
||||
requesterTurnYielded: preserveCompletedRun ? undefined : source.requesterTurnYielded,
|
||||
retireAfterRequesterTurn: preserveCompletedRun
|
||||
? undefined
|
||||
: source.retireAfterRequesterTurn,
|
||||
task: nextTask,
|
||||
generation,
|
||||
createdAt: now,
|
||||
|
|
@ -322,14 +338,22 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
});
|
||||
}
|
||||
postimages.set(nextRunId, next);
|
||||
if (previousRunId !== nextRunId) {
|
||||
if (preserveCompletedRun) {
|
||||
postimages.set(previousRunId, {
|
||||
...source,
|
||||
execution: { ...source.execution, suppressSessionEffects: true },
|
||||
});
|
||||
} else if (previousRunId !== nextRunId) {
|
||||
postimages.set(previousRunId, null);
|
||||
}
|
||||
return { value: { source, next }, postimages };
|
||||
},
|
||||
{
|
||||
runs: this.options.runs,
|
||||
assertCurrent,
|
||||
assertCurrent: () => {
|
||||
assertCurrent();
|
||||
authority?.assertCurrent();
|
||||
},
|
||||
onPublished: (postimages, value) => {
|
||||
const next = postimages.get(nextRunId);
|
||||
if (!value || !next) {
|
||||
|
|
@ -340,7 +364,14 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
next,
|
||||
replaceParams.gatewayContextResolver ?? getGatewayContextResolver(value.source),
|
||||
);
|
||||
subagentRuns.transferCompletionAuthority(value.source, next);
|
||||
if (preserveCompletedRun) {
|
||||
if (authority?.operatorAuthority) {
|
||||
subagentRuns.bindCompletionAuthority(next, authority);
|
||||
custodyTransferred = true;
|
||||
}
|
||||
} else {
|
||||
subagentRuns.transferCompletionAuthority(value.source, next);
|
||||
}
|
||||
subagentRuns.commitOwnership(next);
|
||||
replaceParams.onPublished?.(next);
|
||||
},
|
||||
|
|
@ -356,6 +387,10 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
return false;
|
||||
}
|
||||
throw error;
|
||||
} finally {
|
||||
if (!custodyTransferred) {
|
||||
authority?.release();
|
||||
}
|
||||
}
|
||||
if (!replacement) {
|
||||
return false;
|
||||
|
|
@ -365,12 +400,17 @@ export class SubagentRecoveryManager extends SubagentWaitManager {
|
|||
if (!isSameSubagentRunOwner(this.options.runs.get(nextRunId), next)) {
|
||||
return true;
|
||||
}
|
||||
replaceRequesterCronAuthorityEntry({
|
||||
previous: source,
|
||||
next,
|
||||
preserve: replaceParams.preserveRequesterSettleWake === true,
|
||||
});
|
||||
if (previousRunId !== nextRunId) {
|
||||
if (!preserveCompletedRun) {
|
||||
replaceRequesterCronAuthorityEntry({
|
||||
previous: source,
|
||||
next,
|
||||
preserve: replaceParams.preserveRequesterSettleWake === true,
|
||||
});
|
||||
}
|
||||
if (preserveCompletedRun && !source.cleanupHandled) {
|
||||
this.options.resumedRuns.delete(getSubagentRunRuntimeKey(source));
|
||||
this.options.resumeSubagentRun(previousRunId);
|
||||
} else if (!preserveCompletedRun && previousRunId !== nextRunId) {
|
||||
this.options.clearPendingLifecycleError(previousRunId);
|
||||
this.options.resumedRuns.delete(getSubagentRunRuntimeKey(source));
|
||||
if (this.shouldDeleteAttachments(source)) {
|
||||
|
|
|
|||
|
|
@ -96,8 +96,8 @@ export async function discardSuspendedPendingFinalDelivery(params: {
|
|||
childSessionKey: entry.childSessionKey,
|
||||
requesterSessionKey: entry.requesterSessionKey,
|
||||
});
|
||||
if (entry.cleanup === "delete" || !entry.retainAttachmentsOnKeep) {
|
||||
await safeRemoveAttachmentsDir(entry, isCurrent);
|
||||
if ((entry.cleanup === "delete" || !entry.retainAttachmentsOnKeep) && isHookCurrent()) {
|
||||
await safeRemoveAttachmentsDir(entry, isHookCurrent);
|
||||
}
|
||||
assertCurrent();
|
||||
if (
|
||||
|
|
|
|||
|
|
@ -119,7 +119,7 @@ export function createSubagentSweeperHarness(
|
|||
resumeRequesterSettleWake,
|
||||
startSubagentAnnounceCleanupFlow: vi.fn(() => true),
|
||||
completeCleanupBookkeeping,
|
||||
isEndedHookOwnerCurrent: (runId, selected) =>
|
||||
isCleanupOwnerCurrent: (runId, selected) =>
|
||||
isSameSubagentRunOwner(runs.get(runId), selected) || !runs.has(runId),
|
||||
sessionEffectsHostCurrent: (selected) => selected.execution.suppressSessionEffects !== true,
|
||||
shouldSuppressSessionEffects: async (selected) =>
|
||||
|
|
|
|||
|
|
@ -77,7 +77,7 @@ export function createSubagentRegistrySweeper(params: {
|
|||
resumeRequesterSettleWake: SubagentLifecycleController["resumeRequesterSettleWake"];
|
||||
startSubagentAnnounceCleanupFlow: SubagentLifecycleController["startSubagentAnnounceCleanupFlow"];
|
||||
completeCleanupBookkeeping: SubagentLifecycleController["completeCleanupBookkeeping"];
|
||||
isEndedHookOwnerCurrent: SubagentLifecycleController["isEndedHookOwnerCurrent"];
|
||||
isCleanupOwnerCurrent: SubagentLifecycleController["isCleanupOwnerCurrent"];
|
||||
sessionEffectsHostCurrent: SubagentLifecycleController["sessionEffectsHostCurrent"];
|
||||
shouldSuppressSessionEffects: SubagentLifecycleController["shouldSuppressSessionEffects"];
|
||||
discardTerminalDelivery: typeof SubagentLifecycleController.discardTerminalDelivery;
|
||||
|
|
@ -303,7 +303,7 @@ export function createSubagentRegistrySweeper(params: {
|
|||
clearPendingLifecycleTimeout: params.clearPendingLifecycleTimeout,
|
||||
discardTerminalDelivery: params.discardTerminalDelivery,
|
||||
completeCleanupBookkeeping: params.completeCleanupBookkeeping,
|
||||
isCurrent: () => params.isEndedHookOwnerCurrent(runId, entry),
|
||||
isCurrent: () => params.isCleanupOwnerCurrent(runId, entry),
|
||||
sessionEffectsHostCurrent: params.sessionEffectsHostCurrent,
|
||||
shouldSuppressSessionEffects: params.shouldSuppressSessionEffects,
|
||||
shouldEmitEndedHookForRun: params.shouldEmitEndedHookForRun,
|
||||
|
|
|
|||
|
|
@ -157,11 +157,11 @@ it("hydrates a cold durable source before advancing its replacement generation",
|
|||
).toBe(true);
|
||||
expect(subagentRuns.get("cold-successor")).toMatchObject({
|
||||
generation: 3,
|
||||
taskRunId: "cold-original",
|
||||
taskRunId: "cold-successor",
|
||||
execution: { status: "running" },
|
||||
});
|
||||
const stored = loadSubagentRegistryFromSqlite();
|
||||
expect(stored.has("cold-predecessor")).toBe(false);
|
||||
expect(stored.get("cold-predecessor")?.execution.suppressSessionEffects).toBe(true);
|
||||
expect(stored.get("cold-successor")?.generation).toBe(3);
|
||||
});
|
||||
|
||||
|
|
@ -383,7 +383,7 @@ it.each(["end", "error"] as const)(
|
|||
}),
|
||||
).toBe(true);
|
||||
const successor = subagentRuns.get("timeout-successor")!;
|
||||
expect(successor.taskRunId).toBe(previous.runId);
|
||||
expect(successor.taskRunId).toBe(successor.runId);
|
||||
expect(getAgentRunContext(previous.runId)).toBe(owner);
|
||||
expect(getAgentRunContextOwnerStatus(previous.runId, claimId, lifecycleGeneration)).toBe(
|
||||
"active",
|
||||
|
|
@ -672,7 +672,14 @@ it("admits a child follow-up while its predecessor's browser cleanup is still pe
|
|||
}),
|
||||
).resolves.toBe(true);
|
||||
const stored = loadSubagentRegistryFromSqlite();
|
||||
expect(stored.has("browser-cleanup-predecessor")).toBe(false);
|
||||
expect(stored.get("browser-cleanup-predecessor")).toMatchObject({
|
||||
task: "Finish browser work",
|
||||
execution: {
|
||||
status: "terminal",
|
||||
outcome: { status: "ok" },
|
||||
suppressSessionEffects: true,
|
||||
},
|
||||
});
|
||||
expect(stored.get("browser-cleanup-successor")).toMatchObject({
|
||||
task: "Continue with the next task",
|
||||
execution: { status: "running" },
|
||||
|
|
@ -904,12 +911,19 @@ it.each([
|
|||
);
|
||||
} else {
|
||||
expect(await followup).toEqual({ value: true });
|
||||
expect(loadSubagentRegistryFromSqlite().get("after-ended-hook")).toMatchObject({
|
||||
const stored = loadSubagentRegistryFromSqlite();
|
||||
expect(stored.get("after-ended-hook")).toMatchObject({
|
||||
task: "follow-up work",
|
||||
generation: original.generation! + 1,
|
||||
execution: { status: "running" },
|
||||
});
|
||||
expect(loadSubagentRegistryFromSqlite().has(original.runId)).toBe(false);
|
||||
expect(stored.get(original.runId)).toMatchObject({
|
||||
task: "original work",
|
||||
generation: original.generation,
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
...(lateCleanup ? { cleanupCompletedAt: 3 } : { endedHookEmittedAt: expect.any(Number) }),
|
||||
...(lateWrite ? { label: "prepared" } : {}),
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
releaseFirst.resolve();
|
||||
|
|
|
|||
|
|
@ -460,6 +460,8 @@ describe("registered completion source custody", () => {
|
|||
"revoke",
|
||||
"gateway-close",
|
||||
"replace",
|
||||
"completed-followup-original-revoked",
|
||||
"completed-followup-successor-revoked",
|
||||
"release-rejected",
|
||||
"registration-rejected",
|
||||
"cancelled-by-another-operator",
|
||||
|
|
@ -620,6 +622,87 @@ describe("registered completion source custody", () => {
|
|||
/authority/,
|
||||
);
|
||||
await releaseSubagentRun("successor");
|
||||
} else if (
|
||||
ending === "completed-followup-original-revoked" ||
|
||||
ending === "completed-followup-successor-revoked"
|
||||
) {
|
||||
await updateRun(entry.runId, (draft) => {
|
||||
draft.execution = { status: "terminal", endedAt: 1, outcome: { status: "ok" } };
|
||||
draft.delivery = { status: "pending" };
|
||||
});
|
||||
const successorClient = createOperatorClient({
|
||||
profileName: "followup-owner",
|
||||
scopes: ["operator.write"],
|
||||
});
|
||||
const successorRevoked = new AbortController();
|
||||
const successorSource = (await operatorCapture.captureGatewayOperatorRunAuthority({
|
||||
client: successorClient,
|
||||
context,
|
||||
sourceAuthority: {
|
||||
signal: successorRevoked.signal,
|
||||
assertCurrent: () => successorRevoked.signal.throwIfAborted(),
|
||||
},
|
||||
}))!;
|
||||
successorClient.internal = { operatorRunAuthority: successorSource.authority };
|
||||
try {
|
||||
expect(
|
||||
await withPluginRuntimeGatewayRequestScope(
|
||||
{
|
||||
client: successorClient,
|
||||
context,
|
||||
resolveGatewayContext,
|
||||
isWebchatConnect: () => false,
|
||||
},
|
||||
() =>
|
||||
replaceSubagentRunAfterSteerCore({
|
||||
previousRunId: entry.runId,
|
||||
nextRunId: "successor",
|
||||
expected: subagentRuns.get(entry.runId),
|
||||
allowEndedSource: true,
|
||||
preserveCompletedRun: true,
|
||||
task: "second turn",
|
||||
}),
|
||||
),
|
||||
).toBe(true);
|
||||
successorSource.release();
|
||||
const retained = subagentRuns.get(entry.runId)!;
|
||||
const successor = subagentRuns.get("successor")!;
|
||||
const stored = loadSubagentRegistryFromSqlite();
|
||||
expect(stored.get(entry.runId)).toMatchObject({
|
||||
execution: { status: "terminal", suppressSessionEffects: true },
|
||||
delivery: { status: "pending" },
|
||||
requesterTurnRunId: "parent",
|
||||
});
|
||||
expect(stored.get(successor.runId)).toMatchObject({
|
||||
taskRunId: successor.runId,
|
||||
task: "second turn",
|
||||
execution: { status: "running" },
|
||||
});
|
||||
expect(successor.requesterTurnRunId).toBeUndefined();
|
||||
expect(successor.requesterTurnYielded).toBeUndefined();
|
||||
const completionSource = (run: SubagentRunRecord) =>
|
||||
subagentRuns.runWithCompletionAuthority(
|
||||
run,
|
||||
() =>
|
||||
getPluginRuntimeGatewayRequestScope()?.client?.internal?.operatorRunAuthority
|
||||
?.source,
|
||||
);
|
||||
// Suppressing old session effects must not retire its pending result's custody.
|
||||
expect(completionSource(retained)).toBe(source.authority.source);
|
||||
expect(completionSource(successor)).toBe(successorSource.authority.source);
|
||||
const revokeOriginal = ending === "completed-followup-original-revoked";
|
||||
(revokeOriginal ? revoked : successorRevoked).abort(new Error("operator revoked"));
|
||||
expect(() => completionSource(revokeOriginal ? retained : successor)).toThrow(
|
||||
/authority/,
|
||||
);
|
||||
expect(completionSource(revokeOriginal ? successor : retained)).toBe(
|
||||
revokeOriginal ? successorSource.authority.source : source.authority.source,
|
||||
);
|
||||
await releaseSubagentRun(entry.runId);
|
||||
await releaseSubagentRun(successor.runId);
|
||||
} finally {
|
||||
successorSource.release();
|
||||
}
|
||||
} else if (ending === "release-rejected") {
|
||||
rejectNextRegistryWrite("write refused");
|
||||
await expect(releaseSubagentRun(entry.runId)).rejects.toThrow("write refused");
|
||||
|
|
|
|||
|
|
@ -498,6 +498,19 @@ function retireSupersededSubagentRun(runId: string, expected: SubagentRunRecord)
|
|||
if (!entry || !isSameSubagentRunOwner(entry, expected)) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
if (
|
||||
entry.execution.status === "terminal" &&
|
||||
entry.expectsCompletionMessage === true &&
|
||||
entry.suppressCompletionDelivery !== true &&
|
||||
!entry.killIntent &&
|
||||
!entry.killReconciliation &&
|
||||
!entry.requesterTurnRunId &&
|
||||
!entry.requesterSettleWake &&
|
||||
!entry.cleanupCompletedAt
|
||||
) {
|
||||
startSubagentAnnounceCleanupFlow(runId, entry);
|
||||
return Promise.resolve();
|
||||
}
|
||||
const wake = entry.requesterSettleWake;
|
||||
const cohort = [...getSubagentRunsForChildSession(entry.childSessionKey)].filter((candidate) =>
|
||||
entry.requesterTurnRunId
|
||||
|
|
@ -507,7 +520,7 @@ function retireSupersededSubagentRun(runId: string, expected: SubagentRunRecord)
|
|||
);
|
||||
const isCurrent = () =>
|
||||
isSameSubagentRunOwner(subagentRuns.get(runId), entry) &&
|
||||
isRequesterCompletionCohortCurrent(entry, cohort, getLatestLiveSubagentRunByChildSessionKey);
|
||||
isRequesterCompletionCohortCurrent(entry, getLatestLiveSubagentRunByChildSessionKey);
|
||||
if (
|
||||
isCurrent() &&
|
||||
entry.expectsCompletionMessage === true &&
|
||||
|
|
@ -551,7 +564,7 @@ const subagentSweeper = createSubagentRegistrySweeper({
|
|||
resumeRequesterSettleWake,
|
||||
startSubagentAnnounceCleanupFlow,
|
||||
completeCleanupBookkeeping,
|
||||
isEndedHookOwnerCurrent: subagentLifecycleController.isEndedHookOwnerCurrent,
|
||||
isCleanupOwnerCurrent: subagentLifecycleController.isCleanupOwnerCurrent,
|
||||
sessionEffectsHostCurrent: (entry) =>
|
||||
subagentLifecycleController.sessionEffectsHostCurrent(entry),
|
||||
shouldSuppressSessionEffects: (entry, effects) =>
|
||||
|
|
|
|||
|
|
@ -157,10 +157,9 @@ export function hasRequesterCompletionCohort(entry: SubagentRunRecord): boolean
|
|||
);
|
||||
}
|
||||
|
||||
/** A frozen completion cohort can own distinct tasks that share one child session. */
|
||||
/** A newer task cannot revoke another task's exact completion custody. */
|
||||
export function isRequesterCompletionCohortCurrent(
|
||||
entry: SubagentRunRecord,
|
||||
cohort: readonly SubagentRunRecord[],
|
||||
latestForSession: (
|
||||
sessionKey: string,
|
||||
matches?: (candidate: SubagentRunRecord) => boolean,
|
||||
|
|
@ -171,26 +170,9 @@ export function isRequesterCompletionCohortCurrent(
|
|||
entry.childSessionKey,
|
||||
(candidate) => (candidate.taskRunId ?? candidate.runId) === taskRunId,
|
||||
);
|
||||
if (
|
||||
entry.killReconciliation?.supersededAt !== undefined ||
|
||||
(task && compareSubagentRunGeneration(task, entry) > 0)
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
const latest = latestForSession(entry.childSessionKey);
|
||||
return (
|
||||
!latest ||
|
||||
compareSubagentRunGeneration(latest, entry) <= 0 ||
|
||||
cohort.some(
|
||||
(candidate) =>
|
||||
candidate.runId === latest.runId &&
|
||||
candidate.generation === latest.generation &&
|
||||
candidate.requesterSessionKey === entry.requesterSessionKey &&
|
||||
candidate.requesterAgentId === entry.requesterAgentId &&
|
||||
candidate.requesterStorePath === entry.requesterStorePath &&
|
||||
candidate.requesterTurnRunId === entry.requesterTurnRunId &&
|
||||
(candidate.taskRunId ?? candidate.runId) !== taskRunId,
|
||||
)
|
||||
entry.killReconciliation?.supersededAt === undefined &&
|
||||
(!task || compareSubagentRunGeneration(task, entry) <= 0)
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -2,17 +2,28 @@
|
|||
import path from "node:path";
|
||||
import { MAX_TIMER_TIMEOUT_MS } from "@openclaw/normalization-core/number-coercion";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { createDeferred } from "../../../test/helpers/promise.js";
|
||||
import { createDeferred, withinTest } from "../../../test/helpers/promise.js";
|
||||
import * as announceDelivery from "../../agents/subagents/announce/subagent-announce-delivery.js";
|
||||
import { sourceOwnerChangedResult } from "../../agents/subagents/announce/subagent-announce-dispatch.js";
|
||||
import { subagentRuns } from "../../agents/subagents/registry/subagent-registry-memory.js";
|
||||
import { restoreSubagentRunsFromDisk } from "../../agents/subagents/registry/subagent-registry-persistence.js";
|
||||
import { subscribeSubagentRunChanges } from "../../agents/subagents/registry/subagent-registry-publication.js";
|
||||
import {
|
||||
getSubagentRunByRunId,
|
||||
initSubagentRegistry,
|
||||
resumeSubagentRun,
|
||||
registerSubagentRun,
|
||||
} from "../../agents/subagents/registry/subagent-registry.js";
|
||||
import { settleSubagentRegistryPersistenceWork } from "../../agents/subagents/registry/subagent-registry.persistence.test-support.js";
|
||||
import { upsertSubagentRunRowInDatabase } from "../../agents/subagents/registry/subagent-registry.store.kernel.js";
|
||||
import { loadSubagentRegistryFromSqlite } from "../../agents/subagents/registry/subagent-registry.store.sqlite.js";
|
||||
import {
|
||||
addSubagentRunForTests,
|
||||
getSubagentRunByChildSessionKey,
|
||||
resetSubagentRegistryForTests,
|
||||
testing as registryTesting,
|
||||
} from "../../agents/subagents/registry/subagent-registry.test-helpers.js";
|
||||
import * as sessionAccessor from "../../config/sessions/session-accessor.js";
|
||||
import { runOpenClawStateWriteTransaction } from "../../state/openclaw-state-db.js";
|
||||
import * as admissionController from "../agent-turn/agent-admission-controller.js";
|
||||
import { resolveAgentRunExpiresAtMs } from "../chat-abort.js";
|
||||
|
|
@ -398,3 +409,299 @@ describe("gateway agent follow-up activity", () => {
|
|||
},
|
||||
);
|
||||
});
|
||||
|
||||
describe("gateway agent completed-child delivery", () => {
|
||||
afterEach(describe0AfterEach0);
|
||||
|
||||
for (const boundary of [
|
||||
"transcript read",
|
||||
"delivery admission",
|
||||
"restored transcript read",
|
||||
] as const) {
|
||||
it(`announces both executions when follow-up starts during the first ${boundary}`, async ({
|
||||
signal,
|
||||
}) => {
|
||||
await withPluginSubagentTestState("openclaw-followup-delivery-", async ({ stateDir }) => {
|
||||
const registry = await vi.importActual<
|
||||
typeof import("../../agents/subagents/registry/subagent-registry.js")
|
||||
>("../../agents/subagents/registry/subagent-registry.js");
|
||||
const reads = await vi.importActual<
|
||||
typeof import("../../agents/subagents/registry/subagent-registry-read.js")
|
||||
>("../../agents/subagents/registry/subagent-registry-read.js");
|
||||
const announce = await vi.importActual<
|
||||
typeof import("../../agents/subagents/announce/subagent-announce.js")
|
||||
>("../../agents/subagents/announce/subagent-announce.js");
|
||||
mocks.getLatestSubagentRunByChildSessionKey.mockImplementation(
|
||||
reads.getLatestSubagentRunByChildSessionKey,
|
||||
);
|
||||
mocks.replaceSubagentRunAfterSteer.mockImplementation(
|
||||
registry.replaceSubagentRunAfterSteerCore,
|
||||
);
|
||||
const childSessionKey = "agent:main:dashboard:kept-child";
|
||||
const requesterSessionKey = "agent:main:main";
|
||||
const firstRunId = "first-execution";
|
||||
const secondRunId = "second-execution";
|
||||
const storePath = path.join(stateDir, "agents", "main", "sessions", "sessions.json");
|
||||
const entry = {
|
||||
sessionId: "kept-child-session",
|
||||
lifecycleRevision: "kept-child-revision",
|
||||
spawnedBy: requesterSessionKey,
|
||||
updatedAt: Date.now(),
|
||||
};
|
||||
sessionAccessor.ensureSessionEntrySync(
|
||||
{ agentId: "main", storePath, sessionKey: childSessionKey },
|
||||
entry,
|
||||
);
|
||||
sessionAccessor.ensureSessionEntrySync(
|
||||
{ agentId: "main", storePath, sessionKey: requesterSessionKey },
|
||||
{
|
||||
sessionId: "parent-session",
|
||||
lifecycleRevision: "parent-revision",
|
||||
updatedAt: Date.now(),
|
||||
},
|
||||
);
|
||||
mocks.userTurnStorePath = storePath;
|
||||
mocks.loadSessionEntry.mockReturnValue({
|
||||
cfg: {},
|
||||
storePath,
|
||||
entry,
|
||||
canonicalKey: childSessionKey,
|
||||
});
|
||||
mocks.updateSessionStore.mockResolvedValue(undefined);
|
||||
await addSubagentRunForTests({
|
||||
runId: firstRunId,
|
||||
childSessionKey,
|
||||
childSessionIdentity: entry,
|
||||
requesterSessionKey,
|
||||
requesterDisplayKey: requesterSessionKey,
|
||||
task: "First task",
|
||||
cleanup: "keep",
|
||||
spawnMode: "session",
|
||||
expectsCompletionMessage: true,
|
||||
completionTarget: "parent",
|
||||
generation: 1,
|
||||
createdAt: Date.now() - 20,
|
||||
execution: {
|
||||
status: "terminal",
|
||||
startedAt: Date.now() - 20,
|
||||
endedAt: Date.now() - 10,
|
||||
outcome: { status: "ok" },
|
||||
transcriptTarget: {
|
||||
agentId: "main",
|
||||
storePath,
|
||||
sessionKey: childSessionKey,
|
||||
sessionId: entry.sessionId,
|
||||
},
|
||||
},
|
||||
completion: {
|
||||
required: true,
|
||||
terminalReply: { disposition: "visible", text: "FIRST_ONLY" },
|
||||
},
|
||||
delivery: { status: "pending" },
|
||||
});
|
||||
const paused = createDeferred();
|
||||
const release = createDeferred();
|
||||
const firstAnnounceFinished = createDeferred();
|
||||
const secondAnnounceFinished = createDeferred();
|
||||
const receipts = new Map([
|
||||
[firstRunId, createDeferred()],
|
||||
[secondRunId, createDeferred()],
|
||||
]);
|
||||
const unsubscribe = subscribeSubagentRunChanges("persistence", () => {
|
||||
for (const [runId, receipt] of receipts) {
|
||||
if (getSubagentRunByRunId(runId)?.cleanupCompletedAt) {
|
||||
receipt.resolve();
|
||||
}
|
||||
}
|
||||
});
|
||||
mocks.registryAnnounce.mockImplementation(async (params) => {
|
||||
try {
|
||||
return await announce.runSubagentAnnounceFlow(params);
|
||||
} finally {
|
||||
(params.childRunId === firstRunId
|
||||
? firstAnnounceFinished
|
||||
: secondAnnounceFinished
|
||||
).resolve();
|
||||
}
|
||||
});
|
||||
const transcript = vi
|
||||
.spyOn(sessionAccessor, "findTranscriptEvent")
|
||||
.mockImplementation(async (_target, match) => {
|
||||
expect(match.kind).toBe("visible-final");
|
||||
if (match.kind !== "visible-final") {
|
||||
throw new Error("Expected an exact-run visible transcript read");
|
||||
}
|
||||
if (match.runId === firstRunId && boundary !== "delivery admission") {
|
||||
paused.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
return {
|
||||
event: {
|
||||
type: "message",
|
||||
message: {
|
||||
role: "assistant",
|
||||
stopReason: "stop",
|
||||
content: [
|
||||
{
|
||||
type: "text",
|
||||
text: match.runId === firstRunId ? "FIRST_ONLY" : "SECOND_ONLY",
|
||||
},
|
||||
],
|
||||
__openclaw: { runId: match.runId },
|
||||
},
|
||||
},
|
||||
};
|
||||
});
|
||||
const handoffs: Parameters<typeof announceDelivery.deliverSubagentAnnouncement>[0][] = [];
|
||||
const deliver = vi
|
||||
.spyOn(announceDelivery, "deliverSubagentAnnouncement")
|
||||
.mockImplementation(async (params) => {
|
||||
if (params.sourceRunId === firstRunId && boundary === "delivery admission") {
|
||||
paused.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
if (params.isSourceSessionEffectsAllowed?.() === false) {
|
||||
return sourceOwnerChangedResult();
|
||||
}
|
||||
handoffs.push(params);
|
||||
return {
|
||||
delivered: true,
|
||||
disposition: "delivered",
|
||||
path: "direct",
|
||||
deliveredAt: Date.now(),
|
||||
};
|
||||
});
|
||||
const provider = createDeferred<{
|
||||
payloads: { text: string }[];
|
||||
meta: { durationMs: number };
|
||||
}>();
|
||||
const secondWait = createDeferred<{
|
||||
status: "ok";
|
||||
startedAt: number;
|
||||
endedAt: number;
|
||||
terminalReply: { disposition: "visible"; text: string };
|
||||
}>();
|
||||
const commandStarted = createDeferred();
|
||||
mocks.agentCommand.mockImplementationOnce(() => {
|
||||
commandStarted.resolve();
|
||||
return provider.promise;
|
||||
});
|
||||
mocks.registryCallGateway.mockImplementation(async (request) => {
|
||||
if (request.method !== "agent.wait") {
|
||||
throw new Error(`Unexpected child-session effect: ${request.method}`);
|
||||
}
|
||||
expect(request.params).toMatchObject({ runId: secondRunId });
|
||||
return await secondWait.promise;
|
||||
});
|
||||
const context = makeContext();
|
||||
const client = backendGatewayClient();
|
||||
const terminal = createDeferred();
|
||||
const respond = vi.fn((ok, payload) => {
|
||||
if (!ok || payload?.status === "ok" || payload?.status === "error") {
|
||||
terminal.resolve();
|
||||
}
|
||||
});
|
||||
let admitted = false;
|
||||
try {
|
||||
if (boundary === "restored transcript read") {
|
||||
await restoreSubagentRunsFromDisk({ runs: subagentRuns });
|
||||
}
|
||||
resumeSubagentRun(firstRunId);
|
||||
await withinTest(paused.promise, signal);
|
||||
await invokeAgent(
|
||||
{
|
||||
sessionKey: childSessionKey,
|
||||
message: "Second task",
|
||||
idempotencyKey: secondRunId,
|
||||
},
|
||||
{ context, respond, reqId: secondRunId, client, flushDispatch: false },
|
||||
);
|
||||
admitted = true;
|
||||
await withinTest(commandStarted.promise, signal);
|
||||
expect(getSubagentRunByRunId(secondRunId)).toMatchObject({
|
||||
generation: 2,
|
||||
task: "Second task",
|
||||
execution: { status: "running" },
|
||||
});
|
||||
if (boundary === "restored transcript read") {
|
||||
expect(loadSubagentRegistryFromSqlite().get(firstRunId)).toMatchObject({
|
||||
runId: firstRunId,
|
||||
generation: 1,
|
||||
execution: { status: "terminal", outcome: { status: "ok" } },
|
||||
completion: { terminalReply: { disposition: "visible", text: "FIRST_ONLY" } },
|
||||
delivery: { status: "pending" },
|
||||
});
|
||||
}
|
||||
release.resolve();
|
||||
await withinTest(firstAnnounceFinished.promise, signal);
|
||||
expect(handoffs).toHaveLength(1);
|
||||
expect(handoffs[0]).toMatchObject({
|
||||
sourceRunId: firstRunId,
|
||||
directIdempotencyKey: `announce:v1:${childSessionKey}:${firstRunId}`,
|
||||
internalEvents: [{ taskLabel: "First task", status: "ok", result: "FIRST_ONLY" }],
|
||||
});
|
||||
await withinTest(receipts.get(firstRunId)!.promise, signal);
|
||||
expect(getSubagentRunByRunId(firstRunId)?.delivery).toMatchObject({
|
||||
status: "delivered",
|
||||
deliveredAt: expect.any(Number),
|
||||
});
|
||||
expect(getSubagentRunByRunId(secondRunId)).toMatchObject({
|
||||
execution: { status: "running" },
|
||||
delivery: { status: "pending" },
|
||||
cleanupHandled: false,
|
||||
});
|
||||
expect(
|
||||
sessionAccessor.loadSessionEntryReadOnly({
|
||||
agentId: "main",
|
||||
storePath,
|
||||
sessionKey: childSessionKey,
|
||||
}),
|
||||
).toMatchObject(entry);
|
||||
provider.resolve({ payloads: [{ text: "SECOND_ONLY" }], meta: { durationMs: 1 } });
|
||||
secondWait.resolve({
|
||||
status: "ok",
|
||||
startedAt: Date.now() - 1,
|
||||
endedAt: Date.now(),
|
||||
terminalReply: { disposition: "visible", text: "SECOND_ONLY" },
|
||||
});
|
||||
await withinTest(terminal.promise, signal);
|
||||
await withinTest(secondAnnounceFinished.promise, signal);
|
||||
await withinTest(receipts.get(secondRunId)!.promise, signal);
|
||||
expect(handoffs).toHaveLength(2);
|
||||
expect(handoffs[1]).toMatchObject({
|
||||
sourceRunId: secondRunId,
|
||||
directIdempotencyKey: `announce:v1:${childSessionKey}:${secondRunId}`,
|
||||
internalEvents: [{ taskLabel: "Second task", status: "ok", result: "SECOND_ONLY" }],
|
||||
});
|
||||
expect(getSubagentRunByRunId(secondRunId)?.delivery).toMatchObject({
|
||||
status: "delivered",
|
||||
deliveredAt: expect.any(Number),
|
||||
});
|
||||
resumeSubagentRun(firstRunId);
|
||||
resumeSubagentRun(secondRunId);
|
||||
await registryTesting.sweepOnceForTests();
|
||||
await settleSubagentRegistryPersistenceWork();
|
||||
expect(handoffs).toHaveLength(2);
|
||||
} finally {
|
||||
release.resolve();
|
||||
provider.resolve({ payloads: [{ text: "SECOND_ONLY" }], meta: { durationMs: 1 } });
|
||||
secondWait.resolve({
|
||||
status: "ok",
|
||||
startedAt: Date.now() - 1,
|
||||
endedAt: Date.now(),
|
||||
terminalReply: { disposition: "visible", text: "SECOND_ONLY" },
|
||||
});
|
||||
if (admitted) {
|
||||
await terminal.promise;
|
||||
}
|
||||
await resetSubagentRegistryForTests({ persist: false });
|
||||
unsubscribe();
|
||||
transcript.mockRestore();
|
||||
deliver.mockRestore();
|
||||
mocks.getLatestSubagentRunByChildSessionKey.mockReset();
|
||||
mocks.replaceSubagentRunAfterSteer.mockReset();
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
});
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ export function expectSubagentFollowupReactivation(params: {
|
|||
expect(params.replaceSubagentRunAfterSteerMock).toHaveBeenCalledWith({
|
||||
previousRunId: "run-old",
|
||||
nextRunId: "run-new",
|
||||
preserveCompletedRun: true,
|
||||
assertCurrent: expect.any(Function),
|
||||
runTimeoutSeconds: 0,
|
||||
...(params.task ? { task: params.task } : {}),
|
||||
|
|
|
|||
|
|
@ -1,7 +1,8 @@
|
|||
import { randomUUID } from "node:crypto";
|
||||
import { expectDefined } from "@openclaw/normalization-core/expect";
|
||||
import { isRecord } from "@openclaw/normalization-core/record-coerce";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { createDeferred } from "../../test/helpers/promise.js";
|
||||
import { createDeferred, withinTest } from "../../test/helpers/promise.js";
|
||||
import { createAgentHarnessCompletionScope } from "../agents/agent-harness-completion-scope.js";
|
||||
import { buildAnnounceIdempotencyKey } from "../agents/announce-idempotency.js";
|
||||
import type { AgentCommandOpts } from "../agents/command/types.js";
|
||||
|
|
@ -20,21 +21,28 @@ import {
|
|||
} from "../agents/sessions/agent-session-loop-correctness.test-support.js";
|
||||
import { createResourceLoader } from "../agents/sessions/agent-session-loop-resource-loader.test-support.js";
|
||||
import { SessionManager } from "../agents/sessions/session-manager.js";
|
||||
import { resumeSubagentRun } from "../agents/subagents/registry/subagent-registry.js";
|
||||
import { loadSubagentRegistryFromSqlite } from "../agents/subagents/registry/subagent-registry.store.sqlite.js";
|
||||
import * as sessionAccessor from "../config/sessions/session-accessor.js";
|
||||
import { listSessionPendingInputs } from "../config/sessions/session-accessor.pending-inputs.js";
|
||||
import { readMessageIdempotencyKey } from "../config/sessions/transcript-message-identity.js";
|
||||
import { createAssistantMessageEventStream } from "../llm/utils/event-stream.js";
|
||||
import {
|
||||
captureAgentHarnessCompletionCustody,
|
||||
deliverAgentHarnessCompletion,
|
||||
type AgentHarnessCompletionCustody,
|
||||
} from "../plugin-sdk/agent-harness-completion.js";
|
||||
import { withPluginRuntimeGatewayRequestScope } from "../plugins/runtime/gateway-request-scope.js";
|
||||
import {
|
||||
getPluginRuntimeGatewayRequestScope,
|
||||
withPluginRuntimeGatewayRequestScope,
|
||||
} from "../plugins/runtime/gateway-request-scope.js";
|
||||
import { tryBeginGatewayRootWorkAdmission } from "../process/gateway-work-admission.js";
|
||||
import * as userTurnTranscript from "../sessions/user-turn-transcript.js";
|
||||
import { captureGatewayOperatorRunAuthority } from "./operator-run-authority.js";
|
||||
import type { GatewayRequestContext } from "./server-methods/types.js";
|
||||
import { createOperatorClient } from "./server-plugin-in-process-dispatch.test-support.js";
|
||||
import { startGatewayServerHarness, type GatewayServerHarness } from "./server.e2e-ws-harness.js";
|
||||
import { createRegisteredCompletionPair } from "./server.subagent-completion-authority.test-support.js";
|
||||
import { loadSessionEntry } from "./session-utils.js";
|
||||
import {
|
||||
agentCommandMock,
|
||||
|
|
@ -175,6 +183,125 @@ describe("native completion final-effect authority", () => {
|
|||
registerAgentSessionLoopTestLifecycle();
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
|
||||
it.for(["neither", "predecessor", "successor"] as const)(
|
||||
"keeps registered turn delivery independent when %s source is revoked in flight",
|
||||
async (revokedSource, { signal }) => {
|
||||
await prepareGatewayReplyRuntimeForTest();
|
||||
const context = kernel.gatewayRequestContext;
|
||||
const pair = await createRegisteredCompletionPair(context);
|
||||
const createGate = () => ({ entered: createDeferred(), resume: createDeferred() });
|
||||
const gates = [createGate(), createGate()] as const;
|
||||
const sources: unknown[] = [];
|
||||
const attempts: string[] = [];
|
||||
let requesterCommands = 0;
|
||||
const release = () => gates.forEach((gate) => gate.resume.resolve());
|
||||
signal.addEventListener("abort", release, { once: true });
|
||||
const requestWork = vi.spyOn(context, "trackExecution");
|
||||
const settleRequests = () =>
|
||||
Promise.allSettled(
|
||||
requestWork.mock.results.flatMap((result) =>
|
||||
result.type === "return" ? [result.value] : [],
|
||||
),
|
||||
);
|
||||
const stage = sessionAccessor.stageSessionPendingInput;
|
||||
const stageSpy = vi
|
||||
.spyOn(sessionAccessor, "stageSessionPendingInput")
|
||||
.mockImplementation(async (...args) => {
|
||||
const index = pair.runs.findIndex(
|
||||
(run) =>
|
||||
args[0].sessionKey === pair.requesterScope.sessionKey &&
|
||||
readMessageIdempotencyKey(args[1].message) === `${run.idempotencyKey}:user`,
|
||||
);
|
||||
if (index === 0 || index === 1) {
|
||||
sources[index] =
|
||||
getPluginRuntimeGatewayRequestScope()?.client?.internal?.operatorRunAuthority?.source;
|
||||
attempts.push(pair.runs[index].runId);
|
||||
gates[index].entered.resolve();
|
||||
await gates[index].resume.promise;
|
||||
}
|
||||
return stage(...args);
|
||||
});
|
||||
agentCommandMock.mockImplementation(async (input) => {
|
||||
const command = input as AgentCommandOpts;
|
||||
const child = pair.handleChildCommand(command);
|
||||
if (child) {
|
||||
return child;
|
||||
}
|
||||
requesterCommands += 1;
|
||||
const recorder = expectDefined(
|
||||
command.userTurnTranscriptRecorder,
|
||||
"Expected real registered completion input recorder",
|
||||
);
|
||||
expect(await recorder.persistApproved()).toMatchObject({ appended: true });
|
||||
return { payloads: [{ text: "Child received", mediaUrl: null }], meta: { durationMs: 1 } };
|
||||
});
|
||||
try {
|
||||
await pair.complete(0);
|
||||
await withinTest(
|
||||
reachBoundary(gates[0].entered.promise, pair.runs[0].settled.promise),
|
||||
signal,
|
||||
);
|
||||
await pair.admitSuccessor();
|
||||
expect(attempts).toEqual([pair.runs[0].runId]);
|
||||
await pair.complete(1);
|
||||
await withinTest(
|
||||
reachBoundary(gates[1].entered.promise, pair.runs[1].settled.promise),
|
||||
signal,
|
||||
);
|
||||
const revokedIndex =
|
||||
revokedSource === "neither" ? -1 : revokedSource === "predecessor" ? 0 : 1;
|
||||
if (revokedIndex !== -1) {
|
||||
pair.runs[revokedIndex].revoked.abort(new Error("operator completion authority revoked"));
|
||||
}
|
||||
release();
|
||||
await pair.settle();
|
||||
await settleRequests();
|
||||
expect(attempts.toSorted()).toEqual(pair.runs.map((run) => run.runId).toSorted());
|
||||
const allowed = pair.runs.filter((_, index) => index !== revokedIndex);
|
||||
expect(requesterCommands).toBe(allowed.length);
|
||||
expect(agentCommandMock).toHaveBeenCalledTimes(2 + allowed.length);
|
||||
const events = sessionAccessor.loadTranscriptEventsSync(pair.requesterScope);
|
||||
const inputKeys = events.flatMap((event) =>
|
||||
isRecord(event) && isRecord(event.message) && event.message.role === "user"
|
||||
? [readMessageIdempotencyKey(event.message)]
|
||||
: [],
|
||||
);
|
||||
expect(inputKeys).toHaveLength(allowed.length);
|
||||
expect(inputKeys).toEqual(
|
||||
expect.arrayContaining(allowed.map((run) => `${run.idempotencyKey}:user`)),
|
||||
);
|
||||
expect(listSessionPendingInputs(pair.requesterScope).total).toBe(0);
|
||||
const stored = loadSubagentRegistryFromSqlite();
|
||||
for (const [index, run] of pair.runs.entries()) {
|
||||
// Exact source identity also catches borrowing when both sources remain live.
|
||||
expect(sources[index]).toBe(run.source.authority.source);
|
||||
if (index === revokedIndex) {
|
||||
expect(stored.get(run.runId)?.delivery?.status).not.toBe("delivered");
|
||||
expect(context.dedupe.has(`agent:${run.idempotencyKey}`)).toBe(false);
|
||||
expect(JSON.stringify(events)).not.toContain(run.result);
|
||||
} else {
|
||||
expect(stored.get(run.runId)?.delivery?.status).toBe("delivered");
|
||||
expect(context.dedupe.get(`agent:${run.idempotencyKey}`)).toMatchObject({ ok: true });
|
||||
expect(JSON.stringify(events)).toContain(run.result);
|
||||
}
|
||||
resumeSubagentRun(run.runId);
|
||||
}
|
||||
await pair.settle();
|
||||
await settleRequests();
|
||||
expect(requesterCommands).toBe(allowed.length);
|
||||
expect(agentCommandMock).toHaveBeenCalledTimes(2 + allowed.length);
|
||||
expect(sessionAccessor.loadTranscriptEventsSync(pair.requesterScope)).toEqual(events);
|
||||
} finally {
|
||||
release();
|
||||
await pair.dispose();
|
||||
await settleRequests();
|
||||
stageSpy.mockRestore();
|
||||
requestWork.mockRestore();
|
||||
signal.removeEventListener("abort", release);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it.for(["live", "operator-revoked", "requester-replaced"] as const)(
|
||||
"revalidates %s authority at real Gateway input staging",
|
||||
async (change, { signal }) => {
|
||||
|
|
|
|||
279
src/gateway/server.subagent-completion-authority.test-support.ts
Normal file
279
src/gateway/server.subagent-completion-authority.test-support.ts
Normal file
|
|
@ -0,0 +1,279 @@
|
|||
import { randomUUID } from "node:crypto";
|
||||
import { expectDefined } from "@openclaw/normalization-core/expect";
|
||||
import { awaitGateBeforeSettlement, createDeferred } from "../../test/helpers/promise.js";
|
||||
import {
|
||||
buildAnnounceIdFromChildRun,
|
||||
buildAnnounceIdempotencyKey,
|
||||
} from "../agents/announce-idempotency.js";
|
||||
import type { AgentCommandOpts } from "../agents/command/types.js";
|
||||
import { SessionManager } from "../agents/sessions/session-manager.js";
|
||||
import { subagentRuns } from "../agents/subagents/registry/subagent-registry-memory.js";
|
||||
import { subscribeSubagentRunChanges } from "../agents/subagents/registry/subagent-registry-publication.js";
|
||||
import { registerSubagentRun } from "../agents/subagents/registry/subagent-registry.js";
|
||||
import { resetSubagentRegistryForTests } from "../agents/subagents/registry/subagent-registry.test-helpers.js";
|
||||
import { createZeroUsageFixture } from "../agents/test-helpers/usage-fixtures.js";
|
||||
import * as sessionAccessor from "../config/sessions/session-accessor.js";
|
||||
import type { AssistantMessage } from "../llm/types.js";
|
||||
import { withPluginRuntimeGatewayRequestScope } from "../plugins/runtime/gateway-request-scope.js";
|
||||
import { tryBeginGatewayRootWorkAdmission } from "../process/gateway-work-admission.js";
|
||||
import { captureGatewayOperatorRunAuthority } from "./operator-run-authority.js";
|
||||
import type { GatewayRequestContext } from "./server-methods/types.js";
|
||||
import { dispatchGatewayMethodInProcess } from "./server-plugin-in-process-dispatch.js";
|
||||
import { createOperatorClient } from "./server-plugin-in-process-dispatch.test-support.js";
|
||||
import { loadSessionEntry } from "./session-utils.js";
|
||||
|
||||
/** Real registered runs share a child session but retain independent completion sources. */
|
||||
export async function createRegisteredCompletionPair(context: GatewayRequestContext) {
|
||||
const id = randomUUID();
|
||||
const requesterSessionKey = `agent:main:completion-pair:${id}`;
|
||||
const childSessionKey = `agent:main:dashboard:completion-pair:${id}`;
|
||||
const requesterSessionId = `requester-${id}`;
|
||||
const requesterLifecycleRevision = `requester-revision-${id}`;
|
||||
const childSessionId = `child-${id}`;
|
||||
const childLifecycleRevision = `child-revision-${id}`;
|
||||
await sessionAccessor.upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: requesterSessionKey },
|
||||
{
|
||||
sessionId: requesterSessionId,
|
||||
lifecycleRevision: requesterLifecycleRevision,
|
||||
updatedAt: Date.now(),
|
||||
},
|
||||
);
|
||||
await sessionAccessor.upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: childSessionKey },
|
||||
{
|
||||
sessionId: childSessionId,
|
||||
lifecycleRevision: childLifecycleRevision,
|
||||
spawnedBy: requesterSessionKey,
|
||||
spawnDepth: 1,
|
||||
updatedAt: Date.now(),
|
||||
},
|
||||
);
|
||||
const requesterScope = {
|
||||
agentId: "main",
|
||||
sessionKey: requesterSessionKey,
|
||||
sessionId: requesterSessionId,
|
||||
storePath: loadSessionEntry(requesterSessionKey, { agentId: "main" }).storePath,
|
||||
};
|
||||
const childScope = {
|
||||
agentId: "main",
|
||||
sessionKey: childSessionKey,
|
||||
sessionId: childSessionId,
|
||||
storePath: loadSessionEntry(childSessionKey, { agentId: "main" }).storePath,
|
||||
};
|
||||
const createRun = async (index: 0 | 1) => {
|
||||
const runId = `completion-pair-${id}-${index}`;
|
||||
const revoked = new AbortController();
|
||||
const client = createOperatorClient({
|
||||
profileName: `completion-pair-${id}-${index}`,
|
||||
scopes: ["operator.write"],
|
||||
});
|
||||
const source = expectDefined(
|
||||
await captureGatewayOperatorRunAuthority({
|
||||
client,
|
||||
context,
|
||||
sourceAuthority: {
|
||||
signal: revoked.signal,
|
||||
assertCurrent: () => revoked.signal.throwIfAborted(),
|
||||
},
|
||||
}),
|
||||
"Expected an independent registered completion source",
|
||||
);
|
||||
client.internal = { operatorRunAuthority: source.authority };
|
||||
return {
|
||||
runId,
|
||||
result: index === 0 ? "FIRST_ONLY" : "SECOND_ONLY",
|
||||
idempotencyKey: buildAnnounceIdempotencyKey(
|
||||
buildAnnounceIdFromChildRun({ childSessionKey, childRunId: runId }),
|
||||
),
|
||||
revoked,
|
||||
source,
|
||||
client,
|
||||
commandEntered: createDeferred(),
|
||||
releaseModel: createDeferred(),
|
||||
terminal: createDeferred(),
|
||||
settled: createDeferred(),
|
||||
};
|
||||
};
|
||||
const first = await createRun(0);
|
||||
let second: Awaited<ReturnType<typeof createRun>>;
|
||||
try {
|
||||
second = await createRun(1);
|
||||
} catch (error) {
|
||||
first.source.release();
|
||||
throw error;
|
||||
}
|
||||
const runs = [first, second] as const;
|
||||
const dispatches = new Map<0 | 1, Promise<unknown>>();
|
||||
const unsubscribe = subscribeSubagentRunChanges("persistence", () => {
|
||||
for (const run of runs) {
|
||||
const entry = subagentRuns.get(run.runId);
|
||||
if (entry?.execution.status === "terminal") {
|
||||
run.terminal.resolve();
|
||||
}
|
||||
if (
|
||||
typeof entry?.cleanupCompletedAt === "number" ||
|
||||
(entry?.cleanupHandled === false && entry.delivery?.nextAttemptAt !== undefined)
|
||||
) {
|
||||
run.settled.resolve();
|
||||
}
|
||||
}
|
||||
});
|
||||
const runScoped = async <T>(index: 0 | 1, operation: () => Promise<T>): Promise<T> => {
|
||||
const root = expectDefined(
|
||||
tryBeginGatewayRootWorkAdmission("test:registered-completion-pair"),
|
||||
"Expected completion fixture root admission",
|
||||
);
|
||||
try {
|
||||
return await root.run(() =>
|
||||
withPluginRuntimeGatewayRequestScope(
|
||||
{
|
||||
client: runs[index].client,
|
||||
context,
|
||||
resolveGatewayContext: context.resolveGatewayContext,
|
||||
isWebchatConnect: () => false,
|
||||
},
|
||||
operation,
|
||||
),
|
||||
);
|
||||
} finally {
|
||||
root.release();
|
||||
}
|
||||
};
|
||||
const startDispatch = (index: 0 | 1) => {
|
||||
const existing = dispatches.get(index);
|
||||
if (existing) {
|
||||
return existing;
|
||||
}
|
||||
const run = runs[index];
|
||||
const dispatch = runScoped(index, () =>
|
||||
dispatchGatewayMethodInProcess(
|
||||
"agent",
|
||||
{
|
||||
sessionKey: childSessionKey,
|
||||
idempotencyKey: run.runId,
|
||||
message: index === 0 ? "First task" : "Second task",
|
||||
deliver: false,
|
||||
},
|
||||
{ expectFinal: true, resolveGatewayContext: context.resolveGatewayContext },
|
||||
),
|
||||
);
|
||||
dispatches.set(index, dispatch);
|
||||
void dispatch.catch(() => {});
|
||||
return dispatch;
|
||||
};
|
||||
const admit = async (index: 0 | 1) => {
|
||||
try {
|
||||
await awaitGateBeforeSettlement(
|
||||
runs[index].commandEntered.promise,
|
||||
startDispatch(index),
|
||||
"Child dispatch settled before entering its controlled model",
|
||||
);
|
||||
} finally {
|
||||
// The admitted execution and registered completion retain their own source custody.
|
||||
runs[index].source.release();
|
||||
}
|
||||
};
|
||||
const settle = async () => {
|
||||
await Promise.all(dispatches.values());
|
||||
await Promise.all([...dispatches.keys()].map((index) => runs[index].settled.promise));
|
||||
};
|
||||
let disposal: Promise<void> | undefined;
|
||||
const dispose = () =>
|
||||
(disposal ??= (async () => {
|
||||
for (const run of runs) {
|
||||
run.releaseModel.resolve();
|
||||
run.revoked.abort(new Error("Registered completion fixture disposed"));
|
||||
}
|
||||
try {
|
||||
await Promise.allSettled(dispatches.values());
|
||||
await Promise.all(
|
||||
[...dispatches.keys()].flatMap((index) =>
|
||||
subagentRuns.get(runs[index].runId)?.execution.status === "terminal"
|
||||
? [runs[index].settled.promise]
|
||||
: [],
|
||||
),
|
||||
);
|
||||
} finally {
|
||||
await resetSubagentRegistryForTests({ persist: false });
|
||||
unsubscribe();
|
||||
for (const run of runs) {
|
||||
run.source.release();
|
||||
}
|
||||
}
|
||||
})());
|
||||
try {
|
||||
await runScoped(0, () =>
|
||||
registerSubagentRun({
|
||||
runId: first.runId,
|
||||
childSessionKey,
|
||||
requesterSessionKey,
|
||||
requesterAgentId: "main",
|
||||
requesterDisplayKey: requesterSessionKey,
|
||||
task: "First task",
|
||||
cleanup: "keep",
|
||||
spawnMode: "session",
|
||||
expectsCompletionMessage: true,
|
||||
completionTarget: "parent",
|
||||
completionRequesterSessionId: requesterSessionId,
|
||||
completionRequesterLifecycleRevision: requesterLifecycleRevision,
|
||||
gatewayContextResolver: context.resolveGatewayContext,
|
||||
}),
|
||||
);
|
||||
return {
|
||||
requesterScope,
|
||||
childScope,
|
||||
runs,
|
||||
handleChildCommand(command: AgentCommandOpts) {
|
||||
const run = runs.find(
|
||||
(candidate) =>
|
||||
candidate.runId === command.runId && command.sessionKey === childSessionKey,
|
||||
);
|
||||
if (!run) {
|
||||
return undefined;
|
||||
}
|
||||
return (async () => {
|
||||
const recorder = expectDefined(
|
||||
command.userTurnTranscriptRecorder,
|
||||
"Expected real child execution input recorder",
|
||||
);
|
||||
await recorder.persistApproved();
|
||||
run.commandEntered.resolve();
|
||||
await run.releaseModel.promise;
|
||||
command.abortSignal?.throwIfAborted();
|
||||
const message: AssistantMessage & { __openclaw: { runId: string } } = {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text: run.result }],
|
||||
api: "openai-responses",
|
||||
provider: "openai",
|
||||
model: "synthetic-completion",
|
||||
usage: createZeroUsageFixture(),
|
||||
stopReason: "stop",
|
||||
timestamp: Date.now(),
|
||||
__openclaw: { runId: run.runId },
|
||||
};
|
||||
SessionManager.open(childScope).appendMessage(message);
|
||||
return {
|
||||
payloads: [{ text: run.result, mediaUrl: null }],
|
||||
meta: {
|
||||
durationMs: 1,
|
||||
terminalReply: { disposition: "visible" as const, text: run.result },
|
||||
},
|
||||
};
|
||||
})();
|
||||
},
|
||||
async complete(index: 0 | 1) {
|
||||
await admit(index);
|
||||
runs[index].releaseModel.resolve();
|
||||
await dispatches.get(index);
|
||||
await runs[index].terminal.promise;
|
||||
},
|
||||
admitSuccessor: () => admit(1),
|
||||
settle,
|
||||
dispose,
|
||||
};
|
||||
} catch (error) {
|
||||
await dispose();
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
|
@ -90,6 +90,7 @@ describe("reactivateCompletedSubagentSession", () => {
|
|||
expect(replaceSubagentRunAfterSteerMock).toHaveBeenCalledWith({
|
||||
previousRunId: "run-current-ended",
|
||||
nextRunId: "run-next",
|
||||
preserveCompletedRun: true,
|
||||
assertCurrent: expect.any(Function),
|
||||
runTimeoutSeconds: 0,
|
||||
gatewayContextResolver: resolveGatewayContext,
|
||||
|
|
@ -139,6 +140,7 @@ describe("reactivateCompletedSubagentSession", () => {
|
|||
expect(replaceSubagentRunAfterSteerMock).toHaveBeenCalledWith({
|
||||
previousRunId: "run-prev-ended",
|
||||
nextRunId: "run-next",
|
||||
preserveCompletedRun: true,
|
||||
assertCurrent: expect.any(Function),
|
||||
runTimeoutSeconds: 0,
|
||||
task: " follow-up prompt text ",
|
||||
|
|
|
|||
|
|
@ -100,6 +100,7 @@ export async function reactivateCompletedSubagentSession(params: {
|
|||
: await runtime.replaceSubagentRunAfterSteerCore({
|
||||
previousRunId: source.runId,
|
||||
nextRunId: runId,
|
||||
preserveCompletedRun: true,
|
||||
runTimeoutSeconds: source.runTimeoutSeconds ?? 0,
|
||||
...(hasTask ? { task } : {}),
|
||||
assertCurrent: assertOriginalOwnerCurrent,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue