Commit 70944fb6 authored by ThinhNC's avatar ThinhNC

chore: remove crawler-related files and library dependencies, set up intern...

chore: remove crawler-related files and library dependencies, set up intern management database schema
parent 30fc505b
NODE_ENV=development NODE_ENV=development
PORT=3000 PORT=3000
# Database Configuration DATABASE_URL=postgresql://postgres.[project-ref]:[password]@aws-1-[region].pooler.supabase.com:6543/postgres
DB_HOST=27.74.255.96 DIRECT_URL=postgresql://postgres:[password]@db.[project-ref].supabase.co:5432/postgres
DB_PORT=5430
DB_USER=postgres SUPABASE_URL=https://[project-ref].supabase.co
DB_PASSWORD="your_password_here" SUPABASE_PUBLISHABLE_KEY=your_supabase_publishable_key
DB_NAME=datacrawler SUPABASE_SECRET_KEY=your_supabase_secret_key
SUPABASE_JWKS_URL=https://[project-ref].supabase.co/auth/v1/.well-known/jwks.json
JWT_ACCESS_SECRET=change_me_access_secret JWT_ACCESS_SECRET=change_me_access_secret
JWT_REFRESH_SECRET=change_me_refresh_secret JWT_REFRESH_SECRET=change_me_refresh_secret
JWT_ACCESS_EXPIRES_IN=1d JWT_ACCESS_EXPIRES_IN=1d
JWT_REFRESH_EXPIRES_IN=7d JWT_REFRESH_EXPIRES_IN=7d
FIRECRAWL_API_KEY=your_firecrawl_api_key
FIRECRAWL_BASE_URL=https://api.firecrawl.dev
REDIS_HOST=127.0.0.1 REDIS_HOST=127.0.0.1
REDIS_PORT=6379 REDIS_PORT=6379
\ No newline at end of file
STORAGE_DRIVER=local
STORAGE_EXPORT_DIR=storage/exports
MAX_CRAWL_PAGES=100
MAX_CRAWL_DEPTH=3
...@@ -6,7 +6,6 @@ ...@@ -6,7 +6,6 @@
"packageManager": "pnpm@9.15.0", "packageManager": "pnpm@9.15.0",
"scripts": { "scripts": {
"dev": "ts-node-dev --respawn --transpile-only src/server.ts", "dev": "ts-node-dev --respawn --transpile-only src/server.ts",
"worker": "ts-node-dev --respawn --transpile-only src/queues/crawl.worker.ts",
"build": "tsc", "build": "tsc",
"start": "node dist/server.js", "start": "node dist/server.js",
...@@ -26,31 +25,22 @@ ...@@ -26,31 +25,22 @@
}, },
"dependencies": { "dependencies": {
"@prisma/client": "^5.22.0", "@prisma/client": "^5.22.0",
"archiver": "^7.0.1",
"bcryptjs": "^2.4.3", "bcryptjs": "^2.4.3",
"bullmq": "^5.34.0",
"cors": "^2.8.5", "cors": "^2.8.5",
"dotenv": "^16.4.7", "dotenv": "^16.4.7",
"exceljs": "^4.4.0",
"express": "^4.21.2", "express": "^4.21.2",
"@mendable/firecrawl-js": "^1.19.0",
"helmet": "^8.0.0", "helmet": "^8.0.0",
"ioredis": "^5.4.2",
"json2csv": "^6.0.0-alpha.2",
"jsonwebtoken": "^9.0.2", "jsonwebtoken": "^9.0.2",
"morgan": "^1.10.0", "morgan": "^1.10.0",
"turndown": "^7.2.0",
"zod": "^3.24.1" "zod": "^3.24.1"
}, },
"devDependencies": { "devDependencies": {
"@types/archiver": "^6.0.3",
"@types/bcryptjs": "^2.4.6", "@types/bcryptjs": "^2.4.6",
"@types/cors": "^2.8.17", "@types/cors": "^2.8.17",
"@types/express": "^4.17.21", "@types/express": "^4.17.21",
"@types/jsonwebtoken": "^9.0.7", "@types/jsonwebtoken": "^9.0.7",
"@types/morgan": "^1.9.9", "@types/morgan": "^1.9.9",
"@types/node": "^22.10.2", "@types/node": "^22.10.2",
"@types/turndown": "^5.0.5",
"eslint": "^9.17.0", "eslint": "^9.17.0",
"prettier": "^3.4.2", "prettier": "^3.4.2",
"prisma": "^5.22.0", "prisma": "^5.22.0",
...@@ -58,5 +48,8 @@ ...@@ -58,5 +48,8 @@
"ts-node-dev": "^2.0.0", "ts-node-dev": "^2.0.0",
"tsx": "^4.19.2", "tsx": "^4.19.2",
"typescript": "^5.7.2" "typescript": "^5.7.2"
},
"prisma": {
"seed": "tsx prisma/seed.ts"
} }
} }
This diff is collapsed.
-- CreateEnum
CREATE TYPE "UserRole" AS ENUM ('ADMIN', 'CRAWLER_USER', 'VIEWER');
-- CreateEnum
CREATE TYPE "CrawlJobStatus" AS ENUM ('PENDING', 'RUNNING', 'COMPLETED', 'PARTIAL_COMPLETED', 'FAILED', 'CANCELED', 'BLOCKED');
-- CreateEnum
CREATE TYPE "CrawlMode" AS ENUM ('SCRAPE', 'CRAWL', 'SITEMAP', 'URL_LIST');
-- CreateEnum
CREATE TYPE "CrawlPageStatus" AS ENUM ('PENDING', 'SUCCESS', 'FAILED', 'BLOCKED', 'REQUIRES_LOGIN', 'CAPTCHA_DETECTED', 'PAYWALL_DETECTED', 'TIMEOUT', 'SKIPPED');
-- CreateEnum
CREATE TYPE "ExportType" AS ENUM ('JSON', 'CSV', 'XLSX', 'MARKDOWN', 'MARKDOWN_ZIP', 'FULL_ZIP');
-- CreateEnum
CREATE TYPE "ExportStatus" AS ENUM ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED');
-- CreateEnum
CREATE TYPE "AssetType" AS ENUM ('IMAGE', 'LINK', 'PDF', 'FILE', 'VIDEO', 'OTHER');
-- CreateTable
CREATE TABLE "users" (
"id" UUID NOT NULL,
"email" TEXT NOT NULL,
"password_hash" TEXT NOT NULL,
"full_name" TEXT,
"role" "UserRole" NOT NULL DEFAULT 'CRAWLER_USER',
"is_active" BOOLEAN NOT NULL DEFAULT true,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "users_pkey" PRIMARY KEY ("id")
);
-- CreateTable
CREATE TABLE "crawl_jobs" (
"id" UUID NOT NULL,
"user_id" UUID NOT NULL,
"start_url" TEXT NOT NULL,
"domain" TEXT,
"mode" "CrawlMode" NOT NULL DEFAULT 'SCRAPE',
"status" "CrawlJobStatus" NOT NULL DEFAULT 'PENDING',
"max_pages" INTEGER NOT NULL DEFAULT 20,
"max_depth" INTEGER NOT NULL DEFAULT 1,
"total_pages" INTEGER NOT NULL DEFAULT 0,
"success_pages" INTEGER NOT NULL DEFAULT 0,
"failed_pages" INTEGER NOT NULL DEFAULT 0,
"error_message" TEXT,
"started_at" TIMESTAMP(3),
"finished_at" TIMESTAMP(3),
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "crawl_jobs_pkey" PRIMARY KEY ("id")
);
-- CreateTable
CREATE TABLE "crawl_pages" (
"id" UUID NOT NULL,
"job_id" UUID NOT NULL,
"url" TEXT NOT NULL,
"title" TEXT,
"description" TEXT,
"markdown_content" TEXT,
"html_content_path" TEXT,
"status" "CrawlPageStatus" NOT NULL DEFAULT 'PENDING',
"status_code" INTEGER,
"error_message" TEXT,
"crawled_at" TIMESTAMP(3),
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "crawl_pages_pkey" PRIMARY KEY ("id")
);
-- CreateTable
CREATE TABLE "crawl_assets" (
"id" UUID NOT NULL,
"job_id" UUID NOT NULL,
"page_id" UUID,
"asset_type" "AssetType" NOT NULL,
"url" TEXT NOT NULL,
"source_url" TEXT,
"alt_text" TEXT,
"mime_type" TEXT,
"order_index" INTEGER,
"css_selector" TEXT,
"dom_path" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "crawl_assets_pkey" PRIMARY KEY ("id")
);
-- CreateTable
CREATE TABLE "crawl_exports" (
"id" UUID NOT NULL,
"job_id" UUID NOT NULL,
"export_type" "ExportType" NOT NULL,
"status" "ExportStatus" NOT NULL DEFAULT 'PENDING',
"file_name" TEXT NOT NULL,
"file_path" TEXT NOT NULL,
"file_size" INTEGER,
"mime_type" TEXT,
"error_message" TEXT,
"created_at" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updated_at" TIMESTAMP(3) NOT NULL,
CONSTRAINT "crawl_exports_pkey" PRIMARY KEY ("id")
);
-- CreateIndex
CREATE UNIQUE INDEX "users_email_key" ON "users"("email");
-- CreateIndex
CREATE INDEX "crawl_jobs_user_id_idx" ON "crawl_jobs"("user_id");
-- CreateIndex
CREATE INDEX "crawl_jobs_status_idx" ON "crawl_jobs"("status");
-- CreateIndex
CREATE INDEX "crawl_jobs_created_at_idx" ON "crawl_jobs"("created_at");
-- CreateIndex
CREATE INDEX "crawl_pages_job_id_idx" ON "crawl_pages"("job_id");
-- CreateIndex
CREATE INDEX "crawl_pages_status_idx" ON "crawl_pages"("status");
-- CreateIndex
CREATE INDEX "crawl_assets_job_id_idx" ON "crawl_assets"("job_id");
-- CreateIndex
CREATE INDEX "crawl_assets_page_id_idx" ON "crawl_assets"("page_id");
-- CreateIndex
CREATE INDEX "crawl_assets_asset_type_idx" ON "crawl_assets"("asset_type");
-- CreateIndex
CREATE INDEX "crawl_exports_job_id_idx" ON "crawl_exports"("job_id");
-- CreateIndex
CREATE INDEX "crawl_exports_export_type_idx" ON "crawl_exports"("export_type");
-- AddForeignKey
ALTER TABLE "crawl_jobs" ADD CONSTRAINT "crawl_jobs_user_id_fkey" FOREIGN KEY ("user_id") REFERENCES "users"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "crawl_pages" ADD CONSTRAINT "crawl_pages_job_id_fkey" FOREIGN KEY ("job_id") REFERENCES "crawl_jobs"("id") ON DELETE CASCADE ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "crawl_assets" ADD CONSTRAINT "crawl_assets_job_id_fkey" FOREIGN KEY ("job_id") REFERENCES "crawl_jobs"("id") ON DELETE CASCADE ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "crawl_assets" ADD CONSTRAINT "crawl_assets_page_id_fkey" FOREIGN KEY ("page_id") REFERENCES "crawl_pages"("id") ON DELETE SET NULL ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "crawl_exports" ADD CONSTRAINT "crawl_exports_job_id_fkey" FOREIGN KEY ("job_id") REFERENCES "crawl_jobs"("id") ON DELETE CASCADE ON UPDATE CASCADE;
# Please do not edit this file manually
# It should be added in your version-control system (i.e. Git)
provider = "postgresql"
\ No newline at end of file
This diff is collapsed.
import { PrismaClient, UserRole } from '@prisma/client'; import { PrismaClient } from '@prisma/client';
import bcrypt from 'bcryptjs'; import bcrypt from 'bcryptjs';
const prisma = new PrismaClient(); const prisma = new PrismaClient();
async function main() { async function main() {
const passwordHash = await bcrypt.hash('Admin@123456', 10); // 1. Seed Roles
const roles = ['ADMIN', 'LEADER', 'INTERN'];
const roleMap: Record<string, string> = {};
await prisma.user.upsert({ for (const roleName of roles) {
where: { email: 'admin@crawl.local' }, const role = await prisma.role.upsert({
update: {}, where: { name: roleName },
update: {},
create: { name: roleName },
});
roleMap[roleName] = role.id;
console.log(`Role ${roleName} upserted with ID ${role.id}`);
}
// 2. Seed default admin User
const adminPassword = await bcrypt.hash('Admin@123456', 10);
const adminEmail = 'admin@nexcampus.local';
const adminUser = await prisma.user.upsert({
where: { email: adminEmail },
update: {
password: adminPassword,
roleId: roleMap['ADMIN'],
isActive: true,
},
create: { create: {
email: 'admin@crawl.local', email: adminEmail,
passwordHash, password: adminPassword,
fullName: 'System Admin', roleId: roleMap['ADMIN'],
role: UserRole.ADMIN,
isActive: true, isActive: true,
}, },
}); });
console.log('Seed completed'); console.log(`Default admin user seeded: ${adminUser.email}`);
console.log('Seed completed successfully');
} }
main() main()
......
require('dotenv').config(); require('dotenv').config();
const { spawn } = require('child_process'); const { spawn } = require('child_process');
const password = encodeURIComponent(process.env.DB_PASSWORD || ''); if (!process.env.DATABASE_URL) {
const user = encodeURIComponent(process.env.DB_USER || 'postgres'); const password = encodeURIComponent(process.env.DB_PASSWORD || '');
process.env.DATABASE_URL = `postgresql://${user}:${password}@${process.env.DB_HOST || 'localhost'}:${process.env.DB_PORT || '5432'}/${process.env.DB_NAME || 'datacrawler'}?schema=public`; const user = encodeURIComponent(process.env.DB_USER || 'postgres');
process.env.DATABASE_URL = `postgresql://${user}:${password}@${process.env.DB_HOST || 'localhost'}:${process.env.DB_PORT || '5432'}/${process.env.DB_NAME || 'datacrawler'}?schema=public`;
}
const args = process.argv.slice(2); const args = process.argv.slice(2);
const cmd = process.platform === 'win32' ? 'npx.cmd' : 'npx'; const cmd = process.platform === 'win32' ? 'npx.cmd' : 'npx';
......
import { UserRole } from '@prisma/client';
declare global { declare global {
namespace Express { namespace Express {
interface Request { interface Request {
user: { user: {
id: string; id: string;
email: string; email: string;
role: UserRole; role: string;
}; };
} }
} }
} }
export {};
import { envConfig } from './env.config';
export const databaseConfig = {
url: envConfig.databaseUrl,
host: envConfig.database.host,
port: envConfig.database.port,
user: envConfig.database.user,
password: envConfig.database.password,
name: envConfig.database.name,
};
export const envConfig = { export const envConfig = {
nodeEnv: process.env.NODE_ENV || 'development', nodeEnv: process.env.NODE_ENV || 'development',
port: parseInt(process.env.PORT || '3000', 10), port: parseInt(process.env.PORT || '3000', 10),
database: {
host: process.env.DB_HOST || 'localhost',
port: parseInt(process.env.DB_PORT || '5432', 10),
user: process.env.DB_USER || 'postgres',
password: process.env.DB_PASSWORD || '',
name: process.env.DB_NAME || 'datacrawler',
},
get databaseUrl() {
return `postgresql://${this.database.user}:${encodeURIComponent(this.database.password)}@${this.database.host}:${this.database.port}/${this.database.name}?schema=public`;
},
jwt: { jwt: {
accessSecret: process.env.JWT_ACCESS_SECRET || 'default_access_secret', accessSecret: process.env.JWT_ACCESS_SECRET || 'default_access_secret',
refreshSecret: process.env.JWT_REFRESH_SECRET || 'default_refresh_secret', refreshSecret: process.env.JWT_REFRESH_SECRET || 'default_refresh_secret',
accessExpiresIn: process.env.JWT_ACCESS_EXPIRES_IN || '1d', accessExpiresIn: process.env.JWT_ACCESS_EXPIRES_IN || '1d',
refreshExpiresIn: process.env.JWT_REFRESH_EXPIRES_IN || '7d', refreshExpiresIn: process.env.JWT_REFRESH_EXPIRES_IN || '7d',
}, },
firecrawl: {
apiKey: process.env.FIRECRAWL_API_KEY || '',
baseUrl: process.env.FIRECRAWL_BASE_URL || 'https://api.firecrawl.dev',
},
redis: {
host: process.env.REDIS_HOST || '127.0.0.1',
port: parseInt(process.env.REDIS_PORT || '6379', 10),
},
storage: {
driver: process.env.STORAGE_DRIVER || 'local',
exportDir: process.env.STORAGE_EXPORT_DIR || 'storage/exports',
},
crawl: {
maxPages: parseInt(process.env.MAX_CRAWL_PAGES || '100', 10),
maxDepth: parseInt(process.env.MAX_CRAWL_DEPTH || '3', 10),
},
}; };
import { envConfig } from './env.config';
export const firecrawlConfig = envConfig.firecrawl;
import { envConfig } from './env.config';
export const storageConfig = envConfig.storage;
import { Request, Response, NextFunction } from 'express'; import { Request, Response, NextFunction } from 'express';
import { UserRole } from '@prisma/client';
import { AppError } from '../common/errors/app-error'; import { AppError } from '../common/errors/app-error';
import { ERROR_CODE } from '../common/errors/error-code'; import { ERROR_CODE } from '../common/errors/error-code';
export function requireRole(...roles: UserRole[]) { export function requireRole(...roles: string[]) {
return (req: Request, res: Response, next: NextFunction): void => { return (req: Request, res: Response, next: NextFunction): void => {
if (!req.user) { if (!req.user) {
next(new AppError('Unauthorized', 401, ERROR_CODE.UNAUTHORIZED)); next(new AppError('Unauthorized', 401, ERROR_CODE.UNAUTHORIZED));
......
...@@ -4,12 +4,14 @@ export class AuthRepository { ...@@ -4,12 +4,14 @@ export class AuthRepository {
findByEmail(email: string) { findByEmail(email: string) {
return prisma.user.findUnique({ return prisma.user.findUnique({
where: { email }, where: { email },
include: { role: true },
}); });
} }
findById(id: string) { findById(id: string) {
return prisma.user.findUnique({ return prisma.user.findUnique({
where: { id }, where: { id },
include: { role: true },
}); });
} }
} }
...@@ -19,13 +19,13 @@ export class AuthService { ...@@ -19,13 +19,13 @@ export class AuthService {
throw new AppError('Account is inactive', 403, ERROR_CODE.USER_INACTIVE); throw new AppError('Account is inactive', 403, ERROR_CODE.USER_INACTIVE);
} }
const isPasswordValid = await bcrypt.compare(password, user.passwordHash); const isPasswordValid = await bcrypt.compare(password, user.password);
if (!isPasswordValid) { if (!isPasswordValid) {
throw new AppError('Invalid credentials', 401, ERROR_CODE.INVALID_CREDENTIALS); throw new AppError('Invalid credentials', 401, ERROR_CODE.INVALID_CREDENTIALS);
} }
const payload = { id: user.id, email: user.email, role: user.role }; const payload = { id: user.id, email: user.email, role: user.role.name };
const accessToken = jwt.sign(payload, jwtConfig.accessSecret, { const accessToken = jwt.sign(payload, jwtConfig.accessSecret, {
expiresIn: jwtConfig.accessExpiresIn as any, expiresIn: jwtConfig.accessExpiresIn as any,
...@@ -41,8 +41,7 @@ export class AuthService { ...@@ -41,8 +41,7 @@ export class AuthService {
user: { user: {
id: user.id, id: user.id,
email: user.email, email: user.email,
fullName: user.fullName, role: user.role.name,
role: user.role,
}, },
}; };
} }
...@@ -57,8 +56,7 @@ export class AuthService { ...@@ -57,8 +56,7 @@ export class AuthService {
return { return {
id: user.id, id: user.id,
email: user.email, email: user.email,
fullName: user.fullName, role: user.role.name,
role: user.role,
isActive: user.isActive, isActive: user.isActive,
createdAt: user.createdAt, createdAt: user.createdAt,
}; };
......
import { AssetType } from '@prisma/client';
export interface CreateCrawlAssetDto {
jobId: string;
pageId?: string;
assetType: AssetType;
url: string;
sourceUrl?: string;
altText?: string;
mimeType?: string;
orderIndex?: number;
cssSelector?: string;
domPath?: string;
}
import { prisma } from '../../database/prisma.client';
import { AssetType } from '@prisma/client';
export class CrawlAssetRepository {
create(data: {
jobId: string;
pageId?: string;
assetType: AssetType;
url: string;
sourceUrl?: string;
altText?: string;
mimeType?: string;
orderIndex?: number;
cssSelector?: string;
domPath?: string;
}) {
return prisma.crawlAsset.create({ data });
}
createMany(assets: Array<{
jobId: string;
pageId?: string;
assetType: AssetType;
url: string;
sourceUrl?: string;
altText?: string;
mimeType?: string;
}>) {
return prisma.crawlAsset.createMany({ data: assets });
}
findByJobId(jobId: string) {
return prisma.crawlAsset.findMany({
where: { jobId },
orderBy: { createdAt: 'asc' },
});
}
findByPageId(pageId: string) {
return prisma.crawlAsset.findMany({
where: { pageId },
});
}
}
import { CrawlAssetRepository } from './crawl-asset.repository';
import { AssetType } from '@prisma/client';
export class CrawlAssetService {
private readonly repository = new CrawlAssetRepository();
async findByJobId(jobId: string) {
return this.repository.findByJobId(jobId);
}
async create(data: {
jobId: string;
pageId?: string;
assetType: AssetType;
url: string;
sourceUrl?: string;
altText?: string;
mimeType?: string;
}) {
return this.repository.create(data);
}
async createMany(assets: Array<{
jobId: string;
pageId?: string;
assetType: AssetType;
url: string;
sourceUrl?: string;
altText?: string;
mimeType?: string;
}>) {
return this.repository.createMany(assets);
}
}
import { Request, Response, NextFunction } from 'express';
import fs from 'fs';
import { CrawlExportService } from './crawl-export.service';
import { AppError } from '../../common/errors/app-error';
export class CrawlExportController {
private readonly service = new CrawlExportService();
download = async (req: Request, res: Response, next: NextFunction) => {
try {
const exportRecord = await this.service.findById(req.params.exportId);
if (!fs.existsSync(exportRecord.filePath)) {
next(new AppError('Export file not found on disk', 404));
return;
}
res.download(exportRecord.filePath, exportRecord.fileName);
} catch (error) {
next(error);
}
};
}
import { ExportType, ExportStatus } from '@prisma/client';
export interface CreateCrawlExportDto {
jobId: string;
exportType: ExportType;
fileName: string;
filePath: string;
fileSize?: number;
mimeType?: string;
}
export interface UpdateCrawlExportDto {
status?: ExportStatus;
fileSize?: number;
errorMessage?: string;
}
import { prisma } from '../../database/prisma.client';
import { ExportStatus, ExportType } from '@prisma/client';
export class CrawlExportRepository {
create(data: {
jobId: string;
exportType: ExportType;
fileName: string;
filePath: string;
fileSize?: number;
mimeType?: string;
}) {
return prisma.crawlExport.create({ data });
}
findByJobId(jobId: string) {
return prisma.crawlExport.findMany({
where: { jobId },
orderBy: { createdAt: 'desc' },
});
}
findById(id: string) {
return prisma.crawlExport.findUnique({
where: { id },
});
}
update(id: string, data: { status?: ExportStatus; fileSize?: number; errorMessage?: string }) {
return prisma.crawlExport.update({
where: { id },
data,
});
}
}
import { Router } from 'express';
import { CrawlExportController } from './crawl-export.controller';
import { authMiddleware } from '../../middlewares/auth.middleware';
const router = Router();
const controller = new CrawlExportController();
router.get('/:exportId/download', authMiddleware, controller.download);
export default router;
import { CrawlExportRepository } from './crawl-export.repository';
import { CrawlJobRepository } from '../crawl-jobs/crawl-job.repository';
import { AppError } from '../../common/errors/app-error';
import { ERROR_CODE } from '../../common/errors/error-code';
import { ExportType } from '@prisma/client';
export class CrawlExportService {
private readonly repository = new CrawlExportRepository();
private readonly jobRepository = new CrawlJobRepository();
async findByJobId(jobId: string) {
return this.repository.findByJobId(jobId);
}
async findById(id: string) {
const exportRecord = await this.repository.findById(id);
if (!exportRecord) {
throw new AppError('Export not found', 404, ERROR_CODE.NOT_FOUND);
}
return exportRecord;
}
async createExport(userId: string, jobId: string, exportType: ExportType) {
const job = await this.jobRepository.findById(jobId);
if (!job || job.userId !== userId) {
throw new AppError('Crawl job not found', 404, ERROR_CODE.CRAWL_JOB_NOT_FOUND);
}
const { ExportService } = await import('../exports/export.service');
const exportService = new ExportService();
return exportService.generate(job, exportType);
}
}
import { Request, Response, NextFunction } from 'express';
import { CrawlJobService } from './crawl-job.service';
import { CrawlPageService } from '../crawl-pages/crawl-page.service';
import { CrawlExportService } from '../crawl-exports/crawl-export.service';
export class CrawlJobController {
private readonly service = new CrawlJobService();
private readonly pageService = new CrawlPageService();
private readonly exportService = new CrawlExportService();
create = async (req: Request, res: Response, next: NextFunction) => {
try {
const userId = req.user.id;
const result = await this.service.create(userId, req.body);
res.status(201).json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
findAll = async (req: Request, res: Response, next: NextFunction) => {
try {
const userId = req.user.id;
const result = await this.service.findAllByUser(userId, req.query);
res.json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
findById = async (req: Request, res: Response, next: NextFunction) => {
try {
const userId = req.user.id;
const result = await this.service.findById(userId, req.params.id);
res.json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
cancel = async (req: Request, res: Response, next: NextFunction) => {
try {
const userId = req.user.id;
const result = await this.service.cancel(userId, req.params.id);
res.json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
getPages = async (req: Request, res: Response, next: NextFunction) => {
try {
await this.service.findById(req.user.id, req.params.id);
const result = await this.pageService.findByJobId(req.params.id);
res.json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
getExports = async (req: Request, res: Response, next: NextFunction) => {
try {
await this.service.findById(req.user.id, req.params.id);
const result = await this.exportService.findByJobId(req.params.id);
res.json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
createExport = async (req: Request, res: Response, next: NextFunction) => {
try {
const result = await this.exportService.createExport(
req.user.id,
req.params.id,
req.body.exportType,
);
res.status(201).json({
success: true,
data: result,
});
} catch (error) {
next(error);
}
};
}
import { CrawlJobStatus, CrawlMode } from '@prisma/client';
export interface CreateCrawlJobDto {
startUrl: string;
mode?: CrawlMode;
maxPages?: number;
maxDepth?: number;
}
export interface CrawlJobQueryDto {
status?: CrawlJobStatus;
page?: number;
limit?: number;
}
export interface UpdateCrawlJobStatusDto {
status: CrawlJobStatus;
errorMessage?: string;
startedAt?: Date;
finishedAt?: Date;
}
import { prisma } from '../../database/prisma.client';
import { CrawlJobStatus } from '@prisma/client';
export class CrawlJobRepository {
create(data: {
userId: string;
startUrl: string;
domain?: string;
mode: any;
maxPages?: number;
maxDepth?: number;
}) {
return prisma.crawlJob.create({
data: {
userId: data.userId,
startUrl: data.startUrl,
domain: data.domain,
mode: data.mode,
maxPages: data.maxPages ?? 20,
maxDepth: data.maxDepth ?? 1,
},
});
}
findAllByUser(userId: string, query: any) {
return prisma.crawlJob.findMany({
where: { userId },
orderBy: { createdAt: 'desc' },
include: {
exports: true,
},
});
}
findAll() {
return prisma.crawlJob.findMany({
orderBy: { createdAt: 'desc' },
});
}
findById(id: string) {
return prisma.crawlJob.findUnique({
where: { id },
include: {
pages: true,
exports: true,
assets: true,
},
});
}
updateStatus(id: string, status: CrawlJobStatus, extra?: {
errorMessage?: string;
startedAt?: Date;
finishedAt?: Date;
totalPages?: number;
successPages?: number;
failedPages?: number;
}) {
return prisma.crawlJob.update({
where: { id },
data: { status, ...extra },
});
}
}
import { Router } from 'express';
import { CrawlJobController } from './crawl-job.controller';
import { authMiddleware } from '../../middlewares/auth.middleware';
import { validate } from '../../middlewares/validate.middleware';
import { createCrawlJobSchema } from './crawl-job.validation';
const router = Router();
const controller = new CrawlJobController();
router.post('/', authMiddleware, validate(createCrawlJobSchema), controller.create);
router.get('/', authMiddleware, controller.findAll);
router.get('/:id', authMiddleware, controller.findById);
router.post('/:id/cancel', authMiddleware, controller.cancel);
router.get('/:id/pages', authMiddleware, controller.getPages);
router.get('/:id/exports', authMiddleware, controller.getExports);
router.post('/:id/exports', authMiddleware, controller.createExport);
export default router;
import { CrawlJobRepository } from './crawl-job.repository';
import { AppError } from '../../common/errors/app-error';
import { ERROR_CODE } from '../../common/errors/error-code';
import { validateUrl, extractDomain } from '../../common/helpers/url.helper';
import { crawlQueue } from '../../queues/crawl.queue';
export class CrawlJobService {
private readonly repository = new CrawlJobRepository();
async create(userId: string, payload: any) {
const parsed = validateUrl(payload.startUrl);
const domain = extractDomain(payload.startUrl);
const job = await this.repository.create({
userId,
startUrl: parsed.href,
domain,
mode: payload.mode ?? 'SCRAPE',
maxPages: payload.maxPages,
maxDepth: payload.maxDepth,
});
await crawlQueue.add('crawl-job', {
jobId: job.id,
});
return job;
}
async findAllByUser(userId: string, query: any) {
return this.repository.findAllByUser(userId, query);
}
async findById(userId: string, jobId: string) {
const job = await this.repository.findById(jobId);
if (!job || job.userId !== userId) {
throw new AppError('Crawl job not found', 404, ERROR_CODE.CRAWL_JOB_NOT_FOUND);
}
return job;
}
async cancel(userId: string, jobId: string) {
const job = await this.findById(userId, jobId);
if (job.status === 'COMPLETED') {
throw new AppError('Completed job cannot be canceled', 400, ERROR_CODE.CRAWL_JOB_ALREADY_COMPLETED);
}
return this.repository.updateStatus(jobId, 'CANCELED');
}
}
import { z } from 'zod';
export const createCrawlJobSchema = z.object({
startUrl: z.string().url('Invalid URL format'),
mode: z.enum(['SCRAPE', 'CRAWL', 'SITEMAP', 'URL_LIST']).optional(),
maxPages: z.number().int().min(1).max(1000).optional(),
maxDepth: z.number().int().min(1).max(10).optional(),
});
import { CrawlPageStatus } from '@prisma/client';
export interface CreateCrawlPageDto {
jobId: string;
url: string;
title?: string;
description?: string;
markdownContent?: string;
htmlContentPath?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}
export interface UpdateCrawlPageDto {
title?: string;
description?: string;
markdownContent?: string;
htmlContentPath?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}
import { prisma } from '../../database/prisma.client';
import { CrawlPageStatus } from '@prisma/client';
export class CrawlPageRepository {
create(data: {
jobId: string;
url: string;
title?: string;
description?: string;
markdownContent?: string;
htmlContentPath?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}) {
return prisma.crawlPage.create({ data });
}
findByJobId(jobId: string) {
return prisma.crawlPage.findMany({
where: { jobId },
orderBy: { createdAt: 'asc' },
});
}
findById(id: string) {
return prisma.crawlPage.findUnique({
where: { id },
include: { assets: true },
});
}
update(id: string, data: {
title?: string;
description?: string;
markdownContent?: string;
htmlContentPath?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}) {
return prisma.crawlPage.update({
where: { id },
data,
});
}
countByJobId(jobId: string) {
return prisma.crawlPage.count({ where: { jobId } });
}
countByJobIdAndStatus(jobId: string, status: CrawlPageStatus) {
return prisma.crawlPage.count({ where: { jobId, status } });
}
}
import { CrawlPageRepository } from './crawl-page.repository';
import { CrawlPageStatus } from '@prisma/client';
export class CrawlPageService {
private readonly repository = new CrawlPageRepository();
async findByJobId(jobId: string) {
return this.repository.findByJobId(jobId);
}
async create(data: {
jobId: string;
url: string;
title?: string;
description?: string;
markdownContent?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}) {
return this.repository.create(data);
}
async update(id: string, data: {
title?: string;
description?: string;
markdownContent?: string;
status?: CrawlPageStatus;
statusCode?: number;
errorMessage?: string;
crawledAt?: Date;
}) {
return this.repository.update(id, data);
}
}
import fs from 'fs';
import path from 'path';
import { CrawlJob, CrawlPage } from '@prisma/client';
import { ensureDirExists } from '../../common/helpers/file.helper';
import { storageConfig } from '../../config/storage.config';
export class CsvExportService {
async export(job: CrawlJob & { pages: CrawlPage[] }): Promise<{ fileName: string; filePath: string }> {
const exportDir = storageConfig.exportDir;
ensureDirExists(exportDir);
const fileName = `crawl-${job.id}-${Date.now()}.csv`;
const filePath = path.join(exportDir, fileName);
const headers = ['url', 'title', 'description', 'status', 'statusCode', 'crawledAt'];
const rows = job.pages.map((page) => [
page.url,
this.escapeCsv(page.title ?? ''),
this.escapeCsv(page.description ?? ''),
page.status,
page.statusCode ?? '',
page.crawledAt?.toISOString() ?? '',
]);
const csv = [headers, ...rows].map((r) => r.join(',')).join('\n');
fs.writeFileSync(filePath, csv, 'utf-8');
return { fileName, filePath };
}
private escapeCsv(value: string): string {
if (value.includes(',') || value.includes('"') || value.includes('\n')) {
return `"${value.replace(/"/g, '""')}"`;
}
return value;
}
}
import { CrawlJob, CrawlPage, ExportType } from '@prisma/client';
import { prisma } from '../../database/prisma.client';
import { CrawlExportRepository } from '../crawl-exports/crawl-export.repository';
import { JsonExportService } from './json-export.service';
import { CsvExportService } from './csv-export.service';
import { XlsxExportService } from './xlsx-export.service';
import { MarkdownExportService } from './markdown-export.service';
import { getFileSizeBytes } from '../../common/helpers/file.helper';
const MIME_TYPE: Record<string, string> = {
JSON: 'application/json',
CSV: 'text/csv',
XLSX: 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
MARKDOWN: 'text/markdown',
MARKDOWN_ZIP: 'application/zip',
FULL_ZIP: 'application/zip',
};
export class ExportService {
private readonly exportRepository = new CrawlExportRepository();
async generate(job: CrawlJob, exportType: ExportType) {
const jobWithPages = await prisma.crawlJob.findUnique({
where: { id: job.id },
include: { pages: true },
});
const fullJob = jobWithPages as CrawlJob & { pages: CrawlPage[] };
let fileName: string;
let filePath: string;
switch (exportType) {
case 'JSON': {
const service = new JsonExportService();
({ fileName, filePath } = await service.export(fullJob));
break;
}
case 'CSV': {
const service = new CsvExportService();
({ fileName, filePath } = await service.export(fullJob));
break;
}
case 'XLSX': {
const service = new XlsxExportService();
({ fileName, filePath } = await service.export(fullJob));
break;
}
case 'MARKDOWN': {
const service = new MarkdownExportService();
({ fileName, filePath } = await service.exportSingleMd(fullJob));
break;
}
case 'MARKDOWN_ZIP':
case 'FULL_ZIP': {
const service = new MarkdownExportService();
({ fileName, filePath } = await service.exportMarkdownZip(fullJob));
break;
}
default:
throw new Error(`Unsupported export type: ${exportType}`);
}
const fileSize = getFileSizeBytes(filePath);
const exportRecord = await this.exportRepository.create({
jobId: job.id,
exportType,
fileName,
filePath,
fileSize,
mimeType: MIME_TYPE[exportType],
});
await this.exportRepository.update(exportRecord.id, { status: 'COMPLETED' });
return exportRecord;
}
}
import fs from 'fs';
import path from 'path';
import { CrawlJob, CrawlPage } from '@prisma/client';
import { ensureDirExists } from '../../common/helpers/file.helper';
import { storageConfig } from '../../config/storage.config';
export class JsonExportService {
async export(job: CrawlJob & { pages: CrawlPage[] }): Promise<{ fileName: string; filePath: string }> {
const exportDir = storageConfig.exportDir;
ensureDirExists(exportDir);
const fileName = `crawl-${job.id}-${Date.now()}.json`;
const filePath = path.join(exportDir, fileName);
const data = {
jobId: job.id,
startUrl: job.startUrl,
domain: job.domain,
mode: job.mode,
status: job.status,
totalPages: job.totalPages,
exportedAt: new Date().toISOString(),
pages: job.pages.map((page) => ({
url: page.url,
title: page.title,
description: page.description,
markdownContent: page.markdownContent,
status: page.status,
statusCode: page.statusCode,
crawledAt: page.crawledAt,
})),
};
fs.writeFileSync(filePath, JSON.stringify(data, null, 2), 'utf-8');
return { fileName, filePath };
}
}
import fs from 'fs';
import path from 'path';
import archiver from 'archiver';
import { CrawlJob, CrawlPage } from '@prisma/client';
import { ensureDirExists } from '../../common/helpers/file.helper';
import { storageConfig } from '../../config/storage.config';
export class MarkdownExportService {
async exportSingleMd(job: CrawlJob & { pages: CrawlPage[] }): Promise<{ fileName: string; filePath: string }> {
const exportDir = storageConfig.exportDir;
ensureDirExists(exportDir);
const fileName = `crawl-${job.id}-${Date.now()}.md`;
const filePath = path.join(exportDir, fileName);
const lines: string[] = [`# Crawl Report: ${job.startUrl}`, ''];
for (const page of job.pages) {
lines.push(`## ${page.title ?? page.url}`);
lines.push(`**URL:** ${page.url}`);
if (page.description) lines.push(`**Description:** ${page.description}`);
lines.push('');
if (page.markdownContent) {
lines.push(page.markdownContent);
lines.push('');
}
lines.push('---', '');
}
fs.writeFileSync(filePath, lines.join('\n'), 'utf-8');
return { fileName, filePath };
}
async exportMarkdownZip(job: CrawlJob & { pages: CrawlPage[] }): Promise<{ fileName: string; filePath: string }> {
const exportDir = storageConfig.exportDir;
ensureDirExists(exportDir);
const fileName = `crawl-${job.id}-${Date.now()}.zip`;
const filePath = path.join(exportDir, fileName);
await new Promise<void>((resolve, reject) => {
const output = fs.createWriteStream(filePath);
const archive = archiver('zip', { zlib: { level: 9 } });
output.on('close', resolve);
archive.on('error', reject);
archive.pipe(output);
job.pages.forEach((page, idx) => {
const mdContent = page.markdownContent ?? `# ${page.url}\n\nNo content available.`;
const mdFileName = `page-${String(idx + 1).padStart(3, '0')}.md`;
archive.append(mdContent, { name: mdFileName });
});
archive.finalize();
});
return { fileName, filePath };
}
}
import path from 'path';
import ExcelJS from 'exceljs';
import { CrawlJob, CrawlPage } from '@prisma/client';
import { ensureDirExists } from '../../common/helpers/file.helper';
import { storageConfig } from '../../config/storage.config';
export class XlsxExportService {
async export(job: CrawlJob & { pages: CrawlPage[] }): Promise<{ fileName: string; filePath: string }> {
const exportDir = storageConfig.exportDir;
ensureDirExists(exportDir);
const fileName = `crawl-${job.id}-${Date.now()}.xlsx`;
const filePath = path.join(exportDir, fileName);
const workbook = new ExcelJS.Workbook();
const sheet = workbook.addWorksheet('Pages');
sheet.columns = [
{ header: 'URL', key: 'url', width: 50 },
{ header: 'Title', key: 'title', width: 40 },
{ header: 'Description', key: 'description', width: 50 },
{ header: 'Status', key: 'status', width: 20 },
{ header: 'Status Code', key: 'statusCode', width: 15 },
{ header: 'Crawled At', key: 'crawledAt', width: 25 },
];
for (const page of job.pages) {
sheet.addRow({
url: page.url,
title: page.title ?? '',
description: page.description ?? '',
status: page.status,
statusCode: page.statusCode ?? '',
crawledAt: page.crawledAt?.toISOString() ?? '',
});
}
await workbook.xlsx.writeFile(filePath);
return { fileName, filePath };
}
}
import FirecrawlApp from '@mendable/firecrawl-js';
import { firecrawlConfig } from '../../config/firecrawl.config';
let firecrawlClient: FirecrawlApp | null = null;
export function getFirecrawlClient(): FirecrawlApp {
if (!firecrawlClient) {
firecrawlClient = new FirecrawlApp({
apiKey: firecrawlConfig.apiKey,
});
}
return firecrawlClient;
}
export interface FirecrawlScrapeDto {
url: string;
formats?: string[];
}
export interface FirecrawlCrawlDto {
url: string;
maxPages?: number;
maxDepth?: number;
}
export interface FirecrawlPageResult {
url: string;
title?: string;
description?: string;
markdown?: string;
html?: string;
statusCode?: number;
links?: string[];
images?: Array<{ url: string; alt?: string }>;
}
import { getFirecrawlClient } from './firecrawl.client';
import { FirecrawlPageResult } from './firecrawl.dto';
export class FirecrawlService {
async scrapePage(url: string): Promise<FirecrawlPageResult> {
const client = getFirecrawlClient();
const result = await client.scrapeUrl(url, {
formats: ['markdown', 'html'],
}) as any;
return {
url,
title: result.metadata?.title,
description: result.metadata?.description,
markdown: result.markdown,
html: result.html,
statusCode: result.metadata?.statusCode,
};
}
async crawlSite(url: string, maxPages: number = 20, maxDepth: number = 1) {
const client = getFirecrawlClient();
const result = await client.crawlUrl(url, {
limit: maxPages,
maxDepth,
scrapeOptions: {
formats: ['markdown', 'html'],
},
} as any);
return result;
}
}
export interface CreateUserDto { export interface CreateUserDto {
email: string; email: string;
password: string; password: string;
fullName?: string; roleId: string;
role?: string;
} }
export interface UpdateUserDto { export interface UpdateUserDto {
fullName?: string;
isActive?: boolean; isActive?: boolean;
role?: string; roleId?: string;
} }
export interface UserResponseDto { export interface UserResponseDto {
id: string; id: string;
email: string; email: string;
fullName: string | null; roleId: string;
role: string; role?: {
id: string;
name: string;
};
isActive: boolean; isActive: boolean;
createdAt: Date; createdAt: Date;
updatedAt: Date; updatedAt: Date;
......
import { prisma } from '../../database/prisma.client'; import { prisma } from '../../database/prisma.client';
import { UserRole } from '@prisma/client';
export class UserRepository { export class UserRepository {
findAll() { findAll() {
return prisma.user.findMany({ return prisma.user.findMany({
include: { role: true },
orderBy: { createdAt: 'desc' }, orderBy: { createdAt: 'desc' },
}); });
} }
...@@ -11,35 +11,37 @@ export class UserRepository { ...@@ -11,35 +11,37 @@ export class UserRepository {
findById(id: string) { findById(id: string) {
return prisma.user.findUnique({ return prisma.user.findUnique({
where: { id }, where: { id },
include: { role: true },
}); });
} }
findByEmail(email: string) { findByEmail(email: string) {
return prisma.user.findUnique({ return prisma.user.findUnique({
where: { email }, where: { email },
include: { role: true },
}); });
} }
create(data: { create(data: {
email: string; email: string;
passwordHash: string; passwordHash: string;
fullName?: string; roleId: string;
role?: UserRole;
}) { }) {
return prisma.user.create({ return prisma.user.create({
data: { data: {
email: data.email, email: data.email,
passwordHash: data.passwordHash, password: data.passwordHash,
fullName: data.fullName, roleId: data.roleId,
role: data.role ?? 'CRAWLER_USER',
}, },
include: { role: true },
}); });
} }
update(id: string, data: { fullName?: string; isActive?: boolean; role?: UserRole }) { update(id: string, data: { isActive?: boolean; roleId?: string }) {
return prisma.user.update({ return prisma.user.update({
where: { id }, where: { id },
data, data,
include: { role: true },
}); });
} }
} }
...@@ -2,7 +2,6 @@ import bcrypt from 'bcryptjs'; ...@@ -2,7 +2,6 @@ import bcrypt from 'bcryptjs';
import { UserRepository } from './user.repository'; import { UserRepository } from './user.repository';
import { AppError } from '../../common/errors/app-error'; import { AppError } from '../../common/errors/app-error';
import { ERROR_CODE } from '../../common/errors/error-code'; import { ERROR_CODE } from '../../common/errors/error-code';
import { UserRole } from '@prisma/client';
export class UserService { export class UserService {
private readonly repository = new UserRepository(); private readonly repository = new UserRepository();
...@@ -21,7 +20,7 @@ export class UserService { ...@@ -21,7 +20,7 @@ export class UserService {
return user; return user;
} }
async create(data: { email: string; password: string; fullName?: string; role?: string }) { async create(data: { email: string; password: string; roleId: string }) {
const existing = await this.repository.findByEmail(data.email); const existing = await this.repository.findByEmail(data.email);
if (existing) { if (existing) {
...@@ -33,18 +32,16 @@ export class UserService { ...@@ -33,18 +32,16 @@ export class UserService {
return this.repository.create({ return this.repository.create({
email: data.email, email: data.email,
passwordHash, passwordHash,
fullName: data.fullName, roleId: data.roleId,
role: data.role as UserRole | undefined,
}); });
} }
async update(id: string, data: { fullName?: string; isActive?: boolean; role?: string }) { async update(id: string, data: { isActive?: boolean; roleId?: string }) {
await this.findById(id); await this.findById(id);
return this.repository.update(id, { return this.repository.update(id, {
fullName: data.fullName,
isActive: data.isActive, isActive: data.isActive,
role: data.role as UserRole | undefined, roleId: data.roleId,
}); });
} }
} }
...@@ -3,12 +3,10 @@ import { z } from 'zod'; ...@@ -3,12 +3,10 @@ import { z } from 'zod';
export const createUserSchema = z.object({ export const createUserSchema = z.object({
email: z.string().email('Invalid email format'), email: z.string().email('Invalid email format'),
password: z.string().min(8, 'Password must be at least 8 characters'), password: z.string().min(8, 'Password must be at least 8 characters'),
fullName: z.string().optional(), roleId: z.string().uuid('Invalid roleId format'),
role: z.enum(['ADMIN', 'CRAWLER_USER', 'VIEWER']).optional(),
}); });
export const updateUserSchema = z.object({ export const updateUserSchema = z.object({
fullName: z.string().optional(),
isActive: z.boolean().optional(), isActive: z.boolean().optional(),
role: z.enum(['ADMIN', 'CRAWLER_USER', 'VIEWER']).optional(), roleId: z.string().uuid('Invalid roleId format').optional(),
}); });
import { Queue } from 'bullmq';
import { envConfig } from '../config/env.config';
export const crawlQueue = new Queue('crawl-jobs', {
connection: {
host: envConfig.redis.host,
port: envConfig.redis.port,
},
defaultJobOptions: {
attempts: 3,
backoff: {
type: 'exponential',
delay: 5000,
},
},
});
import 'dotenv/config';
import { Worker, Job } from 'bullmq';
import { envConfig } from '../config/env.config';
import { CrawlJobRepository } from '../modules/crawl-jobs/crawl-job.repository';
import { CrawlPageRepository } from '../modules/crawl-pages/crawl-page.repository';
import { CrawlAssetRepository } from '../modules/crawl-assets/crawl-asset.repository';
import { FirecrawlService } from '../modules/firecrawl/firecrawl.service';
const jobRepository = new CrawlJobRepository();
const pageRepository = new CrawlPageRepository();
const assetRepository = new CrawlAssetRepository();
const firecrawlService = new FirecrawlService();
async function processCrawlJob(job: Job<{ jobId: string }>) {
const { jobId } = job.data;
await jobRepository.updateStatus(jobId, 'RUNNING', {
startedAt: new Date(),
});
try {
const crawlJob = await jobRepository.findById(jobId);
if (!crawlJob) {
throw new Error(`Job ${jobId} not found`);
}
if (crawlJob.mode === 'SCRAPE') {
const result = await firecrawlService.scrapePage(crawlJob.startUrl);
const page = await pageRepository.create({
jobId,
url: result.url,
title: result.title,
description: result.description,
markdownContent: result.markdown,
status: 'SUCCESS',
statusCode: result.statusCode,
crawledAt: new Date(),
});
if (result.images && result.images.length > 0) {
await assetRepository.createMany(
result.images.map((img) => ({
jobId,
pageId: page.id,
assetType: 'IMAGE' as const,
url: img.url,
altText: img.alt,
})),
);
}
await jobRepository.updateStatus(jobId, 'COMPLETED', {
finishedAt: new Date(),
totalPages: 1,
successPages: 1,
failedPages: 0,
});
} else {
const result = await firecrawlService.crawlSite(
crawlJob.startUrl,
crawlJob.maxPages,
crawlJob.maxDepth,
) as any;
const pages = result?.data ?? [];
let successCount = 0;
let failedCount = 0;
for (const item of pages) {
try {
await pageRepository.create({
jobId,
url: item.metadata?.url ?? item.url ?? crawlJob.startUrl,
title: item.metadata?.title,
description: item.metadata?.description,
markdownContent: item.markdown,
status: 'SUCCESS',
statusCode: item.metadata?.statusCode,
crawledAt: new Date(),
});
successCount++;
} catch {
failedCount++;
}
}
await jobRepository.updateStatus(jobId, 'COMPLETED', {
finishedAt: new Date(),
totalPages: pages.length,
successPages: successCount,
failedPages: failedCount,
});
}
} catch (error: any) {
await jobRepository.updateStatus(jobId, 'FAILED', {
finishedAt: new Date(),
errorMessage: error?.message ?? 'Unknown error',
});
throw error;
}
}
const worker = new Worker('crawl-jobs', processCrawlJob, {
connection: {
host: envConfig.redis.host,
port: envConfig.redis.port,
},
concurrency: 3,
});
worker.on('completed', (job) => {
console.log(`[Worker] Job ${job.id} completed`);
});
worker.on('failed', (job, err) => {
console.error(`[Worker] Job ${job?.id} failed: ${err.message}`);
});
console.log('[Worker] Crawl worker started');
import { Router } from 'express'; import { Router } from 'express';
import authRoute from '../modules/auth/auth.route'; import authRoute from '../modules/auth/auth.route';
import userRoute from '../modules/users/user.route'; import userRoute from '../modules/users/user.route';
import crawlJobRoute from '../modules/crawl-jobs/crawl-job.route';
import crawlExportRoute from '../modules/crawl-exports/crawl-export.route';
const router = Router(); const router = Router();
...@@ -12,7 +10,5 @@ router.get('/health', (req, res) => { ...@@ -12,7 +10,5 @@ router.get('/health', (req, res) => {
router.use('/auth', authRoute); router.use('/auth', authRoute);
router.use('/users', userRoute); router.use('/users', userRoute);
router.use('/crawl-jobs', crawlJobRoute);
router.use('/exports', crawlExportRoute);
export default router; export default router;
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment