Bot & Automation

AI Creator Studio Telegram

/root/hermes-projects/AI Creator Studio Telegram

apps/worker/src/index.ts text
import { QueueEvents, Worker } from "bullmq";
import { logger } from "@ai-creator/core";
import { mockProvider } from "@ai-creator/provider-adapters";
import { QUEUE_NAMES } from "@ai-creator/queue";

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const parsedRedisUrl = new URL(redisUrl);
const connection = {
  host: parsedRedisUrl.hostname,
  port: Number(parsedRedisUrl.port || 6379),
  username: parsedRedisUrl.username || undefined,
  password: parsedRedisUrl.password || undefined,
  db: Number(parsedRedisUrl.pathname.replace("/", "") || 0),
  maxRetriesPerRequest: null,
};

const supportedQueues = [
  QUEUE_NAMES.imageGeneration,
  QUEUE_NAMES.imageEnhancement,
  QUEUE_NAMES.videoGeneration,
  QUEUE_NAMES.videoEnhancement,
];

for (const queueName of supportedQueues) {
  new Worker(
    queueName,
    async (job) => {
      logger.info("job_started", { queueName, jobId: job.id, name: job.name });

      await job.updateProgress(25);
      const result = await mockProvider.generate(job.data);

      await job.updateProgress(90);
      logger.info("job_completed", {
        queueName,
        jobId: job.id,
        providerJobId: result.providerJobId,
      });

      return result;
    },
    { connection, concurrency: Number(process.env.WORKER_CONCURRENCY ?? 2) },
  );

  const events = new QueueEvents(queueName, { connection });
  events.on("failed", ({ jobId, failedReason }) => {
    logger.error("job_failed", { queueName, jobId, failedReason });
  });
}

logger.info("worker_started", { queues: supportedQueues });