mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
* fix(subagents): preserve completion during concurrent registry writes Replace speculative row staging with per-row FIFO planning, immutable post-ACK publication, digest CAS, and bounded conflict replanning. Route registry and native completion writes through the same owner, and commit eligible terminal signals with their terminal rows. Retain runtime custody across metadata publication and acknowledged collector rekeys. Fence retired Gateway incarnations and unknown writes until canonical reconciliation. This local checkpoint has independent review clean through P2, core and selected test type checks, typed caller lint, and focused worker/lifecycle proof. Main integration and complete validation remain in progress; descriptorless recovery fixture setup and warm restoration failures remain recorded in the task notes for the next checkpoint. * fix(subagents): preserve reset and stop ownership after publication Prepare reset revocation after session generation checks and revalidate current child captures before commit. Reacquire queued Stop rows through their runtime owner after terminal publication. Adapt remaining fixtures to immutable registry observations and awaited acknowledgements. * fix(subagents): retain cancellation and completion custody across acknowledgements Preserve semantic kill ownership during replanning, allow confirmation of the same cancellation, and retain grace callbacks through their own terminal acknowledgement. Keep compact projections noncustodial and accept current runtime owners in Stop. Complete immutable-row fixture cutovers and join acknowledged terminal work during teardown. * test(subagents): assert committed cancellation failure boundaries * fix(agents): retain retirement receipts and lifecycle attempt custody Keep acknowledged quiet deletion receipts through publication recovery and bind delayed terminal grace to its admitted execution attempt. Adapt immutable publication fixtures and requester reply proof to the integrated yield policy. * refactor(agents): centralize canonical subagent restoration Keep restore admission, immutable row installation, and uncertain write recovery in the registry mutation owner. Retain cache and notification ownership in state without the reverse import, and migrate all callers and fixtures to the owning contracts. * fix(agents): preserve child ownership during collector settlement Resolve collector usage through the existing child-session owner before terminal publication. Cover recorded child agents and configured defaults for unscoped session keys through real queued registration failure settlement. * fix(subagents): retain Gateway custody after lifecycle commits Return the acknowledged published row from lifecycle mutations instead of their private planning draft. The draft has no Gateway binding, so cleanup could leave an owned managed-worktree session behind after falling back to raw RPC. The real production-boundary deletion case failed before this fix and now deletes the session/worktree while preserving its result snapshot. Separate raw row hashing and mutation contracts from persisted record types to remove static import cycles. Finish the cutover by removing obsolete exports, eliminate a redundant assertion, and give sweep deletion a unique export name. Keep fixture checks on actual published rows and provide the scheduler logger's trace method. No test guard, budget, baseline, or schema change is needed. Final focused proof: five affected files / 311 cases, including all 25 spawn production-boundary cases. Additional unchanged-contract CI repair tests retain their passing replays. Core and all 27 affected test graphs, typed lint, architecture/dependency and assertion guards, formatting, diff checks, and isolated review through P2 pass. * test(agents): align registry proof with owner lifetimes Keep requester replay guards until isolated database teardown, avoiding schema writes while worker cleanup is active. Update benchmark imports to the extracted registry contract and row owners.
125 lines
5.1 KiB
TypeScript
125 lines
5.1 KiB
TypeScript
import { mock } from "node:test";
|
|
import type { SubagentRunMutation } from "../src/agents/subagents/registry/subagent-registry-mutation.types.js";
|
|
import type { SubagentRunMutationOptions } from "../src/agents/subagents/registry/subagent-registry-persistence.types.js";
|
|
import type { SubagentRunRecord } from "../src/agents/subagents/registry/subagent-registry.types.js";
|
|
import type { callGateway } from "../src/gateway/call.js";
|
|
|
|
/** Install before the worker imports the registry; keep these bindings across its samples. */
|
|
export async function installBenchmarkRegistryRuntime(mode: "memory" | "durable") {
|
|
const handles: Array<{ restore(): void }> = [];
|
|
const replaceModule = (specifier: string, namedExports: object) => {
|
|
handles.push(mock.module(new URL(specifier, import.meta.url), { namedExports }));
|
|
};
|
|
let closed = false;
|
|
let clearConfig: (() => void) | undefined;
|
|
let request: typeof callGateway = async () => {
|
|
throw new Error("Benchmark Gateway wait barrier is not installed");
|
|
};
|
|
const close = () => {
|
|
if (closed) {
|
|
return;
|
|
}
|
|
closed = true;
|
|
try {
|
|
clearConfig?.();
|
|
} finally {
|
|
for (const handle of handles.toReversed()) {
|
|
handle.restore();
|
|
}
|
|
}
|
|
};
|
|
|
|
try {
|
|
replaceModule("../src/agents/subagents/announce/subagent-announce.ts", {
|
|
captureSubagentCompletionReply: async () => undefined,
|
|
} satisfies Pick<
|
|
typeof import("../src/agents/subagents/announce/subagent-announce.js"),
|
|
"captureSubagentCompletionReply"
|
|
>);
|
|
replaceModule("../src/agents/subagents/announce/subagent-announce.requester-settle-wake.ts", {
|
|
maybeWakeRequesterAfterAllChildrenSettled: async () => false,
|
|
} satisfies typeof import("../src/agents/subagents/announce/subagent-announce.requester-settle-wake.js"));
|
|
const persistence =
|
|
await import("../src/agents/subagents/registry/subagent-registry-persistence.js");
|
|
if (mode === "memory") {
|
|
const { bindSubagentRunRecord } =
|
|
await import("../src/agents/subagents/registry/subagent-registry.store.codec.js");
|
|
const { subagentRunRowVersion } =
|
|
await import("../src/agents/subagents/registry/subagent-registry.store.row.js");
|
|
const mutate = async <P extends SubagentRunMutation<unknown>>(
|
|
runIds: readonly string[],
|
|
plan: (rows: ReadonlyMap<string, SubagentRunRecord>) => P,
|
|
options: SubagentRunMutationOptions<P> = {},
|
|
): Promise<P["value"]> => {
|
|
if (options.commit) {
|
|
throw new Error("Memory benchmark reached a native completion transaction");
|
|
}
|
|
return await persistence.mutateSubagentRuns(runIds, plan, {
|
|
...options,
|
|
onPublished: (postimages, value) => {
|
|
if (postimages.size > 0) {
|
|
options.onPublished?.(postimages, value);
|
|
}
|
|
},
|
|
commit: async (planned, _versions, authority) => {
|
|
authority.assertCurrent();
|
|
if (planned.terminalEvents?.length) {
|
|
throw new Error("Memory benchmark attempted durable terminal signals");
|
|
}
|
|
const postimages = new Map(
|
|
[...(planned.postimages ?? [])].map(
|
|
([runId, row]) => [runId, row ? structuredClone(row) : null] as const,
|
|
),
|
|
);
|
|
const versions = new Map(
|
|
[...postimages].map(
|
|
([runId, row]) =>
|
|
[runId, row ? subagentRunRowVersion(bindSubagentRunRecord(row)) : null] as const,
|
|
),
|
|
);
|
|
return { value: planned.value, postimages, versions };
|
|
},
|
|
});
|
|
};
|
|
replaceModule("../src/agents/subagents/registry/subagent-registry-persistence.ts", {
|
|
...persistence,
|
|
mutateSubagentRuns: mutate,
|
|
} satisfies typeof persistence);
|
|
}
|
|
|
|
const gatewayRuntime = await import("../src/gateway/server-recovery-runtime-context.js");
|
|
replaceModule("../src/gateway/server-recovery-runtime-context.ts", {
|
|
...gatewayRuntime,
|
|
bindGatewayLifecycleRequest: () => request,
|
|
} satisfies typeof gatewayRuntime);
|
|
|
|
const completion =
|
|
await import("../src/agents/subagents/registry/subagent-registry-lifecycle-completion.js");
|
|
replaceModule("../src/agents/subagents/registry/subagent-registry-lifecycle-completion.ts", {
|
|
...completion,
|
|
completeSubagentRunAttempt: async (context, params) =>
|
|
completion.completeSubagentRunAttempt(context, {
|
|
...params,
|
|
// These synthetic runs have no child session or transcript to update.
|
|
suppressSessionEffects: true,
|
|
completionSnapshot: { resultText: null, capturedAt: params.endedAt ?? Date.now() },
|
|
}),
|
|
} satisfies typeof completion);
|
|
|
|
const config = await import("../src/config/runtime-snapshot.js");
|
|
clearConfig = config.clearRuntimeConfigSnapshot;
|
|
config.setRuntimeConfigSnapshot({});
|
|
return {
|
|
setCallGateway(call: typeof callGateway) {
|
|
if (closed) {
|
|
throw new Error("Benchmark registry runtime is closed");
|
|
}
|
|
request = call;
|
|
},
|
|
close,
|
|
};
|
|
} catch (error) {
|
|
close();
|
|
throw error;
|
|
}
|
|
}
|