Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
141 changes: 141 additions & 0 deletions convex/httpApiV1.handlers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16899,6 +16899,147 @@ describe("httpApiV1 handlers", () => {
);
});

it("delete unpublished package blobs when loose-file publish action throws after stores", async () => {
vi.mocked(getOptionalApiTokenUserId).mockResolvedValue("users:1" as never);
vi.mocked(requirePackagePublishAuth).mockResolvedValue({
kind: "user",
userId: "users:1",
user: { _id: "users:1", handle: "p" },
} as never);
const runMutation = vi.fn().mockResolvedValue(okRate());
const runAction = vi.fn().mockRejectedValue(new Error("GitHub account is too new to publish"));
const storedIds: string[] = [];
const storageStore = vi.fn(async () => {
const id = `storage:loose-${storedIds.length + 1}`;
storedIds.push(id);
return id;
});
const storageDelete = vi.fn(async () => {});
const form = packagePublishForm(
packagePublishMetadata({
ownerHandle: "openclaw",
bundle: { hostTargets: ["desktop"] },
}),
);
form.append("files", new File(["{}"], "openclaw.plugin.json", { type: "application/json" }));
form.append("files", new File(["readme"], "README.md", { type: "text/markdown" }));

const response = await __handlers.publishPackageV1Handler(
makeCtx({
runAction,
runMutation,
storage: { store: storageStore, delete: storageDelete },
}),
new Request("https://example.com/api/v1/packages", {
method: "POST",
headers: { Authorization: "Bearer clh_test" },
body: form,
}),
);

expect(response.status).toBe(400);
expect(await response.text()).toBe("GitHub account is too new to publish");
expect(runAction).toHaveBeenCalledOnce();
expect(storedIds).toEqual(["storage:loose-1", "storage:loose-2"]);
expect(storageDelete).toHaveBeenCalledTimes(storedIds.length);
for (const id of storedIds) {
expect(storageDelete).toHaveBeenCalledWith(id);
}
});

it("delete unpublished package blobs when ClawPack extracted-file store throws after tarball store", async () => {
vi.mocked(getOptionalApiTokenUserId).mockResolvedValue("users:1" as never);
vi.mocked(requirePackagePublishAuth).mockResolvedValue({
kind: "user",
userId: "users:1",
user: { _id: "users:1", handle: "p" },
} as never);
const runMutation = vi.fn().mockResolvedValue(okRate());
const runAction = vi.fn();
const storageStore = vi.fn(async () => {
if (storageStore.mock.calls.length > 1) {
throw new Error("extracted store failed");
}
return "storage:tarball";
});
const storageDelete = vi.fn(async () => {});
const pack = npmPackFixture({
"package/package.json": JSON.stringify({ name: "demo-plugin", version: "1.0.0" }),
"package/openclaw.plugin.json": JSON.stringify({ id: "demo.plugin" }),
"package/dist/index.js": "export const demo = true;\n",
});
const form = packagePublishForm(packagePublishMetadata({ family: "code-plugin" }));
form.append(
"clawpack",
new File([bytesToArrayBuffer(pack)], "demo-plugin-1.0.0.tgz", {
type: "application/octet-stream",
}),
);

const response = await __handlers.publishPackageV1Handler(
makeCtx({
runAction,
runMutation,
storage: { store: storageStore, delete: storageDelete },
}),
new Request("https://example.com/api/v1/packages", {
method: "POST",
headers: { Authorization: "Bearer clh_test" },
body: form,
}),
);

expect(response.status).toBe(400);
expect(await response.text()).toBe("extracted store failed");
expect(runAction).not.toHaveBeenCalled();
expect(storageStore).toHaveBeenCalled();
expect(storageDelete).toHaveBeenCalledWith("storage:tarball");
});

it("delete unpublished package blobs when a later loose-file store throws", async () => {
vi.mocked(getOptionalApiTokenUserId).mockResolvedValue("users:1" as never);
vi.mocked(requirePackagePublishAuth).mockResolvedValue({
kind: "user",
userId: "users:1",
user: { _id: "users:1", handle: "p" },
} as never);
const runMutation = vi.fn().mockResolvedValue(okRate());
const runAction = vi.fn();
const storageStore = vi.fn(async () => {
if (storageStore.mock.calls.length > 1) {
throw new Error("second file store failed");
}
return "storage:first";
});
const storageDelete = vi.fn(async () => {});
const form = packagePublishForm(
packagePublishMetadata({
ownerHandle: "openclaw",
bundle: { hostTargets: ["desktop"] },
}),
);
form.append("files", new File(["{}"], "openclaw.plugin.json", { type: "application/json" }));
form.append("files", new File(["readme"], "README.md", { type: "text/markdown" }));

const response = await __handlers.publishPackageV1Handler(
makeCtx({
runAction,
runMutation,
storage: { store: storageStore, delete: storageDelete },
}),
new Request("https://example.com/api/v1/packages", {
method: "POST",
headers: { Authorization: "Bearer clh_test" },
body: form,
}),
);

expect(response.status).toBe(400);
expect(await response.text()).toBe("second file store failed");
expect(runAction).not.toHaveBeenCalled();
expect(storageDelete).toHaveBeenCalledWith("storage:first");
});

