mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-04 02:00:10 +00:00
perf(gateway): reuse state admission for session events (#158627)
* perf(state): retain warm database read admissions Reuse immutable admissions per resolved path and generation at the lifecycle owner. Preserve first-creation aliases, close custody, replacement exclusion, and caller-specific schema scope guards. Keep warm capture allocation-free by constructing token closures only on the cold path. * test(state): avoid shadowing the read context fixture * perf(gateway): carry admitted state facts through session projections Bind session-row readers to their existing captured state source and retain schema lifetime facts. Share same-file admission renewal with deletion reads, fence replaced sources and late publications, and keep warm readiness checks free of path resolution and allocation. * test(gateway): match suite-owned ACP read context The projection now retains its creation-time state environment. Match the suite-owned ACP read in the nested cleanup regression while preserving both cleanup ordering and disposal failure assertions. CI run 36216060049 reproduced both mismatched interceptions; the corrected file passes all 10 tests. * fix(test): advance denied Windows port reservations Treat Windows EACCES at initial automatic listener reservation consistently with candidate probing, after provisional owners drain. Keep explicit ports and later pinned reacquisition unchanged. Add deterministic reservation and release coverage; the original hosted lifecycle denial did not retain its stage and remains a separate uncertainty. * test(gateway): avoid shadowing boxed progress value
This commit is contained in:
parent
017d224a02
commit
59100c2655
27 changed files with 713 additions and 191 deletions
|
|
@ -20,6 +20,7 @@ import { resolveOpenClawStateSqlitePath } from "../../../state/openclaw-state-db
|
|||
import {
|
||||
captureOpenClawStateReadContext,
|
||||
captureOpenClawStateWorkerContext,
|
||||
type OpenClawStateReadContext,
|
||||
} from "../../../state/openclaw-state-worker-context.js";
|
||||
import type { OpenClawStateWorkerContext } from "../../../state/openclaw-state-worker-context.types.js";
|
||||
import {
|
||||
|
|
@ -144,7 +145,7 @@ export function retainUnpublishedSubagentChanges<T>(
|
|||
/** Selecting a read must not consume another database owner's publication. */
|
||||
function selectSubagentCacheStateForRead<T extends SubagentRunReadRecord>(
|
||||
state: SubagentRunsCacheState<T>,
|
||||
context?: OpenClawStateWorkerContext,
|
||||
context?: OpenClawStateReadContext,
|
||||
): SubagentRunsCacheState<T> {
|
||||
const identity = state.retiredPublicationIdentity ?? state.admission?.identity;
|
||||
const matches = context
|
||||
|
|
@ -259,10 +260,11 @@ export function rememberSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
|
||||
export function getPersistedSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
||||
cache: SubagentRunsCache<T>,
|
||||
prepared?: OpenClawStateReadContext,
|
||||
): Map<string, T> | null {
|
||||
let admission: OpenClawStateDatabaseReadAdmission | undefined;
|
||||
if (!cache.load) {
|
||||
const context = captureOpenClawStateReadContext();
|
||||
const context = prepared ?? captureOpenClawStateReadContext();
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
context.admission.assertCurrent();
|
||||
admission = context.admission;
|
||||
|
|
@ -282,7 +284,9 @@ export function getPersistedSubagentRunsSnapshot<T extends SubagentRunReadRecord
|
|||
!matchesSubagentCacheAdmission(cache.state.admission, admission) ||
|
||||
(admission && cache.state.sourceIdentity !== admission.identity.key)
|
||||
) {
|
||||
cache.state = { admission, sourceIdentity: admission?.identity.key };
|
||||
if (!prepared) {
|
||||
cache.state = { admission, sourceIdentity: admission?.identity.key };
|
||||
}
|
||||
return null;
|
||||
}
|
||||
return cache.state.snapshot ?? null;
|
||||
|
|
@ -290,8 +294,9 @@ export function getPersistedSubagentRunsSnapshot<T extends SubagentRunReadRecord
|
|||
|
||||
export function loadPersistedSubagentRunsForRead<T extends SubagentRunReadRecord>(
|
||||
cache: SubagentRunsCache<T>,
|
||||
prepared?: OpenClawStateReadContext,
|
||||
): Map<string, T> {
|
||||
const cached = getPersistedSubagentRunsSnapshot(cache);
|
||||
const cached = getPersistedSubagentRunsSnapshot(cache, prepared);
|
||||
if (cached) {
|
||||
return cached;
|
||||
}
|
||||
|
|
@ -323,6 +328,7 @@ export function getSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
inMemoryRuns: Map<string, SubagentRunRecord>,
|
||||
cache: SubagentRunsCache<T>,
|
||||
scope?: {
|
||||
context?: OpenClawStateReadContext;
|
||||
load?: () => Iterable<T>;
|
||||
selectCached?: (lookup: SubagentSessionReadLookup) => readonly string[];
|
||||
fresh?: boolean;
|
||||
|
|
@ -333,7 +339,7 @@ export function getSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
if (
|
||||
shouldReadPersistedSubagentRuns() &&
|
||||
!cache.load &&
|
||||
!getPersistedSubagentRunsSnapshot(cache)
|
||||
!getPersistedSubagentRunsSnapshot(cache, scope?.context)
|
||||
) {
|
||||
throw new Error("Subagent session-list facts must be prepared before synchronous reads");
|
||||
}
|
||||
|
|
@ -341,7 +347,8 @@ export function getSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
if (shouldReadPersistedSubagentRuns()) {
|
||||
try {
|
||||
// Scoped reads use indexed SQL until a complete owner snapshot is available.
|
||||
const cached = scope?.load && !scope.fresh ? getPersistedSubagentRunsSnapshot(cache) : null;
|
||||
const cached =
|
||||
scope?.load && !scope.fresh ? getPersistedSubagentRunsSnapshot(cache, scope.context) : null;
|
||||
const cachedRows =
|
||||
cached && scope?.selectCached
|
||||
? indexedSnapshotRows(
|
||||
|
|
@ -353,7 +360,7 @@ export function getSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
: cached?.values();
|
||||
const persisted = scope?.load
|
||||
? (cachedRows ?? scope.load())
|
||||
: loadPersistedSubagentRunsForRead(cache).values();
|
||||
: loadPersistedSubagentRunsForRead(cache, scope?.context).values();
|
||||
for (const entry of persisted) {
|
||||
if (!scope || scope.matches(entry)) {
|
||||
merged.set(
|
||||
|
|
@ -367,7 +374,7 @@ export function getSubagentRunsSnapshot<T extends SubagentRunReadRecord>(
|
|||
}
|
||||
}
|
||||
if (shouldReadPersistedSubagentRuns()) {
|
||||
const state = selectSubagentCacheStateForRead(cache.state);
|
||||
const state = selectSubagentCacheStateForRead(cache.state, scope?.context);
|
||||
for (const [runId, { entry }] of state.changes ?? []) {
|
||||
if (entry && (!scope || scope.matches(entry))) {
|
||||
merged.set(runId, scope?.load && !scope.borrowPersisted ? structuredClone(entry) : entry);
|
||||
|
|
@ -393,6 +400,7 @@ export async function readCompactSubagentRuns(context: OpenClawStateWorkerContex
|
|||
const reply = await executeExistingOpenClawStateRead(
|
||||
{ path: context.admission.databasePath, env: context.environment },
|
||||
{ type: "subagents.sessionList" },
|
||||
{ context },
|
||||
);
|
||||
if (!reply) {
|
||||
return new Map<string, SubagentRunReadRecord>();
|
||||
|
|
@ -435,19 +443,34 @@ export async function readFullSubagentRuns(
|
|||
export async function prepareSubagentRunsCache<T extends SubagentRunReadRecord>(
|
||||
cache: SubagentRunsCache<T>,
|
||||
load: (context: OpenClawStateWorkerContext) => Promise<Map<string, T>>,
|
||||
prepared?: OpenClawStateWorkerContext,
|
||||
): Promise<Map<string, T>> {
|
||||
const context = captureOpenClawStateWorkerContext();
|
||||
assertSubagentReadContext(context);
|
||||
if (getActiveOpenClawStateDatabaseReadSnapshot()) {
|
||||
const context = prepared ?? captureOpenClawStateWorkerContext();
|
||||
const assertCurrent = () => {
|
||||
if (!prepared) {
|
||||
assertSubagentReadContext(context);
|
||||
return;
|
||||
}
|
||||
getAsyncWorkSignal()?.throwIfAborted();
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
context.admission.assertCurrent();
|
||||
};
|
||||
assertCurrent();
|
||||
if (
|
||||
getActiveOpenClawStateDatabaseReadSnapshot({
|
||||
path: context.admission.databasePath,
|
||||
env: context.environment,
|
||||
})
|
||||
) {
|
||||
// Private snapshot bytes never become canonical resident facts.
|
||||
const runs = await load(context);
|
||||
assertSubagentReadContext(context);
|
||||
assertCurrent();
|
||||
return runs;
|
||||
}
|
||||
const callerAbortSignal = getAsyncWorkSignal();
|
||||
let retriedCanceledFill = false;
|
||||
while (true) {
|
||||
assertSubagentReadContext(context);
|
||||
assertCurrent();
|
||||
let state = cache.state;
|
||||
if (
|
||||
!matchesSubagentCacheAdmission(state.admission, context.admission) ||
|
||||
|
|
@ -472,7 +495,7 @@ export async function prepareSubagentRunsCache<T extends SubagentRunReadRecord>(
|
|||
promise: Promise.resolve().then(async () => {
|
||||
const runs = await load(context);
|
||||
try {
|
||||
assertSubagentReadContext(context);
|
||||
assertCurrent();
|
||||
} catch (error) {
|
||||
// Classify only after successful read cleanup, before waiters observe rejection.
|
||||
fill.cleanCancellation =
|
||||
|
|
@ -514,13 +537,20 @@ export async function prepareSubagentRunsCache<T extends SubagentRunReadRecord>(
|
|||
) {
|
||||
throw error;
|
||||
}
|
||||
assertSubagentReadContext(context);
|
||||
assertCurrent();
|
||||
retriedCanceledFill = true;
|
||||
} finally {
|
||||
if (cache.state.pending === fill) {
|
||||
cache.state.pending = undefined;
|
||||
}
|
||||
}
|
||||
if (
|
||||
prepared &&
|
||||
cache.state.sourceIdentity !== undefined &&
|
||||
cache.state.sourceIdentity !== context.admission.identity.key
|
||||
) {
|
||||
throw new Error("Subagent registry database changed during preparation");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import { renameSync } from "node:fs";
|
||||
import { expect, it, vi } from "vitest";
|
||||
import { AsyncWorkScope } from "../../../shared/async-work-scope.js";
|
||||
import { createDeferredCore } from "../../../shared/deferred.js";
|
||||
|
|
@ -54,43 +55,56 @@ it("does not retain a removed live-only row as a prepared durable payload", asyn
|
|||
});
|
||||
});
|
||||
|
||||
it("rejects a retired published snapshot after reopening the same database file", async () => {
|
||||
await withPersistedReads(async () => {
|
||||
const entry = retainedRun();
|
||||
persistSubagentRunsToDiskOrThrow(new Map([[entry.runId, entry]]));
|
||||
const context = captureOpenClawStateWorkerContext();
|
||||
const release = createDeferredCore();
|
||||
const unregister = registerOpenClawStateDatabaseAsyncResource({ close: () => release.promise });
|
||||
const closing = closeOpenClawStateDatabaseByPathAsync(context.admission.databasePath);
|
||||
try {
|
||||
expect(() => captureOpenClawStateWorkerContext()).toThrow("read admission is closed");
|
||||
const retired = {
|
||||
...entry,
|
||||
completion: { required: false, resultText: "retired publication" },
|
||||
};
|
||||
persistSubagentRunsToDiskOrThrow(new Map([[entry.runId, retired]]), [entry.runId]);
|
||||
release.resolve();
|
||||
await closing;
|
||||
const reopened = {
|
||||
...entry,
|
||||
completion: { required: false, resultText: "reopened durable result" },
|
||||
};
|
||||
saveSubagentRegistryToSqlite(new Map([[entry.runId, reopened]]));
|
||||
expect(captureOpenClawStateWorkerContext().admission.identity.key).toBe(
|
||||
context.admission.identity.key,
|
||||
);
|
||||
const prepared = await prepareSubagentRunsSnapshotForRunIds(new Map(), ["collector"]);
|
||||
expect(prepared.consume((runs) => runs.get(entry.runId)?.completion?.resultText)).toEqual({
|
||||
ready: true,
|
||||
value: "reopened durable result",
|
||||
it.each(["reopening", "replacing"] as const)(
|
||||
"rejects a retired published snapshot after %s the database file",
|
||||
async (change) => {
|
||||
await withPersistedReads(async () => {
|
||||
const entry = retainedRun();
|
||||
persistSubagentRunsToDiskOrThrow(new Map([[entry.runId, entry]]));
|
||||
const context = captureOpenClawStateWorkerContext();
|
||||
const original = await prepareSubagentRunsSnapshotForRunIds(new Map(), ["collector"]);
|
||||
const release = createDeferredCore();
|
||||
const unregister = registerOpenClawStateDatabaseAsyncResource({
|
||||
close: () => release.promise,
|
||||
});
|
||||
} finally {
|
||||
release.resolve();
|
||||
await closing;
|
||||
unregister();
|
||||
}
|
||||
});
|
||||
});
|
||||
const closing = closeOpenClawStateDatabaseByPathAsync(context.admission.databasePath);
|
||||
try {
|
||||
expect(() => captureOpenClawStateWorkerContext()).toThrow("read admission is closed");
|
||||
const retired = {
|
||||
...entry,
|
||||
completion: { required: false, resultText: "retired publication" },
|
||||
};
|
||||
persistSubagentRunsToDiskOrThrow(new Map([[entry.runId, retired]]), [entry.runId]);
|
||||
release.resolve();
|
||||
await closing;
|
||||
if (change === "replacing") {
|
||||
renameSync(context.admission.databasePath, `${context.admission.databasePath}.retired`);
|
||||
}
|
||||
const reopened = {
|
||||
...entry,
|
||||
completion: { required: false, resultText: "reopened durable result" },
|
||||
};
|
||||
saveSubagentRegistryToSqlite(new Map([[entry.runId, reopened]]));
|
||||
const current = captureOpenClawStateWorkerContext();
|
||||
expect(current.admission.identity.key === context.admission.identity.key).toBe(
|
||||
change === "reopening",
|
||||
);
|
||||
const consume = vi.fn();
|
||||
expect(() => original.consume(consume)).toThrow("read admission changed");
|
||||
expect(consume).not.toHaveBeenCalled();
|
||||
const prepared = await prepareSubagentRunsSnapshotForRunIds(new Map(), ["collector"]);
|
||||
expect(prepared.consume((runs) => runs.get(entry.runId)?.completion?.resultText)).toEqual({
|
||||
ready: true,
|
||||
value: "reopened durable result",
|
||||
});
|
||||
} finally {
|
||||
release.resolve();
|
||||
await closing;
|
||||
unregister();
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["delete", "replace", "move alias", "full publication"] as const)(
|
||||
"applies a %s in the prepared read's consuming frame",
|
||||
|
|
|
|||
|
|
@ -59,10 +59,13 @@ export {
|
|||
export function buildSubagentSessionListReadIndex(
|
||||
now = Date.now(),
|
||||
sessionKeys?: readonly string[],
|
||||
preparedRuns?: Map<string, SubagentRunReadRecord>,
|
||||
): SubagentRunReadIndex<SubagentRunReadRecord> {
|
||||
const runs = sessionKeys
|
||||
? getSubagentSessionListRunsSnapshotForSessions(subagentRuns, sessionKeys)
|
||||
: getSubagentSessionListRunsSnapshotForRead(subagentRuns);
|
||||
const runs =
|
||||
preparedRuns ??
|
||||
(sessionKeys
|
||||
? getSubagentSessionListRunsSnapshotForSessions(subagentRuns, sessionKeys)
|
||||
: getSubagentSessionListRunsSnapshotForRead(subagentRuns));
|
||||
return buildSubagentRunReadIndexFromRuns({
|
||||
runs,
|
||||
inMemoryRuns: sessionKeys
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
// Real-storage proof that committed and best-effort publications survive read-owner retirement.
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import path from "node:path";
|
||||
import { afterEach, beforeEach, expect, it, vi } from "vitest";
|
||||
import { readSqliteUserVersion } from "../../../infra/sqlite-user-version.js";
|
||||
import { AsyncWorkScope } from "../../../shared/async-work-scope.js";
|
||||
|
|
@ -24,6 +25,7 @@ import { createSubagentRunRecord } from "../../subagent-test-fixtures.test-helpe
|
|||
import { buildControlledSubagentRunsReadContext } from "./subagent-control-scope.js";
|
||||
import {
|
||||
clearSubagentRunsReadCacheForTest,
|
||||
createSubagentSessionListReadView,
|
||||
getSubagentRunsSnapshotForChildSession,
|
||||
getSubagentRunsSnapshotForRead,
|
||||
getSubagentRunsSnapshotForSessions,
|
||||
|
|
@ -67,6 +69,8 @@ function runs(model: string, runId = "one") {
|
|||
|
||||
it("reads resident session-list facts without enumerating the process environment", () => {
|
||||
persistSubagentRunsToDiskOrThrow(runs("retained"));
|
||||
const prepared = createSubagentSessionListReadView({ env: state.env });
|
||||
const preparedIdentity = prepared.snapshotIdentity();
|
||||
const originalEnv = process.env;
|
||||
let enumerations = 0;
|
||||
let identity: object | undefined;
|
||||
|
|
@ -86,6 +90,14 @@ it("reads resident session-list facts without enumerating the process environmen
|
|||
expect(identity).toBeDefined();
|
||||
expect(model).toBe("retained");
|
||||
expect(enumerations).toBe(0);
|
||||
const resolutions = vi.spyOn(path, "resolve");
|
||||
let current: object | undefined;
|
||||
for (let index = 0; index < 100; index++) {
|
||||
current = prepared.snapshotIdentity();
|
||||
}
|
||||
expect(current).toBe(preparedIdentity);
|
||||
expect(resolutions).not.toHaveBeenCalled();
|
||||
resolutions.mockRestore();
|
||||
});
|
||||
|
||||
it.each(["maintenance", "schema"] as const)(
|
||||
|
|
@ -94,7 +106,14 @@ it.each(["maintenance", "schema"] as const)(
|
|||
persistSubagentRunsToDiskOrThrow(runs("retained"));
|
||||
const identity = getSubagentSessionListReadSnapshotIdentity();
|
||||
const maintenance = createOpenClawDatabaseMaintenanceScope();
|
||||
const bind = () => AsyncLocalStorage.bind(getSubagentSessionListReadSnapshotIdentity);
|
||||
const bind = () => {
|
||||
const prepared = createSubagentSessionListReadView({ env: state.env });
|
||||
expect(prepared.snapshotIdentity()).toBe(identity);
|
||||
return {
|
||||
current: AsyncLocalStorage.bind(getSubagentSessionListReadSnapshotIdentity),
|
||||
prepared,
|
||||
};
|
||||
};
|
||||
const read =
|
||||
kind === "maintenance"
|
||||
? maintenance.run(bind)
|
||||
|
|
@ -103,7 +122,10 @@ it.each(["maintenance", "schema"] as const)(
|
|||
bind,
|
||||
);
|
||||
await maintenance.close();
|
||||
expect(read).toThrow(kind === "maintenance" ? "scope is closed" : "admission has ended");
|
||||
const failure = kind === "maintenance" ? "scope is closed" : "admission has ended";
|
||||
expect(read.current).toThrow(failure);
|
||||
expect(read.prepared.snapshotIdentity).toThrow(failure);
|
||||
await expect(read.prepared.prepare()).rejects.toThrow(failure);
|
||||
expect(getSubagentSessionListReadSnapshotIdentity()).toBe(identity);
|
||||
},
|
||||
);
|
||||
|
|
@ -139,6 +161,38 @@ function holdFirstCompactRead(
|
|||
return { entered: entered.promise, release: release.resolve, read };
|
||||
}
|
||||
|
||||
it("keeps a newer source publication when a prepared projection read settles", async () => {
|
||||
store.saveSubagentRegistryToSqlite(runs("original"));
|
||||
const prepared = createSubagentSessionListReadView({ env: state.env });
|
||||
const other = await createOpenClawTestState({ scenario: "minimal", applyEnv: false });
|
||||
const gate = holdFirstCompactRead();
|
||||
const pending = prepared.prepare();
|
||||
const outcome = pending.catch((error: unknown) => error);
|
||||
let published: object | undefined;
|
||||
try {
|
||||
await gate.entered;
|
||||
await withEnvAsync({ OPENCLAW_STATE_DIR: other.stateDir }, async () => {
|
||||
persistSubagentRunsToDiskOrThrow(runs("newer source"));
|
||||
published = getSubagentSessionListReadSnapshotIdentity();
|
||||
});
|
||||
expect(prepared.snapshotIdentity()).toBeUndefined();
|
||||
gate.release();
|
||||
expect(await outcome).toMatchObject({
|
||||
message: "Subagent registry database changed during preparation",
|
||||
});
|
||||
await withEnvAsync({ OPENCLAW_STATE_DIR: other.stateDir }, async () => {
|
||||
expect(getSubagentSessionListReadSnapshotIdentity()).toBe(published);
|
||||
expect(getSubagentSessionListRunsSnapshotForRead(new Map()).get("one")?.model).toBe(
|
||||
"newer source",
|
||||
);
|
||||
});
|
||||
} finally {
|
||||
gate.release();
|
||||
await outcome;
|
||||
await other.cleanup();
|
||||
}
|
||||
});
|
||||
|
||||
it.each(["cancel", "query failure", "cleanup failure"] as const)(
|
||||
"lets a current waiter recover only a clean canceled fill (%s)",
|
||||
async (kind) => {
|
||||
|
|
|
|||
|
|
@ -4,7 +4,11 @@ import {
|
|||
} from "../../../sessions/session-lifecycle-events.js";
|
||||
import { isStateDatabaseReadAdmissionInvalidatedError } from "../../../state/openclaw-state-db-async-lifecycle.js";
|
||||
import { getActiveOpenClawStateDatabaseReadSnapshot } from "../../../state/openclaw-state-db-readonly.js";
|
||||
import { captureOpenClawStateWorkerContext } from "../../../state/openclaw-state-worker-context.js";
|
||||
import { resolveOpenClawStateSqlitePath } from "../../../state/openclaw-state-db.paths.js";
|
||||
import {
|
||||
captureOpenClawStateWorkerContext,
|
||||
prepareOpenClawStateReadSource,
|
||||
} from "../../../state/openclaw-state-worker-context.js";
|
||||
import type { OpenClawStateWorkerContext } from "../../../state/openclaw-state-worker-context.types.js";
|
||||
import {
|
||||
projectSubagentRunForMaintenance,
|
||||
|
|
@ -226,6 +230,56 @@ export function getSubagentSessionListReadSnapshotIdentity(): object | undefined
|
|||
}
|
||||
}
|
||||
|
||||
export type SubagentSessionListReadView = {
|
||||
snapshotIdentity(this: void): object | undefined;
|
||||
runs(this: void): Map<string, SubagentRunReadRecord>;
|
||||
prepare(this: void): Promise<void>;
|
||||
};
|
||||
|
||||
/** A long-lived projection retains its source; registry publications still own the facts. */
|
||||
export function createSubagentSessionListReadView(options: {
|
||||
env: NodeJS.ProcessEnv;
|
||||
path?: string;
|
||||
}): SubagentSessionListReadView {
|
||||
const path = options.path ?? resolveOpenClawStateSqlitePath(options.env);
|
||||
const source = prepareOpenClawStateReadSource({ path, env: options.env });
|
||||
const cache = persistedSubagentSessionListRunsReadCache;
|
||||
const readPersisted = shouldReadPersistedSubagentRuns();
|
||||
const matches = () => true;
|
||||
const prepare = (context: OpenClawStateWorkerContext) =>
|
||||
prepareSubagentRunsCache(cache, readCompactSubagentRuns, context);
|
||||
return {
|
||||
snapshotIdentity() {
|
||||
if (!readPersisted) {
|
||||
return subagentRuns;
|
||||
}
|
||||
try {
|
||||
return getPersistedSubagentRunsSnapshot(cache, source.current()) ?? undefined;
|
||||
} catch (error) {
|
||||
if (!isStateDatabaseReadAdmissionInvalidatedError(error)) {
|
||||
throw error;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
},
|
||||
runs() {
|
||||
return getSubagentRunsSnapshot(subagentRuns, cache, {
|
||||
context: readPersisted ? source.current() : undefined,
|
||||
matches,
|
||||
});
|
||||
},
|
||||
async prepare() {
|
||||
if (!readPersisted) {
|
||||
return;
|
||||
}
|
||||
if (getActiveOpenClawStateDatabaseReadSnapshot({ path, env: options.env })) {
|
||||
throw new Error("Resident subagent preparation cannot adopt a private database snapshot");
|
||||
}
|
||||
await source.withCurrent(prepare);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export async function prepareSubagentSessionListReadCache(): Promise<void> {
|
||||
if (!shouldReadPersistedSubagentRuns()) {
|
||||
return;
|
||||
|
|
|
|||
|
|
@ -603,7 +603,14 @@ describe("createChatRunState", () => {
|
|||
|
||||
it("retains native boxed values, shared containers, and a proxy revoked by its last getter", () => {
|
||||
const state = createChatRunState();
|
||||
const shared = { value: Object(3), text: Object("é"), enabled: Object(false) };
|
||||
const numberHints: string[] = [];
|
||||
const boxedNumber = Object.assign(Object(3), {
|
||||
[Symbol.toPrimitive]: (hint: string) => {
|
||||
numberHints.push(hint);
|
||||
return -3.25;
|
||||
},
|
||||
});
|
||||
const shared = { value: boxedNumber, text: Object("é"), enabled: Object(false) };
|
||||
const revocable = Proxy.revocable(["last"], {
|
||||
get(target, key, receiver) {
|
||||
const value: unknown = Reflect.get(target, key, receiver);
|
||||
|
|
@ -626,10 +633,11 @@ describe("createChatRunState", () => {
|
|||
});
|
||||
const event = state.runs.get("run-1")?.progressSnapshot?.events[0];
|
||||
expect(event?.data.result).toEqual({
|
||||
first: { value: 3, text: "é", enabled: false },
|
||||
again: { value: 3, text: "é", enabled: false },
|
||||
first: { value: -3.25, text: "é", enabled: false },
|
||||
again: { value: -3.25, text: "é", enabled: false },
|
||||
last: ["last"],
|
||||
});
|
||||
expect(numberHints).toEqual(["number", "number"]);
|
||||
expect(state.runs.get("run-1")?.progressSnapshot?.byteLength).toBe(
|
||||
Buffer.byteLength(JSON.stringify(event)),
|
||||
);
|
||||
|
|
|
|||
|
|
@ -45,6 +45,7 @@ const { createSessionStoreDir, openClient, withSessionTestState } =
|
|||
test.each([false, true])(
|
||||
"nested state cleanup joins suite ACP reads (disposal fails=%s)",
|
||||
async (disposalFails) => {
|
||||
const suiteStateDir = runtimePaths.resolveStateDir();
|
||||
const entered = createDeferredCore();
|
||||
const release = createDeferredCore();
|
||||
const boundary = createDeferredCore<"joined" | "closed">();
|
||||
|
|
@ -88,7 +89,7 @@ test.each([false, true])(
|
|||
if (
|
||||
!held &&
|
||||
args[1].type === "acpSessions.metadata" &&
|
||||
args[0].env?.OPENCLAW_STATE_DIR === state.stateDir &&
|
||||
args[0].env?.OPENCLAW_STATE_DIR === suiteStateDir &&
|
||||
getAsyncWorkSignal() !== fixtureSignal
|
||||
) {
|
||||
held = true;
|
||||
|
|
@ -105,7 +106,7 @@ test.each([false, true])(
|
|||
await Promise.race([
|
||||
entered.promise,
|
||||
preparing.then(() => {
|
||||
throw new Error("Suite projection completed without the fixture's ACP metadata read");
|
||||
throw new Error("Suite projection completed without its ACP metadata read");
|
||||
}),
|
||||
]);
|
||||
const ensure = projection.ensureMaterialized;
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import { expect, it } from "vitest";
|
||||
import { createSubagentSessionListReadView } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import type { SessionEntry } from "../config/sessions.js";
|
||||
import { runSynchronousWork } from "../shared/synchronous-work.js";
|
||||
import { filterSessionEntries } from "./session-list-filters.js";
|
||||
|
|
@ -8,7 +9,9 @@ import { createSessionRowProjectionContext } from "./session-row-projection-cont
|
|||
|
||||
it("benchmarks warm identity filtering across viewers", () => {
|
||||
const cfg = { agents: { list: [{ id: "main", default: true }] } };
|
||||
const context = createSessionRowProjectionContext().current;
|
||||
const context = createSessionRowProjectionContext(
|
||||
createSubagentSessionListReadView({ env: process.env }),
|
||||
).current;
|
||||
const identities = context.userProfileIdentityById;
|
||||
for (let index = 0; index < 50; index++) {
|
||||
const id = `profile-${index}`;
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ import { subagentRuns } from "../agents/subagents/registry/subagent-registry-mem
|
|||
import { publishSubagentRunChanges } from "../agents/subagents/registry/subagent-registry-publication.js";
|
||||
import { buildSubagentRunReadIndexFromRuns } from "../agents/subagents/registry/subagent-registry-queries.js";
|
||||
import * as registryRead from "../agents/subagents/registry/subagent-registry-read.js";
|
||||
import { createSubagentSessionListReadView } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import { registerAgentRunCapacityWait } from "../infra/agent-run-capacity-wait.js";
|
||||
import {
|
||||
buildProjectedAgentRunIndex,
|
||||
|
|
@ -35,7 +36,9 @@ const runContext = {
|
|||
};
|
||||
|
||||
function fixture() {
|
||||
const context = createSessionRowProjectionContext();
|
||||
const context = createSessionRowProjectionContext(
|
||||
createSubagentSessionListReadView({ env: process.env }),
|
||||
);
|
||||
const prepare = (epoch: number) =>
|
||||
context.prepare(
|
||||
epoch,
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
import { getSubagentRegistryPublicationRevision } from "../agents/subagents/registry/subagent-registry-publication.js";
|
||||
import { buildSubagentSessionListReadIndex } from "../agents/subagents/registry/subagent-registry-read.js";
|
||||
import { getSubagentSessionListReadSnapshotIdentity } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import type { SubagentSessionListReadView } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import {
|
||||
buildProjectedAgentRunIndex,
|
||||
readAgentRunIndexVersion,
|
||||
|
|
@ -18,23 +18,27 @@ import {
|
|||
import { refreshSessionRowProfiles } from "./session-utils-row.js";
|
||||
|
||||
/** Registry and display facts have their own lifecycle, independent of stored row acquisition. */
|
||||
export function createSessionRowProjectionContext() {
|
||||
export function createSessionRowProjectionContext(subagents: SubagentSessionListReadView) {
|
||||
let preparedEpoch = -1;
|
||||
let registryRevision: number | undefined = getSubagentRegistryPublicationRevision();
|
||||
let registrySnapshot = getSubagentSessionListReadSnapshotIdentity();
|
||||
let registrySnapshot = subagents.snapshotIdentity();
|
||||
let agentRunRevision = readAgentRunIndexVersion();
|
||||
let profileRevision = 0;
|
||||
let subagentRevision = 0;
|
||||
let parentRevision = 0;
|
||||
let modelFactsDirty = false;
|
||||
const identityProjection = createSessionIdentityProjection();
|
||||
const initialNow = Date.now();
|
||||
let current = {
|
||||
...buildSessionListRowMetadataContext({ now: Date.now() }),
|
||||
...buildSessionListRowMetadataContext({
|
||||
now: initialNow,
|
||||
subagentRuns: buildSubagentSessionListReadIndex(initialNow, undefined, subagents.runs()),
|
||||
}),
|
||||
identityProjection,
|
||||
};
|
||||
const subagentInputs = current.subagentRuns.inputs;
|
||||
function prepare(epoch: number) {
|
||||
const snapshot = getSubagentSessionListReadSnapshotIdentity();
|
||||
const snapshot = subagents.snapshotIdentity();
|
||||
const revision = getSubagentRegistryPublicationRevision();
|
||||
const runRevision = readAgentRunIndexVersion();
|
||||
if (
|
||||
|
|
@ -56,7 +60,7 @@ export function createSessionRowProjectionContext() {
|
|||
const subagentRuns =
|
||||
registryRevision === revision
|
||||
? current.subagentRuns.atTime(now)
|
||||
: buildSubagentSessionListReadIndex(now);
|
||||
: buildSubagentSessionListReadIndex(now, undefined, subagents.runs());
|
||||
// Keep maps local to this projection; independent builders still own fresh indexes.
|
||||
const projectedAgentRuns =
|
||||
agentRunRevision === runRevision ? current.projectedAgentRuns : buildProjectedAgentRunIndex();
|
||||
|
|
@ -95,7 +99,7 @@ export function createSessionRowProjectionContext() {
|
|||
parentRevision === subagentRevision &&
|
||||
agentRunRevision === readAgentRunIndexVersion() &&
|
||||
registryRevision === getSubagentRegistryPublicationRevision() &&
|
||||
registrySnapshot === getSubagentSessionListReadSnapshotIdentity()
|
||||
registrySnapshot === subagents.snapshotIdentity()
|
||||
? current
|
||||
: undefined;
|
||||
},
|
||||
|
|
|
|||
|
|
@ -1,7 +1,5 @@
|
|||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import { readAcpSessionMetaForEntries } from "../acp/runtime/session-meta-readonly.js";
|
||||
import { getSubagentSessionListReadSnapshotIdentity } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js";
|
||||
import { captureCanonicalSessionReaderContinuation } from "../config/sessions/session-canonical-key.js";
|
||||
import {
|
||||
assertSessionStoreReadCandidate,
|
||||
|
|
@ -9,7 +7,6 @@ import {
|
|||
} from "../config/sessions/session-store-read-candidates.js";
|
||||
import { withSessionHistoryWorkerDatabases } from "../config/sessions/session-transcript-worker-runtime.js";
|
||||
import { MAX_SESSION_ROW_FACTS_KEYS } from "../config/sessions/session-transcript-worker.types.js";
|
||||
import { resolveStateDir } from "../config/state-dir.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { normalizeAgentId } from "../routing/session-key.js";
|
||||
import { retainOpenClawAgentDatabaseReadCandidates } from "../state/openclaw-agent-db.js";
|
||||
|
|
@ -28,6 +25,8 @@ export async function withSessionRowDatabaseFacts(
|
|||
rows: ReadonlyMap<string, Row>;
|
||||
dirty: ReadonlySet<string>;
|
||||
revision: () => number | undefined;
|
||||
registrySnapshot: () => object | undefined;
|
||||
env: NodeJS.ProcessEnv;
|
||||
cfg: OpenClawConfig;
|
||||
selected?: ReadonlySet<string>;
|
||||
},
|
||||
|
|
@ -40,7 +39,7 @@ export async function withSessionRowDatabaseFacts(
|
|||
},
|
||||
): Promise<void> {
|
||||
const revision = owner.revision();
|
||||
const registrySnapshot = getSubagentSessionListReadSnapshotIdentity();
|
||||
const registrySnapshot = owner.registrySnapshot();
|
||||
const ids: string[] = [];
|
||||
for (const id of owner.selected ?? owner.dirty) {
|
||||
ids.push(id);
|
||||
|
|
@ -66,8 +65,7 @@ export async function withSessionRowDatabaseFacts(
|
|||
}
|
||||
const rows = ids.flatMap((id) => owner.rows.get(id) ?? []);
|
||||
const rowRevisions = new Map(rows.map((row) => [identity(row), row.databaseFactsRevision]));
|
||||
const env = cloneEnvWithPlatformSemantics(process.env);
|
||||
env.OPENCLAW_STATE_DIR = resolveStateDir(env);
|
||||
const env = owner.env;
|
||||
const groups = new Map<
|
||||
string,
|
||||
{
|
||||
|
|
@ -178,7 +176,7 @@ export async function withSessionRowDatabaseFacts(
|
|||
if (
|
||||
revision !== undefined &&
|
||||
owner.revision() === revision &&
|
||||
registrySnapshot === getSubagentSessionListReadSnapshotIdentity()
|
||||
registrySnapshot === owner.registrySnapshot()
|
||||
) {
|
||||
const currentIds = rows
|
||||
.filter(
|
||||
|
|
|
|||
|
|
@ -28,6 +28,8 @@ export function createSessionRowRefresh(
|
|||
registryPrepared: boolean;
|
||||
};
|
||||
databaseRevision: () => number;
|
||||
registrySnapshot: () => object | undefined;
|
||||
env: NodeJS.ProcessEnv;
|
||||
runAsOwner: <T>(operation: () => T) => T;
|
||||
lookup: (query: records.Lookup) => records.Row | undefined;
|
||||
prepareRegistryFacts: () => Promise<void> | undefined;
|
||||
|
|
@ -165,7 +167,15 @@ export function createSessionRowRefresh(
|
|||
}
|
||||
function readExactRows(selected: ReadonlySet<string>) {
|
||||
return withSessionRowDatabaseFacts(
|
||||
{ rows: owner.rows, dirty: owner.dirty, selected, cfg: owner.state().cfg, revision },
|
||||
{
|
||||
rows: owner.rows,
|
||||
dirty: owner.dirty,
|
||||
selected,
|
||||
cfg: owner.state().cfg,
|
||||
revision,
|
||||
registrySnapshot: owner.registrySnapshot,
|
||||
env: owner.env,
|
||||
},
|
||||
{
|
||||
refreshPending: materializer.refreshPending,
|
||||
accept: (ids, facts) => materializer.accept(ids, facts, true),
|
||||
|
|
@ -253,7 +263,15 @@ export function createSessionRowRefresh(
|
|||
}
|
||||
try {
|
||||
await withSessionRowDatabaseFacts(
|
||||
{ rows: owner.rows, dirty: owner.dirty, selected, cfg: owner.state().cfg, revision },
|
||||
{
|
||||
rows: owner.rows,
|
||||
dirty: owner.dirty,
|
||||
selected,
|
||||
cfg: owner.state().cfg,
|
||||
revision,
|
||||
registrySnapshot: owner.registrySnapshot,
|
||||
env: owner.env,
|
||||
},
|
||||
materializer,
|
||||
);
|
||||
} finally {
|
||||
|
|
|
|||
|
|
@ -1,9 +1,6 @@
|
|||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import { listAgentIds, withAgentRosterFactsBatch } from "../agents/agent-scope-config.js";
|
||||
import {
|
||||
getSubagentSessionListReadSnapshotIdentity,
|
||||
prepareSubagentSessionListReadCache,
|
||||
} from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import { createSubagentSessionListReadView } from "../agents/subagents/registry/subagent-registry-state.js";
|
||||
import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js";
|
||||
import { isInternalSessionEffectsKey } from "../config/sessions/internal-session-key.js";
|
||||
import type { InternalSessionEntry as SessionEntry } from "../config/sessions/types.js";
|
||||
|
|
@ -62,8 +59,9 @@ export async function createSessionRowProjection(params: records.ProjectionOptio
|
|||
const env = cloneEnvWithPlatformSemantics(process.env);
|
||||
env.OPENCLAW_STATE_DIR = resolveStateDir(env);
|
||||
const discoveryRead = prepareAgentDatabaseDeletionSnapshotRead({ env }, "runtime");
|
||||
while (!getSubagentSessionListReadSnapshotIdentity()) {
|
||||
await prepareSubagentSessionListReadCache();
|
||||
const subagents = createSubagentSessionListReadView({ env });
|
||||
while (!subagents.snapshotIdentity()) {
|
||||
await subagents.prepare();
|
||||
}
|
||||
let cfg = params.getConfig?.() ?? params.cfg;
|
||||
const getPolicyConfig = () => params.getPolicyConfig?.() ?? cfg;
|
||||
|
|
@ -82,8 +80,8 @@ export async function createSessionRowProjection(params: records.ProjectionOptio
|
|||
let topologyDirty = true,
|
||||
disposed = false;
|
||||
const prepareRegistryFacts = (): Promise<void> | undefined =>
|
||||
!disposed && !inOwnerContext(getSubagentSessionListReadSnapshotIdentity)
|
||||
? inOwnerContext(prepareSubagentSessionListReadCache)
|
||||
!disposed && !inOwnerContext(subagents.snapshotIdentity)
|
||||
? inOwnerContext(subagents.prepare)
|
||||
: undefined;
|
||||
const placementFacts = createSessionRowPlacementProjection(params.placementFactsReader, () =>
|
||||
!disposed && topologyDirty ? topology() : prepareRegistryFacts(),
|
||||
|
|
@ -125,7 +123,7 @@ export async function createSessionRowProjection(params: records.ProjectionOptio
|
|||
void ensureMaterialized().catch(() => {});
|
||||
},
|
||||
});
|
||||
const metadata = createSessionRowProjectionContext();
|
||||
const metadata = createSessionRowProjectionContext(subagents);
|
||||
const backfill = createSessionRowProjectionBackfill({
|
||||
ready: ensureMaterialized,
|
||||
read: (id) => rows.get(id),
|
||||
|
|
@ -315,7 +313,7 @@ export async function createSessionRowProjection(params: records.ProjectionOptio
|
|||
} else if (!presentationOnly) {
|
||||
const query = { ...change, key: change.sessionKey };
|
||||
const exact = matching(query);
|
||||
const registryFactsReady = inOwnerContext(getSubagentSessionListReadSnapshotIdentity);
|
||||
const registryFactsReady = inOwnerContext(subagents.snapshotIdentity);
|
||||
for (const previous of new Set([...exact, ...matching(query, "id")])) {
|
||||
records.invalidateDatabaseFacts(previous);
|
||||
if (previous.entry && change.scope !== "session-entry") {
|
||||
|
|
@ -429,11 +427,13 @@ export async function createSessionRowProjection(params: records.ProjectionOptio
|
|||
} = createSessionRowRefresh({
|
||||
rows,
|
||||
dirty,
|
||||
registrySnapshot: subagents.snapshotIdentity,
|
||||
env,
|
||||
state: () => ({
|
||||
cfg,
|
||||
disposed,
|
||||
topologyDirty,
|
||||
registryPrepared: Boolean(inOwnerContext(getSubagentSessionListReadSnapshotIdentity)),
|
||||
registryPrepared: Boolean(inOwnerContext(subagents.snapshotIdentity)),
|
||||
}),
|
||||
runAsOwner: inOwnerContext,
|
||||
lookup,
|
||||
|
|
|
|||
|
|
@ -199,7 +199,7 @@ export function installGatewaySessionsTestResources(
|
|||
runQaGatewayFixture(
|
||||
() => run(state),
|
||||
disposeSessionReadContexts,
|
||||
// The suite projection also reads this state, but its store lives outside state.root.
|
||||
// Nested fixtures can leave suite projection work pending outside state.root.
|
||||
() => releaseGatewaySessionStoreFixture(requireSharedSessionStoreDir()),
|
||||
),
|
||||
);
|
||||
|
|
|
|||
|
|
@ -1,4 +1,3 @@
|
|||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import path from "node:path";
|
||||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js";
|
||||
|
|
@ -11,7 +10,6 @@ import {
|
|||
} from "../infra/kysely-sync.js";
|
||||
import { isSqliteCorruptionError } from "../infra/sqlite-error-diagnostics.js";
|
||||
import { runSqliteDeferredTransactionSync } from "../infra/sqlite-transaction.js";
|
||||
import { assertExistingDatabaseIdentity } from "../infra/sqlite-worker-identity.js";
|
||||
import { normalizeAgentId } from "../routing/session-key.js";
|
||||
import { sessionChanges } from "../sessions/session-row-changes.js";
|
||||
import { readAgentDeletionRecoveryHolds } from "./agent-deletion-journal-recovery.js";
|
||||
|
|
@ -31,10 +29,7 @@ import {
|
|||
import { tableExists } from "./openclaw-state-db-schema-helpers.js";
|
||||
import type { DB } from "./openclaw-state-db.generated.js";
|
||||
import { resolveOpenClawStateSqlitePath } from "./openclaw-state-db.paths.js";
|
||||
import {
|
||||
captureOpenClawStateReadContext,
|
||||
captureOpenClawStateWorkerContext,
|
||||
} from "./openclaw-state-worker-context.js";
|
||||
import { prepareOpenClawStateReadSource } from "./openclaw-state-worker-context.js";
|
||||
|
||||
/** Completed cleanup still retains a deletion tombstone. */
|
||||
export function readAgentDeletionJournalStatusInDatabase(
|
||||
|
|
@ -186,8 +181,8 @@ export function prepareAgentDatabaseDeletionSnapshotRead(
|
|||
env,
|
||||
path: path.resolve(inputOptions.path ?? resolveOpenClawStateSqlitePath(env)),
|
||||
};
|
||||
const inSourceContext = AsyncLocalStorage.snapshot();
|
||||
const context = captureOpenClawStateWorkerContext(options);
|
||||
const source = prepareOpenClawStateReadSource(options);
|
||||
const context = source.workerContext();
|
||||
const assertCurrent = () => {
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
context.admission.assertCurrent();
|
||||
|
|
@ -213,28 +208,7 @@ export function prepareAgentDatabaseDeletionSnapshotRead(
|
|||
return {
|
||||
read,
|
||||
async readWithCurrentAdmission() {
|
||||
return inSourceContext(() => {
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
const original = context.admission.identity;
|
||||
if (original.key.startsWith("file:")) {
|
||||
assertExistingDatabaseIdentity(options.path, original.key, original.birthtime);
|
||||
} else {
|
||||
// A still-current absent source can bind its first canonical creation.
|
||||
context.admission.assertCurrent();
|
||||
}
|
||||
const current = captureOpenClawStateReadContext(options.path);
|
||||
const source = context.admission.identity;
|
||||
if (
|
||||
current.admission.identity.key !== source.key ||
|
||||
current.admission.identity.birthtime !== source.birthtime ||
|
||||
current.maintenanceScope !== context.maintenanceScope ||
|
||||
current.existingSchemaPath !== context.existingSchemaPath
|
||||
) {
|
||||
throw new Error("Deletion snapshot source changed before read admission");
|
||||
}
|
||||
// A new read may follow handle retirement; an earlier reply retains its own revoked admission.
|
||||
return readSnapshot({ ...context, ...current });
|
||||
});
|
||||
return source.withCurrent(readSnapshot);
|
||||
},
|
||||
async withCurrentSnapshot(consume) {
|
||||
let changed: boolean;
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
import { existsSync, linkSync, mkdirSync, writeFileSync } from "node:fs";
|
||||
import { existsSync, linkSync, mkdirSync, renameSync, writeFileSync } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
|
||||
|
|
@ -16,7 +16,12 @@ import {
|
|||
registerOpenClawStateDatabaseAsyncResource,
|
||||
} from "./openclaw-state-db-cache.js";
|
||||
import { openOpenClawStateReadConnection } from "./openclaw-state-db-read-connection.js";
|
||||
import { withExistingOpenClawStateSchema } from "./openclaw-state-db-schema-policy.js";
|
||||
import { openOpenClawStateDatabase } from "./openclaw-state-db.js";
|
||||
import {
|
||||
captureOpenClawStateReadContext,
|
||||
prepareOpenClawStateReadSource,
|
||||
} from "./openclaw-state-worker-context.js";
|
||||
|
||||
const dirs = useAutoCleanupTempDirTracker((cleanup) =>
|
||||
afterEach(async () => {
|
||||
|
|
@ -30,6 +35,104 @@ function databasePath(name = "state") {
|
|||
}
|
||||
|
||||
describe("canonical shared-state resource drainage", () => {
|
||||
it.each(["ordinary", "existing"] as const)(
|
||||
"reuses prepared %s source facts without resolving or re-admitting the schema",
|
||||
(scope) => {
|
||||
const pathname = databasePath();
|
||||
writeFileSync(pathname, "");
|
||||
const consume = () => {
|
||||
const source = prepareOpenClawStateReadSource({ path: pathname });
|
||||
const context = source.current();
|
||||
const resolve = vi.spyOn(path, "resolve");
|
||||
let reused = true;
|
||||
let resolutions: number;
|
||||
try {
|
||||
for (let index = 0; index < 100; index++) {
|
||||
reused &&= source.current() === context;
|
||||
}
|
||||
resolutions = resolve.mock.calls.length;
|
||||
} finally {
|
||||
resolve.mockRestore();
|
||||
}
|
||||
expect(resolutions).toBe(0);
|
||||
expect(reused).toBe(true);
|
||||
};
|
||||
if (scope === "existing") {
|
||||
withExistingOpenClawStateSchema({ path: pathname }, consume);
|
||||
} else {
|
||||
consume();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("renews prepared reads of the same file without adopting its replacement", async () => {
|
||||
const pathname = databasePath();
|
||||
const source = prepareOpenClawStateReadSource({ path: pathname });
|
||||
const absent = source.current();
|
||||
writeFileSync(pathname, "");
|
||||
const created = source.current();
|
||||
expect(created.admission.identity.key).toMatch(/^file:/);
|
||||
absent.admission.assertCurrent();
|
||||
await closeOpenClawStateDatabaseByPathAsync(pathname);
|
||||
const renewed = source.current();
|
||||
expect(created.admission.assertCurrent).toThrow(/admission changed/);
|
||||
expect(renewed.admission.identity.key).toBe(created.admission.identity.key);
|
||||
await closeOpenClawStateDatabaseByPathAsync(pathname);
|
||||
renameSync(pathname, `${pathname}.retired`);
|
||||
writeFileSync(pathname, "");
|
||||
const replacement = captureOpenClawStateReadContext(pathname);
|
||||
expect(replacement.admission.identity.key).not.toBe(created.admission.identity.key);
|
||||
expect(source.current).toThrow(/identity changed/);
|
||||
expect(source.workerContext).toThrow(/identity changed/);
|
||||
});
|
||||
|
||||
it("reuses warm admission without resolving paths or allocating replacement tokens", () => {
|
||||
const lifecycle = createOpenClawStateDatabaseAsyncLifecycle();
|
||||
const pathname = databasePath();
|
||||
writeFileSync(pathname, "");
|
||||
const retained = lifecycle.capture(pathname);
|
||||
const resolve = vi.spyOn(path, "resolve");
|
||||
let reused = true;
|
||||
let resolutions: number;
|
||||
try {
|
||||
for (let index = 0; index < 100; index++) {
|
||||
const admission = lifecycle.capture(pathname);
|
||||
admission.assertCurrent();
|
||||
reused &&= admission === retained;
|
||||
}
|
||||
resolutions = resolve.mock.calls.length;
|
||||
} finally {
|
||||
resolve.mockRestore();
|
||||
}
|
||||
expect(resolutions).toBe(0);
|
||||
expect(reused).toBe(true);
|
||||
lifecycle.invalidate(pathname);
|
||||
expect(retained.assertCurrent).toThrow(/admission changed/);
|
||||
const renewed = lifecycle.capture(pathname);
|
||||
expect(renewed).not.toBe(retained);
|
||||
renewed.assertCurrent();
|
||||
});
|
||||
|
||||
it("keeps captured schema scope lifetime separate from shared physical admission", async () => {
|
||||
const pathname = databasePath();
|
||||
const ordinary = captureOpenClawStateReadContext(pathname);
|
||||
const restricted = withExistingOpenClawStateSchema({ path: pathname }, () => {
|
||||
const context = captureOpenClawStateReadContext(pathname);
|
||||
writeFileSync(pathname, "");
|
||||
const created = captureOpenClawStateDatabaseReadAdmission(pathname);
|
||||
expect(context.admission.identity.key).toBe(created.identity.key);
|
||||
expect(context.admission.identity.key).toMatch(/^file:/);
|
||||
context.admission.assertCurrent();
|
||||
return { context, source: prepareOpenClawStateReadSource({ path: pathname }) };
|
||||
});
|
||||
expect(restricted.context.admission.assertCurrent).toThrow(/schema admission has ended/);
|
||||
expect(restricted.source.current).toThrow(/schema admission has ended/);
|
||||
ordinary.admission.assertCurrent();
|
||||
captureOpenClawStateReadContext(pathname).admission.assertCurrent();
|
||||
await closeOpenClawStateDatabaseByPathAsync(pathname);
|
||||
expect(restricted.source.current).toThrow(/schema admission has ended/);
|
||||
});
|
||||
|
||||
it.each(["missing", "directory"] as const)(
|
||||
"keeps unrelated owners while closing a never-admitted %s path",
|
||||
async (kind) => {
|
||||
|
|
|
|||
|
|
@ -35,6 +35,7 @@ type IdentityRecord = {
|
|||
identity: DatabasePathIdentity;
|
||||
paths: Set<string>;
|
||||
generation: object;
|
||||
admissions: Map<string, OpenClawStateDatabaseReadAdmission>;
|
||||
};
|
||||
type ReadSeal = { record?: IdentityRecord };
|
||||
type CloseAttempt = {
|
||||
|
|
@ -327,18 +328,41 @@ export function createOpenClawDatabaseMaintenanceScope(
|
|||
export function createOpenClawStateDatabaseAsyncLifecycle() {
|
||||
const resources = new Set<OpenClawStateDatabaseAsyncResource>();
|
||||
const records = new Map<string, IdentityRecord>();
|
||||
const recordsByPath = new Map<string, IdentityRecord>();
|
||||
const seals = new Set<ReadSeal>();
|
||||
const attempts = new Map<IdentityRecord | undefined, CloseAttempt>();
|
||||
let tail = Promise.resolve();
|
||||
|
||||
// Internal lookups reuse the path normalized at the lifecycle boundary.
|
||||
const known = (resolvedPath: string) =>
|
||||
[...records.values()].find((record) => record.paths.has(resolvedPath));
|
||||
const overlaps = (left: IdentityRecord, right: IdentityRecord) =>
|
||||
left.identity.key === right.identity.key ||
|
||||
[...left.paths].some((pathname) => right.paths.has(pathname));
|
||||
const isSealed = (record: IdentityRecord) =>
|
||||
[...seals].some((held) => held.record === undefined || overlaps(held.record, record));
|
||||
const known = (pathname: string) =>
|
||||
recordsByPath.get(pathname) ?? recordsByPath.get(path.resolve(pathname));
|
||||
const bindPath = (record: IdentityRecord, pathname: string) => {
|
||||
record.paths.add(pathname);
|
||||
if (!recordsByPath.has(pathname)) {
|
||||
recordsByPath.set(pathname, record);
|
||||
}
|
||||
};
|
||||
const overlaps = (left: IdentityRecord, right: IdentityRecord) => {
|
||||
if (left.identity.key === right.identity.key) {
|
||||
return true;
|
||||
}
|
||||
for (const pathname of left.paths) {
|
||||
if (right.paths.has(pathname)) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
};
|
||||
const isSealed = (record: IdentityRecord) => {
|
||||
if (seals.size === 0) {
|
||||
return false;
|
||||
}
|
||||
for (const held of seals) {
|
||||
if (held.record === undefined || overlaps(held.record, record)) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
};
|
||||
const assertOpen = (record: IdentityRecord) => {
|
||||
if (isSealed(record)) {
|
||||
throw new StateDatabaseReadAdmissionInvalidatedError(
|
||||
|
|
@ -373,7 +397,7 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
resolvedPath: string,
|
||||
preparedIdentity?: DatabasePathIdentity,
|
||||
): IdentityRecord => {
|
||||
const cached = known(resolvedPath);
|
||||
const cached = recordsByPath.get(resolvedPath);
|
||||
if (cached && (!preparedIdentity || cached.identity.key === preparedIdentity.key)) {
|
||||
// Resolve first creation without replacing an established file's admission.
|
||||
return !preparedIdentity && cached.identity.key.startsWith("path:")
|
||||
|
|
@ -406,14 +430,16 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
}
|
||||
}
|
||||
if (!record) {
|
||||
record = { identity, paths: new Set(), generation: {} };
|
||||
record = { identity, paths: new Set(), generation: {}, admissions: new Map() };
|
||||
records.set(identity.key, record);
|
||||
}
|
||||
record.paths.add(resolvedPath).add(identity.canonicalPath);
|
||||
bindPath(record, resolvedPath);
|
||||
bindPath(record, identity.canonicalPath);
|
||||
return record;
|
||||
};
|
||||
const resolveForNative = (resolvedPath: string): IdentityRecord | undefined => {
|
||||
const cached = known(resolvedPath);
|
||||
const resolveForNative = (pathname: string): IdentityRecord | undefined => {
|
||||
const resolvedPath = path.resolve(pathname);
|
||||
const cached = recordsByPath.get(resolvedPath);
|
||||
if (cached) {
|
||||
return cached;
|
||||
}
|
||||
|
|
@ -423,6 +449,7 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
const invalidate = (record?: IdentityRecord) => {
|
||||
for (const current of record ? [record] : records.values()) {
|
||||
current.generation = {};
|
||||
current.admissions.clear();
|
||||
}
|
||||
};
|
||||
const seal = (record?: IdentityRecord): ReadSeal => {
|
||||
|
|
@ -434,20 +461,60 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
const forget = (record: IdentityRecord) => {
|
||||
if (!isSealed(record) && records.get(record.identity.key) === record) {
|
||||
records.delete(record.identity.key);
|
||||
for (const pathname of record.paths) {
|
||||
if (recordsByPath.get(pathname) === record) {
|
||||
recordsByPath.delete(pathname);
|
||||
// A sealed predecessor keeps close custody until retirement, even if
|
||||
// publication has already recorded a replacement at the same path.
|
||||
for (const replacement of records.values()) {
|
||||
if (replacement.paths.has(pathname)) {
|
||||
recordsByPath.set(pathname, replacement);
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
// Keep closure creation off capture's warm branch: V8 otherwise allocates its
|
||||
// captured environment even when returning an already retained admission.
|
||||
const captureResolved = (databasePath: string): OpenClawStateDatabaseReadAdmission => {
|
||||
const record = resolve(databasePath);
|
||||
assertOpen(record);
|
||||
const previous = record.admissions.get(databasePath);
|
||||
if (previous) {
|
||||
return previous;
|
||||
}
|
||||
const generation = record.generation;
|
||||
const admission: OpenClawStateDatabaseReadAdmission = Object.freeze({
|
||||
databasePath,
|
||||
get identity() {
|
||||
return record.identity;
|
||||
},
|
||||
assertCurrent() {
|
||||
assertOpen(record);
|
||||
if (records.get(record.identity.key) !== record || record.generation !== generation) {
|
||||
throw new StateDatabaseReadAdmissionInvalidatedError(
|
||||
"OpenClaw state database read admission changed",
|
||||
);
|
||||
}
|
||||
},
|
||||
});
|
||||
record.admissions.set(databasePath, admission);
|
||||
return admission;
|
||||
};
|
||||
|
||||
return {
|
||||
identity(pathname: string): DatabasePathIdentity | undefined {
|
||||
return known(path.resolve(pathname))?.identity ?? inspectDatabasePathIdentitySync(pathname);
|
||||
return known(pathname)?.identity ?? inspectDatabasePathIdentitySync(pathname);
|
||||
},
|
||||
knownIdentity(this: void, pathname: string): DatabasePathIdentity | undefined {
|
||||
return known(path.resolve(pathname))?.identity;
|
||||
return known(pathname)?.identity;
|
||||
},
|
||||
publish(pathname: string): DatabasePathIdentity {
|
||||
const resolvedPath = path.resolve(pathname);
|
||||
const identity = readDatabasePathIdentitySync(resolvedPath);
|
||||
const previous = known(resolvedPath);
|
||||
const previous = recordsByPath.get(resolvedPath);
|
||||
let record = findPhysicalRecord(identity);
|
||||
if (previous && previous.identity.key !== identity.key) {
|
||||
if (previous.identity.key.startsWith("path:") && !record) {
|
||||
|
|
@ -464,14 +531,15 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
if (!record) {
|
||||
record = resolve(resolvedPath, identity);
|
||||
}
|
||||
record.paths.add(resolvedPath).add(identity.canonicalPath);
|
||||
bindPath(record, resolvedPath);
|
||||
bindPath(record, identity.canonicalPath);
|
||||
return identity;
|
||||
},
|
||||
invalidate(pathname?: string): void {
|
||||
if (pathname === undefined) {
|
||||
invalidate();
|
||||
} else {
|
||||
const record = known(path.resolve(pathname));
|
||||
const record = known(pathname);
|
||||
if (record) {
|
||||
invalidate(record);
|
||||
}
|
||||
|
|
@ -487,24 +555,13 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
};
|
||||
},
|
||||
capture(this: void, pathname: string): OpenClawStateDatabaseReadAdmission {
|
||||
const databasePath = path.resolve(pathname);
|
||||
const record = resolve(databasePath);
|
||||
assertOpen(record);
|
||||
const generation = record.generation;
|
||||
return {
|
||||
databasePath,
|
||||
get identity() {
|
||||
return record.identity;
|
||||
},
|
||||
assertCurrent() {
|
||||
assertOpen(record);
|
||||
if (records.get(record.identity.key) !== record || record.generation !== generation) {
|
||||
throw new StateDatabaseReadAdmissionInvalidatedError(
|
||||
"OpenClaw state database read admission changed",
|
||||
);
|
||||
}
|
||||
},
|
||||
};
|
||||
const cached = recordsByPath.get(pathname);
|
||||
const retained = cached?.admissions.get(pathname);
|
||||
if (cached && retained && cached.identity.key.startsWith("file:")) {
|
||||
retained.assertCurrent();
|
||||
return retained;
|
||||
}
|
||||
return captureResolved(path.resolve(pathname));
|
||||
},
|
||||
holdExclusion(pathname: string): () => void {
|
||||
const record = resolve(path.resolve(pathname));
|
||||
|
|
@ -523,7 +580,7 @@ export function createOpenClawStateDatabaseAsyncLifecycle() {
|
|||
pathname: string | undefined,
|
||||
retireNative: (identity?: DatabasePathIdentity) => boolean,
|
||||
): Promise<boolean> {
|
||||
const record = pathname === undefined ? undefined : resolveForNative(path.resolve(pathname));
|
||||
const record = pathname === undefined ? undefined : resolveForNative(pathname);
|
||||
if (pathname !== undefined && !record) {
|
||||
// No worker could enter a non-file target. Retire only the caller's exact
|
||||
// native path; undefined must not reach resource.close as a global drain.
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
import { renameSync, writeFileSync } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { afterEach, expect, it } from "vitest";
|
||||
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
|
||||
|
|
@ -5,6 +6,32 @@ import { createDeferredCore } from "../shared/deferred.js";
|
|||
import { createOpenClawStateDatabaseAsyncLifecycle } from "./openclaw-state-db-async-lifecycle.js";
|
||||
|
||||
const dirs = useAutoCleanupTempDirTracker(afterEach);
|
||||
|
||||
it("coalesces close custody while a sealed path is rebound to a replacement", async () => {
|
||||
const owner = createOpenClawStateDatabaseAsyncLifecycle();
|
||||
const file = path.join(dirs.make("pending-state-replacement-"), "state.sqlite");
|
||||
writeFileSync(file, "original");
|
||||
const admission = owner.capture(file);
|
||||
const finish = createDeferredCore();
|
||||
owner.register({ close: () => finish.promise });
|
||||
const closing = owner.close(file, () => false);
|
||||
try {
|
||||
renameSync(file, `${file}.retired`);
|
||||
writeFileSync(file, "replacement");
|
||||
const replacement = owner.publish(file);
|
||||
expect(replacement.key).not.toBe(admission.identity.key);
|
||||
expect(owner.close(file, () => false)).toBe(closing);
|
||||
expect(() => owner.capture(file)).toThrow(/admission is closed/);
|
||||
} finally {
|
||||
finish.resolve();
|
||||
await closing;
|
||||
}
|
||||
expect(admission.assertCurrent).toThrow(/admission changed/);
|
||||
const replacement = owner.capture(file);
|
||||
expect(replacement.identity.key).not.toBe(admission.identity.key);
|
||||
replacement.assertCurrent();
|
||||
});
|
||||
|
||||
it("joins a shared resource registered during a pending close before native retirement", async () => {
|
||||
const owner = createOpenClawStateDatabaseAsyncLifecycle();
|
||||
const file = path.join(dirs.make("pending-state-owner-"), "state.sqlite");
|
||||
|
|
|
|||
|
|
@ -62,14 +62,28 @@ export function withExistingOpenClawStateSchema<T>(
|
|||
}
|
||||
}
|
||||
|
||||
export function getExistingOpenClawStateSchemaPath(): string | undefined {
|
||||
const scope = schemaPolicies.scopes.getStore();
|
||||
function assertSchemaScopeActive(scope: ExistingSchemaScope | undefined): void {
|
||||
if (scope && !scope.active) {
|
||||
throw new Error("Existing shared-state schema admission has ended.");
|
||||
}
|
||||
}
|
||||
|
||||
export function getExistingOpenClawStateSchemaPath(): string | undefined {
|
||||
const scope = schemaPolicies.scopes.getStore();
|
||||
assertSchemaScopeActive(scope);
|
||||
return scope?.path;
|
||||
}
|
||||
|
||||
/** Admit this source once; retained reads only need the original scope's live lifetime. */
|
||||
export function captureOpenClawStateSchemaReadAdmission(pathname: string) {
|
||||
const scope = schemaPolicies.scopes.getStore();
|
||||
if (!scope) {
|
||||
return undefined;
|
||||
}
|
||||
isExistingOpenClawStateSchema(pathname);
|
||||
return { path: scope.path, assertCurrent: () => assertSchemaScopeActive(scope) };
|
||||
}
|
||||
|
||||
/** Check supplied and cached handles before exposing them to another admission policy. */
|
||||
export function isExistingOpenClawStateSchema(pathname: string, database?: DatabaseSync): boolean {
|
||||
const scopedPath = getExistingOpenClawStateSchemaPath();
|
||||
|
|
|
|||
|
|
@ -174,8 +174,11 @@ it("preserves the caller's typed admission refusal through lease acquisition", a
|
|||
const capture = workerContext.captureOpenClawStateWorkerContext;
|
||||
vi.spyOn(workerContext, "captureOpenClawStateWorkerContext").mockImplementation((options) => {
|
||||
const context = capture(options);
|
||||
context.admission.assertCurrent = () => {
|
||||
throw refusal;
|
||||
context.admission = {
|
||||
...context.admission,
|
||||
assertCurrent() {
|
||||
throw refusal;
|
||||
},
|
||||
};
|
||||
return context;
|
||||
});
|
||||
|
|
|
|||
|
|
@ -672,6 +672,13 @@ it.each(["single", "union"] as const)(
|
|||
async (shape) => {
|
||||
const { options } = source();
|
||||
const context = captureOpenClawStateWorkerContext(options);
|
||||
const admission = context.admission;
|
||||
context.admission = {
|
||||
...admission,
|
||||
get identity() {
|
||||
return admission.identity;
|
||||
},
|
||||
};
|
||||
const selector = "任务🦞".repeat(512);
|
||||
const scope = {
|
||||
taskId: selector,
|
||||
|
|
|
|||
|
|
@ -4,35 +4,41 @@ import { cloneEnvWithPlatformSemantics } from "../config/config-env-vars.js";
|
|||
import { resolveStateDir } from "../config/state-dir.js";
|
||||
import { isGatewayExternallySupervised } from "../infra/gateway-supervision.js";
|
||||
import { mergeProcessEnv } from "../infra/process-env.js";
|
||||
import { assertExistingDatabaseIdentity } from "../infra/sqlite-worker-identity.js";
|
||||
import { captureStateDatabaseCoordinatorRuntime } from "../infra/state-database-coordinator.js";
|
||||
import { getOpenClawDatabaseMaintenanceScope } from "./openclaw-state-db-async-lifecycle.js";
|
||||
import { captureOpenClawStateDatabaseReadAdmission } from "./openclaw-state-db-cache.js";
|
||||
import {
|
||||
getExistingOpenClawStateSchemaPath,
|
||||
isExistingOpenClawStateSchema,
|
||||
} from "./openclaw-state-db-schema-policy.js";
|
||||
getOpenClawDatabaseMaintenanceScope,
|
||||
isStateDatabaseReadAdmissionInvalidatedError,
|
||||
} from "./openclaw-state-db-async-lifecycle.js";
|
||||
import { captureOpenClawStateDatabaseReadAdmission } from "./openclaw-state-db-cache.js";
|
||||
import { captureOpenClawStateSchemaReadAdmission } from "./openclaw-state-db-schema-policy.js";
|
||||
import { resolveOpenClawStateSqlitePath } from "./openclaw-state-db.paths.js";
|
||||
import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js";
|
||||
|
||||
export type OpenClawStateReadContext = Pick<
|
||||
OpenClawStateWorkerContext,
|
||||
"admission" | "maintenanceScope" | "existingSchemaPath" | "runInCapturedSchemaScope"
|
||||
>;
|
||||
|
||||
/** Capture read authority without constructing a worker environment. */
|
||||
export function captureOpenClawStateReadContext(
|
||||
pathname = resolveOpenClawStateSqlitePath(),
|
||||
): Pick<
|
||||
OpenClawStateWorkerContext,
|
||||
"admission" | "maintenanceScope" | "existingSchemaPath" | "runInCapturedSchemaScope"
|
||||
> {
|
||||
const databasePath = path.resolve(pathname);
|
||||
isExistingOpenClawStateSchema(databasePath);
|
||||
const existingSchemaPath = getExistingOpenClawStateSchemaPath();
|
||||
const admission = captureOpenClawStateDatabaseReadAdmission(databasePath);
|
||||
): OpenClawStateReadContext {
|
||||
const schema = captureOpenClawStateSchemaReadAdmission(pathname);
|
||||
const capturedAdmission = captureOpenClawStateDatabaseReadAdmission(pathname);
|
||||
let admission = capturedAdmission;
|
||||
let runInCapturedSchemaScope: OpenClawStateWorkerContext["runInCapturedSchemaScope"];
|
||||
if (existingSchemaPath !== undefined) {
|
||||
if (schema) {
|
||||
const inCapturedScope = AsyncLocalStorage.snapshot();
|
||||
const assertCurrent = admission.assertCurrent;
|
||||
admission.assertCurrent = () => {
|
||||
assertCurrent();
|
||||
// Queued dispatch may run outside this caller, but its captured scope must still be active.
|
||||
inCapturedScope(getExistingOpenClawStateSchemaPath);
|
||||
admission = {
|
||||
databasePath: capturedAdmission.databasePath,
|
||||
get identity() {
|
||||
return capturedAdmission.identity;
|
||||
},
|
||||
assertCurrent() {
|
||||
capturedAdmission.assertCurrent();
|
||||
schema.assertCurrent();
|
||||
},
|
||||
};
|
||||
runInCapturedSchemaScope = (operation) =>
|
||||
inCapturedScope(() => {
|
||||
|
|
@ -43,11 +49,73 @@ export function captureOpenClawStateReadContext(
|
|||
return {
|
||||
maintenanceScope: getOpenClawDatabaseMaintenanceScope(),
|
||||
admission,
|
||||
existingSchemaPath,
|
||||
existingSchemaPath: schema?.path,
|
||||
runInCapturedSchemaScope,
|
||||
};
|
||||
}
|
||||
|
||||
/** Resident readers retain their source and schema policy without re-admitting each publication. */
|
||||
export function prepareOpenClawStateReadSource(input: { path: string; env?: NodeJS.ProcessEnv }) {
|
||||
const env = cloneEnvWithPlatformSemantics(input.env ?? process.env);
|
||||
env.OPENCLAW_STATE_DIR = resolveStateDir(env);
|
||||
const options = { path: path.resolve(input.path), env };
|
||||
const inSourceContext = AsyncLocalStorage.snapshot();
|
||||
const original = captureOpenClawStateReadContext(options.path);
|
||||
let context = original;
|
||||
let worker: OpenClawStateWorkerContext | undefined;
|
||||
|
||||
const refresh = () => {
|
||||
original.maintenanceScope?.assertAdmission();
|
||||
const identity = original.admission.identity;
|
||||
if (identity.key.startsWith("file:")) {
|
||||
assertExistingDatabaseIdentity(options.path, identity.key, identity.birthtime);
|
||||
} else {
|
||||
original.admission.assertCurrent();
|
||||
}
|
||||
const next = captureOpenClawStateReadContext(options.path);
|
||||
const source = original.admission.identity;
|
||||
if (
|
||||
next.admission.identity.key !== source.key ||
|
||||
next.admission.identity.birthtime !== source.birthtime ||
|
||||
next.maintenanceScope !== original.maintenanceScope ||
|
||||
next.existingSchemaPath !== original.existingSchemaPath
|
||||
) {
|
||||
throw new Error("Prepared state read source changed before read admission");
|
||||
}
|
||||
return (context = next);
|
||||
};
|
||||
const current = () => {
|
||||
context.maintenanceScope?.assertAdmission();
|
||||
try {
|
||||
context.admission.assertCurrent();
|
||||
if (context.admission.identity.key.startsWith("file:")) {
|
||||
return context;
|
||||
}
|
||||
} catch (error) {
|
||||
if (!isStateDatabaseReadAdmissionInvalidatedError(error)) {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
return inSourceContext(refresh);
|
||||
};
|
||||
const prepareWorker = () => {
|
||||
// Actual reads verify the file even when no local publication revoked its admission.
|
||||
const next = refresh();
|
||||
worker ??= captureOpenClawStateWorkerContext(options);
|
||||
if (worker.admission !== next.admission) {
|
||||
worker = { ...worker, ...next };
|
||||
}
|
||||
return worker;
|
||||
};
|
||||
return {
|
||||
current,
|
||||
workerContext: () => inSourceContext(prepareWorker),
|
||||
withCurrent<T>(consume: (context: OpenClawStateWorkerContext) => T): T {
|
||||
return inSourceContext(() => consume(prepareWorker()));
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** Capture host facts before asynchronous work, without opening SQLite. */
|
||||
export function captureOpenClawStateWorkerContext(
|
||||
options: {
|
||||
|
|
|
|||
|
|
@ -224,6 +224,13 @@ describe("restored task flow synchronization", () => {
|
|||
store: {
|
||||
...first.store,
|
||||
async syncTaskFlowAsync(this: TaskRegistryStore, context, params) {
|
||||
const admission = context.admission;
|
||||
context.admission = {
|
||||
...admission,
|
||||
get identity() {
|
||||
return admission.identity;
|
||||
},
|
||||
};
|
||||
let failure: unknown;
|
||||
try {
|
||||
retryCalls += 1;
|
||||
|
|
|
|||
|
|
@ -527,10 +527,19 @@ it.each(["current row", "retired store", "retired admission"] as const)(
|
|||
async () => readRows(),
|
||||
);
|
||||
const releaseRead = createDeferred();
|
||||
const asyncRead = vi.spyOn(store, "loadMutationSnapshotAsync").mockImplementation(async () => {
|
||||
await releaseRead.promise;
|
||||
return readRows();
|
||||
});
|
||||
const asyncRead = vi
|
||||
.spyOn(store, "loadMutationSnapshotAsync")
|
||||
.mockImplementation(async (readContext) => {
|
||||
const admission = readContext.admission;
|
||||
readContext.admission = {
|
||||
...admission,
|
||||
get identity() {
|
||||
return admission.identity;
|
||||
},
|
||||
};
|
||||
await releaseRead.promise;
|
||||
return readRows();
|
||||
});
|
||||
const syncRead = vi.spyOn(store, "loadSnapshot").mockImplementation(() => {
|
||||
throw new Error("Unexpected synchronous task refresh");
|
||||
});
|
||||
|
|
|
|||
|
|
@ -384,6 +384,13 @@ describe("asynchronous registry restoration", () => {
|
|||
);
|
||||
const retirementError = new Error("Synthetic task publication admission retired");
|
||||
const context = captureOpenClawStateWorkerContext();
|
||||
const admission = context.admission;
|
||||
context.admission = {
|
||||
...admission,
|
||||
get identity() {
|
||||
return admission.identity;
|
||||
},
|
||||
};
|
||||
const observed: string[] = [];
|
||||
let loads = 0;
|
||||
configureTaskFlowRegistryRuntime({
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import type { Server } from "node:net";
|
||||
import { platform } from "node:os";
|
||||
import { runQaGatewayFixture } from "../../test/helpers/qa-gateway-cleanup.js";
|
||||
import { hasErrnoCode } from "../infra/errno.js";
|
||||
import { FILE_LOCK_TIMEOUT_ERROR_CODE } from "../infra/file-lock.js";
|
||||
|
|
@ -111,13 +112,16 @@ export async function reserveTestPortListener<T extends Server>(params: {
|
|||
() => verifyCleanup(claim.release),
|
||||
);
|
||||
} catch (rollbackError) {
|
||||
// An unrelated listener can win after the free-port probe closes. Only
|
||||
// initial automatic selection may move, after both provisional owners drain.
|
||||
// Windows can deny a candidate after its probe too. Only initial automatic
|
||||
// selection may move, after both provisional owners drain.
|
||||
if (
|
||||
rollbackError !== error ||
|
||||
params.port !== undefined ||
|
||||
error !== bindError ||
|
||||
!hasErrnoCode(error, "EADDRINUSE")
|
||||
!(
|
||||
hasErrnoCode(error, "EADDRINUSE") ||
|
||||
(platform() === "win32" && hasErrnoCode(error, "EACCES"))
|
||||
)
|
||||
) {
|
||||
throw rollbackError;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import fs from "node:fs/promises";
|
||||
import { syncBuiltinESMExports } from "node:module";
|
||||
import net from "node:net";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
|
|
@ -21,6 +22,57 @@ import { createDeferred, withTestTimeout } from "./promise.js";
|
|||
import { runQaGatewayFixture } from "./qa-gateway-cleanup.js";
|
||||
|
||||
describe("createOpenClawTestInstance acquisition", () => {
|
||||
it.each([
|
||||
{ platform: "win32", explicit: false, advances: true },
|
||||
{ platform: "win32", explicit: true, advances: false },
|
||||
{ platform: "darwin", explicit: false, advances: false },
|
||||
] as const)(
|
||||
"preserves reservation policy after $platform EACCES (explicit=$explicit)",
|
||||
async ({ platform, explicit, advances }) => {
|
||||
const port = await testPorts.getDeterministicFreePortBlock({ offsets: [0] });
|
||||
const denied = Object.assign(new Error("candidate listener denied"), { code: "EACCES" });
|
||||
const platformSpy = vi.spyOn(os, "platform").mockReturnValue(platform);
|
||||
syncBuiltinESMExports();
|
||||
const pickerSpy = vi
|
||||
.spyOn(testPorts, "getDeterministicFreePortBlock")
|
||||
.mockResolvedValueOnce(port);
|
||||
let first = true;
|
||||
let reservation: Awaited<ReturnType<typeof reserveTestPortListener>> | undefined;
|
||||
try {
|
||||
const pending = reserveTestPortListener({
|
||||
offsets: [0],
|
||||
...(explicit ? { port } : {}),
|
||||
createListener: () => {
|
||||
const listener = net.createServer();
|
||||
if (first) {
|
||||
first = false;
|
||||
vi.spyOn(listener, "listen").mockImplementationOnce(() => {
|
||||
queueMicrotask(() => listener.emit("error", denied));
|
||||
return listener;
|
||||
});
|
||||
}
|
||||
return listener;
|
||||
},
|
||||
});
|
||||
if (advances) {
|
||||
reservation = await pending;
|
||||
expect(reservation.claim.port).not.toBe(port);
|
||||
expect(reservation.listener.listening).toBe(true);
|
||||
} else {
|
||||
await expect(pending).rejects.toBe(denied);
|
||||
}
|
||||
const released = await acquireTestPortBlock({ port, offsets: [0] });
|
||||
await released.release();
|
||||
} finally {
|
||||
platformSpy.mockRestore();
|
||||
syncBuiltinESMExports();
|
||||
pickerSpy.mockRestore();
|
||||
await reservation?.releaseListener();
|
||||
await reservation?.claim.release();
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it.each([
|
||||
{
|
||||
name: "child-process",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue