chore: remove registry artifact backup jobs

Stop custom registry artifact backup/backfill/restore behavior now that Convex backups with file storage are the recovery source of truth. Legacy backup schema tables remain inert until a separate verified cleanup removes stored rows.
This commit is contained in:
Patrick Erichsen
2026-06-19 13:35:44 -07:00
committed by GitHub
parent 2ae30d70d5
commit 7accfb71c7
23 changed files with 12 additions and 6148 deletions
+5 -6
View File
@@ -175,12 +175,11 @@ Without `OPENAI_API_KEY`, public corpus import still works, but semantic search
These features degrade gracefully without their keys:
| Variable | Purpose |
| ---------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------- |
| `OPENAI_API_KEY` | Embeddings and vector search (falls back to zero vectors) |
| `VT_API_KEY` | VirusTotal malware scanning |
| `DISCORD_WEBHOOK_URL` | Discord notifications |
| `REGISTRY_BACKUP_R2_ACCOUNT_ID` / `REGISTRY_BACKUP_BUCKET` / `REGISTRY_BACKUP_ACCESS_KEY_ID` / `REGISTRY_BACKUP_SECRET_ACCESS_KEY` | Registry artifact publish backup and seed/backfill |
| Variable | Purpose |
| --------------------- | --------------------------------------------------------- |
| `OPENAI_API_KEY` | Embeddings and vector search (falls back to zero vectors) |
| `VT_API_KEY` | VirusTotal malware scanning |
| `DISCORD_WEBHOOK_URL` | Discord notifications |
## CLI Development
-10
View File
@@ -93,7 +93,6 @@ import type * as lib_publisherStats from "../lib/publisherStats.js";
import type * as lib_publishers from "../lib/publishers.js";
import type * as lib_rateLimitConfig from "../lib/rateLimitConfig.js";
import type * as lib_recommendationScore from "../lib/recommendationScore.js";
import type * as lib_registryArtifactBackup from "../lib/registryArtifactBackup.js";
import type * as lib_reporting from "../lib/reporting.js";
import type * as lib_reservedHandles from "../lib/reservedHandles.js";
import type * as lib_reservedSlugs from "../lib/reservedSlugs.js";
@@ -131,10 +130,6 @@ import type * as publisherAbuse from "../publisherAbuse.js";
import type * as publisherAbuseDevSeed from "../publisherAbuseDevSeed.js";
import type * as publishers from "../publishers.js";
import type * as rateLimits from "../rateLimits.js";
import type * as registryArtifactBackups from "../registryArtifactBackups.js";
import type * as registryArtifactBackupsNode from "../registryArtifactBackupsNode.js";
import type * as registryArtifactRestore from "../registryArtifactRestore.js";
import type * as registryArtifactRestoreMutations from "../registryArtifactRestoreMutations.js";
import type * as retention from "../retention.js";
import type * as search from "../search.js";
import type * as securityDataset from "../securityDataset.js";
@@ -245,7 +240,6 @@ declare const fullApi: ApiFromModules<{
"lib/publishers": typeof lib_publishers;
"lib/rateLimitConfig": typeof lib_rateLimitConfig;
"lib/recommendationScore": typeof lib_recommendationScore;
"lib/registryArtifactBackup": typeof lib_registryArtifactBackup;
"lib/reporting": typeof lib_reporting;
"lib/reservedHandles": typeof lib_reservedHandles;
"lib/reservedSlugs": typeof lib_reservedSlugs;
@@ -283,10 +277,6 @@ declare const fullApi: ApiFromModules<{
publisherAbuseDevSeed: typeof publisherAbuseDevSeed;
publishers: typeof publishers;
rateLimits: typeof rateLimits;
registryArtifactBackups: typeof registryArtifactBackups;
registryArtifactBackupsNode: typeof registryArtifactBackupsNode;
registryArtifactRestore: typeof registryArtifactRestore;
registryArtifactRestoreMutations: typeof registryArtifactRestoreMutations;
retention: typeof retention;
search: typeof search;
securityDataset: typeof securityDataset;
-16
View File
@@ -4,7 +4,6 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => {
const interval = vi.fn();
const githubSkillSyncRef = Symbol("github-skill-source-sync");
const registryArtifactBackupRetryRef = Symbol("registry-artifact-backup-retry");
const installTelemetryDedupePruneRef = Symbol("install-telemetry-dedupe-prune");
const rateLimitCountersPruneRef = Symbol("rate-limit-counters-prune");
const skillStatEventPruneRef = Symbol("skill-stat-event-prune");
@@ -13,7 +12,6 @@ const mocks = vi.hoisted(() => {
return {
interval,
githubSkillSyncRef,
registryArtifactBackupRetryRef,
installTelemetryDedupePruneRef,
rateLimitCountersPruneRef,
skillStatEventPruneRef,
@@ -30,9 +28,6 @@ vi.mock("convex/server", () => ({
vi.mock("./_generated/api", () => ({
internal: {
registryArtifactBackupsNode: {
processRegistryArtifactBackupRetriesInternal: mocks.registryArtifactBackupRetryRef,
},
githubSkillSyncNode: { syncGitHubSkillSourcesInternal: mocks.githubSkillSyncRef },
leaderboards: { rebuildTrendingLeaderboardAction: Symbol("trending-leaderboard") },
statsMaintenance: {
@@ -93,17 +88,6 @@ describe("crons", () => {
expect(mocks.interval).not.toHaveBeenCalled();
});
it("drains registry artifact backup retries frequently enough for publish bursts", async () => {
await import("./crons");
expect(mocks.interval).toHaveBeenCalledWith(
"registry-artifact-backup-retries",
{ minutes: 5 },
mocks.registryArtifactBackupRetryRef,
{},
);
});
it("runs GitHub skill source sync every 15 minutes", async () => {
await import("./crons");
-7
View File
@@ -5,13 +5,6 @@ import { RETENTION_STANDARD_BATCH_SIZE } from "./lib/retentionPolicy";
const crons = cronJobs();
if (process.env.CLAWHUB_DISABLE_CRONS !== "1") {
crons.interval(
"registry-artifact-backup-retries",
{ minutes: 5 },
internal.registryArtifactBackupsNode.processRegistryArtifactBackupRetriesInternal,
{},
);
crons.interval(
"github-skill-source-sync",
{ minutes: 15 },
-94
View File
@@ -366,100 +366,6 @@ describe("httpApiV1 handlers", () => {
expect(runAction).not.toHaveBeenCalled();
});
it("users/restore forbids non-admin api tokens", async () => {
const runQuery = vi.fn();
const runAction = vi.fn();
const runMutation = vi.fn(async (_mutation: unknown, args: Record<string, unknown>) => {
if (isRateLimitArgs(args)) return okRate();
if (args.ownerHandle === "me") return { publisherId: "publishers:me" };
return okRate();
});
vi.mocked(requireApiTokenUser).mockResolvedValue({
userId: "users:actor",
user: { _id: "users:actor", role: "user" },
} as never);
const response = await __handlers.usersPostRouterV1Handler(
makeCtx({ runQuery, runAction, runMutation }),
new Request("https://example.com/api/v1/users/restore", {
method: "POST",
body: JSON.stringify({ handle: "target", slugs: ["a"] }),
}),
);
expect(response.status).toBe(403);
expect(runQuery).not.toHaveBeenCalled();
expect(runAction).not.toHaveBeenCalled();
});
it("users/restore calls restore action for admin", async () => {
const runAction = vi.fn().mockResolvedValue({ ok: true, totalRestored: 1, results: [] });
const runMutation = vi.fn(async (_mutation: unknown, args: Record<string, unknown>) => {
if (isRateLimitArgs(args)) return okRate();
return { ok: true };
});
const runQuery = vi.fn(async (_query: unknown, args: Record<string, unknown>) => {
if ("handle" in args) return { _id: "users:target" };
return null;
});
vi.mocked(requireApiTokenUser).mockResolvedValue({
userId: "users:admin",
user: { _id: "users:admin", role: "admin" },
} as never);
const response = await __handlers.usersPostRouterV1Handler(
makeCtx({ runQuery, runAction, runMutation }),
new Request("https://example.com/api/v1/users/restore", {
method: "POST",
body: JSON.stringify({
handle: "Target",
slugs: ["a", "b"],
versionsBySlug: { a: "1.0.0", b: "1.1.0" },
forceOverwriteSquatter: true,
}),
}),
);
if (response.status !== 200) throw new Error(await response.text());
expect(runAction).toHaveBeenCalledWith(expect.anything(), {
actorUserId: "users:admin",
ownerHandle: "target",
ownerUserId: "users:target",
slugs: ["a", "b"],
versionsBySlug: { a: "1.0.0", b: "1.1.0" },
forceOverwriteSquatter: true,
});
});
it("users/restore requires a backup version for every slug", async () => {
const runAction = vi.fn();
const runMutation = vi.fn(async (_mutation: unknown, args: Record<string, unknown>) => {
if (isRateLimitArgs(args)) return okRate();
return { ok: true };
});
const runQuery = vi.fn();
vi.mocked(requireApiTokenUser).mockResolvedValue({
userId: "users:admin",
user: { _id: "users:admin", role: "admin" },
} as never);
const response = await __handlers.usersPostRouterV1Handler(
makeCtx({ runQuery, runAction, runMutation }),
new Request("https://example.com/api/v1/users/restore", {
method: "POST",
body: JSON.stringify({
handle: "Target",
slugs: ["a", "b"],
versionsBySlug: { a: "1.0.0" },
forceOverwriteSquatter: true,
}),
}),
);
expect(response.status).toBe(400);
expect(await response.text()).toBe("Missing backup version for slug b");
expect(runQuery).not.toHaveBeenCalled();
expect(runAction).not.toHaveBeenCalled();
});
it("skills export allows authenticated non-admin users at the key rate limit", async () => {
vi.mocked(requireApiTokenUser).mockResolvedValue({
userId: "users:actor",
-70
View File
@@ -86,7 +86,6 @@ export async function usersPostRouterV1Handler(ctx: ActionCtx, request: Request)
action !== "ban" &&
action !== "unban" &&
action !== "role" &&
action !== "restore" &&
action !== "reclassify-ban" &&
action !== "ban-appeal-unban" &&
action !== "reclaim" &&
@@ -114,13 +113,6 @@ export async function usersPostRouterV1Handler(ctx: ActionCtx, request: Request)
const actorUserId = authResult.userId;
const actorUser = authResult.user;
// Restore and reclaim have different parameter shapes, handle them separately
if (action === "restore") {
const admin = requireAdminOrResponse(actorUser, rate.headers);
if (!admin.ok) return admin.response;
return handleAdminRestore(ctx, request, payload, actorUserId, rate.headers);
}
if (action === "reclassify-ban") {
const admin = requireAdminOrResponse(actorUser, rate.headers);
if (!admin.ok) return admin.response;
@@ -554,68 +546,6 @@ export async function usersGetRouterV1Handler(ctx: ActionCtx, request: Request)
}
}
/**
* POST /api/v1/users/restore
* Admin-only: restore skills from registry artifact backup for a user.
* Body: { handle: string, slugs: string[], versionsBySlug: Record<string, string>, forceOverwriteSquatter?: boolean }
*/
async function handleAdminRestore(
ctx: ActionCtx,
_request: Request,
payload: Record<string, unknown>,
actorUserId: Id<"users">,
headers: HeadersInit,
) {
const handle = typeof payload.handle === "string" ? payload.handle.trim().toLowerCase() : "";
if (!handle) return text("Missing handle", 400, headers);
const slugs = Array.isArray(payload.slugs)
? payload.slugs.filter((s): s is string => typeof s === "string")
: [];
if (slugs.length === 0) return text("Missing slugs array", 400, headers);
if (slugs.length > 100) return text("Too many slugs (max 100)", 400, headers);
const versionsBySlug =
payload.versionsBySlug && typeof payload.versionsBySlug === "object"
? Object.fromEntries(
Object.entries(payload.versionsBySlug).filter(
(entry): entry is [string, string] =>
typeof entry[0] === "string" && typeof entry[1] === "string",
),
)
: undefined;
if (!versionsBySlug) return text("Missing versionsBySlug", 400, headers);
const missingVersionSlug = slugs.find((slug) => !versionsBySlug[slug]?.trim());
if (missingVersionSlug) {
return text(`Missing backup version for slug ${missingVersionSlug}`, 400, headers);
}
const forceOverwriteSquatter = Boolean(payload.forceOverwriteSquatter);
const targetUser = await ctx.runQuery(api.users.getByHandle, { handle });
if (!targetUser?._id) return text("User not found", 404, headers);
try {
const result = await ctx.runAction(
internal.registryArtifactRestore.restoreUserSkillsFromBackup,
{
actorUserId,
ownerHandle: handle,
ownerUserId: targetUser._id,
slugs,
versionsBySlug,
forceOverwriteSquatter,
},
);
return json(result, 200, headers);
} catch (error) {
const message = error instanceof Error ? error.message : "Restore failed";
if (message.toLowerCase().includes("forbidden")) {
return text("Forbidden", 403, headers);
}
return text(message, 400, headers);
}
}
/**
* POST /api/v1/users/reclaim
* Admin-only: reclaim root slugs for the rightful owner.
-486
View File
@@ -1,486 +0,0 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import type { Id } from "../_generated/dataModel";
import {
__registryArtifactBackupTestInternals,
backupPackageReleaseToObjectStorage,
backupSkillVersionToObjectStorage,
buildPackageReleaseBackupManifest,
buildSkillVersionBackupManifest,
getRegistryArtifactBackupSettings,
readRegistryArtifactBackupObject,
} from "./registryArtifactBackup";
describe("registry artifact backup settings", () => {
const originalEnv = {
endpoint: process.env.REGISTRY_BACKUP_S3_ENDPOINT,
accountId: process.env.REGISTRY_BACKUP_R2_ACCOUNT_ID,
bucket: process.env.REGISTRY_BACKUP_BUCKET,
accessKeyId: process.env.REGISTRY_BACKUP_ACCESS_KEY_ID,
secretAccessKey: process.env.REGISTRY_BACKUP_SECRET_ACCESS_KEY,
region: process.env.REGISTRY_BACKUP_S3_REGION,
skillsRoot: process.env.REGISTRY_BACKUP_SKILLS_ROOT,
packagesRoot: process.env.REGISTRY_BACKUP_PACKAGES_ROOT,
skillFileUploadConcurrency: process.env.REGISTRY_BACKUP_SKILL_FILE_UPLOAD_CONCURRENCY,
};
afterEach(() => {
setEnv("REGISTRY_BACKUP_S3_ENDPOINT", originalEnv.endpoint);
setEnv("REGISTRY_BACKUP_R2_ACCOUNT_ID", originalEnv.accountId);
setEnv("REGISTRY_BACKUP_BUCKET", originalEnv.bucket);
setEnv("REGISTRY_BACKUP_ACCESS_KEY_ID", originalEnv.accessKeyId);
setEnv("REGISTRY_BACKUP_SECRET_ACCESS_KEY", originalEnv.secretAccessKey);
setEnv("REGISTRY_BACKUP_S3_REGION", originalEnv.region);
setEnv("REGISTRY_BACKUP_SKILLS_ROOT", originalEnv.skillsRoot);
setEnv("REGISTRY_BACKUP_PACKAGES_ROOT", originalEnv.packagesRoot);
setEnv("REGISTRY_BACKUP_SKILL_FILE_UPLOAD_CONCURRENCY", originalEnv.skillFileUploadConcurrency);
});
it("defaults registry artifact backups to skills and packages object roots", () => {
delete process.env.REGISTRY_BACKUP_S3_ENDPOINT;
process.env.REGISTRY_BACKUP_R2_ACCOUNT_ID = "account-id";
process.env.REGISTRY_BACKUP_BUCKET = "clawhub-registry-backup";
process.env.REGISTRY_BACKUP_ACCESS_KEY_ID = "access-key";
process.env.REGISTRY_BACKUP_SECRET_ACCESS_KEY = "secret-key";
delete process.env.REGISTRY_BACKUP_S3_REGION;
delete process.env.REGISTRY_BACKUP_SKILLS_ROOT;
delete process.env.REGISTRY_BACKUP_PACKAGES_ROOT;
expect(getRegistryArtifactBackupSettings()).toEqual({
endpoint: "https://account-id.r2.cloudflarestorage.com",
bucket: "clawhub-registry-backup",
accessKeyId: "access-key",
secretAccessKey: "secret-key",
region: "auto",
skillsRoot: "skills",
packagesRoot: "packages",
});
});
it("builds versioned skill backup paths and restore metadata", () => {
const manifest = buildSkillVersionBackupManifest({
root: "skills",
ownerHandle: "OpenClaw Team",
skillId: "skills:demo" as Id<"skills">,
versionId: "skillVersions:demo-1" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
files: [
{
path: "SKILL.md",
size: 42,
storageId: "storage:skill" as Id<"_storage">,
sha256: "sha256:skill",
contentType: "text/markdown",
},
],
});
expect(manifest).toMatchObject({
skillRoot: "skills/openclaw-team/demo-skill",
versionRoot: "skills/openclaw-team/demo-skill/1%2E2%2E3",
metaPath: "skills/openclaw-team/demo-skill/1%2E2%2E3/_meta.json",
fileObjects: [
{
key: "skills/openclaw-team/demo-skill/1%2E2%2E3/SKILL.md",
path: "SKILL.md",
sha256: "sha256:skill",
contentType: "text/markdown",
},
],
meta: {
kind: "skillVersion",
owner: "openclaw-team",
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.2.3",
restore: {
skillId: "skills:demo",
versionId: "skillVersions:demo-1",
},
},
});
});
it("rejects unsafe skill file paths before writing backup object keys", () => {
expect(() =>
buildSkillVersionBackupManifest({
root: "skills",
ownerHandle: "OpenClaw Team",
versionId: "skillVersions:demo-1" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
files: [
{
path: "../SKILL.md",
size: 42,
storageId: "storage:skill" as Id<"_storage">,
sha256: "sha256:skill",
},
],
}),
).toThrow("Invalid skill backup file path");
});
it("builds package release backup paths and restore metadata", () => {
const manifest = buildPackageReleaseBackupManifest({
root: "packages",
ownerHandle: "OpenClaw Team",
packageId: "packages:demo" as Id<"packages">,
releaseId: "packageReleases:demo-1" as Id<"packageReleases">,
packageName: "@openclaw/demo-plugin",
normalizedName: "@openclaw/demo-plugin",
displayName: "Demo Plugin",
family: "code-plugin",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
artifactKind: "npm-pack",
artifactFileName: "demo-plugin-1.2.3.tgz",
artifactSha256: "sha256:artifact",
artifactSize: 42,
artifactFormat: "tgz",
npmIntegrity: "sha512-demo",
npmShasum: "abc123",
files: [{ path: "package.json", size: 10, sha256: "sha256:package-json" }],
});
expect(manifest).toMatchObject({
packageRoot: "packages/openclaw-team/%40openclaw%2Fdemo-plugin",
releaseRoot: "packages/openclaw-team/%40openclaw%2Fdemo-plugin/1%2E2%2E3",
artifactPath:
"packages/openclaw-team/%40openclaw%2Fdemo-plugin/1%2E2%2E3/demo-plugin-1.2.3.tgz",
metaPath: "packages/openclaw-team/%40openclaw%2Fdemo-plugin/1%2E2%2E3/_meta.json",
meta: {
kind: "packageRelease",
restore: {
packageId: "packages:demo",
releaseId: "packageReleases:demo-1",
},
artifact: {
path: "demo-plugin-1.2.3.tgz",
sha256: "sha256:artifact",
size: 42,
format: "tgz",
npmIntegrity: "sha512-demo",
npmShasum: "abc123",
},
},
});
});
it("rejects unsafe package artifact filenames before writing backup object keys", () => {
expect(() =>
buildPackageReleaseBackupManifest({
root: "packages",
ownerHandle: "OpenClaw Team",
packageId: "packages:demo" as Id<"packages">,
releaseId: "packageReleases:demo-1" as Id<"packageReleases">,
packageName: "@openclaw/demo-plugin",
normalizedName: "@openclaw/demo-plugin",
displayName: "Demo Plugin",
family: "code-plugin",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
artifactKind: "npm-pack",
artifactFileName: "../evil.tgz",
artifactSha256: "sha256:artifact",
artifactSize: 42,
artifactFormat: "tgz",
files: [],
}),
).toThrow("Invalid package backup artifact filename");
});
it("uses lossless path encoding to avoid package and version collisions", () => {
expect(__registryArtifactBackupTestInternals.encodeBackupPathSegment("@openclaw/demo")).toBe(
"%40openclaw%2Fdemo",
);
expect(__registryArtifactBackupTestInternals.encodeBackupPathSegment("foo.bar")).toBe(
"foo%2Ebar",
);
expect(__registryArtifactBackupTestInternals.encodeBackupPathSegment("foo_bar")).toBe(
"foo_bar",
);
});
it("uses strict AWS URI encoding for object keys before signing R2 requests", () => {
expect(
__registryArtifactBackupTestInternals.encodeObjectKey(
"skills/smartpeopleconnected/token-optimizer/1%2E0%2E0/infomaterial/4_github_publish_Alles ist fertig!.txt",
),
).toBe(
"skills/smartpeopleconnected/token-optimizer/1%252E0%252E0/infomaterial/4_github_publish_Alles%20ist%20fertig%21.txt",
);
});
it("preserves valid owner handle punctuation in backup paths", () => {
const dotted = buildSkillVersionBackupManifest({
root: "skills",
ownerHandle: "foo.bar",
versionId: "skillVersions:dotted" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.0.0",
publishedAt: 1_700_000_000_000,
files: [],
});
const underscored = buildSkillVersionBackupManifest({
root: "skills",
ownerHandle: "foo_bar",
versionId: "skillVersions:underscored" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.0.0",
publishedAt: 1_700_000_000_000,
files: [],
});
const dashed = buildSkillVersionBackupManifest({
root: "skills",
ownerHandle: "foo-bar",
versionId: "skillVersions:dashed" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.0.0",
publishedAt: 1_700_000_000_000,
files: [],
});
expect([dotted.skillRoot, underscored.skillRoot, dashed.skillRoot]).toEqual([
"skills/foo.bar/demo-skill",
"skills/foo_bar/demo-skill",
"skills/foo-bar/demo-skill",
]);
});
it("reads object bytes from object storage", async () => {
vi.stubGlobal(
"fetch",
vi.fn(async (url: URL | string, init?: RequestInit) => {
const key = objectKey(String(url));
if (
init?.method === "GET" &&
key === "skills/openclaw-team/demo-skill/1%2E2%2E3/SKILL.md"
) {
return response(200, "hello skill");
}
return response(404, "");
}),
);
const bytes = await readRegistryArtifactBackupObject(
makeContext(),
"skills/openclaw-team/demo-skill/1%2E2%2E3/SKILL.md",
);
expect(Buffer.from(bytes!).toString("utf8")).toBe("hello skill");
});
it("writes skill files and version metadata to object storage", async () => {
const calls: Array<{ method: string; url: string; body: string }> = [];
vi.stubGlobal(
"fetch",
vi.fn(async (url: URL | string, init?: RequestInit) => {
const method = init?.method ?? "GET";
const body = await requestBodyText(init?.body);
calls.push({ method, url: String(url), body });
return response(200, "");
}),
);
await backupSkillVersionToObjectStorage(
makeStorageCtx({ "storage:skill": "hello skill" }) as never,
{
root: "skills",
ownerHandle: "OpenClaw Team",
versionId: "skillVersions:demo-1" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
files: [
{
path: "SKILL.md",
size: 11,
storageId: "storage:skill" as Id<"_storage">,
sha256: "sha256:skill",
contentType: "text/markdown",
},
],
},
makeContext(),
);
expect(calls.map((call) => [call.method, objectKey(call.url)])).toEqual([
["PUT", "skills/openclaw-team/demo-skill/1%2E2%2E3/SKILL.md"],
["PUT", "skills/openclaw-team/demo-skill/1%2E2%2E3/_meta.json"],
]);
expect(JSON.parse(calls[1].body)).toMatchObject({
kind: "skillVersion",
version: "1.2.3",
metadata: { files: [{ path: "SKILL.md", sha256: "sha256:skill" }] },
});
});
it("uploads skill files with bounded parallelism before writing version metadata", async () => {
process.env.REGISTRY_BACKUP_SKILL_FILE_UPLOAD_CONCURRENCY = "3";
const calls: Array<{ method: string; key: string }> = [];
let activeFileUploads = 0;
let maxActiveFileUploads = 0;
let completedFileUploads = 0;
let completedWhenMetaStarted = 0;
vi.stubGlobal(
"fetch",
vi.fn(async (url: URL | string, init?: RequestInit) => {
const key = objectKey(String(url));
calls.push({ method: init?.method ?? "GET", key });
if (!key.endsWith("/_meta.json")) {
activeFileUploads += 1;
maxActiveFileUploads = Math.max(maxActiveFileUploads, activeFileUploads);
await new Promise((resolve) => setTimeout(resolve, 20));
activeFileUploads -= 1;
completedFileUploads += 1;
} else {
completedWhenMetaStarted = completedFileUploads;
}
return response(200, "");
}),
);
const files = Array.from({ length: 6 }, (_, index) => ({
path: `file-${index}.txt`,
size: 6,
storageId: `storage:file-${index}` as Id<"_storage">,
sha256: `sha256:file-${index}`,
contentType: "text/plain",
}));
await backupSkillVersionToObjectStorage(
makeStorageCtx(
Object.fromEntries(files.map((file) => [file.storageId, `body-${file.path}`])),
) as never,
{
root: "skills",
ownerHandle: "OpenClaw Team",
versionId: "skillVersions:demo-1" as Id<"skillVersions">,
slug: "demo-skill",
displayName: "Demo Skill",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
files,
},
makeContext(),
);
expect(maxActiveFileUploads).toBe(3);
expect(completedWhenMetaStarted).toBe(6);
expect(calls.at(-1)).toEqual({
method: "PUT",
key: "skills/openclaw-team/demo-skill/1%2E2%2E3/_meta.json",
});
});
it("writes package artifacts and version metadata to object storage", async () => {
const calls: Array<{ method: string; url: string; body: string }> = [];
vi.stubGlobal(
"fetch",
vi.fn(async (url: URL | string, init?: RequestInit) => {
const method = init?.method ?? "GET";
const body = await requestBodyText(init?.body);
calls.push({ method, url: String(url), body });
return response(200, "");
}),
);
await backupPackageReleaseToObjectStorage(
makeStorageCtx({ "storage:artifact": "tgz bytes" }) as never,
{
root: "packages",
ownerHandle: "OpenClaw Team",
packageId: "packages:demo" as Id<"packages">,
releaseId: "packageReleases:demo-1" as Id<"packageReleases">,
packageName: "@openclaw/demo-plugin",
normalizedName: "@openclaw/demo-plugin",
displayName: "Demo Plugin",
family: "code-plugin",
version: "1.2.3",
publishedAt: 1_700_000_000_000,
artifactStorageId: "storage:artifact" as Id<"_storage">,
artifactFileName: "demo-plugin-1.2.3.tgz",
artifactSha256: "sha256:artifact",
artifactSize: 9,
files: [],
},
makeContext(),
);
expect(calls.map((call) => [call.method, objectKey(call.url)])).toEqual([
["PUT", "packages/openclaw-team/%40openclaw%2Fdemo-plugin/1%2E2%2E3/demo-plugin-1.2.3.tgz"],
["PUT", "packages/openclaw-team/%40openclaw%2Fdemo-plugin/1%2E2%2E3/_meta.json"],
]);
expect(JSON.parse(calls[1].body)).toMatchObject({
kind: "packageRelease",
artifact: { path: "demo-plugin-1.2.3.tgz", sha256: "sha256:artifact" },
});
});
});
function setEnv(name: keyof NodeJS.ProcessEnv, value: string | undefined) {
if (value === undefined) {
delete process.env[name];
return;
}
process.env[name] = value;
}
function makeContext() {
return {
endpoint: "https://account.r2.cloudflarestorage.com",
bucket: "clawhub-registry-backup",
accessKeyId: "access-key",
secretAccessKey: "secret-key",
region: "auto",
skillsRoot: "skills",
packagesRoot: "packages",
};
}
function makeStorageCtx(contents: Record<string, string>) {
return {
storage: {
get: async (id: Id<"_storage">) => {
const value = contents[id];
return value === undefined ? null : new Blob([value]);
},
},
};
}
function response(status: number, body: string, headers: Record<string, string> = {}) {
return {
ok: status >= 200 && status < 300,
status,
headers: new Headers(headers),
text: async () => body,
json: async () => JSON.parse(body),
arrayBuffer: async () => {
const buffer = Buffer.from(body);
return buffer.buffer.slice(buffer.byteOffset, buffer.byteOffset + buffer.byteLength);
},
};
}
function objectKey(url: string) {
const parsed = new URL(url);
const prefix = "/clawhub-registry-backup/";
return decodeURIComponent(parsed.pathname.slice(prefix.length));
}
async function requestBodyText(body: BodyInit | null | undefined) {
if (!body) return "";
return Buffer.from(await new Response(body).arrayBuffer()).toString("utf8");
}
-525
View File
@@ -1,525 +0,0 @@
"use node";
import { createHash, createHmac } from "node:crypto";
import type { Id } from "../_generated/dataModel";
import type { ActionCtx } from "../_generated/server";
import { validateFilePath } from "./skillZip";
const DEFAULT_SKILLS_ROOT = "skills";
const DEFAULT_PACKAGES_ROOT = "packages";
const DEFAULT_SKILL_FILE_UPLOAD_CONCURRENCY = 16;
const MAX_SKILL_FILE_UPLOAD_CONCURRENCY = 64;
const META_FILENAME = "_meta.json";
type BackupFile = {
path: string;
size: number;
storageId: Id<"_storage">;
sha256: string;
contentType?: string;
};
type SkillBackupParams = {
skillId?: Id<"skills">;
versionId?: Id<"skillVersions">;
slug: string;
version: string;
isLatest?: boolean;
displayName: string;
ownerHandle: string;
files: BackupFile[];
publishedAt: number;
};
type PackageBackupParams = {
ownerHandle: string;
packageId: Id<"packages">;
releaseId: Id<"packageReleases">;
packageName: string;
normalizedName: string;
displayName: string;
family: "code-plugin" | "bundle-plugin";
version: string;
isLatest?: boolean;
publishedAt: number;
artifactKind?: "legacy-zip" | "npm-pack";
artifactFileName?: string;
artifactSha256?: string;
artifactSize?: number;
artifactFormat?: "tgz";
npmIntegrity?: string;
npmShasum?: string;
npmUnpackedSize?: number;
npmFileCount?: number;
runtimeId?: string;
sourceRepo?: string;
compatibility?: unknown;
extractedPackageJson?: unknown;
extractedPluginManifest?: unknown;
normalizedBundleManifest?: unknown;
files: Array<{ path: string; size: number; sha256: string }>;
};
export type RegistryArtifactBackupContext = RegistryArtifactBackupSettings;
export type RegistryArtifactBackupSettings = {
endpoint: string;
bucket: string;
accessKeyId: string;
secretAccessKey: string;
region: string;
skillsRoot: string;
packagesRoot: string;
};
export function isRegistryArtifactBackupConfigured() {
return Boolean(
(process.env.REGISTRY_BACKUP_S3_ENDPOINT || process.env.REGISTRY_BACKUP_R2_ACCOUNT_ID) &&
process.env.REGISTRY_BACKUP_BUCKET &&
process.env.REGISTRY_BACKUP_ACCESS_KEY_ID &&
process.env.REGISTRY_BACKUP_SECRET_ACCESS_KEY,
);
}
export function getRegistryArtifactBackupSettings(): RegistryArtifactBackupSettings {
const endpoint =
process.env.REGISTRY_BACKUP_S3_ENDPOINT ??
r2EndpointFromAccountId(process.env.REGISTRY_BACKUP_R2_ACCOUNT_ID);
if (!endpoint) {
throw new Error("REGISTRY_BACKUP_S3_ENDPOINT or REGISTRY_BACKUP_R2_ACCOUNT_ID is required");
}
const bucket = requiredEnv("REGISTRY_BACKUP_BUCKET");
const accessKeyId = requiredEnv("REGISTRY_BACKUP_ACCESS_KEY_ID");
const secretAccessKey = requiredEnv("REGISTRY_BACKUP_SECRET_ACCESS_KEY");
return {
endpoint,
bucket,
accessKeyId,
secretAccessKey,
region: process.env.REGISTRY_BACKUP_S3_REGION ?? "auto",
skillsRoot: process.env.REGISTRY_BACKUP_SKILLS_ROOT ?? DEFAULT_SKILLS_ROOT,
packagesRoot: process.env.REGISTRY_BACKUP_PACKAGES_ROOT ?? DEFAULT_PACKAGES_ROOT,
};
}
export function getRegistryArtifactBackupContext(): RegistryArtifactBackupContext {
return getRegistryArtifactBackupSettings();
}
export async function backupSkillVersionToObjectStorage(
ctx: Pick<ActionCtx, "storage">,
params: SkillBackupParams & { root?: string },
context: RegistryArtifactBackupContext = getRegistryArtifactBackupContext(),
) {
const planned = buildSkillVersionBackupManifest({
root: params.root ?? context.skillsRoot,
...params,
});
await runWithConcurrency(planned.fileObjects, getSkillFileUploadConcurrency(), async (file) => {
const blob = await readStorageBlob(ctx, file.storageId);
await putObject(context, file.key, new Uint8Array(await blob.arrayBuffer()), {
contentType: file.contentType,
});
});
await putJsonObject(context, planned.metaPath, planned.meta);
}
export async function backupPackageReleaseToObjectStorage(
ctx: Pick<ActionCtx, "storage">,
params: PackageBackupParams & { artifactStorageId: Id<"_storage">; root?: string },
context: RegistryArtifactBackupContext = getRegistryArtifactBackupContext(),
) {
const planned = buildPackageReleaseBackupManifest({
root: params.root ?? context.packagesRoot,
...params,
});
const artifact = await readStorageBlob(ctx, params.artifactStorageId);
await putObject(context, planned.artifactPath, new Uint8Array(await artifact.arrayBuffer()), {
contentType: packageArtifactContentType(params.artifactFormat),
});
await putJsonObject(context, planned.metaPath, planned.meta);
}
export async function fetchSkillVersionBackupMeta(
context: RegistryArtifactBackupContext,
ownerHandle: string,
slug: string,
version: string,
) {
const owner = normalizeOwner(ownerHandle);
const path = `${context.skillsRoot}/${owner}/${slug}/${encodeBackupPathSegment(
version,
)}/${META_FILENAME}`;
return getJsonObject<ReturnType<typeof buildSkillVersionBackupManifest>["meta"]>(context, path);
}
export async function fetchPackageReleaseBackupMeta(
context: RegistryArtifactBackupContext,
ownerHandle: string,
normalizedName: string,
version: string,
) {
const owner = normalizeOwner(ownerHandle);
const path = `${context.packagesRoot}/${owner}/${encodeBackupPathSegment(
normalizedName,
)}/${encodeBackupPathSegment(version)}/${META_FILENAME}`;
return getJsonObject<ReturnType<typeof buildPackageReleaseBackupManifest>["meta"]>(context, path);
}
export async function readRegistryArtifactBackupObject(
context: RegistryArtifactBackupContext,
key: string,
) {
const response = await signedFetch(context, "GET", key);
if (response.status === 404) return null;
if (!response.ok) {
const body = await response.text();
throw new Error(`Registry artifact backup GET ${key} failed: ${body}`);
}
return new Uint8Array(await response.arrayBuffer());
}
export function buildSkillVersionBackupManifest(params: SkillBackupParams & { root: string }) {
const owner = normalizeOwner(params.ownerHandle);
const versionSegment = encodeBackupPathSegment(params.version);
const skillRoot = `${params.root}/${owner}/${params.slug}`;
const versionRoot = `${skillRoot}/${versionSegment}`;
const metaPath = `${versionRoot}/${META_FILENAME}`;
const files = params.files.map((file) => {
if (!validateFilePath(file.path)) {
throw new Error(`Invalid skill backup file path: ${file.path}`);
}
return file;
});
const fileObjects = files.map((file) => ({
...file,
key: `${versionRoot}/${file.path}`,
}));
const meta = {
kind: "skillVersion" as const,
owner,
slug: params.slug,
displayName: params.displayName,
version: params.version,
isLatest: params.isLatest,
publishedAt: params.publishedAt,
restore: {
skillId: params.skillId,
versionId: params.versionId,
},
metadata: {
files: files.map(({ path, size, sha256, contentType }) => ({
path,
size,
sha256,
contentType,
})),
},
};
return {
skillRoot,
versionRoot,
metaPath,
fileObjects,
meta,
};
}
export function buildPackageReleaseBackupManifest(params: PackageBackupParams & { root: string }) {
const owner = normalizeOwner(params.ownerHandle);
const packageSegment = encodeBackupPathSegment(params.normalizedName || params.packageName);
const artifactFileName = validatePackageArtifactFileName(
params.artifactFileName ?? defaultPackageArtifactFileName(params),
);
const packageRoot = `${params.root}/${owner}/${packageSegment}`;
const releaseRoot = `${packageRoot}/${encodeBackupPathSegment(params.version)}`;
const meta = {
kind: "packageRelease" as const,
owner,
packageName: params.packageName,
normalizedName: params.normalizedName,
displayName: params.displayName,
family: params.family,
version: params.version,
isLatest: params.isLatest,
publishedAt: params.publishedAt,
runtimeId: params.runtimeId,
sourceRepo: params.sourceRepo,
artifactKind: params.artifactKind,
artifact: {
path: artifactFileName,
sha256: params.artifactSha256,
size: params.artifactSize,
format: params.artifactFormat,
npmIntegrity: params.npmIntegrity,
npmShasum: params.npmShasum,
npmUnpackedSize: params.npmUnpackedSize,
npmFileCount: params.npmFileCount,
},
restore: {
packageId: params.packageId,
releaseId: params.releaseId,
},
metadata: {
compatibility: params.compatibility,
extractedPackageJson: params.extractedPackageJson,
extractedPluginManifest: params.extractedPluginManifest,
normalizedBundleManifest: params.normalizedBundleManifest,
files: params.files,
},
};
return {
packageRoot,
releaseRoot,
artifactPath: `${releaseRoot}/${artifactFileName}`,
metaPath: `${releaseRoot}/${META_FILENAME}`,
meta,
};
}
export const __registryArtifactBackupTestInternals = {
encodeBackupPathSegment,
encodeObjectKey,
getSkillFileUploadConcurrency,
};
async function runWithConcurrency<T>(
items: T[],
concurrency: number,
worker: (item: T, index: number) => Promise<void>,
) {
if (items.length === 0) return;
let nextIndex = 0;
let firstError: unknown;
const workerCount = Math.min(concurrency, items.length);
async function runWorker() {
while (firstError === undefined) {
const index = nextIndex;
nextIndex += 1;
if (index >= items.length) return;
try {
await worker(items[index]!, index);
} catch (error) {
firstError ??= error;
return;
}
}
}
await Promise.allSettled(Array.from({ length: workerCount }, () => runWorker()));
if (firstError !== undefined) throw firstError;
}
function getSkillFileUploadConcurrency() {
const raw = process.env.REGISTRY_BACKUP_SKILL_FILE_UPLOAD_CONCURRENCY;
const parsed = raw ? Number.parseInt(raw, 10) : DEFAULT_SKILL_FILE_UPLOAD_CONCURRENCY;
if (!Number.isFinite(parsed) || parsed < 1) return DEFAULT_SKILL_FILE_UPLOAD_CONCURRENCY;
return Math.min(parsed, MAX_SKILL_FILE_UPLOAD_CONCURRENCY);
}
export function normalizeOwner(value: string) {
const normalized = value
.trim()
.toLowerCase()
.replace(/^@+/, "")
.replace(/[^a-z0-9._-]/g, "-")
.replace(/-+/g, "-")
.replace(/^[._-]+|[._-]+$/g, "");
return normalized || "unknown";
}
function encodeBackupPathSegment(value: string) {
return encodeURIComponent(value.trim()).replace(/\./g, "%2E");
}
function normalizePackagePathSegment(value: string) {
return normalizeOwner(value.replace(/^@/, "").replace("/", "-"));
}
function defaultPackageArtifactFileName(
params: Pick<PackageBackupParams, "normalizedName" | "version">,
) {
return `${normalizePackagePathSegment(params.normalizedName)}-${encodeBackupPathSegment(
params.version,
)}.tgz`;
}
function validatePackageArtifactFileName(value: string) {
const artifactFileName = value.trim();
if (
!artifactFileName ||
artifactFileName === "." ||
artifactFileName === ".." ||
artifactFileName.includes("/") ||
artifactFileName.includes("\\") ||
artifactFileName.includes("\0")
) {
throw new Error("Invalid package backup artifact filename");
}
return artifactFileName;
}
async function readStorageBlob(ctx: Pick<ActionCtx, "storage">, storageId: Id<"_storage">) {
const blob = await ctx.storage.get(storageId);
if (!blob) throw new Error("File missing in storage");
return blob;
}
async function putJsonObject(context: RegistryArtifactBackupContext, key: string, value: unknown) {
await putObject(context, key, `${JSON.stringify(value, null, 2)}\n`, {
contentType: "application/json; charset=utf-8",
});
}
async function getJsonObject<T>(context: RegistryArtifactBackupContext, key: string) {
const response = await signedFetch(context, "GET", key);
if (response.status === 404) return null;
if (!response.ok) {
const body = await response.text();
throw new Error(`Registry artifact backup GET ${key} failed: ${body}`);
}
return (await response.json()) as T;
}
async function putObject(
context: RegistryArtifactBackupContext,
key: string,
body: string | Uint8Array,
options: {
contentType?: string;
} = {},
) {
const response = await signedFetch(context, "PUT", key, body, options);
if (!response.ok) {
const responseBody = await response.text();
throw new Error(`Registry artifact backup PUT ${key} failed: ${responseBody}`);
}
}
async function signedFetch(
context: RegistryArtifactBackupContext,
method: "GET" | "PUT",
key: string,
body?: string | Uint8Array,
options: { contentType?: string } = {},
) {
const now = new Date();
const bodyBytes = body === undefined ? new Uint8Array() : toBytes(body);
const payloadHash = sha256Hex(bodyBytes);
const url = objectUrl(context, key);
const headers = new Headers();
headers.set("host", url.host);
headers.set("x-amz-content-sha256", payloadHash);
headers.set("x-amz-date", amzDate(now));
if (options.contentType) headers.set("content-type", options.contentType);
headers.set(
"authorization",
authorizationHeader(context, method, url, headers, payloadHash, now),
);
const init: RequestInit = { method, headers };
if (method === "PUT") {
init.body = toArrayBuffer(bodyBytes);
}
return fetch(url, init);
}
function toArrayBuffer(bytes: Uint8Array) {
return bytes.buffer.slice(bytes.byteOffset, bytes.byteOffset + bytes.byteLength) as ArrayBuffer;
}
function authorizationHeader(
context: RegistryArtifactBackupContext,
method: string,
url: URL,
headers: Headers,
payloadHash: string,
now: Date,
) {
const date = amzDate(now).slice(0, 8);
const credentialScope = `${date}/${context.region}/s3/aws4_request`;
const signedHeaders = Array.from(headers.keys())
.map((name) => name.toLowerCase())
.sort()
.join(";");
const canonicalHeaders = signedHeaders
.split(";")
.map((name) => `${name}:${headers.get(name)?.trim() ?? ""}\n`)
.join("");
const canonicalRequest = [
method,
url.pathname,
url.search.slice(1),
canonicalHeaders,
signedHeaders,
payloadHash,
].join("\n");
const stringToSign = [
"AWS4-HMAC-SHA256",
amzDate(now),
credentialScope,
sha256Hex(canonicalRequest),
].join("\n");
const signingKey = hmac(
hmac(hmac(hmac(`AWS4${context.secretAccessKey}`, date), context.region), "s3"),
"aws4_request",
);
const signature = hmacHex(signingKey, stringToSign);
return `AWS4-HMAC-SHA256 Credential=${context.accessKeyId}/${credentialScope}, SignedHeaders=${signedHeaders}, Signature=${signature}`;
}
function objectUrl(context: RegistryArtifactBackupContext, key: string) {
const endpoint = context.endpoint.replace(/\/+$/, "");
return new URL(`${endpoint}/${encodePathSegment(context.bucket)}/${encodeObjectKey(key)}`);
}
function encodeObjectKey(key: string) {
return key.split("/").map(encodePathSegment).join("/");
}
function encodePathSegment(value: string) {
return encodeURIComponent(value).replace(
/[!'()*]/g,
(char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`,
);
}
function amzDate(date: Date) {
return date.toISOString().replace(/[:-]|\.\d{3}/g, "");
}
function toBytes(value: string | Uint8Array) {
return typeof value === "string" ? new TextEncoder().encode(value) : value;
}
function sha256Hex(value: string | Uint8Array) {
return createHash("sha256").update(value).digest("hex");
}
function hmac(key: string | Buffer, value: string) {
return createHmac("sha256", key).update(value).digest();
}
function hmacHex(key: Buffer, value: string) {
return createHmac("sha256", key).update(value).digest("hex");
}
function packageArtifactContentType(format: PackageBackupParams["artifactFormat"]) {
return format === "tgz" ? "application/gzip" : "application/octet-stream";
}
function r2EndpointFromAccountId(accountId: string | undefined) {
return accountId ? `https://${accountId}.r2.cloudflarestorage.com` : undefined;
}
function requiredEnv(name: string) {
const value = process.env[name];
if (!value) throw new Error(`${name} is required`);
return value;
}
+2 -2
View File
@@ -230,8 +230,8 @@ export const RETENTION_POLICIES = {
retention: "Slug reservation cooldown.",
}),
reservedHandles: permanent("Reserved handles are explicit policy records until released."),
registryArtifactBackupSyncState: permanent("Registry artifact backup cursor state."),
registryArtifactBackupJobs: permanent("Registry artifact backup job history and retry state."),
registryArtifactBackupSyncState: permanent("Legacy registry artifact backup cursor state."),
registryArtifactBackupJobs: permanent("Legacy registry artifact backup job history."),
userSkillInstalls: permanent("Current user install records."),
skillOwnershipTransfers: ephemeral("Ownership transfer invitations expire.", {
expirationField: "expiresAt",
-3
View File
@@ -78,7 +78,6 @@ description: Automation workflow for recurring reports.
{
bypassGitHubAccountAge: true,
bypassQualityGate: true,
skipBackup: true,
skipWebhook: true,
},
);
@@ -157,7 +156,6 @@ description: Research helper for literature reviews.
{
bypassGitHubAccountAge: true,
bypassQualityGate: true,
skipBackup: true,
skipWebhook: true,
},
);
@@ -236,7 +234,6 @@ description: Research helper for literature reviews.
{
bypassGitHubAccountAge: true,
bypassQualityGate: true,
skipBackup: true,
skipWebhook: true,
},
);
-46
View File
@@ -99,7 +99,6 @@ export type PublishOptions = {
bypassGitHubAccountAge?: boolean;
bypassNewSkillRateLimit?: boolean;
bypassQualityGate?: boolean;
skipBackup?: boolean;
skipWebhook?: boolean;
ownerPublisherId?: Id<"publishers">;
sourceOwnerPublisherId?: Id<"publishers">;
@@ -143,10 +142,6 @@ export async function publishVersionForUser(
migrateOwner: options.migrateOwner,
})) as Doc<"skills"> | null;
const isNewSkill = !existingSkill;
const publishedVersionIsLatest = shouldPublishVersionBecomeLatest(
version,
existingSkill?.latestVersionSummary?.version,
);
// For new skills, enforce the full write-path rules (length, pattern,
// reserved-word blocklist). For existing skills the slug is already
@@ -390,35 +385,6 @@ export async function publishVersionForUser(
const ownerHandle =
targetPublisher?.handle ?? owner?.handle ?? owner?.displayName ?? owner?.name ?? "unknown";
if (!options.skipBackup) {
await ctx.scheduler
.runAfter(0, internal.registryArtifactBackupsNode.backupSkillForPublishInternal, {
skillId: publishResult.skillId,
versionId: publishResult.versionId,
slug,
version,
isLatest: publishedVersionIsLatest,
displayName,
ownerHandle,
files: publishFiles,
publishedAt: Date.now(),
})
.catch((error) => {
const message = errorMessage(error);
console.error("registry artifact backup scheduling failed", error);
return ctx
.runMutation(internal.registryArtifactBackups.enqueueRegistryArtifactBackupJobInternal, {
targetKind: "skillVersion",
skillVersionId: publishResult.versionId,
reason: "publish",
error: message,
})
.catch((enqueueError) => {
console.error("registry artifact backup retry enqueue failed", enqueueError);
});
});
}
if (!options.skipWebhook) {
void schedulePublishWebhook(ctx, {
slug,
@@ -431,18 +397,6 @@ export async function publishVersionForUser(
return publishResult;
}
function errorMessage(error: unknown) {
return error instanceof Error ? error.message : String(error);
}
function shouldPublishVersionBecomeLatest(version: string, previousLatestVersion?: string) {
return (
!previousLatestVersion ||
!semver.valid(previousLatestVersion) ||
semver.gt(version, previousLatestVersion)
);
}
function mergeSourceIntoMetadata(
metadata: unknown,
source: PublishVersionArgs["source"],
+1 -39
View File
@@ -7161,16 +7161,7 @@ describe("packages public queries", () => {
runMutation,
runAction: vi.fn(async () => makeCleanPackageInspectorResult()),
scheduler: {
runAfter: vi.fn(async (_delayMs: number, _ref: unknown, args: unknown) => {
if (
typeof args === "object" &&
args !== null &&
"artifactStorageId" in args &&
args.artifactStorageId === "storage:clawpack"
) {
throw new Error("scheduler unavailable");
}
}),
runAfter: vi.fn(async () => {}),
},
storage: {
get: vi.fn(async (storageId: string) => {
@@ -7274,35 +7265,6 @@ describe("packages public queries", () => {
]),
}),
);
expect(ctx.scheduler.runAfter).toHaveBeenCalledWith(
0,
expect.anything(),
expect.objectContaining({
releaseId: "releases:demo-1",
packageName: "demo-plugin",
artifactStorageId: "storage:clawpack",
artifactSha256: "clawpack",
artifactFileName: "demo-plugin-1.0.0.tgz",
}),
);
await vi.waitFor(() => {
const retryArgs = runMutation.mock.calls
.map(([, args]) => args)
.find(
(args): args is Record<string, unknown> =>
typeof args === "object" &&
args !== null &&
"targetKind" in args &&
args.targetKind === "packageRelease",
);
expect(retryArgs).toEqual(
expect.objectContaining({
packageReleaseId: "releases:demo-1",
reason: "publish",
error: "scheduler unavailable",
}),
);
});
});
it("rejects trusted publish tokens after trusted publisher rotation or deletion", async () => {
-75
View File
@@ -391,12 +391,6 @@ const internalRefs = internal as unknown as {
packageInspectorNode: {
runPackageInspectorForPublishInternal: unknown;
};
registryArtifactBackupsNode: {
backupPackageForPublishInternal: unknown;
};
registryArtifactBackups: {
enqueueRegistryArtifactBackupJobInternal: unknown;
};
packagePublishTokens: {
createInternal: unknown;
getByIdInternal: unknown;
@@ -6989,78 +6983,9 @@ async function publishPackageImpl(
source: "publish",
});
if (payload.artifact?.storageId) {
const backupIsLatest = (
payload.tags?.map((tag: string) => tag.trim()).filter(Boolean) ?? ["latest"]
).includes("latest");
const backupOwner =
ownerPublisher ??
((await runQueryRef<Doc<"users"> | null>(ctx, internalRefs.users.getByIdInternal, {
userId: ownerUserId,
})) as Doc<"users"> | null);
const ownerHandle = backupOwner?.handle ?? String(ownerPublisherId ?? ownerUserId);
await runAfterRef(
ctx,
0,
internalRefs.registryArtifactBackupsNode.backupPackageForPublishInternal,
{
ownerHandle,
packageId: publishResult.packageId,
releaseId: publishResult.releaseId,
packageName: name,
normalizedName: name,
displayName,
family,
version,
isLatest: backupIsLatest,
publishedAt: Date.now(),
artifactKind: payload.artifact.kind ?? "legacy-zip",
artifactStorageId: payload.artifact.storageId,
artifactFileName: payload.artifact.npmTarballName,
artifactSha256: payload.artifact.sha256,
artifactSize: payload.artifact.size,
artifactFormat: payload.artifact.format,
npmIntegrity: payload.artifact.npmIntegrity,
npmShasum: payload.artifact.npmShasum,
npmUnpackedSize: payload.artifact.npmUnpackedSize,
npmFileCount: payload.artifact.npmFileCount,
runtimeId: codeArtifacts?.runtimeId ?? bundleArtifacts?.runtimeId,
sourceRepo: effectiveSource?.repo || effectiveSource?.url,
compatibility: codeArtifacts?.compatibility ?? bundleArtifacts?.compatibility,
extractedPackageJson: storedPackageJson,
extractedPluginManifest: storedPluginManifest,
normalizedBundleManifest: family === "bundle-plugin" ? storedBundleManifest : undefined,
files: files.map((file) => ({
path: file.path,
size: file.size,
sha256: file.sha256,
})),
},
).catch((error) => {
const message = errorMessage(error);
console.error("registry artifact package backup scheduling failed", error);
return runMutationRef(
ctx,
internalRefs.registryArtifactBackups.enqueueRegistryArtifactBackupJobInternal,
{
targetKind: "packageRelease",
packageReleaseId: publishResult.releaseId,
reason: "publish",
error: message,
},
).catch((enqueueError) => {
console.error("registry artifact package backup retry enqueue failed", enqueueError);
});
});
}
return inspectorFindings.length > 0 ? { ...publishResult, inspectorFindings } : publishResult;
}
function errorMessage(error: unknown) {
return error instanceof Error ? error.message : String(error);
}
function toPackageInspectorPublishResponseFinding(
finding: PackageInspectorFinding,
metadata: PackageInspectorPublishResult["metadata"],
File diff suppressed because it is too large Load Diff
-816
View File
@@ -1,816 +0,0 @@
import { v } from "convex/values";
import { internal } from "./_generated/api";
import type { Doc, Id } from "./_generated/dataModel";
import type { MutationCtx, QueryCtx } from "./_generated/server";
import { action, internalMutation, internalQuery } from "./functions";
import { assertRole, requireUserFromAction } from "./lib/access";
import { isPublicSkillDoc } from "./lib/globalStats";
import { getOwnerPublisher } from "./lib/publishers";
const DEFAULT_BATCH_SIZE = 50;
const MAX_BATCH_SIZE = 200;
const SYNC_STATE_KEY = "default";
const PACKAGE_SYNC_STATE_KEY = "packageReleases";
const RETRY_LEASE_KEY = "retryLease";
const MAX_BACKUP_JOB_ERROR_LENGTH = 4000;
const DEFAULT_BACKUP_HEALTH_SAMPLE_LIMIT = 500;
const MAX_BACKUP_HEALTH_SAMPLE_LIMIT = 1000;
const DEFAULT_BACKUP_JOB_LIMIT = 25;
const MAX_BACKUP_JOB_LIMIT = 500;
const DEFAULT_BACKUP_JOB_REPAIR_ATTEMPTS = 16;
const DEFAULT_RETRY_LEASE_TTL_MS = 20 * 60 * 1000;
const MAX_RETRY_LEASE_TTL_MS = 60 * 60 * 1000;
const DEFAULT_BACKUP_JOB_LEASE_TTL_MS = 20 * 60 * 1000;
const MAX_BACKUP_JOB_LEASE_TTL_MS = 60 * 60 * 1000;
type BackupPageItem =
| {
kind: "ok";
skillId: Id<"skills">;
versionId: Id<"skillVersions">;
slug: string;
displayName: string;
version: string;
isLatest: boolean;
ownerHandle: string;
publishedAt: number;
}
| { kind: "missingOwner"; skillId: Id<"skills">; ownerUserId: Id<"users"> };
type BackupPageResult = {
items: BackupPageItem[];
cursor: string | null;
isDone: boolean;
};
type PackageBackupPageItem =
| {
kind: "ok";
packageId: Id<"packages">;
releaseId: Id<"packageReleases">;
ownerHandle: string;
packageName: string;
normalizedName: string;
displayName: string;
family: "code-plugin" | "bundle-plugin";
version: string;
isLatest: boolean;
publishedAt: number;
artifactKind?: "legacy-zip" | "npm-pack";
artifactStorageId: Id<"_storage">;
artifactFileName?: string;
artifactSha256?: string;
artifactSize?: number;
artifactFormat?: "tgz";
npmIntegrity?: string;
npmShasum?: string;
npmUnpackedSize?: number;
npmFileCount?: number;
runtimeId?: string;
sourceRepo?: string;
compatibility?: unknown;
extractedPackageJson?: unknown;
extractedPluginManifest?: unknown;
normalizedBundleManifest?: unknown;
files: Array<{ path: string; size: number; sha256: string }>;
}
| { kind: "missingPackage"; releaseId: Id<"packageReleases">; packageId: Id<"packages"> }
| { kind: "missingOwner"; releaseId: Id<"packageReleases">; packageId: Id<"packages"> }
| { kind: "missingArtifact"; releaseId: Id<"packageReleases">; packageId: Id<"packages"> };
type PackageBackupPageResult = {
items: PackageBackupPageItem[];
cursor: string | null;
isDone: boolean;
};
type BackupSyncState = {
cursor: string | null;
isDone: boolean;
};
export type SeedRegistryArtifactBackupsResult = {
stats: {
skillsScanned: number;
skillsSkipped: number;
skillsBackedUp: number;
skillsMissingVersion: number;
skillsMissingOwner: number;
packagesScanned: number;
packagesSkipped: number;
packagesBackedUp: number;
packagesMissingArtifact: number;
packagesMissingPackage: number;
packagesMissingOwner: number;
skillsEnqueued: number;
packagesEnqueued: number;
retryJobsProcessed: number;
retryJobsSucceeded: number;
retryJobsFailed: number;
staleJobs: number;
exhaustedJobs: number;
errors: number;
};
cursor: string | null;
packageCursor: string | null;
skillsIsDone: boolean;
packageIsDone: boolean;
isDone: boolean;
};
export const getRegistryArtifactBackupPageInternal = internalQuery({
args: {
cursor: v.optional(v.string()),
batchSize: v.optional(v.number()),
},
handler: async (ctx, args): Promise<BackupPageResult> => {
const batchSize = clampInt(args.batchSize ?? DEFAULT_BATCH_SIZE, 1, MAX_BATCH_SIZE);
let pageResult;
try {
pageResult = await ctx.db
.query("skillVersions")
.withIndex("by_active_created", (q) => q.eq("softDeletedAt", undefined))
.order("asc")
.paginate({ cursor: args.cursor ?? null, numItems: batchSize });
} catch (error) {
if (!args.cursor || !isStaleCursorError(error)) throw error;
pageResult = await ctx.db
.query("skillVersions")
.withIndex("by_active_created", (q) => q.eq("softDeletedAt", undefined))
.order("asc")
.paginate({ cursor: null, numItems: batchSize });
}
const items: BackupPageItem[] = [];
for (const version of pageResult.page) {
const item = await toSkillVersionBackupPageItem(ctx, version);
if (item) items.push(item);
}
return { items, cursor: pageResult.continueCursor, isDone: pageResult.isDone };
},
});
export const getPackageRegistryArtifactBackupPageInternal = internalQuery({
args: {
cursor: v.optional(v.string()),
batchSize: v.optional(v.number()),
},
handler: async (ctx, args): Promise<PackageBackupPageResult> => {
const batchSize = clampInt(args.batchSize ?? DEFAULT_BATCH_SIZE, 1, MAX_BATCH_SIZE);
const pageResult = await ctx.db
.query("packageReleases")
.withIndex("by_active_created", (q) => q.eq("softDeletedAt", undefined))
.order("asc")
.paginate({ cursor: args.cursor ?? null, numItems: batchSize });
const items: PackageBackupPageItem[] = [];
for (const release of pageResult.page) {
const item = await toPackageBackupPageItem(ctx, release);
if (item) items.push(item);
}
return { items, cursor: pageResult.continueCursor, isDone: pageResult.isDone };
},
});
async function toSkillVersionBackupPageItem(
ctx: Parameters<typeof getOwnerPublisher>[0],
version: Doc<"skillVersions">,
): Promise<BackupPageItem | null> {
const skill = await ctx.db.get(version.skillId);
if (!skill || !isPublicSkillDoc(skill)) return null;
const owner = await getOwnerPublisher(ctx, {
ownerPublisherId: skill.ownerPublisherId,
ownerUserId: skill.ownerUserId,
});
if (!owner || owner.deletedAt || owner.deactivatedAt) {
return { kind: "missingOwner", skillId: skill._id, ownerUserId: skill.ownerUserId };
}
return {
kind: "ok",
skillId: skill._id,
versionId: version._id,
slug: skill.slug,
displayName: skill.displayName,
version: version.version,
isLatest: skill.latestVersionId === version._id,
ownerHandle: owner.handle ?? String(skill.ownerPublisherId ?? skill.ownerUserId),
publishedAt: version.createdAt,
};
}
async function toPackageBackupPageItem(
ctx: Parameters<typeof getOwnerPublisher>[0],
release: Doc<"packageReleases">,
): Promise<PackageBackupPageItem | null> {
const pkg = await ctx.db.get(release.packageId);
if (!pkg || pkg.softDeletedAt) {
return { kind: "missingPackage", releaseId: release._id, packageId: release.packageId };
}
if (pkg.family !== "code-plugin" && pkg.family !== "bundle-plugin") return null;
if (!release.clawpackStorageId) {
return { kind: "missingArtifact", releaseId: release._id, packageId: release.packageId };
}
const owner = await getOwnerPublisher(ctx, {
ownerPublisherId: pkg.ownerPublisherId,
ownerUserId: pkg.ownerUserId,
});
if (!owner || owner.deletedAt || owner.deactivatedAt) {
return { kind: "missingOwner", releaseId: release._id, packageId: release.packageId };
}
return {
kind: "ok",
packageId: pkg._id,
releaseId: release._id,
ownerHandle: owner.handle,
packageName: pkg.name,
normalizedName: pkg.normalizedName,
displayName: pkg.displayName,
family: pkg.family,
version: release.version,
isLatest: pkg.latestReleaseId === release._id,
publishedAt: release.createdAt,
artifactKind: release.artifactKind,
artifactStorageId: release.clawpackStorageId,
artifactFileName: release.npmTarballName,
artifactSha256: release.clawpackSha256,
artifactSize: release.clawpackSize,
artifactFormat: release.clawpackFormat,
npmIntegrity: release.npmIntegrity,
npmShasum: release.npmShasum,
npmUnpackedSize: release.npmUnpackedSize,
npmFileCount: release.npmFileCount,
runtimeId: release.runtimeId,
sourceRepo: release.sourceRepo,
compatibility: release.compatibility,
extractedPackageJson: release.extractedPackageJson,
extractedPluginManifest: release.extractedPluginManifest,
normalizedBundleManifest: release.normalizedBundleManifest,
files: release.files.map((file) => ({
path: file.path,
size: file.size,
sha256: file.sha256,
})),
};
}
function isStaleCursorError(error: unknown) {
const message =
typeof error === "string"
? error
: error && typeof error === "object" && "message" in error
? String((error as { message?: unknown }).message)
: "";
return (
message.includes("Failed to parse cursor") ||
message.includes("cursor is from a different query")
);
}
export const getRegistryArtifactBackupSyncStateInternal = internalQuery({
args: {},
handler: async (ctx): Promise<BackupSyncState> => {
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", SYNC_STATE_KEY))
.unique();
return { cursor: state?.cursor ?? null, isDone: state?.isDone === true };
},
});
export const setRegistryArtifactBackupSyncStateInternal = internalMutation({
args: {
cursor: v.optional(v.string()),
isDone: v.optional(v.boolean()),
},
handler: async (ctx, args) => {
const now = Date.now();
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", SYNC_STATE_KEY))
.unique();
if (!state) {
await ctx.db.insert("registryArtifactBackupSyncState", {
key: SYNC_STATE_KEY,
cursor: args.cursor,
isDone: args.isDone ?? false,
updatedAt: now,
});
return { ok: true as const };
}
await ctx.db.patch(state._id, {
cursor: args.cursor,
isDone: args.isDone ?? false,
updatedAt: now,
});
return { ok: true as const };
},
});
export const getPackageRegistryArtifactBackupSyncStateInternal = internalQuery({
args: {},
handler: async (ctx): Promise<BackupSyncState> => {
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", PACKAGE_SYNC_STATE_KEY))
.unique();
return { cursor: state?.cursor ?? null, isDone: state?.isDone === true };
},
});
export const setPackageRegistryArtifactBackupSyncStateInternal = internalMutation({
args: {
cursor: v.optional(v.string()),
isDone: v.optional(v.boolean()),
},
handler: async (ctx, args) => {
const now = Date.now();
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", PACKAGE_SYNC_STATE_KEY))
.unique();
if (!state) {
await ctx.db.insert("registryArtifactBackupSyncState", {
key: PACKAGE_SYNC_STATE_KEY,
cursor: args.cursor,
isDone: args.isDone ?? false,
updatedAt: now,
});
return { ok: true as const };
}
await ctx.db.patch(state._id, {
cursor: args.cursor,
isDone: args.isDone ?? false,
updatedAt: now,
});
return { ok: true as const };
},
});
export async function tryAcquireRegistryArtifactBackupRetryLeaseHandler(
ctx: Pick<MutationCtx, "db">,
args: { now?: number; token: string; ttlMs?: number },
) {
const now = args.now ?? Date.now();
const ttlMs = clampInt(args.ttlMs ?? DEFAULT_RETRY_LEASE_TTL_MS, 1_000, MAX_RETRY_LEASE_TTL_MS);
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", RETRY_LEASE_KEY))
.unique();
if (state?.cursor && state.updatedAt + ttlMs > now) {
return { acquired: false as const, holderUpdatedAt: state.updatedAt };
}
if (!state) {
await ctx.db.insert("registryArtifactBackupSyncState", {
key: RETRY_LEASE_KEY,
cursor: args.token,
updatedAt: now,
});
return { acquired: true as const };
}
await ctx.db.patch(state._id, {
cursor: args.token,
updatedAt: now,
});
return { acquired: true as const };
}
export const tryAcquireRegistryArtifactBackupRetryLeaseInternal = internalMutation({
args: {
now: v.optional(v.number()),
token: v.string(),
ttlMs: v.optional(v.number()),
},
handler: tryAcquireRegistryArtifactBackupRetryLeaseHandler,
});
export async function releaseRegistryArtifactBackupRetryLeaseHandler(
ctx: Pick<MutationCtx, "db">,
args: { now?: number; token: string },
) {
const now = args.now ?? Date.now();
const state = await ctx.db
.query("registryArtifactBackupSyncState")
.withIndex("by_key", (q) => q.eq("key", RETRY_LEASE_KEY))
.unique();
if (!state || state.cursor !== args.token) return { released: false as const };
await ctx.db.patch(state._id, {
cursor: undefined,
updatedAt: now,
});
return { released: true as const };
}
export const releaseRegistryArtifactBackupRetryLeaseInternal = internalMutation({
args: {
now: v.optional(v.number()),
token: v.string(),
},
handler: releaseRegistryArtifactBackupRetryLeaseHandler,
});
const registryArtifactBackupTargetKindValidator = v.union(
v.literal("skillVersion"),
v.literal("packageRelease"),
);
const registryArtifactBackupReasonValidator = v.union(
v.literal("publish"),
v.literal("seed"),
v.literal("retry"),
v.literal("sync"),
);
const registryArtifactBackupStatusValidator = v.union(
v.literal("pending"),
v.literal("running"),
v.literal("succeeded"),
v.literal("exhausted"),
v.literal("missingArtifact"),
);
export const enqueueRegistryArtifactBackupJobInternal = internalMutation({
args: {
targetKind: registryArtifactBackupTargetKindValidator,
skillVersionId: v.optional(v.id("skillVersions")),
packageReleaseId: v.optional(v.id("packageReleases")),
reason: registryArtifactBackupReasonValidator,
status: v.optional(registryArtifactBackupStatusValidator),
preserveTerminal: v.optional(v.boolean()),
error: v.optional(v.string()),
now: v.optional(v.number()),
},
handler: enqueueRegistryArtifactBackupJobHandler,
});
export async function claimRegistryArtifactBackupJobsHandler(
ctx: Pick<MutationCtx, "db">,
args: {
now?: number;
limit?: number;
leaseToken: string;
leaseTtlMs?: number;
forceDue?: boolean;
includeExhaustedRepair?: boolean;
maxRepairAttempts?: number;
},
) {
const now = args.now ?? Date.now();
const limit = clampInt(args.limit ?? DEFAULT_BACKUP_JOB_LIMIT, 1, MAX_BACKUP_JOB_LIMIT);
const leaseTtlMs = clampInt(
args.leaseTtlMs ?? DEFAULT_BACKUP_JOB_LEASE_TTL_MS,
1_000,
MAX_BACKUP_JOB_LEASE_TTL_MS,
);
const pending = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_nextRunAt", (q) => {
const byStatus = q.eq("status", "pending");
return args.forceDue ? byStatus : byStatus.lte("nextRunAt", now);
})
.take(limit);
let claimed = pending;
if (claimed.length < limit) {
const expiredRunning = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_leaseExpiresAt", (q) =>
q.eq("status", "running").lte("leaseExpiresAt", now),
)
.take(limit - claimed.length);
claimed = [...claimed, ...expiredRunning];
}
if (args.includeExhaustedRepair && claimed.length < limit) {
const maxRepairAttempts = Math.max(
1,
Math.floor(args.maxRepairAttempts ?? DEFAULT_BACKUP_JOB_REPAIR_ATTEMPTS),
);
const exhausted = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_attempts", (q) =>
q.eq("status", "exhausted").lt("attempts", maxRepairAttempts),
)
.take(limit - claimed.length);
claimed = [...claimed, ...exhausted];
}
const leaseExpiresAt = now + leaseTtlMs;
for (const job of claimed) {
await ctx.db.patch(job._id, {
status: "running",
leaseToken: args.leaseToken,
leaseExpiresAt,
claimedAt: now,
lastAttemptAt: now,
updatedAt: now,
});
}
return claimed;
}
export const claimRegistryArtifactBackupJobsInternal = internalMutation({
args: {
now: v.optional(v.number()),
limit: v.optional(v.number()),
leaseToken: v.string(),
leaseTtlMs: v.optional(v.number()),
forceDue: v.optional(v.boolean()),
includeExhaustedRepair: v.optional(v.boolean()),
maxRepairAttempts: v.optional(v.number()),
},
handler: claimRegistryArtifactBackupJobsHandler,
});
export async function markRegistryArtifactBackupJobSucceededHandler(
ctx: Pick<MutationCtx, "db">,
args: { jobId: Id<"registryArtifactBackupJobs">; leaseToken?: string; now?: number },
) {
const now = args.now ?? Date.now();
const job = await ctx.db.get(args.jobId);
if (!job) return { missing: true as const, stale: false as const };
if (args.leaseToken && (job.status !== "running" || job.leaseToken !== args.leaseToken)) {
return { missing: false as const, stale: true as const };
}
await ctx.db.patch(args.jobId, {
status: "succeeded",
completedAt: now,
lastError: undefined,
leaseToken: undefined,
leaseExpiresAt: undefined,
claimedAt: undefined,
updatedAt: now,
});
return { missing: false as const, stale: false as const };
}
export const markRegistryArtifactBackupJobSucceededInternal = internalMutation({
args: {
jobId: v.id("registryArtifactBackupJobs"),
leaseToken: v.optional(v.string()),
now: v.optional(v.number()),
},
handler: markRegistryArtifactBackupJobSucceededHandler,
});
export const markRegistryArtifactBackupJobFailedInternal = internalMutation({
args: {
jobId: v.id("registryArtifactBackupJobs"),
error: v.string(),
leaseToken: v.optional(v.string()),
now: v.optional(v.number()),
maxAttempts: v.optional(v.number()),
},
handler: async (ctx, args) => {
const now = args.now ?? Date.now();
const maxAttempts = Math.max(1, Math.floor(args.maxAttempts ?? 8));
const job = await ctx.db.get(args.jobId);
if (!job) return { missing: true as const };
if (args.leaseToken && (job.status !== "running" || job.leaseToken !== args.leaseToken)) {
return { missing: false as const, stale: true as const };
}
const attempts = job.attempts + 1;
const exhausted = attempts >= maxAttempts;
await ctx.db.patch(args.jobId, {
status: exhausted ? "exhausted" : "pending",
attempts,
lastAttemptAt: now,
lastError: truncateBackupJobError(args.error),
nextRunAt: exhausted ? now : now + retryDelayMs(attempts),
exhaustedAt: exhausted ? now : undefined,
leaseToken: undefined,
leaseExpiresAt: undefined,
claimedAt: undefined,
updatedAt: now,
});
return { missing: false as const, stale: false as const, exhausted, attempts };
},
});
export const getDueRegistryArtifactBackupJobsInternal = internalQuery({
args: {
includeExhaustedRepair: v.optional(v.boolean()),
ignoreNextRunAt: v.optional(v.boolean()),
maxRepairAttempts: v.optional(v.number()),
now: v.optional(v.number()),
limit: v.optional(v.number()),
},
handler: async (ctx, args) => {
const now = args.now ?? Date.now();
const limit = clampInt(args.limit ?? DEFAULT_BACKUP_JOB_LIMIT, 1, MAX_BACKUP_JOB_LIMIT);
const pending = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_nextRunAt", (q) => {
const byStatus = q.eq("status", "pending");
return args.ignoreNextRunAt ? byStatus : byStatus.lte("nextRunAt", now);
})
.take(limit);
if (!args.includeExhaustedRepair || pending.length >= limit) return pending;
const maxRepairAttempts = Math.max(
1,
Math.floor(args.maxRepairAttempts ?? DEFAULT_BACKUP_JOB_REPAIR_ATTEMPTS),
);
const remaining = limit - pending.length;
const exhausted = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_attempts", (q) =>
q.eq("status", "exhausted").lt("attempts", maxRepairAttempts),
)
.take(remaining);
return [...pending, ...exhausted];
},
});
export const getRegistryArtifactBackupHealthInternal = internalQuery({
args: {
now: v.optional(v.number()),
staleAfterMs: v.optional(v.number()),
sampleLimit: v.optional(v.number()),
},
handler: getRegistryArtifactBackupHealthHandler,
});
export const seedRegistryArtifactBackups: ReturnType<typeof action> = action({
args: {
dryRun: v.optional(v.boolean()),
batchSize: v.optional(v.number()),
maxBatches: v.optional(v.number()),
queueOnly: v.optional(v.boolean()),
resetCursor: v.optional(v.boolean()),
},
handler: async (ctx, args): Promise<SeedRegistryArtifactBackupsResult> => {
const { user } = await requireUserFromAction(ctx);
assertRole(user, ["admin"]);
if (args.resetCursor && !args.dryRun) {
await ctx.runMutation(
internal.registryArtifactBackups.setRegistryArtifactBackupSyncStateInternal,
{
cursor: undefined,
isDone: false,
},
);
await ctx.runMutation(
internal.registryArtifactBackups.setPackageRegistryArtifactBackupSyncStateInternal,
{
cursor: undefined,
isDone: false,
},
);
}
return ctx.runAction(internal.registryArtifactBackupsNode.seedRegistryArtifactBackupsInternal, {
dryRun: args.dryRun,
batchSize: args.batchSize,
maxBatches: args.maxBatches,
queueOnly: args.queueOnly,
}) as Promise<SeedRegistryArtifactBackupsResult>;
},
});
function clampInt(value: number, min: number, max: number) {
return Math.max(min, Math.min(max, Math.floor(value)));
}
export async function enqueueRegistryArtifactBackupJobHandler(
ctx: Pick<MutationCtx, "db">,
args: {
targetKind: "skillVersion" | "packageRelease";
skillVersionId?: Id<"skillVersions">;
packageReleaseId?: Id<"packageReleases">;
reason: "publish" | "seed" | "retry" | "sync";
status?: "pending" | "running" | "succeeded" | "exhausted" | "missingArtifact";
preserveTerminal?: boolean;
error?: string;
now?: number;
},
) {
const now = args.now ?? Date.now();
const status = args.status ?? "pending";
const existing =
args.targetKind === "skillVersion" && args.skillVersionId
? await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_skill_version", (q) => q.eq("skillVersionId", args.skillVersionId))
.unique()
: args.targetKind === "packageRelease" && args.packageReleaseId
? await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_package_release", (q) => q.eq("packageReleaseId", args.packageReleaseId))
.unique()
: null;
const lastError = truncateBackupJobError(args.error);
if (existing) {
if (
args.preserveTerminal &&
(existing.status === "succeeded" || existing.status === "missingArtifact")
) {
return { jobId: existing._id, created: false as const, preserved: true as const };
}
await ctx.db.patch(existing._id, {
status,
reason: args.reason,
attempts: 0,
lastError,
nextRunAt: now,
leaseToken: undefined,
leaseExpiresAt: undefined,
claimedAt: undefined,
createdAt: now,
updatedAt: now,
exhaustedAt: undefined,
completedAt: status === "missingArtifact" || status === "succeeded" ? now : undefined,
});
return { jobId: existing._id, created: false as const };
}
const jobId = await ctx.db.insert("registryArtifactBackupJobs", {
targetKind: args.targetKind,
skillVersionId: args.skillVersionId,
packageReleaseId: args.packageReleaseId,
status,
reason: args.reason,
attempts: 0,
nextRunAt: now,
lastError,
completedAt: status === "missingArtifact" || status === "succeeded" ? now : undefined,
createdAt: now,
updatedAt: now,
});
return { jobId, created: true as const };
}
export async function getRegistryArtifactBackupHealthHandler(
ctx: Pick<QueryCtx, "db">,
args: { now?: number; staleAfterMs?: number; sampleLimit?: number },
) {
const now = args.now ?? Date.now();
const staleAfterMs = args.staleAfterMs ?? 24 * 60 * 60 * 1000;
const sampleLimit = clampInt(
args.sampleLimit ?? DEFAULT_BACKUP_HEALTH_SAMPLE_LIMIT,
1,
MAX_BACKUP_HEALTH_SAMPLE_LIMIT,
);
const pending = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_nextRunAt", (q) => q.eq("status", "pending").lte("nextRunAt", now))
.take(sampleLimit + 1);
const exhausted = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_nextRunAt", (q) => q.eq("status", "exhausted"))
.take(sampleLimit + 1);
const running = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_leaseExpiresAt", (q) => q.eq("status", "running"))
.take(sampleLimit + 1);
const expiredRunning = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_status_leaseExpiresAt", (q) =>
q.eq("status", "running").lte("leaseExpiresAt", now),
)
.take(sampleLimit + 1);
const pendingSample = pending.slice(0, sampleLimit);
const exhaustedSample = exhausted.slice(0, sampleLimit);
const runningSample = running.slice(0, sampleLimit);
const expiredRunningSample = expiredRunning.slice(0, sampleLimit);
const oldestPendingAgeMs = pendingSample.reduce(
(max: number, job: { createdAt: number }) => Math.max(max, now - job.createdAt),
0,
);
const stalePending = pendingSample.filter(
(job: { createdAt: number }) => now - job.createdAt >= staleAfterMs,
).length;
const stale = stalePending + expiredRunningSample.length;
return {
pending: pendingSample.length,
running: runningSample.length,
expiredRunning: expiredRunningSample.length,
stale,
exhausted: exhaustedSample.length,
oldestPendingAgeMs,
pendingCapped: pending.length > sampleLimit,
runningCapped: running.length > sampleLimit,
expiredRunningCapped: expiredRunning.length > sampleLimit,
exhaustedCapped: exhausted.length > sampleLimit,
};
}
function truncateBackupJobError(error: string | undefined) {
if (!error) return undefined;
return error.slice(0, MAX_BACKUP_JOB_ERROR_LENGTH);
}
function retryDelayMs(attempts: number) {
const minutes = Math.min(60, 2 ** Math.min(attempts, 6));
return minutes * 60 * 1000;
}
File diff suppressed because it is too large Load Diff
-283
View File
@@ -1,283 +0,0 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { restoreSkillFromBackup } from "./registryArtifactRestore";
const registryBackupMocks = vi.hoisted(() => ({
fetchSkillVersionBackupMeta: vi.fn(),
getRegistryArtifactBackupContext: vi.fn(),
isRegistryArtifactBackupConfigured: vi.fn(),
normalizeOwner: vi.fn((value: string) => value.toLowerCase()),
readRegistryArtifactBackupObject: vi.fn(),
}));
const skillPublishMocks = vi.hoisted(() => ({
publishVersionForUser: vi.fn(),
}));
vi.mock("./lib/registryArtifactBackup", () => registryBackupMocks);
vi.mock("./lib/skillPublish", () => skillPublishMocks);
const restoreHandler = (restoreSkillFromBackup as unknown as { _handler: Function })._handler;
describe("restoreSkillFromBackup", () => {
beforeEach(() => {
vi.resetAllMocks();
registryBackupMocks.normalizeOwner.mockImplementation((value: string) => value.toLowerCase());
registryBackupMocks.getRegistryArtifactBackupContext.mockReturnValue({
endpoint: "https://account.r2.cloudflarestorage.com",
bucket: "clawhub-registry-backup",
accessKeyId: "access-key",
secretAccessKey: "secret-key",
region: "auto",
skillsRoot: "skills",
packagesRoot: "packages",
});
registryBackupMocks.isRegistryArtifactBackupConfigured.mockReturnValue(true);
});
it("blocks restore when the current slug row is not public", async () => {
const result = await restoreHandler(
{
runQuery: vi
.fn()
.mockResolvedValueOnce({ _id: "users:admin", role: "admin" })
.mockResolvedValueOnce({
_id: "skills:hidden",
ownerUserId: "users:owner",
slug: "demo-skill",
softDeletedAt: undefined,
moderationStatus: "hidden",
}),
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
version: "1.0.0",
},
);
expect(result).toEqual({
slug: "demo-skill",
status: "error",
detail: "Existing skill is not public; restore blocked",
});
expect(registryBackupMocks.fetchSkillVersionBackupMeta).not.toHaveBeenCalled();
});
it("requires an explicit version before any forced slug eviction", async () => {
const runMutation = vi.fn();
const result = await restoreHandler(
{
runQuery: vi.fn().mockResolvedValueOnce({ _id: "users:admin", role: "admin" }),
runMutation,
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
forceOverwriteSquatter: true,
},
);
expect(result).toEqual({
slug: "demo-skill",
status: "no_backup",
detail: "Restore requires an explicit backup version",
});
expect(runMutation).not.toHaveBeenCalled();
expect(registryBackupMocks.fetchSkillVersionBackupMeta).not.toHaveBeenCalled();
});
it("validates the requested version before forced slug eviction", async () => {
registryBackupMocks.fetchSkillVersionBackupMeta.mockResolvedValueOnce(null);
const runMutation = vi.fn();
const result = await restoreHandler(
{
runQuery: vi
.fn()
.mockResolvedValueOnce({ _id: "users:admin", role: "admin" })
.mockResolvedValueOnce({
_id: "skills:squatter",
ownerUserId: "users:other",
slug: "demo-skill",
softDeletedAt: undefined,
moderationStatus: "active",
}),
runMutation,
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
version: "typo-version",
forceOverwriteSquatter: true,
},
);
expect(result).toEqual({
slug: "demo-skill",
status: "no_backup",
detail: "No version backup found",
});
expect(runMutation).not.toHaveBeenCalled();
});
it("reactivates the same owner's soft-deleted skill row without republishing a duplicate version", async () => {
registryBackupMocks.fetchSkillVersionBackupMeta.mockResolvedValueOnce({
version: "1.0.0",
displayName: "Demo Skill",
metadata: {
files: [
{
path: "SKILL.md",
size: 5,
sha256: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824",
},
],
},
});
registryBackupMocks.readRegistryArtifactBackupObject.mockResolvedValueOnce(
new TextEncoder().encode("hello"),
);
const storage = { store: vi.fn().mockResolvedValue("storage:restored") };
const runMutation = vi.fn();
const result = await restoreHandler(
{
runQuery: vi
.fn()
.mockResolvedValueOnce({ _id: "users:admin", role: "admin" })
.mockResolvedValueOnce({
_id: "skills:deleted",
ownerUserId: "users:owner",
slug: "demo-skill",
softDeletedAt: 123,
moderationStatus: "active",
})
.mockResolvedValueOnce({
_id: "skillVersions:existing",
skillId: "skills:deleted",
version: "1.0.0",
softDeletedAt: undefined,
}),
runMutation,
storage,
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
version: "1.0.0",
},
);
expect(result).toEqual({ slug: "demo-skill", status: "restored" });
expect(skillPublishMocks.publishVersionForUser).not.toHaveBeenCalled();
expect(runMutation).toHaveBeenNthCalledWith(
1,
expect.anything(),
expect.objectContaining({
userId: "users:admin",
slug: "demo-skill",
deleted: false,
}),
);
expect(runMutation).toHaveBeenNthCalledWith(
2,
expect.anything(),
expect.objectContaining({
actorUserId: "users:admin",
skillId: "skills:deleted",
versionId: "skillVersions:existing",
files: [
expect.objectContaining({
path: "SKILL.md",
size: 5,
sha256: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824",
storageId: "storage:restored",
}),
],
}),
);
});
it("fails restore when a manifest file is missing from backup storage", async () => {
registryBackupMocks.fetchSkillVersionBackupMeta.mockResolvedValueOnce({
version: "1.0.0",
displayName: "Demo Skill",
metadata: {
files: [{ path: "SKILL.md", size: 5, sha256: "2cf24dba5fb0a30e26e83b2ac5b9e29e" }],
},
});
registryBackupMocks.readRegistryArtifactBackupObject.mockResolvedValueOnce(null);
const storage = { store: vi.fn() };
const result = await restoreHandler(
{
runQuery: vi
.fn()
.mockResolvedValueOnce({ _id: "users:admin", role: "admin" })
.mockResolvedValueOnce(null),
storage,
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
version: "1.0.0",
},
);
expect(result).toEqual({
slug: "demo-skill",
status: "error",
detail: "Backup missing file SKILL.md",
});
expect(storage.store).not.toHaveBeenCalled();
expect(skillPublishMocks.publishVersionForUser).not.toHaveBeenCalled();
});
it("fails restore when a manifest file checksum does not match backup storage", async () => {
registryBackupMocks.fetchSkillVersionBackupMeta.mockResolvedValueOnce({
version: "1.0.0",
displayName: "Demo Skill",
metadata: {
files: [{ path: "SKILL.md", size: 5, sha256: "wrong-sha256" }],
},
});
registryBackupMocks.readRegistryArtifactBackupObject.mockResolvedValueOnce(
new TextEncoder().encode("hello"),
);
const storage = { store: vi.fn() };
const result = await restoreHandler(
{
runQuery: vi
.fn()
.mockResolvedValueOnce({ _id: "users:admin", role: "admin" })
.mockResolvedValueOnce(null),
storage,
} as never,
{
actorUserId: "users:admin",
ownerHandle: "alice",
ownerUserId: "users:owner",
slug: "demo-skill",
version: "1.0.0",
},
);
expect(result).toEqual({
slug: "demo-skill",
status: "error",
detail: "Backup file checksum mismatch for SKILL.md",
});
expect(storage.store).not.toHaveBeenCalled();
expect(skillPublishMocks.publishVersionForUser).not.toHaveBeenCalled();
});
});
-350
View File
@@ -1,350 +0,0 @@
"use node";
import { v } from "convex/values";
import { internal } from "./_generated/api";
import type { Doc, Id } from "./_generated/dataModel";
import { internalAction } from "./functions";
import { assertAdmin } from "./lib/access";
import { guessContentTypeForPath } from "./lib/contentTypes";
import { isPublicSkillDoc } from "./lib/globalStats";
import {
fetchSkillVersionBackupMeta,
getRegistryArtifactBackupContext,
isRegistryArtifactBackupConfigured,
normalizeOwner,
readRegistryArtifactBackupObject,
} from "./lib/registryArtifactBackup";
import { publishVersionForUser } from "./lib/skillPublish";
import { validateFilePath } from "./lib/skillZip";
type RestoreResult = {
slug: string;
status: "restored" | "slug_conflict" | "already_exists" | "no_backup" | "error";
detail?: string;
};
type BulkRestoreResult = {
results: RestoreResult[];
totalRestored: number;
totalConflicts: number;
totalSkipped: number;
totalErrors: number;
};
type SkillBackupMeta = NonNullable<Awaited<ReturnType<typeof fetchSkillVersionBackupMeta>>>;
type VerifiedSkillBackup = {
meta: SkillBackupMeta;
files: Array<{
path: string;
size: number;
sha256: string;
contentType: string;
content: Uint8Array;
}>;
};
/**
* Admin-only: restore a single skill from registry artifact backup.
* Reads backed-up objects and re-creates the skill in the database.
*/
export const restoreSkillFromBackup = internalAction({
args: {
actorUserId: v.id("users"),
ownerHandle: v.string(),
ownerUserId: v.id("users"),
slug: v.string(),
version: v.optional(v.string()),
forceOverwriteSquatter: v.optional(v.boolean()),
},
handler: async (ctx, args): Promise<RestoreResult> => {
try {
const actor = await ctx.runQuery(internal.users.getByIdInternal, {
userId: args.actorUserId,
});
if (!actor || actor.deletedAt || actor.deactivatedAt) {
return { slug: args.slug, status: "error", detail: "Actor not found" };
}
assertAdmin(actor as Doc<"users">);
if (!isRegistryArtifactBackupConfigured()) {
return { slug: args.slug, status: "error", detail: "Registry backup not configured" };
}
const backupContext = getRegistryArtifactBackupContext();
if (!args.version) {
return {
slug: args.slug,
status: "no_backup",
detail: "Restore requires an explicit backup version",
};
}
let verifiedBackup: VerifiedSkillBackup | null = null;
const loadVerifiedBackup = async (): Promise<VerifiedSkillBackup | RestoreResult> => {
if (verifiedBackup) return verifiedBackup;
const meta = await fetchSkillVersionBackupMeta(
backupContext,
args.ownerHandle,
args.slug,
args.version!,
);
if (!meta) {
return { slug: args.slug, status: "no_backup", detail: "No version backup found" };
}
const backupFiles = meta.metadata.files;
if (backupFiles.length === 0) {
return { slug: args.slug, status: "no_backup", detail: "Backup has no files" };
}
const owner = normalizeOwner(args.ownerHandle);
const files: VerifiedSkillBackup["files"] = [];
for (const file of backupFiles) {
if (!validateFilePath(file.path)) {
return { slug: args.slug, status: "error", detail: "Backup contains unsafe file path" };
}
const fileContent = await readRegistryArtifactBackupObject(
backupContext,
`${backupContext.skillsRoot}/${owner}/${args.slug}/${encodeBackupPathSegment(
meta.version,
)}/${file.path}`,
);
if (!fileContent) {
return {
slug: args.slug,
status: "error",
detail: `Backup missing file ${file.path}`,
};
}
if (fileContent.byteLength !== file.size) {
return {
slug: args.slug,
status: "error",
detail: `Backup file size mismatch for ${file.path}`,
};
}
const sha256 = await sha256Hex(fileContent);
if (sha256 !== file.sha256) {
return {
slug: args.slug,
status: "error",
detail: `Backup file checksum mismatch for ${file.path}`,
};
}
files.push({
path: file.path,
size: fileContent.byteLength,
sha256,
contentType: file.contentType ?? guessContentTypeForPath(file.path),
content: fileContent,
});
}
verifiedBackup = { meta, files };
return verifiedBackup;
};
// Check if skill already exists in the DB
const existingSkill = (await ctx.runQuery(
internal.skills.getSkillBySlugIncludingSoftDeletedInternal,
{
slug: args.slug,
},
)) as Doc<"skills"> | null;
const sameOwnerSoftDeletedSkill =
existingSkill?.ownerUserId === args.ownerUserId && existingSkill.softDeletedAt
? existingSkill
: null;
if (existingSkill) {
const sameOwner = existingSkill.ownerUserId === args.ownerUserId;
if (sameOwner && existingSkill.softDeletedAt) {
// Continue: if the backed-up version already exists, restore by
// reactivating that row instead of republishing a duplicate version.
} else if (!isPublicSkillDoc(existingSkill)) {
return {
slug: args.slug,
status: "error",
detail: "Existing skill is not public; restore blocked",
};
} else if (sameOwner) {
return {
slug: args.slug,
status: "already_exists",
detail: "Skill already owned by user",
};
} else if (!args.forceOverwriteSquatter) {
return {
slug: args.slug,
status: "slug_conflict",
detail: `Slug occupied by another user. Set forceOverwriteSquatter=true to reclaim.`,
};
} else {
const backup = await loadVerifiedBackup();
if ("status" in backup) return backup;
// Free the slug in-transaction by renaming the squatter, then enqueue cleanup.
await ctx.runMutation(
internal.registryArtifactRestoreMutations.evictSquatterSkillForRestoreInternal,
{
actorUserId: args.actorUserId,
slug: args.slug,
rightfulOwnerUserId: args.ownerUserId,
},
);
}
}
const backup = verifiedBackup ?? (await loadVerifiedBackup());
if ("status" in backup) return backup;
const { meta } = backup;
// Download and store each file in Convex storage
const storedFiles: Array<{
path: string;
size: number;
storageId: Id<"_storage">;
sha256: string;
contentType: string;
}> = [];
for (const file of backup.files) {
const blob = new Blob([Buffer.from(file.content)], { type: file.contentType });
const storageId = await ctx.storage.store(blob);
storedFiles.push({
path: file.path,
size: file.size,
storageId,
sha256: file.sha256,
contentType: file.contentType,
});
}
if (storedFiles.length === 0) {
return { slug: args.slug, status: "error", detail: "Could not download any backup files" };
}
if (sameOwnerSoftDeletedSkill) {
const existingVersion = (await ctx.runQuery(
internal.skills.getVersionBySkillAndVersionInternal,
{
skillId: sameOwnerSoftDeletedSkill._id,
version: meta.version,
},
)) as Doc<"skillVersions"> | null;
if (existingVersion && !existingVersion.softDeletedAt) {
await ctx.runMutation(internal.skills.setSkillSoftDeletedInternal, {
userId: args.actorUserId,
slug: args.slug,
deleted: false,
reason: "Restored from registry artifact backup",
});
await ctx.runMutation(
internal.registryArtifactRestoreMutations.refreshRestoredSkillVersionInternal,
{
actorUserId: args.actorUserId,
skillId: sameOwnerSoftDeletedSkill._id,
versionId: existingVersion._id,
files: storedFiles,
},
);
return { slug: args.slug, status: "restored" };
}
}
await publishVersionForUser(
ctx,
args.ownerUserId,
{
slug: args.slug,
displayName: meta.displayName,
version: meta.version,
changelog: "Restored from registry artifact backup",
files: storedFiles,
},
{
bypassGitHubAccountAge: true,
bypassNewSkillRateLimit: true,
bypassQualityGate: true,
skipBackup: true,
skipWebhook: true,
},
);
return { slug: args.slug, status: "restored" };
} catch (error) {
const message = error instanceof Error ? error.message : "Unknown error";
console.error(`[restore] Failed to restore ${args.slug}:`, message);
return { slug: args.slug, status: "error", detail: message };
}
},
});
/**
* Admin-only: bulk restore all skills for a user from registry artifact backup.
*/
export const restoreUserSkillsFromBackup = internalAction({
args: {
actorUserId: v.id("users"),
ownerHandle: v.string(),
ownerUserId: v.id("users"),
slugs: v.array(v.string()),
versionsBySlug: v.optional(v.record(v.string(), v.string())),
forceOverwriteSquatter: v.optional(v.boolean()),
},
handler: async (ctx, args): Promise<BulkRestoreResult> => {
const results: RestoreResult[] = [];
let totalRestored = 0;
let totalConflicts = 0;
let totalSkipped = 0;
let totalErrors = 0;
for (const slug of args.slugs) {
const result = (await ctx.runAction(internal.registryArtifactRestore.restoreSkillFromBackup, {
actorUserId: args.actorUserId,
ownerHandle: args.ownerHandle,
ownerUserId: args.ownerUserId,
slug,
version: args.versionsBySlug?.[slug],
forceOverwriteSquatter: args.forceOverwriteSquatter,
})) as RestoreResult;
results.push(result);
switch (result.status) {
case "restored":
totalRestored += 1;
break;
case "slug_conflict":
totalConflicts += 1;
break;
case "already_exists":
case "no_backup":
totalSkipped += 1;
break;
case "error":
totalErrors += 1;
break;
}
}
return { results, totalRestored, totalConflicts, totalSkipped, totalErrors };
},
});
async function sha256Hex(bytes: Uint8Array) {
const { createHash } = await import("node:crypto");
const hash = createHash("sha256");
hash.update(bytes);
return hash.digest("hex");
}
function encodeBackupPathSegment(value: string) {
return encodeURIComponent(value.trim()).replace(/\./g, "%2E");
}
// guessContentTypeForPath in lib/contentTypes.ts
-168
View File
@@ -1,168 +0,0 @@
import { v } from "convex/values";
import { internal } from "./_generated/api";
import type { Doc } from "./_generated/dataModel";
import { internalMutation } from "./functions";
import { assertAdmin } from "./lib/access";
const restoredSkillFileValidator = v.object({
path: v.string(),
size: v.number(),
storageId: v.id("_storage"),
sha256: v.string(),
contentType: v.optional(v.string()),
});
export const evictSquatterSkillForRestoreInternal = internalMutation({
args: {
actorUserId: v.id("users"),
slug: v.string(),
rightfulOwnerUserId: v.id("users"),
},
handler: async (ctx, args) => {
const actor = await ctx.db.get(args.actorUserId);
if (!actor || actor.deletedAt || actor.deactivatedAt) throw new Error("Actor not found");
assertAdmin(actor);
const slug = args.slug.trim().toLowerCase();
if (!slug) throw new Error("Slug required");
const now = Date.now();
const existingSkill = await ctx.db
.query("skills")
.withIndex("by_slug", (q) => q.eq("slug", slug))
.unique();
if (!existingSkill) return { ok: true as const, action: "noop" as const };
if (existingSkill.ownerUserId === args.rightfulOwnerUserId) {
return { ok: true as const, action: "already_owned" as const };
}
const evictedSlug = buildEvictedSlug(slug, now);
// Free the slug immediately (same transaction) by renaming the squatter's skill.
await ctx.db.patch(existingSkill._id, {
slug: evictedSlug,
softDeletedAt: now,
hiddenAt: existingSkill.hiddenAt ?? now,
hiddenBy: existingSkill.hiddenBy ?? actor._id,
updatedAt: now,
});
// Remove from vector search ASAP.
const embeddings = await ctx.db
.query("skillEmbeddings")
.withIndex("by_skill", (q) => q.eq("skillId", existingSkill._id))
.collect();
for (const embedding of embeddings) {
await ctx.db.patch(embedding._id, {
visibility: "deleted",
updatedAt: now,
});
}
// Cleanup the rest asynchronously (versions, fingerprints, installs, etc.)
await ctx.scheduler.runAfter(0, internal.skills.hardDeleteInternal, {
skillId: existingSkill._id,
actorUserId: actor._id,
phase: "versions",
});
await ctx.db.insert("auditLogs", {
actorUserId: actor._id,
action: "slug.reclaim.sync",
targetType: "skill",
targetId: existingSkill._id,
metadata: {
slug,
evictedSlug,
squatterUserId: existingSkill.ownerUserId,
rightfulOwnerUserId: args.rightfulOwnerUserId,
reason: "Synchronous eviction during registry artifact restore",
},
createdAt: now,
});
return { ok: true as const, action: "evicted" as const, evictedSlug };
},
});
export const refreshRestoredSkillVersionInternal = internalMutation({
args: {
actorUserId: v.id("users"),
skillId: v.id("skills"),
versionId: v.id("skillVersions"),
files: v.array(restoredSkillFileValidator),
},
handler: async (ctx, args) => {
const actor = await ctx.db.get(args.actorUserId);
if (!actor || actor.deletedAt || actor.deactivatedAt) throw new Error("Actor not found");
assertAdmin(actor);
const [skill, version] = await Promise.all([
ctx.db.get(args.skillId),
ctx.db.get(args.versionId),
]);
if (!skill) throw new Error("Skill not found");
if (!version || version.skillId !== skill._id || version.softDeletedAt) {
throw new Error("Skill version not found");
}
const now = Date.now();
await ctx.db.patch(version._id, {
files: args.files,
});
await ctx.db.patch(skill._id, {
latestVersionId: version._id,
latestVersionSummary: latestVersionSummaryFromVersion(version),
tags: { ...skill.tags, latest: version._id },
softDeletedAt: undefined,
moderationStatus: "active",
hiddenAt: undefined,
hiddenBy: undefined,
unpublishedSlugReservedUntil: undefined,
unpublishedSlugReleasedAt: undefined,
unpublishedOriginalSlug: undefined,
updatedAt: now,
});
await ctx.db.insert("auditLogs", {
actorUserId: actor._id,
action: "skill.restore.registry_artifact",
targetType: "skill",
targetId: skill._id,
metadata: {
slug: skill.slug,
version: version.version,
versionId: version._id,
},
createdAt: now,
});
return { ok: true as const, skillId: skill._id, versionId: version._id };
},
});
function buildEvictedSlug(slug: string, now: number) {
const suffix = now.toString(36);
return `${slug}-evicted-${suffix}`;
}
function latestVersionSummaryFromVersion(
version: Pick<
Doc<"skillVersions">,
"version" | "createdAt" | "changelog" | "changelogSource" | "parsed"
>,
): NonNullable<Doc<"skills">["latestVersionSummary"]> {
return {
version: version.version,
createdAt: version.createdAt,
changelog: version.changelog,
changelogSource: version.changelogSource,
description: frontmatterString(version.parsed?.frontmatter?.description),
clawdis: version.parsed?.clawdis,
};
}
function frontmatterString(value: unknown) {
return typeof value === "string" ? value.trim() || undefined : undefined;
}
+1 -14
View File
@@ -1205,11 +1205,10 @@ describe("skills ownership", () => {
createdAt: 1_700_000_000_000,
softDeletedAt: undefined,
};
const runAfter = vi.fn(async () => {});
const result = await renameOwnedSkillInternalHandler(
{
scheduler: { runAfter },
scheduler: { runAfter: vi.fn(async () => {}) },
db: {
normalizeId: vi.fn(() => null),
get: vi.fn(async (id: string) => {
@@ -1301,18 +1300,6 @@ describe("skills ownership", () => {
ownerPublisherId: "publishers:org",
}),
);
expect(runAfter).toHaveBeenCalledWith(
0,
expect.anything(),
expect.objectContaining({
skillId: "skills:source",
versionId: "skillVersions:latest",
slug: "new-name",
version: "1.0.0",
isLatest: true,
ownerHandle: "org",
}),
);
});
it("allows publisher admins to move a skill into an org they administer", async () => {
-138
View File
@@ -9883,113 +9883,9 @@ export const updateTags = mutation({
if (latestEntry) {
await setSkillEmbeddingsLatestVersion(ctx, skill._id, latestEntry.versionId, now);
}
if (latestEntry && latestEntry.versionId !== skill.latestVersionId) {
const version = versionsById.get(latestEntry.versionId);
const owner = await getOwnerPublisher(ctx, {
ownerPublisherId: skill.ownerPublisherId,
ownerUserId: skill.ownerUserId,
});
if (version && owner) {
await scheduleRegistryArtifactSkillVersionBackupRefresh(ctx, {
skill,
version,
isLatest: true,
ownerHandle: owner.handle ?? String(skill.ownerPublisherId ?? skill.ownerUserId),
logContext: "latest refresh",
});
}
}
},
});
async function scheduleRegistryArtifactSkillVersionBackupRefresh(
ctx: MutationCtx,
params: {
skill: Pick<Doc<"skills">, "_id" | "slug" | "displayName" | "ownerUserId" | "ownerPublisherId">;
version: Doc<"skillVersions">;
ownerHandle: string;
isLatest: boolean;
logContext: string;
},
) {
try {
await ctx.scheduler.runAfter(
0,
internal.registryArtifactBackupsNode.backupSkillForPublishInternal,
{
skillId: params.skill._id,
versionId: params.version._id,
slug: params.skill.slug,
version: params.version.version,
isLatest: params.isLatest,
displayName: params.skill.displayName,
ownerHandle: params.ownerHandle,
files: params.version.files,
publishedAt: params.version.createdAt,
},
);
} catch (error) {
console.error(`registry artifact backup ${params.logContext} scheduling failed`, error);
await enqueueSkillVersionBackupRetryJob(ctx, {
versionId: params.version._id,
error: errorMessageForBackupJob(error),
}).catch((enqueueError) => {
console.error(
`registry artifact backup ${params.logContext} retry enqueue failed`,
enqueueError,
);
});
}
}
async function enqueueSkillVersionBackupRetryJob(
ctx: Pick<MutationCtx, "db">,
args: { versionId: Id<"skillVersions">; error?: string },
) {
const now = Date.now();
const lastError = truncateBackupJobError(args.error);
const existing = await ctx.db
.query("registryArtifactBackupJobs")
.withIndex("by_skill_version", (q) => q.eq("skillVersionId", args.versionId))
.unique();
if (existing) {
await ctx.db.patch(existing._id, {
status: "pending",
reason: "retry",
attempts: 0,
lastError,
nextRunAt: now,
createdAt: now,
updatedAt: now,
exhaustedAt: undefined,
completedAt: undefined,
});
return;
}
await ctx.db.insert("registryArtifactBackupJobs", {
targetKind: "skillVersion",
skillVersionId: args.versionId,
packageReleaseId: undefined,
status: "pending",
reason: "retry",
attempts: 0,
nextRunAt: now,
lastError,
createdAt: now,
updatedAt: now,
});
}
function errorMessageForBackupJob(error: unknown) {
return error instanceof Error ? error.message : String(error);
}
function truncateBackupJobError(error: string | undefined) {
if (!error) return undefined;
return error.length > 4000 ? `${error.slice(0, 3997)}...` : error;
}
export const deleteTags = mutation({
args: {
skillId: v.id("skills"),
@@ -10613,23 +10509,6 @@ async function renameOwnedSkillByActor(
createdAt: now,
});
if (skill.latestVersionId) {
const latestVersion = await ctx.db.get(skill.latestVersionId);
const owner = await getOwnerPublisher(ctx, {
ownerPublisherId: skill.ownerPublisherId,
ownerUserId: skill.ownerUserId,
});
if (latestVersion && !latestVersion.softDeletedAt && owner) {
await scheduleRegistryArtifactSkillVersionBackupRefresh(ctx, {
skill: { ...skill, slug: newSlug },
version: latestVersion,
isLatest: true,
ownerHandle: owner.handle ?? String(skill.ownerPublisherId ?? skill.ownerUserId),
logContext: "rename refresh",
});
}
}
return { ok: true as const, slug: newSlug, previousSlug: skill.slug };
}
@@ -11063,23 +10942,6 @@ export const transferSkillOwnerForUserInternal = internalMutation({
createdAt: now,
});
if (skill.latestVersionId) {
const latestVersion = await ctx.db.get(skill.latestVersionId);
if (latestVersion && !latestVersion.softDeletedAt) {
await scheduleRegistryArtifactSkillVersionBackupRefresh(ctx, {
skill: {
...skill,
ownerUserId: nextOwner._id,
ownerPublisherId: destinationPublisher._id,
},
version: latestVersion,
isLatest: true,
ownerHandle: destinationPublisher.handle,
logContext: "owner transfer refresh",
});
}
}
return {
ok: true as const,
transferred: true as const,
-1
View File
@@ -108,7 +108,6 @@ describe("permission boundary e2e", () => {
method: "DELETE",
path: `${ApiRoutes.packages}/e2e-nonexistent-permission/trusted-publisher`,
},
{ method: "POST", path: `${ApiRoutes.users}/restore`, body: { handle: "nobody" } },
{ method: "POST", path: `${ApiRoutes.users}/reclaim`, body: { handle: "nobody" } },
{ method: "POST", path: `${ApiRoutes.users}/reserve`, body: { handle: "nobody" } },
{ method: "POST", path: `${ApiRoutes.users}/publisher`, body: { handle: "nobody" } },
+3 -3
View File
@@ -12,8 +12,8 @@ read_when:
- Minimal, fast SPA for browsing and publishing agent skills.
- Skills stored in Convex (files + metadata + versions + stats).
- GitHub OAuth login; R2/object storage backs up hosted registry artifacts for
disaster recovery.
- GitHub OAuth login; Convex backups with file storage are the source of truth
for hosted registry artifact disaster recovery.
- Vector-based search over skill text + metadata.
- Versioning, tags (`latest` + user tags), changelog, rollback (tag movement).
- Public read access; upload requires auth.
@@ -149,7 +149,7 @@ Local fixture data lives in `convex/devSeed.ts` and `fixtures/public-corpus/`.
## Vercel
- Env vars: Convex deployment URLs + GitHub OAuth client + OpenAI key (if used) + registry artifact backup R2 credentials.
- Env vars: Convex deployment URLs + GitHub OAuth client + OpenAI key (if used).
- SPA feel: client-side transitions, prefetching, optimistic UI.
## Open questions (carry forward)