it("returns package publish attempt status to the exact API token actor", async () => {
vi.mocked(requirePackagePublishAuth).mockResolvedValue({
kind: "user",
Expand Down
97 changes: 79 additions & 18 deletions convex/httpApiV1/packagesV1.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1357,6 +1357,27 @@ async function storeClawPackFile(
const CLAWPACK_STORE_BATCH_BYTES = 8 * 1024 * 1024;
const CLAWPACK_STORE_BATCH_FILES = 16;

async function deleteUnpublishedPackagePublishBlobs(
ctx: ActionCtx,
stored: {
files?: Array<{ storageId: string }>;
artifact?: { storageId?: string } | null;
},
) {
const remove = ctx.storage.delete?.bind(ctx.storage);
if (!remove) {
return;
}
const storageIds: string[] = [];
for (const file of stored.files ?? []) {
storageIds.push(file.storageId);
}
if (stored.artifact?.storageId) {
storageIds.push(stored.artifact.storageId);
}
await Promise.allSettled(storageIds.map((storageId) => remove(storageId as Id<"_storage">)));
}

async function storeClawPackFiles(
ctx: ActionCtx,
entries: Array<{ path: string; bytes: Uint8Array }>,
Expand All @@ -1365,20 +1386,34 @@ async function storeClawPackFiles(
let batch: Array<{ path: string; bytes: Uint8Array }> = [];
let batchBytes = 0;
const flush = async () => {
files.push(...(await Promise.all(batch.map((entry) => storeClawPackFile(ctx, entry)))));
const results = await Promise.allSettled(batch.map((entry) => storeClawPackFile(ctx, entry)));
for (const result of results) {
if (result.status === "fulfilled") {
files.push(result.value);
}
}
const rejected = results.find((result) => result.status === "rejected");
if (rejected?.status === "rejected") {
throw rejected.reason;
}
batch = [];
batchBytes = 0;
};
for (const entry of entries) {
const overflows =
batch.length >= CLAWPACK_STORE_BATCH_FILES ||
batchBytes + entry.bytes.byteLength > CLAWPACK_STORE_BATCH_BYTES;
if (batch.length > 0 && overflows) await flush();
batch.push(entry);
batchBytes += entry.bytes.byteLength;
try {
for (const entry of entries) {
const overflows =
batch.length >= CLAWPACK_STORE_BATCH_FILES ||
batchBytes + entry.bytes.byteLength > CLAWPACK_STORE_BATCH_BYTES;
if (batch.length > 0 && overflows) await flush();
batch.push(entry);
batchBytes += entry.bytes.byteLength;
}
if (batch.length > 0) await flush();
return files;
} catch (error) {
await deleteUnpublishedPackagePublishBlobs(ctx, { files });
throw error;
}
if (batch.length > 0) await flush();
return files;
}

async function storeUploadedPackageFile(
Expand Down Expand Up @@ -1489,8 +1524,14 @@ async function buildPackagePublishRequestFromClawPack(
npmUnpackedSize: parsed.unpackedSize,
npmFileCount: parsed.fileCount,
};
const files = await storeClawPackFiles(ctx, parsed.entries);
return { ...metadata, files, artifact };
try {
const files = await storeClawPackFiles(ctx, parsed.entries);
return { ...metadata, files, artifact };
} catch (error) {
// Extracted-file store failed after the tarball was already written.
await deleteUnpublishedPackagePublishBlobs(ctx, { artifact });
throw error;
}
}

function assertClawPackPublicationIdentity(
Expand Down Expand Up @@ -1630,11 +1671,26 @@ async function parseMultipartPackagePublish(
}

const packageFileParts = fileParts.filter((entry) => !isMacJunkPath(entry.name));
const files = await Promise.all(
packageFileParts.map((entry) => storeUploadedPackageFile(ctx, entry)),
);
if (files.length === 0) throw new Error("files required");
return { ...metadata, files };
const files: StoredPackagePublishFile[] = [];
try {
const results = await Promise.allSettled(
packageFileParts.map((entry) => storeUploadedPackageFile(ctx, entry)),
);
for (const result of results) {
if (result.status === "fulfilled") {
files.push(result.value);
}
}
const rejected = results.find((result) => result.status === "rejected");
if (rejected?.status === "rejected") {
throw rejected.reason;
}
if (files.length === 0) throw new Error("files required");
return { ...metadata, files };
} catch (error) {
await deleteUnpublishedPackagePublishBlobs(ctx, { files });
throw error;
}
}

async function listPackages(
Expand Down Expand Up @@ -2521,12 +2577,13 @@ export async function publishPackageV1Handler(ctx: ActionCtx, request: Request)
const auth = await requirePackagePublishAuthOrResponse(ctx, request, rate.headers);
if (!auth.ok) return auth.response;

let payload: ServerPackagePublishRequest | undefined;
try {
const contentType = request.headers.get("content-type") ?? "";
if (!contentType.includes("multipart/form-data")) {
return text("Package publish requires multipart/form-data", 415, rate.headers);
}
const payload = await parseMultipartPackagePublish(ctx, auth.auth, request);
payload = await parseMultipartPackagePublish(ctx, auth.auth, request);
const result =
auth.auth.kind === "user"
? await runActionRef(ctx, internalRefs.packages.publishPackageForUserInternal, {
Expand All @@ -2539,6 +2596,10 @@ export async function publishPackageV1Handler(ctx: ActionCtx, request: Request)
});
return json(result, 200, rate.headers);
} catch (error) {
// Parse succeeded, so these blobs have no release owner yet.
if (payload) {
await deleteUnpublishedPackagePublishBlobs(ctx, payload);
}
return packagePublishErrorToResponse(error, rate.headers);
}
}
Expand Down
Loading