const nodemailer = require('nodemailer'); const { db } = require('../database/db'); const logger = require('../utils/logger'); let transporter = null; let lastConfigHash = null; // Generate hash from config for change detection function generateConfigHash(config) { const crypto = require('crypto'); const configString = `${config.smtp_host}:${config.smtp_port}:${config.smtp_user}:${config.smtp_pass}:${config.smtp_secure}`; return crypto.createHash('md5').update(configString).digest('hex'); } // Initialize transporter from database config async function initializeTransporter(forceReinit = false) { try { const config = await db('email_configs').first(); if (!config) { logger.warn('No email configuration found'); return null; } // Check if configuration has changed const currentConfigHash = generateConfigHash(config); if (!forceReinit && transporter && currentConfigHash === lastConfigHash) { // Configuration hasn't changed, return existing transporter return transporter; } // Configuration has changed or first initialization logger.info('Initializing email transporter' + (lastConfigHash && currentConfigHash !== lastConfigHash ? ' (configuration changed)' : '')); transporter = nodemailer.createTransport({ host: config.smtp_host, port: config.smtp_port, secure: config.smtp_secure, auth: config.smtp_user ? { user: config.smtp_user, pass: config.smtp_pass } : undefined }); // Verify configuration await transporter.verify(); logger.info('Email transporter initialized successfully'); // Update the config hash lastConfigHash = currentConfigHash; return transporter; } catch (error) { logger.error('Failed to initialize email transporter:', error); transporter = null; lastConfigHash = null; return null; } } // Get the appropriate language for a recipient async function getRecipientLanguage(email) { // For now, check if the email domain ends with .de // In the future, this could check user preferences if (email && email.endsWith('.de')) { return 'de'; } // Check if there's a saved preference for this email // This could be expanded to check user preferences in the database return 'en'; // Default to English } // Process email template with variables async function processTemplate(template, variables, language = 'en') { // Get the appropriate language fields const subjectField = language === 'de' ? 'subject_de' : 'subject_en'; const htmlField = language === 'de' ? 'body_html_de' : 'body_html_en'; const textField = language === 'de' ? 'body_text_de' : 'body_text_en'; // Fall back to non-language-specific fields for backward compatibility let subject = template[subjectField] || template.subject || ''; let htmlBody = template[htmlField] || template.body_html || ''; let textBody = template[textField] || template.body_text || ''; // Get branding settings for logo let logoUrl = ''; let companyName = 'PicPeak'; try { const brandingSettings = await db('app_settings') .whereIn('setting_key', ['branding_logo_url', 'branding_company_name']) .select('setting_key', 'setting_value'); brandingSettings.forEach(setting => { if (setting.setting_key === 'branding_logo_url' && setting.setting_value) { try { logoUrl = JSON.parse(setting.setting_value); } catch (e) { logoUrl = setting.setting_value; } } else if (setting.setting_key === 'branding_company_name' && setting.setting_value) { try { companyName = JSON.parse(setting.setting_value); } catch (e) { companyName = setting.setting_value; } } }); } catch (error) { logger.error('Error fetching branding settings:', error); } // If no custom logo, use default PicPeak logo const apiUrl = process.env.API_URL || 'http://localhost:3001'; const frontendUrl = process.env.FRONTEND_URL || 'http://localhost:3005'; const logoFullUrl = logoUrl ? `${apiUrl}${logoUrl}` : `${frontendUrl}/picpeak-logo-transparent.png`; // Process welcome message section if present let welcomeMessageSection = ''; if (variables.welcome_message && variables.welcome_message.trim() !== '') { const welcomeTitle = language === 'de' ? 'Persönliche Nachricht:' : 'Personal Message:'; welcomeMessageSection = `

${welcomeTitle}

${variables.welcome_message}

