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
4 changes: 2 additions & 2 deletions deno/api/mod.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { SSESource } from "@jsr/planigale__sse";
import { SSESource, type SSESourceInit } from "@jsr/planigale__sse";
import {
ApiErrorResponse,
Channel,
Expand Down Expand Up @@ -167,7 +167,7 @@ class API extends EventTarget {
headers: {
Authorization: `Bearer ${this.token}`,
},
});
} as SSESourceInit);
if (!this.source) return;
this.emit(new CustomEvent("con:open", { detail: {} }));
// @ts-ignore For some reason the AsyncIterator is not recognized
Expand Down
13 changes: 13 additions & 0 deletions deno/storage/src/core/env.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
export const getEnvInt = (name: string, fallback: number): number => {
let raw: string | undefined;
try {
raw = Deno.env.get(name);
} catch {
return fallback;
}
if (raw === undefined) return fallback;
const value = Number(raw);
if (Number.isInteger(value) && value > 0) return value;
console.warn(`[storage] invalid ${name}="${raw}", using ${fallback}`);
return fallback;
};
44 changes: 20 additions & 24 deletions deno/storage/src/core/mod.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { PhotonImage, resize, SamplingFilter } from "@cf-wasm/photon/node";
import type { Config } from "@quack/config";
import type { FileData, FileOpts } from "./types.ts";
import { files } from "./store/mod.ts";
import { getResizePool } from "./resizePool.ts";
import { ApiError } from "@planigale/planigale";

type ScalingOpts = {
Expand Down Expand Up @@ -87,33 +87,29 @@ class Files {
width?: number,
height?: number,
): Promise<ReadableStream<Uint8Array> | null> {
let img: PhotonImage | undefined;
let out: PhotonImage | undefined;
const pool = getResizePool();
if (!pool.hasCapacity()) {
console.warn("[storage] resize pool saturated, serving original");
return null;
}
let bytes: Uint8Array<ArrayBuffer>;
try {
const bytes = new Uint8Array(
await new Response(file.stream).arrayBuffer(),
);
img = PhotonImage.new_from_byteslice(bytes);

const ow = img.get_width();
const oh = img.get_height();
let w = width || 0;
let h = height || 0;
if (!w) w = Math.max(1, Math.round((ow / oh) * h));
if (!h) h = Math.max(1, Math.round((oh / ow) * w));

out = resize(img, w, h, SamplingFilter.Lanczos3);
const result = file.contentType === "image/png"
? out.get_bytes()
: out.get_bytes_jpeg(90);
return new Blob([new Uint8Array(result)]).stream();
bytes = new Uint8Array(await new Response(file.stream).arrayBuffer());
} catch (e) {
console.warn("[storage] thumbnail resize failed, serving original", e);
console.warn("[storage] reading image to resize failed", e);
return null;
}
const resized = await pool.resize(
bytes,
width || 0,
height || 0,
file.contentType === "image/png",
);
if (!resized) {
console.warn("[storage] thumbnail resize failed, serving original");
return null;
} finally {
img?.free();
out?.free();
}
return new Blob([resized]).stream();
}
}

Expand Down
154 changes: 154 additions & 0 deletions deno/storage/src/core/resizePool.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
import type { ResizeRequest, ResizeResponse } from "./resizeWorker.ts";
import { getEnvInt } from "./env.ts";

type Pending = (bytes: Uint8Array<ArrayBuffer> | null) => void;

type Task = {
request: ResizeRequest;
resolve: Pending;
};

type Active = {
resolve: Pending;
timer: number;
};

const DEFAULT_WORKERS = 2;
const DEFAULT_TIMEOUT_MS = 30_000;

