refactor(agents): move workspace state operations off the Gateway thread (#163496)

* perf(agents): move workspace state operations into shared workers

Keep canonical snapshots, conditional alias registration, first-writer setup merges, and exact vanished-workspace expiry off the Gateway thread. Preserve live filesystem and caller grants at transaction and commit, recent/future protection, and acknowledged cache retirement through native settlement.

* fix(agents): evaluate workspace recovery holds in worker transactions

Split workspace guards into SQL-free host authority and a transaction-local recovery predicate. Preserve the recovery reader and typed refusal across worker transport, and prove creation grants avoid host SQL and snapshot children.

Route bootstrap and Bun Talk fixtures through the host broker, join Windows worker teardown, and keep the assertion baseline shrink-only.

* fix(test): preserve sorted Gateway worker ownership

Keep the Talk worker fixture in canonical order so CI planning and CLI ownership select the same files.

* fix: preserve workspace mutation guards across worker admission

Keep the released plugin callback compatible and deprecated, with actionable synchronous database access refusal inside worker grants. Recheck recovery holds before host filesystem effects and preserve callback refusals through auxiliary attestation.

* fix: allow legacy workspace guards before worker dispatch

Preserve synchronous database access in released plugin callbacks with one deprecation warning. Keep typed guards at commit admission and recovery checks before filesystem effects. Extract persisted state-owner record admission without changing ownership policy.

* fix(agents): avoid snapshots in workspace recovery guards

Keep current host recovery checks native outside worker grants while preserving private snapshots for artifact-preserving and schema-admission reads. Prove that real workspace preparation during creation spawns no synchronous snapshots, retaining the grant SQL boundary and late-hold regressions.

* fix: declare the ownership guard callback as a function

Use a function-valued callback contract and avoid shadowing the workspace test namespace. Both corrections preserve runtime behavior and pass focused type-aware lint and independent P2 review.

* fix(test): route Talk workers and settle cleanup fixtures
This commit is contained in:
Peter Steinberger 2026-10-02 16:26:06 -07:00 • committed by GitHub
parent 73fe145a1e
commit 6c1e1667bd
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
45 changed files with 1717 additions and 601 deletions

View file

@ -1614,7 +1614,6 @@ src/agents/transcript-redact.ts 2
src/agents/transport-message-transform.ts 2
src/agents/utils/frontmatter.ts 2
src/agents/workspace-legacy-state.ts 2
src/agents/workspace-state-store.ts 1
src/agents/workspace.ts 5
src/agents/worktrees/run-lease.ts 1
src/audit/execution-identity-admission.ts 3

View file

@ -9,6 +9,26 @@ sidebarTitle: "How to migrate"
The ordered migration steps. Work through them in order; each step is self-contained. Part of the [Plugin SDK migration](/plugins/sdk-migration) guide.
## Workspace mutation guards
Await `api.runtime.agent.ensureAgentWorkspace({ dir, guard: { assertHost } })`.
Prepare database-derived inputs asynchronously before calling it; `assertHost`
must synchronously check current caller authority without accessing SQLite.
Core-owned recovery predicates execute on the worker's transaction connection.
The released `beforePersistentApply: () => void` option remains supported for
TypeScript and JavaScript plugins until the next Plugin SDK major. It runs on the
host once immediately before each worker mutation dispatch, outside admission
grants, and at host filesystem mutation boundaries. Throwing stops that apply.
Synchronous OpenClaw database access in the callback is allowed and deprecated;
a warning explains the timing and typed replacement once per process.
There is no compatibility break for legacy callbacks or their database reads.
The timing nuance is that the legacy check runs just before dispatch, while
`guard.assertHost` is also rechecked inside transaction and commit grants.
Prefer the typed guard for live revocation at commit. Callback errors continue
to propagate. No schema, retention, durability, or update migration is required.
## Await session transcript persistence
Use the awaited `SessionManager` methods from

View file

@ -528,6 +528,28 @@ owners, so a nested scope cannot close an executor retained by its parent.
## Carry facts, publish after commit
Workspace snapshots and conditional alias registration, first-writer setup merges,
and exact expired-state deletion use the shared-state writer. Read-only snapshots
retain the existing reader. The host captures the physical database and filesystem
evidence before waiting, rechecks current authority and workspace identity at
transaction and commit admission, and validates the evidence after delivery.
Workspace guards separate SQL-free host authority from a serialized recovery-hold
predicate. The shared recovery reader evaluates that predicate on the worker's
transaction connection before commit; refusal preserves the caller's duplicate-agent
error. Creation guards never recursively read that database from a host grant.
Host filesystem mutations retain a separate recovery-aware callback after awaited
preparation and immediately before each effect. It uses the same recovery kernel
through the existing current read-only connection, outside all worker grants.
These host guards explicitly allow a native read when the worker owns the cached
writer, avoiding a snapshot subprocess per file mutation. Artifact-preserving
scopes and schema-admission reads still select private snapshots; current guards
never reuse an inherited discovery snapshot.
Expiry rereads current setup and attestation rows, preserving the 24-hour and
future-timestamp protections. Native commit receipts retire the stored workspace's
file cache even if ordinary result delivery fails; uncertain writes are never
replayed. Explicit agent deletion and Doctor relocation retain their existing
transaction owners. Schemas, retention, durability, and update behavior are unchanged.
Session branch summaries retain compact counts and headlines in the transcript
read worker, keyed by physical database identity and the transcript rewrite/append
watermark. After a complete scan verifies unique indexed identities and backward

View file

@ -197,6 +197,7 @@ scripts/watch-pr-ci.mts
scripts/windows-cmd-helpers.mjs
src/acp/event-ledger-bytes.ts
src/agents/agent-bundle-mcp-names.ts
src/agents/agent-create-error.ts
src/agents/agent-dir-registry.ts
src/agents/agent-roster.ts
src/agents/agent-run-terminal-receipt.ts
@ -749,6 +750,7 @@ src/infra/gateway-owner-lease.read.ts
src/infra/gateway-owner-lease.ts
src/infra/gateway-lock-payload.ts
src/infra/gateway-state-owner-directory.ts
src/infra/gateway-state-owner-record.ts
src/infra/gateway-state-owner.ts
src/infra/gateway-process-argv.ts
src/infra/gateway-scheduler.ts
@ -1069,6 +1071,7 @@ src/plugins/compat/plugin-sdk-subpath-records.ts
src/plugins/compat/registry-records.ts
src/plugins/compat/registry.ts
src/plugins/compat/session-persistence-records.ts
src/plugins/compat/workspace-mutation-guard.ts
src/plugins/config-activation-shared.ts
src/plugins/config-contract-matches.ts
src/plugins/config-normalization-shared.ts
@ -1403,6 +1406,7 @@ src/state/agent-database-startup.ts
src/state/agent-deletion-cleanup.ts
src/state/agent-deletion-discovery.ts
src/state/agent-deletion-journal-history.ts
src/state/agent-deletion-journal-recovery.kernel.ts
src/state/agent-deletion-journal-recovery.ts
src/state/agent-deletion-journal.read.ts
src/state/agent-deletion-journal.ts

View file

@ -0,0 +1 @@
export class DuplicateAgentError extends Error {}

View file

@ -3,10 +3,8 @@ import fs from "node:fs/promises";
import path from "node:path";
import { vi } from "vitest";
import { writeSessionEntry } from "../config/sessions/session-accessor.sqlite-entry-store.js";
import {
readAgentDeletionRecoveryHolds,
reconstructAgentDeletionJournal,
} from "../state/agent-deletion-journal-recovery.js";
import { reconstructAgentDeletionJournal } from "../state/agent-deletion-journal-recovery.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.kernel.js";
import {
closeOpenClawAgentDatabasesForTest,
runOpenClawAgentWriteTransaction,

View file

@ -12,7 +12,6 @@ import { migrateLegacyConfig } from "../commands/doctor/shared/legacy-config-mig
import { ensureOnboardingAgent } from "../commands/onboard-agent.js";
import {
mutateConfigFileWithRetry,
readConfigFileSnapshotForWrite,
transformConfigFileWithRetry,
withConfigMutationExclusive,
} from "../config/config.js";
@ -68,11 +67,6 @@ import { resolveSharedAuthStorePath } from "./auth-profiles/path-resolve.js";
import { resolveAuthProfileDatabasePath } from "./auth-profiles/sqlite.js";
import { ensureAuthProfileStore } from "./auth-profiles/store-runtime.js";
import { readWorkspaceStateSnapshot } from "./workspace-state-store.js";
import {
DEFAULT_IDENTITY_FILENAME,
ensureAgentWorkspace,
isWorkspaceBootstrapPending,
} from "./workspace.js";
it("restores only the configured held store after explicit creation, never through bootstrap or retargeting", async () => {
const state = await createOpenClawTestState({ scenario: "minimal", label: "held-agent-restore" });
@ -735,110 +729,6 @@ it.each(["commit", "rollback", "close-failure"] as const)(
},
);
it("preserves env references from guided staging when preparation changes the environment", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "minimal",
label: "guided-stage-env",
});
const oldToken = process.env.GUIDED_STAGE_TOKEN;
try {
process.env.GUIDED_STAGE_TOKEN = "synthetic-read-value";
const config = JSON.parse(await fs.readFile(state.configPath, "utf8")) as OpenClawConfig;
await state.writeConfig({
...config,
gateway: { ...config.gateway, auth: { mode: "token", token: "${GUIDED_STAGE_TOKEN}" } },
});
const writeSnapshot = await readConfigFileSnapshotForWrite();
const staged = writeSnapshot.snapshot.sourceConfig;
expect(staged.gateway?.auth?.token).toBe("synthetic-read-value");
await Promise.resolve();
process.env.GUIDED_STAGE_TOKEN = "synthetic-after-guided-await";
const created = await createAgent({
name: "guided",
workspace: state.path("guided-workspace"),
stagedConfig: { config: staged, writeSnapshot },
prepareConfigCommit: async () => {
await Promise.resolve();
process.env.GUIDED_STAGE_TOKEN = "synthetic-after-preparation";
},
});
expect(created).toMatchObject({ status: "created", agentId: "guided" });
const saved = JSON.parse(await fs.readFile(state.configPath, "utf8")) as OpenClawConfig;
expect(saved.gateway?.auth?.token).toBe("${GUIDED_STAGE_TOKEN}");
expect(saved.agents?.entries?.guided).toBeDefined();
} finally {
if (oldToken === undefined) {
delete process.env.GUIDED_STAGE_TOKEN;
} else {
process.env.GUIDED_STAGE_TOKEN = oldToken;
}
closeOpenClawStateDatabaseForTest();
await state.cleanup();
}
});
it("keeps a fresh named workspace pending through the first run setup", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "minimal",
label: "named-agent-hatch",
});
const workspace = state.path("named-workspace");
try {
const created = await createAgent({ name: "Researcher", workspace });
expect(created).toMatchObject({ status: "created", bootstrapPending: true });
expect(await isWorkspaceBootstrapPending(workspace)).toBe(true);
const firstRunWorkspace = await ensureAgentWorkspace({
dir: workspace,
ensureBootstrapFiles: true,
});
expect(firstRunWorkspace.bootstrapPending).toBe(true);
expect(await isWorkspaceBootstrapPending(workspace)).toBe(true);
expect(
await fs.readFile(path.join(workspace, DEFAULT_IDENTITY_FILENAME), "utf8"),
).not.toContain("Researcher");
} finally {
closeOpenClawStateDatabaseForTest();
await state.cleanup();
}
});
it("records operator and agent creation provenance after roster commits", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "empty",
label: "agent-creation-provenance",
});
try {
await createAgent({ name: "Operator Child", workspace: state.path("operator-child") });
await createAgent({
name: "Agent Child",
workspace: state.path("agent-child"),
provenance: { createdVia: "agent", creatorAgentId: "main" },
});
expect(readAgentProvenance("operator-child", { env: state.env })).toMatchObject({
agentId: "operator-child",
createdVia: "operator",
creatorAgentId: null,
createdAtMs: expect.any(Number),
});
expect(readAgentProvenance("agent-child", { env: state.env })).toMatchObject({
agentId: "agent-child",
createdVia: "agent",
creatorAgentId: "main",
createdAtMs: expect.any(Number),
});
} finally {
closeOpenClawStateDatabaseForTest();
await state.cleanup();
}
});
describe("agent roster persistence", () => {
async function addWorkerToConfig(config: unknown): Promise<OpenClawConfig> {
const state = await createOpenClawTestState({

View file

@ -23,10 +23,11 @@ import type { OpenClawConfig } from "../config/types.openclaw.js";
import { FsSafeError, root } from "../infra/fs-safe.js";
import { normalizeAgentId, normalizeAgentIdStrict } from "../routing/session-key.js";
import { runWithAgentCreationClaim } from "../state/agent-creation-claim.js";
import { resolveAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.js";
import {
assertAgentDeletionRecoveryHoldPredicate,
readAgentDeletionRecoveryHolds,
resolveAgentDeletionRecoveryHolds,
} from "../state/agent-deletion-journal-recovery.js";
} from "../state/agent-deletion-journal-recovery.kernel.js";
import { readAgentDeletionJournal } from "../state/agent-deletion-journal.js";
import type { HeldAgentDatabase } from "../state/agent-deletion-journal.types.js";
import { recordAgentProvenance, type AgentCreatedVia } from "../state/agent-provenance.js";
@ -35,6 +36,7 @@ import { withExistingOpenClawStateDatabaseCurrentReadOnly } from "../state/openc
import { runOpenClawStateWriteTransaction } from "../state/openclaw-state-db.js";
import { isReservedSystemAgentId } from "../system-agent/agent-id.js";
import { resolveUserPath } from "../utils.js";
import { DuplicateAgentError } from "./agent-create-error.js";
import { normalizeAgentDirRegistryPath } from "./agent-dir-registry.js";
import { claimCompletedAgentDeletion } from "./agent-lifecycle-registry.js";
import { listAgentRoles, loadAgentRole } from "./agent-roles.js";
@ -46,6 +48,8 @@ import {
mergeIdentityMarkdownContent,
sanitizeAgentIdentityLine,
} from "./identity-file.js";
import { createWorkspaceFileMutationGuard } from "./workspace-file-mutation-guard.js";
import type { WorkspaceStateGuard } from "./workspace-state-store.worker-contract.js";
import {
DEFAULT_IDENTITY_FILENAME,
ensureAgentWorkspace,
@ -125,7 +129,6 @@ type CreateAgentParams = {
provenance?: { createdVia: AgentCreatedVia; creatorAgentId?: string };
};
class DuplicateAgentError extends Error {}
class InvalidAgentBindingsError extends Error {}
class UnfinishedRoleBootstrapError extends Error {}
@ -273,8 +276,9 @@ export async function checkAgentCreationGate(agentId: string): Promise<CreateErr
async function writeIdentityFile(params: {
workspaceDir: string;
identity: NonNullable<ReturnType<typeof createAgentIdentityConfig>>;
beforePersistentApply?: () => void;
guard?: WorkspaceStateGuard;
}): Promise<void> {
const beforeFileMutation = createWorkspaceFileMutationGuard(params.guard);
const workspaceRoot = await root(params.workspaceDir);
let existing: string | undefined;
try {
@ -288,11 +292,11 @@ async function writeIdentityFile(params: {
}
}
const content = mergeIdentityMarkdownContent(existing, params.identity);
params.beforePersistentApply?.();
beforeFileMutation?.();
// Root.write rechecks after its own async preparation and before each mutation.
await workspaceRoot.write(DEFAULT_IDENTITY_FILENAME, content, {
encoding: "utf8",
assertBeforeMutation: params.beforePersistentApply,
assertBeforeMutation: beforeFileMutation,
});
}
@ -343,34 +347,39 @@ export async function createAgent(params: CreateAgentParams): Promise<CreateAgen
let identityPublished = false;
let held: HeldAgentDatabase[] = [];
const readCurrentHolds = () =>
withExistingOpenClawStateDatabaseCurrentReadOnly(readAgentDeletionRecoveryHolds) ?? [];
withExistingOpenClawStateDatabaseCurrentReadOnly(readAgentDeletionRecoveryHolds, {
allowNativeRead: true,
}) ?? [];
const recoveryPathMatcher = createOpenClawAgentDatabasePathMatcher();
const assertRecoveryCurrent = () => {
const recoveryHoldPredicate = () => ({ agentId, held, applies: creating || !automaticBootstrap });
const assertRecoveryPathCurrent = () => {
if (!recoveryPathMatcher.isCurrent()) {
throw new DuplicateAgentError(
`Agent ${agentId} preserved database changed during restoration; its hold remains. Restore the original database and retry agents add.`,
);
}
if (
(creating || !automaticBootstrap) &&
readCurrentHolds().some(
(entry) =>
entry.agentId === agentId &&
!held.some((previous) => previous.agentId === agentId && previous.path === entry.path),
)
) {
throw new DuplicateAgentError(
`Agent ${agentId} has held databases. Restore its original agentDir and session.store configuration, then run agents add explicitly to restore the preserved store.`,
};
const assertRecoveryCurrent = () => {
assertRecoveryPathCurrent();
const predicate = recoveryHoldPredicate();
if (predicate.applies) {
withExistingOpenClawStateDatabaseCurrentReadOnly(
(database) => assertAgentDeletionRecoveryHoldPredicate(database, predicate),
{ allowNativeRead: true },
);
}
};
const hasBootstrapHold = () =>
automaticBootstrap && readCurrentHolds().some((entry) => entry.agentId === agentId);
const beforePersistentApply = () => {
const assertHost = () => {
params.beforePersistentApply?.();
if (!identityPublished) {
params.assertIdentityInputAllowed?.();
}
assertRecoveryPathCurrent();
};
const beforePersistentApply = () => {
assertHost();
assertRecoveryCurrent();
};
@ -573,7 +582,7 @@ export async function createAgent(params: CreateAgentParams): Promise<CreateAgen
beforePersistentApply();
const workspace = await ensureAgentWorkspace({
dir: workspaceDir,
beforePersistentApply,
guard: { assertHost, recoveryHoldPredicate: recoveryHoldPredicate() },
ensureBootstrapFiles: !skipBootstrap,
purpose: params.purpose,
...(template
@ -619,7 +628,7 @@ export async function createAgent(params: CreateAgentParams): Promise<CreateAgen
await writeIdentityFile({
workspaceDir: workspace.dir,
identity,
beforePersistentApply,
guard: { assertHost, recoveryHoldPredicate: recoveryHoldPredicate() },
});
// Publish the config projection of these accepted bytes even if new
// uploads close; delegated and recovery authority remain live above.

View file

@ -0,0 +1,237 @@
import fs from "node:fs/promises";
import path from "node:path";
import { expect, it, vi } from "vitest";
import { readConfigFileSnapshotForWrite } from "../config/config.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import * as fsSafe from "../infra/fs-safe.js";
import * as snapshots from "../infra/sqlite-readonly-worker.js";
import * as workerAdmission from "../infra/sqlite-worker-operation-admission.js";
import { reconstructAgentDeletionJournal } from "../state/agent-deletion-journal-recovery.js";
import { readAgentProvenance } from "../state/agent-provenance.js";
import { runOpenClawStateWriteTransaction } from "../state/openclaw-state-db.js";
import { observeMainThreadSql } from "../test-utils/main-thread-sql-spies.test-support.js";
import { createOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import { createAgent } from "./agent-create.js";
import * as workspaceModule from "./workspace.js";
import {
DEFAULT_IDENTITY_FILENAME,
ensureAgentWorkspace,
isWorkspaceBootstrapPending,
} from "./workspace.js";
function addRecoveryHold(agentId: string, heldPath: string) {
runOpenClawStateWriteTransaction((database) => {
database.db.exec("DROP TABLE agent_deletion_journal");
reconstructAgentDeletionJournal(database, [{ agentId, path: heldPath }]);
});
}
it("preserves IDENTITY.md when a recovery hold arrives during its awaited read", async () => {
const state = await createOpenClawTestState({ scenario: "minimal" });
const identityPath = path.join(state.workspaceDir, DEFAULT_IDENTITY_FILENAME);
const original = "# Identity\n- **Name:** Kept\n";
let inserted = false;
try {
await fs.writeFile(identityPath, original);
await ensureAgentWorkspace({ dir: state.workspaceDir, ensureBootstrapFiles: true });
const configBefore = await fs.readFile(state.configPath, "utf8");
const root = fsSafe.root;
vi.spyOn(fsSafe, "root").mockImplementation(async (...args) => {
const handle = await root(...args);
const read = handle.read.bind(handle);
vi.spyOn(handle, "read").mockImplementation(async (...readArgs) => {
const result = await read(...readArgs);
if (
args[0] === state.workspaceDir &&
readArgs[0] === DEFAULT_IDENTITY_FILENAME &&
!inserted
) {
inserted = true;
addRecoveryHold("guarded", state.path("held.sqlite"));
}
return result;
});
return handle;
});
const result = await createAgent({
entry: { id: "guarded", identity: { name: "Replacement" } },
workspace: state.workspaceDir,
});
expect(inserted).toBe(true);
expect(result).toMatchObject({
status: "error",
reason: "already-exists",
message: expect.stringContaining("held databases"),
});
expect(await fs.readFile(identityPath, "utf8")).toBe(original);
expect(await fs.readFile(state.configPath, "utf8")).toBe(configBefore);
} finally {
vi.restoreAllMocks();
await state.cleanup();
}
});
it("refuses bootstrap publication when a recovery hold arrives during root preparation", async () => {
const state = await createOpenClawTestState({ scenario: "minimal" });
let inserted = false;
try {
await fs.writeFile(path.join(state.workspaceDir, "AGENTS.md"), "Synthetic instructions.\n");
const root = fsSafe.root;
vi.spyOn(fsSafe, "root").mockImplementation(async (...args) => {
const handle = await root(...args);
if (args[0] === state.workspaceDir && !inserted) {
inserted = true;
addRecoveryHold("guarded", state.path("held.sqlite"));
}
return handle;
});
await expect(
ensureAgentWorkspace({
dir: state.workspaceDir,
ensureBootstrapFiles: true,
guard: { recoveryHoldPredicate: { agentId: "guarded", held: [], applies: true } },
}),
).rejects.toThrow("held databases");
expect(inserted).toBe(true);
expect(await fs.readdir(state.workspaceDir)).toEqual(["AGENTS.md"]);
} finally {
vi.restoreAllMocks();
await state.cleanup();
}
});
it("records operator and agent creation provenance after roster commits", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "empty",
label: "agent-creation-provenance",
});
const admission = workerAdmission.createSqliteWorkerOperationAdmission;
const ensureWorkspace = ensureAgentWorkspace;
const preparation = vi
.spyOn(workspaceModule, "ensureAgentWorkspace")
.mockImplementation(async (params) => {
const snapshot = vi.spyOn(snapshots, "runSqliteReadOnlyWorkerSync").mockImplementation(() => {
throw new Error("Workspace creation must not spawn synchronous SQLite snapshots");
});
try {
const result = await ensureWorkspace(params);
expect(snapshot).not.toHaveBeenCalled();
return result;
} finally {
snapshot.mockRestore();
}
});
let grants = 0;
const spy = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
admission((request, grant) => {
const sql = observeMainThreadSql();
try {
admit(request, grant);
sql.expectIdle();
grants++;
} finally {
sql.restore();
}
}, attachment),
);
try {
await createAgent({ name: "Operator Child", workspace: state.path("operator-child") });
await createAgent({
name: "Agent Child",
workspace: state.path("agent-child"),
provenance: { createdVia: "agent", creatorAgentId: "main" },
});
expect(preparation).toHaveBeenCalled();
expect(grants).toBeGreaterThan(0);
expect(readAgentProvenance("operator-child", { env: state.env })).toMatchObject({
agentId: "operator-child",
createdVia: "operator",
creatorAgentId: null,
createdAtMs: expect.any(Number),
});
expect(readAgentProvenance("agent-child", { env: state.env })).toMatchObject({
agentId: "agent-child",
createdVia: "agent",
creatorAgentId: "main",
createdAtMs: expect.any(Number),
});
} finally {
spy.mockRestore();
preparation.mockRestore();
await state.cleanup();
}
});
it("preserves env references from guided staging when preparation changes the environment", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "minimal",
label: "guided-stage-env",
});
const oldToken = process.env.GUIDED_STAGE_TOKEN;
try {
process.env.GUIDED_STAGE_TOKEN = "synthetic-read-value";
const config = JSON.parse(await fs.readFile(state.configPath, "utf8")) as OpenClawConfig;
await state.writeConfig({
...config,
gateway: { ...config.gateway, auth: { mode: "token", token: "${GUIDED_STAGE_TOKEN}" } },
});
const writeSnapshot = await readConfigFileSnapshotForWrite();
const staged = writeSnapshot.snapshot.sourceConfig;
expect(staged.gateway?.auth?.token).toBe("synthetic-read-value");
await Promise.resolve();
process.env.GUIDED_STAGE_TOKEN = "synthetic-after-guided-await";
const created = await createAgent({
name: "guided",
workspace: state.path("guided-workspace"),
stagedConfig: { config: staged, writeSnapshot },
prepareConfigCommit: async () => {
await Promise.resolve();
process.env.GUIDED_STAGE_TOKEN = "synthetic-after-preparation";
},
});
expect(created).toMatchObject({ status: "created", agentId: "guided" });
const saved = JSON.parse(await fs.readFile(state.configPath, "utf8")) as OpenClawConfig;
expect(saved.gateway?.auth?.token).toBe("${GUIDED_STAGE_TOKEN}");
expect(saved.agents?.entries?.guided).toBeDefined();
} finally {
if (oldToken === undefined) {
delete process.env.GUIDED_STAGE_TOKEN;
} else {
process.env.GUIDED_STAGE_TOKEN = oldToken;
}
await state.cleanup();
}
});
it("keeps a fresh named workspace pending through the first run setup", async () => {
const state = await createOpenClawTestState({
layout: "state-only",
scenario: "minimal",
label: "named-agent-hatch",
});
const workspace = state.path("named-workspace");
try {
const created = await createAgent({ name: "Researcher", workspace });
expect(created).toMatchObject({ status: "created", bootstrapPending: true });
expect(await isWorkspaceBootstrapPending(workspace)).toBe(true);
const firstRunWorkspace = await ensureAgentWorkspace({
dir: workspace,
ensureBootstrapFiles: true,
});
expect(firstRunWorkspace.bootstrapPending).toBe(true);
expect(await isWorkspaceBootstrapPending(workspace)).toBe(true);
expect(
await fs.readFile(path.join(workspace, DEFAULT_IDENTITY_FILENAME), "utf8"),
).not.toContain("Researcher");
} finally {
await state.cleanup();
}
});

View file

@ -7,10 +7,8 @@ import { observeHostDataSql } from "../../test/helpers/sqlite-statement-executio
import { clearCronJobActive, markCronJobActive } from "../cron/active-jobs.js";
import { registerActiveCronTaskRun } from "../cron/service/active-run-cancellation.js";
import { createDeferredCore } from "../shared/deferred.js";
import {
readAgentDeletionRecoveryHolds,
reconstructAgentDeletionJournal,
} from "../state/agent-deletion-journal-recovery.js";
import { reconstructAgentDeletionJournal } from "../state/agent-deletion-journal-recovery.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.kernel.js";
import {
beginAgentDeletionJournal,
claimCompletedAgentDeletionJournal,

View file

@ -0,0 +1,23 @@
import { assertAgentDeletionRecoveryHoldPredicate } from "../state/agent-deletion-journal-recovery.kernel.js";
import { withExistingOpenClawStateDatabaseCurrentReadOnly } from "../state/openclaw-state-db-readonly.js";
import type { WorkspaceStateGuard } from "./workspace-state-store.worker-contract.js";
/** Filesystem effects run outside worker grants and retain their current native recovery check. */
export function createWorkspaceFileMutationGuard(
guard?: WorkspaceStateGuard,
): (() => void) | undefined {
if (!guard) {
return undefined;
}
return () => {
guard.assertHost?.();
guard.beforeLegacyApply?.();
const predicate = guard.recoveryHoldPredicate;
if (predicate?.applies) {
withExistingOpenClawStateDatabaseCurrentReadOnly(
(database) => assertAgentDeletionRecoveryHoldPredicate(database, predicate),
{ allowNativeRead: true },
);
}
};
}

View file

@ -0,0 +1,49 @@
import fs from "node:fs/promises";
import path from "node:path";
import { runCommandWithTimeout } from "../process/exec.js";
import { createLazyPromise, getOrCreatePromise } from "../shared/lazy-promise.js";
const gitInitializationInFlight = new Map<string, Promise<void>>();
// Git availability is process-stable; cache the probe result, including failure, until restart.
const isGitAvailable = createLazyPromise(async () => {
try {
const result = await runCommandWithTimeout(["git", "--version"], { timeoutMs: 2_000 });
return result.code === 0;
} catch {
return false;
}
});
export async function ensureGitRepo(
dir: string,
isBrandNewWorkspace: boolean,
beforePersistentApply?: () => void,
) {
if (!isBrandNewWorkspace) {
return;
}
// Concurrent first turns can all observe missing Git metadata. Join only the
// current initialization; later calls must inspect the workspace again.
beforePersistentApply?.();
await getOrCreatePromise(
gitInitializationInFlight,
dir,
async () => {
if (await fs.stat(path.join(dir, ".git")).catch(() => undefined)) {
return;
}
if (!(await isGitAvailable())) {
return;
}
// Only the initializer's owner admits Git; joining callers cannot cancel it.
beforePersistentApply?.();
try {
await runCommandWithTimeout(["git", "init"], { cwd: dir, timeoutMs: 10_000 });
} catch {
// Ignore git init failures; workspace creation should still succeed.
}
},
{ evictOnSettled: true },
);
}

View file

@ -0,0 +1,94 @@
import fs from "node:fs";
import path from "node:path";
import { isDeepStrictEqual } from "node:util";
import { hasErrnoCode } from "../infra/errno.js";
import { resolveUserPath } from "../utils.js";
import { WORKSPACE_BOOTSTRAP_FILENAMES } from "./workspace-bootstrap-policy.js";
import {
resolveCanonicalWorkspacePath,
WorkspaceAliasRepointedError,
} from "./workspace-state-identity.js";
function statFact(file: string, follow: boolean, content: boolean): string | undefined {
const stat = follow
? fs.statSync(file, { bigint: true, throwIfNoEntry: false })
: fs.lstatSync(file, { bigint: true, throwIfNoEntry: false });
if (!stat) {
return undefined;
}
return content
? `${stat.dev}:${stat.ino}:${stat.mode}:${stat.size}:${stat.mtimeNs}:${stat.ctimeNs}`
: `${stat.dev}:${stat.ino}:${stat.mode}`;
}
/** Request-local filesystem evidence; stored NFC keys never select filesystem paths. */
export function captureWorkspaceStateFilesystemGuard(
workspaceDir: string,
content = true,
): () => void {
const dir = path.resolve(resolveUserPath(workspaceDir));
const canonical = resolveCanonicalWorkspacePath(dir);
const root = [statFact(dir, false, false), statFact(dir, true, false)];
const entries = (directory: string) => {
try {
return fs.readdirSync(directory, { withFileTypes: true });
} catch (error) {
if (!hasErrnoCode(error, "ENOENT") && !hasErrnoCode(error, "ENOTDIR")) {
throw error;
}
return [];
}
};
const observe = () => {
const paths: string[] = [];
if (content) {
paths.push(
...[...WORKSPACE_BOOTSTRAP_FILENAMES, "memory", ".git", "skills"].map((name) =>
path.join(dir, name),
),
);
for (const entry of entries(path.join(dir, "skills"))) {
if (entry.isDirectory()) {
paths.push(
path.join(dir, "skills", entry.name),
path.join(dir, "skills", entry.name, "SKILL.md"),
);
}
}
}
return {
names: content
? entries(dir)
.map((entry) => entry.name)
.toSorted()
: [],
files: paths
.toSorted()
.map((file) => [file, statFact(file, false, content), statFact(file, true, content)]),
};
};
const initial = observe();
return () => {
const current = resolveCanonicalWorkspacePath(dir);
if (current !== canonical) {
throw new WorkspaceAliasRepointedError({
aliasPath: dir,
storedWorkspacePath: canonical,
currentWorkspacePath: current,
});
}
// Concurrent first-time provisioning can create the same empty directory.
// Existing roots remain pinned; new content still invalidates the observation.
const currentRoot = [statFact(dir, false, false), statFact(dir, true, false)];
if (
(root.some((fact) => fact !== undefined) && !isDeepStrictEqual(root, currentRoot)) ||
(!root.some((fact) => fact !== undefined) &&
fs.statSync(dir, { throwIfNoEntry: false })?.isFile()) ||
!isDeepStrictEqual(observe(), initial)
) {
throw new Error(
`Workspace filesystem changed while preparing state for ${dir}; retry workspace setup.`,
);
}
};
}

View file

@ -1,6 +1,8 @@
import fs from "node:fs";
import path from "node:path";
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import * as workerAdmission from "../infra/sqlite-worker-operation-admission.js";
import { reconstructAgentDeletionJournal } from "../state/agent-deletion-journal-recovery.js";
import {
withArtifactPreservingStateReads,
withOpenClawStateDatabaseReadSnapshot,
@ -8,20 +10,26 @@ import {
import {
closeOpenClawStateDatabaseAsync,
openOpenClawStateDatabase,
runOpenClawStateWriteTransaction,
} from "../state/openclaw-state-db.js";
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
import * as stateWorker from "../state/openclaw-state-worker-store.js";
import { observeMainThreadSql } from "../test-utils/main-thread-sql-spies.test-support.js";
import {
createOpenClawTestState,
type OpenClawTestState,
} from "../test-utils/openclaw-test-state.js";
import { DuplicateAgentError } from "./agent-create-error.js";
import { resolveBootstrapFilesForPreparation } from "./bootstrap-files.js";
import { readWorkspaceFileCache, writeWorkspaceFileCache } from "./workspace-file-cache.js";
import { assertConfiguredWorkspaceStateReady } from "./workspace-state-dirs.js";
import { WorkspaceAliasRepointedError } from "./workspace-state-identity.js";
import {
clearExpiredWorkspaceStateForVanishedWorkspace,
mergeWorkspaceSetupState,
readWorkspaceStateSnapshot,
replaceWorkspaceAttestation,
WORKSPACE_ATTESTATION_RECENT_MS,
} from "./workspace-state-store.js";
let state: OpenClawTestState;
@ -62,6 +70,242 @@ async function withoutMainThreadSql<T>(read: () => Promise<T>): Promise<T> {
}
}
it("snapshots, registers aliases, merges setup, and expires exact state without caller-thread SQL", async () => {
await withoutMainThreadSql(seed);
await closeOpenClawStateDatabaseAsync();
const alias = state.path("runtime-alias");
fs.symlinkSync(state.workspaceDir, alias, process.platform === "win32" ? "junction" : "dir");
await withoutMainThreadSql(async () => {
expect((await readWorkspaceStateSnapshot(alias)).setup.bootstrapSeededAt).toBe(
"2026-07-16T01:00:00.000Z",
);
expect(
await mergeWorkspaceSetupState(
alias,
{ bootstrapSeededAt: "2026-07-17T01:00:00.000Z" },
2_000,
),
).toEqual({ version: 1, bootstrapSeededAt: "2026-07-16T01:00:00.000Z" });
expect(await clearExpiredWorkspaceStateForVanishedWorkspace(alias, 2_000)).toBe(false);
fs.unlinkSync(alias);
expect(
await clearExpiredWorkspaceStateForVanishedWorkspace(
alias,
WORKSPACE_ATTESTATION_RECENT_MS + 2_001,
),
).toBe(true);
expect((await readWorkspaceStateSnapshot(state.workspaceDir)).setupExists).toBe(false);
});
});
it("rolls back setup when a recovery hold arrives after caller preparation", async () => {
const before = await seed();
const recoveryHoldPredicate = { agentId: "new", held: [], applies: true };
runOpenClawStateWriteTransaction((database) => {
database.db.exec("DROP TABLE agent_deletion_journal");
reconstructAgentDeletionJournal(database, [
{ agentId: "new", path: state.path("held.sqlite") },
]);
});
await withoutMainThreadSql(async () => {
const refusal = await mergeWorkspaceSetupState(
state.workspaceDir,
{ setupCompletedAt: "2026-07-16T02:00:00.000Z" },
2_000,
{ recoveryHoldPredicate },
).catch((error: unknown) => error);
expect(refusal).toBeInstanceOf(DuplicateAgentError);
expect(refusal).toMatchObject({
message:
"Agent new has held databases. Restore its original agentDir and session.store configuration, then run agents add explicitly to restore the preserved store.",
});
expect(await readWorkspaceStateSnapshot(state.workspaceDir)).toEqual(before);
});
});
it.each(["snapshot", "merge", "expire"] as const)(
"rejects revoked authority at transaction and commit for %s",
async (operation) => {
const before = await seed();
const alias = state.path("grant-alias");
fs.symlinkSync(state.workspaceDir, alias, process.platform === "win32" ? "junction" : "dir");
const originalAdmission = workerAdmission.createSqliteWorkerOperationAdmission;
for (const stage of ["transaction", "commit"] as const) {
let retired = false;
const refusal = new Error("workspace owner retired");
const spy = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
originalAdmission((request, grant) => {
retired ||= request.stage === stage;
admit(request, grant);
}, attachment),
);
const options = {
assertCurrent: () => {
if (retired) {
throw refusal;
}
},
};
const pending =
operation === "snapshot"
? readWorkspaceStateSnapshot(alias, options)
: operation === "merge"
? mergeWorkspaceSetupState(
alias,
{ setupCompletedAt: "2026-07-16T02:00:00.000Z" },
2_000,
options,
)
: clearExpiredWorkspaceStateForVanishedWorkspace(
alias,
WORKSPACE_ATTESTATION_RECENT_MS + 2_001,
options,
);
await expect(pending).rejects.toBe(refusal);
spy.mockRestore();
expect(retired).toBe(true);
expect(await readWorkspaceStateSnapshot(state.workspaceDir)).toEqual(before);
}
fs.unlinkSync(alias);
expect((await readWorkspaceStateSnapshot(alias)).setupExists).toBe(false);
},
);
it.each(["transaction", "commit"] as const)(
"refuses expiry when workspace files reappear at %s admission",
async (stage) => {
const before = await seed();
const filePath = path.join(state.workspaceDir, "AGENTS.md");
writeWorkspaceFileCache({ filePath, content: "cached", identity: "identity" });
const originalAdmission = workerAdmission.createSqliteWorkerOperationAdmission;
let restored = false;
const spy = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
originalAdmission((request, grant) => {
if (!restored && request.stage === stage) {
fs.writeFileSync(filePath, "Restored workspace instructions.");
restored = true;
}
admit(request, grant);
}, attachment),
);
await expect(
clearExpiredWorkspaceStateForVanishedWorkspace(
state.workspaceDir,
WORKSPACE_ATTESTATION_RECENT_MS + 2_001,
),
).rejects.toThrow("Workspace filesystem changed");
spy.mockRestore();
expect(restored).toBe(true);
expect(await readWorkspaceStateSnapshot(state.workspaceDir)).toEqual(before);
expect(readWorkspaceFileCache(filePath, "identity")).toBe("cached");
},
);
it("retires committed expiry cache entries when result delivery fails without replay", async () => {
await seed();
const filePath = path.join(state.workspaceDir, "AGENTS.md");
writeWorkspaceFileCache({ filePath, content: "cached", identity: "identity" });
const execute = stateWorker.runOpenClawStateWorkerOperation;
const spy = vi
.spyOn(stateWorker, "runOpenClawStateWorkerOperation")
.mockImplementationOnce(async (...args) => {
await execute(...args);
throw new Error("expiry reply lost");
});
await expect(
clearExpiredWorkspaceStateForVanishedWorkspace(
state.workspaceDir,
WORKSPACE_ATTESTATION_RECENT_MS + 2_001,
),
).rejects.toThrow("expiry reply lost");
expect(spy).toHaveBeenCalledTimes(1);
spy.mockRestore();
expect(readWorkspaceFileCache(filePath, "identity")).toBeUndefined();
expect((await readWorkspaceStateSnapshot(state.workspaceDir)).setupExists).toBe(false);
});
it("keeps captured first-writer setup milestones across queued merges and reopen", async () => {
await seed();
const next = { setupCompletedAt: "2026-07-16T02:00:00.000Z" };
const first = mergeWorkspaceSetupState(state.workspaceDir, next, 2_000);
next.setupCompletedAt = "2026-07-17T02:00:00.000Z";
const second = mergeWorkspaceSetupState(state.workspaceDir, next, 3_000);
expect((await first).setupCompletedAt).toBe("2026-07-16T02:00:00.000Z");
expect((await second).setupCompletedAt).toBe("2026-07-16T02:00:00.000Z");
await closeOpenClawStateDatabaseAsync();
expect((await readWorkspaceStateSnapshot(state.workspaceDir)).setup.setupCompletedAt).toBe(
"2026-07-16T02:00:00.000Z",
);
});
it("rolls back alias registration when its symlink is repointed before commit", async () => {
const before = await seed();
const alias = state.path("pending-alias");
const replacement = state.path("new-target");
fs.mkdirSync(replacement);
fs.symlinkSync(state.workspaceDir, alias, process.platform === "win32" ? "junction" : "dir");
const originalAdmission = workerAdmission.createSqliteWorkerOperationAdmission;
const spy = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
originalAdmission((request, grant) => {
if (request.stage === "commit") {
fs.unlinkSync(alias);
fs.symlinkSync(replacement, alias, process.platform === "win32" ? "junction" : "dir");
}
admit(request, grant);
}, attachment),
);
await expect(readWorkspaceStateSnapshot(alias)).rejects.toBeInstanceOf(
WorkspaceAliasRepointedError,
);
spy.mockRestore();
fs.unlinkSync(alias);
expect((await readWorkspaceStateSnapshot(alias)).setupExists).toBe(false);
expect(await readWorkspaceStateSnapshot(state.workspaceDir)).toEqual(before);
});
it("joins a granted workspace mutation before closing its database", async () => {
await seed();
let close: Promise<void> | undefined;
let closed = false;
const originalAdmission = workerAdmission.createSqliteWorkerOperationAdmission;
const spy = vi
.spyOn(workerAdmission, "createSqliteWorkerOperationAdmission")
.mockImplementation((admit, attachment) =>
originalAdmission((request, grant) => {
admit(request, () => {
const granted = grant();
if (request.stage === "commit") {
close = closeOpenClawStateDatabaseAsync().then(() => {
closed = true;
});
expect(closed).toBe(false);
}
return granted;
});
}, attachment),
);
await expect(
mergeWorkspaceSetupState(
state.workspaceDir,
{ setupCompletedAt: "2026-07-16T02:00:00.000Z" },
2_000,
),
).rejects.toThrow("read admission is closed");
expect(close).toBeDefined();
await close;
spy.mockRestore();
expect(closed).toBe(true);
expect((await readWorkspaceStateSnapshot(state.workspaceDir)).setup.setupCompletedAt).toBe(
"2026-07-16T02:00:00.000Z",
);
});
it.each([false, true])(
"prepares bootstrap files for completed=%s without main-thread SQL",
async (completed) => {

View file

@ -1,5 +1,6 @@
import fs from "node:fs";
import path from "node:path";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import {
executeSqliteQuerySync,
executeSqliteQueryTakeFirstSync,
@ -410,3 +411,95 @@ export function replaceWorkspaceAttestationInDatabase(
generatedHashes: new Map(sortedHashes),
};
}
export function deleteWorkspaceStateRowsInDatabase(
database: WorkspaceStateDatabaseHandle,
{ workspaceKey }: WorkspaceStateIdentity,
): void {
const kysely = getNodeSqliteKysely<WorkspaceStateDatabase>(database.db);
const receiptRows = executeSqliteQuerySync(
database.db,
kysely
.selectFrom("migration_sources")
.select(["source_key", "last_run_id", "report_json"])
.where("migration_kind", "in", [
WORKSPACE_LEGACY_STATE_MIGRATION_KIND,
WORKSPACE_CONTENT_RELOCATION_MIGRATION_KIND,
]),
).rows.filter((row) => {
try {
const report: unknown = JSON.parse(row.report_json);
return isRecord(report) && report.workspaceKey === workspaceKey;
} catch {
return false;
}
});
if (receiptRows.length > 0) {
const receiptKeys = receiptRows.map((row) => row.source_key);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("migration_sources").where("source_key", "in", receiptKeys),
);
const runIds = [...new Set(receiptRows.map((row) => row.last_run_id))];
const referencedRunIds = new Set(
executeSqliteQuerySync(
database.db,
kysely
.selectFrom("migration_sources")
.select("last_run_id")
.where("last_run_id", "in", runIds),
).rows.map((row) => row.last_run_id),
);
const orphanedRunIds = runIds.filter((runId) => !referencedRunIds.has(runId));
if (orphanedRunIds.length > 0) {
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("migration_runs").where("id", "in", orphanedRunIds),
);
}
}
executeSqliteQuerySync(
database.db,
kysely
.deleteFrom("workspace_generated_bootstrap_hashes")
.where("workspace_key", "=", workspaceKey),
);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("workspace_setup_state").where("workspace_key", "=", workspaceKey),
);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("workspace_path_aliases").where("workspace_key", "=", workspaceKey),
);
}
export function recentWorkspaceAttestation(
attestation: WorkspaceAttestation | undefined,
nowMs = Date.now(),
): WorkspaceAttestation | undefined {
if (!attestation) {
return undefined;
}
const ageMs = nowMs - attestation.attestedAtMs;
// Clock rollback must not turn disappearance protection into permission to
// reseed. A healthy workspace refreshes the future-dated row below.
if (ageMs > WORKSPACE_ATTESTATION_RECENT_MS) {
return undefined;
}
return attestation;
}
export function hasWorkspaceSetupStateMarker(state: WorkspaceSetupState): boolean {
return Boolean(state.bootstrapSeededAt || state.setupCompletedAt);
}
export function hasRecentWorkspaceSetupState(
snapshot: WorkspaceStateSnapshot,
nowMs = Date.now(),
): boolean {
if (!hasWorkspaceSetupStateMarker(snapshot.setup) || snapshot.setupUpdatedAtMs === undefined) {
return false;
}
return nowMs - snapshot.setupUpdatedAtMs <= WORKSPACE_ATTESTATION_RECENT_MS;
}

View file

@ -156,7 +156,7 @@ describe("workspace state store", () => {
});
it.each(["merge", "expire", "delete", "register-alias"] as const)(
"checks current ownership inside the %s transaction before changing state",
"refuses retired ownership before changing %s state",
async (operation) => {
const dir = workspaceDir();
const alias = testState!.path("workspace-link");
@ -174,7 +174,9 @@ describe("workspace state store", () => {
writeWorkspaceFileCache({ filePath, content: "cached", identity: "identity" });
const retired = new Error("workspace owner retired");
const assertCurrent = () => {
expect(db.isTransaction).toBe(true);
if (operation === "delete") {
expect(db.isTransaction).toBe(true);
}
throw retired;
};
const operations = {

View file

@ -6,11 +6,9 @@ import {
getNodeSqliteKysely,
} from "../infra/kysely-sync.js";
import { deferSqlitePostCommitPublication } from "../infra/sqlite-post-commit.js";
import { runSqliteDeferredTransactionSync } from "../infra/sqlite-transaction.js";
import { createSqliteWorkerOperationAdmission } from "../infra/sqlite-worker-operation-admission.js";
import { executeExistingOpenClawStateRead } from "../state/openclaw-state-db-readonly.js";
import {
openOpenClawStateDatabase,
runOpenClawStateWriteTransaction,
type OpenClawStateDatabaseOptions,
} from "../state/openclaw-state-db.js";
@ -19,23 +17,20 @@ import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worke
import { runOpenClawStateWorkerOperation } from "../state/openclaw-state-worker-store.js";
import { resolveUserPath } from "../utils.js";
import { retireWorkspaceFileCache } from "./workspace-file-cache.js";
import { captureWorkspaceStateFilesystemGuard } from "./workspace-state-guard.js";
import {
createWorkspaceStateIdentity,
resolveCanonicalWorkspacePath,
resolveWorkspaceStateAliases,
resolveWorkspaceStateIdentity,
WorkspaceAliasRepointedError,
type WorkspaceStateIdentity,
} from "./workspace-state-identity.js";
import {
assertCanonicalIntegerTimestamp,
assertCanonicalTimestamp,
readWorkspaceStateSnapshotFromDatabase,
registerWorkspaceStateAliasIdentitiesInTransaction,
deleteWorkspaceStateRowsInDatabase,
resolveWorkspaceIdentityFromDatabase,
WORKSPACE_ATTESTATION_RECENT_MS,
WORKSPACE_CONTENT_RELOCATION_MIGRATION_KIND,
WORKSPACE_LEGACY_STATE_MIGRATION_KIND,
WORKSPACE_SETUP_STATE_VERSION,
workspacePathEntryExists,
type WorkspaceAttestation,
@ -45,8 +40,15 @@ import {
type WorkspaceStateDatabaseHandle,
type WorkspaceStateSnapshot,
} from "./workspace-state-store.kernel.js";
import type {
WorkspaceStateGuard,
WorkspaceStateWorkerOperations,
} from "./workspace-state-store.worker-contract.js";
export {
hasRecentWorkspaceSetupState,
hasWorkspaceSetupStateMarker,
recentWorkspaceAttestation,
isSafeWorkspaceAttestationFilename,
readWorkspaceStateSnapshotFromDatabase,
registerWorkspaceStateAliasIdentitiesInTransaction,
@ -60,7 +62,10 @@ export {
type WorkspaceStateSnapshot,
} from "./workspace-state-store.kernel.js";
type WorkspaceStateOperationOptions = { assertCurrent?: () => void };
type WorkspaceStateOperationOptions = { assertCurrent?: () => void } & Pick<
WorkspaceStateGuard,
"recoveryHoldPredicate" | "beforeLegacyApply"
>;
type WorkspaceStateDeletionPlan = {
cacheRoot: string;
@ -69,16 +74,90 @@ type WorkspaceStateDeletionPlan = {
pathEntryExisted: boolean;
};
async function runWorkspaceStateOperation<K extends keyof WorkspaceStateWorkerOperations>(
command: { type: K; input: WorkspaceStateWorkerOperations[K]["input"] },
options: OpenClawStateDatabaseOptions & WorkspaceStateOperationOptions,
): Promise<WorkspaceStateWorkerOperations[K]["output"]> {
const capturedCommand = {
...command,
input: {
...command.input,
recoveryHoldPredicate: structuredClone(options.recoveryHoldPredicate),
},
};
const context = captureOpenClawStateWorkerContext({
...options,
path: options.database?.path ?? options.path,
});
const assertFilesystem = captureWorkspaceStateFilesystemGuard(
command.input.workspaceDir,
command.type !== "workspace.snapshotAndRegister",
);
const assertCurrent = () => {
context.admission.assertCurrent();
options.assertCurrent?.();
assertFilesystem();
};
let expiryResult: string | false | undefined;
let publication: Promise<void> | undefined;
try {
const result = await runOpenClawStateWorkerOperation(
context,
async (scope) => {
try {
options.beforeLegacyApply?.();
return await scope.execute(capturedCommand);
} finally {
// Native settlement and cache retirement stay inside the writer's FIFO interval.
await publication;
}
},
{
assertCurrent,
createAdmission: (operation) => {
const admission = createSqliteWorkerOperationAdmission((request, grant) => {
if (request.stage !== "transaction" && request.stage !== "commit") {
throw new Error("Workspace state requires transaction admission");
}
assertCurrent();
if (command.type === "workspace.expire" && request.stage === "commit") {
if (typeof request.facts !== "string" && request.facts !== false) {
throw new Error("Workspace expiry has no admitted result");
}
expiryResult = request.facts;
}
grant();
});
const accepted = admission;
publication = operation.settled.then(() => {
if (typeof expiryResult === "string" && accepted.committed?.facts === expiryResult) {
retireWorkspaceFileCache(expiryResult);
}
});
return { admission, nativeLocations: [context.admission.databasePath] };
},
},
);
assertCurrent();
return result;
} finally {
await publication;
}
}
export async function readWorkspaceStateSnapshot(
workspaceDir: string,
options: OpenClawStateDatabaseOptions & WorkspaceStateOperationOptions = {},
): Promise<WorkspaceStateSnapshot> {
const capturedWorkspaceDir = path.resolve(resolveUserPath(workspaceDir));
if (options.readOnly) {
const capturedWorkspaceDir = path.resolve(resolveUserPath(workspaceDir));
const assertFilesystem = captureWorkspaceStateFilesystemGuard(capturedWorkspaceDir, false);
const reply = await executeExistingOpenClawStateRead(options, {
type: "workspace.snapshot",
workspaceDir: capturedWorkspaceDir,
});
options.assertCurrent?.();
assertFilesystem();
if (reply && (!reply.ok || reply.type !== "workspace.snapshot")) {
throw new Error("Unexpected workspace state snapshot result");
}
@ -90,56 +169,10 @@ export async function readWorkspaceStateSnapshot(
}
);
}
const database = openOpenClawStateDatabase(options);
const initial = runSqliteDeferredTransactionSync(database.db, () => {
const resolution = resolveWorkspaceIdentityFromDatabase({ workspaceDir, database });
return {
resolution,
snapshot: readWorkspaceStateSnapshotFromDatabase({ identity: resolution.identity, database }),
};
});
if (
initial.resolution.missingAliasKeys.length === 0 ||
(!initial.snapshot.setupExists && !initial.snapshot.attestation)
) {
return initial.snapshot;
}
// Register a newly observed configured spelling once state proves the target
// identity. Later disappearance must still find the same safety evidence.
return runOpenClawStateWriteTransaction((writeDatabase) => {
options.assertCurrent?.();
const currentAliases = resolveWorkspaceStateAliases(workspaceDir);
const currentCanonicalIdentity = currentAliases.at(-1)!;
if (
workspacePathEntryExists(workspaceDir) &&
currentCanonicalIdentity.workspaceKey !== initial.resolution.identity.workspaceKey
) {
throw new WorkspaceAliasRepointedError({
aliasPath: currentAliases[0]!.workspacePath,
storedWorkspacePath: initial.resolution.identity.workspacePath,
currentWorkspacePath: currentCanonicalIdentity.workspacePath,
});
}
const snapshot = readWorkspaceStateSnapshotFromDatabase({
identity: initial.resolution.identity,
database: writeDatabase,
});
if (snapshot.setupExists || snapshot.attestation) {
const aliases = new Map(
[...initial.resolution.aliases, ...currentAliases].map((alias) => [
alias.workspaceKey,
alias,
]),
);
registerWorkspaceStateAliasIdentitiesInTransaction({
database: writeDatabase,
identity: initial.resolution.identity,
aliases: [...aliases.values()],
updatedAtMs: Date.now(),
});
}
return snapshot;
}, options);
return runWorkspaceStateOperation(
{ type: "workspace.snapshotAndRegister", input: { workspaceDir: capturedWorkspaceDir } },
options,
);
}
export async function mergeWorkspaceSetupState(
@ -155,49 +188,17 @@ export async function mergeWorkspaceSetupState(
if (next.setupCompletedAt) {
assertCanonicalTimestamp(next.setupCompletedAt, "setup completed");
}
return runOpenClawStateWriteTransaction((database) => {
options.assertCurrent?.();
const resolution = resolveWorkspaceIdentityFromDatabase({ workspaceDir, database });
const identity = resolution.identity;
const snapshot = readWorkspaceStateSnapshotFromDatabase({ identity, database });
const bootstrapSeededAt = snapshot.setup.bootstrapSeededAt ?? next.bootstrapSeededAt;
const setupCompletedAt = snapshot.setup.setupCompletedAt ?? next.setupCompletedAt;
const merged: WorkspaceSetupState = {
version: WORKSPACE_SETUP_STATE_VERSION,
...(bootstrapSeededAt ? { bootstrapSeededAt } : {}),
...(setupCompletedAt ? { setupCompletedAt } : {}),
};
const kysely = getNodeSqliteKysely<WorkspaceStateDatabase>(database.db);
executeSqliteQuerySync(
database.db,
kysely
.insertInto("workspace_setup_state")
.values({
workspace_key: identity.workspaceKey,
workspace_path: identity.workspacePath,
version: WORKSPACE_SETUP_STATE_VERSION,
bootstrap_seeded_at: merged.bootstrapSeededAt ?? null,
setup_completed_at: merged.setupCompletedAt ?? null,
updated_at: nowMs,
})
.onConflict((conflict) =>
conflict.column("workspace_key").doUpdateSet({
workspace_path: identity.workspacePath,
version: WORKSPACE_SETUP_STATE_VERSION,
bootstrap_seeded_at: merged.bootstrapSeededAt ?? null,
setup_completed_at: merged.setupCompletedAt ?? null,
updated_at: nowMs,
}),
),
);
registerWorkspaceStateAliasIdentitiesInTransaction({
database,
identity,
aliases: resolution.aliases,
updatedAtMs: nowMs,
});
return merged;
}, options);
return runWorkspaceStateOperation(
{
type: "workspace.mergeSetup",
input: {
workspaceDir: path.resolve(resolveUserPath(workspaceDir)),
next: { ...next },
nowMs,
},
},
options,
);
}
export async function replaceWorkspaceAttestation(
@ -205,15 +206,19 @@ export async function replaceWorkspaceAttestation(
): Promise<WorkspaceAttestation> {
const context = captureOpenClawStateWorkerContext();
const { assertCurrent } = params;
const input: WorkspaceAttestationInput = {
const input = {
workspaceDir: path.resolve(resolveUserPath(params.workspaceDir)),
attestedAtMs: params.attestedAtMs,
generatedHashes: new Map(params.generatedHashes),
nowMs: params.nowMs,
recoveryHoldPredicate: structuredClone(params.recoveryHoldPredicate),
};
return runOpenClawStateWorkerOperation(
context,
(scope) => scope.execute({ type: "workspace.replaceAttestation", input }),
(scope) => {
params.beforeLegacyApply?.();
return scope.execute({ type: "workspace.replaceAttestation", input });
},
{
assertCurrent,
createAdmission: () => ({
@ -233,67 +238,13 @@ export async function replaceWorkspaceAttestation(
function deleteWorkspaceRows(
database: WorkspaceStateDatabaseHandle,
{ workspaceKey, workspacePath }: WorkspaceStateIdentity,
identity: WorkspaceStateIdentity,
): void {
const kysely = getNodeSqliteKysely<WorkspaceStateDatabase>(database.db);
const receiptRows = executeSqliteQuerySync(
database.db,
kysely
.selectFrom("migration_sources")
.select(["source_key", "last_run_id", "report_json"])
.where("migration_kind", "in", [
WORKSPACE_LEGACY_STATE_MIGRATION_KIND,
WORKSPACE_CONTENT_RELOCATION_MIGRATION_KIND,
]),
).rows.filter((row) => {
try {
const report = JSON.parse(row.report_json) as Record<string, unknown>;
return report.workspaceKey === workspaceKey;
} catch {
return false;
}
});
if (receiptRows.length > 0) {
const receiptKeys = receiptRows.map((row) => row.source_key);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("migration_sources").where("source_key", "in", receiptKeys),
);
const runIds = [...new Set(receiptRows.map((row) => row.last_run_id))];
const referencedRunIds = new Set(
executeSqliteQuerySync(
database.db,
kysely
.selectFrom("migration_sources")
.select("last_run_id")
.where("last_run_id", "in", runIds),
).rows.map((row) => row.last_run_id),
);
const orphanedRunIds = runIds.filter((runId) => !referencedRunIds.has(runId));
if (orphanedRunIds.length > 0) {
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("migration_runs").where("id", "in", orphanedRunIds),
);
}
}
executeSqliteQuerySync(
database.db,
kysely
.deleteFrom("workspace_generated_bootstrap_hashes")
.where("workspace_key", "=", workspaceKey),
deleteWorkspaceStateRowsInDatabase(database, identity);
// Explicit agent deletion retains its native transaction and publication owner.
deferSqlitePostCommitPublication(database.db, () =>
retireWorkspaceFileCache(identity.workspacePath),
);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("workspace_setup_state").where("workspace_key", "=", workspaceKey),
);
executeSqliteQuerySync(
database.db,
kysely.deleteFrom("workspace_path_aliases").where("workspace_key", "=", workspaceKey),
);
// Both deletion paths retire the actual stored identity only after the outer
// commit; a failed transaction must retain content for the surviving workspace.
deferSqlitePostCommitPublication(database.db, () => retireWorkspaceFileCache(workspacePath));
}
/** The migration owner has verified the same workspace and every relocated byte before this commit. */
@ -327,38 +278,17 @@ export async function clearExpiredWorkspaceStateForVanishedWorkspace(
options: WorkspaceStateOperationOptions = {},
): Promise<boolean> {
assertCanonicalIntegerTimestamp(nowMs, "workspace expiry check");
return runOpenClawStateWriteTransaction((database) => {
options.assertCurrent?.();
const resolution = resolveWorkspaceIdentityFromDatabase({ workspaceDir, database });
const identity = resolution.identity;
const snapshot = readWorkspaceStateSnapshotFromDatabase({ identity, database });
const preserveRecentState = () => {
registerWorkspaceStateAliasIdentitiesInTransaction({
database,
identity,
aliases: resolution.aliases,
updatedAtMs: nowMs,
});
return false;
};
if (snapshot.attestation) {
const ageMs = nowMs - snapshot.attestation.attestedAtMs;
if (ageMs <= WORKSPACE_ATTESTATION_RECENT_MS) {
return preserveRecentState();
}
}
if (
(snapshot.setup.bootstrapSeededAt || snapshot.setup.setupCompletedAt) &&
snapshot.setupUpdatedAtMs !== undefined
) {
const ageMs = nowMs - snapshot.setupUpdatedAtMs;
if (ageMs <= WORKSPACE_ATTESTATION_RECENT_MS) {
return preserveRecentState();
}
}
deleteWorkspaceRows(database, identity);
return true;
});
const result = await runWorkspaceStateOperation(
{
type: "workspace.expire",
input: {
workspaceDir: path.resolve(resolveUserPath(workspaceDir)),
nowMs,
},
},
options,
);
return result !== false;
}
/** Capture workspace identity before the filesystem entry is removed. */

View file

@ -0,0 +1,40 @@
import type { AgentDeletionRecoveryHoldPredicate } from "../state/agent-deletion-journal-recovery.kernel.js";
import type {
WorkspaceSetupState,
WorkspaceStateSnapshot,
} from "./workspace-state-store.kernel.js";
export type WorkspaceStateGuard = {
/** Host lifecycle and filesystem authority only; never reads SQLite. */
assertHost?: () => void;
recoveryHoldPredicate?: AgentDeletionRecoveryHoldPredicate;
/** Released SDK compatibility: before worker dispatch or host file mutation, never a grant. */
beforeLegacyApply?: () => void;
};
type WorkspaceStateInput = {
workspaceDir: string;
recoveryHoldPredicate?: AgentDeletionRecoveryHoldPredicate;
};
export type WorkspaceStateWorkerOperations = {
"workspace.snapshotAndRegister": {
input: WorkspaceStateInput;
output: WorkspaceStateSnapshot;
};
"workspace.mergeSetup": {
input: WorkspaceStateInput & {
next: Partial<Omit<WorkspaceSetupState, "version">>;
nowMs: number;
};
output: WorkspaceSetupState;
};
"workspace.expire": { input: WorkspaceStateInput & { nowMs: number }; output: string | false };
};
export type WorkspaceStateWorkerCommand = {
[K in keyof WorkspaceStateWorkerOperations]: {
type: K;
input: WorkspaceStateWorkerOperations[K]["input"];
};
}[keyof WorkspaceStateWorkerOperations];

View file

@ -0,0 +1,179 @@
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../infra/kysely-sync.js";
import { runSqliteDeferredTransactionSync } from "../infra/sqlite-transaction.js";
import {
deferSqliteWorkerCommitReceipt,
requestSqliteWorkerOperationAdmission,
} from "../infra/sqlite-worker-operation-admission.js";
import { assertAgentDeletionRecoveryHoldPredicate } from "../state/agent-deletion-journal-recovery.kernel.js";
import {
runOpenClawStateWriteTransaction,
type OpenClawStateDatabaseOptions,
} from "../state/openclaw-state-db.js";
import {
assertCanonicalIntegerTimestamp,
assertCanonicalTimestamp,
deleteWorkspaceStateRowsInDatabase,
readWorkspaceStateSnapshotFromDatabase,
registerWorkspaceStateAliasIdentitiesInTransaction,
resolveWorkspaceIdentityFromDatabase,
hasRecentWorkspaceSetupState,
recentWorkspaceAttestation,
WORKSPACE_SETUP_STATE_VERSION,
type WorkspaceSetupState,
type WorkspaceStateDatabase,
type WorkspaceStateDatabaseHandle,
} from "./workspace-state-store.kernel.js";
import type { WorkspaceStateWorkerCommand } from "./workspace-state-store.worker-contract.js";
function mergeSetup(
database: WorkspaceStateDatabaseHandle,
workspaceDir: string,
next: Partial<Omit<WorkspaceSetupState, "version">>,
nowMs: number,
): WorkspaceSetupState {
assertCanonicalIntegerTimestamp(nowMs, "setup update");
if (next.bootstrapSeededAt) {
assertCanonicalTimestamp(next.bootstrapSeededAt, "bootstrap seeded");
}
if (next.setupCompletedAt) {
assertCanonicalTimestamp(next.setupCompletedAt, "setup completed");
}
const resolution = resolveWorkspaceIdentityFromDatabase({ workspaceDir, database });
const identity = resolution.identity;
const snapshot = readWorkspaceStateSnapshotFromDatabase({ identity, database });
const bootstrapSeededAt = snapshot.setup.bootstrapSeededAt ?? next.bootstrapSeededAt;
const setupCompletedAt = snapshot.setup.setupCompletedAt ?? next.setupCompletedAt;
const merged: WorkspaceSetupState = {
version: WORKSPACE_SETUP_STATE_VERSION,
...(bootstrapSeededAt ? { bootstrapSeededAt } : {}),
...(setupCompletedAt ? { setupCompletedAt } : {}),
};
const kysely = getNodeSqliteKysely<WorkspaceStateDatabase>(database.db);
executeSqliteQuerySync(
database.db,
kysely
.insertInto("workspace_setup_state")
.values({
workspace_key: identity.workspaceKey,
workspace_path: identity.workspacePath,
version: WORKSPACE_SETUP_STATE_VERSION,
bootstrap_seeded_at: merged.bootstrapSeededAt ?? null,
setup_completed_at: merged.setupCompletedAt ?? null,
updated_at: nowMs,
})
.onConflict((conflict) =>
conflict.column("workspace_key").doUpdateSet({
workspace_path: identity.workspacePath,
version: WORKSPACE_SETUP_STATE_VERSION,
bootstrap_seeded_at: merged.bootstrapSeededAt ?? null,
setup_completed_at: merged.setupCompletedAt ?? null,
updated_at: nowMs,
}),
),
);
registerWorkspaceStateAliasIdentitiesInTransaction({
database,
identity,
aliases: resolution.aliases,
updatedAtMs: nowMs,
});
return merged;
}
function expire(
database: WorkspaceStateDatabaseHandle,
workspaceDir: string,
nowMs: number,
): string | false {
assertCanonicalIntegerTimestamp(nowMs, "workspace expiry check");
const resolution = resolveWorkspaceIdentityFromDatabase({ workspaceDir, database });
const identity = resolution.identity;
const snapshot = readWorkspaceStateSnapshotFromDatabase({ identity, database });
const preserveRecentState = () => {
registerWorkspaceStateAliasIdentitiesInTransaction({
database,
identity,
aliases: resolution.aliases,
updatedAtMs: nowMs,
});
return false as const;
};
if (
recentWorkspaceAttestation(snapshot.attestation, nowMs) ||
hasRecentWorkspaceSetupState(snapshot, nowMs)
) {
return preserveRecentState();
}
deleteWorkspaceStateRowsInDatabase(database, identity);
return identity.workspacePath;
}
export function executeWorkspaceStateCommand(
command: WorkspaceStateWorkerCommand,
database: WorkspaceStateDatabaseHandle,
options: OpenClawStateDatabaseOptions,
) {
if (command.type === "workspace.snapshotAndRegister") {
const initial = runSqliteDeferredTransactionSync(database.db, () => {
assertAgentDeletionRecoveryHoldPredicate(database, command.input.recoveryHoldPredicate);
const resolution = resolveWorkspaceIdentityFromDatabase({
workspaceDir: command.input.workspaceDir,
database,
});
return {
resolution,
snapshot: readWorkspaceStateSnapshotFromDatabase({
identity: resolution.identity,
database,
}),
};
});
if (
initial.resolution.missingAliasKeys.length === 0 ||
(!initial.snapshot.setupExists && !initial.snapshot.attestation)
) {
return initial.snapshot;
}
return runOpenClawStateWriteTransaction((writer) => {
requestSqliteWorkerOperationAdmission({ stage: "transaction", facts: undefined });
const resolution = resolveWorkspaceIdentityFromDatabase({
workspaceDir: command.input.workspaceDir,
database: writer,
});
if (resolution.identity.workspaceKey !== initial.resolution.identity.workspaceKey) {
throw new Error("Workspace state identity changed before alias registration");
}
const snapshot = readWorkspaceStateSnapshotFromDatabase({
identity: resolution.identity,
database: writer,
});
if (snapshot.setupExists || snapshot.attestation) {
registerWorkspaceStateAliasIdentitiesInTransaction({
database: writer,
identity: resolution.identity,
aliases: resolution.aliases,
updatedAtMs: Date.now(),
});
}
requestSqliteWorkerOperationAdmission({ stage: "commit", facts: undefined });
assertAgentDeletionRecoveryHoldPredicate(writer, command.input.recoveryHoldPredicate);
return snapshot;
}, options);
}
return runOpenClawStateWriteTransaction((writer) => {
requestSqliteWorkerOperationAdmission({ stage: "transaction", facts: undefined });
const result =
command.type === "workspace.mergeSetup"
? mergeSetup(writer, command.input.workspaceDir, command.input.next, command.input.nowMs)
: expire(writer, command.input.workspaceDir, command.input.nowMs);
requestSqliteWorkerOperationAdmission({
stage: "commit",
facts: command.type === "workspace.expire" ? result : undefined,
});
assertAgentDeletionRecoveryHoldPredicate(writer, command.input.recoveryHoldPredicate);
if (command.type === "workspace.expire") {
deferSqliteWorkerCommitReceipt(writer.db, result);
}
return result;
}, options);
}

View file

@ -540,10 +540,12 @@ describe("workspace completion persistence", () => {
const pending = ensureAgentWorkspace({
dir,
ensureBootstrapFiles: true,
beforePersistentApply: () => {
if (!current) {
throw new Error("workspace owner retired");
}
guard: {
assertHost: () => {
if (!current) {
throw new Error("workspace owner retired");
}
},
},
});
void pending.then(

View file

@ -12,12 +12,12 @@ import { isPathInside } from "../infra/path-guards.js";
import { retryAsync } from "../infra/retry.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import { exactWorkspaceEntryExists } from "../memory/root-memory-files.js";
import { runCommandWithTimeout } from "../process/exec.js";
import { isCronSessionKey, isSubagentSessionKey } from "../routing/session-key.js";
import { deriveSessionChatTypeFromKey } from "../sessions/session-chat-type-shared.js";
import { createLazyPromise, getOrCreatePromise } from "../shared/lazy-promise.js";
import { getOrCreatePromise } from "../shared/lazy-promise.js";
import type { OpenClawStateDatabaseOptions } from "../state/openclaw-state-db.js";
import { resolveUserPath } from "../utils.js";
import { DuplicateAgentError } from "./agent-create-error.js";
import { getAgentWorkspaceAccess } from "./workspace-access.js";
import {
DEFAULT_AGENTS_FILENAME,
@ -43,27 +43,33 @@ import {
} from "./workspace-bootstrap-publish.js";
import { MAX_WORKSPACE_BOOTSTRAP_FILE_BYTES } from "./workspace-bootstrap-read.js";
import { DEFAULT_AGENT_WORKSPACE_DIR } from "./workspace-default.js";
import { createWorkspaceFileMutationGuard } from "./workspace-file-mutation-guard.js";
import {
isTransientWorkspaceReadError,
readWorkspaceFileWithGuards,
setWorkspaceFileSourceIdentity,
} from "./workspace-file-read.js";
import { ensureGitRepo } from "./workspace-git.js";
import {
assertNoUnmigratedWorkspaceState,
LEGACY_WORKSPACE_STATE_CURRENT_FILENAME,
LEGACY_WORKSPACE_STATE_DIRNAME,
} from "./workspace-legacy-state.js";
import { captureWorkspaceStateFilesystemGuard } from "./workspace-state-guard.js";
import { WorkspaceVanishedError } from "./workspace-state-identity.js";
import {
clearExpiredWorkspaceStateForVanishedWorkspace,
hasRecentWorkspaceSetupState,
hasWorkspaceSetupStateMarker,
recentWorkspaceAttestation,
mergeWorkspaceSetupState,
readWorkspaceStateSnapshot,
replaceWorkspaceAttestation,
WORKSPACE_ATTESTATION_RECENT_MS,
type WorkspaceAttestation,
type WorkspaceStateSnapshot,
type WorkspaceSetupState,
} from "./workspace-state-store.js";
import type { WorkspaceStateGuard } from "./workspace-state-store.worker-contract.js";
import { resolveWorkspaceTemplateSearchDirs } from "./workspace-templates.js";
export {
DEFAULT_AGENTS_FILENAME,
@ -99,7 +105,6 @@ const WORKSPACE_ONBOARDING_PROFILE_FILENAMES = [
const workspaceLogger = createSubsystemLogger("workspace");
const workspaceTemplateCache = new Map<string, Promise<string>>();
const gitInitializationInFlight = new Map<string, Promise<void>>();
function stripFrontMatter(content: string): string {
return extractFrontmatterBlock(content)?.body.replace(/^\s+/, "") ?? content;
@ -340,8 +345,10 @@ async function reconcileWorkspaceBootstrapCompletionState(params: {
bootstrapPath: string;
state: WorkspaceSetupState;
bootstrapExists?: boolean;
beforePersistentApply?: () => void;
guard?: WorkspaceStateGuard;
}): Promise<WorkspaceBootstrapCompletionReconcileResult> {
const assertEvidence = captureWorkspaceStateFilesystemGuard(params.dir);
const beforeFileMutation = createWorkspaceFileMutationGuard(params.guard);
const bootstrapExists = params.bootstrapExists ?? (await pathExists(params.bootstrapPath));
if (
typeof params.state.setupCompletedAt === "string" &&
@ -364,14 +371,21 @@ async function reconcileWorkspaceBootstrapCompletionState(params: {
bootstrapSeededAt: params.state.bootstrapSeededAt ?? now,
setupCompletedAt: now,
};
params.beforePersistentApply?.();
params.guard?.assertHost?.();
const persistedState = await mergeWorkspaceSetupState(params.dir, repairedState, undefined, {
assertCurrent: params.beforePersistentApply,
recoveryHoldPredicate: params.guard?.recoveryHoldPredicate,
beforeLegacyApply: params.guard?.beforeLegacyApply,
assertCurrent: () => {
params.guard?.assertHost?.();
assertEvidence();
},
});
params.guard?.assertHost?.();
assertEvidence();
if (!bootstrapExists) {
return { repaired: true, bootstrapExists: false, state: persistedState };
}
params.beforePersistentApply?.();
beforeFileMutation?.();
try {
await fs.rm(params.bootstrapPath, { force: true });
return { repaired: true, bootstrapExists: false, state: persistedState };
@ -396,58 +410,34 @@ async function collectGeneratedBootstrapHashes(dir: string): Promise<Map<string,
return hashes;
}
function recentWorkspaceAttestation(
attestation: WorkspaceAttestation | undefined,
nowMs = Date.now(),
): WorkspaceAttestation | undefined {
if (!attestation) {
return undefined;
}
const ageMs = nowMs - attestation.attestedAtMs;
// Clock rollback must not turn disappearance protection into permission to
// reseed. A healthy workspace refreshes the future-dated row below.
if (ageMs > WORKSPACE_ATTESTATION_RECENT_MS) {
return undefined;
}
return attestation;
}
async function maybeWriteWorkspaceAttestation(
dir: string,
beforePersistentApply?: () => void,
guard?: WorkspaceStateGuard,
): Promise<void> {
const assertHost = guard?.assertHost;
// Order snapshots by when their filesystem observation starts. The store
// compares against a separate lock-time clock, so a newer committed scan
// wins when this async collection finishes later.
const attestedAtMs = Date.now();
const generatedHashes = await collectGeneratedBootstrapHashes(dir);
beforePersistentApply?.();
assertHost?.();
try {
await replaceWorkspaceAttestation({
workspaceDir: dir,
attestedAtMs,
generatedHashes,
assertCurrent: beforePersistentApply,
assertCurrent: assertHost,
recoveryHoldPredicate: guard?.recoveryHoldPredicate,
beforeLegacyApply: guard?.beforeLegacyApply,
});
} catch {
} catch (error) {
if (error instanceof DuplicateAgentError) {
throw error;
}
// Attestation is a lifecycle guard; setup should not fail solely because
// the auxiliary disappearance evidence could not be refreshed.
}
beforePersistentApply?.();
}
function hasWorkspaceSetupStateMarker(state: WorkspaceSetupState): boolean {
return Boolean(state.bootstrapSeededAt || state.setupCompletedAt);
}
function hasRecentWorkspaceSetupState(
snapshot: WorkspaceStateSnapshot,
nowMs = Date.now(),
): boolean {
if (!hasWorkspaceSetupStateMarker(snapshot.setup) || snapshot.setupUpdatedAtMs === undefined) {
return false;
}
return nowMs - snapshot.setupUpdatedAtMs <= WORKSPACE_ATTESTATION_RECENT_MS;
assertHost?.();
}
async function workspaceAttestationHasSurvivalEvidence(params: {
@ -479,7 +469,7 @@ async function workspaceSetupStateHasSurvivalEvidence(params: {
dir: string;
bootstrapPath: string;
initialState: WorkspaceStateSnapshot;
beforePersistentApply?: () => void;
guard?: WorkspaceStateGuard;
}): Promise<boolean> {
if (await pathExists(params.bootstrapPath)) {
return true;
@ -490,7 +480,7 @@ async function workspaceSetupStateHasSurvivalEvidence(params: {
const currentState = await readCanonicalWorkspaceStateSnapshot(
params.dir,
undefined,
params.beforePersistentApply,
params.guard,
);
if (
currentState.setup.bootstrapSeededAt !== params.initialState.setup.bootstrapSeededAt ||
@ -505,9 +495,15 @@ async function workspaceSetupStateHasSurvivalEvidence(params: {
async function readCanonicalWorkspaceStateSnapshot(
dir: string,
options: OpenClawStateDatabaseOptions = {},
assertCurrent?: () => void,
guard?: WorkspaceStateGuard,
): Promise<WorkspaceStateSnapshot> {
const snapshot = await readWorkspaceStateSnapshot(dir, { ...options, assertCurrent });
const snapshot = await readWorkspaceStateSnapshot(dir, {
...options,
assertCurrent: guard?.assertHost,
recoveryHoldPredicate: guard?.recoveryHoldPredicate,
beforeLegacyApply: guard?.beforeLegacyApply,
});
guard?.assertHost?.();
assertNoUnmigratedWorkspaceState({
workspaceDir: dir,
});
@ -664,49 +660,6 @@ export async function isWorkspaceBootstrapPending(dir: string): Promise<boolean>
return (await resolveWorkspaceBootstrapStatus(dir)) === "pending";
}
// Git availability is process-stable; cache the probe result, including failure, until restart.
const isGitAvailable = createLazyPromise(async () => {
try {
const result = await runCommandWithTimeout(["git", "--version"], { timeoutMs: 2_000 });
return result.code === 0;
} catch {
return false;
}
});
async function ensureGitRepo(
dir: string,
isBrandNewWorkspace: boolean,
beforePersistentApply?: () => void,
) {
if (!isBrandNewWorkspace) {
return;
}
// Concurrent first turns can all observe missing Git metadata. Join only the
// current initialization; later calls must inspect the workspace again.
beforePersistentApply?.();
await getOrCreatePromise(
gitInitializationInFlight,
dir,
async () => {
if (await fs.stat(path.join(dir, ".git")).catch(() => undefined)) {
return;
}
if (!(await isGitAvailable())) {
return;
}
// Only the initializer's owner admits Git; joining callers cannot cancel it.
beforePersistentApply?.();
try {
await runCommandWithTimeout(["git", "init"], { cwd: dir, timeoutMs: 10_000 });
} catch {
// Ignore git init failures; workspace creation should still succeed.
}
},
{ evictOnSettled: true },
);
}
export async function ensureAgentWorkspace(params?: {
dir?: string;
ensureBootstrapFiles?: boolean;
@ -715,7 +668,7 @@ export async function ensureAgentWorkspace(params?: {
/** Creation-time role content; existing workspace files are still preserved. */
templates?: Partial<Record<"AGENTS.md" | "SOUL.md" | "IDENTITY.md", string>>;
/** Guard each new mutation after async preparation; admitted effects may settle. */
beforePersistentApply?: () => void;
guard?: WorkspaceStateGuard;
/**
* List of optional bootstrap filenames to skip writing.
* Applies only to SOUL.md, USER.md, IDENTITY.md.
@ -741,16 +694,26 @@ export async function ensureAgentWorkspace(params?: {
}> {
const rawDir = params?.dir?.trim() ? params.dir.trim() : DEFAULT_AGENT_WORKSPACE_DIR;
const dir = resolveUserPath(rawDir);
const beforePersistentApply = params?.beforePersistentApply;
const guard = params?.guard;
const assertHost = guard?.assertHost;
const beforeFileMutation = createWorkspaceFileMutationGuard(guard);
let assertExpiryEvidence: () => void;
const clearExpiredState = async () => {
beforePersistentApply?.();
assertHost?.();
if (
!(await clearExpiredWorkspaceStateForVanishedWorkspace(dir, undefined, {
assertCurrent: beforePersistentApply,
recoveryHoldPredicate: guard?.recoveryHoldPredicate,
beforeLegacyApply: guard?.beforeLegacyApply,
assertCurrent: () => {
assertHost?.();
assertExpiryEvidence();
},
}))
) {
throw new WorkspaceVanishedError({ workspaceDir: dir });
}
assertHost?.();
assertExpiryEvidence();
};
const purpose = params?.purpose?.trim();
if (purpose && (params?.templates || getAgentWorkspaceAccess(dir))) {
@ -767,18 +730,15 @@ export async function ensureAgentWorkspace(params?: {
// The workspace belongs to a runtime-managed agent with a distinct cwd.
// Provision the directory (cwd fallback, media staging) without scaffolding
// bootstrap files, setup state, or a nested git repository (#92015).
beforePersistentApply?.();
beforeFileMutation?.();
await fs.mkdir(dir, { recursive: true });
return { dir, bootstrapPending: false };
}
let initialState = await readCanonicalWorkspaceStateSnapshot(
dir,
undefined,
beforePersistentApply,
);
let initialState = await readCanonicalWorkspaceStateSnapshot(dir, undefined, guard);
let reseedingExpiredWorkspaceState = false;
const recentAttestation = recentWorkspaceAttestation(initialState.attestation);
const recentSetupState = hasRecentWorkspaceSetupState(initialState);
assertExpiryEvidence = captureWorkspaceStateFilesystemGuard(dir);
const workspaceExists = await pathExists(dir);
if (!workspaceExists) {
@ -791,8 +751,9 @@ export async function ensureAgentWorkspace(params?: {
await clearExpiredState();
}
beforePersistentApply?.();
beforeFileMutation?.();
await fs.mkdir(dir, { recursive: true });
assertExpiryEvidence = captureWorkspaceStateFilesystemGuard(dir);
const bootstrapPath = path.join(dir, DEFAULT_BOOTSTRAP_FILENAME);
if (!params?.ensureBootstrapFiles) {
@ -807,7 +768,7 @@ export async function ensureAgentWorkspace(params?: {
dir,
bootstrapPath,
initialState,
beforePersistentApply,
guard,
}))
) {
if (recentSetupState) {
@ -820,11 +781,11 @@ export async function ensureAgentWorkspace(params?: {
path.join(dir, DEFAULT_AGENTS_FILENAME),
await loadTemplate(DEFAULT_AGENTS_FILENAME),
purpose,
beforePersistentApply,
beforeFileMutation,
);
}
if (hasContentEvidence || purpose) {
await maybeWriteWorkspaceAttestation(dir, beforePersistentApply);
await maybeWriteWorkspaceAttestation(dir, guard);
}
return { dir, bootstrapPending: false };
}
@ -878,7 +839,7 @@ export async function ensureAgentWorkspace(params?: {
dir,
bootstrapPath,
initialState,
beforePersistentApply,
guard,
}))
) {
// Setup can outlive a best-effort attestation write or arrive alone from
@ -901,7 +862,7 @@ export async function ensureAgentWorkspace(params?: {
const userTemplate = await loadTemplate(DEFAULT_USER_FILENAME);
// Template and filesystem checks above are async. Another process may have
// completed setup while they ran, so optional-file policy needs fresh state.
initialState = await readCanonicalWorkspaceStateSnapshot(dir, undefined, beforePersistentApply);
initialState = await readCanonicalWorkspaceStateSnapshot(dir, undefined, guard);
const skipOptionalBootstrapFiles = new Set(params?.skipOptionalBootstrapFiles ?? []);
// When the workspace is already configured, skip optional bootstrap files to
// prevent subagent spawns from recreating root-level SOUL.md, USER.md, or
@ -912,19 +873,18 @@ export async function ensureAgentWorkspace(params?: {
skipOptionalBootstrapFiles.add(filename);
}
}
await publishAgentInstructions(agentsPath, defaultAgentsTemplate, purpose, beforePersistentApply);
await publishAgentInstructions(agentsPath, defaultAgentsTemplate, purpose, beforeFileMutation);
if (!skipOptionalBootstrapFiles.has(DEFAULT_SOUL_FILENAME)) {
await publishBootstrapFile(soulPath, soulTemplate, beforePersistentApply);
await publishBootstrapFile(soulPath, soulTemplate, beforeFileMutation);
}
const identityPathCreated = !skipOptionalBootstrapFiles.has(DEFAULT_IDENTITY_FILENAME)
? await publishBootstrapFile(identityPath, identityTemplate, beforePersistentApply)
? await publishBootstrapFile(identityPath, identityTemplate, beforeFileMutation)
: false;
if (!skipOptionalBootstrapFiles.has(DEFAULT_USER_FILENAME)) {
await publishBootstrapFile(userPath, userTemplate, beforePersistentApply);
await publishBootstrapFile(userPath, userTemplate, beforeFileMutation);
}
let state = (await readCanonicalWorkspaceStateSnapshot(dir, undefined, beforePersistentApply))
.setup;
let state = (await readCanonicalWorkspaceStateSnapshot(dir, undefined, guard)).setup;
let stateDirty = false;
const markState = (next: Partial<WorkspaceSetupState>) => {
state = { ...state, ...next };
@ -943,7 +903,7 @@ export async function ensureAgentWorkspace(params?: {
bootstrapPath,
state,
bootstrapExists,
beforePersistentApply,
guard,
});
if (repair.repaired) {
state = repair.state;
@ -975,7 +935,7 @@ export async function ensureAgentWorkspace(params?: {
const wroteBootstrap = await publishBootstrapFile(
bootstrapPath,
bootstrapTemplate,
beforePersistentApply,
beforeFileMutation,
);
bootstrapExists = wroteBootstrap || (await pathExists(bootstrapPath));
if (bootstrapExists && !state.bootstrapSeededAt) {
@ -985,13 +945,15 @@ export async function ensureAgentWorkspace(params?: {
}
if (stateDirty) {
beforePersistentApply?.();
assertHost?.();
state = await mergeWorkspaceSetupState(dir, state, undefined, {
assertCurrent: beforePersistentApply,
assertCurrent: assertHost,
recoveryHoldPredicate: guard?.recoveryHoldPredicate,
beforeLegacyApply: guard?.beforeLegacyApply,
});
}
await ensureGitRepo(dir, isBrandNewWorkspace, beforePersistentApply);
await maybeWriteWorkspaceAttestation(dir, beforePersistentApply);
await ensureGitRepo(dir, isBrandNewWorkspace, beforeFileMutation);
await maybeWriteWorkspaceAttestation(dir, guard);
return {
dir,

View file

@ -14,7 +14,7 @@ import { extractSqliteTableSchema } from "../infra/sqlite-schema-sql.js";
import { createLegacyDatabaseFixture } from "../infra/state-migrations.media-persistence.test-support.js";
import { setLoggerOverride } from "../logging/logger.js";
import { testApi } from "../logging/logger.test-support.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.kernel.js";
import { OPENCLAW_AGENT_SCHEMA_VERSION } from "../state/openclaw-agent-db-contract.js";
import { unregisterOpenClawAgentDatabase } from "../state/openclaw-agent-db-registry.js";
import { recordOpenClawDatabaseQuarantine } from "../state/openclaw-quarantine-store.js";

View file

@ -8,6 +8,7 @@ import { normalizeOptionalString } from "@openclaw/normalization-core/string-coe
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
import { stylePromptTitle } from "../../packages/terminal-core/src/prompt-style.js";
import { resolveAgentEffectiveModelPrimary, resolveDefaultAgentId } from "../agents/agent-scope.js";
import type { WorkspaceStateGuard } from "../agents/workspace-state-store.worker-contract.js";
import { DEFAULT_AGENT_WORKSPACE_DIR, ensureAgentWorkspace } from "../agents/workspace.js";
import { printClawBanner } from "../cli/claw-banner.js";
import { readSourceConfigBestEffort } from "../config/config.js";
@ -201,18 +202,18 @@ export async function ensureWorkspaceAndSessions(
skipBootstrap?: boolean;
skipOptionalBootstrapFiles?: OptionalBootstrapFileName[];
agentId: string;
beforePersistentApply?: () => void;
guard?: WorkspaceStateGuard;
},
): Promise<{ bootstrapPending: boolean }> {
const ws = await ensureAgentWorkspace({
dir: workspaceDir,
ensureBootstrapFiles: !options.skipBootstrap,
skipOptionalBootstrapFiles: options.skipOptionalBootstrapFiles,
beforePersistentApply: options.beforePersistentApply,
guard: options.guard,
});
runtime.log(`Workspace OK: ${shortenHomePath(ws.dir)}`);
const sessionsDir = resolveSessionTranscriptsDirForAgent(options.agentId);
options.beforePersistentApply?.();
options.guard?.assertHost?.();
await fs.mkdir(sessionsDir, { recursive: true });
runtime.log(`Sessions OK: ${shortenHomePath(sessionsDir)}`);
return { bootstrapPending: ws.bootstrapPending === true };

View file

@ -5,7 +5,7 @@ import { expect, it, vi } from "vitest";
import { prepareDoctorDatabasePreflight } from "../commands/doctor-database-preflight.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { createLegacyDatabaseFixture } from "../infra/state-migrations.media-persistence.test-support.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.js";
import { readAgentDeletionRecoveryHolds } from "../state/agent-deletion-journal-recovery.kernel.js";
import { unregisterOpenClawAgentDatabase } from "../state/openclaw-agent-db-registry.js";
import {
closeOpenClawStateDatabaseForTest,

View file

@ -537,7 +537,7 @@ export const agentsHandlers: GatewayRequestHandlers = {
const skipBootstrap = Boolean(nextConfig.agents?.defaults?.skipBootstrap);
ensuredWorkspace = await ensureAgentWorkspace({
dir: workspaceDir,
beforePersistentApply: assertUploadAllowed,
guard: { assertHost: assertUploadAllowed },
ensureBootstrapFiles: !skipBootstrap,
skipOptionalBootstrapFiles: nextConfig.agents?.defaults?.skipOptionalBootstrapFiles,
});

View file

@ -0,0 +1,79 @@
import fs from "node:fs";
import { isMainThread } from "node:worker_threads";
import { extractErrorCode } from "@openclaw/normalization-core/error-coercion";
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
import { isPidAlive } from "../shared/pid-alive.js";
import { parseGatewayLockPayload } from "./gateway-lock-payload.js";
import { isLockOwnerDefinitelyStale } from "./stale-lock-file.js";
export const StateDatabaseAdmissionPendingError = resolveGlobalSingleton(
Symbol.for("openclaw.stateDatabaseAdmissionPendingError"),
() =>
class extends Error {
constructor(
readonly databasePath: string,
message: string,
) {
super(message);
}
},
);
/** Persisted ownership can identify a holder but cannot lend this process maintenance authority. */
export function assertPersistedStateDatabaseAccessAllowed(params: {
databasePath: string;
ownerPath: string;
assertMaintenance: () => void;
}): void {
const { databasePath, ownerPath, assertMaintenance } = params;
const unavailable = `OpenClaw state ownership at ${databasePath} could not be verified; retry after maintenance finishes.`;
let raw: string;
try {
raw = fs.readFileSync(ownerPath, "utf8");
} catch (error) {
if (extractErrorCode(error) === "ENOENT") {
return;
}
throw new Error(unavailable, { cause: error });
}
const owner = parseGatewayLockPayload(raw);
if (!owner) {
// Native exclusive creation precedes the payload write; cold admission may wait for publication.
throw new StateDatabaseAdmissionPendingError(databasePath, unavailable);
}
if (!Number.isSafeInteger(owner.pid) || owner.pid <= 0) {
throw new Error(unavailable);
}
if (
isLockOwnerDefinitelyStale({
payload: { pid: owner.pid, starttime: owner.startTime },
})
) {
return;
}
// Workers share the process PID, but their schema authority still comes from
// the host operation's retained lease and is never inferred from this record.
if (owner.pid === process.pid && !isMainThread) {
return;
}
if (!isPidAlive(owner.pid)) {
throw new Error(unavailable);
}
const role = owner.role ?? "gateway";
if (role === "gateway" || role === "agent-embedded") {
return;
}
if (owner.pid === process.pid) {
assertMaintenance();
return;
}
if (owner.stateOwnerKind === "schema" && owner.role === "sqlite-maintenance") {
throw new StateDatabaseAdmissionPendingError(
databasePath,
`OpenClaw state at ${databasePath} is undergoing offline maintenance; retry when it finishes.`,
);
}
throw new Error(
`OpenClaw state at ${databasePath} is undergoing offline maintenance; retry when it finishes.`,
);
}

View file

@ -3,7 +3,6 @@ import { randomUUID } from "node:crypto";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { isMainThread } from "node:worker_threads";
import { extractErrorCode } from "@openclaw/normalization-core/error-coercion";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import {
@ -31,6 +30,10 @@ import {
removeCreatedProjectionDirectories,
type StateOwnerDirectoryIdentity,
} from "./gateway-state-owner-directory.js";
import {
assertPersistedStateDatabaseAccessAllowed,
StateDatabaseAdmissionPendingError,
} from "./gateway-state-owner-record.js";
import { normalizeSqliteNonNegativeInteger } from "./sqlite-busy-timeout.js";
import { runWithSqliteCleanup } from "./sqlite-lifecycle-errors.js";
import { isLockOwnerDefinitelyStale } from "./stale-lock-file.js";
@ -180,19 +183,6 @@ export type GatewayStateOwnerContentionError = InstanceType<
typeof GatewayStateOwnerContentionError
>;
const StateDatabaseAdmissionPendingError = resolveGlobalSingleton(
Symbol.for("openclaw.stateDatabaseAdmissionPendingError"),
() =>
class extends Error {
constructor(
readonly databasePath: string,
message: string,
) {
super(message);
}
},
);
/** Retry cold admission only; the same budget covers opening and its first unentered write. */
export function withStateDatabaseColdAdmission<T>(
params: { databasePath: string; busyTimeoutMs: number; canRetry?: () => boolean },
@ -663,55 +653,11 @@ export function assertStateDatabaseAccessAllowed(
}
return;
}
let raw: string;
try {
raw = fs.readFileSync(pathname, "utf8");
} catch (error) {
if (extractErrorCode(error) === "ENOENT") {
return;
}
throw new Error(unavailable, { cause: error });
}
const owner = parseGatewayLockPayload(raw);
if (!owner) {
// Native exclusive creation precedes the payload write; cold admission may wait for publication.
throw new StateDatabaseAdmissionPendingError(databasePath, unavailable);
}
if (!Number.isSafeInteger(owner.pid) || owner.pid <= 0) {
throw new Error(unavailable);
}
if (
isLockOwnerDefinitelyStale({
payload: { pid: owner.pid, starttime: owner.startTime },
})
) {
return;
}
// Workers share the process PID, but their schema authority still comes from
// the host operation's retained lease and is never inferred from this record.
if (owner.pid === process.pid && !isMainThread) {
return;
}
if (!isPidAlive(owner.pid)) {
throw new Error(unavailable);
}
const role = owner.role ?? "gateway";
if (role === "gateway" || role === "agent-embedded") {
return;
}
if (owner.pid === process.pid) {
assertMaintenance();
return;
}
if (owner.stateOwnerKind === "schema" && owner.role === "sqlite-maintenance") {
throw new StateDatabaseAdmissionPendingError(
databasePath,
`OpenClaw state at ${databasePath} is undergoing offline maintenance; retry when it finishes.`,
);
}
throw new Error(
`OpenClaw state at ${databasePath} is undergoing offline maintenance; retry when it finishes.`,
);
assertPersistedStateDatabaseAccessAllowed({
databasePath,
ownerPath: pathname,
assertMaintenance,
});
}
/** Cleanup must compete with local roots too; it cannot borrow a live Gateway's authority. */

View file

@ -3,7 +3,10 @@ import path from "node:path";
import { afterEach } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { closeOpenClawStateDatabaseForTest } from "../state/openclaw-state-db.js";
import {
closeOpenClawStateDatabaseAsync,
closeOpenClawStateDatabaseForTest,
} from "../state/openclaw-state-db.js";
import { captureEnv, setTestEnvValue } from "../test-utils/env.js";
import {
detectLegacyWorkspaceState,
@ -13,7 +16,8 @@ import {
export function useWorkspaceMigrationTestFixture() {
let envSnapshot: ReturnType<typeof captureEnv> | undefined;
const tempDirs = useAutoCleanupTempDirTracker((cleanup) => {
afterEach(() => {
afterEach(async () => {
await closeOpenClawStateDatabaseAsync();
closeOpenClawStateDatabaseForTest();
envSnapshot?.restore();
envSnapshot = undefined;

View file

@ -6,6 +6,7 @@ import {
} from "./plugin-sdk-subpath-records.js";
import { SESSION_PERSISTENCE_COMPAT_RECORDS } from "./session-persistence-records.js";
import type { PluginCompatRecord } from "./types.js";
import { WORKSPACE_MUTATION_GUARD_COMPAT_RECORD } from "./workspace-mutation-guard.js";
const ACTIVATION_HINT_METADATA = {
status: "active",
@ -17,6 +18,7 @@ const ACTIVATION_HINT_METADATA = {
} as const;
export const PLUGIN_COMPAT_RECORDS = [
WORKSPACE_MUTATION_GUARD_COMPAT_RECORD,
...SESSION_PERSISTENCE_COMPAT_RECORDS,
{
code: "memory-session-sync-inventory",

View file

@ -0,0 +1,25 @@
import type { PluginCompatRecord } from "./types.js";
export const WORKSPACE_MUTATION_GUARD_COMPAT_RECORD = {
code: "workspace-mutation-guard-callback",
status: "deprecated",
owner: "sdk",
introduced: "2026-10-02",
deprecated: "2026-10-02",
warningStarts: "2026-10-02",
removalGate: "next-plugin-sdk-major",
replacement:
"Await database preparation before ensureAgentWorkspace and use SQL-free guard.assertHost for live authority. Internal recovery predicates run on the worker transaction connection.",
docsPath: "/plugins/sdk-migration/how-to-migrate#workspace-mutation-guards",
surfaces: ["api.runtime.agent.ensureAgentWorkspace.beforePersistentApply"],
diagnostics: [
"@deprecated JSDoc and one DEP_WORKSPACE_MUTATION_GUARD warning per process",
"Warning explains before-dispatch timing and deprecated but allowed synchronous OpenClaw DB access",
],
tests: [
"src/plugins/runtime/runtime-agent.workspace.test.ts",
"src/plugins/compat/registry.test.ts",
],
releaseNote:
"Released callbacks and their synchronous OpenClaw DB access remain supported. The legacy check runs once before worker dispatch, outside admission grants; use the typed guard for live revocation at commit. Removal requires the next Plugin SDK major.",
} as const satisfies PluginCompatRecord;

View file

@ -0,0 +1,45 @@
import { ensureAgentWorkspace } from "../../agents/workspace.js";
import { resolveGlobalSingleton } from "../../shared/global-singleton.js";
type PluginWorkspaceParams = Omit<
NonNullable<Parameters<typeof ensureAgentWorkspace>[0]>,
"guard"
> & {
guard?: { assertHost?: () => void };
/** @deprecated Runs before dispatch. Use SQL-free guard.assertHost for live commit authority. */
beforePersistentApply?: () => void;
};
export function ensurePluginAgentWorkspace(params?: PluginWorkspaceParams) {
const legacy = params?.beforePersistentApply;
if (!legacy) {
return ensureAgentWorkspace(params);
}
// Auxiliary attestation may ignore storage errors; a callback refusal stays fatal.
let refusal: { error: unknown } | undefined;
const check = (assertAllowed?: () => void) => {
if (refusal) {
throw refusal.error;
}
try {
assertAllowed?.();
} catch (error) {
refusal = { error };
throw error;
}
};
resolveGlobalSingleton(Symbol.for("openclaw.workspaceGuardDeprecation"), () => {
process.emitWarning(
"ensureAgentWorkspace.beforePersistentApply runs before dispatch; synchronous OpenClaw DB access in this callback is deprecated. Use SQL-free guard.assertHost for live revocation at commit. Removal: next Plugin SDK major.",
{ code: "DEP_WORKSPACE_MUTATION_GUARD", type: "DeprecationWarning" },
);
return true;
});
return ensureAgentWorkspace({
...params,
guard: {
assertHost: () => check(params?.guard?.assertHost),
beforeLegacyApply: () => check(legacy),
},
});
}

View file

@ -13,7 +13,6 @@ import {
resolveEffectiveAgentRuntime,
} from "../../agents/thinking-runtime.js";
import { resolveAgentTimeoutMs } from "../../agents/timeout.js";
import { ensureAgentWorkspace } from "../../agents/workspace.js";
import { normalizeThinkLevel, resolveThinkingProfile } from "../../auto-reply/thinking.js";
import { getRuntimeConfig } from "../../config/config.js";
import * as session from "../../config/sessions/lifecycle.js";
@ -44,6 +43,7 @@ import {
import { createLazyRuntimeMethod, createLazyRuntimeModule } from "../../shared/lazy-runtime.js";
import { getPluginRuntimeGatewayRequestScope } from "./gateway-request-scope.js";
import { resolveAgentCatalogCreateTarget } from "./runtime-agent-session-catalog.js";
import { ensurePluginAgentWorkspace } from "./runtime-agent-workspace.js";
import { defineCachedValue } from "./runtime-cache.js";
import type { PluginRuntime } from "./types.js";
@ -666,7 +666,7 @@ export function createRuntimeAgent(): PluginRuntime["agent"] {
},
resolveAgentTimeoutMs,
resolveCliBackendDispatchEligibility: resolveEmbeddedCliBackendDispatchEligibility,
ensureAgentWorkspace,
ensureAgentWorkspace: ensurePluginAgentWorkspace,
} satisfies Omit<
PluginRuntime["agent"],
"runCommandFromIngress" | "runEmbeddedAgent" | "session"

View file

@ -0,0 +1,146 @@
import fs from "node:fs";
import fsPromises from "node:fs/promises";
import { expect, it, vi } from "vitest";
import { readWorkspaceStateSnapshot } from "../../agents/workspace-state-store.js";
import * as admission from "../../infra/sqlite-worker-operation-admission.js";
import { withExistingOpenClawStateDatabaseCurrentReadOnly } from "../../state/openclaw-state-db-readonly.js";
import { runOpenClawStateWriteTransaction } from "../../state/openclaw-state-db.js";
import * as workerStore from "../../state/openclaw-state-worker-store.js";
import { createOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
import { createRuntimeAgent } from "./runtime-agent.js";
import type { PluginRuntime } from "./types.js";
it("allows deprecated plugin SQL checks once before dispatch while typed guards retain commit authority", async () => {
type Params = NonNullable<Parameters<PluginRuntime["agent"]["ensureAgentWorkspace"]>[0]>;
const state = await createOpenClawTestState({ layout: "state-only" });
const warning = vi.spyOn(process, "emitWarning").mockImplementation(() => {});
const ensure = createRuntimeAgent().ensureAgentWorkspace;
const originalAdmission = admission.createSqliteWorkerOperationAdmission;
const originalOperation = workerStore.runOpenClawStateWorkerOperation;
const originalMkdir = fsPromises.mkdir;
let phase: string | undefined;
let workspace: string;
const events: string[] = [];
vi.spyOn(admission, "createSqliteWorkerOperationAdmission").mockImplementation((admit, data) =>
originalAdmission((request, grant) => {
phase = request.stage;
try {
admit(request, () => {
events.push(`grant:${phase}`);
return grant();
});
} finally {
phase = undefined;
}
}, data),
);
vi.spyOn(workerStore, "runOpenClawStateWorkerOperation").mockImplementation(
(context, operation, options) =>
originalOperation(
context,
(scope) => {
const execute: typeof scope.execute = (...args) => {
events.push(`dispatch:${args[0].type}`);
return scope.execute(...args);
};
return operation({ execute });
},
options,
),
);
vi.spyOn(fsPromises, "mkdir").mockImplementation((dir, options) => {
if (dir === workspace) {
events.push("file:mkdir");
}
return originalMkdir(dir, options);
});
try {
runOpenClawStateWriteTransaction(() => undefined);
for (const mode of [
"initial-refusal",
"allowed-sql",
"legacy-refusal",
"typed-revocation",
] as const) {
workspace = state.path(mode);
events.length = 0;
if (mode !== "initial-refusal") {
fs.mkdirSync(workspace);
fs.writeFileSync(`${workspace}/AGENTS.md`, "Synthetic workspace instructions.\n");
}
const before = await readWorkspaceStateSnapshot(workspace, { readOnly: true });
const refusal = new Error("plugin authority revoked");
const params: Params = {
dir: workspace,
ensureBootstrapFiles: false,
guard: {
assertHost() {
if (phase) {
events.push("typed");
if (mode === "typed-revocation" && phase === "commit") {
throw refusal;
}
}
},
},
beforePersistentApply() {
expect(phase).toBeUndefined();
const previous = events.at(-1);
events.push("legacy");
if (
mode === "initial-refusal" ||
(mode === "legacy-refusal" && previous === "file:mkdir")
) {
throw refusal;
}
expect(
withExistingOpenClawStateDatabaseCurrentReadOnly(({ db }) =>
db.prepare("SELECT 1 AS value").get(),
),
).toEqual({ value: 1 });
},
};
const pending = ensure(params);
if (mode === "allowed-sql") {
await pending;
expect(events.filter((event) => event !== "typed" && !event.startsWith("grant:"))).toEqual([
"legacy",
"dispatch:workspace.snapshotAndRegister",
"legacy",
"file:mkdir",
"legacy",
"dispatch:workspace.replaceAttestation",
]);
for (const [index, event] of events.entries()) {
if (event.startsWith("grant:")) {
expect(events[index - 1]).toBe("typed");
}
}
expect(events).toContain("grant:commit");
expect(
(await readWorkspaceStateSnapshot(workspace, { readOnly: true })).attestation,
).toBeDefined();
} else {
await expect(pending).rejects.toBe(refusal);
expect(await readWorkspaceStateSnapshot(workspace, { readOnly: true })).toEqual(before);
if (mode === "initial-refusal") {
expect(events).toEqual(["legacy"]);
expect(fs.existsSync(workspace)).toBe(false);
} else if (mode === "legacy-refusal") {
expect(events.filter((event) => event.startsWith("dispatch:"))).toEqual([
"dispatch:workspace.snapshotAndRegister",
]);
}
}
}
expect(warning).toHaveBeenCalledExactlyOnceWith(
expect.stringMatching(
/before dispatch.*synchronous OpenClaw DB access.*deprecated.*guard.assertHost/,
),
{ code: "DEP_WORKSPACE_MUTATION_GUARD", type: "DeprecationWarning" },
);
} finally {
vi.restoreAllMocks();
await state.cleanup();
}
});

View file

@ -388,7 +388,7 @@ export type PluginRuntimeCore = {
* budget timeouts for the run that will actually execute.
*/
resolveCliBackendDispatchEligibility: typeof import("../../agents/embedded-agent-runner/cli-backend-dispatch-eligibility.js").resolveEmbeddedCliBackendDispatchEligibility;
ensureAgentWorkspace: typeof import("../../agents/workspace.js").ensureAgentWorkspace;
ensureAgentWorkspace: typeof import("./runtime-agent-workspace.js").ensurePluginAgentWorkspace;
session: {
resolveStorePath: typeof import("../../config/sessions/paths.js").resolveSessionStorePathCore;
createSessionEntry: (

View file

@ -0,0 +1,76 @@
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { DuplicateAgentError } from "../agents/agent-create-error.js";
import { readLegacyMigrationReceiptFromDatabase } from "../infra/state-migrations.receipts.js";
import { normalizeAgentId } from "../routing/session-key.js";
import type { HeldAgentDatabase } from "./agent-deletion-journal.types.js";
import type { OpenClawStateDatabase } from "./openclaw-state-db-contract.js";
import { resolveOpenClawRegisteredAgentDatabasePath } from "./openclaw-state-db.paths.js";
export type RecoveryDatabase = Pick<OpenClawStateDatabase, "db" | "path">;
export type RecoveryReport = { description: string; held: HeldAgentDatabase[] };
export type AgentDeletionRecoveryHoldPredicate = {
agentId: string;
held: readonly HeldAgentDatabase[];
applies: boolean;
};
export const AGENT_DELETION_RECOVERY_SOURCE_KEY = "agent-deletion-journal-reconstruction";
export function readReport(database: RecoveryDatabase): RecoveryReport | undefined {
const receipt = readLegacyMigrationReceiptFromDatabase(
database.db,
AGENT_DELETION_RECOVERY_SOURCE_KEY,
);
if (!receipt) {
return undefined;
}
const report: unknown = JSON.parse(receipt.reportJson);
if (!isRecord(report) || typeof report.description !== "string" || !Array.isArray(report.held)) {
throw new Error("Invalid agent deletion journal reconstruction receipt.");
}
const held = report.held.map((entry: unknown) => {
if (
!isRecord(entry) ||
typeof entry.agentId !== "string" ||
entry.agentId !== normalizeAgentId(entry.agentId) ||
typeof entry.path !== "string" ||
!entry.path.trim()
) {
throw new Error("Invalid agent database hold in deletion journal reconstruction receipt.");
}
return { agentId: entry.agentId, path: entry.path };
});
return { description: report.description, held };
}
export function decodeHolds(database: RecoveryDatabase, held: readonly HeldAgentDatabase[]) {
return held.map((entry) => ({
agentId: entry.agentId,
path: resolveOpenClawRegisteredAgentDatabasePath(database.path, entry.path),
}));
}
/** Read recovery facts from this exact shared-state generation without opening another database. */
export function readAgentDeletionRecoveryHolds(database: RecoveryDatabase): HeldAgentDatabase[] {
return decodeHolds(database, readReport(database)?.held ?? []);
}
export function assertAgentDeletionRecoveryHoldPredicate(
database: RecoveryDatabase,
predicate?: AgentDeletionRecoveryHoldPredicate,
): void {
if (
predicate?.applies &&
readAgentDeletionRecoveryHolds(database).some(
(entry) =>
entry.agentId === predicate.agentId &&
!predicate.held.some(
(previous) => previous.agentId === entry.agentId && previous.path === entry.path,
),
)
) {
throw new DuplicateAgentError(
`Agent ${predicate.agentId} has held databases. Restore its original agentDir and session.store configuration, then run agents add explicitly to restore the preserved store.`,
);
}
}

View file

@ -1,69 +1,30 @@
import { randomUUID } from "node:crypto";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../infra/kysely-sync.js";
import {
readLegacyMigrationReceiptFromDatabase,
recordLegacyMigrationReceipt,
} from "../infra/state-migrations.receipts.js";
import { recordLegacyMigrationReceipt } from "../infra/state-migrations.receipts.js";
import { normalizeAgentId } from "../routing/session-key.js";
import {
AGENT_DELETION_RECOVERY_SOURCE_KEY as SOURCE_KEY,
decodeHolds,
readAgentDeletionRecoveryHolds,
readReport,
type RecoveryDatabase,
type RecoveryReport,
} from "./agent-deletion-journal-recovery.kernel.js";
import type { HeldAgentDatabase } from "./agent-deletion-journal.types.js";
import { createOpenClawAgentDatabasePathMatcher } from "./openclaw-agent-db.paths.js";
import type { OpenClawStateDatabase } from "./openclaw-state-db-contract.js";
import {
assertAgentDeletionJournalAvailable,
reconstructAgentDeletionJournalSchema,
} from "./openclaw-state-db-schema-additive.js";
import type { DB } from "./openclaw-state-db.generated.js";
import {
resolveOpenClawAgentDatabaseStoredPath,
resolveOpenClawRegisteredAgentDatabasePath,
} from "./openclaw-state-db.paths.js";
import { resolveOpenClawAgentDatabaseStoredPath } from "./openclaw-state-db.paths.js";
type RecoveryDatabase = Pick<OpenClawStateDatabase, "db" | "path">;
type RecoveryReport = { description: string; held: HeldAgentDatabase[] };
type RecoveryTables = Pick<DB, "migration_sources">;
const SOURCE_KEY = "agent-deletion-journal-reconstruction";
const JOURNAL_TABLE = "agent_deletion_journal";
const DESCRIPTION =
"Reconstructed the missing agent deletion journal; held databases require explicit agent restore or delete.";
function readReport(database: RecoveryDatabase): RecoveryReport | undefined {
const receipt = readLegacyMigrationReceiptFromDatabase(database.db, SOURCE_KEY);
if (!receipt) {
return undefined;
}
const report: unknown = JSON.parse(receipt.reportJson);
if (!isRecord(report) || typeof report.description !== "string" || !Array.isArray(report.held)) {
throw new Error("Invalid agent deletion journal reconstruction receipt.");
}
const held = report.held.map((entry: unknown) => {
if (
!isRecord(entry) ||
typeof entry.agentId !== "string" ||
entry.agentId !== normalizeAgentId(entry.agentId) ||
typeof entry.path !== "string" ||
!entry.path.trim()
) {
throw new Error("Invalid agent database hold in deletion journal reconstruction receipt.");
}
return { agentId: entry.agentId, path: entry.path };
});
return { description: report.description, held };
}
function decodeHolds(database: RecoveryDatabase, held: readonly HeldAgentDatabase[]) {
return held.map((entry) => ({
agentId: entry.agentId,
path: resolveOpenClawRegisteredAgentDatabasePath(database.path, entry.path),
}));
}
/** Read recovery facts from this exact shared-state generation without opening another database. */
export function readAgentDeletionRecoveryHolds(database: RecoveryDatabase): HeldAgentDatabase[] {
return decodeHolds(database, readReport(database)?.held ?? []);
}
/** Schema and Doctor producers cannot infer permission to repair an unverified store. */
export function assertAgentDeletionRecoveryAllowsMutation(
database: RecoveryDatabase,

View file

@ -13,7 +13,7 @@ import { runSqliteDeferredTransactionSync } from "../infra/sqlite-transaction.js
import { normalizeAgentId } from "../routing/session-key.js";
import { isSessionStoreTopologyChange, sessionChanges } from "../sessions/session-row-changes.js";
import { hasPreJournalStateSchema } from "./agent-deletion-journal-history.js";
import { readAgentDeletionRecoveryHolds } from "./agent-deletion-journal-recovery.js";
import { readAgentDeletionRecoveryHolds } from "./agent-deletion-journal-recovery.kernel.js";
import type {
AgentDatabaseDeletionSnapshot,
AgentDeletionJournalAuthority,

View file

@ -548,7 +548,10 @@ export function withExistingOpenClawStateDatabaseArtifactPreservingReadOnly<T>(
/** Publication guards need current rows, never an inherited discovery snapshot. */
export function withExistingOpenClawStateDatabaseCurrentReadOnly<T>(
operation: (database: OpenClawStateReadOnlyDatabase) => T,
options: OpenClawStateDatabaseOptions = {},
options: OpenClawStateDatabaseOptions & {
/** Existing host mutation guards may read natively outside worker admission grants. */
allowNativeRead?: true;
} = {},
openStateSchemaReadAdmission?: OpenClawStateSchemaReadAdmission,
): T | undefined {
const pathname = resolveReadOnlyPath(options);
@ -570,7 +573,9 @@ export function withExistingOpenClawStateDatabaseCurrentReadOnly<T>(
return withOpenClawStateReadOnlyLocation(
operation,
pathname,
prepareSqliteReadOnlyLocationSync(pathname),
options.allowNativeRead && !isArtifactPreservingStateRead() && !openStateSchemaReadAdmission
? pathname
: prepareSqliteReadOnlyLocationSync(pathname),
openStateSchemaReadAdmission,
);
});

View file

@ -10,6 +10,10 @@ import type {
WorkspaceAttestation,
WorkspaceAttestationInput,
} from "../agents/workspace-state-store.kernel.js";
import type {
WorkspaceStateGuard,
WorkspaceStateWorkerOperations,
} from "../agents/workspace-state-store.worker-contract.js";
import type { ClawInstallSchemaVersionRow } from "../claws/provenance-runtime-read.kernel.js";
import type { ConfigHealthPatch } from "../config/io.health-state.kernel.js";
import type {
@ -68,6 +72,7 @@ export type OpenClawStateWorkerOpenPreparation = { type: "deviceIdentity"; ident
/** Commands share one physical shared-state actor; bindings belong to commands, not open input. */
export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations &
WorkspaceStateWorkerOperations &
UpdateRunReconciliationOperations &
UpdateRunWriteOperations &
CaptureWorkerOperations &
@ -85,7 +90,7 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations &
"sandboxRegistry.insertIfMissing": { input: SandboxRegistryInsert; output: void };
"sandboxRegistry.write": { input: SandboxRegistryWrite; output: void };
"workspace.replaceAttestation": {
input: WorkspaceAttestationInput;
input: WorkspaceAttestationInput & Pick<WorkspaceStateGuard, "recoveryHoldPredicate">;
output: WorkspaceAttestation;
};
"updateRuns.reconcileInterrupted": {

View file

@ -1,3 +1,4 @@
import { DuplicateAgentError } from "../agents/agent-create-error.js";
import { McpOAuthStoreCorruptionError } from "../agents/mcp-oauth-store-error.js";
import { WorkspaceAliasRepointedError } from "../agents/workspace-state-identity.js";
import { WorkerSessionAlreadyAttachedError } from "../gateway/worker-environments/session-attachment.js";
@ -47,6 +48,7 @@ export type ErrorIdentity =
| "range-error"
| "syntax-error"
| "type-error"
| "duplicate-agent"
| "skill-upload-request"
| "mcp-oauth-corruption";
}
@ -70,6 +72,9 @@ export type ErrorIdentity =
| { type: "agent-media-migration"; pathname: string; schemaVersion: number };
export function identifyError(error: Error): ErrorIdentity {
if (error instanceof DuplicateAgentError) {
return { type: "duplicate-agent" };
}
if (error instanceof WorkerSessionAlreadyAttachedError) {
return {
type: "worker-session-already-attached",
@ -204,6 +209,7 @@ export function parseIdentity(node: Record<string, unknown>): ErrorIdentity | un
case "range-error":
case "syntax-error":
case "type-error":
case "duplicate-agent":
case "skill-upload-request":
case "mcp-oauth-corruption":
return { type: node.type };
@ -265,6 +271,8 @@ function unreachableErrorNode(node: never): never {
export function createError(node: ErrorIdentity & { message: string }): Error {
switch (node.type) {
case "duplicate-agent":
return new DuplicateAgentError(node.message);
case "worker-session-already-attached":
return new WorkerSessionAlreadyAttachedError(node.sessionId, node.environmentId);
case "workspace-alias-repointed":

View file

@ -2,6 +2,7 @@ import { importSandboxRegistryRow } from "../agents/sandbox/registry-import.work
import { writeSandboxRegistry } from "../agents/sandbox/registry-write.worker.js";
import { persistSubagentRunChangesInWorker } from "../agents/subagents/registry/subagent-registry.store.worker.js";
import { replaceWorkspaceAttestationInDatabase } from "../agents/workspace-state-store.kernel.js";
import { executeWorkspaceStateCommand } from "../agents/workspace-state-store.worker.js";
import { readClawInstallSchemaVersionRows } from "../claws/provenance-runtime-read.kernel.js";
import {
patchConfigHealthEntryInDatabase,
@ -31,6 +32,7 @@ import { listWatchedSessionUpstreamLinksInDatabase } from "../sessions/session-u
import { executeSessionUpstreamCommand } from "../sessions/session-upstream-links.worker.js";
import { executeTranscriptRead } from "../transcripts/store-worker-read.js";
import { clearRetiredTuiPointers } from "../tui/tui-last-session.kernel.js";
import { assertAgentDeletionRecoveryHoldPredicate } from "./agent-deletion-journal-recovery.kernel.js";
import {
listAgentProvenanceInDatabase,
readAgentProvenanceBatchInDatabase,
@ -204,9 +206,17 @@ export function executeSharedStateCommand(
requestSqliteWorkerOperationAdmission({ stage: "transaction", facts: undefined });
const result = replaceWorkspaceAttestationInDatabase(writer, command.input);
requestSqliteWorkerOperationAdmission({ stage: "commit", facts: undefined });
assertAgentDeletionRecoveryHoldPredicate(writer, command.input.recoveryHoldPredicate);
return result;
}, writeOptions);
}
if (
command.type === "workspace.snapshotAndRegister" ||
command.type === "workspace.mergeSetup" ||
command.type === "workspace.expire"
) {
return executeWorkspaceStateCommand(command, database, writeOptions);
}
if (command.type === "sandboxRegistry.write") {
return writeSandboxRegistry(command.input, writeOptions);
}

View file

@ -385,16 +385,19 @@ describe("custodian role creation through persisted configuration", () => {
: failure === "post-commit-later"
? ["coordinator", "researcher", "writer"]
: ["coordinator", "researcher"];
if (failure === "post-commit-first" || failure === "post-commit-later") {
const failingAgentId = failure === "post-commit-first" ? "coordinator" : "writer";
const recordProvenance = agentProvenance.recordAgentProvenance;
vi.spyOn(agentProvenance, "recordAgentProvenance").mockImplementation((...args) => {
if (args[0] === failingAgentId) {
throw new Error("provenance unavailable");
}
return recordProvenance(...args);
});
}
let researcherRecorded = false;
const recordProvenance = agentProvenance.recordAgentProvenance;
vi.spyOn(agentProvenance, "recordAgentProvenance").mockImplementation((...args) => {
if (
(failure === "post-commit-first" && args[0] === "coordinator") ||
(failure === "post-commit-later" && args[0] === "writer")
) {
throw new Error("provenance unavailable");
}
const result = recordProvenance(...args);
researcherRecorded ||= args[0] === "researcher";
return result;
});
if (failure === "unfinished-bootstrap") {
const unfinished = await ensureAgentWorkspace({
dir: path.join(workspaceRoot, "writer"),
@ -410,10 +413,7 @@ describe("custodian role creation through persisted configuration", () => {
{
approved: true,
beforePersistentApply: () => {
if (
failure === "authority-revoked" &&
agentProvenance.readAgentProvenance("researcher")
) {
if (failure === "authority-revoked" && researcherRecorded) {
throw new Error("authority closed");
}
},

View file

@ -502,7 +502,7 @@ export async function applySystemAgentSetup(
agentId: effectiveAgentId,
skipBootstrap: Boolean(nextConfig.agents?.defaults?.skipBootstrap),
skipOptionalBootstrapFiles: nextConfig.agents?.defaults?.skipOptionalBootstrapFiles,
beforePersistentApply,
guard: { assertHost: beforePersistentApply },
}),
(error) => lines.push(`Workspace files: ${formatErrorMessage(error)}`),
);

View file

@ -159,8 +159,10 @@ export const databaseWorkerCoreTestFiles = [
"src/agents/agent-tools.safe-bins.test.ts",
"src/agents/agent-tools.workspace-paths.test.ts",
"src/agents/agent-create.integration.test.ts",
"src/agents/agent-create.workspace-worker.test.ts",
"src/agents/sandbox.resolveSandboxContext.test.ts",
"src/agents/workspace-alias-rebind.test.ts",
"src/agents/bootstrap-files.test.ts",
"src/agents/workspace-attestation.worker.test.ts",
"src/agents/workspace-bootstrap-publish.test.ts",
"src/agents/workspace-sqlite-safety.test.ts",
@ -752,6 +754,8 @@ export const databaseWorkerCoreTestFiles = [
"src/agents/context.opencode-go.test.ts",
"src/agents/simple-completion-runtime.selected-model.test.ts",
"src/agents/tools/pdf-tool.resources.test.ts",
"src/talk/agent-consult-runtime.lineage.test.ts",
"src/talk/agent-consult-runtime.storage.test.ts",
"src/tts/tts-summary.static-catalog.test.ts",
"src/tts/tts-summary.selection.test.ts",
"src/agents/prepared-model-catalog.resources.test.ts",

View file

@ -243,6 +243,9 @@ export const gatewayDatabaseWorkerTestFiles = [
"src/gateway/startup-local-cli-pairing.test.ts",
"src/gateway/talk/client-authority.test.ts",
"src/gateway/talk/handlers/client-native-actions.test.ts",
"src/gateway/talk/handlers/client-native-control.test.ts",
"src/gateway/talk/handlers/client.test.ts",
"src/gateway/talk/handlers/native-consult-target.test.ts",
"src/gateway/talk/relay/index.test.ts",
"src/gateway/test-helpers.acquisition.test.ts",
"src/gateway/tool-resolution.cron-capture.test.ts",