Files
plunk/apps/api/src/controllers/Contacts.ts
T

516 lines
15 KiB
TypeScript

import {Controller, Delete, Get, Middleware, Patch, Post} from '@overnightjs/core';
import type {NextFunction, Request, Response} from 'express';
import multer from 'multer';
import signale from 'signale';
import type {AuthResponse} from '../middleware/auth.js';
import {requireAuth, requireEmailVerified} from '../middleware/auth.js';
import {ContactService} from '../services/ContactService.js';
import {QueueService} from '../services/QueueService.js';
import {CatchAsync} from '../utils/asyncHandler.js';
// Configure multer for file uploads (memory storage)
const upload = multer({
storage: multer.memoryStorage(),
limits: {
fileSize: 5 * 1024 * 1024, // 5MB max file size
},
fileFilter: (_req, file, cb) => {
// Only accept CSV files
if (file.mimetype === 'text/csv' || file.originalname.endsWith('.csv')) {
cb(null, true);
} else {
cb(new Error('Only CSV files are allowed'));
}
},
});
@Controller('contacts')
export class Contacts {
/**
* GET /contacts
* List all contacts for the authenticated project with cursor-based pagination
*/
@Get('')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async list(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const limit = Math.min(parseInt(req.query.limit as string) || 20, 100);
const cursor = req.query.cursor as string | undefined;
const search = req.query.search as string | undefined;
const result = await ContactService.list(auth.projectId!, limit, cursor, search);
return res.status(200).json(result);
}
/**
* GET /contacts/fields
* Get all available contact fields (both standard and custom fields from data JSON)
* Returns field names with inferred types (string, number, boolean, date)
*/
@Get('fields')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async getAvailableFields(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
try {
const fieldsWithTypes = await ContactService.getAvailableFields(auth.projectId!);
return res.status(200).json({
fields: fieldsWithTypes,
count: fieldsWithTypes.length,
});
} catch (error) {
signale.error('[CONTACTS] Failed to get available fields:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to get available fields',
});
}
}
/**
* GET /contacts/fields/:field/values
* Get unique values for a contact field (for workflow conditions, segment filters, etc.)
* Example: /contacts/fields/data.plan/values or /contacts/fields/subscribed/values
*/
@Get('fields/:field/values')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async getFieldValues(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const field = req.params.field;
const limit = Math.min(parseInt(req.query.limit as string) || 100, 200);
if (!field) {
return res.status(400).json({error: 'Field is required'});
}
try {
const values = await ContactService.getUniqueFieldValues(auth.projectId!, field, limit);
return res.status(200).json({
field,
values,
count: values.length,
limit,
});
} catch (error) {
signale.error('[CONTACTS] Failed to get field values:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to get field values',
});
}
}
/**
* GET /contacts/:id
* Get a specific contact by ID
*/
@Get(':id')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async get(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const contactId = req.params.id;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
const contact = await ContactService.get(auth.projectId!, contactId);
return res.status(200).json(contact);
}
/**
* POST /contacts
* Create or update a contact (upsert)
*/
@Post('')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async create(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const {email, data, subscribed} = req.body;
if (!email) {
return res.status(400).json({error: 'Email is required'});
}
// Check if contact exists before upserting
const existingContact = await ContactService.findByEmail(auth.projectId!, email);
const isUpdate = !!existingContact;
const contact = await ContactService.upsert(auth.projectId!, email, data, subscribed);
return res.status(isUpdate ? 200 : 201).json({
...contact,
_meta: {
isNew: !isUpdate,
isUpdate,
},
});
}
/**
* PATCH /contacts/:id
* Update a contact
*/
@Patch(':id')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async update(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const contactId = req.params.id;
const {email, data, subscribed} = req.body;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
const contact = await ContactService.update(auth.projectId!, contactId, {email, data, subscribed});
return res.status(200).json(contact);
}
/**
* DELETE /contacts/:id
* Delete a contact
*/
@Delete(':id')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async delete(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const contactId = req.params.id;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
await ContactService.delete(auth.projectId!, contactId);
return res.status(204).send();
}
/**
* GET /contacts/public/:id
* PUBLIC: Get contact information (no auth required)
*/
@Get('public/:id')
@CatchAsync
public async getPublic(req: Request, res: Response, _next: NextFunction) {
const contactId = req.params.id;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
const contact = await ContactService.getById(contactId);
return res.status(200).json({
id: contact.id,
email: contact.email,
subscribed: contact.subscribed,
});
}
/**
* POST /contacts/public/:id/subscribe
* PUBLIC: Subscribe a contact (no auth required)
*/
@Post('public/:id/subscribe')
@CatchAsync
public async subscribePublic(req: Request, res: Response, _next: NextFunction) {
const contactId = req.params.id;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
const contact = await ContactService.subscribe(contactId);
return res.status(200).json({
id: contact.id,
email: contact.email,
subscribed: contact.subscribed,
});
}
/**
* POST /contacts/public/:id/unsubscribe
* PUBLIC: Unsubscribe a contact (no auth required)
*/
@Post('public/:id/unsubscribe')
@CatchAsync
public async unsubscribePublic(req: Request, res: Response, _next: NextFunction) {
const contactId = req.params.id;
if (!contactId) {
return res.status(400).json({error: 'Contact ID is required'});
}
const contact = await ContactService.unsubscribe(contactId);
return res.status(200).json({
id: contact.id,
email: contact.email,
subscribed: contact.subscribed,
});
}
/**
* POST /contacts/import
* Import contacts from CSV file
*/
@Post('import')
@Middleware([requireAuth, upload.single('file')])
@CatchAsync
public async importCsv(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
if (!req.file) {
return res.status(400).json({error: 'CSV file is required'});
}
try {
// Convert file buffer to base64 for storage in queue
const csvData = req.file.buffer.toString('base64');
const filename = req.file.originalname;
// Queue import job
const job = await QueueService.queueImport(auth.projectId!, csvData, filename);
return res.status(202).json({
message: 'Import queued successfully',
jobId: job.id,
});
} catch (error) {
signale.error('[CONTACTS] Failed to queue import:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to queue import',
});
}
}
/**
* GET /contacts/import/:jobId
* Get import job status
*/
@Get('import/:jobId')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async getImportStatus(req: Request, res: Response, _next: NextFunction) {
const jobId = req.params.jobId;
if (!jobId) {
return res.status(400).json({error: 'Job ID is required'});
}
try {
const status = await QueueService.getImportJobStatus(jobId);
if (!status) {
return res.status(404).json({error: 'Import job not found'});
}
return res.status(200).json(status);
} catch (error) {
signale.error('[CONTACTS] Failed to get import status:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to get import status',
});
}
}
/**
* GET /contacts/fields/:field/usage
* Check if a field is used in segments/campaigns and get usage statistics
* Returns information about where the field is used and whether it can be safely deleted
*/
@Get('fields/:field/usage')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async getFieldUsage(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const field = req.params.field;
if (!field) {
return res.status(400).json({error: 'Field is required'});
}
try {
const usage = await ContactService.getFieldUsage(auth.projectId!, field);
return res.status(200).json(usage);
} catch (error) {
signale.error('[CONTACTS] Failed to get field usage:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to get field usage',
});
}
}
/**
* DELETE /contacts/fields/:field
* Delete a custom field from all contacts
* Only works if the field is not used in any segments or campaigns
*/
@Delete('fields/:field')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async deleteField(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const field = req.params.field;
if (!field) {
return res.status(400).json({error: 'Field is required'});
}
try {
const result = await ContactService.deleteField(auth.projectId!, field);
return res.status(200).json(result);
} catch (error) {
signale.error('[CONTACTS] Failed to delete field:', error);
return res.status(error instanceof Error && error.message.includes('Cannot delete') ? 400 : 500).json({
error: error instanceof Error ? error.message : 'Failed to delete field',
});
}
}
/**
* POST /contacts/bulk-subscribe
* Queue bulk subscribe operation
*/
@Post('bulk-subscribe')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async bulkSubscribe(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const {contactIds} = req.body;
if (!Array.isArray(contactIds) || contactIds.length === 0) {
return res.status(400).json({error: 'contactIds array is required'});
}
// Validate limit
if (contactIds.length > 1000) {
return res.status(400).json({error: 'Maximum 1000 contacts can be processed at once'});
}
try {
const job = await QueueService.queueBulkContactAction(auth.projectId!, contactIds, 'subscribe');
return res.status(202).json({
message: 'Bulk subscribe queued successfully',
jobId: job.id,
});
} catch (error) {
signale.error('[CONTACTS] Failed to queue bulk subscribe:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to queue bulk subscribe',
});
}
}
/**
* POST /contacts/bulk-unsubscribe
* Queue bulk unsubscribe operation
*/
@Post('bulk-unsubscribe')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async bulkUnsubscribe(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const {contactIds} = req.body;
if (!Array.isArray(contactIds) || contactIds.length === 0) {
return res.status(400).json({error: 'contactIds array is required'});
}
if (contactIds.length > 1000) {
return res.status(400).json({error: 'Maximum 1000 contacts can be processed at once'});
}
try {
const job = await QueueService.queueBulkContactAction(auth.projectId!, contactIds, 'unsubscribe');
return res.status(202).json({
message: 'Bulk unsubscribe queued successfully',
jobId: job.id,
});
} catch (error) {
signale.error('[CONTACTS] Failed to queue bulk unsubscribe:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to queue bulk unsubscribe',
});
}
}
/**
* POST /contacts/bulk-delete
* Queue bulk delete operation
*/
@Post('bulk-delete')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async bulkDelete(req: Request, res: Response, _next: NextFunction) {
const auth = res.locals.auth as AuthResponse;
const {contactIds} = req.body;
if (!Array.isArray(contactIds) || contactIds.length === 0) {
return res.status(400).json({error: 'contactIds array is required'});
}
if (contactIds.length > 1000) {
return res.status(400).json({error: 'Maximum 1000 contacts can be processed at once'});
}
try {
const job = await QueueService.queueBulkContactAction(auth.projectId!, contactIds, 'delete');
return res.status(202).json({
message: 'Bulk delete queued successfully',
jobId: job.id,
});
} catch (error) {
signale.error('[CONTACTS] Failed to queue bulk delete:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to queue bulk delete',
});
}
}
/**
* GET /contacts/bulk/:jobId
* Get bulk action job status
*/
@Get('bulk/:jobId')
@Middleware([requireAuth, requireEmailVerified])
@CatchAsync
public async getBulkActionStatus(req: Request, res: Response, _next: NextFunction) {
const jobId = req.params.jobId;
if (!jobId) {
return res.status(400).json({error: 'Job ID is required'});
}
try {
const status = await QueueService.getBulkActionJobStatus(jobId);
if (!status) {
return res.status(404).json({error: 'Bulk action job not found'});
}
return res.status(200).json(status);
} catch (error) {
signale.error('[CONTACTS] Failed to get bulk action status:', error);
return res.status(500).json({
error: error instanceof Error ? error.message : 'Failed to get bulk action status',
});
}
}
}