Commit e796add5 authored by Grant's avatar Grant
Browse files

dedicate a pg client to each worker

parent 99c63ff7
Loading
Loading
Loading
Loading
+14 −5
Original line number Diff line number Diff line
@@ -4,15 +4,15 @@ export const client = new Client({
  connectionString: process.env.DATABASE_URL!,
});

export const query = {
const makeQuery = (cl: Client) => ({
  pixelsWithinTimeframe: (start: Date, end: Date) => {
    return client.query<DBPixel>(
    return cl.query<DBPixel>(
      `SELECT * FROM "Pixel" WHERE ("createdAt" >= $1 AND "createdAt" < $2) or ("deletedAt" >= $1 AND "deletedAt" < $2) ORDER BY "createdAt" ASC`,
      [start, end]
    );
  },
  existingPixelsBeforeTime: (end: Date) => {
    return client.query<DBPixel>(
    return cl.query<DBPixel>(
      `SELECT * FROM "Pixel" WHERE "createdAt" < $1 AND "deletedAt" IS NULL ORDER BY "createdAt" ASC`,
      [end]
    );
@@ -28,11 +28,20 @@ export const query = {
        return `($${params.length - 1}::int, $${params.length}::int)`;
      })
      .join(", ");
    return client.query<DBPixel>(
    return cl.query<DBPixel>(
      `SELECT * FROM "Pixel" WHERE (x, y) IN (${placeholders}) ORDER BY x, y, "createdAt" ASC`,
      params
    );
  },
});

export const query = makeQuery(client);

/** Create an independent database connection + query functions for a worker thread.
 *  Worker threads must NOT share the main thread's pg.Client — it's not thread-safe. */
export const createWorkerDb = () => {
  const cl = new Client({ connectionString: process.env.DATABASE_URL! });
  return { client: cl, query: makeQuery(cl) };
};

export type DBPixel = {
+3 −1
Original line number Diff line number Diff line
import fs from "fs";
import { parentPort } from "worker_threads";
import Canvas from "@napi-rs/canvas";
import { DBPixel, client, query } from "./lib/postgres";
import { createWorkerDb, DBPixel } from "./lib/postgres";

const { client, query } = createWorkerDb();

const canvas_size_parts = process.env.CANVAS_SIZE!.split(",");
const CANVAS_SIZE = [