Fix cleanup and simulate batch
This commit is contained in:
@@ -8,8 +8,9 @@ import type {
|
||||
Task,
|
||||
TaskList,
|
||||
TaskSpec,
|
||||
WorkerUtils,
|
||||
} from "graphile-worker";
|
||||
import { run as graphileRun, parseCronItems } from "graphile-worker";
|
||||
import { run as graphileRun, makeWorkerUtils, parseCronItems } from "graphile-worker";
|
||||
|
||||
import omit from "lodash.omit";
|
||||
import { z } from "zod";
|
||||
@@ -138,6 +139,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
#reporter?: ZodWorkerReporter;
|
||||
#shutdownTimeoutInMs?: number;
|
||||
#shuttingDown = false;
|
||||
#workerUtils?: WorkerUtils;
|
||||
|
||||
constructor(options: ZodWorkerOptions<TMessageCatalog>) {
|
||||
this.#name = options.name;
|
||||
@@ -166,6 +168,8 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
|
||||
const parsedCronItems = parseCronItems(this.#createCronItemsFromRecurringTasks());
|
||||
|
||||
this.#workerUtils = await makeWorkerUtils(this.#runnerOptions);
|
||||
|
||||
this.#runner = await graphileRun({
|
||||
...this.#runnerOptions,
|
||||
noHandleSignals: true,
|
||||
@@ -330,11 +334,11 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
>(
|
||||
identifier: K,
|
||||
payload: TPayload extends any[] ? TPayload : never,
|
||||
options?: ZodWorkerBatchEnqueueOptions
|
||||
options: ZodWorkerBatchEnqueueOptions
|
||||
): Promise<GraphileJob> {
|
||||
const task = this.#tasks[identifier];
|
||||
|
||||
const optionsWithoutTx = removeUndefinedKeys(omit(options ?? {}, ["tx"]));
|
||||
const optionsWithoutTx = removeUndefinedKeys(omit(options, ["tx"]));
|
||||
|
||||
// Make sure options passed in to enqueue take precedence over task options
|
||||
const spec = {
|
||||
@@ -364,7 +368,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
spec,
|
||||
});
|
||||
|
||||
const job = await this.#addBatchJob(
|
||||
const { job, durationInMs } = await this.#addBatchJob(
|
||||
identifier as string,
|
||||
payload,
|
||||
spec as BatchTaskSpec,
|
||||
@@ -376,6 +380,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
payload,
|
||||
spec,
|
||||
job,
|
||||
durationInMs,
|
||||
});
|
||||
|
||||
return job;
|
||||
@@ -444,6 +449,8 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
spec: BatchTaskSpec,
|
||||
tx: PrismaClientOrTransaction
|
||||
) {
|
||||
const now = performance.now();
|
||||
|
||||
const results = await tx.$queryRawUnsafe(
|
||||
`SELECT * FROM add_batch_job(
|
||||
identifier => $1::text,
|
||||
@@ -452,10 +459,10 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
queue_name => $4::text,
|
||||
run_at => $5::timestamptz,
|
||||
max_attempts => $6::int,
|
||||
max_payloads => $7::int,
|
||||
priority => $8::int,
|
||||
flags => $9::text[],
|
||||
job_key_mode => $10::text
|
||||
priority => $7::int,
|
||||
flags => $8::text[],
|
||||
job_key_mode => $9::text,
|
||||
max_payloads => $10::int
|
||||
)`,
|
||||
identifier,
|
||||
spec.jobKey,
|
||||
@@ -463,12 +470,14 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
spec.queueName || null,
|
||||
spec.runAt || null,
|
||||
spec.maxAttempts || null,
|
||||
spec.maxPayloads || null,
|
||||
spec.priority || null,
|
||||
spec.flags || null,
|
||||
spec.jobKeyMode || "preserve_run_at"
|
||||
spec.jobKeyMode || "preserve_run_at",
|
||||
spec.maxPayloads || null
|
||||
);
|
||||
|
||||
const durationInMs = performance.now() - now;
|
||||
|
||||
const rows = AddJobResultsSchema.safeParse(results);
|
||||
|
||||
if (!rows.success) {
|
||||
@@ -479,7 +488,7 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
|
||||
const job = rows.data[0];
|
||||
|
||||
return job as GraphileJob;
|
||||
return { job: job as GraphileJob, durationInMs: Math.floor(durationInMs) };
|
||||
}
|
||||
|
||||
async #removeJob(jobKey: string, tx: PrismaClientOrTransaction) {
|
||||
@@ -670,6 +679,10 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
return;
|
||||
}
|
||||
|
||||
if (!this.#workerUtils) {
|
||||
throw new Error("WorkerUtils need to be initialized before running job cleanup.");
|
||||
}
|
||||
|
||||
const job = helpers.job;
|
||||
|
||||
logger.debug("Received cleanup task", {
|
||||
@@ -696,22 +709,37 @@ export class ZodWorker<TMessageCatalog extends MessageCatalogSchema> {
|
||||
});
|
||||
|
||||
const rawResults = await this.#prisma.$queryRawUnsafe(
|
||||
`WITH rows AS (SELECT id FROM ${this.graphileWorkerSchema}.jobs WHERE run_at < $1 AND locked_at IS NULL AND max_attempts = attempts LIMIT $2 FOR UPDATE) DELETE FROM ${this.graphileWorkerSchema}.jobs WHERE id IN (SELECT id FROM rows) RETURNING id`,
|
||||
`SELECT id
|
||||
FROM ${this.graphileWorkerSchema}.jobs
|
||||
WHERE run_at > $1
|
||||
AND locked_at IS NULL
|
||||
AND max_attempts = attempts
|
||||
LIMIT $2`,
|
||||
expirationDate,
|
||||
this.#cleanup.maxCount
|
||||
);
|
||||
|
||||
const results = Array.isArray(rawResults) ? rawResults : [];
|
||||
const results = z
|
||||
.array(
|
||||
z.object({
|
||||
id: z.coerce.string(),
|
||||
})
|
||||
)
|
||||
.parse(rawResults);
|
||||
|
||||
const completedJobs = await this.#workerUtils.completeJobs(results.map((job) => job.id));
|
||||
|
||||
logger.debug("Cleaned up old jobs", {
|
||||
count: results.length,
|
||||
found: results.length,
|
||||
deleted: completedJobs.length,
|
||||
expirationDate,
|
||||
payload,
|
||||
});
|
||||
|
||||
if (this.#reporter) {
|
||||
await this.#reporter("cleanup_stats", {
|
||||
count: results.length,
|
||||
found: results.length,
|
||||
deleted: completedJobs.length,
|
||||
expirationDate,
|
||||
ts: payload._cron.ts,
|
||||
});
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
import { ActionArgs, json, redirect } from "@remix-run/server-runtime";
|
||||
import { ActionFunctionArgs, json } from "@remix-run/server-runtime";
|
||||
import { prisma } from "~/db.server";
|
||||
import { authenticateApiRequest } from "~/services/apiAuth.server";
|
||||
import { workerQueue } from "~/services/worker.server";
|
||||
|
||||
export async function action({ request }: ActionArgs) {
|
||||
export async function action({ request }: ActionFunctionArgs) {
|
||||
// Next authenticate the request
|
||||
const authenticationResult = await authenticateApiRequest(request);
|
||||
|
||||
@@ -36,7 +36,7 @@ export async function action({ request }: ActionArgs) {
|
||||
{
|
||||
jobKey: "simulateBatch",
|
||||
runAt: new Date(Date.now() + (body.deliverAfter ?? 0) * 1000),
|
||||
maxPayloads: body.maxPayloads ?? null
|
||||
maxPayloads: body.maxPayloads ?? null,
|
||||
}
|
||||
);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user