| 54 | } |
| 55 | |
| 56 | export class TencentCmqQueueProvider implements QueueProvider { |
| 57 | readonly providerName = 'tdmq-cmq' as const; |
| 58 | private readonly receipts = new Map<string, string>(); |
| 59 | |
| 60 | constructor( |
| 61 | private readonly store: JobStore = jobStore, |
| 62 | private readonly client: TencentCmqClient = createTencentCmqClient() |
| 63 | ) {} |
| 64 | |
| 65 | async publish(message: QueueJobMessage, options: PublishOptions = {}): Promise<void> { |
| 66 | await this.client.sendMessage(message, options.delaySeconds); |
| 67 | } |
| 68 | |
| 69 | async claimNext(input: ClaimJobInput): Promise<CloudJob | undefined> { |
| 70 | const received = await this.client.receiveMessage(); |
| 71 | if (!received) return undefined; |
| 72 | |
| 73 | const existing = await this.store.get(received.body.jobId); |
| 74 | if (!existing) { |
| 75 | await this.client.deleteMessage(received.receiptHandle); |
| 76 | return undefined; |
| 77 | } |
| 78 | const job = await this.store.claim(received.body.jobId, input); |
| 79 | if (!job) { |
| 80 | const latest = (await this.store.get(received.body.jobId)) || existing; |
| 81 | if (isTerminalJobStatus(latest.status)) { |
| 82 | await this.client.deleteMessage(received.receiptHandle); |
| 83 | } else if (latest.status === 'processing') { |
| 84 | await this.deferWatchdogMessage(received.receiptHandle, latest); |
| 85 | } else if (latest.status === 'retry_wait') { |
| 86 | await this.deferWatchdogMessage(received.receiptHandle, latest); |
| 87 | } else { |
| 88 | await this.client.deleteMessage(received.receiptHandle); |
| 89 | } |
| 90 | return undefined; |
| 91 | } |
| 92 | this.receipts.set(job.id, received.receiptHandle); |
| 93 | return job; |
| 94 | } |
| 95 | |
| 96 | async complete(job: CloudJob): Promise<void> { |
| 97 | await this.deleteReceipt(job.id); |
| 98 | } |
| 99 | |
| 100 | async release(job: CloudJob, options: PublishOptions = {}): Promise<void> { |
| 101 | await this.deleteReceipt(job.id); |
| 102 | if (job.status === 'retry_wait') { |
| 103 | await this.publish( |
| 104 | { |
| 105 | jobId: job.id, |
| 106 | tenantId: job.tenantId, |
| 107 | taskType: job.taskType, |
| 108 | attempt: job.attempts, |
| 109 | traceId: job.externalJobId, |
| 110 | }, |
| 111 | options |
| 112 | ); |
| 113 | } |
nothing calls this directly
no outgoing calls
no test coverage detected