diff --git a/packages/backends/backend-test/src/createNewJob.ts b/packages/backends/backend-test/src/createNewJob.ts index 2885bc1..73ba7f3 100644 --- a/packages/backends/backend-test/src/createNewJob.ts +++ b/packages/backends/backend-test/src/createNewJob.ts @@ -1,4 +1,5 @@ import { NewJobData } from "@sidequest/backend"; +import { DuplicatedJobError } from "@sidequest/core"; import { describe, it } from "vitest"; import { backend } from "./backend"; @@ -124,7 +125,7 @@ export default function defineCreateNewJobTestSuite() { }; await backend.createNewJob(job); - await expect(backend.createNewJob(job2)).rejects.toThrow(); + await expect(backend.createNewJob(job2)).rejects.toThrow(DuplicatedJobError); }); }); } diff --git a/packages/backends/mongo/src/mongo-backend.ts b/packages/backends/mongo/src/mongo-backend.ts index c69552d..57316da 100644 --- a/packages/backends/mongo/src/mongo-backend.ts +++ b/packages/backends/mongo/src/mongo-backend.ts @@ -10,7 +10,7 @@ import { UpdateJobData, UpdateQueueData, } from "@sidequest/backend"; -import { JobData, JobState, QueueConfig } from "@sidequest/core"; +import { DuplicatedJobError, JobData, JobState, QueueConfig } from "@sidequest/core"; import { Collection, Db, Filter, MongoClient } from "mongodb"; import { addCoalescedField, generateTimeBuckets, getTimeRangeConfig, matchDateRange, parseTimeRange } from "./utils"; @@ -156,7 +156,20 @@ export default class MongoBackend implements Backend { inserted_at: now, available_at: job.available_at ?? now, }; - await this.jobs.insertOne(doc); + try { + await this.jobs.insertOne(doc); + } catch (error) { + if ( + error instanceof Error && + (("code" in error && error.code === 11000) || + error.message?.includes("E11000") || + error.message?.includes("unique_digest")) + ) { + throw new DuplicatedJobError(doc); + } + + throw error; + } return doc; } diff --git a/packages/backends/mysql/src/mysql-backend.ts b/packages/backends/mysql/src/mysql-backend.ts index e6e7dc5..f26242e 100644 --- a/packages/backends/mysql/src/mysql-backend.ts +++ b/packages/backends/mysql/src/mysql-backend.ts @@ -25,6 +25,31 @@ const defaultKnexConfig = { }, }; +const MYSQL_ER_DUP_ENTRY = 1062; + +function isUniqueDigestDuplicateError(error: unknown): boolean { + if (!(error instanceof Error)) { + return false; + } + + const code = "code" in error ? error.code : undefined; + const errno = "errno" in error ? error.errno : undefined; + const sqlMessage = "sqlMessage" in error && typeof error.sqlMessage === "string" ? error.sqlMessage : ""; + const constraint = "constraint" in error && typeof error.constraint === "string" ? error.constraint : ""; + const details = `${error.message}\n${sqlMessage}\n${constraint}`; + + if (!details.includes("unique_digest")) { + return false; + } + + return ( + code === "ER_DUP_ENTRY" || + errno === MYSQL_ER_DUP_ENTRY || + constraint === "sidequest_jobs_unique_digest_active_idx" || + /duplicate entry/i.test(details) + ); +} + export default class MysqlBackend extends SQLBackend { constructor(dbConfig: string | SQLDriverConfig) { const knexConfig: Knex.Config = { @@ -119,11 +144,7 @@ export default class MysqlBackend extends SQLBackend { return insertedJob; } catch (error) { - if ( - error instanceof Error && - (error.message?.includes("sidequest_jobs.unique_digest") || - ("constraint" in error && error.constraint === "sidequest_jobs_unique_digest_active_idx")) - ) { + if (isUniqueDigestDuplicateError(error)) { throw new DuplicatedJobError(job as JobData); }