Initial commit
This commit is contained in:
@@ -0,0 +1,72 @@
|
||||
const fs = require("fs").promises;
|
||||
const path = require("path");
|
||||
const config = require("../config");
|
||||
|
||||
class AttachmentHandler {
|
||||
constructor() {
|
||||
this.attachmentsDir = path.join(__dirname, "..", "attachments");
|
||||
}
|
||||
|
||||
async getAttachments() {
|
||||
if (!config.attachments.enabled || config.attachments.files.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const attachments = [];
|
||||
|
||||
for (const filename of config.attachments.files) {
|
||||
try {
|
||||
const filePath = path.join(this.attachmentsDir, filename);
|
||||
|
||||
// Check if file exists
|
||||
await fs.access(filePath);
|
||||
|
||||
// Get file stats for size validation
|
||||
const stats = await fs.stat(filePath);
|
||||
|
||||
// Skip files larger than 10MB to avoid email size issues
|
||||
if (stats.size > 10 * 1024 * 1024) {
|
||||
console.warn(`Skipping ${filename}: File too large (>10MB)`);
|
||||
continue;
|
||||
}
|
||||
|
||||
attachments.push({
|
||||
filename: filename,
|
||||
path: filePath,
|
||||
});
|
||||
} catch (error) {
|
||||
console.warn(`Unable to attach ${filename}: ${error.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
return attachments;
|
||||
}
|
||||
|
||||
// Helper to validate attachment types
|
||||
isValidFileType(filename) {
|
||||
const validExtensions = [
|
||||
".pdf",
|
||||
".doc",
|
||||
".docx",
|
||||
".txt",
|
||||
".jpg",
|
||||
".jpeg",
|
||||
".png",
|
||||
".gif",
|
||||
];
|
||||
const ext = path.extname(filename).toLowerCase();
|
||||
return validExtensions.includes(ext);
|
||||
}
|
||||
|
||||
// Get list of available attachments for logging
|
||||
async listAvailableFiles() {
|
||||
try {
|
||||
const files = await fs.readdir(this.attachmentsDir);
|
||||
return files.filter((file) => this.isValidFileType(file));
|
||||
} catch (error) {
|
||||
return [];
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new AttachmentHandler();
|
||||
+398
@@ -0,0 +1,398 @@
|
||||
const sqlite3 = require("sqlite3").verbose();
|
||||
const { open } = require("sqlite");
|
||||
const path = require("path");
|
||||
const logger = require("./logger");
|
||||
|
||||
class Database {
|
||||
constructor() {
|
||||
this.db = null;
|
||||
this.dbPath = path.join(__dirname, "..", "data", "outreach.db");
|
||||
}
|
||||
|
||||
async init() {
|
||||
try {
|
||||
// Create data directory if it doesn't exist
|
||||
const fs = require("fs");
|
||||
const dataDir = path.dirname(this.dbPath);
|
||||
if (!fs.existsSync(dataDir)) {
|
||||
fs.mkdirSync(dataDir, { recursive: true });
|
||||
}
|
||||
|
||||
// Open database connection
|
||||
this.db = await open({
|
||||
filename: this.dbPath,
|
||||
driver: sqlite3.Database,
|
||||
});
|
||||
|
||||
logger.info("Database connected", { dbPath: this.dbPath });
|
||||
|
||||
// Create tables
|
||||
await this.createTables();
|
||||
|
||||
return this.db;
|
||||
} catch (error) {
|
||||
logger.error("Database initialization failed", { error: error.message });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async createTables() {
|
||||
const tables = [
|
||||
// Firms table
|
||||
`CREATE TABLE IF NOT EXISTS firms (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
firm_name TEXT NOT NULL,
|
||||
location TEXT,
|
||||
website TEXT,
|
||||
contact_email TEXT NOT NULL,
|
||||
state TEXT,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
updated_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
UNIQUE(contact_email)
|
||||
)`,
|
||||
|
||||
// Email campaigns table
|
||||
`CREATE TABLE IF NOT EXISTS campaigns (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT NOT NULL,
|
||||
subject TEXT NOT NULL,
|
||||
template_name TEXT NOT NULL,
|
||||
test_mode BOOLEAN DEFAULT 0,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
started_at DATETIME,
|
||||
completed_at DATETIME,
|
||||
total_emails INTEGER DEFAULT 0,
|
||||
sent_emails INTEGER DEFAULT 0,
|
||||
failed_emails INTEGER DEFAULT 0
|
||||
)`,
|
||||
|
||||
// Email sends table
|
||||
`CREATE TABLE IF NOT EXISTS email_sends (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
campaign_id INTEGER,
|
||||
firm_id INTEGER,
|
||||
recipient_email TEXT NOT NULL,
|
||||
subject TEXT NOT NULL,
|
||||
status TEXT NOT NULL, -- 'sent', 'failed', 'retry', 'permanent_failure'
|
||||
error_type TEXT,
|
||||
error_message TEXT,
|
||||
retry_count INTEGER DEFAULT 0,
|
||||
tracking_id TEXT UNIQUE, -- Add tracking ID for email tracking
|
||||
sent_at DATETIME,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
FOREIGN KEY (campaign_id) REFERENCES campaigns (id),
|
||||
FOREIGN KEY (firm_id) REFERENCES firms (id)
|
||||
)`,
|
||||
|
||||
// Tracking events table for email opens and clicks
|
||||
`CREATE TABLE IF NOT EXISTS tracking_events (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
tracking_id TEXT NOT NULL,
|
||||
event_type TEXT NOT NULL, -- 'open', 'click'
|
||||
event_data TEXT, -- JSON data (link_id, target_url, etc.)
|
||||
ip_address TEXT,
|
||||
user_agent TEXT,
|
||||
referer TEXT,
|
||||
created_at DATETIME DEFAULT CURRENT_TIMESTAMP,
|
||||
FOREIGN KEY (tracking_id) REFERENCES email_sends (tracking_id)
|
||||
)`,
|
||||
|
||||
// Create indexes for better performance
|
||||
`CREATE INDEX IF NOT EXISTS idx_firms_email ON firms(contact_email)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_firms_state ON firms(state)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_email_sends_campaign ON email_sends(campaign_id)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_email_sends_status ON email_sends(status)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_email_sends_tracking ON email_sends(tracking_id)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_tracking_events_tracking_id ON tracking_events(tracking_id)`,
|
||||
`CREATE INDEX IF NOT EXISTS idx_tracking_events_type ON tracking_events(event_type)`,
|
||||
];
|
||||
|
||||
for (const table of tables) {
|
||||
await this.db.exec(table);
|
||||
}
|
||||
|
||||
logger.info("Database tables created/verified");
|
||||
}
|
||||
|
||||
// Firm operations
|
||||
async insertFirm(firmData) {
|
||||
const { firmName, location, website, contactEmail, state } = firmData;
|
||||
|
||||
try {
|
||||
const result = await this.db.run(
|
||||
`INSERT OR IGNORE INTO firms (firm_name, location, website, contact_email, state)
|
||||
VALUES (?, ?, ?, ?, ?)`,
|
||||
[firmName, location, website, contactEmail, state]
|
||||
);
|
||||
|
||||
return result.lastID;
|
||||
} catch (error) {
|
||||
logger.error("Failed to insert firm", { firmData, error: error.message });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async insertFirmsBatch(firms) {
|
||||
const stmt = await this.db.prepare(
|
||||
`INSERT OR IGNORE INTO firms (firm_name, location, website, contact_email, state)
|
||||
VALUES (?, ?, ?, ?, ?)`
|
||||
);
|
||||
|
||||
let inserted = 0;
|
||||
for (const firm of firms) {
|
||||
try {
|
||||
const result = await stmt.run([
|
||||
firm.firmName,
|
||||
firm.location,
|
||||
firm.website,
|
||||
firm.contactEmail || firm.email,
|
||||
firm.state,
|
||||
]);
|
||||
if (result.changes > 0) inserted++;
|
||||
} catch (error) {
|
||||
logger.warn("Failed to insert firm", {
|
||||
firm: firm.firmName,
|
||||
error: error.message,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
await stmt.finalize();
|
||||
logger.info(`Inserted ${inserted} firms into database`);
|
||||
return inserted;
|
||||
}
|
||||
|
||||
async getFirms(limit = null, offset = 0) {
|
||||
const sql = limit
|
||||
? `SELECT * FROM firms ORDER BY id LIMIT ? OFFSET ?`
|
||||
: `SELECT * FROM firms ORDER BY id`;
|
||||
|
||||
const params = limit ? [limit, offset] : [];
|
||||
return await this.db.all(sql, params);
|
||||
}
|
||||
|
||||
async getFirmByEmail(email) {
|
||||
return await this.db.get("SELECT * FROM firms WHERE contact_email = ?", [
|
||||
email,
|
||||
]);
|
||||
}
|
||||
|
||||
async getFirmById(id) {
|
||||
return await this.db.get("SELECT * FROM firms WHERE id = ?", [id]);
|
||||
}
|
||||
|
||||
async removeDuplicateFirms() {
|
||||
// Remove duplicates keeping the first occurrence
|
||||
const result = await this.db.run(`
|
||||
DELETE FROM firms
|
||||
WHERE id NOT IN (
|
||||
SELECT MIN(id)
|
||||
FROM firms
|
||||
GROUP BY contact_email
|
||||
)
|
||||
`);
|
||||
|
||||
logger.info(`Removed ${result.changes} duplicate firms`);
|
||||
return result.changes;
|
||||
}
|
||||
|
||||
// Campaign operations
|
||||
async createCampaign(campaignData) {
|
||||
const { name, subject, templateName, testMode } = campaignData;
|
||||
|
||||
const result = await this.db.run(
|
||||
`INSERT INTO campaigns (name, subject, template_name, test_mode)
|
||||
VALUES (?, ?, ?, ?)`,
|
||||
[name, subject, templateName, testMode ? 1 : 0]
|
||||
);
|
||||
|
||||
logger.info("Campaign created", { campaignId: result.lastID, name });
|
||||
return result.lastID;
|
||||
}
|
||||
|
||||
async startCampaign(campaignId, totalEmails) {
|
||||
await this.db.run(
|
||||
`UPDATE campaigns
|
||||
SET started_at = CURRENT_TIMESTAMP, total_emails = ?
|
||||
WHERE id = ?`,
|
||||
[totalEmails, campaignId]
|
||||
);
|
||||
}
|
||||
|
||||
async completeCampaign(campaignId, stats) {
|
||||
await this.db.run(
|
||||
`UPDATE campaigns
|
||||
SET completed_at = CURRENT_TIMESTAMP,
|
||||
sent_emails = ?,
|
||||
failed_emails = ?
|
||||
WHERE id = ?`,
|
||||
[stats.sent, stats.failed, campaignId]
|
||||
);
|
||||
}
|
||||
|
||||
// Email send tracking
|
||||
async logEmailSend(emailData) {
|
||||
const {
|
||||
campaignId,
|
||||
firmId,
|
||||
recipientEmail,
|
||||
subject,
|
||||
status,
|
||||
errorType,
|
||||
errorMessage,
|
||||
retryCount,
|
||||
trackingId,
|
||||
} = emailData;
|
||||
|
||||
const result = await this.db.run(
|
||||
`INSERT INTO email_sends
|
||||
(campaign_id, firm_id, recipient_email, subject, status, error_type, error_message, retry_count, tracking_id, sent_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||||
[
|
||||
campaignId,
|
||||
firmId,
|
||||
recipientEmail,
|
||||
subject,
|
||||
status,
|
||||
errorType,
|
||||
errorMessage,
|
||||
retryCount,
|
||||
trackingId,
|
||||
status === "sent" ? new Date().toISOString() : null,
|
||||
]
|
||||
);
|
||||
|
||||
return result.lastID;
|
||||
}
|
||||
|
||||
async updateEmailSendStatus(
|
||||
emailSendId,
|
||||
status,
|
||||
errorType = null,
|
||||
errorMessage = null
|
||||
) {
|
||||
await this.db.run(
|
||||
`UPDATE email_sends
|
||||
SET status = ?, error_type = ?, error_message = ?,
|
||||
sent_at = CASE WHEN ? = 'sent' THEN CURRENT_TIMESTAMP ELSE sent_at END
|
||||
WHERE id = ?`,
|
||||
[status, errorType, errorMessage, status, emailSendId]
|
||||
);
|
||||
}
|
||||
|
||||
// Statistics and reporting
|
||||
async getCampaignStats(campaignId) {
|
||||
const stats = await this.db.all(
|
||||
`
|
||||
SELECT
|
||||
status,
|
||||
COUNT(*) as count
|
||||
FROM email_sends
|
||||
WHERE campaign_id = ?
|
||||
GROUP BY status
|
||||
`,
|
||||
[campaignId]
|
||||
);
|
||||
|
||||
const result = {
|
||||
sent: 0,
|
||||
failed: 0,
|
||||
retry: 0,
|
||||
permanent_failure: 0,
|
||||
};
|
||||
|
||||
stats.forEach((stat) => {
|
||||
result[stat.status] = stat.count;
|
||||
});
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
async getFailedEmails(campaignId) {
|
||||
return await this.db.all(
|
||||
`
|
||||
SELECT es.*, f.firm_name, f.contact_email
|
||||
FROM email_sends es
|
||||
JOIN firms f ON es.firm_id = f.id
|
||||
WHERE es.campaign_id = ? AND es.status IN ('failed', 'permanent_failure')
|
||||
ORDER BY es.created_at DESC
|
||||
`,
|
||||
[campaignId]
|
||||
);
|
||||
}
|
||||
|
||||
async close() {
|
||||
if (this.db) {
|
||||
await this.db.close();
|
||||
this.db = null;
|
||||
logger.info("Database connection closed");
|
||||
}
|
||||
}
|
||||
|
||||
// Utility methods
|
||||
async getTableCounts() {
|
||||
const counts = {};
|
||||
const tables = ["firms", "campaigns", "email_sends", "tracking_events"];
|
||||
|
||||
for (const table of tables) {
|
||||
const result = await this.db.get(
|
||||
`SELECT COUNT(*) as count FROM ${table}`
|
||||
);
|
||||
counts[table] = result.count;
|
||||
}
|
||||
|
||||
return counts;
|
||||
}
|
||||
|
||||
// Tracking events methods
|
||||
async storeTrackingEvent(trackingId, eventType, eventData = {}) {
|
||||
const query = `
|
||||
INSERT INTO tracking_events (tracking_id, event_type, event_data, ip_address, user_agent, referer)
|
||||
VALUES (?, ?, ?, ?, ?, ?)
|
||||
`;
|
||||
|
||||
await this.db.run(query, [
|
||||
trackingId,
|
||||
eventType,
|
||||
JSON.stringify(eventData),
|
||||
eventData.ip || null,
|
||||
eventData.userAgent || null,
|
||||
eventData.referer || null,
|
||||
]);
|
||||
}
|
||||
|
||||
async getTrackingEvents(trackingId) {
|
||||
const query = `
|
||||
SELECT * FROM tracking_events
|
||||
WHERE tracking_id = ?
|
||||
ORDER BY created_at ASC
|
||||
`;
|
||||
|
||||
const events = await this.db.all(query, [trackingId]);
|
||||
return events.map((event) => ({
|
||||
...event,
|
||||
event_data: JSON.parse(event.event_data || "{}"),
|
||||
}));
|
||||
}
|
||||
|
||||
async getTrackingStats(campaignId) {
|
||||
const query = `
|
||||
SELECT
|
||||
es.tracking_id,
|
||||
es.recipient_email,
|
||||
COUNT(CASE WHEN te.event_type = 'open' THEN 1 END) as opens,
|
||||
COUNT(CASE WHEN te.event_type = 'click' THEN 1 END) as clicks,
|
||||
MIN(CASE WHEN te.event_type = 'open' THEN te.created_at END) as first_open,
|
||||
MIN(CASE WHEN te.event_type = 'click' THEN te.created_at END) as first_click
|
||||
FROM email_sends es
|
||||
LEFT JOIN tracking_events te ON es.tracking_id = te.tracking_id
|
||||
WHERE es.campaign_id = ? AND es.status = 'sent'
|
||||
GROUP BY es.tracking_id, es.recipient_email
|
||||
ORDER BY es.sent_at ASC
|
||||
`;
|
||||
|
||||
return await this.db.all(query, [campaignId]);
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new Database();
|
||||
@@ -0,0 +1,242 @@
|
||||
const config = require("../config");
|
||||
const logger = require("./logger");
|
||||
|
||||
class ErrorHandler {
|
||||
constructor() {
|
||||
this.failedEmails = [];
|
||||
this.retryAttempts = new Map(); // Track retry attempts per email
|
||||
}
|
||||
|
||||
// Classify error types
|
||||
classifyError(error) {
|
||||
const errorMessage = error.message.toLowerCase();
|
||||
|
||||
if (
|
||||
errorMessage.includes("invalid login") ||
|
||||
errorMessage.includes("authentication") ||
|
||||
errorMessage.includes("unauthorized")
|
||||
) {
|
||||
return "AUTH_ERROR";
|
||||
}
|
||||
|
||||
if (
|
||||
errorMessage.includes("rate limit") ||
|
||||
errorMessage.includes("too many requests") ||
|
||||
errorMessage.includes("quota exceeded")
|
||||
) {
|
||||
return "RATE_LIMIT";
|
||||
}
|
||||
|
||||
if (
|
||||
errorMessage.includes("network") ||
|
||||
errorMessage.includes("connection") ||
|
||||
errorMessage.includes("timeout") ||
|
||||
errorMessage.includes("econnrefused")
|
||||
) {
|
||||
return "NETWORK_ERROR";
|
||||
}
|
||||
|
||||
if (
|
||||
errorMessage.includes("invalid recipient") ||
|
||||
errorMessage.includes("mailbox unavailable") ||
|
||||
errorMessage.includes("user unknown")
|
||||
) {
|
||||
return "RECIPIENT_ERROR";
|
||||
}
|
||||
|
||||
if (
|
||||
errorMessage.includes("message too large") ||
|
||||
errorMessage.includes("attachment")
|
||||
) {
|
||||
return "MESSAGE_ERROR";
|
||||
}
|
||||
|
||||
return "UNKNOWN_ERROR";
|
||||
}
|
||||
|
||||
// Determine if error is retryable
|
||||
isRetryable(errorType) {
|
||||
const retryableErrors = ["RATE_LIMIT", "NETWORK_ERROR", "UNKNOWN_ERROR"];
|
||||
|
||||
const nonRetryableErrors = [
|
||||
"AUTH_ERROR",
|
||||
"RECIPIENT_ERROR",
|
||||
"MESSAGE_ERROR",
|
||||
];
|
||||
|
||||
return retryableErrors.includes(errorType);
|
||||
}
|
||||
|
||||
// Calculate exponential backoff delay
|
||||
getRetryDelay(attemptNumber) {
|
||||
// Base delay: 1 minute, exponentially increasing
|
||||
const baseDelay = 60 * 1000; // 1 minute in ms
|
||||
const exponentialDelay = baseDelay * Math.pow(2, attemptNumber - 1);
|
||||
|
||||
// Add jitter (±25%)
|
||||
const jitter = exponentialDelay * 0.25 * (Math.random() - 0.5);
|
||||
|
||||
// Cap at maximum delay (30 minutes)
|
||||
const maxDelay = 30 * 60 * 1000;
|
||||
|
||||
return Math.min(exponentialDelay + jitter, maxDelay);
|
||||
}
|
||||
|
||||
// Handle email sending error
|
||||
async handleError(email, recipient, error, transporter) {
|
||||
const errorType = this.classifyError(error);
|
||||
const emailKey = `${recipient}_${Date.now()}`;
|
||||
|
||||
logger.emailFailed(
|
||||
recipient,
|
||||
error,
|
||||
errorType,
|
||||
email.firmName || "Unknown"
|
||||
);
|
||||
console.error(`❌ Error sending to ${recipient}: ${error.message}`);
|
||||
console.error(` Error Type: ${errorType}`);
|
||||
|
||||
// Get current retry count
|
||||
const currentAttempts = this.retryAttempts.get(emailKey) || 0;
|
||||
const maxRetries = config.errorHandling?.maxRetries || 3;
|
||||
|
||||
if (this.isRetryable(errorType) && currentAttempts < maxRetries) {
|
||||
// Schedule retry
|
||||
const retryDelay = this.getRetryDelay(currentAttempts + 1);
|
||||
this.retryAttempts.set(emailKey, currentAttempts + 1);
|
||||
|
||||
logger.emailRetry(
|
||||
recipient,
|
||||
currentAttempts + 1,
|
||||
maxRetries,
|
||||
retryDelay,
|
||||
errorType
|
||||
);
|
||||
console.warn(
|
||||
`🔄 Scheduling retry ${
|
||||
currentAttempts + 1
|
||||
}/${maxRetries} for ${recipient} in ${Math.round(retryDelay / 1000)}s`
|
||||
);
|
||||
|
||||
// Add to retry queue
|
||||
this.failedEmails.push({
|
||||
email,
|
||||
recipient,
|
||||
error: errorType,
|
||||
attempts: currentAttempts + 1,
|
||||
retryAt: Date.now() + retryDelay,
|
||||
originalError: error.message,
|
||||
});
|
||||
|
||||
return true; // Indicates retry scheduled
|
||||
} else {
|
||||
// Max retries reached or non-retryable error
|
||||
logger.emailPermanentFailure(recipient, errorType, currentAttempts);
|
||||
console.error(
|
||||
`💀 Permanent failure for ${recipient}: ${errorType} (${currentAttempts} attempts)`
|
||||
);
|
||||
|
||||
// Log permanently failed email
|
||||
this.logPermanentFailure(
|
||||
email,
|
||||
recipient,
|
||||
error,
|
||||
errorType,
|
||||
currentAttempts
|
||||
);
|
||||
|
||||
return false; // Indicates permanent failure
|
||||
}
|
||||
}
|
||||
|
||||
// Log permanent failures for later review
|
||||
logPermanentFailure(email, recipient, error, errorType, attempts) {
|
||||
const failure = {
|
||||
timestamp: new Date().toISOString(),
|
||||
recipient,
|
||||
errorType,
|
||||
attempts,
|
||||
error: error.message,
|
||||
emailData: {
|
||||
subject: email.subject,
|
||||
firmName: email.firmName,
|
||||
},
|
||||
};
|
||||
|
||||
// You could write this to a file or database
|
||||
console.error("🚨 PERMANENT FAILURE:", JSON.stringify(failure, null, 2));
|
||||
}
|
||||
|
||||
// Process retry queue
|
||||
async processRetries(transporter, sendEmailFunction) {
|
||||
const now = Date.now();
|
||||
const readyToRetry = this.failedEmails.filter(
|
||||
(item) => item.retryAt <= now
|
||||
);
|
||||
|
||||
if (readyToRetry.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
console.log(`🔄 Processing ${readyToRetry.length} email retries...`);
|
||||
|
||||
for (const retryItem of readyToRetry) {
|
||||
try {
|
||||
console.log(
|
||||
`🔄 Retrying ${retryItem.recipient} (attempt ${retryItem.attempts})`
|
||||
);
|
||||
|
||||
// Attempt to send email again
|
||||
await sendEmailFunction(
|
||||
retryItem.email,
|
||||
retryItem.recipient,
|
||||
transporter
|
||||
);
|
||||
|
||||
// Success - remove from retry queue
|
||||
this.failedEmails = this.failedEmails.filter(
|
||||
(item) => item !== retryItem
|
||||
);
|
||||
console.log(`✅ Retry successful for ${retryItem.recipient}`);
|
||||
} catch (error) {
|
||||
// Handle retry failure
|
||||
await this.handleError(
|
||||
retryItem.email,
|
||||
retryItem.recipient,
|
||||
error,
|
||||
transporter
|
||||
);
|
||||
|
||||
// Remove the processed item from queue
|
||||
this.failedEmails = this.failedEmails.filter(
|
||||
(item) => item !== retryItem
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Get retry queue status
|
||||
getRetryStats() {
|
||||
const now = Date.now();
|
||||
const pending = this.failedEmails.filter((item) => item.retryAt > now);
|
||||
const ready = this.failedEmails.filter((item) => item.retryAt <= now);
|
||||
|
||||
return {
|
||||
totalFailed: this.failedEmails.length,
|
||||
pendingRetries: pending.length,
|
||||
readyToRetry: ready.length,
|
||||
nextRetryIn:
|
||||
pending.length > 0
|
||||
? Math.min(...pending.map((item) => item.retryAt - now))
|
||||
: 0,
|
||||
};
|
||||
}
|
||||
|
||||
// Clear retry queue (for testing)
|
||||
clearRetries() {
|
||||
this.failedEmails = [];
|
||||
this.retryAttempts.clear();
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new ErrorHandler();
|
||||
+159
@@ -0,0 +1,159 @@
|
||||
const winston = require("winston");
|
||||
const path = require("path");
|
||||
const config = require("../config");
|
||||
|
||||
// Create logs directory if it doesn't exist
|
||||
const fs = require("fs");
|
||||
const logsDir = path.join(__dirname, "..", "logs");
|
||||
if (!fs.existsSync(logsDir)) {
|
||||
fs.mkdirSync(logsDir);
|
||||
}
|
||||
|
||||
// Custom format for console output
|
||||
const consoleFormat = winston.format.combine(
|
||||
winston.format.timestamp({ format: "HH:mm:ss" }),
|
||||
winston.format.colorize(),
|
||||
winston.format.printf(({ timestamp, level, message, ...meta }) => {
|
||||
let metaStr = "";
|
||||
if (Object.keys(meta).length > 0) {
|
||||
metaStr = " " + JSON.stringify(meta);
|
||||
}
|
||||
return `${timestamp} [${level}] ${message}${metaStr}`;
|
||||
})
|
||||
);
|
||||
|
||||
// Custom format for file output
|
||||
const fileFormat = winston.format.combine(
|
||||
winston.format.timestamp(),
|
||||
winston.format.errors({ stack: true }),
|
||||
winston.format.json()
|
||||
);
|
||||
|
||||
// Create the logger
|
||||
const logger = winston.createLogger({
|
||||
level: config.logging.level,
|
||||
format: fileFormat,
|
||||
defaultMeta: {
|
||||
service: "outreach-engine",
|
||||
environment: config.env,
|
||||
},
|
||||
transports: [
|
||||
// Error log file
|
||||
new winston.transports.File({
|
||||
filename: path.join(logsDir, "error.log"),
|
||||
level: "error",
|
||||
maxsize: 5242880, // 5MB
|
||||
maxFiles: 5,
|
||||
}),
|
||||
|
||||
// Combined log file
|
||||
new winston.transports.File({
|
||||
filename: path.join(logsDir, "combined.log"),
|
||||
maxsize: 5242880, // 5MB
|
||||
maxFiles: 5,
|
||||
}),
|
||||
|
||||
// Email activity log
|
||||
new winston.transports.File({
|
||||
filename: path.join(logsDir, "email-activity.log"),
|
||||
level: "info",
|
||||
maxsize: 10485760, // 10MB
|
||||
maxFiles: 10,
|
||||
}),
|
||||
],
|
||||
});
|
||||
|
||||
// Add console transport for development
|
||||
if (config.isDevelopment) {
|
||||
logger.add(
|
||||
new winston.transports.Console({
|
||||
format: consoleFormat,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
// Custom logging methods for email activities
|
||||
logger.emailSent = (recipient, subject, firmName, testMode = false) => {
|
||||
logger.info("Email sent successfully", {
|
||||
event: "email_sent",
|
||||
recipient,
|
||||
subject,
|
||||
firmName,
|
||||
testMode,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.emailFailed = (recipient, error, errorType, firmName) => {
|
||||
logger.error("Email failed to send", {
|
||||
event: "email_failed",
|
||||
recipient,
|
||||
error: error.message,
|
||||
errorType,
|
||||
firmName,
|
||||
stack: error.stack,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.emailRetry = (recipient, attempt, maxRetries, retryDelay, errorType) => {
|
||||
logger.warn("Email retry scheduled", {
|
||||
event: "email_retry",
|
||||
recipient,
|
||||
attempt,
|
||||
maxRetries,
|
||||
retryDelay,
|
||||
errorType,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.emailPermanentFailure = (recipient, errorType, attempts) => {
|
||||
logger.error("Email permanent failure", {
|
||||
event: "email_permanent_failure",
|
||||
recipient,
|
||||
errorType,
|
||||
attempts,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.campaignStart = (totalEmails, testMode) => {
|
||||
logger.info("Email campaign started", {
|
||||
event: "campaign_start",
|
||||
totalEmails,
|
||||
testMode,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.campaignComplete = (stats) => {
|
||||
logger.info("Email campaign completed", {
|
||||
event: "campaign_complete",
|
||||
...stats,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
logger.rateLimitPause = (duration, reason) => {
|
||||
logger.warn("Rate limit pause activated", {
|
||||
event: "rate_limit_pause",
|
||||
duration,
|
||||
reason,
|
||||
timestamp: new Date().toISOString(),
|
||||
});
|
||||
};
|
||||
|
||||
// Helper method to log with email context
|
||||
logger.withEmail = (recipient, firmName) => {
|
||||
return {
|
||||
info: (message, meta = {}) =>
|
||||
logger.info(message, { ...meta, recipient, firmName }),
|
||||
warn: (message, meta = {}) =>
|
||||
logger.warn(message, { ...meta, recipient, firmName }),
|
||||
error: (message, meta = {}) =>
|
||||
logger.error(message, { ...meta, recipient, firmName }),
|
||||
};
|
||||
};
|
||||
|
||||
module.exports = logger;
|
||||
@@ -0,0 +1,107 @@
|
||||
const config = require("../config");
|
||||
const logger = require("./logger");
|
||||
|
||||
class RateLimiter {
|
||||
constructor() {
|
||||
this.sentCount = 0;
|
||||
this.startTime = Date.now();
|
||||
this.lastSentTime = null;
|
||||
}
|
||||
|
||||
// Record a successful send
|
||||
recordSuccess() {
|
||||
this.sentCount++;
|
||||
this.lastSentTime = Date.now();
|
||||
}
|
||||
|
||||
// Get randomized delay between min and max
|
||||
getRandomDelay() {
|
||||
const baseDelay = config.app.delayMinutes * 60 * 1000; // Convert to ms
|
||||
const minDelay = baseDelay * 0.8; // 20% less than base
|
||||
const maxDelay = baseDelay * 1.2; // 20% more than base
|
||||
|
||||
// Add additional randomization
|
||||
const randomFactor = Math.random() * (maxDelay - minDelay) + minDelay;
|
||||
|
||||
// Add jitter to avoid patterns
|
||||
const jitter = (Math.random() - 0.5) * 30000; // +/- 30 seconds
|
||||
|
||||
return Math.floor(randomFactor + jitter);
|
||||
}
|
||||
|
||||
// Check if we should pause based on sent count
|
||||
shouldPause() {
|
||||
const hoursSinceStart = (Date.now() - this.startTime) / (1000 * 60 * 60);
|
||||
const emailsPerHour = this.sentCount / hoursSinceStart;
|
||||
|
||||
// Gmail limits: ~500/day = ~20/hour
|
||||
// Be conservative: pause if exceeding 15/hour
|
||||
if (emailsPerHour > 15) {
|
||||
return true;
|
||||
}
|
||||
|
||||
// Pause every 50 emails for 30 minutes
|
||||
if (this.sentCount > 0 && this.sentCount % 50 === 0) {
|
||||
return true;
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
// Get pause duration
|
||||
getPauseDuration() {
|
||||
// Standard pause: 30 minutes
|
||||
const basePause = 30 * 60 * 1000;
|
||||
|
||||
// Add randomization to pause duration
|
||||
const randomPause = basePause + Math.random() * 10 * 60 * 1000; // +0-10 minutes
|
||||
|
||||
return randomPause;
|
||||
}
|
||||
|
||||
// Calculate next send time
|
||||
async getNextSendDelay() {
|
||||
if (this.shouldPause()) {
|
||||
const pauseDuration = this.getPauseDuration();
|
||||
const pauseMinutes = Math.round(pauseDuration / 60000);
|
||||
|
||||
logger.rateLimitPause(
|
||||
pauseDuration,
|
||||
`Automatic pause after ${this.sentCount} emails`
|
||||
);
|
||||
console.log(`Rate limit pause: ${pauseMinutes} minutes`);
|
||||
return pauseDuration;
|
||||
}
|
||||
|
||||
return this.getRandomDelay();
|
||||
}
|
||||
|
||||
// Get human-readable time
|
||||
formatDelay(ms) {
|
||||
const minutes = Math.floor(ms / 60000);
|
||||
const seconds = Math.floor((ms % 60000) / 1000);
|
||||
return `${minutes}m ${seconds}s`;
|
||||
}
|
||||
|
||||
// Reset counters (for testing)
|
||||
reset() {
|
||||
this.sentCount = 0;
|
||||
this.startTime = Date.now();
|
||||
this.lastSentTime = null;
|
||||
}
|
||||
|
||||
// Get current stats
|
||||
getStats() {
|
||||
const runtime = (Date.now() - this.startTime) / 1000; // seconds
|
||||
const avgRate = this.sentCount / (runtime / 3600); // emails per hour
|
||||
|
||||
return {
|
||||
sentCount: this.sentCount,
|
||||
runtime: Math.floor(runtime / 60), // minutes
|
||||
averageRate: avgRate.toFixed(1),
|
||||
nextDelay: this.formatDelay(this.getRandomDelay()),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new RateLimiter();
|
||||
@@ -0,0 +1,130 @@
|
||||
const fs = require("fs").promises;
|
||||
const path = require("path");
|
||||
const Handlebars = require("handlebars");
|
||||
const config = require("../config");
|
||||
|
||||
class TemplateEngine {
|
||||
constructor() {
|
||||
this.templatesDir = path.join(__dirname, "..", "templates");
|
||||
this.compiledTemplates = new Map();
|
||||
}
|
||||
|
||||
// Convert HTML to plain text by removing tags and formatting
|
||||
htmlToText(html) {
|
||||
return html
|
||||
.replace(/<style[^>]*>.*?<\/style>/gis, "") // Remove style blocks
|
||||
.replace(/<script[^>]*>.*?<\/script>/gis, "") // Remove script blocks
|
||||
.replace(/<br\s*\/?>/gi, "\n") // Convert <br> to newlines
|
||||
.replace(/<\/p>/gi, "\n\n") // Convert </p> to double newlines
|
||||
.replace(/<\/div>/gi, "\n") // Convert </div> to newlines
|
||||
.replace(/<\/h[1-6]>/gi, "\n\n") // Convert headings to double newlines
|
||||
.replace(/<li[^>]*>/gi, "• ") // Convert <li> to bullet points
|
||||
.replace(/<\/li>/gi, "\n") // End list items with newlines
|
||||
.replace(/<[^>]*>/g, "") // Remove all other HTML tags
|
||||
.replace(/ /g, " ") // Convert to spaces
|
||||
.replace(/&/g, "&") // Convert & to &
|
||||
.replace(/</g, "<") // Convert < to <
|
||||
.replace(/>/g, ">") // Convert > to >
|
||||
.replace(/"/g, '"') // Convert " to "
|
||||
.replace(/'/g, "'") // Convert ' to '
|
||||
.replace(/\n\s*\n\s*\n/g, "\n\n") // Reduce multiple newlines to double
|
||||
.replace(/^\s+|\s+$/gm, "") // Trim whitespace from lines
|
||||
.trim();
|
||||
}
|
||||
|
||||
async loadTemplate(templateName) {
|
||||
const cacheKey = templateName;
|
||||
|
||||
if (this.compiledTemplates.has(cacheKey)) {
|
||||
return this.compiledTemplates.get(cacheKey);
|
||||
}
|
||||
|
||||
try {
|
||||
const htmlPath = path.join(this.templatesDir, `${templateName}.html`);
|
||||
const htmlContent = await fs.readFile(htmlPath, "utf-8");
|
||||
|
||||
// Automatically inject GIF if enabled
|
||||
let finalHtmlContent = htmlContent;
|
||||
if (config.gif.enabled && templateName === "outreach") {
|
||||
finalHtmlContent = this.injectGifIntoHtml(htmlContent);
|
||||
}
|
||||
|
||||
const htmlTemplate = Handlebars.compile(finalHtmlContent);
|
||||
|
||||
// Generate text version from HTML
|
||||
const textTemplate = Handlebars.compile(
|
||||
this.htmlToText(finalHtmlContent)
|
||||
);
|
||||
|
||||
const compiledTemplate = {
|
||||
html: htmlTemplate,
|
||||
text: textTemplate,
|
||||
};
|
||||
|
||||
this.compiledTemplates.set(cacheKey, compiledTemplate);
|
||||
return compiledTemplate;
|
||||
} catch (error) {
|
||||
throw new Error(
|
||||
`Failed to load template ${templateName}: ${error.message}`
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Inject GIF into HTML template automatically
|
||||
injectGifIntoHtml(htmlContent) {
|
||||
const gifHtml = `
|
||||
{{#if gifUrl}}
|
||||
<div style="text-align: center; margin: 20px 0;">
|
||||
<img src="{{gifUrl}}" alt="{{gifAlt}}" style="max-width: 100%; height: auto; border-radius: 5px;" />
|
||||
</div>
|
||||
{{/if}}
|
||||
`;
|
||||
|
||||
// Insert GIF after the header or at the beginning of content
|
||||
if (htmlContent.includes('<div class="content">')) {
|
||||
return htmlContent.replace(
|
||||
'<div class="content">',
|
||||
`<div class="content">${gifHtml}`
|
||||
);
|
||||
} else if (htmlContent.includes("<body>")) {
|
||||
return htmlContent.replace("<body>", `<body>${gifHtml}`);
|
||||
} else {
|
||||
// Fallback: add at the beginning
|
||||
return gifHtml + htmlContent;
|
||||
}
|
||||
}
|
||||
|
||||
async render(templateName, data) {
|
||||
const template = await this.loadTemplate(templateName);
|
||||
|
||||
// Add default sender information from config
|
||||
const defaultData = {
|
||||
senderName: "John Smith",
|
||||
senderTitle: "Business Development Manager",
|
||||
senderCompany: "Legal Solutions Inc.",
|
||||
fromEmail: config.email.user,
|
||||
gifUrl: config.gif.enabled ? config.gif.url : null,
|
||||
gifAlt: config.gif.enabled ? config.gif.alt : null,
|
||||
...data,
|
||||
};
|
||||
|
||||
return {
|
||||
html: template.html(defaultData),
|
||||
text: template.text(defaultData),
|
||||
};
|
||||
}
|
||||
|
||||
// Helper to format firm data for template
|
||||
formatFirmData(firm) {
|
||||
return {
|
||||
firmName: firm.firmName || "your firm",
|
||||
location: firm.location,
|
||||
website: firm.website,
|
||||
email: firm.contactEmail || firm.email,
|
||||
greeting: firm.name || "Legal Professional",
|
||||
// Additional fields can be mapped here
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new TemplateEngine();
|
||||
@@ -0,0 +1,266 @@
|
||||
const http = require("http");
|
||||
const url = require("url");
|
||||
const path = require("path");
|
||||
const logger = require("./logger");
|
||||
const database = require("./database");
|
||||
|
||||
class TrackingServer {
|
||||
constructor() {
|
||||
this.server = null;
|
||||
this.port = process.env.TRACKING_PORT || 3000;
|
||||
this.trackingDomain =
|
||||
process.env.TRACKING_DOMAIN || `http://localhost:${this.port}`;
|
||||
}
|
||||
|
||||
// 1x1 transparent pixel GIF
|
||||
get trackingPixel() {
|
||||
return Buffer.from(
|
||||
"R0lGODlhAQABAIAAAAAAAP///yH5BAEAAAAALAAAAAABAAEAAAIBRAA7",
|
||||
"base64"
|
||||
);
|
||||
}
|
||||
|
||||
async start() {
|
||||
this.server = http.createServer((req, res) => {
|
||||
this.handleRequest(req, res);
|
||||
});
|
||||
|
||||
return new Promise((resolve, reject) => {
|
||||
this.server.listen(this.port, (err) => {
|
||||
if (err) {
|
||||
reject(err);
|
||||
} else {
|
||||
console.log(`📊 Tracking server started on ${this.trackingDomain}`);
|
||||
logger.info("Tracking server started", {
|
||||
port: this.port,
|
||||
domain: this.trackingDomain,
|
||||
});
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
async stop() {
|
||||
if (this.server) {
|
||||
return new Promise((resolve) => {
|
||||
this.server.close(() => {
|
||||
console.log("📊 Tracking server stopped");
|
||||
logger.info("Tracking server stopped");
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async handleRequest(req, res) {
|
||||
const parsedUrl = url.parse(req.url, true);
|
||||
const pathname = parsedUrl.pathname;
|
||||
|
||||
try {
|
||||
if (pathname.startsWith("/track/open/")) {
|
||||
await this.handleOpenTracking(req, res, parsedUrl);
|
||||
} else if (pathname.startsWith("/track/click/")) {
|
||||
await this.handleClickTracking(req, res, parsedUrl);
|
||||
} else if (pathname === "/health") {
|
||||
this.handleHealthCheck(req, res);
|
||||
} else {
|
||||
this.handle404(req, res);
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error("Tracking request error", {
|
||||
error: error.message,
|
||||
url: req.url,
|
||||
userAgent: req.headers["user-agent"],
|
||||
});
|
||||
this.handle500(req, res);
|
||||
}
|
||||
}
|
||||
|
||||
async handleOpenTracking(req, res, parsedUrl) {
|
||||
const pathParts = parsedUrl.pathname.split("/");
|
||||
const trackingId = pathParts[3]; // /track/open/{trackingId}
|
||||
|
||||
if (!trackingId) {
|
||||
return this.handle400(req, res, "Missing tracking ID");
|
||||
}
|
||||
|
||||
// Log the email open
|
||||
const trackingData = {
|
||||
trackingId,
|
||||
event: "email_open",
|
||||
timestamp: new Date().toISOString(),
|
||||
ip: req.connection.remoteAddress || req.headers["x-forwarded-for"],
|
||||
userAgent: req.headers["user-agent"],
|
||||
referer: req.headers.referer,
|
||||
};
|
||||
|
||||
logger.info("Email opened", trackingData);
|
||||
|
||||
// Store in database
|
||||
try {
|
||||
await database.storeTrackingEvent(trackingId, "open", trackingData);
|
||||
} catch (error) {
|
||||
logger.error("Failed to store tracking event", {
|
||||
trackingId,
|
||||
event: "open",
|
||||
error: error.message,
|
||||
});
|
||||
}
|
||||
|
||||
// Return 1x1 transparent pixel
|
||||
res.writeHead(200, {
|
||||
"Content-Type": "image/gif",
|
||||
"Content-Length": this.trackingPixel.length,
|
||||
"Cache-Control": "no-cache, no-store, must-revalidate",
|
||||
Pragma: "no-cache",
|
||||
Expires: "0",
|
||||
});
|
||||
res.end(this.trackingPixel);
|
||||
}
|
||||
|
||||
async handleClickTracking(req, res, parsedUrl) {
|
||||
const pathParts = parsedUrl.pathname.split("/");
|
||||
const trackingId = pathParts[3]; // /track/click/{trackingId}
|
||||
const linkId = pathParts[4]; // /track/click/{trackingId}/{linkId}
|
||||
const targetUrl = parsedUrl.query.url;
|
||||
|
||||
if (!trackingId || !targetUrl) {
|
||||
return this.handle400(req, res, "Missing tracking ID or target URL");
|
||||
}
|
||||
|
||||
// Log the click
|
||||
const trackingData = {
|
||||
trackingId,
|
||||
linkId,
|
||||
event: "email_click",
|
||||
targetUrl,
|
||||
timestamp: new Date().toISOString(),
|
||||
ip: req.connection.remoteAddress || req.headers["x-forwarded-for"],
|
||||
userAgent: req.headers["user-agent"],
|
||||
referer: req.headers.referer,
|
||||
};
|
||||
|
||||
logger.info("Email link clicked", trackingData);
|
||||
|
||||
// Store in database
|
||||
try {
|
||||
await database.storeTrackingEvent(trackingId, "click", trackingData);
|
||||
} catch (error) {
|
||||
logger.error("Failed to store tracking event", {
|
||||
trackingId,
|
||||
event: "click",
|
||||
error: error.message,
|
||||
});
|
||||
}
|
||||
|
||||
// Redirect to target URL
|
||||
res.writeHead(302, {
|
||||
Location: targetUrl,
|
||||
"Cache-Control": "no-cache",
|
||||
});
|
||||
res.end();
|
||||
}
|
||||
|
||||
handleHealthCheck(req, res) {
|
||||
res.writeHead(200, { "Content-Type": "application/json" });
|
||||
res.end(
|
||||
JSON.stringify({
|
||||
status: "healthy",
|
||||
timestamp: new Date().toISOString(),
|
||||
uptime: process.uptime(),
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
handle400(req, res, message) {
|
||||
res.writeHead(400, { "Content-Type": "text/plain" });
|
||||
res.end(`Bad Request: ${message}`);
|
||||
}
|
||||
|
||||
handle404(req, res) {
|
||||
res.writeHead(404, { "Content-Type": "text/plain" });
|
||||
res.end("Not Found");
|
||||
}
|
||||
|
||||
handle500(req, res) {
|
||||
res.writeHead(500, { "Content-Type": "text/plain" });
|
||||
res.end("Internal Server Error");
|
||||
}
|
||||
|
||||
// Generate tracking URLs
|
||||
generateOpenTrackingUrl(emailId) {
|
||||
return `${this.trackingDomain}/track/open/${emailId}`;
|
||||
}
|
||||
|
||||
generateClickTrackingUrl(emailId, linkId, targetUrl) {
|
||||
const encodedUrl = encodeURIComponent(targetUrl);
|
||||
return `${this.trackingDomain}/track/click/${emailId}/${linkId}?url=${encodedUrl}`;
|
||||
}
|
||||
|
||||
// Add tracking to email content
|
||||
addTrackingToEmail(htmlContent, emailId) {
|
||||
if (!htmlContent) return htmlContent;
|
||||
|
||||
// Add tracking pixel just before closing body tag
|
||||
const trackingPixel = `<img src="${this.generateOpenTrackingUrl(
|
||||
emailId
|
||||
)}" width="1" height="1" style="display:none;" alt="" />`;
|
||||
|
||||
if (htmlContent.includes("</body>")) {
|
||||
return htmlContent.replace("</body>", `${trackingPixel}</body>`);
|
||||
} else {
|
||||
// If no body tag, append at the end
|
||||
return htmlContent + trackingPixel;
|
||||
}
|
||||
}
|
||||
|
||||
// Replace links with tracking URLs
|
||||
addClickTrackingToEmail(htmlContent, emailId) {
|
||||
if (!htmlContent) return htmlContent;
|
||||
|
||||
let linkId = 0;
|
||||
return htmlContent.replace(
|
||||
/<a\s+([^>]*href=["']([^"']+)["'][^>]*)>/gi,
|
||||
(match, attributes, href) => {
|
||||
linkId++;
|
||||
|
||||
// Skip if already a tracking URL or mailto/tel links
|
||||
if (
|
||||
href.includes("/track/click/") ||
|
||||
href.startsWith("mailto:") ||
|
||||
href.startsWith("tel:")
|
||||
) {
|
||||
return match;
|
||||
}
|
||||
|
||||
const trackingUrl = this.generateClickTrackingUrl(
|
||||
emailId,
|
||||
linkId,
|
||||
href
|
||||
);
|
||||
return `<a ${attributes.replace(
|
||||
/href=["'][^"']+["']/i,
|
||||
`href="${trackingUrl}"`
|
||||
)}`;
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
// Combined method to add all tracking
|
||||
addEmailTracking(htmlContent, emailId, skipTracking = false) {
|
||||
if (!htmlContent) return htmlContent;
|
||||
|
||||
// If skip tracking is enabled, return original content without tracking
|
||||
if (skipTracking) {
|
||||
return htmlContent;
|
||||
}
|
||||
|
||||
let trackedContent = this.addClickTrackingToEmail(htmlContent, emailId);
|
||||
trackedContent = this.addTrackingToEmail(trackedContent, emailId);
|
||||
|
||||
return trackedContent;
|
||||
}
|
||||
}
|
||||
|
||||
module.exports = new TrackingServer();
|
||||
Reference in New Issue
Block a user