export class ResizePool {
private workers = new Set<Worker>();
private idle: Worker[] = [];
private queue: Task[] = [];
private active = new Map<Worker, Active>();
private started = false;
private closed = false;

constructor(
private readonly size: number,
private readonly maxPending: number = size * 4,
private readonly timeoutMs: number = DEFAULT_TIMEOUT_MS,
) {}

private start() {
if (this.started) return;
this.started = true;
for (let i = 0; i < this.size; i++) {
this.idle.push(this.spawn());
}
}

private spawn(): Worker {
const worker = new Worker(
new URL("./resizeWorker.ts", import.meta.url),
{ type: "module" },
);
this.workers.add(worker);
worker.onmessage = (event: MessageEvent<ResizeResponse>) => {
const entry = this.active.get(worker);
this.active.delete(worker);
if (entry) {
clearTimeout(entry.timer);
entry.resolve(event.data.ok ? event.data.bytes : null);
}
this.release(worker);
};
worker.onerror = (event) => {
event.preventDefault();
this.discard(worker);
};
return worker;
}

private discard(worker: Worker) {
const entry = this.active.get(worker);
this.active.delete(worker);
if (entry) {
clearTimeout(entry.timer);
entry.resolve(null);
}
this.workers.delete(worker);
worker.terminate();
if (this.closed) return;
this.idle.push(this.spawn());
this.drain();
}

private release(worker: Worker) {
if (this.closed) return;
this.idle.push(worker);
this.drain();
}

private drain() {
while (this.queue.length > 0 && this.idle.length > 0) {
const worker = this.idle.shift()!;
const task = this.queue.shift()!;
const timer = setTimeout(() => {
console.warn("[storage] thumbnail resize timed out, serving original");
this.discard(worker);
}, this.timeoutMs);
this.active.set(worker, { resolve: task.resolve, timer });
worker.postMessage(task.request, [task.request.bytes.buffer]);
}
}

hasCapacity(): boolean {
return !this.closed &&
this.active.size + this.queue.length < this.maxPending;
}

resize(
bytes: Uint8Array<ArrayBuffer>,
width: number,
height: number,
png: boolean,
): Promise<Uint8Array<ArrayBuffer> | null> {
if (this.closed) return Promise.resolve(null);
this.start();
return new Promise<Uint8Array<ArrayBuffer> | null>((resolve) => {
const request: ResizeRequest = { bytes, width, height, png };
this.queue.push({ request, resolve });
this.drain();
});
}

close() {
this.closed = true;
for (const task of this.queue) {
task.resolve(null);
}
this.queue = [];
for (const entry of this.active.values()) {
clearTimeout(entry.timer);
entry.resolve(null);
}
this.active.clear();
for (const worker of this.workers) {
worker.terminate();
}
this.workers.clear();
this.idle = [];
}
}

let pool: ResizePool | undefined;

export const getResizePool = (): ResizePool => {
if (!pool) {
const workers = getEnvInt("STORAGE_RESIZE_WORKERS", DEFAULT_WORKERS);
pool = new ResizePool(
workers,
getEnvInt("STORAGE_RESIZE_MAX_PENDING", workers * 4),
getEnvInt("STORAGE_RESIZE_TIMEOUT_MS", DEFAULT_TIMEOUT_MS),
);
}
return pool;
};

export const closeResizePool = () => {
pool?.close();
pool = undefined;
};

globalThis.addEventListener("unload", () => closeResizePool());
40 changes: 40 additions & 0 deletions deno/storage/src/core/resizeWorker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
/// <reference lib="deno.worker" />
import { PhotonImage, resize, SamplingFilter } from "@cf-wasm/photon/node";

export type ResizeRequest = {
bytes: Uint8Array<ArrayBuffer>;
width: number;
height: number;
png: boolean;
};

export type ResizeResponse =
| { ok: true; bytes: Uint8Array<ArrayBuffer> }
| { ok: false };

self.onmessage = (event: MessageEvent<ResizeRequest>) => {
const { bytes, width, height, png } = event.data;
let img: PhotonImage | undefined;
let out: PhotonImage | undefined;
try {
img = PhotonImage.new_from_byteslice(bytes);

const ow = img.get_width();
const oh = img.get_height();
let w = width || 0;
let h = height || 0;
if (!w) w = Math.max(1, Math.round((ow / oh) * h));
if (!h) h = Math.max(1, Math.round((oh / ow) * w));

out = resize(img, w, h, SamplingFilter.Lanczos3);
const result = new Uint8Array(
png ? out.get_bytes() : out.get_bytes_jpeg(90),
);
self.postMessage({ ok: true, bytes: result }, [result.buffer]);
} catch {
self.postMessage({ ok: false });
} finally {
img?.free();
out?.free();
}
};
Loading
Loading