Skip to content
Merged
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
49 changes: 49 additions & 0 deletions apps/worker/src/tasks/sync-author/processor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import { s3 } from "@playfulprogramming/s3";
import * as github from "@playfulprogramming/github-api";
import { Readable } from "node:stream";
import { eq } from "drizzle-orm";
import { uploadProcessedImage } from "../../utils/uploadProcessedImage.ts";

test("Creates an example profile successfully", async () => {
const insertProfilesValues = vi.fn().mockReturnValue({
Expand Down Expand Up @@ -251,3 +252,51 @@ test("Deletes a profile record if it no longer exists", async () => {
// The profile was deleted from the database
expect(deleteWhere).toBeCalledWith(eq(profiles.slug, "example"));
});

test("Rejects the profile image upload when the signal is already aborted", async () => {
const controller = new AbortController();
controller.abort();

const buffer = Buffer.from(
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAACklEQVR42mMAAQAABQABoIJXOQAAAABJRU5ErkJggg==",
"base64",
);
const stream = Readable.toWeb(
Readable.from(buffer),
) as ReadableStream<Uint8Array>;

await expect(
uploadProcessedImage(
stream,
"profiles/example.jpeg",
2048,
controller.signal,
),
).rejects.toThrow();
});

test("Aborts the in-flight pipeline when the s3 upload rejects first", async () => {
vi.mocked(s3.upload).mockImplementation(async () => {
throw new Error("Upload failed");
});

const buffer = Buffer.from(
"iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAACklEQVR42mMAAQAABQABoIJXOQAAAABJRU5ErkJggg==",
"base64",
);
const nodeSource = Readable.from(buffer);
const stream = Readable.toWeb(nodeSource) as ReadableStream<Uint8Array>;

await expect(
uploadProcessedImage(
stream,
"profiles/example.jpeg",
2048,
new AbortController().signal,
),
).rejects.toThrow("Upload failed");

// The upload rejecting should also unwind the pipeline reading from it,
// rather than leaving it hanging with nothing left to consume the stream
await vi.waitFor(() => expect(nodeSource.destroyed).toBe(true));
});
31 changes: 7 additions & 24 deletions apps/worker/src/tasks/sync-author/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,38 +7,16 @@ import {
authorRoles,
} from "@playfulprogramming/db";
import * as github from "@playfulprogramming/github-api";
import { s3 } from "@playfulprogramming/s3";
import { createProcessor } from "../../createProcessor.ts";
import matter from "gray-matter";
import { AuthorMetaSchema } from "./types.ts";
import { Value } from "typebox/value";
import sharp from "sharp";
import { Readable } from "node:stream";
import { and, eq, inArray } from "drizzle-orm";
import { MANUAL_ACHIEVEMENT_IDS } from "../grant-author-achievements/achievement-ids.ts";
import { uploadProcessedImage } from "../../utils/uploadProcessedImage.ts";

const PROFILE_IMAGE_SIZE_MAX = 2048;

async function processProfileImg(
stream: ReadableStream<Uint8Array>,
uploadKey: string,
) {
const pipeline = sharp()
.resize({
width: PROFILE_IMAGE_SIZE_MAX,
height: PROFILE_IMAGE_SIZE_MAX,
fit: "inside",
})
.jpeg({ mozjpeg: true });

const source = Readable.fromWeb(stream as never);
source.on("error", (err) => pipeline.destroy(err));
source.pipe(pipeline);

const bucket = await s3.ensureBucket(env.S3_BUCKET);
await s3.upload(bucket, uploadKey, undefined, pipeline, "image/jpeg");
}

export default createProcessor(Tasks.SYNC_AUTHOR, async (job, { signal }) => {
const authorId = job.data.author;
const authorMetaUrl = new URL(
Expand Down Expand Up @@ -87,7 +65,12 @@ export default createProcessor(Tasks.SYNC_AUTHOR, async (job, { signal }) => {
}

profileImgKey = `profiles/${authorId}.jpeg`;
await processProfileImg(profileImgStream, profileImgKey);
await uploadProcessedImage(
profileImgStream,
profileImgKey,
PROFILE_IMAGE_SIZE_MAX,
signal,
);
}

const result = {
Expand Down
20 changes: 20 additions & 0 deletions apps/worker/src/tasks/sync-collection/processor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { s3 } from "@playfulprogramming/s3";
import * as github from "@playfulprogramming/github-api";
import { Readable } from "node:stream";
import { eq } from "drizzle-orm";
import { uploadProcessedImage } from "../../utils/uploadProcessedImage.ts";

const mockImage = `iVBORw0KGgoAAAANSUhEUgAAAPIAAADOCAYAAAAE0F9yAAAACXBIWXMAABYZAAAWGQFJGZrZAAAN+ElEQVR4Aeyd65mjOBpGeSaU3c2hOpWqzqErlq4ccKfSymFm8tjdH7t6sbHxBSGEBLqceUYFBl0+nU8HsKu7/ce///Pf/1FgwBooew380fEfBCBQPAFELj6FTAACXYfIrAIIVECgZpErSA9TgIAfAUT240QtCGRNAJGzTg/BQcCPACL7caIWBLImgMhZp2c2OE5A4I4AIt/h4AUEyiSAyGXmjaghcEcAke9w8AICZRJA5DLzVnPUzC2AACIHQKMJBHIjgMi5ZYR4IBBAAJEDoNEEArkRQOTcMkI8NRNINjdEToaWjiGwHwFE3o81I0EgGQFEToaWjiGwHwFE3o81I0EgGYEMRE42NzqGQDMEELmZVDPRmgkgcs3ZZW7NEEDkZlLNRGsmgMhJs0vnENiHACLvw5lRIJCUACInxUvnENiHACLvw5lRIJCUACInxVtz58wtJwKInFM2iAUCgQQQORAczSCQEwFEzikbxAKBQAKIHAiOZjUTKG9uiFxezogYAk8EEPkJCQcgUB4BRC4vZ0QMgScCiPyEhAMQKI+Av8jlzY2IIdAMAURuJtVMtGYCiFxzdplbMwQQuZlUM9GaCSCyskuBQOEEELnwBBI+BEQAkUWBAoHCCSBy4QkkfAiIACKLQs2FuTVBAJGbSDOTrJ0AIteeYebXBAFEbiLNTLJ2Aohce4Zrnh9zuxJA5CsKdiBQLgFELjd3RA6BKwFEvqJgBwLlEkDkcnNH5DUTWDk3RF4JjOoQyJEAIueYFWKCwEoCiLwSGNUhkCMBRM4xK8QEgZUEihJ55dyoDoFmCCByZqk2xnQqX18/u+8fH522KjqWS6iKRUXxqSg+lVziazEORM4k6xLhX//8h5X3fShfP39aoX932qp8/3jvdF71jgh5FFcxKBYVY37fxahzKkfFeASXXMZE5IMzYewdeLirWXF9QpHUe8siMUdx18SoufnUp852Aoi8nWFwD1NB1nYiodV+bbu19TWGxlrbTvUlv9prn5KWACKn5TvbuxZ4qCBjp2qvfsbXsbfqW2Ns6VftuTNvIejXFpH9OEWvpQUeo1P1k0KUGBKP81OM4z7bNAQQOQ1XZ6+SxFlh5ckUosTs09gPxfQ5wMppUX0FAUReAStW1SVJ3t6+df3p11B+fH52Kq6xJYqxH5q56qw559PXmhg1trEya0tJQwCR03Cd7dUsCCdp+9Ope3t7G8qPH1ZkFSv0bKeRT5gF6frhIrM+RrMw98jTaKo7RN453cYhiSSWuK9C0nGdf3VOx5bu8qoToygGXWRe9aUYdad+dY5jaQkgclq+UXvfSxLXRUGyuiYl0efOu/qda8NxPwKI7McpWi3XYl6SZO5OqOBcd3qdp6QhkEuviLxzJlx3VcN7yJ2zUc9wiFxQLl2/tnJdINZO0dXX0sXGOD4DePv2tjYU6nsSQGRPULGqud5D6o80Gsdd2fw2s2HsJYnrrYFid52fDZ4Tmwkg8maEcTuQzI93XgmiP1BhHHe7mFG4LjaKQbE8jmfsBUixPx6fvl76DGBal/11BFKIvC6CxmrrAyvXo6tw6K6mv+EkYc7b9+GvC+rcXNlTEmMvKIpLZRrjXGw67ro46DxlGwFE3sYvqLXvopYwPgP49ufTl+roYuPbp2+M6peSjgAip2M72/MaUWY7uZyQcCnuxupz6cnhEsLiJlWMiwM3VAGRD0p2LFHUT6opSMCtfauPlDFuja+W9oi8LpNRa/en0+JfiJgbUHfLP//6e+50lON6ctAYGiukQyQOoRbWBpHDuEVrpbuVFvyaDlVfF4E1bbbU1Vgac00fqq+5rWlD3XACiBzOLlpLLXjd+YbF//nZPd4B9VqlH/7W0a9O9aMN7tmRxvSNcaj349OzZ6rFIIDIMShG6kOyqPT2kVsy9Fbc8/bU9faYHnVVIg0X1I3iU1E859h+deftLcagjmm0iQAib8KXtvHO0gZNpoQYgyZWWCNELixhhAuBVwQQ+RUVjkGgMAKIXFjCCBcCrwgg8isqHKuNQPXzQeTqU8wEWyCAyC1kmTlWTwCRq08xE2yBACK3kGXmWDOBYW6IPGDgBwTKJoDIZeeP6CEwEEDkAQM/IFA2AUQuO39ED4GBQKUiD3PjBwSaIYDIzaSaidZMAJFrzi5za4YAIjeTaiZaMwFELi67BAyBZwKI/MyEIxAojgAiF5cyAobAMwFEfmbCEQgURwCRi0tZzQEzt1ACTYtsjAnlRrtMCCiH+kbIr6+fnUomYe0eRpMij8n//vHe6atBW14Au6+4iAMqb8qhMb87fRXtUKzQEYcopqsmRR6TP2ZJC0BX9fE12/wJSGLl7TFSHTMNPmk1J7IWwGPy9droqt7o1VzzL61I2LmYlcu5c8cdTztycyK7cGpxmAav5i4mOZ6buxjnGOteMTUnsr4MzQVXj92u85w7loAk1gXXFYW+m8p1vsZzDYr89vRth4+J5f3yI5F8Xi9K/Nnmt0A2J7KWpL6+VNu5Yni/PIfm0OO6G7sCUF5bvBuLybEiK4IDir5BUEl3Da0rv+H9sgvRrucksXLiGrRVicWkSZE1cSV96f3y0sJRP5T0BHRBXcrF0oU5fZTHjtCsyMLen07O98uGR2xhOrQY+1S09AGkJNaF+dBADx68aZHFXotA27miO4Ee6+bOczwtgSWJNXrrEotB8yL7vl9eLbPoUjYR8PntwdKFeFMABTVuXmTlSld0n/fLyCxa+xRJbOxbG9dokli5c9Vp5RwiXzKtRXHZnd3wmD2LJuoJJF6PE5EvzPSI3Z9+XV7Nb5B5nk2MMz4S6+mJO/E9bUSe8EDmCQzXbqJzvhL39rcNiUIotltEfkgdMj8A2eGlGX7F9NGZhffEuhMj8euEIPILLmtkPt9F+JdGXmD0OqQPEPUrJiT2wjVbCZFn0Ehmnw/AtAC1ELUgZ7ri8AwBMdNnDjOn7w5zJ77D8fQCkZ+Q3A7oAxUfmdVCC1ILU/sUNwFzeZQWM3fN89ne40PIc80IPwvtApEXErdWZj1qL3TZ9OmzxO+L74dHSJJYT0fja7avCSDyay53R9fIbOwHNvyDfnf4ri/0xKK3IdcDjh19sPXnX393SOyANDmFyBMYrl3JfF5Y31zVruf02KiFa+xj5PVgoztioCcVMfFBoLczPb9i8kF1rYPIVxR+O1pgWmg+tbVwdQeS0D71a6szCiwGxj6p+Myvt++HddH0qUudGwEvkW/V2RMBLTRfmVVfQrf0uG3sU4juwGsE1qO0JOZRWitmfUHk9cyGFpL5vPD8HrXVqHahQwQWF10Ue/sojcSiEVYQOYzb0EoLTwtQC3E44PljKnQNj92hAguX2OmiqH1KOAFEDmd3bamFqAV5PeC5I6FVSnzsHuVV7GseoUc046O02I3H2IYTaF7kcHT3LbUg9al2iNDqaSp0znfpUeAQeTVPld5+oNXzKC0U0QoiR0N57khC93ahbhH6UWrJc+59358aV+X8wdXH8IV3WwTWXVgXO70l2Xcm9Y+GyAlyrIUqoSWzFm/oEBJaRfKcH2E/hq8OTXXHlrQqo7gaV8XYXx2phM5DDHRx6+1dOLQP2rkJILKbz6azklmLV0Jv6ujSWDJJbBWJrXKT7ia5RFcx9tdA06JjY5m20/65r/cuhriXcLupwLq4jcfZxieAyPGZPvUoofVIKaFVnipsOGAud0ttJfi0SMppmZ5T/WnZEMJTUwR+QpL8ACInR3wbQEKrjFJrwd/Olr2nufT2swHNrbeP0NyB980nIu/L+zqahNaC7+3ij32Xvg6SeEfyKnbNoUfexLTd3SOym0/ys7pzSWrdySSFSvJBNwwgeVV6ewHqrbyKXXPY0CVNIxBA5AgQY3UhKVQkdW9FkdQqsfoP6UfSqiieczl1vRX4aHlD5lJzG0TONLsSRVKrSGwVSa0isVKFrb41Rm8vJBqzt9KqKB6VVOPS7zYCiLyN366tJbWKxJJkKr0VbiwS8LFIzLFMz41ttFU/Y+mtuBoDaXdN7ebBEHkzwmM7kHBjkYCPRWKOZXpubKPtsTNg9BgEEDkGRfqAQEQCIV0hcgg12kAgMwKInFlCCAcCIQQQOYQabSCQGQFEziwhhAOBEAKliBwyN9pAoBkCiNxMqplozQQQuebsMrdmCCByM6lmojUTQOTjs0sEENhMAJE3I6QDCBxPAJGPzwERQGAzAUTejJAOIHA8AUQ+Pgc1R8DcdiKAyDuBZhgIpCSAyCnp0jcEdiKAyDuBZhgIpCSAyCnp0nfNBLKaGyJnlQ6CgUAYAUQO40YrCGRFAJGzSgfBQCCMACKHcaMVBLIiEFnkrOZGMBBohgAiN5NqJlozAUSuObvMrRkCiNxMqplozQQQ2Tu7VIRAvgQQOd/cEBkEvAkgsjcqKkIgXwKInG9uiAwC3gQQ2RtVzRWZW+kEELn0DBI/BCwBRLYQ+B8CpRNA5NIzSPwQsAQQ2ULg/5oJtDE3RG4jz8yycgKIXHmCmV4bBBC5jTwzy8oJIHLlCWZ6NRO4zQ2RbyzYg0CxBBC52NQROARuBBD5xoI9CBRLAJGLTR2BQ+BGoD6Rb3NjDwLNEEDkZlLNRGsmgMg1Z5e5NUMAkZtJNROtmQAil5RdYoXADAFEngHDYQiURACRS8oWsUJghgAiz4DhMARKIoDIJWWr5liZ2yYCiLwJH40hkAcBRM4jD0QBgU0EEHkTPhpDIA8CiJxHHoiiZgI7zA2Rd4DMEBBITQCRUxOmfwjsQACRd4DMEBBITeD/AAAA//+PgPcJAAAABklEQVQDAEL72xBBlWtZAAAAAElFTkSuQmCC`;

Expand Down Expand Up @@ -698,3 +699,22 @@ tags:
);
expect(deleteCollectionAuthorsWhere).toBeCalledTimes(2);
});

test("Rejects the image upload when the signal is already aborted", async () => {
const controller = new AbortController();
controller.abort();

const buffer = Buffer.from(mockImage, "base64");
const stream = Readable.toWeb(
Readable.from(buffer),
) as ReadableStream<Uint8Array>;

await expect(
uploadProcessedImage(
stream,
"collections/example-collection/en/cover.jpg",
2048,
controller.signal,
),
).rejects.toThrow();
});
36 changes: 13 additions & 23 deletions apps/worker/src/tasks/sync-collection/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,31 +13,11 @@ import { eq } from "drizzle-orm";
import matter from "gray-matter";
import { CollectionMetaSchema } from "./types.ts";
import { Value } from "typebox/value";
import sharp from "sharp";
import { Readable } from "node:stream";
import { s3 } from "@playfulprogramming/s3";
import { extractLocale } from "../../utils/extractLocale.ts";
import { uploadProcessedImage } from "../../utils/uploadProcessedImage.ts";

const IMAGE_SIZE_MAX = 2048;

async function processImg(
stream: ReadableStream<Uint8Array>,
uploadKey: string,
) {
const pipeline = sharp()
.resize({
width: IMAGE_SIZE_MAX,
height: IMAGE_SIZE_MAX,
fit: "inside",
})
.jpeg({ mozjpeg: true });

Readable.fromWeb(stream as never).pipe(pipeline);

const bucket = await s3.ensureBucket(env.S3_BUCKET);
await s3.upload(bucket, uploadKey, undefined, pipeline, "image/jpeg");
}

export default createProcessor(
Tasks.SYNC_COLLECTION,
async (job, { signal }) => {
Expand Down Expand Up @@ -150,7 +130,12 @@ export default createProcessor(
}

coverImgKey = `collections/${collectionId}/${locale}/cover.jpg`;
await processImg(coverImgStream, coverImgKey);
await uploadProcessedImage(
coverImgStream,
coverImgKey,
IMAGE_SIZE_MAX,
signal,
);
}

if (collectionParsedData.socialImg) {
Expand All @@ -176,7 +161,12 @@ export default createProcessor(
}

socialImgKey = `collections/${collectionId}/${locale}/social.jpg`;
await processImg(socialImgStream, socialImgKey);
await uploadProcessedImage(
socialImgStream,
socialImgKey,
IMAGE_SIZE_MAX,
signal,
);
}

const result = {
Expand Down
39 changes: 39 additions & 0 deletions apps/worker/src/utils/uploadProcessedImage.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
import { env } from "@playfulprogramming/common";
import { s3 } from "@playfulprogramming/s3";
import sharp from "sharp";
import { Readable } from "node:stream";
import { pipeline } from "node:stream/promises";

export async function uploadProcessedImage(
stream: ReadableStream<Uint8Array>,
uploadKey: string,
maxSize: number,
signal: AbortSignal,
) {
const transform = sharp()
.resize({ width: maxSize, height: maxSize, fit: "inside" })
.jpeg({ mozjpeg: true });

const source = Readable.fromWeb(stream as never);

const bucket = await s3.ensureBucket(env.S3_BUCKET);

// Abort the pipeline too if upload rejects first - nothing else destroys the streams
const uploadFailureController = new AbortController();
const combinedSignal = AbortSignal.any([
signal,
uploadFailureController.signal,
]);

const upload = s3
.upload(bucket, uploadKey, undefined, transform, "image/jpeg")
.catch((err: unknown) => {
uploadFailureController.abort(err);
throw err;
});

await Promise.all([
pipeline(source, transform, { signal: combinedSignal }),
upload,
]);
}
15 changes: 14 additions & 1 deletion apps/worker/test-utils/setup.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,25 @@
import "./server.ts";
import { vi, afterEach } from "vitest";
import { vi, afterEach, beforeEach } from "vitest";
import "@playfulprogramming/test-fixtures";
import { s3 } from "@playfulprogramming/s3";
import { Readable } from "node:stream";

afterEach(() => {
vi.clearAllMocks();
vi.setSystemTime(new Date("2025-05-05"));
});

beforeEach(() => {
// pipeline() won't resolve until the transform's readable side is drained
vi.mocked(s3.upload).mockImplementation(async (_bucket, _key, _tag, file) => {
if (file instanceof Readable) {
for await (const _chunk of file) {
// drain
}
}
});
});

vi.mock("@playfulprogramming/bullmq", async () => {
const tasks = await import("@playfulprogramming/bullmq/src/tasks/index.ts");
return {
Expand Down