import { InjectQueue } from '@nestjs/bullmq'; import { Injectable, Logger } from '@nestjs/common'; import { Cron, CronExpression } from '@nestjs/schedule'; import { Queue } from 'bullmq'; import { env } from '../config/env'; import { MerchantsService } from '../merchants/merchants.service'; import { GenerationTrigger } from './generation.service'; export const GENERATION_QUEUE = 'keyword-generation'; export interface GenerationJob { merchantId: string; trigger: GenerationTrigger; } @Injectable() export class GenerationQueue { private readonly logger = new Logger(GenerationQueue.name); constructor( @InjectQueue(GENERATION_QUEUE) private readonly queue: Queue, private readonly merchants: MerchantsService, ) {} async enqueue(merchantId: string, trigger: GenerationTrigger): Promise { // 짧은 시간 내 같은 업체가 여러 번 발행돼도 한 번만 처리 (60초 dedupe 창) const job = await this.queue.add( 'generate', { merchantId, trigger }, { deduplication: { id: `${merchantId}-${trigger}`, ttl: 60_000 }, removeOnComplete: 100, removeOnFail: 500, attempts: 3, backoff: { type: 'exponential', delay: 5_000 }, }, ); return String(job.id); } /** 주기 리프레시: 매일 03:00, N일 지난 업체를 큐에 적재 */ @Cron(CronExpression.EVERY_DAY_AT_3AM) async scheduleRefresh() { const stale = await this.merchants.findStale(env.generation.refreshIntervalDays, 200); for (const m of stale) { await this.enqueue(m.id, 'scheduled'); } if (stale.length) this.logger.log(`scheduled refresh queued: ${stale.length} merchants`); } }