diff --git a/server/src/services/company-skills.ts b/server/src/services/company-skills.ts index 4b432dcd4f..25ecb3bc1f 100644 --- a/server/src/services/company-skills.ts +++ b/server/src/services/company-skills.ts @@ -1,3 +1,5 @@ +import { logger } from "../middleware/logger.js"; +import { removeRuntimeSkillCache, resolveRuntimeSkillCache, runtimeSkillCacheSpec } from "./runtime-skill-cache.js"; import { createHash, randomUUID } from "node:crypto"; import { promises as fs } from "node:fs"; import path from "node:path"; @@ -3149,10 +3151,10 @@ export function companySkillService(db: Db) { continue; } - await db - .delete(companySkills) - .where(eq(companySkills.id, skill.id)); - await fs.rm(resolveRuntimeSkillMaterializedPath(companyId, skill), { recursive: true, force: true }); + await removeRuntimeSkillCache(resolveManagedSkillsRoot(companyId), skill.id, async () => { + await fs.rm(resolveRuntimeSkillMaterializedPath(companyId, skill), { recursive: true, force: true }); + await db.delete(companySkills).where(eq(companySkills.id, skill.id)); + }); } } @@ -4152,12 +4154,14 @@ export function companySkillService(db: Db) { throw error; } - // Remove the stale runtime materialization so runtime sync recreates it - // under the new key/slug. - await fs.rm( - path.resolve(managedRoot, "__runtime__", buildSkillRuntimeName(previousKey, previousSlug)), - { recursive: true, force: true }, - ); + // The rename has committed. Cache cleanup must not make a successful rename appear to fail. + try { + await fs.rm(path.resolve(managedRoot, "__runtime__", buildSkillRuntimeName(previousKey, previousSlug)), + { recursive: true, force: true }); + await removeRuntimeSkillCache(managedRoot, skill.id); + } catch (error) { + logger.warn({ err: error, companyId, skillId: skill.id }, "Skill renamed; obsolete runtime cache cleanup failed"); + } const renamed = await getById(companyId, skill.id); if (!renamed) throw notFound("Renamed skill not found"); @@ -4273,6 +4277,10 @@ export function companySkillService(db: Db) { const skill = await getById(companyId, skillId); if (!skill) return null; + return readLoadedSkillFile(skill, relativePath); + } + + async function readLoadedSkillFile(skill: CompanySkill, relativePath: string): Promise { const normalizedPath = normalizePortablePath(relativePath || "SKILL.md"); const fileEntry = skill.fileInventory.find((entry) => entry.path === normalizedPath); if (!fileEntry) { @@ -5672,10 +5680,12 @@ export function companySkillService(db: Db) { let wroteSkillFile = false; for (const entry of skill.fileInventory) { const normalizedPath = normalizePortablePath(entry.path); - const detail = await readFile(companyId, skill.id, normalizedPath).catch(() => null); + const detail = await readLoadedSkillFile(skill, normalizedPath); const content = detail?.content ?? (normalizedPath === "SKILL.md" ? skill.markdown : null); - if (content === null) continue; - const targetPath = path.resolve(skillDir, entry.path); + if (content === null) throw unprocessable("Declared skill file is unavailable"); + const resolved = resolveVersionSnapshotPath(skillDir, entry.path); + if (!resolved) throw unprocessable("Invalid skill file path"); + const targetPath = resolved.targetPath; await fs.mkdir(path.dirname(targetPath), { recursive: true }); await fs.writeFile(targetPath, content, "utf8"); if (normalizedPath === "SKILL.md") wroteSkillFile = true; @@ -5817,6 +5827,24 @@ export function companySkillService(db: Db) { const source = await resolveExistingSkillDirectory(normalizeSkillDirectory(skill)); if (source) return { status: "available", source }; + try { + const cache = runtimeSkillCacheSpec(resolveManagedSkillsRoot(companyId), skill); + if (cache) { + const cachedSource = await resolveRuntimeSkillCache(cache, + async (relativePath) => (await readLoadedSkillFile(skill, relativePath)).content, + options.materializeMissing !== false, + async () => (await getById(companyId, skill.id))?.key === skill.key); + return cachedSource + ? { status: "available", source: cachedSource } + : { status: "missing", source: path.join(cache.entry, "files"), detail: buildMissingRuntimeSourceDetail(skill) }; + } + } catch (error) { + return { + status: "missing", source: resolveRuntimeSkillMaterializedPath(companyId, skill), + detail: `Failed to materialize skill files: ${error instanceof Error ? error.message : String(error)}`, + }; + } + if (options.materializeMissing === false) { const materializedPath = resolveRuntimeSkillMaterializedPath(companyId, skill); const materializedSource = await resolveExistingSkillDirectory(materializedPath); @@ -6937,13 +6965,12 @@ export function companySkillService(db: Db) { ); } - // Delete DB row - await db - .delete(companySkills) - .where(eq(companySkills.id, skillId)); - - // Clean up materialized runtime files - await fs.rm(resolveRuntimeSkillMaterializedPath(companyId, skill), { recursive: true, force: true }); + // Take the cache lifecycle lock before deleting the row. A busy publisher must not + // turn a committed deletion into an apparent API failure, nor recreate its cache. + await removeRuntimeSkillCache(resolveManagedSkillsRoot(companyId), skill.id, async () => { + await fs.rm(resolveRuntimeSkillMaterializedPath(companyId, skill), { recursive: true, force: true }); + await db.delete(companySkills).where(eq(companySkills.id, skillId)); + }); return skill; } diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 6d85f99e8c..350877a00e 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -158,7 +158,7 @@ import { assertManagedProfileRecoveryBinding, resolvePaperclipRunnerNativeProviderInput, } from "./native-runtime/provider-profile.js"; -import type { NativeRunHistoricalSpan } from "./native-runtime/native-run-trace.js"; +import { recordFailedSkillPreparation, type NativeRunHistoricalSpan } from "./native-runtime/native-run-trace.js"; import { parseNativeExecutionInput, type NativeExecutionInput, @@ -18916,18 +18916,39 @@ export function heartbeatService( const runtimeSkillPreference = readPaperclipSkillSyncPreference( effectiveResolvedConfig, ); - const runtimeSkillEntries = await companySkills.listRuntimeSkillEntries( - agent.companyId, - { - versionSelections: skillVersionSelectionMap( - runtimeSkillPreference.desiredSkillEntries, + const nativeRunnerPreparationSpans: NativeRunHistoricalSpan[] = []; + const skillsPrepareStartedAtMs = Date.now(); + const runtimeSkillEntries = await (async () => { + try { + return await companySkills.listRuntimeSkillEntries( + agent.companyId, { - versionPinsEnabled: - resolvedInstanceSettings.experimental.enableBetaSkills === true, + versionSelections: skillVersionSelectionMap( + runtimeSkillPreference.desiredSkillEntries, + { + versionPinsEnabled: + resolvedInstanceSettings.experimental.enableBetaSkills === true, + }, + ), }, - ), - }, - ); + ); + } catch (error) { + if (agent.adapterType === "paperclip_runner") { + await recordFailedSkillPreparation({ + runId: run.id, + startedAtMs: skillsPrepareStartedAtMs, + onEvent: async (event) => { await appendRunEvent(run, event); }, + }); + } + throw error; + } + })(); + nativeRunnerPreparationSpans.push({ + name: "skills.prepare", + parentName: "task.prepare", + startedAtMs: skillsPrepareStartedAtMs, + endedAtMs: Date.now(), + }); let runtimeConfig: Record = { ...effectiveResolvedConfig, paperclipRuntimeSkills: runtimeSkillEntries, @@ -19611,7 +19632,6 @@ export function heartbeatService( }) .where(eq(heartbeatRuns.id, run.id)); } - const nativeRunnerPreparationSpans: NativeRunHistoricalSpan[] = []; const environmentAcquireStartedAtMs = Date.now(); let acquiredEnvironment: Awaited< ReturnType diff --git a/server/src/services/native-runtime/native-run-trace.ts b/server/src/services/native-runtime/native-run-trace.ts index 796533ee7e..e923bdd5b0 100644 --- a/server/src/services/native-runtime/native-run-trace.ts +++ b/server/src/services/native-runtime/native-run-trace.ts @@ -431,3 +431,22 @@ export function createNativeRunTrace(input: { } export type NativeRunTrace = ReturnType; + +/** Emit preparation failure even when execution aborts before a native session exists. */ +export async function recordFailedSkillPreparation(input: { + runId: string; + startedAtMs: number; + onEvent?: NativeRunTraceSink; + traceContext?: StartupTraceContextHandle; +}): Promise { + try { + const trace = createNativeRunTrace(input); + const preparation = trace.start("task.prepare", { parentName: "task.run", startedAtMs: input.startedAtMs }); + const endedAtMs = Date.now(); + await trace.record({ name: "skills.prepare", parentName: "task.prepare", startedAtMs: input.startedAtMs, endedAtMs, outcome: "failed" }); + await trace.end(preparation, { endedAtMs, outcome: "failed" }); + await trace.finish("failed"); + } catch { + // Diagnostics must not replace the original preparation error. + } +} diff --git a/server/src/services/runtime-skill-cache.ts b/server/src/services/runtime-skill-cache.ts new file mode 100644 index 0000000000..76503a6565 --- /dev/null +++ b/server/src/services/runtime-skill-cache.ts @@ -0,0 +1,227 @@ +import { createHash, randomUUID } from "node:crypto"; +import { constants, promises as fs } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import type { CompanySkill } from "@paperclipai/shared"; + +const FORMAT = 1; +const inFlight = new Map>(); +const digest = (value: string | Buffer) => createHash("sha256").update(value).digest("hex"); + +type FileRecord = { path: string; size: number; digest: string }; +type CacheSpec = { root: string; entry: string; fingerprint: string; paths: string[] }; + +function filePath(value: string): string { + const normalized = value.replace(/\\/g, "/"); + if (!normalized || normalized.startsWith("/") || /^[a-z]:/i.test(normalized) + || normalized.split("/").some((part) => !part || part === "." || part === "..") + || normalized.includes("\0")) throw new Error("Invalid runtime skill file path"); + return normalized; +} + +export function runtimeSkillCacheRoot(managedRoot: string, skillId: string): string { + if (!/^[a-zA-Z0-9_-]+$/.test(skillId)) throw new Error("Invalid runtime skill ID"); + // A sibling of __runtime__: old runtime cleanup cannot remove published revisions. + return path.resolve(managedRoot, `__runtime_cache_v${FORMAT}__`, skillId); +} + +export function runtimeSkillCacheSpec(managedRoot: string, skill: CompanySkill): CacheSpec | null { + if ((skill.sourceType === "github" || skill.sourceType === "skills_sh") + && !/^[a-f0-9]{40}$/i.test(skill.sourceRef ?? "")) return null; + const metadata = skill.metadata ?? {}; + const inventory = skill.fileInventory.map((entry) => ({ path: filePath(entry.path), kind: entry.kind })) + .sort((a, b) => a.path.localeCompare(b.path)); + const paths = inventory.map((entry) => entry.path); + if (!paths.includes("SKILL.md")) throw new Error("Company skill could not be materialized because its stored SKILL.md copy is missing."); + if (new Set(paths).size !== paths.length) throw new Error("Invalid runtime skill file inventory"); + const fingerprint = digest(JSON.stringify({ + format: FORMAT, companyId: skill.companyId, skillId: skill.id, + sourceType: skill.sourceType, sourceLocator: skill.sourceLocator, sourceRef: skill.sourceRef, + // These are the only metadata fields used by the source reader. slug is its fallback directory. + source: { owner: metadata.owner, repo: metadata.repo, hostname: metadata.hostname, + ref: skill.sourceRef ? undefined : metadata.ref, repoSkillDir: metadata.repoSkillDir, + fallbackDirectory: typeof metadata.repoSkillDir === "string" ? undefined : skill.slug }, + markdown: digest(skill.markdown), inventory, + })); + const root = runtimeSkillCacheRoot(managedRoot, skill.id); + return { root, entry: path.join(root, fingerprint), fingerprint, paths }; +} + +// Check every ancestor before traversing it, including the configured cache root. +async function assertDirectories(directory: string, trustedRoot: string, create = false): Promise { + const absolute = path.resolve(directory); + let cursor = path.resolve(trustedRoot); + if (!absolute.startsWith(`${cursor}${path.sep}`)) throw new Error("Runtime cache escaped its root"); + // Ancestors of the configured storage root may be system aliases (e.g. macOS /var). + if (create) await fs.mkdir(cursor, { recursive: true }); + const rootStat = await fs.lstat(cursor); + if (!rootStat.isDirectory() || rootStat.isSymbolicLink()) throw new Error("Unsafe runtime skill storage root"); + for (const part of absolute.slice(cursor.length).split(path.sep).filter(Boolean)) { + cursor = path.join(cursor, part); + if (create) await fs.mkdir(cursor).catch((error: NodeJS.ErrnoException) => { + if (error.code !== "EEXIST") throw error; + }); + const stat = await fs.lstat(cursor); + if (!stat.isDirectory() || stat.isSymbolicLink()) throw new Error("Unsafe runtime skill cache directory"); + } +} + +async function readRegularFile(filename: string): Promise { + const handle = await fs.open(filename, constants.O_RDONLY | constants.O_NOFOLLOW); + try { + if (!(await handle.stat()).isFile()) throw new Error("Unsafe runtime skill cache file"); + return await handle.readFile(); + } finally { await handle.close(); } +} + +async function inventory(directory: string, base = ""): Promise { + const out: string[] = []; + if ((await fs.lstat(directory)).mode & 0o222) throw new Error("Writable runtime skill cache directory"); + for (const entry of await fs.readdir(directory, { withFileTypes: true })) { + const relative = base ? `${base}/${entry.name}` : entry.name; + if (entry.isDirectory()) out.push(...await inventory(path.join(directory, entry.name), relative)); + else if (entry.isFile()) out.push(relative); + else throw new Error("Symlink or special file in runtime skill cache"); + } + return out.sort(); +} + +async function matches(spec: CacheSpec, entry = spec.entry): Promise { + try { + await assertDirectories(path.join(entry, "files"), path.dirname(path.dirname(spec.root))); + const manifest = JSON.parse((await readRegularFile(path.join(entry, "manifest.json"))).toString("utf8")); + if (manifest.format !== FORMAT || manifest.fingerprint !== spec.fingerprint || !Array.isArray(manifest.files) + || manifest.files.length !== spec.paths.length) return false; + const actual = await inventory(path.join(entry, "files")); + if (JSON.stringify(actual) !== JSON.stringify([...spec.paths].sort())) return false; + const seen = new Set(); + for (const record of manifest.files as FileRecord[]) { + if (!record || typeof record.path !== "string" || filePath(record.path) !== record.path + || !spec.paths.includes(record.path) || seen.has(record.path) + || !Number.isSafeInteger(record.size) || record.size < 0 || !/^[a-f0-9]{64}$/.test(record.digest)) return false; + seen.add(record.path); + const content = await readRegularFile(path.join(entry, "files", record.path)); + if ((await fs.lstat(path.join(entry, "files", record.path))).mode & 0o222) return false; + if (content.length !== record.size || digest(content) !== record.digest) return false; + } + return true; + } catch { return false; } +} + +// Serialize builds and cleanup for one skill across processes as well as callers. +// A hard link publishes complete lock ownership atomically; crashed owners are reported without stealing another publisher’s lock. +async function publishLocked(root: string, fingerprint: string, action: () => Promise): Promise { + const lock = path.join(root, `${fingerprint}.lock`); + const owner = path.join(root, `.owner-${randomUUID()}`); + await fs.writeFile(owner, JSON.stringify({ pid: process.pid, host: os.hostname() }), { flag: "wx" }); + let acquired = false; + try { + const deadline = Date.now() + 60_000; + while (!acquired) { + try { await fs.link(owner, lock); acquired = true; } + catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + const lockContent = await readRegularFile(lock).catch((readError: NodeJS.ErrnoException) => { + if (readError.code === "ENOENT") return null; + throw readError; + }); + if (!lockContent) continue; + const holder = JSON.parse(lockContent.toString("utf8")); + if (holder.host === os.hostname() && Number.isSafeInteger(holder.pid) && holder.pid > 0) { + try { process.kill(holder.pid, 0); } + catch (probeError) { + if ((probeError as NodeJS.ErrnoException).code === "ESRCH") { + throw new Error("Runtime skill cache publisher exited; remove its stale publication lock before retrying"); + } + } + } + if (Date.now() >= deadline) throw new Error("Runtime skill cache publisher is busy; retry preparation"); + await new Promise((resolve) => setTimeout(resolve, 20)); + } + } + return await action(); + } finally { + if (acquired) await fs.unlink(lock).catch(() => {}); + await fs.unlink(owner).catch(() => {}); + } +} + +async function setTreeMode(directory: string, readonly: boolean): Promise { + const stat = await fs.lstat(directory).catch((error: NodeJS.ErrnoException) => { + if (error.code === "ENOENT") return null; + throw error; + }); + if (!stat || stat.isSymbolicLink()) return; + if (!stat.isDirectory()) { + if (readonly && stat.isFile()) await fs.chmod(directory, 0o444); + return; + } + if (!readonly) await fs.chmod(directory, 0o700); + for (const entry of await fs.readdir(directory)) await setTreeMode(path.join(directory, entry), readonly); + if (readonly) await fs.chmod(directory, 0o555); +} + +async function removeTree(directory: string): Promise { + await setTreeMode(directory, false); + await fs.rm(directory, { recursive: true, force: true }); +} + +export async function resolveRuntimeSkillCache( + spec: CacheSpec, read: (relativePath: string) => Promise, materialize = true, + stillInstalled: () => Promise = async () => true, +): Promise { + if (await matches(spec)) return path.join(spec.entry, "files"); + if (!materialize) return null; + const active = inFlight.get(spec.entry); + if (active) return active; + const build = (async () => { + const namespace = path.dirname(spec.root); + await assertDirectories(namespace, path.dirname(namespace), true); + // The lock lives outside the skill directory, so cleanup cannot unlink an active lock. + return publishLocked(namespace, path.basename(spec.root), async () => { + if (!await stillInstalled()) throw new Error("Skill was renamed or removed during preparation"); + await assertDirectories(spec.root, path.dirname(namespace), true); + if (await matches(spec)) return path.join(spec.entry, "files"); + const staging = await fs.mkdtemp(path.join(spec.root, ".staging-")); + try { + await fs.mkdir(path.join(staging, "files")); + const files: FileRecord[] = []; + for (const relative of spec.paths) { + const content = Buffer.from(await read(relative), "utf8"); + const target = path.join(staging, "files", relative); + await fs.mkdir(path.dirname(target), { recursive: true }); + await fs.writeFile(target, content, { flag: "wx" }); + files.push({ path: relative, size: content.length, digest: digest(content) }); + } + await fs.writeFile(path.join(staging, "manifest.json"), JSON.stringify({ format: FORMAT, fingerprint: spec.fingerprint, files })); + await setTreeMode(staging, true); + if (!await matches(spec, staging)) throw new Error("Runtime skill cache validation failed"); + // Lifecycle mutations can update the DB while this builder owns the filesystem lock. + if (!await stillInstalled()) throw new Error("Skill was renamed or removed during preparation"); + await fs.rename(spec.entry, path.join(spec.root, `.invalid-${spec.fingerprint}-${randomUUID()}`)) + .catch((error: NodeJS.ErrnoException) => { if (error.code !== "ENOENT") throw error; }); + await fs.rename(staging, spec.entry); + return path.join(spec.entry, "files"); + } finally { await removeTree(staging); } + }); + })(); + inFlight.set(spec.entry, build); + try { return await build; } finally { if (inFlight.get(spec.entry) === build) inFlight.delete(spec.entry); } +} + +export async function removeRuntimeSkillCache( + managedRoot: string, skillId: string, afterRemove?: () => Promise, +): Promise { + const root = runtimeSkillCacheRoot(managedRoot, skillId); + const namespace = path.dirname(root); + try { await assertDirectories(namespace, managedRoot, Boolean(afterRemove)); } + catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT" && !afterRemove) return; throw error; } + await publishLocked(namespace, skillId, async () => { + try { + await assertDirectories(root, managedRoot); + await removeTree(root); + } catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + // Commit deletion while builders remain excluded. Lock/cleanup failures leave the row intact. + await afterRemove?.(); + }); +}