{"version":3,"sources":["../../../../src/jobs/types/redis.ts"],"names":["config","redisCache","RedisCache","#queue","Queue","job","#callbacks","#cleanup","#crons","cron","#addCron","name","worker","callbacks","crons","#getNewId","Random","data","delayInMs","tz","jobId","failedJobs","RedisJob","repeatableJobs"],"mappings":"AAAA,6vCAA8B,8DAGrB,2DAEc,6DAKtB,IAAA,CAAA,CAAA,CAAA,CAAgB,EAAA,CAAA,CAAA,CAAA,OAAA,CAAA,SAChB,CAAA,CAAA,CAAA,aAAa,CAAA,eAYP,CAAA,CAAA,CAAA,UAEqB,CAAC,YAG5B,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,EAAYA,CAAAA,CAAwB,CACnC,CAAA,MAAMC,CAAa,CAAA,CAAA,CAAA,CAAIC,CAAAA,CAAWF,CAAAA,CAAO,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,WACxC,CAAA,CAAA,CAAA,CAAA,MAAA,CAAA,CAAA,IAAsB,4BAAA,CAAA,CAAA,CACtB,WAAA,CAAA,CAAA,oBAEgC,CAAA,IAAA,CAAA,gBAAqB,CAAA,CAAS,CAAA,CAC/D,CAAA,CAAA,CAAA,CAAA,qBAAKG,CAAAA,GAAS,CAAA,CAAA,CAAIC,aAAmB,CAAA,CAAA,CAAYH,SAAkB,CAAA,CAAA,IAAA,CAAA,CAAS,CAAA,CAAA,IAAA,kBAAA,CAAA,CAAA,CAAA,CAAA,UAAwB,CACpG,CAAA,CAAA,MAAe,CAAA,OAEd,CAAA,gBACSI,CAAAA,CAAI,CAAA,CAAA,CAAA,CAAA,MACN,CAAA,CAAA,IAAA,mBAAA,CAAA,CAAA,CAAA,MACJ,CAAA,EAAA,CAAQ,MAAKC,CAAAA,CAAW,CAAA,IAAA,CAAA,CAAA,IAAA,YACzB,CAAK,uBAAA,IACJ,mBAAA,CAAA,CAAA,qBAAA,SAAaA,0BAAW,CAAA,CAAA,CAAA,IAAA,GAAiBD,CAAAA,IAAI,SAC9C,CAAA,uBAAK,IAAA,qBAAA,CAAA,CAAA,qBAAA,MACJ,0BAAA,CAAA,CAAQ,CAAA,IAAKC,CAAAA,IAAW,GAAA,CAAA,IAAA,eAE3B,CAAA,uBACE,IAAA,qBAAA,CAAYL,CAAAA,qBAAW,YAAO,0BAAA,CAAS,CAAA,CAAA,IAAA,GAAA,CAAS,CAAA,CAAA,CAAO,UAAA,CAAA,CAAA,CAAA,MAAuB,CACjF,OAGC,CAAA,OACA,CAAA,CAAA,CAAA,CAAA,gBACO,CAAKM,CAAAA,CAAAA,CAAS,CAAA,CACpB,qBAAA,CAAA,EAAA,CAAA,OAAM,CAAA,KAAY,CAAA,CAAA,EAAA,CAAKC,MAAO,IAAO,CAAA,CAAA,CAAA,CAAAC,CAAAA,CAAM,MAAK,OAAM,CAAA,GAAKC,CAASC,IAAW,CAAC,CAAA,CAChFC,CAAAA,GAAO,CAAA,CAAI,CACZ,IAED,CACD,CAEA,CAAA,IAAI,CAAA,CAAA,CAAA,CAAA,EAAA,IAAUC,CAAyB,CACtC,CAAA,CAAA,CAAA,CAAA,CAAKP,CAAAA,CAAAA,CAAaO,CACnB,CAEA,CAAA,GAAI,CAAA,CAAA,CAAA,CAAA,EAAMC,CAAAA,CAAuC,IAChD,SAGD,CAAA,CAAA,CAAA,CAAOC,IAAY,CAClB,CAAA,CAAA,CAAA,CAAA,CAAA,IAAQ,KAAK,CAAA,CAAI,CAAA,CAAGC,IAAO,CAAA,CAAA,CAAA,CAAA,CAAO,CAAC,MAAE,CAAK,CAAA,CAAA,CAAG,CAC9C,MAEA,CAAM,IAAA,CAAA,GAAA,CAAA,CAAA,CAAWC,oBAAAA,CAAuBC,MAQvC,CAAA,CAAA,CAAA,CAPY,IAAA,CAAA,GAAM,CAAA,CAAA,MAAY,UAAI,CAAA,CAAA,CAAA,CAAA,CAAqBD,CAAAA,MACtD,CAAA,MAA0B,IAC1B,CAAA,CAAA,CAAA,CAAOC,GACP,CAAA,YAAA,CAAA,CAAA,CAAkB,CAAA,KAClB,CAAA,CAAA,CAAA,CAAA,CAAS,CAAA,CAAA,CAAA,KACT,CAAA,CAAA,CAAA,gBAEc,CAAA,CAAA,CAAS,CACzB,OAEM,CAAA,GAAA,CAAA,QAAA,CAAcD,CAAAA,CAA0BR,CAAAA,CAAcU,CAAAA,EAQ3D,CAAA,QAPY,CAAA,CAAA,CAAA,MAAWhB,aAAW,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,6DACjC,CAAA,MAAgBY,IAAU,CAC1B,CAAA,CAAA,CAAA,GAAA,CAAQ,eAAqBI,CAAK,CAAE,CAAA,CAAA,KAAQ,CAAG,CAAA,CAC/C,CAAA,CAAA,CAAA,CAAA,CAAA,MAAA,CAAA,CAAA,OACA,CAAA,CAAA,CAAA,GAAA,CAAA,CAAS,CAAA,EAAA,CACT,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,CAAA,gBAEgB,CAAA,CAAA,CAAA,CAAA,OAAe,CACjC,GAEA,CAAA,QAAM,CAAA,CAAA,CAAA,CAAA,CAAA,6BAAA,IAAcC,qCACnB,MAAMf,qCAAM,KAAA,eAAA,IAAM,CAAA,MAAY,aACrB,CAAA,CAAA,CAAA,CAAA,MAAU,CAAA,CAAA,MAGpB,IAAA,CAAM,CAAA,CAAA,CAAA,MAAA,CAAA,CAAA,CAAA,CAAA,CAAA,EAAA,MACL,CAAA,CAAA,MAAmB,CAAA,CAAA,CAAA,MAAM,kBACzB,CAAA,CAAA,CAAA,MAAM,CAAA,CAAA,MAAYgB,IAAW,CAAKhB,CAAAA,CAAAA,CAAQA,SAAW,CAAC,CACvD,CAEA,MAAMK,OAYL,CAAA,GAAA,CAAA,CAAA,CAXY,GAAA,CAAA,CAAA,EAAM,CAAA,CAAA,KAAY,CAAA,CAAA,CAAA,CAC7B,CAAA,KAAA,CAAA,CAAA,CAAA,CACA,CAAE,CAAA,CAAA,CAAA,MAED,CAAA,MAAOY,IAAmB,CAAA,CAC1B,CAAA,CAAA,GAAA,CAAA,SAAU,CAASb,CAAK,IACxB,CAAA,CAAA,CAAA,CAAA,CAAA,KAAA,CAAA,CAAA,CAAA,CAAA,CAAkB,CAAA,CAAA,CAClB,MAAA,CAAA,CAAS,OACT,CAAA,CAAA,CAAA,CAAA,gBAGa,CAAA,CAAA,CAAS,CACzB,OAEMF,CAAAA,GACL,CAAA,QAAM,CAAA,CAAK,CAAA,CAAA,CAAA,CAAA,EAAA,CAAA,QAAA,CAAA,CAAA,CAAA,KACX,CAAA,CAAA,CAAA,CAAMgB,CAAAA,MAAiB,IAAM,CAAA,kBAAY,CAAA,CAAA,CAAiB,MAC1D,CAAA,CAAM,MAAA,IAAQ,CAAA,CAAIA,CAAAA,CAAe,gBAAyB,CAAA,CAAA,CAAA,MAAA,OAAA,CAAA,GAAuB,CAAA,CAAA,CAAG,GACrF,CACD,CAAA,EAAA,IAAA,CAAA,CAAA,CAAA,CAAA,kBAAA,CAAA,CAAA,CAAA,GAAA,CAAA,CAAA,CAAA,CAAA,CAAA,qBAAA","file":"/home/runner/work/equipped/equipped/dist/cjs/jobs/types/redis.min.cjs","sourcesContent":["import { Queue, Worker } from 'bullmq'\n\nimport { RedisCache } from '../../cache/types/redis'\nimport { Instance } from '../../instance'\nimport type { CronTypes, DelayedJobs, RepeatableJobs } from '../../types'\nimport { Random } from '../../utilities'\nimport type { RedisJobConfig } from '../pipes'\n\nenum JobNames {\n\tCronJob = 'CronJob',\n\tRepeatableJob = 'RepeatableJob',\n\tDelayedJob = 'DelayedJob',\n}\n\ntype Cron = CronTypes[keyof CronTypes]\ntype DelayedJobEvent = DelayedJobs[keyof DelayedJobs]\ntype RepeatableJobEvent = RepeatableJobs[keyof RepeatableJobs]\ntype DelayedJobCallback = (data: DelayedJobEvent) => Promise<void> | void\ntype CronJobCallback = (name: CronTypes[keyof CronTypes]) => Promise<void> | void\ntype RepeatableJobCallback = (data: RepeatableJobEvent) => Promise<void> | void\n\ntype JobCallbacks = { onDelayed?: DelayedJobCallback; onCron?: CronJobCallback; onRepeatable?: RepeatableJobCallback }\n\nexport class RedisJob {\n\t#queue: Queue\n\t#callbacks: JobCallbacks = {}\n\t#crons: { name: Cron; cron: string }[] = []\n\n\tconstructor(config: RedisJobConfig) {\n\t\tconst redisCache = new RedisCache(config.redisConfig, {\n\t\t\tmaxRetriesPerRequest: null,\n\t\t\tenableReadyCheck: false,\n\t\t})\n\t\tconst queueName = Instance.get().getScopedName(config.queueName)\n\t\tthis.#queue = new Queue(queueName, { connection: redisCache.client.options, skipVersionCheck: true })\n\t\tconst worker = new Worker(\n\t\t\tqueueName,\n\t\t\tasync (job) => {\n\t\t\t\tswitch (job.name) {\n\t\t\t\t\tcase JobNames.DelayedJob:\n\t\t\t\t\t\treturn (this.#callbacks.onDelayed as any)?.(job.data)\n\t\t\t\t\tcase JobNames.CronJob:\n\t\t\t\t\t\treturn (this.#callbacks.onCron as any)?.(job.data.type)\n\t\t\t\t\tcase JobNames.RepeatableJob:\n\t\t\t\t\t\treturn (this.#callbacks.onRepeatable as any)?.(job.data)\n\t\t\t\t}\n\t\t\t},\n\t\t\t{ connection: redisCache.client.options, autorun: false, skipVersionCheck: true },\n\t\t)\n\n\t\tInstance.on(\n\t\t\t'start',\n\t\t\tasync () => {\n\t\t\t\tawait this.#cleanup()\n\t\t\t\tawait Promise.all(this.#crons.map(({ cron, name }) => this.#addCron(name, cron)))\n\t\t\t\tworker.run()\n\t\t\t},\n\t\t\t10,\n\t\t)\n\t}\n\n\tset callbacks(callbacks: JobCallbacks) {\n\t\tthis.#callbacks = callbacks\n\t}\n\n\tset crons(crons: { name: Cron; cron: string }[]) {\n\t\tthis.#crons = crons\n\t}\n\n\tstatic #getNewId() {\n\t\treturn [Date.now(), Random.string()].join('_')\n\t}\n\n\tasync addDelayed(data: DelayedJobEvent, delayInMs: number): Promise<string> {\n\t\tconst job = await this.#queue.add(JobNames.DelayedJob, data, {\n\t\t\tjobId: RedisJob.#getNewId(),\n\t\t\tdelay: delayInMs,\n\t\t\tremoveOnComplete: true,\n\t\t\tbackoff: 1000,\n\t\t\tattempts: 3,\n\t\t})\n\t\treturn job.id!.toString()\n\t}\n\n\tasync addRepeatable(data: RepeatableJobEvent, cron: string, tz?: string): Promise<string> {\n\t\tconst job = await this.#queue.add(JobNames.RepeatableJob, data, {\n\t\t\tjobId: RedisJob.#getNewId(),\n\t\t\trepeat: { pattern: cron, ...(tz ? { tz } : {}) },\n\t\t\tremoveOnComplete: true,\n\t\t\tbackoff: 1000,\n\t\t\tattempts: 3,\n\t\t})\n\t\treturn job.opts?.repeat?.key ?? ''\n\t}\n\n\tasync removeDelayed(jobId: string) {\n\t\tconst job = await this.#queue.getJob(jobId)\n\t\tif (job) await job.remove()\n\t}\n\n\tasync retryAllFailedJobs() {\n\t\tconst failedJobs = await this.#queue.getFailed()\n\t\tawait Promise.all(failedJobs.map((job) => job.retry()))\n\t}\n\n\tasync #addCron(type: Cron | string, cron: string): Promise<string> {\n\t\tconst job = await this.#queue.add(\n\t\t\tJobNames.CronJob,\n\t\t\t{ type },\n\t\t\t{\n\t\t\t\tjobId: RedisJob.#getNewId(),\n\t\t\t\trepeat: { pattern: cron },\n\t\t\t\tremoveOnComplete: true,\n\t\t\t\tbackoff: 1000,\n\t\t\t\tattempts: 3,\n\t\t\t},\n\t\t)\n\t\treturn job.id!.toString()\n\t}\n\n\tasync #cleanup() {\n\t\tawait this.retryAllFailedJobs()\n\t\tconst repeatableJobs = await this.#queue.getJobSchedulers()\n\t\tawait Promise.all(repeatableJobs.map((job) => this.#queue.removeJobScheduler(job.key)))\n\t}\n}\n"]}