chore: replace standard logging with signale
This commit is contained in:
+1
-1
@@ -309,7 +309,7 @@ server.app.use((error: Error, req: Request, res: Response, _next: NextFunction)
|
|||||||
// Global error handlers to prevent server crashes
|
// Global error handlers to prevent server crashes
|
||||||
process.on('unhandledRejection', (reason, promise) => {
|
process.on('unhandledRejection', (reason, promise) => {
|
||||||
signale.error('Unhandled Promise Rejection:', reason);
|
signale.error('Unhandled Promise Rejection:', reason);
|
||||||
console.error('Promise:', promise);
|
signale.error('Promise:', promise);
|
||||||
// Don't exit the process - just log the error
|
// Don't exit the process - just log the error
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import {Controller, Delete, Get, Middleware, Patch, Post} from '@overnightjs/core';
|
import {Controller, Delete, Get, Middleware, Patch, Post} from '@overnightjs/core';
|
||||||
import type {NextFunction, Request, Response} from 'express';
|
import type {NextFunction, Request, Response} from 'express';
|
||||||
import multer from 'multer';
|
import multer from 'multer';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import type {AuthResponse} from '../middleware/auth.js';
|
import type {AuthResponse} from '../middleware/auth.js';
|
||||||
import {requireAuth} from '../middleware/auth.js';
|
import {requireAuth} from '../middleware/auth.js';
|
||||||
@@ -63,7 +64,7 @@ export class Contacts {
|
|||||||
count: fieldsWithTypes.length,
|
count: fieldsWithTypes.length,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to get available fields:', error);
|
signale.error('[CONTACTS] Failed to get available fields:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get available fields',
|
error: error instanceof Error ? error.message : 'Failed to get available fields',
|
||||||
});
|
});
|
||||||
@@ -97,7 +98,7 @@ export class Contacts {
|
|||||||
limit,
|
limit,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to get field values:', error);
|
signale.error('[CONTACTS] Failed to get field values:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get field values',
|
error: error instanceof Error ? error.message : 'Failed to get field values',
|
||||||
});
|
});
|
||||||
@@ -288,7 +289,7 @@ export class Contacts {
|
|||||||
jobId: job.id,
|
jobId: job.id,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to queue import:', error);
|
signale.error('[CONTACTS] Failed to queue import:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to queue import',
|
error: error instanceof Error ? error.message : 'Failed to queue import',
|
||||||
});
|
});
|
||||||
@@ -318,7 +319,7 @@ export class Contacts {
|
|||||||
|
|
||||||
return res.status(200).json(status);
|
return res.status(200).json(status);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to get import status:', error);
|
signale.error('[CONTACTS] Failed to get import status:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get import status',
|
error: error instanceof Error ? error.message : 'Failed to get import status',
|
||||||
});
|
});
|
||||||
@@ -345,7 +346,7 @@ export class Contacts {
|
|||||||
const usage = await ContactService.getFieldUsage(auth.projectId!, field);
|
const usage = await ContactService.getFieldUsage(auth.projectId!, field);
|
||||||
return res.status(200).json(usage);
|
return res.status(200).json(usage);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to get field usage:', error);
|
signale.error('[CONTACTS] Failed to get field usage:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get field usage',
|
error: error instanceof Error ? error.message : 'Failed to get field usage',
|
||||||
});
|
});
|
||||||
@@ -372,7 +373,7 @@ export class Contacts {
|
|||||||
const result = await ContactService.deleteField(auth.projectId!, field);
|
const result = await ContactService.deleteField(auth.projectId!, field);
|
||||||
return res.status(200).json(result);
|
return res.status(200).json(result);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[CONTACTS] Failed to delete field:', error);
|
signale.error('[CONTACTS] Failed to delete field:', error);
|
||||||
return res.status(error instanceof Error && error.message.includes('Cannot delete') ? 400 : 500).json({
|
return res.status(error instanceof Error && error.message.includes('Cannot delete') ? 400 : 500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to delete field',
|
error: error instanceof Error ? error.message : 'Failed to delete field',
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
import {Controller, Delete, Get, Middleware, Post} from '@overnightjs/core';
|
import {Controller, Delete, Get, Middleware, Post} from '@overnightjs/core';
|
||||||
import type {NextFunction, Request, Response} from 'express';
|
import type {NextFunction, Request, Response} from 'express';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import type {AuthResponse} from '../middleware/auth.js';
|
import type {AuthResponse} from '../middleware/auth.js';
|
||||||
import {requireAuth} from '../middleware/auth.js';
|
import {requireAuth} from '../middleware/auth.js';
|
||||||
@@ -118,7 +119,7 @@ export class Events {
|
|||||||
const usage = await EventService.getEventUsage(auth.projectId!, eventName);
|
const usage = await EventService.getEventUsage(auth.projectId!, eventName);
|
||||||
return res.status(200).json(usage);
|
return res.status(200).json(usage);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[EVENTS] Failed to get event usage:', error);
|
signale.error('[EVENTS] Failed to get event usage:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get event usage',
|
error: error instanceof Error ? error.message : 'Failed to get event usage',
|
||||||
});
|
});
|
||||||
@@ -145,7 +146,7 @@ export class Events {
|
|||||||
const result = await EventService.deleteEvent(auth.projectId!, eventName);
|
const result = await EventService.deleteEvent(auth.projectId!, eventName);
|
||||||
return res.status(200).json(result);
|
return res.status(200).json(result);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[EVENTS] Failed to delete event:', error);
|
signale.error('[EVENTS] Failed to delete event:', error);
|
||||||
return res.status(error instanceof Error && error.message.includes('Cannot delete') ? 400 : 500).json({
|
return res.status(error instanceof Error && error.message.includes('Cannot delete') ? 400 : 500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to delete event',
|
error: error instanceof Error ? error.message : 'Failed to delete event',
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import {Controller, Middleware, Post} from '@overnightjs/core';
|
import {Controller, Middleware, Post} from '@overnightjs/core';
|
||||||
import type {NextFunction, Request, Response} from 'express';
|
import type {NextFunction, Request, Response} from 'express';
|
||||||
import multer from 'multer';
|
import multer from 'multer';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import type {AuthResponse} from '../middleware/auth.js';
|
import type {AuthResponse} from '../middleware/auth.js';
|
||||||
import {requireAuth} from '../middleware/auth.js';
|
import {requireAuth} from '../middleware/auth.js';
|
||||||
@@ -66,7 +67,7 @@ export class Uploads {
|
|||||||
size: req.file.size,
|
size: req.file.size,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[UPLOADS] Failed to upload image:', error);
|
signale.error('[UPLOADS] Failed to upload image:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to upload image',
|
error: error instanceof Error ? error.message : 'Failed to upload image',
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import {Controller, Delete, Get, Middleware, Patch, Post} from '@overnightjs/core';
|
import {Controller, Delete, Get, Middleware, Patch, Post} from '@overnightjs/core';
|
||||||
import {WorkflowExecutionStatus} from '@plunk/db';
|
import {WorkflowExecutionStatus} from '@plunk/db';
|
||||||
import type {NextFunction, Request, Response} from 'express';
|
import type {NextFunction, Request, Response} from 'express';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import type {AuthResponse} from '../middleware/auth.js';
|
import type {AuthResponse} from '../middleware/auth.js';
|
||||||
import {requireAuth} from '../middleware/auth.js';
|
import {requireAuth} from '../middleware/auth.js';
|
||||||
@@ -45,7 +46,7 @@ export class Workflows {
|
|||||||
|
|
||||||
return res.status(200).json(result);
|
return res.status(200).json(result);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[WORKFLOWS] Failed to get available fields:', error);
|
signale.error('[WORKFLOWS] Failed to get available fields:', error);
|
||||||
return res.status(500).json({
|
return res.status(500).json({
|
||||||
error: error instanceof Error ? error.message : 'Failed to get available fields',
|
error: error instanceof Error ? error.message : 'Failed to get available fields',
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -4,6 +4,7 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import {type Job, Worker} from 'bullmq';
|
import {type Job, Worker} from 'bullmq';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {CampaignService} from '../services/CampaignService.js';
|
import {CampaignService} from '../services/CampaignService.js';
|
||||||
import {type CampaignBatchJobData, campaignQueue} from '../services/QueueService.js';
|
import {type CampaignBatchJobData, campaignQueue} from '../services/QueueService.js';
|
||||||
@@ -14,11 +15,11 @@ export function createCampaignWorker() {
|
|||||||
async (job: Job<CampaignBatchJobData>) => {
|
async (job: Job<CampaignBatchJobData>) => {
|
||||||
const {campaignId, batchNumber, offset, limit, cursor} = job.data;
|
const {campaignId, batchNumber, offset, limit, cursor} = job.data;
|
||||||
|
|
||||||
console.log(`[CAMPAIGN-PROCESSOR] Processing batch ${batchNumber} for campaign ${campaignId}`);
|
signale.info(`[CAMPAIGN-PROCESSOR] Processing batch ${batchNumber} for campaign ${campaignId}`);
|
||||||
|
|
||||||
await CampaignService.processBatch(campaignId, batchNumber, offset, limit, cursor);
|
await CampaignService.processBatch(campaignId, batchNumber, offset, limit, cursor);
|
||||||
|
|
||||||
console.log(`[CAMPAIGN-PROCESSOR] Completed batch ${batchNumber} for campaign ${campaignId}`);
|
signale.info(`[CAMPAIGN-PROCESSOR] Completed batch ${batchNumber} for campaign ${campaignId}`);
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
connection: campaignQueue.opts.connection,
|
connection: campaignQueue.opts.connection,
|
||||||
@@ -27,15 +28,15 @@ export function createCampaignWorker() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
worker.on('completed', job => {
|
worker.on('completed', job => {
|
||||||
console.log(`[CAMPAIGN-PROCESSOR] Job ${job.id} completed`);
|
signale.info(`[CAMPAIGN-PROCESSOR] Job ${job.id} completed`);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('failed', (job, err) => {
|
worker.on('failed', (job, err) => {
|
||||||
console.error(`[CAMPAIGN-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
signale.error(`[CAMPAIGN-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('error', err => {
|
worker.on('error', err => {
|
||||||
console.error('[CAMPAIGN-PROCESSOR] Worker error:', err);
|
signale.error('[CAMPAIGN-PROCESSOR] Worker error:', err);
|
||||||
});
|
});
|
||||||
|
|
||||||
return worker;
|
return worker;
|
||||||
|
|||||||
@@ -5,6 +5,7 @@
|
|||||||
|
|
||||||
import {EmailSourceType, EmailStatus} from '@plunk/db';
|
import {EmailSourceType, EmailStatus} from '@plunk/db';
|
||||||
import {type Job, Worker} from 'bullmq';
|
import {type Job, Worker} from 'bullmq';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {EmailService} from '../services/EmailService.js';
|
import {EmailService} from '../services/EmailService.js';
|
||||||
@@ -23,23 +24,23 @@ async function getEmailRateLimit(): Promise<number> {
|
|||||||
|
|
||||||
// If env variable is set, use it (override)
|
// If env variable is set, use it (override)
|
||||||
if (EMAIL_RATE_LIMIT_PER_SECOND !== undefined) {
|
if (EMAIL_RATE_LIMIT_PER_SECOND !== undefined) {
|
||||||
console.log(`[EMAIL-PROCESSOR] Using rate limit from environment: ${EMAIL_RATE_LIMIT_PER_SECOND} emails/second`);
|
signale.info(`[EMAIL-PROCESSOR] Using rate limit from environment: ${EMAIL_RATE_LIMIT_PER_SECOND} emails/second`);
|
||||||
return EMAIL_RATE_LIMIT_PER_SECOND;
|
return EMAIL_RATE_LIMIT_PER_SECOND;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Try to fetch from AWS SES
|
// Try to fetch from AWS SES
|
||||||
console.log('[EMAIL-PROCESSOR] Fetching rate limit from AWS SES...');
|
signale.info('[EMAIL-PROCESSOR] Fetching rate limit from AWS SES...');
|
||||||
const quota = await getSendingQuota();
|
const quota = await getSendingQuota();
|
||||||
|
|
||||||
if (quota) {
|
if (quota) {
|
||||||
console.log(
|
signale.info(
|
||||||
`[EMAIL-PROCESSOR] AWS SES quota: ${quota.maxSendRate} emails/second (${quota.sentLast24Hours}/${quota.max24HourSend} emails sent today)`,
|
`[EMAIL-PROCESSOR] AWS SES quota: ${quota.maxSendRate} emails/second (${quota.sentLast24Hours}/${quota.max24HourSend} emails sent today)`,
|
||||||
);
|
);
|
||||||
return quota.maxSendRate;
|
return quota.maxSendRate;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Fallback to safe default
|
// Fallback to safe default
|
||||||
console.warn(`[EMAIL-PROCESSOR] Failed to fetch AWS quota, using safe default: ${DEFAULT_RATE_LIMIT} emails/second`);
|
signale.warn(`[EMAIL-PROCESSOR] Failed to fetch AWS quota, using safe default: ${DEFAULT_RATE_LIMIT} emails/second`);
|
||||||
return DEFAULT_RATE_LIMIT;
|
return DEFAULT_RATE_LIMIT;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -69,7 +70,7 @@ export async function createEmailWorker() {
|
|||||||
|
|
||||||
// Check if project is disabled
|
// Check if project is disabled
|
||||||
if (email.project.disabled) {
|
if (email.project.disabled) {
|
||||||
console.warn(`[EMAIL-PROCESSOR] Project ${email.projectId} is disabled, cancelling email ${emailId}`);
|
signale.warn(`[EMAIL-PROCESSOR] Project ${email.projectId} is disabled, cancelling email ${emailId}`);
|
||||||
await prisma.email.update({
|
await prisma.email.update({
|
||||||
where: {id: emailId},
|
where: {id: emailId},
|
||||||
data: {
|
data: {
|
||||||
@@ -170,7 +171,7 @@ export async function createEmailWorker() {
|
|||||||
sentAt: new Date().toISOString(),
|
sentAt: new Date().toISOString(),
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[EMAIL-PROCESSOR] Failed to send email ${emailId}:`, error);
|
signale.error(`[EMAIL-PROCESSOR] Failed to send email ${emailId}:`, error);
|
||||||
|
|
||||||
// Mark as failed
|
// Mark as failed
|
||||||
await prisma.email.update({
|
await prisma.email.update({
|
||||||
@@ -195,15 +196,15 @@ export async function createEmailWorker() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
worker.on('completed', job => {
|
worker.on('completed', job => {
|
||||||
console.log(`[EMAIL-PROCESSOR] Job ${job.id} completed`);
|
signale.info(`[EMAIL-PROCESSOR] Job ${job.id} completed`);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('failed', (job, err) => {
|
worker.on('failed', (job, err) => {
|
||||||
console.error(`[EMAIL-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
signale.error(`[EMAIL-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('error', err => {
|
worker.on('error', err => {
|
||||||
console.error('[EMAIL-PROCESSOR] Worker error:', err);
|
signale.error('[EMAIL-PROCESSOR] Worker error:', err);
|
||||||
});
|
});
|
||||||
|
|
||||||
return worker;
|
return worker;
|
||||||
|
|||||||
@@ -5,6 +5,7 @@
|
|||||||
|
|
||||||
import {type Job, Worker} from 'bullmq';
|
import {type Job, Worker} from 'bullmq';
|
||||||
import {parse} from 'csv-parse/sync';
|
import {parse} from 'csv-parse/sync';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {ContactService} from '../services/ContactService.js';
|
import {ContactService} from '../services/ContactService.js';
|
||||||
@@ -28,7 +29,7 @@ export function createImportWorker() {
|
|||||||
async (job: Job<ContactImportJobData>) => {
|
async (job: Job<ContactImportJobData>) => {
|
||||||
const {projectId, csvData, filename} = job.data;
|
const {projectId, csvData, filename} = job.data;
|
||||||
|
|
||||||
console.log(`[IMPORT-PROCESSOR] Processing import for project ${projectId} (${filename})`);
|
signale.info(`[IMPORT-PROCESSOR] Processing import for project ${projectId} (${filename})`);
|
||||||
|
|
||||||
// Fetch project information for notifications
|
// Fetch project information for notifications
|
||||||
const project = await prisma.project.findUnique({
|
const project = await prisma.project.findUnique({
|
||||||
@@ -66,7 +67,7 @@ export function createImportWorker() {
|
|||||||
throw new Error('CSV file is empty');
|
throw new Error('CSV file is empty');
|
||||||
}
|
}
|
||||||
|
|
||||||
console.log(`[IMPORT-PROCESSOR] Parsed ${records.length} rows from CSV`);
|
signale.info(`[IMPORT-PROCESSOR] Parsed ${records.length} rows from CSV`);
|
||||||
|
|
||||||
// Notify that import has started
|
// Notify that import has started
|
||||||
await NtfyService.notifyContactImportStarted(projectName, projectId, filename, result.totalRows);
|
await NtfyService.notifyContactImportStarted(projectName, projectId, filename, result.totalRows);
|
||||||
@@ -151,7 +152,7 @@ export function createImportWorker() {
|
|||||||
await job.updateProgress(progress);
|
await job.updateProgress(progress);
|
||||||
}
|
}
|
||||||
|
|
||||||
console.log(
|
signale.info(
|
||||||
`[IMPORT-PROCESSOR] Import completed: ${result.createdCount} created, ${result.updatedCount} updated, ${result.failureCount} failed`,
|
`[IMPORT-PROCESSOR] Import completed: ${result.createdCount} created, ${result.updatedCount} updated, ${result.failureCount} failed`,
|
||||||
);
|
);
|
||||||
|
|
||||||
@@ -168,7 +169,7 @@ export function createImportWorker() {
|
|||||||
|
|
||||||
return result;
|
return result;
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[IMPORT-PROCESSOR] Failed to process import:`, error);
|
signale.error(`[IMPORT-PROCESSOR] Failed to process import:`, error);
|
||||||
|
|
||||||
// Notify that import has failed
|
// Notify that import has failed
|
||||||
const errorMessage = error instanceof Error ? error.message : 'Unknown error';
|
const errorMessage = error instanceof Error ? error.message : 'Unknown error';
|
||||||
@@ -191,15 +192,15 @@ export function createImportWorker() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
worker.on('completed', job => {
|
worker.on('completed', job => {
|
||||||
console.log(`[IMPORT-PROCESSOR] Job ${job.id} completed`);
|
signale.info(`[IMPORT-PROCESSOR] Job ${job.id} completed`);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('failed', (job, err) => {
|
worker.on('failed', (job, err) => {
|
||||||
console.error(`[IMPORT-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
signale.error(`[IMPORT-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('error', err => {
|
worker.on('error', err => {
|
||||||
console.error('[IMPORT-PROCESSOR] Worker error:', err);
|
signale.error('[IMPORT-PROCESSOR] Worker error:', err);
|
||||||
});
|
});
|
||||||
|
|
||||||
return worker;
|
return worker;
|
||||||
|
|||||||
@@ -5,6 +5,7 @@
|
|||||||
|
|
||||||
import {CampaignStatus} from '@plunk/db';
|
import {CampaignStatus} from '@plunk/db';
|
||||||
import {type Job, Worker} from 'bullmq';
|
import {type Job, Worker} from 'bullmq';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {CampaignService} from '../services/CampaignService.js';
|
import {CampaignService} from '../services/CampaignService.js';
|
||||||
@@ -16,7 +17,7 @@ export function createScheduledCampaignWorker() {
|
|||||||
async (job: Job<ScheduledCampaignJobData>) => {
|
async (job: Job<ScheduledCampaignJobData>) => {
|
||||||
const {campaignId} = job.data;
|
const {campaignId} = job.data;
|
||||||
|
|
||||||
console.log(`[SCHEDULED-PROCESSOR] Processing scheduled campaign ${campaignId}`);
|
signale.info(`[SCHEDULED-PROCESSOR] Processing scheduled campaign ${campaignId}`);
|
||||||
|
|
||||||
// Get campaign with project
|
// Get campaign with project
|
||||||
const campaign = await prisma.campaign.findUnique({
|
const campaign = await prisma.campaign.findUnique({
|
||||||
@@ -29,13 +30,13 @@ export function createScheduledCampaignWorker() {
|
|||||||
});
|
});
|
||||||
|
|
||||||
if (!campaign) {
|
if (!campaign) {
|
||||||
console.warn(`[SCHEDULED-PROCESSOR] Campaign ${campaignId} not found, skipping`);
|
signale.warn(`[SCHEDULED-PROCESSOR] Campaign ${campaignId} not found, skipping`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if project is disabled
|
// Check if project is disabled
|
||||||
if (campaign.project.disabled) {
|
if (campaign.project.disabled) {
|
||||||
console.warn(
|
signale.warn(
|
||||||
`[SCHEDULED-PROCESSOR] Project ${campaign.projectId} (${campaign.project.name}) is disabled, cancelling campaign ${campaignId}`,
|
`[SCHEDULED-PROCESSOR] Project ${campaign.projectId} (${campaign.project.name}) is disabled, cancelling campaign ${campaignId}`,
|
||||||
);
|
);
|
||||||
await prisma.campaign.update({
|
await prisma.campaign.update({
|
||||||
@@ -47,7 +48,7 @@ export function createScheduledCampaignWorker() {
|
|||||||
|
|
||||||
// Verify campaign is still in SCHEDULED status
|
// Verify campaign is still in SCHEDULED status
|
||||||
if (campaign.status !== CampaignStatus.SCHEDULED) {
|
if (campaign.status !== CampaignStatus.SCHEDULED) {
|
||||||
console.warn(
|
signale.warn(
|
||||||
`[SCHEDULED-PROCESSOR] Campaign ${campaignId} is not in SCHEDULED status (${campaign.status}), skipping`,
|
`[SCHEDULED-PROCESSOR] Campaign ${campaignId} is not in SCHEDULED status (${campaign.status}), skipping`,
|
||||||
);
|
);
|
||||||
return;
|
return;
|
||||||
@@ -56,7 +57,7 @@ export function createScheduledCampaignWorker() {
|
|||||||
// Start sending the campaign
|
// Start sending the campaign
|
||||||
await CampaignService.startSending(campaign.projectId, campaignId);
|
await CampaignService.startSending(campaign.projectId, campaignId);
|
||||||
|
|
||||||
console.log(`[SCHEDULED-PROCESSOR] Started sending campaign ${campaignId}`);
|
signale.info(`[SCHEDULED-PROCESSOR] Started sending campaign ${campaignId}`);
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
connection: scheduledQueue.opts.connection,
|
connection: scheduledQueue.opts.connection,
|
||||||
@@ -65,15 +66,15 @@ export function createScheduledCampaignWorker() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
worker.on('completed', job => {
|
worker.on('completed', job => {
|
||||||
console.log(`[SCHEDULED-PROCESSOR] Job ${job.id} completed`);
|
signale.info(`[SCHEDULED-PROCESSOR] Job ${job.id} completed`);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('failed', (job, err) => {
|
worker.on('failed', (job, err) => {
|
||||||
console.error(`[SCHEDULED-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
signale.error(`[SCHEDULED-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('error', err => {
|
worker.on('error', err => {
|
||||||
console.error('[SCHEDULED-PROCESSOR] Worker error:', err);
|
signale.error('[SCHEDULED-PROCESSOR] Worker error:', err);
|
||||||
});
|
});
|
||||||
|
|
||||||
return worker;
|
return worker;
|
||||||
|
|||||||
@@ -4,6 +4,7 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
import {type Job, Worker} from 'bullmq';
|
import {type Job, Worker} from 'bullmq';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {workflowQueue, type WorkflowStepJobData} from '../services/QueueService.js';
|
import {workflowQueue, type WorkflowStepJobData} from '../services/QueueService.js';
|
||||||
import {WorkflowExecutionService} from '../services/WorkflowExecutionService.js';
|
import {WorkflowExecutionService} from '../services/WorkflowExecutionService.js';
|
||||||
@@ -33,15 +34,15 @@ export function createWorkflowWorker() {
|
|||||||
);
|
);
|
||||||
|
|
||||||
worker.on('completed', job => {
|
worker.on('completed', job => {
|
||||||
console.log(`[WORKFLOW-PROCESSOR] Job ${job.id} completed`);
|
signale.info(`[WORKFLOW-PROCESSOR] Job ${job.id} completed`);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('failed', (job, err) => {
|
worker.on('failed', (job, err) => {
|
||||||
console.error(`[WORKFLOW-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
signale.error(`[WORKFLOW-PROCESSOR] Job ${job?.id} failed:`, err.message);
|
||||||
});
|
});
|
||||||
|
|
||||||
worker.on('error', err => {
|
worker.on('error', err => {
|
||||||
console.error('[WORKFLOW-PROCESSOR] Worker error:', err);
|
signale.error('[WORKFLOW-PROCESSOR] Worker error:', err);
|
||||||
});
|
});
|
||||||
|
|
||||||
return worker;
|
return worker;
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import type {Prisma} from '@plunk/db';
|
import type {Prisma} from '@plunk/db';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {redis} from '../database/redis.js';
|
import {redis} from '../database/redis.js';
|
||||||
@@ -178,7 +179,7 @@ export class ActivityService {
|
|||||||
return JSON.parse(cached);
|
return JSON.parse(cached);
|
||||||
}
|
}
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[ACTIVITY] Failed to get stats from cache:', error);
|
signale.warn('[ACTIVITY] Failed to get stats from cache:', error);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Default date range to last 30 days if not specified
|
// Default date range to last 30 days if not specified
|
||||||
@@ -241,7 +242,7 @@ export class ActivityService {
|
|||||||
try {
|
try {
|
||||||
await redis.setex(cacheKey, this.STATS_CACHE_TTL, JSON.stringify(stats));
|
await redis.setex(cacheKey, this.STATS_CACHE_TTL, JSON.stringify(stats));
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[ACTIVITY] Failed to cache stats:', error);
|
signale.warn('[ACTIVITY] Failed to cache stats:', error);
|
||||||
}
|
}
|
||||||
|
|
||||||
return stats;
|
return stats;
|
||||||
@@ -261,7 +262,7 @@ export class ActivityService {
|
|||||||
await redis.del(...keys);
|
await redis.del(...keys);
|
||||||
}
|
}
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[ACTIVITY] Failed to invalidate stats cache:', error);
|
signale.warn('[ACTIVITY] Failed to invalidate stats cache:', error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import type {Campaign, Contact, Prisma} from '@plunk/db';
|
import type {Campaign, Contact, Prisma} from '@plunk/db';
|
||||||
import {CampaignAudienceType, CampaignStatus, EmailSourceType} from '@plunk/db';
|
import {CampaignAudienceType, CampaignStatus, EmailSourceType} from '@plunk/db';
|
||||||
import type {FilterCondition} from '@plunk/types';
|
import type {FilterCondition} from '@plunk/types';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {HttpException} from '../exceptions/index.js';
|
import {HttpException} from '../exceptions/index.js';
|
||||||
@@ -461,7 +462,7 @@ export class CampaignService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if (campaign.status !== CampaignStatus.SENDING) {
|
if (campaign.status !== CampaignStatus.SENDING) {
|
||||||
console.warn(`[CAMPAIGN] Campaign ${campaignId} is not in SENDING status, skipping batch ${batchNumber}`);
|
signale.warn(`[CAMPAIGN] Campaign ${campaignId} is not in SENDING status, skipping batch ${batchNumber}`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -507,7 +508,7 @@ export class CampaignService {
|
|||||||
replyTo: campaign.replyTo || undefined,
|
replyTo: campaign.replyTo || undefined,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[CAMPAIGN] Failed to queue email for contact ${contact.id}:`, error);
|
signale.error(`[CAMPAIGN] Failed to queue email for contact ${contact.id}:`, error);
|
||||||
// Continue with other contacts even if one fails
|
// Continue with other contacts even if one fails
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -407,7 +407,7 @@ export class EmailService {
|
|||||||
sentAt: new Date().toISOString(),
|
sentAt: new Date().toISOString(),
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[EMAIL] Failed to send email ${emailId}:`, error);
|
signale.error(`[EMAIL] Failed to send email ${emailId}:`, error);
|
||||||
|
|
||||||
// Mark as failed
|
// Mark as failed
|
||||||
await prisma.email.update({
|
await prisma.email.update({
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import type {Event} from '@plunk/db';
|
import type {Event} from '@plunk/db';
|
||||||
import {Prisma} from '@plunk/db';
|
import {Prisma} from '@plunk/db';
|
||||||
import type {FilterCondition, FilterGroup} from '@plunk/types';
|
import type {FilterCondition, FilterGroup} from '@plunk/types';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {redis} from '../database/redis.js';
|
import {redis} from '../database/redis.js';
|
||||||
@@ -52,7 +53,7 @@ export class EventService {
|
|||||||
try {
|
try {
|
||||||
await redis.del(cacheKey);
|
await redis.del(cacheKey);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[EVENT] Failed to invalidate workflow cache:', error);
|
signale.warn('[EVENT] Failed to invalidate workflow cache:', error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -331,7 +332,7 @@ export class EventService {
|
|||||||
workflows = JSON.parse(cached);
|
workflows = JSON.parse(cached);
|
||||||
}
|
}
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[EVENT] Failed to get workflows from cache:', error);
|
signale.warn('[EVENT] Failed to get workflows from cache:', error);
|
||||||
}
|
}
|
||||||
|
|
||||||
// If not in cache, fetch from database
|
// If not in cache, fetch from database
|
||||||
@@ -353,7 +354,7 @@ export class EventService {
|
|||||||
try {
|
try {
|
||||||
await redis.setex(cacheKey, 300, JSON.stringify(workflows));
|
await redis.setex(cacheKey, 300, JSON.stringify(workflows));
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.warn('[EVENT] Failed to cache workflows:', error);
|
signale.warn('[EVENT] Failed to cache workflows:', error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -368,7 +369,7 @@ export class EventService {
|
|||||||
} else {
|
} else {
|
||||||
// If event is not contact-specific, you might want different logic
|
// If event is not contact-specific, you might want different logic
|
||||||
// For example, trigger for all contacts, or skip
|
// For example, trigger for all contacts, or skip
|
||||||
console.log(`[EVENT] Event ${eventName} triggered workflow ${workflow.id}, but no contact specified`);
|
signale.info(`[EVENT] Event ${eventName} triggered workflow ${workflow.id}, but no contact specified`);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -394,7 +395,7 @@ export class EventService {
|
|||||||
});
|
});
|
||||||
|
|
||||||
if (!workflow || workflow.steps.length === 0) {
|
if (!workflow || workflow.steps.length === 0) {
|
||||||
console.error(`[EVENT] Workflow ${workflowId} has no trigger step`);
|
signale.error(`[EVENT] Workflow ${workflowId} has no trigger step`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -429,7 +430,7 @@ export class EventService {
|
|||||||
const triggerStep = workflow.steps[0];
|
const triggerStep = workflow.steps[0];
|
||||||
|
|
||||||
if (!triggerStep) {
|
if (!triggerStep) {
|
||||||
console.error(`[EVENT] Workflow ${workflowId} trigger step not found`);
|
signale.error(`[EVENT] Workflow ${workflowId} trigger step not found`);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -444,14 +445,14 @@ export class EventService {
|
|||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|
||||||
console.log(
|
signale.info(
|
||||||
`[EVENT] Started workflow ${workflowId} execution ${execution.id} for contact ${contactId}${workflow.allowReentry ? ' (re-entry allowed)' : ''}`,
|
`[EVENT] Started workflow ${workflowId} execution ${execution.id} for contact ${contactId}${workflow.allowReentry ? ' (re-entry allowed)' : ''}`,
|
||||||
);
|
);
|
||||||
|
|
||||||
// Start executing the workflow
|
// Start executing the workflow
|
||||||
await WorkflowExecutionService.processStepExecution(execution.id, triggerStep.id);
|
await WorkflowExecutionService.processStepExecution(execution.id, triggerStep.id);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[EVENT] Error starting workflow ${workflowId}:`, error);
|
signale.error(`[EVENT] Error starting workflow ${workflowId}:`, error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -6,6 +6,8 @@ import {
|
|||||||
PutBucketPolicyCommand,
|
PutBucketPolicyCommand,
|
||||||
} from '@aws-sdk/client-s3';
|
} from '@aws-sdk/client-s3';
|
||||||
import crypto from 'crypto';
|
import crypto from 'crypto';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {
|
import {
|
||||||
S3_ENDPOINT,
|
S3_ENDPOINT,
|
||||||
S3_ACCESS_KEY_ID,
|
S3_ACCESS_KEY_ID,
|
||||||
@@ -69,13 +71,13 @@ export async function initializeBucket(): Promise<void> {
|
|||||||
Bucket: S3_BUCKET,
|
Bucket: S3_BUCKET,
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
console.log(`[S3] Created bucket: ${S3_BUCKET}`);
|
signale.info(`[S3] Created bucket: ${S3_BUCKET}`);
|
||||||
} catch (createError) {
|
} catch (createError) {
|
||||||
console.error('[S3] Failed to create bucket:', createError);
|
signale.error('[S3] Failed to create bucket:', createError);
|
||||||
throw createError;
|
throw createError;
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
console.error('[S3] Failed to check bucket:', error);
|
signale.error('[S3] Failed to check bucket:', error);
|
||||||
throw error;
|
throw error;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -103,10 +105,10 @@ export async function initializeBucket(): Promise<void> {
|
|||||||
);
|
);
|
||||||
|
|
||||||
if (!bucketExists) {
|
if (!bucketExists) {
|
||||||
console.log(`[S3] Set public read policy for bucket: ${S3_BUCKET}`);
|
signale.info(`[S3] Set public read policy for bucket: ${S3_BUCKET}`);
|
||||||
}
|
}
|
||||||
} catch (policyError) {
|
} catch (policyError) {
|
||||||
console.error('[S3] Failed to set bucket policy:', policyError);
|
signale.error('[S3] Failed to set bucket policy:', policyError);
|
||||||
// Don't throw - bucket was created but policy failed
|
// Don't throw - bucket was created but policy failed
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import {SES} from '@aws-sdk/client-ses';
|
import {SES} from '@aws-sdk/client-ses';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {
|
import {
|
||||||
AWS_SES_ACCESS_KEY_ID,
|
AWS_SES_ACCESS_KEY_ID,
|
||||||
@@ -275,7 +276,7 @@ export const getSendingQuota = async (): Promise<{
|
|||||||
sentLast24Hours: quota.SentLast24Hours ?? 0,
|
sentLast24Hours: quota.SentLast24Hours ?? 0,
|
||||||
};
|
};
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error('[SES] Failed to fetch sending quota:', error);
|
signale.error('[SES] Failed to fetch sending quota:', error);
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
import {type Contact, Prisma, type Segment} from '@plunk/db';
|
import {type Contact, Prisma, type Segment} from '@plunk/db';
|
||||||
import type {FilterCondition, FilterGroup, SegmentFilter} from '@plunk/types';
|
import type {FilterCondition, FilterGroup, SegmentFilter} from '@plunk/types';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {HttpException} from '../exceptions/index.js';
|
import {HttpException} from '../exceptions/index.js';
|
||||||
@@ -276,7 +277,7 @@ export class SegmentService {
|
|||||||
data: {memberCount},
|
data: {memberCount},
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`Failed to update count for segment ${segment.id}:`, error);
|
signale.error(`Failed to update count for segment ${segment.id}:`, error);
|
||||||
}
|
}
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
@@ -383,7 +384,7 @@ export class SegmentService {
|
|||||||
segmentName: segment.name,
|
segmentName: segment.name,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[SEGMENT] Failed to track segment entry event for contact ${contactId}:`, error);
|
signale.error(`[SEGMENT] Failed to track segment entry event for contact ${contactId}:`, error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -414,7 +415,7 @@ export class SegmentService {
|
|||||||
segmentName: segment.name,
|
segmentName: segment.name,
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`[SEGMENT] Failed to track segment exit event for contact ${contactId}:`, error);
|
signale.error(`[SEGMENT] Failed to track segment exit event for contact ${contactId}:`, error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -425,7 +426,7 @@ export class SegmentService {
|
|||||||
data: {memberCount: matchingContactIds.size},
|
data: {memberCount: matchingContactIds.size},
|
||||||
});
|
});
|
||||||
|
|
||||||
console.log(
|
signale.info(
|
||||||
`[SEGMENT] Computed membership for segment ${segmentId}: added ${toAdd.length}, removed ${toRemove.length}, total ${matchingContactIds.size}`,
|
`[SEGMENT] Computed membership for segment ${segmentId}: added ${toAdd.length}, removed ${toRemove.length}, total ${matchingContactIds.size}`,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
import type {Workflow, WorkflowExecution, WorkflowStep, WorkflowStepExecution, WorkflowTransition} from '@plunk/db';
|
import type {Workflow, WorkflowExecution, WorkflowStep, WorkflowStepExecution, WorkflowTransition} from '@plunk/db';
|
||||||
import {Prisma, WorkflowExecutionStatus} from '@plunk/db';
|
import {Prisma, WorkflowExecutionStatus} from '@plunk/db';
|
||||||
|
import signale from 'signale';
|
||||||
|
|
||||||
import {prisma} from '../database/prisma.js';
|
import {prisma} from '../database/prisma.js';
|
||||||
import {HttpException} from '../exceptions/index.js';
|
import {HttpException} from '../exceptions/index.js';
|
||||||
@@ -786,7 +787,7 @@ export class WorkflowService {
|
|||||||
// Start executing the workflow asynchronously
|
// Start executing the workflow asynchronously
|
||||||
// Don't await - let it run in background
|
// Don't await - let it run in background
|
||||||
WorkflowExecutionService.processStepExecution(execution.id, triggerStep.id).catch(error => {
|
WorkflowExecutionService.processStepExecution(execution.id, triggerStep.id).catch(error => {
|
||||||
console.error('Error executing workflow:', error);
|
signale.error('Error executing workflow:', error);
|
||||||
});
|
});
|
||||||
|
|
||||||
return execution;
|
return execution;
|
||||||
|
|||||||
Reference in New Issue
Block a user