`; } // Replace variables Object.entries(variables).forEach(([key, value]) => { const regex = new RegExp(`{{${key}}}`, 'g'); subject = subject.replace(regex, value || ''); htmlBody = htmlBody.replace(regex, value || ''); textBody = textBody.replace(regex, value || ''); }); // Replace welcome message section placeholder htmlBody = htmlBody.replace(/{{welcome_message_section}}/g, welcomeMessageSection); // Wrap HTML body in styled template const styledHtmlBody = ` ${subject}
`; return { subject, htmlBody: styledHtmlBody, textBody }; } // Send email using template async function sendTemplateEmail(to, templateKey, variables) { try { // Always check for configuration changes before sending transporter = await initializeTransporter(); if (!transporter) { throw new Error('Email service not configured'); } // Get email template const template = await db('email_templates') .where('template_key', templateKey) .first(); if (!template) { throw new Error(`Email template '${templateKey}' not found`); } // Get email config for from address const config = await db('email_configs').first(); if (!config) { throw new Error('Email configuration not found'); } // Determine recipient language const language = await getRecipientLanguage(to); // Process template with variables const { subject, htmlBody, textBody } = await processTemplate(template, variables, language); // Send email const info = await transporter.sendMail({ from: `${config.from_name} <${config.from_email}>`, to: to, subject: subject, html: htmlBody, text: textBody || htmlBody.replace(/<[^>]*>/g, '') // Strip HTML if no text version }); logger.info(`Email sent successfully: ${info.messageId} (${language})`); return { success: true, messageId: info.messageId, language }; } catch (error) { logger.error('Error sending template email:', error); throw error; } } // Process email queue async function processEmailQueue() { logger.info('Email queue processor: Checking for pending emails...'); try { // Try to initialize transporter if it's null (in case it failed at startup) if (!transporter) { logger.info('Transporter not initialized, attempting to initialize...'); transporter = await initializeTransporter(); if (!transporter) { logger.warn('Email transporter could not be initialized, skipping queue processing'); return; } } let pendingEmails = []; try { pendingEmails = await db('email_queue') .where('status', 'pending') .where('retry_count', '<', 3) .orderBy('created_at', 'asc') .limit(10); } catch (dbError) { logger.error('Failed to query email queue:', dbError); return; } if (pendingEmails.length === 0) { logger.info('Email queue processor: No pending emails found'); return; } logger.info(`Processing ${pendingEmails.length} emails from queue`); for (const email of pendingEmails) { try { const emailData = typeof email.email_data === 'string' ? JSON.parse(email.email_data || '{}') : email.email_data || {}; await sendTemplateEmail( email.recipient_email, email.email_type, emailData ); // Mark as sent await db('email_queue') .where('id', email.id) .update({ status: 'sent', sent_at: new Date() }); logger.info(`Email ${email.id} sent successfully`); } catch (error) { // Increment retry count try { await db('email_queue') .where('id', email.id) .update({ retry_count: email.retry_count + 1, error_message: error.message }); } catch (updateError) { logger.error(`Failed to update email retry count for ${email.id}:`, updateError); // If update fails due to column issue, try without any potential auto-added fields if (updateError.message && updateError.message.includes('updated_at')) { logger.warn('Detected updated_at column issue, attempting raw query...'); await db.raw( 'UPDATE email_queue SET retry_count = ?, error_message = ? WHERE id = ?', [email.retry_count + 1, error.message, email.id] ); } } logger.error(`Failed to send email ${email.id}:`, error); } } } catch (error) { logger.error('Error processing email queue:', error); } } // Queue an email for sending async function queueEmail(eventId, recipientEmail, emailType, emailData) { try { await db('email_queue').insert({ event_id: eventId, recipient_email: recipientEmail, email_type: emailType, email_data: JSON.stringify(emailData), status: 'pending', retry_count: 0, created_at: new Date() }); logger.info(`Email queued: ${emailType} to ${recipientEmail}`); } catch (error) { logger.error('Error queueing email:', error); throw error; } } // Test email connection async function testEmailConnection() { try { if (!transporter) { await initializeTransporter(); } if (!transporter) { return false; } await transporter.verify(); return true; } catch (error) { logger.error('Email connection test failed:', error); return false; } } // Start email queue processor let emailQueueInterval = null; function startEmailQueueProcessor() { logger.info('Email queue processor: Attempting to start...'); if (!emailQueueInterval) { // Process immediately on start processEmailQueue().catch(err => { logger.error('Email queue processor: Initial processing failed:', err); }); // Then process every minute emailQueueInterval = setInterval(() => { processEmailQueue().catch(err => { logger.error('Email queue processor: Periodic processing failed:', err); }); }, 60000); logger.info('Email queue processor started successfully'); } else { logger.info('Email queue processor: Already running'); } } function stopEmailQueueProcessor() { if (emailQueueInterval) { clearInterval(emailQueueInterval); emailQueueInterval = null; logger.info('Email queue processor stopped'); } } // Initialize on module load - DISABLED for production startup // This will be called from server.js after database is ready // initializeTransporter().then(() => { // startEmailQueueProcessor(); // }); module.exports = { initializeTransporter, startEmailQueueProcessor, sendTemplateEmail, processEmailQueue, queueEmail, stopEmailQueueProcessor, testEmailConnection };