/** * Background Processor Service * Handles batch processing of historical time entries data for analytics */ import { TimeEntry } from '@/lib/types/database'; import { analyticsEngine } from './analytics-engine'; import { llmAnalyzer } from './llm-analyzer'; import { postgresClient } from './postgres-client'; export interface BatchAnalysisJob { id: string; type: 'scoring' | 'llm_analysis' | 'full_analysis'; status: 'pending' | 'running' | 'completed' | 'failed'; progress: number; // 0 to 100 totalRecords: number; processedRecords: number; startTime: Date; endTime?: Date; errorMessage?: string; filters?: { startDate?: Date; endDate?: Date; resourceIds?: number[]; projectIds?: number[]; }; } export interface ProcessingResult { jobId: string; success: boolean; recordsProcessed: number; errors: string[]; processingTime: number; // milliseconds } export class BackgroundProcessor { private jobs: Map = new Map(); private isProcessing: boolean = false; private batchSize: number = 100; // Process 100 records at a time private maxConcurrentJobs: number = 3; private processingQueue: string[] = []; /** * Start a batch analysis job */ async startBatchAnalysis( type: 'scoring' | 'llm_analysis' | 'full_analysis', filters?: { startDate?: Date; endDate?: Date; resourceIds?: number[]; projectIds?: number[]; } ): Promise { const jobId = this.generateJobId(); const job: BatchAnalysisJob = { id: jobId, type, status: 'pending', progress: 0, totalRecords: 0, processedRecords: 0, startTime: new Date(), filters, }; this.jobs.set(jobId, job); this.processingQueue.push(jobId); // Start processing if not already running this.processQueue(); return jobId; } /** * Get job status */ getJobStatus(jobId: string): BatchAnalysisJob | null { return this.jobs.get(jobId) || null; } /** * Get all jobs */ getAllJobs(): BatchAnalysisJob[] { return Array.from(this.jobs.values()); } /** * Cancel a job */ cancelJob(jobId: string): boolean { const job = this.jobs.get(jobId); if (!job || job.status === 'completed' || job.status === 'running') { return false; } job.status = 'failed'; job.errorMessage = 'Job cancelled by user'; job.endTime = new Date(); // Remove from queue const queueIndex = this.processingQueue.indexOf(jobId); if (queueIndex > -1) { this.processingQueue.splice(queueIndex, 1); } return true; } /** * Process the job queue */ private async processQueue(): Promise { if (this.isProcessing || this.processingQueue.length === 0) { return; } this.isProcessing = true; while (this.processingQueue.length > 0) { const runningJobs = Array.from(this.jobs.values()) .filter(job => job.status === 'running').length; if (runningJobs >= this.maxConcurrentJobs) { break; // Wait for current jobs to finish } const jobId = this.processingQueue.shift(); if (!jobId) continue; const job = this.jobs.get(jobId); if (!job || job.status !== 'pending') continue; // Process job in background this.processJob(jobId).catch(error => { console.error(`Background job ${jobId} failed:`, error); }); } this.isProcessing = false; } /** * Process a single job */ private async processJob(jobId: string): Promise { const job = this.jobs.get(jobId); if (!job) return; try { job.status = 'running'; // Get time entries to process const timeEntries = await this.getTimeEntriesForJob(job); job.totalRecords = timeEntries.length; if (timeEntries.length === 0) { job.status = 'completed'; job.progress = 100; job.endTime = new Date(); return; } // Process based on job type switch (job.type) { case 'scoring': await this.processScoringJob(job, timeEntries); break; case 'llm_analysis': await this.processLLMAnalysisJob(job, timeEntries); break; case 'full_analysis': await this.processFullAnalysisJob(job, timeEntries); break; } job.status = 'completed'; job.progress = 100; job.endTime = new Date(); console.log(`Background job ${jobId} completed successfully`); } catch (error) { job.status = 'failed'; job.errorMessage = error instanceof Error ? error.message : 'Unknown error'; job.endTime = new Date(); console.error(`Background job ${jobId} failed:`, error); } // Continue processing queue setTimeout(() => this.processQueue(), 100); } /** * Process scoring job */ private async processScoringJob(job: BatchAnalysisJob, timeEntries: TimeEntry[]): Promise { const analyses = []; for (let i = 0; i < timeEntries.length; i += this.batchSize) { const batch = timeEntries.slice(i, i + this.batchSize); // Analyze batch for (const entry of batch) { const analysis = analyticsEngine.analyzeTimeEntry(entry); analyses.push(analysis); // Store analysis in database (would need to create analysis table) await this.storeTimeEntryAnalysis(entry.id, analysis); job.processedRecords++; job.progress = Math.round((job.processedRecords / job.totalRecords) * 100); } // Small delay to prevent overwhelming the system await new Promise(resolve => setTimeout(resolve, 10)); } console.log(`Scoring job ${job.id}: Analyzed ${analyses.length} time entries`); } /** * Process LLM analysis job */ private async processLLMAnalysisJob(job: BatchAnalysisJob, timeEntries: TimeEntry[]): Promise { // Process in larger batches for LLM to be more cost-effective const llmBatchSize = 500; for (let i = 0; i < timeEntries.length; i += llmBatchSize) { const batch = timeEntries.slice(i, i + llmBatchSize); try { const insights = await llmAnalyzer.generateInsights(batch); // Store insights in database await this.storeLLMInsights(job.id, batch, insights); job.processedRecords += batch.length; job.progress = Math.round((job.processedRecords / job.totalRecords) * 100); // Longer delay for LLM processing await new Promise(resolve => setTimeout(resolve, 1000)); } catch (error) { console.error(`LLM analysis failed for batch ${i}-${i + batch.length}:`, error); // Continue with next batch } } console.log(`LLM analysis job ${job.id}: Processed ${timeEntries.length} time entries`); } /** * Process full analysis job */ private async processFullAnalysisJob(job: BatchAnalysisJob, timeEntries: TimeEntry[]): Promise { // First, run scoring await this.processScoringJob(job, timeEntries); // Then, run LLM analysis on aggregated data const insights = await llmAnalyzer.generateInsights(timeEntries); const aggregateAnalysis = analyticsEngine.analyzeTimeEntries(timeEntries); // Store comprehensive results await this.storeFullAnalysisResults(job.id, aggregateAnalysis, insights); console.log(`Full analysis job ${job.id}: Completed comprehensive analysis`); } /** * Get time entries for job based on filters */ private async getTimeEntriesForJob(job: BatchAnalysisJob): Promise { let query = 'SELECT * FROM time_entries WHERE is_deleted = false'; const params: any[] = []; let paramIndex = 1; if (job.filters?.startDate) { query += ` AND entry_date >= $${paramIndex}`; params.push(job.filters.startDate.toISOString()); paramIndex++; } if (job.filters?.endDate) { query += ` AND entry_date <= $${paramIndex}`; params.push(job.filters.endDate.toISOString()); paramIndex++; } if (job.filters?.resourceIds && job.filters.resourceIds.length > 0) { query += ` AND resource_id = ANY($${paramIndex})`; params.push(job.filters.resourceIds); paramIndex++; } if (job.filters?.projectIds && job.filters.projectIds.length > 0) { query += ` AND project_id = ANY($${paramIndex})`; params.push(job.filters.projectIds); paramIndex++; } query += ' ORDER BY entry_date DESC'; const result = await postgresClient.query(query, params); return result.rows; } /** * Store time entry analysis in database */ private async storeTimeEntryAnalysis(timeEntryId: number, analysis: any): Promise { // This would create/update a time_entry_analyses table // For now, just log the analysis console.log(`Storing analysis for time entry ${timeEntryId}:`, analysis.overallScore); } /** * Store LLM insights in database */ private async storeLLMInsights(jobId: string, timeEntries: TimeEntry[], insights: any[]): Promise { // This would create/update an llm_insights table console.log(`Storing ${insights.length} LLM insights for job ${jobId}`); } /** * Store full analysis results */ private async storeFullAnalysisResults(jobId: string, analysis: any, insights: any[]): Promise { // This would create/update a comprehensive analysis results table console.log(`Storing full analysis results for job ${jobId}`); } /** * Generate unique job ID */ private generateJobId(): string { return `job_${Date.now()}_${Math.random().toString(36).substring(2, 9)}`; } /** * Clean up old completed jobs */ cleanupOldJobs(maxAge: number = 24 * 60 * 60 * 1000): void { // 24 hours const cutoff = Date.now() - maxAge; for (const [jobId, job] of this.jobs.entries()) { if ( (job.status === 'completed' || job.status === 'failed') && job.endTime && job.endTime.getTime() < cutoff ) { this.jobs.delete(jobId); } } } /** * Get processing statistics */ getStats(): { totalJobs: number; pendingJobs: number; runningJobs: number; completedJobs: number; failedJobs: number; queueLength: number; } { const jobs = Array.from(this.jobs.values()); return { totalJobs: jobs.length, pendingJobs: jobs.filter(j => j.status === 'pending').length, runningJobs: jobs.filter(j => j.status === 'running').length, completedJobs: jobs.filter(j => j.status === 'completed').length, failedJobs: jobs.filter(j => j.status === 'failed').length, queueLength: this.processingQueue.length, }; } } // Create singleton instance export const backgroundProcessor = new BackgroundProcessor(); // Schedule cleanup every hour setInterval(() => { backgroundProcessor.cleanupOldJobs(); }, 60 * 60 * 1000);