Bot & Automation

Market AI

/root/hermes-projects/Market AI

apps/worker/src/index.ts text
import { QueueEvents, Worker } from "bullmq";
import { generateMemeSignalAnalysis, logger, memeSignalRequestSchema, scannerCriteriaSchema } from "@market-ai/core";
import { memeMarketDataRouter } from "@market-ai/provider-adapters";
import { QUEUE_NAMES } from "@market-ai/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,
};

new Worker(
  QUEUE_NAMES.memeAnalysis,
  async (job) => {
    const parsed = memeSignalRequestSchema.parse(job.data);
    const candidate = await memeMarketDataRouter.getMemeCandidate(parsed.query, parsed.chainId ?? "solana");
    const signal = generateMemeSignalAnalysis(parsed, candidate);

    logger.info("meme_analysis_job_completed", { jobId: job.id, token: candidate.token.symbol, action: signal.action });
    return signal;
  },
  { connection, concurrency: Number(process.env.WORKER_CONCURRENCY ?? 4) },
);

new Worker(
  QUEUE_NAMES.memeScan,
  async (job) => {
    const parsed = scannerCriteriaSchema.parse(job.data ?? {});
    const candidates = await memeMarketDataRouter.scanMemeCoins(parsed);

    logger.info("meme_scan_job_completed", { jobId: job.id, count: candidates.length });
    return { checkedAt: new Date().toISOString(), candidates };
  },
  { connection, concurrency: 1 },
);

new Worker(
  QUEUE_NAMES.alertScan,
  async (job) => {
    logger.info("alert_scan_started", { jobId: job.id });
    return {
      checkedAt: new Date().toISOString(),
      matchedAlerts: [],
      note: "Meme coin alert scanning is ready for database-backed alert rules.",
    };
  },
  { connection, concurrency: 1 },
);

for (const queueName of [QUEUE_NAMES.memeAnalysis, QUEUE_NAMES.memeScan, QUEUE_NAMES.alertScan]) {
  const events = new QueueEvents(queueName, { connection });
  events.on("failed", ({ jobId, failedReason }) => {
    logger.error("job_failed", { queueName, jobId, failedReason });
  });
}

logger.info("worker_started", { queues: [QUEUE_NAMES.memeAnalysis, QUEUE_NAMES.memeScan, QUEUE_NAMES.alertScan] });