From 609fa2d9cfbfc45d1a105d596a34faa58addfd85 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Mon, 21 Sep 2026 00:47:38 -0700 Subject: [PATCH] refactor(skills): prepare pinned library resources off-thread (#148560) Read immutable pin descriptions and eligible manifests through the canonical shared-state reader, retaining admission, snapshot scope, exact artifact bytes, and first-resource error order. Capture pin and catalog facts before awaits so caller mutation cannot mix revisions or change the selected filesystem host. Await the existing per-root node fixture cleanup before deleting test state. No schema, retention, durability, permission, or dependency change is introduced. --- .../database-schemas/storage-changes.md | 7 + src/gateway/node-claude-skill-runtime.test.ts | 8 +- src/skills/library/resource-read.test.ts | 370 ++++++++++++++++++ src/skills/library/selection-read.kernel.ts | 63 +++ src/skills/library/selection-read.ts | 40 ++ src/skills/library/selection.ts | 185 ++++++--- src/skills/library/service.test.ts | 4 +- src/skills/library/store.ts | 32 +- src/skills/runtime/resource-candidates.ts | 10 +- src/skills/runtime/resources.ts | 88 ++++- src/state/openclaw-state-read-worker.test.ts | 30 ++ src/state/openclaw-state-read-worker.ts | 19 + src/state/openclaw-state-read.types.ts | 15 + src/state/openclaw-state-read.worker.ts | 33 ++ .../vitest.database-worker-core-paths.mjs | 2 + 15 files changed, 787 insertions(+), 119 deletions(-) create mode 100644 src/skills/library/resource-read.test.ts create mode 100644 src/skills/library/selection-read.kernel.ts create mode 100644 src/skills/library/selection-read.ts diff --git a/docs/reference/database-schemas/storage-changes.md b/docs/reference/database-schemas/storage-changes.md index 17f84daaa9e3..5cba82d895e0 100644 --- a/docs/reference/database-schemas/storage-changes.md +++ b/docs/reference/database-schemas/storage-changes.md @@ -1440,6 +1440,13 @@ runtimes discard late discovery results. Registration's alias bootstrap, tab mutations, and the final synchronous ownership check before closing a browser target retain their existing owners. +Selected library resources read cold pin descriptions and eligible manifests +through the shared read-only worker. Resource preparation retains its captured +state root and admission through both reads and file preparation, preserving +snapshot scopes, selected revision bytes, hidden-pin omission, and the first +resource failure. Synchronous discovery and borrowed-database readers keep their +existing contracts. This changes no schema, migration, or persistent data. + ### Preserve the data and concurrency contracts Task, flow, and Cron receipt execution identity bindings run in the shared-state diff --git a/src/gateway/node-claude-skill-runtime.test.ts b/src/gateway/node-claude-skill-runtime.test.ts index 9bb1a3914590..f3ed4dbe1333 100644 --- a/src/gateway/node-claude-skill-runtime.test.ts +++ b/src/gateway/node-claude-skill-runtime.test.ts @@ -24,9 +24,8 @@ import { } from "../skills/library/selection.js"; import { listSkillLibrary, readSkillLibrary, saveSkillLibrary } from "../skills/library/service.js"; import { buildSkillSnapshot } from "../skills/loading/workspace-skill-prompt.js"; -import { closeOpenClawAgentDatabasesForTest } from "../state/openclaw-agent-db.js"; -import { closeOpenClawStateDatabaseForTest } from "../state/openclaw-state-db.js"; import { ensureProfileForEmail } from "../state/user-profiles.js"; +import { cleanupSessionStateForTest } from "../test-utils/session-state-cleanup.js"; import { invokeNodeClaudeCliRun } from "./node-agent-cli-runtime.js"; import { NodeRegistry, type NodeRegistryOptions } from "./node-registry.js"; import { @@ -41,8 +40,9 @@ import { createWorkerSessionPlacementStore } from "./worker-environments/placeme const temps = useAutoCleanupTempDirTracker((cleanup) => afterEach(async () => { vi.restoreAllMocks(); - closeOpenClawAgentDatabasesForTest(); - closeOpenClawStateDatabaseForTest(); + for (const stateDir of temps.dirs) { + await cleanupSessionStateForTest({ stateDir }); + } vi.unstubAllEnvs(); cleanup(); }), diff --git a/src/skills/library/resource-read.test.ts b/src/skills/library/resource-read.test.ts new file mode 100644 index 000000000000..b05adb2d6961 --- /dev/null +++ b/src/skills/library/resource-read.test.ts @@ -0,0 +1,370 @@ +import { createHash } from "node:crypto"; +import fs from "node:fs"; +import path from "node:path"; +import { afterEach, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; +import { requireNodeSqlite } from "../../infra/node-sqlite.js"; +import { createDeferredCore } from "../../shared/deferred.js"; +import * as stateReads from "../../state/openclaw-state-db-readonly.js"; +import { tableExists } from "../../state/openclaw-state-db-schema-helpers.js"; +import { + closeOpenClawStateDatabaseAsync, + closeOpenClawStateDatabaseForTest, + openOpenClawStateDatabase, +} from "../../state/openclaw-state-db.js"; +import { ensureProfileForEmail } from "../../state/user-profiles.js"; +import { withEnvAsync } from "../../test-utils/env.js"; +import { materializeSkill } from "../loading/skill-materializer.js"; +import { prepareSkillResourceDelivery } from "../runtime/resources.js"; +import type { SkillSnapshot } from "../types.js"; +import * as bundles from "./bundle.js"; +import { decodeSkillLibraryFile } from "./bundle.js"; +import * as selectionReads from "./selection-read.js"; +import * as selection from "./selection.js"; +import { loadSkillLibrarySelection } from "./selection.js"; +import { saveSkillLibrary } from "./service.js"; + +const dirs = useAutoCleanupTempDirTracker((cleanup) => + afterEach(async () => { + vi.restoreAllMocks(); + await closeOpenClawStateDatabaseAsync(); + closeOpenClawStateDatabaseForTest(); + cleanup(); + }), +); +const content = + "---\nname: guide\ndescription: Saved procedure\n---\n# Guide\nRead references/data.bin.\n"; + +async function fixture() { + const root = dirs.make("skill-library-resource-read-"); + const options = { env: { OPENCLAW_STATE_DIR: root } }; + const profile = ensureProfileForEmail("reader@example.test", options); + const authority = { + profileId: profile.id, + scopes: ["operator.read", "operator.write"], + getConfig: () => ({}), + assertCurrent() {}, + }; + const saved = await saveSkillLibrary( + authority, + { + slug: "guide", + content, + expectedRevision: null, + files: [{ path: "references/data.bin", content: "AP+A", encoding: "base64" }], + }, + options, + ); + const pin = { + skillId: saved.entry.skillId, + revision: saved.entry.revision, + name: saved.entry.name, + ownerProfileId: saved.entry.ownerProfileId, + }; + const snapshot: SkillSnapshot = { + prompt: "", + skills: [{ name: pin.name }], + librarySelections: [pin], + }; + const databasePath = openOpenClawStateDatabase(options).path; + await closeOpenClawStateDatabaseAsync(); + return { root, options, authority, pin, snapshot, databasePath }; +} + +function observeParentSqlite() { + const native = requireNodeSqlite(); + return [ + vi.spyOn(native.DatabaseSync.prototype, "prepare"), + vi.spyOn(native.DatabaseSync.prototype, "exec"), + ...(["get", "all", "run", "iterate"] as const).map((method) => + vi.spyOn(native.StatementSync.prototype, method), + ), + ]; +} + +function familyHashes(databasePath: string) { + return ["", "-wal", "-shm", "-journal"].map((suffix) => { + const file = `${databasePath}${suffix}`; + return fs.existsSync(file) + ? [suffix, createHash("sha256").update(fs.readFileSync(file)).digest("hex")] + : [suffix, null]; + }); +} + +it.each(["ordinary", "artifact-preserving", "disposable"] as const)( + "delivers cold and warm immutable pins without parent SQL under %s reads", + async (mode) => { + const { root, snapshot, pin, databasePath } = await fixture(); + const before = familyHashes(databasePath); + const execute = vi.spyOn(stateReads, "executeExistingOpenClawStateRead"); + const counters = observeParentSqlite(); + const run = () => + withEnvAsync({ OPENCLAW_STATE_DIR: root }, async () => { + const first = await prepareSkillResourceDelivery(snapshot, () => {}); + expect(first?.skills).toHaveLength(1); + const skill = first!.skills[0]!; + expect(skill.revision).toBe(pin.revision); + expect( + decodeSkillLibraryFile(skill.files.find((file) => file.path === "SKILL.md")!).toString( + "utf8", + ), + ).toBe(content); + expect( + decodeSkillLibraryFile(skill.files.find((file) => file.path === "references/data.bin")!), + ).toEqual(Buffer.from([0, 255, 128])); + // The synchronous caller must reuse the same cache key as async preparation. + expect(loadSkillLibrarySelection([pin])[0]?.skill.name).toBe(pin.name); + expect(await prepareSkillResourceDelivery(structuredClone(snapshot), () => {})).toEqual( + first, + ); + await closeOpenClawStateDatabaseAsync(); + }); + try { + if (mode === "artifact-preserving") { + await stateReads.withArtifactPreservingStateReads(run); + } else if (mode === "disposable") { + await stateReads.withArtifactPreservingStateReads(() => + stateReads.withDisposableOpenClawStateReads(databasePath, run), + ); + } else { + await run(); + } + expect(counters.map((counter) => counter.mock.calls.length)).toEqual([0, 0, 0, 0, 0, 0]); + expect(execute.mock.calls.map(([, command]) => command.type)).toEqual([ + "skills.library.descriptions", + "skills.library.manifests", + "skills.library.manifests", + ]); + if (mode === "artifact-preserving") { + expect(familyHashes(databasePath)).toEqual(before); + } + } finally { + for (const counter of counters) { + counter.mockRestore(); + } + } + }, +); + +it("reads both phases from the active snapshot without fetching hidden manifests", async () => { + const { root, options, authority, pin, snapshot } = await fixture(); + const hidden = await saveSkillLibrary( + authority, + { + slug: "hidden", + content: content.replace("guide", "hidden"), + expectedRevision: null, + }, + options, + ); + snapshot.librarySelections!.push({ + skillId: hidden.entry.skillId, + revision: hidden.entry.revision, + name: hidden.entry.name, + ownerProfileId: hidden.entry.ownerProfileId, + }); + await closeOpenClawStateDatabaseAsync(); + const execute = vi.spyOn(stateReads, "executeExistingOpenClawStateRead"); + await withEnvAsync({ OPENCLAW_STATE_DIR: root }, () => + stateReads.withOpenClawStateDatabaseReadSnapshot(async () => { + const database = openOpenClawStateDatabase(options); + database.db + .prepare( + "UPDATE skill_library_revisions SET description = ?, files_json = ? WHERE skill_id = ?", + ) + .run("Later live description", "invalid live manifest", pin.skillId); + const delivery = await prepareSkillResourceDelivery(snapshot, () => {}); + expect(delivery?.skills).toHaveLength(1); + expect(delivery?.skills[0]?.description).toBe("Saved procedure"); + expect(delivery?.skills[0]?.revision).toBe(pin.revision); + const manifests = execute.mock.calls + .map(([, command]) => command) + .find((command) => command.type === "skills.library.manifests"); + expect(manifests?.input).toEqual([{ skillId: pin.skillId, revision: pin.revision }]); + }, options), + ); +}); + +it.each(["root", "authority", "admission"] as const)( + "retains the captured owner when %s changes while descriptions are pending", + async (change) => { + const { root, snapshot } = await fixture(); + const other = dirs.make("skill-library-other-root-"); + const entered = createDeferredCore(); + const gate = createDeferredCore(); + const original = selectionReads.readSkillLibrarySelectionDescriptions; + vi.spyOn(selectionReads, "readSkillLibrarySelectionDescriptions").mockImplementation( + async (...args) => { + const descriptions = await original(...args); + entered.resolve(); + await gate.promise; + return descriptions; + }, + ); + const closed = new Error("Synthetic resource owner closed"); + let current = true; + await withEnvAsync({ OPENCLAW_STATE_DIR: root }, async () => { + const pending = prepareSkillResourceDelivery(snapshot, () => { + if (!current) { + throw closed; + } + }); + const outcome = pending.then( + (value) => ({ value }), + (error: unknown) => ({ error }), + ); + await entered.promise; + if (change === "root") { + process.env.OPENCLAW_STATE_DIR = other; + } else if (change === "admission") { + await closeOpenClawStateDatabaseAsync(); + } else { + current = false; + } + gate.resolve(); + const result = await outcome; + if (change === "authority") { + expect(result).toEqual({ error: closed }); + } else if (change === "admission") { + expect(result).toHaveProperty("error"); + } else { + expect(result).toMatchObject({ value: { skills: [{ description: "Saved procedure" }] } }); + expect(fs.existsSync(path.join(other, "state", "openclaw.sqlite"))).toBe(false); + } + }); + }, +); + +it.each(["selection", "descriptions", "manifests"] as const)( + "retains pin values and eligibility when the caller mutates them during %s", + async (phase) => { + const { root, pin, snapshot } = await fixture(); + const originalPin = { ...pin }; + const entered = createDeferredCore(); + const gate = createDeferredCore(); + const descriptions = selectionReads.readSkillLibrarySelectionDescriptions; + const manifests = selectionReads.readSkillLibrarySelectionManifests; + async function hold(value: T) { + entered.resolve(); + await gate.promise; + return value; + } + if (phase === "manifests") { + vi.spyOn(selectionReads, "readSkillLibrarySelectionManifests").mockImplementation( + async (...args) => hold(await manifests(...args)), + ); + } else { + vi.spyOn(selectionReads, "readSkillLibrarySelectionDescriptions").mockImplementation( + async (...args) => hold(await descriptions(...args)), + ); + } + await withEnvAsync({ OPENCLAW_STATE_DIR: root }, async () => { + const pending = + phase === "selection" + ? selection + .prepareSkillLibrarySelection(snapshot.librarySelections!, {}, () => {}) + .then((entries) => entries.map((entry) => entry.skill)) + : prepareSkillResourceDelivery(snapshot, () => {}).then((delivery) => delivery?.skills); + const outcome = pending.then( + (value) => ({ value }), + (error: unknown) => ({ error }), + ); + try { + await entered.promise; + Object.assign(pin, { + skillId: "changed", + revision: "0".repeat(64), + name: "changed", + ownerProfileId: null, + }); + snapshot.librarySelections!.push({ ...pin }); + snapshot.skills[0]!.name = "changed"; + gate.resolve(); + expect(await outcome).toMatchObject({ + value: [{ name: originalPin.name, description: "Saved procedure" }], + }); + expect(loadSkillLibrarySelection([originalPin])[0]?.skill.name).toBe(originalPin.name); + } finally { + gate.resolve(); + await outcome; + } + }); + }, +); + +it.each(["file", "table"] as const)( + "does not create missing library %s storage", + async (missing) => { + const root = dirs.make("skill-library-missing-read-"); + const options = { env: { OPENCLAW_STATE_DIR: root } }; + if (missing === "table") { + openOpenClawStateDatabase(options); + await closeOpenClawStateDatabaseAsync(); + } + const snapshot: SkillSnapshot = { + prompt: "", + skills: [{ name: "missing" }], + librarySelections: [ + { skillId: "missing", revision: "0".repeat(64), name: "missing", ownerProfileId: null }, + ], + }; + await withEnvAsync({ OPENCLAW_STATE_DIR: root }, async () => { + await expect(prepareSkillResourceDelivery(snapshot, () => {})).rejects.toMatchObject({ + code: "NOT_FOUND", + }); + }); + if (missing === "file") { + expect(fs.existsSync(path.join(root, "state", "openclaw.sqlite"))).toBe(false); + } else { + expect(tableExists(openOpenClawStateDatabase(options).db, "skill_library_entries")).toBe( + false, + ); + } + }, +); + +it.each(["workspace", "library"] as const)( + "preserves the first %s failure and resource context before reading later candidates", + async (first) => { + const root = dirs.make("skill-library-read-order-"); + const skills = ["workspace", "library"].map((name) => + materializeSkill({ + name, + description: name, + content, + frontmatter: {}, + filePath: path.join(root, name, "SKILL.md"), + baseDir: path.join(root, name), + source: "openclaw-workspace", + sourceOptions: { source: "openclaw-workspace" }, + }), + ); + const [workspace, library] = skills; + const snapshot: SkillSnapshot = { + prompt: "", + skills: skills.map(({ name }) => ({ name })), + resolvedSkills: first === "workspace" ? skills : [library!, workspace!], + librarySelections: [ + { skillId: "library", revision: "0".repeat(64), name: "library", ownerProfileId: null }, + ], + }; + const workspaceError = new Error("Synthetic workspace read failed"); + const libraryError = new Error("Synthetic database read failed"); + vi.spyOn(selection, "prepareSkillLibrarySelection").mockResolvedValue([]); + const readWorkspace = vi + .spyOn(bundles, "readSkillBundleTree") + .mockRejectedValue(workspaceError); + const readManifests = vi + .spyOn(selectionReads, "readSkillLibrarySelectionManifests") + .mockRejectedValue(libraryError); + await withEnvAsync({ OPENCLAW_STATE_DIR: root }, async () => { + await expect(prepareSkillResourceDelivery(snapshot, () => {})).rejects.toMatchObject({ + code: "INVALID_BUNDLE", + message: expect.stringContaining(`skill="${first}"`), + cause: first === "workspace" ? workspaceError : libraryError, + }); + }); + expect(readWorkspace).toHaveBeenCalledTimes(first === "workspace" ? 1 : 0); + expect(readManifests).toHaveBeenCalledTimes(first === "library" ? 1 : 0); + }, +); diff --git a/src/skills/library/selection-read.kernel.ts b/src/skills/library/selection-read.kernel.ts new file mode 100644 index 000000000000..f03a9e65bcb6 --- /dev/null +++ b/src/skills/library/selection-read.kernel.ts @@ -0,0 +1,63 @@ +import type { DatabaseSync } from "node:sqlite"; +import type { SkillLibrarySelection } from "../../../packages/gateway-protocol/src/schema/skill-library.js"; +import { executeSqliteQuerySync, getNodeSqliteKysely } from "../../infra/kysely-sync.js"; +import type { DB as StateDatabase } from "../../state/openclaw-state-db.generated.js"; + +type RevisionSelection = Pick; +export type SkillLibraryReadOnlyOperations = { + "skills.library.descriptions": { + input: readonly RevisionSelection[]; + output: Array<{ description: string } | undefined> | undefined; + }; + "skills.library.manifests": { + input: readonly RevisionSelection[]; + output: Array<{ files_json: string } | undefined> | undefined; + }; +}; + +function revisionBatchQuery(db: DatabaseSync, selections: readonly RevisionSelection[]) { + return getNodeSqliteKysely>(db) + .selectFrom("skill_library_revisions") + .where((eb) => + eb.or( + selections.map((pin) => + eb.and([eb("skill_id", "=", pin.skillId), eb("revision", "=", pin.revision)]), + ), + ), + ); +} + +/** Resolve a bounded session selection in its original order, including repeated pins. */ +export function selectSkillLibraryRevisionMetadataBatch( + db: DatabaseSync, + selections: readonly RevisionSelection[], +) { + const rows = executeSqliteQuerySync( + db, + revisionBatchQuery(db, selections).select(["skill_id", "revision", "description"]), + ).rows; + const metadata = new Map( + rows.map((row) => [ + JSON.stringify([row.skill_id, row.revision]), + { description: row.description }, + ]), + ); + return selections.map((pin) => metadata.get(JSON.stringify([pin.skillId, pin.revision]))); +} + +export function selectSkillLibraryRevisionManifestsBatch( + db: DatabaseSync, + selections: readonly RevisionSelection[], +) { + const rows = executeSqliteQuerySync( + db, + revisionBatchQuery(db, selections).select(["skill_id", "revision", "files_json"]), + ).rows; + const manifests = new Map( + rows.map((row) => [ + JSON.stringify([row.skill_id, row.revision]), + { files_json: row.files_json }, + ]), + ); + return selections.map((pin) => manifests.get(JSON.stringify([pin.skillId, pin.revision]))); +} diff --git a/src/skills/library/selection-read.ts b/src/skills/library/selection-read.ts new file mode 100644 index 000000000000..0d683d905b57 --- /dev/null +++ b/src/skills/library/selection-read.ts @@ -0,0 +1,40 @@ +import type { SkillLibrarySelection } from "../../../packages/gateway-protocol/src/schema/skill-library.js"; +import { executeExistingOpenClawStateRead } from "../../state/openclaw-state-db-readonly.js"; +import type { OpenClawStateDatabaseOptions } from "../../state/openclaw-state-db.js"; + +export async function readSkillLibrarySelectionDescriptions( + selections: readonly Pick[], + options: Pick, +) { + const result = await executeExistingOpenClawStateRead(options, { + type: "skills.library.descriptions", + input: selections.map(({ skillId, revision }) => ({ skillId, revision })), + }); + if (result === undefined) { + return undefined; + } + if (result.ok && result.type === "skills.library.descriptions") { + return result.value; + } + throw new Error("Unexpected skill library descriptions result"); +} + +export async function readSkillLibrarySelectionManifests( + selections: readonly Pick[], + options: Pick, +) { + if (!selections.length) { + return []; + } + const result = await executeExistingOpenClawStateRead(options, { + type: "skills.library.manifests", + input: selections.map(({ skillId, revision }) => ({ skillId, revision })), + }); + if (result === undefined) { + return undefined; + } + if (result.ok && result.type === "skills.library.manifests") { + return result.value; + } + throw new Error("Unexpected skill library manifests result"); +} diff --git a/src/skills/library/selection.ts b/src/skills/library/selection.ts index 8f7d9104dcc4..683e1c4486ca 100644 --- a/src/skills/library/selection.ts +++ b/src/skills/library/selection.ts @@ -19,6 +19,8 @@ import { materializeSkill } from "../loading/skill-materializer.js"; import type { SkillEntry } from "../types.js"; import { readSkillLibraryManifestTree, skillLibraryRevisionDir } from "./bundle.js"; import { SkillLibraryError } from "./errors.js"; +import { readSkillLibrarySelectionDescriptions } from "./selection-read.js"; +import { selectSkillLibraryRevisionMetadataBatch } from "./selection-read.kernel.js"; import { projectSkillLibraryEntry, readSkillLibraryStore, @@ -26,7 +28,6 @@ import { resolveSkillLibraryActor, selectSkillLibraryRevision, selectSkillLibraryRevisionMetadata, - selectSkillLibraryRevisionMetadataBatch, skillLibraryDb, type SkillLibraryAuthority, } from "./store.js"; @@ -43,6 +44,16 @@ export function assertPreparedSkillLibrarySelection( const selectedEntryCache = new Map(); +/** Async preparation retains pin values even when a caller later edits its snapshot. */ +export function captureSkillLibrarySelection(selections: readonly SkillLibrarySelection[]) { + return selections.map(({ skillId, revision, name, ownerProfileId }) => ({ + skillId, + revision, + name, + ownerProfileId, + })); +} + /** The session owner has already authorized this exact immutable pin. */ export async function readSelectedSkillLibraryFiles( selection: SkillLibrarySelection, @@ -208,71 +219,115 @@ export function loadSkillLibrarySelection( if (selections.length > SKILL_LIBRARY_MAX_SELECTIONS) { throw new SkillLibraryError("LIMIT", "Invalid session skill selection."); } - const entries = readSkillLibraryStore((db) => { - const revisions = selectSkillLibraryRevisionMetadataBatch(db, selections); - return selections.map((selection, index) => { - const revision = revisions[index]; - if (!revision) { - throw new SkillLibraryError( - "NOT_FOUND", - "A pinned skill revision is unavailable; restore the library artifact or detach it explicitly.", - ); - } - const baseDir = skillLibraryRevisionDir(selection.skillId, selection.revision, options.env); - const filePath = path.join(baseDir, "SKILL.md"); - const opened = openRootFileSync({ - absolutePath: filePath, - rootPath: baseDir, - boundaryLabel: "skill library revision", - maxBytes: SKILL_LIBRARY_MAX_FILE_BYTES, - rejectHardlinks: true, - symlinks: "reject", - }); - if (!opened.ok) { - throw new SkillLibraryError( - "INVALID_BUNDLE", - "Pinned skill instructions could not be read; restore the library artifact or detach it explicitly.", - undefined, - { cause: opened.error }, - ); - } - let content: string; - try { - content = readFileDescriptorBoundedSync(opened.fd, SKILL_LIBRARY_MAX_FILE_BYTES).toString( - "utf8", - ); - } finally { - fs.closeSync(opened.fd); - } - const frontmatter = parseSkillFrontmatter(content); - const metadata = resolveSkillManifestMetadata(frontmatter); - const invocation = resolveSkillInvocationPolicy(frontmatter); - const name = selection.name; - return { - skill: materializeSkill({ - content, - frontmatter, - name, - description: revision.description, - baseDir, - filePath, - source: "openclaw-library", - sourceOptions: { source: "openclaw-library" }, - }), - frontmatter, - invocation, - // Untrusted frontmatter can constrain executable eligibility, but cannot claim global credentials/config. - metadata: { - skillKey: name, - os: metadata?.os, - requires: metadata?.requires, - }, - disableCommandDispatch: true, - syncSourceDir: baseDir, - syncDirName: `library-${selection.skillId}-${selection.revision}`, - } satisfies SkillEntry; + const entries = readSkillLibraryStore( + (db) => + materializeSkillLibrarySelection( + selections, + selectSkillLibraryRevisionMetadataBatch(db, selections), + options.env, + ), + options, + ); + return cacheSkillLibrarySelection(cacheKey, entries); +} + +/** Prepare immutable pins without borrowing a caller's synchronous database handle. */ +export async function prepareSkillLibrarySelection( + inputSelections: readonly SkillLibrarySelection[], + options: Pick, + assertCurrent: () => void, +): Promise { + assertCurrent(); + if (!inputSelections.length) { + return []; + } + const selections = captureSkillLibrarySelection(inputSelections); + const cacheKey = JSON.stringify([resolveStateDir(options.env), options.path, selections]); + const cached = selectedEntryCache.get(cacheKey); + if (cached) { + return [...cached]; + } + if (selections.length > SKILL_LIBRARY_MAX_SELECTIONS) { + throw new SkillLibraryError("LIMIT", "Invalid session skill selection."); + } + const revisions = await readSkillLibrarySelectionDescriptions(selections, options); + assertCurrent(); + return cacheSkillLibrarySelection( + cacheKey, + revisions && materializeSkillLibrarySelection(selections, revisions, options.env), + ); +} + +function materializeSkillLibrarySelection( + selections: readonly SkillLibrarySelection[], + revisions: ReturnType, + env: NodeJS.ProcessEnv | undefined, +): SkillEntry[] { + return selections.map((selection, index) => { + const revision = revisions[index]; + if (!revision) { + throw new SkillLibraryError( + "NOT_FOUND", + "A pinned skill revision is unavailable; restore the library artifact or detach it explicitly.", + ); + } + const baseDir = skillLibraryRevisionDir(selection.skillId, selection.revision, env); + const filePath = path.join(baseDir, "SKILL.md"); + const opened = openRootFileSync({ + absolutePath: filePath, + rootPath: baseDir, + boundaryLabel: "skill library revision", + maxBytes: SKILL_LIBRARY_MAX_FILE_BYTES, + rejectHardlinks: true, + symlinks: "reject", }); - }, options); + if (!opened.ok) { + throw new SkillLibraryError( + "INVALID_BUNDLE", + "Pinned skill instructions could not be read; restore the library artifact or detach it explicitly.", + undefined, + { cause: opened.error }, + ); + } + let content: string; + try { + content = readFileDescriptorBoundedSync(opened.fd, SKILL_LIBRARY_MAX_FILE_BYTES).toString( + "utf8", + ); + } finally { + fs.closeSync(opened.fd); + } + const frontmatter = parseSkillFrontmatter(content); + const metadata = resolveSkillManifestMetadata(frontmatter); + const invocation = resolveSkillInvocationPolicy(frontmatter); + const name = selection.name; + return { + skill: materializeSkill({ + content, + frontmatter, + name, + description: revision.description, + baseDir, + filePath, + source: "openclaw-library", + sourceOptions: { source: "openclaw-library" }, + }), + frontmatter, + invocation, + // Untrusted frontmatter can constrain executable eligibility, but cannot claim global credentials/config. + metadata: { + skillKey: name, + os: metadata?.os, + requires: metadata?.requires, + }, + disableCommandDispatch: true, + syncSourceDir: baseDir, + syncDirName: `library-${selection.skillId}-${selection.revision}`, + } satisfies SkillEntry; + }); +} + +function cacheSkillLibrarySelection(cacheKey: string, entries: SkillEntry[] | undefined) { if (!entries) { throw new SkillLibraryError( "NOT_FOUND", diff --git a/src/skills/library/service.test.ts b/src/skills/library/service.test.ts index be6a7c1a0e51..016acc2be754 100644 --- a/src/skills/library/service.test.ts +++ b/src/skills/library/service.test.ts @@ -10,6 +10,7 @@ import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js" import { declareAgentWorkspaceAccess } from "../../agents/workspace-access.js"; import { tableExists, tableHasColumn } from "../../state/openclaw-state-db-schema-helpers.js"; import { + closeOpenClawStateDatabaseAsync, closeOpenClawStateDatabaseForTest, openOpenClawStateDatabase, } from "../../state/openclaw-state-db.js"; @@ -39,7 +40,8 @@ import { import type { SkillLibraryAuthority } from "./store.js"; const tempDirs = useAutoCleanupTempDirTracker((cleanup) => - afterEach(() => { + afterEach(async () => { + await closeOpenClawStateDatabaseAsync(); closeOpenClawStateDatabaseForTest(); cleanup(); }), diff --git a/src/skills/library/store.ts b/src/skills/library/store.ts index c0ee8b6995f4..2ddcb979dd67 100644 --- a/src/skills/library/store.ts +++ b/src/skills/library/store.ts @@ -1,10 +1,7 @@ import { randomUUID } from "node:crypto"; import type { DatabaseSync } from "node:sqlite"; import type { SelectQueryBuilder } from "kysely"; -import type { - SkillLibraryEntry, - SkillLibrarySelection, -} from "../../../packages/gateway-protocol/src/schema/skill-library.js"; +import type { SkillLibraryEntry } from "../../../packages/gateway-protocol/src/schema/skill-library.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { authorizeOperatorScopesForRequiredScope } from "../../gateway/method-scopes.js"; import { resolveOperatorRolePolicyForAssignment } from "../../gateway/operator-role-policy.js"; @@ -225,33 +222,6 @@ export function selectSkillLibraryRevisionMetadata( ); } -/** Resolve a bounded session selection in its original order, including repeated pins. */ -export function selectSkillLibraryRevisionMetadataBatch( - db: DatabaseSync, - selections: readonly Pick[], -) { - const rows = executeSqliteQuerySync( - db, - skillLibraryDb(db) - .selectFrom("skill_library_revisions") - .select(["skill_id", "revision", "description"]) - .where((eb) => - eb.or( - selections.map((pin) => - eb.and([eb("skill_id", "=", pin.skillId), eb("revision", "=", pin.revision)]), - ), - ), - ), - ).rows; - const metadata = new Map( - rows.map((row) => [ - JSON.stringify([row.skill_id, row.revision]), - { description: row.description }, - ]), - ); - return selections.map((pin) => metadata.get(JSON.stringify([pin.skillId, pin.revision]))); -} - export function selectSkillLibraryOwner(db: DatabaseSync, profileId: string) { // Actor resolution stays separate because existing profile tables may omit its optional role. return selectResolvedUserProfile( diff --git a/src/skills/runtime/resource-candidates.ts b/src/skills/runtime/resource-candidates.ts index c682675d98bf..8368f3e4b94b 100644 --- a/src/skills/runtime/resource-candidates.ts +++ b/src/skills/runtime/resource-candidates.ts @@ -1,13 +1,17 @@ import { loadSkillLibrarySelection } from "../library/selection.js"; -import type { SkillSnapshot } from "../types.js"; +import type { SkillSnapshot, SkillEntry } from "../types.js"; /** Resource access includes eligible immutable pins, not just the prompt projection. */ -export function resolveSkillResourceCandidates(snapshot: SkillSnapshot | undefined) { +export function resolveSkillResourceCandidates( + snapshot: SkillSnapshot | undefined, + libraryEntries?: readonly SkillEntry[], +) { if (!snapshot) { return undefined; } const candidates = [...(snapshot.resolvedSkills ?? [])]; - for (const entry of loadSkillLibrarySelection(snapshot.librarySelections ?? [])) { + for (const entry of libraryEntries ?? + loadSkillLibrarySelection(snapshot.librarySelections ?? [])) { if ( snapshot.skills.some((skill) => skill.name === entry.skill.name) && !candidates.some((skill) => skill.name === entry.skill.name) diff --git a/src/skills/runtime/resources.ts b/src/skills/runtime/resources.ts index 166d44362839..e8c014f7d85d 100644 --- a/src/skills/runtime/resources.ts +++ b/src/skills/runtime/resources.ts @@ -18,13 +18,20 @@ import { import { isMissingPathError } from "../../infra/errors.js"; import { removeTemporaryArtifacts } from "../../infra/temp-artifact-cleanup.js"; import { createSubsystemLogger } from "../../logging/subsystem.js"; +import { captureOpenClawStateWorkerContext } from "../../state/openclaw-state-worker-context.js"; import { prepareSkillBundle, readSkillBundleTree, + readSkillLibraryManifestTree, + skillLibraryRevisionDir, SkillTreeDirectoryError, } from "../library/bundle.js"; import { SkillLibraryError } from "../library/errors.js"; -import { readSelectedSkillLibraryFiles } from "../library/selection.js"; +import { readSkillLibrarySelectionManifests } from "../library/selection-read.js"; +import { + captureSkillLibrarySelection, + prepareSkillLibrarySelection, +} from "../library/selection.js"; import { loadSingleSkillDirectory } from "../loading/local-loader.js"; import { createSyntheticSourceInfo, type Skill } from "../loading/skill-contract.js"; import { shouldSyncSkillPath } from "../loading/skill-paths.js"; @@ -103,22 +110,30 @@ const localSkillResourceReader: SkillResourceSourceReader = { // The caller retains these bytes for its turn. Catalog versions do not version supporting files. export async function prepareSkillResourceDelivery( - snapshot: SkillSnapshot | undefined, + inputSnapshot: SkillSnapshot | undefined, assertCurrent: () => void, - explicitSelections: readonly ExplicitSkillSelection[] = [], + inputExplicitSelections: readonly ExplicitSkillSelection[] = [], workspaceDir?: string, ): Promise { - if (!snapshot) { + if (!inputSnapshot) { return undefined; } assertCurrent(); if ( - !snapshot.resolvedSkills?.length && - !snapshot.librarySelections?.length && - !explicitSelections.length + !inputSnapshot.resolvedSkills?.length && + !inputSnapshot.librarySelections?.length && + !inputExplicitSelections.length ) { return undefined; } + const snapshot = { + ...inputSnapshot, + librarySelections: captureSkillLibrarySelection(inputSnapshot.librarySelections ?? []), + skills: inputSnapshot.skills.map((skill) => ({ ...skill })), + resolvedSkills: inputSnapshot.resolvedSkills?.map((skill) => ({ ...skill })), + skillRoots: inputSnapshot.skillRoots && { ...inputSnapshot.skillRoots }, + }; + const explicitSelections = inputExplicitSelections.map((selection) => ({ ...selection })); // Library-only and node-native catalogs need no workspace filesystem access. let sourceReader: SkillResourceSourceReader | undefined; const getSourceReader = (skill?: Skill) => { @@ -142,7 +157,24 @@ export async function prepareSkillResourceDelivery( }; const skills: SkillResourceDelivery["skills"] = []; let total = 0; - const candidates = resolveSkillResourceCandidates(snapshot)!; + const libraryContext = snapshot.librarySelections?.length + ? captureOpenClawStateWorkerContext() + : undefined; + const assertLibraryCurrent = () => { + assertCurrent(); + libraryContext?.maintenanceScope?.assertAdmission(); + libraryContext?.admission.assertCurrent(); + }; + const libraryOptions = { env: libraryContext?.environment }; + const libraryEntries = libraryContext + ? await prepareSkillLibrarySelection( + snapshot.librarySelections ?? [], + libraryOptions, + assertLibraryCurrent, + ) + : []; + assertLibraryCurrent(); + const candidates = resolveSkillResourceCandidates(snapshot, libraryEntries)!; for (const selected of explicitSelections) { if ( selected.path.startsWith("node://") || @@ -165,7 +197,7 @@ export async function prepareSkillResourceDelivery( } catch (error) { throw contextualizeSkillResourceError({ name: selected.name, baseDir: skillDir }, error); } - assertCurrent(); + assertLibraryCurrent(); if ( !loaded || loaded.filePath !== selected.path || @@ -180,6 +212,17 @@ export async function prepareSkillResourceDelivery( } candidates.push(loaded); } + const pins = [ + ...new Set( + candidates.flatMap((skill) => { + const pin = + !skill.filePath.startsWith("node://") && + snapshot.librarySelections?.find((selection) => selection.name === skill.name); + return pin ? [pin] : []; + }), + ), + ]; + let manifests: Awaited>; for (const skill of candidates) { if (skill.filePath.startsWith("node://")) { continue; @@ -190,15 +233,30 @@ export async function prepareSkillResourceDelivery( ); let files: SkillLibraryFile[] | null; try { - files = pin - ? await readSelectedSkillLibraryFiles(pin) - : await getSourceReader(skill).readSkillFiles(skill, { - allowMissingRoot: !explicitlySelected, - }); + if (pin) { + assertLibraryCurrent(); + if (!manifests) { + manifests = await readSkillLibrarySelectionManifests(pins, libraryOptions); + assertLibraryCurrent(); + } + const manifest = manifests?.[pins.indexOf(pin)]; + if (!manifest) { + throw new SkillLibraryError("NOT_FOUND", "Selected skill revision is unavailable."); + } + files = await readSkillLibraryManifestTree( + skillLibraryRevisionDir(pin.skillId, pin.revision, libraryContext?.environment), + manifest.files_json, + pin.revision, + ); + } else { + files = await getSourceReader(skill).readSkillFiles(skill, { + allowMissingRoot: !explicitlySelected, + }); + } } catch (error) { throw contextualizeSkillResourceError(skill, error); } - assertCurrent(); + assertLibraryCurrent(); if (files === null) { if (explicitlySelected) { throw contextualizeSkillResourceError( diff --git a/src/state/openclaw-state-read-worker.test.ts b/src/state/openclaw-state-read-worker.test.ts index b4b0efb8319b..03ed7033bf23 100644 --- a/src/state/openclaw-state-read-worker.test.ts +++ b/src/state/openclaw-state-read-worker.test.ts @@ -737,6 +737,36 @@ it("captures update history filters and charges retained selectors before dispat } }); +it.each(["skills.library.descriptions", "skills.library.manifests"] as const)( + "retains library pins and their byte charge while dispatch waits (%s)", + async (type) => { + const { options } = source(); + const input = [{ skillId: "技能🦞".repeat(512), revision: "版本🦞".repeat(512) }]; + const expected = structuredClone(input); + const dispatch = createDeferredCore(); + const task = queueTask(dispatch.promise); + const result = executeExistingOpenClawStateRead(options, { type, input }); + const returned: OpenClawStateReadReply = { ok: true, type, sourceAdmitted: true, value: [] }; + try { + const submitted = await task.submitted; + input[0]!.skillId = "changed"; + input[0]!.revision = "changed"; + input.push({ skillId: "extra", revision: "extra" }); + expect(submitted.inputBytes).toBeGreaterThanOrEqual( + Buffer.byteLength(expected[0]!.skillId) + Buffer.byteLength(expected[0]!.revision), + ); + dispatch.resolve(); + expect((await task.captured).command).toEqual({ type, input: expected }); + task.result.resolve(returned); + expect(await result).toEqual(returned); + } finally { + dispatch.resolve(); + task.result.resolve(returned); + await Promise.allSettled([result]); + } + }, +); + it.each([ { input: { diff --git a/src/state/openclaw-state-read-worker.ts b/src/state/openclaw-state-read-worker.ts index 99684d2a8063..f31d111c0e9c 100644 --- a/src/state/openclaw-state-read-worker.ts +++ b/src/state/openclaw-state-read-worker.ts @@ -96,6 +96,15 @@ function captureCommand(command: OpenClawStateReadCommand): OpenClawStateReadCom if (command.type === "updateRuns.list") { return { ...command, input: { ...command.input } }; } + if ( + command.type === "skills.library.descriptions" || + command.type === "skills.library.manifests" + ) { + return { + type: command.type, + input: command.input.map(({ skillId, revision }) => ({ skillId, revision })), + }; + } if (command.type === "audit.run.inspect") { const input = command.input; const common = { @@ -151,6 +160,16 @@ function commandBytes(command: OpenClawStateReadRequest["command"]): number { (command.input.active === undefined ? 0 : 1) ); } + if ( + command.type === "skills.library.descriptions" || + command.type === "skills.library.manifests" + ) { + return command.input.reduce( + (total, pin) => + total + Buffer.byteLength(pin.skillId, "utf8") + Buffer.byteLength(pin.revision, "utf8"), + bytes, + ); + } if (command.type === "fleet.get") { return bytes + Buffer.byteLength(command.tenantId, "utf8"); } diff --git a/src/state/openclaw-state-read.types.ts b/src/state/openclaw-state-read.types.ts index e50e29b927e0..0dbbf808d8b4 100644 --- a/src/state/openclaw-state-read.types.ts +++ b/src/state/openclaw-state-read.types.ts @@ -22,6 +22,7 @@ import type { PluginBlobReadReply, } from "../plugin-state/plugin-blob-worker-contract.js"; import type { AsyncWorkScope } from "../shared/async-work-scope.js"; +import type { SkillLibraryReadOnlyOperations } from "../skills/library/selection-read.kernel.js"; import type { OnboardingRecommendationsRecord } from "./onboarding-recommendations.contract.js"; import type { OpenClawAgentDatabaseRegistryReadResult } from "./openclaw-agent-db-contract.js"; import type { ConfigMachineState } from "./openclaw-state-db.generated.js"; @@ -45,6 +46,12 @@ export type OpenClawStateReadAuthority = { export type OpenClawStateReadCommand = | PluginBlobReadCommand | { type: "exec-approvals.read" } + | { + [Kind in keyof SkillLibraryReadOnlyOperations]: { + type: Kind; + input: SkillLibraryReadOnlyOperations[Kind]["input"]; + }; + }[keyof SkillLibraryReadOnlyOperations] | { type: "agentDatabaseRegistry.read" } | { type: "onboardingRecommendations.read"; configKey: string } | { type: "userProfiles.avatar.reconcile"; profileId: string } @@ -70,6 +77,14 @@ export type OpenClawStateReadRequest = { }; export type OpenClawStateReadReply = ( | PluginBlobReadReply + | { + [Kind in keyof SkillLibraryReadOnlyOperations]: { + ok: true; + type: Kind; + sourceAdmitted: true; + value: SkillLibraryReadOnlyOperations[Kind]["output"]; + }; + }[keyof SkillLibraryReadOnlyOperations] | { ok: true; type: "agentDatabaseRegistry.read"; diff --git a/src/state/openclaw-state-read.worker.ts b/src/state/openclaw-state-read.worker.ts index 4fb243188568..3c0d5e275832 100644 --- a/src/state/openclaw-state-read.worker.ts +++ b/src/state/openclaw-state-read.worker.ts @@ -1,5 +1,6 @@ import { toStringifiedError } from "@openclaw/normalization-core/error-coercion"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { SKILL_LIBRARY_MAX_SELECTIONS } from "../../packages/gateway-protocol/src/schema/skill-library.js"; import { readSandboxBrowserRegistryInDatabase, readSandboxRegistryEntryInDatabase, @@ -21,6 +22,10 @@ import { pluginBlobEntriesInDatabase, } from "../plugin-state/plugin-blob-store.sqlite.js"; import { isPluginBlobReadCommand } from "../plugin-state/plugin-blob-worker-contract.js"; +import { + selectSkillLibraryRevisionMetadataBatch, + selectSkillLibraryRevisionManifestsBatch, +} from "../skills/library/selection-read.kernel.js"; import { readConfigMachineStateRowInDatabase } from "./config-machine-state.js"; import { readOnboardingRecommendationsInDatabase } from "./onboarding-recommendations.kernel.js"; import { readRegisteredAgentDatabaseRows } from "./openclaw-agent-db-registry.read.js"; @@ -61,6 +66,14 @@ function isReadRequest(input: unknown): input is OpenClawStateReadRequest { (isPluginBlobReadCommand(input.command) || input.command.type === "admit" || input.command.type === "exec-approvals.read" || + ((input.command.type === "skills.library.descriptions" || + input.command.type === "skills.library.manifests") && + Array.isArray(input.command.input) && + input.command.input.length <= SKILL_LIBRARY_MAX_SELECTIONS && + input.command.input.every( + (pin) => + isRecord(pin) && typeof pin.skillId === "string" && typeof pin.revision === "string", + )) || input.command.type === "agentDatabaseRegistry.read" || (input.command.type === "userProfiles.avatar.reconcile" && typeof input.command.profileId === "string") || @@ -200,6 +213,26 @@ serveOwnedWorkerTasks( row: readExecApprovalsConfigRow(db), }; } + if (command.type === "skills.library.descriptions") { + return { + ok: true, + type: command.type, + sourceAdmitted, + value: tableExists(db, "skill_library_entries") + ? selectSkillLibraryRevisionMetadataBatch(db, command.input) + : undefined, + }; + } + if (command.type === "skills.library.manifests") { + return { + ok: true, + type: command.type, + sourceAdmitted, + value: tableExists(db, "skill_library_entries") + ? selectSkillLibraryRevisionManifestsBatch(db, command.input) + : undefined, + }; + } if (command.type === "onboardingRecommendations.read") { return { ok: true, diff --git a/test/vitest/vitest.database-worker-core-paths.mjs b/test/vitest/vitest.database-worker-core-paths.mjs index 8c048a3603ce..5cf0c66aaed0 100644 --- a/test/vitest/vitest.database-worker-core-paths.mjs +++ b/test/vitest/vitest.database-worker-core-paths.mjs @@ -124,6 +124,8 @@ export const databaseWorkerCoreTestFiles = [ "src/hooks/hooks-install.test.ts", "src/infra/state-migrations.audit-logs.windows.test.ts", "src/infra/state-migrations.media-persistence.lifecycle-recovery.test.ts", + "src/skills/library/resource-read.test.ts", + "src/skills/library/service.test.ts", "src/skills/workshop/collection-restore.test.ts", "src/skills/workshop/experience-review.apply.test.ts", "src/skills/workshop/policy.test.ts",