File size: 4,325 Bytes
3144483
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
import {
  RastermillError,
  RastermillUnavailableError,
  type ImageInput,
  type Rastermill,
} from "rastermill";
import { runtimeProcessEntrypoints } from "../infra/runtime-process-entrypoints.js";
import { resolveRuntimeWorkerUrl } from "../infra/runtime-worker-url.js";
import { WorkerTaskPool } from "../infra/worker-task-pool.js";
import {
  createLocalImageProcessor,
  MAX_IMAGE_INPUT_PIXELS,
  type ImageProcessorPixelLimits,
} from "./image-processor-config.js";
import type {
  ImageProcessorOperation,
  ImageProcessorReply,
  ImageProcessorRequest,
} from "./image-processor.types.js";

const pool = new WorkerTaskPool<ImageProcessorRequest, ImageProcessorReply>({
  workerUrl: resolveRuntimeWorkerUrl(runtimeProcessEntrypoints.imageProcessor),
  // Each Photon instance retains a WASM heap; serialize transforms rather than multiply decodes.
  maxWorkers: 1,
  sharedCompute: true,
});

async function runImageTask(
  input: ImageInput,
  operation: ImageProcessorOperation,
  signal?: AbortSignal,
  limits?: ImageProcessorPixelLimits,
): Promise<ImageProcessorReply> {
  const reply = await pool.run(
    () => ({
      ...operation,
      ...(limits ? { limits } : {}),
      // The caller may reuse its Buffer. Transfer a dedicated copy only after admission.
      input: Uint8Array.from(input instanceof ArrayBuffer ? new Uint8Array(input) : input),
    }),
    {
      timeoutMs: 180_000,
      signal,
      inputBytes: input.byteLength,
      transferList: (request) => [request.input.buffer],
    },
  );
  // Structured cloning preserves Error messages but drops Rastermill's public error codes.
  if (reply.kind === "failed" && reply.code && !reply.unavailable) {
    reply.error = new RastermillError(reply.code, reply.error.message, { cause: reply.error });
  }
  return reply;
}

/** Keep cheap probes local and move in-process image computation off the caller's event loop. */
export function createImageProcessor(): Rastermill {
  return createImageProcessorWithPixelLimits({
    inputPixels: MAX_IMAGE_INPUT_PIXELS,
    outputPixels: MAX_IMAGE_INPUT_PIXELS,
  });
}

/** Internal operation-specific admission uses the same worker and native fallback owners. */
export function createImageProcessorWithPixelLimits(params: ImageProcessorPixelLimits): Rastermill {
  const limits = { inputPixels: params.inputPixels, outputPixels: params.outputPixels };
  const local = createLocalImageProcessor("auto", limits);
  const workerLimits =
    limits.inputPixels === MAX_IMAGE_INPUT_PIXELS && limits.outputPixels === MAX_IMAGE_INPUT_PIXELS
      ? undefined
      : limits;
  return {
    probe: (input) => local.probe(input),
    transparency: async (input) => {
      const reply = await runImageTask(input, { kind: "transparency" }, undefined, workerLimits);
      if (reply.kind === "failed") {
        throw reply.unavailable
          ? new RastermillUnavailableError("transparency", reply.error.message, [reply.error])
          : reply.error;
      }
      if (reply.kind !== "transparency") {
        throw new Error("Unexpected image worker result");
      }
      return reply.value;
    },
    encode: async (input, options) => {
      const { signal, ...workerOptions } = options ?? {};
      const reply = await runImageTask(
        input,
        { kind: "encode", options: workerOptions },
        signal,
        workerLimits,
      );
      signal?.throwIfAborted();
      if (reply.kind === "failed") {
        if (!reply.unavailable) {
          throw reply.error;
        }
        // Preserve Rastermill's native-codec and alpha policy when its internal backend declines.
        return local.encode(input, options);
      }
      if (reply.kind !== "encode") {
        throw new Error("Unexpected image worker result");
      }
      const { data, ...value } = reply.value;
      return { ...value, data: Buffer.from(data.buffer, data.byteOffset, data.byteLength) };
    },
  };
}

export async function convertBmpToPngWithWorker(input: Buffer): Promise<Buffer> {
  const reply = await runImageTask(input, { kind: "bmpToPng" });
  if (reply.kind === "failed") {
    throw reply.error;
  }
  if (reply.kind !== "bmpToPng") {
    throw new Error("Unexpected image worker result");
  }
  return Buffer.from(reply.value.buffer, reply.value.byteOffset, reply.value.byteLength);
}