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