返回 AiToEarn
queue-processor.decorator.ts
根目录 / project / aitoearn-backend / libs / aitoearn-queue / src / decorators / queue-processor.decorator.ts
1 import type { ProcessorOptions } from '@nestjs/bullmq'
2 import type { NestWorkerOptions } from '@nestjs/bullmq/dist/interfaces/worker-options.interface'
3 import type { Job } from 'bullmq'
4
5 import { inspect } from 'node:util'
6 import { Processor } from '@nestjs/bullmq'
7 import { Logger } from '@nestjs/common'
8 import { PinoLogger } from 'nestjs-pino'
9 import { storage, Store } from 'nestjs-pino/storage'
10
11 const ON_WORKER_EVENT_METADATA = 'bullmq:worker_events_metadata'
12
13 type ProcessFn = (job: Job, token?: string) => Promise<unknown>
14 type EventHandlerFn = (job: Job, ...args: unknown[]) => unknown
15 interface WorkerEventMetadata { eventName: string }
16
17 function ensureStore(bindings: Record<string, unknown>): Store {
18 const existing = storage.getStore()
19 if (existing) {
20 existing.logger = existing.logger.child(bindings)
21 return existing
22 }
23 return new Store(PinoLogger.root.child(bindings))
24 }
25
26 function wrapProcess(proto: Record<string, unknown>, name: string): void {
27 const originalProcess = proto['process'] as ProcessFn | undefined
28 if (!originalProcess)
29 return
30
31 const logger = new Logger(name)
32
33 proto['process'] = function (
34 this: unknown,
35 job: Job,
36 token?: string,
37 ): Promise<unknown> {
38 const bindings = { jobId: job.id, queue: job.queueName, jobName: job.name }
39 const store = ensureStore(bindings)
40
41 const run = async (): Promise<unknown> => {
42 const attemptsMade = job.attemptsMade
43 const attempt = attemptsMade + 1
44 const maxAttempts = job.opts.attempts ?? 1
45 const startedAt = Date.now()
46
47 try {
48 logger.log({
49 data: job.data,
50 attempt,
51 maxAttempts,
52 }, 'Job started')
53
54 const result = await originalProcess.call(this, job, token)
55
56 logger.log({
57 attempt,
58 maxAttempts,
59 durationMs: Date.now() - startedAt,
60 }, 'Job completed')
61
62 return result
63 }
64 catch (error) {
65 const durationMs = Date.now() - startedAt
66 const details = `data=${inspect(job.data, { depth: 5, breakLength: Infinity })} attempt=${attempt} maxAttempts=${maxAttempts} durationMs=${durationMs} isLastAttempt=${attempt >= maxAttempts}`
67
68 if (attempt >= maxAttempts) {
69 logger.error(error, `Job failed, no more retries ${details}`)
70 }
71 else {
72 logger.error(error, `Job failed, will retry ${details}`)
73 }
74
75 throw error
76 }
77 }
78
79 if (storage.getStore() === store) {
80 return run()
81 }
82 return storage.run(store, run) as Promise<unknown>
83 }
84 }
85
86 function wrapWorkerEvents(proto: Record<string, unknown>): void {
87 const methodNames = Object.getOwnPropertyNames(proto)
88
89 for (const methodName of methodNames) {
90 if (methodName === 'constructor' || methodName === 'process')
91 continue
92
93 const method = proto[methodName]
94 if (typeof method !== 'function')
95 continue
96
97 const metadata: WorkerEventMetadata | undefined = Reflect.getMetadata(
98 ON_WORKER_EVENT_METADATA,
99 method,
100 )
101 if (!metadata)
102 continue
103
104 const original = method as EventHandlerFn
105 const wrapped = function (this: unknown, job: Job, ...args: unknown[]): unknown {
106 const bindings: Record<string, unknown> = {
107 event: metadata.eventName,
108 }
109 if (job?.id) {
110 bindings['jobId'] = job.id
111 }
112 if (job?.queueName) {
113 bindings['queue'] = job.queueName
114 }
115
116 const logger = PinoLogger.root.child(bindings)
117 const store = new Store(logger)
118 return storage.run(store, () => original.call(this, job, ...args))
119 }
120
121 const keys: string[] = Reflect.getMetadataKeys(method)
122 for (const key of keys) {
123 Reflect.defineMetadata(key, Reflect.getMetadata(key, method), wrapped)
124 }
125
126 proto[methodName] = wrapped
127 }
128 }
129
130 export function QueueProcessor(queueName: string): ClassDecorator
131 export function QueueProcessor(queueName: string, workerOptions: NestWorkerOptions): ClassDecorator
132 export function QueueProcessor(processorOptions: ProcessorOptions): ClassDecorator
133 export function QueueProcessor(processorOptions: ProcessorOptions, workerOptions: NestWorkerOptions): ClassDecorator
134 export function QueueProcessor(queueNameOrOptions: string | ProcessorOptions, workerOptions?: NestWorkerOptions): ClassDecorator {
135 return (target) => {
136 if (workerOptions) {
137 Processor(queueNameOrOptions as string, workerOptions)(target)
138 }
139 else {
140 Processor(queueNameOrOptions as string)(target)
141 }
142
143 const proto = target.prototype as Record<string, unknown>
144 wrapProcess(proto, target.name)
145 wrapWorkerEvents(proto)
146 }
147 }
148
148 lines TYPESCRIPT