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 });