diff --git a/apps/api/src/db/models/suppression.model.ts b/apps/api/src/db/models/suppression.model.ts index 6c9901c1..4fa45c1f 100644 --- a/apps/api/src/db/models/suppression.model.ts +++ b/apps/api/src/db/models/suppression.model.ts @@ -1,9 +1,13 @@ import mongoose, { Schema } from 'mongoose'; import { workspacePlugin, type WorkspaceScopedDocument } from '../plugins/index.js'; -import { SuppressionReason } from '@leadforge/schema'; +import { SuppressionReason, SuppressionTargetType } from '@leadforge/schema'; export interface SuppressionDocument extends mongoose.Document, WorkspaceScopedDocument { - email: string; + targetType: SuppressionTargetType; + targetId: string; + email?: string | null; + companyId?: string | null; + domain?: string | null; reason: SuppressionReason; source: string; evidence?: Record | null; @@ -17,7 +21,16 @@ export interface SuppressionDocument extends mongoose.Document, WorkspaceScopedD const suppressionSchema = new Schema( { workspaceId: { type: String, required: true, index: true }, - email: { type: String, required: true, trim: true, lowercase: true }, + targetType: { + type: String, + enum: Object.values(SuppressionTargetType), + default: SuppressionTargetType.RECIPIENT, + required: true + }, + targetId: { type: String, required: true, trim: true }, + email: { type: String, default: null, trim: true, lowercase: true }, + companyId: { type: String, default: null, trim: true }, + domain: { type: String, default: null, trim: true, lowercase: true }, reason: { type: String, enum: Object.values(SuppressionReason), @@ -36,7 +49,10 @@ const suppressionSchema = new Schema( ); suppressionSchema.plugin(workspacePlugin); -suppressionSchema.index({ workspaceId: 1, email: 1 }, { unique: true }); +suppressionSchema.index({ workspaceId: 1, targetType: 1, targetId: 1 }, { unique: true }); +suppressionSchema.index({ workspaceId: 1, targetType: 1, companyId: 1 }, { sparse: true }); +suppressionSchema.index({ workspaceId: 1, targetType: 1, domain: 1 }, { sparse: true }); +suppressionSchema.index({ workspaceId: 1, email: 1 }, { sparse: true }); suppressionSchema.index({ workspaceId: 1, reason: 1 }); export const SuppressionModel = diff --git a/apps/api/src/repositories/suppression/suppression.repository.ts b/apps/api/src/repositories/suppression/suppression.repository.ts index 6fa82d47..1b06f1f4 100644 --- a/apps/api/src/repositories/suppression/suppression.repository.ts +++ b/apps/api/src/repositories/suppression/suppression.repository.ts @@ -3,30 +3,190 @@ import { type SuppressionDocument } from '../../db/models/suppression.model.js'; import { ContactModel } from '../../db/models/contact.model.js'; +import { SequenceExecutionModel } from '../../db/models/sequence-execution.model.js'; import { SuppressionReason, + SuppressionTargetType, compareSuppressionPrecedence, + normalizeDomain, ContactStatus, ContactEmailStatus } from '@leadforge/schema'; import { logger } from '../../config/index.js'; +export interface EffectiveSuppressionResult { + suppressed: boolean; + isRecipientSuppressed: boolean; + isCompanySuppressed: boolean; + isDomainSuppressed: boolean; + reasons: Array<{ + targetType: SuppressionTargetType; + targetId: string; + reason: SuppressionReason; + record: SuppressionDocument; + }>; + primaryReason?: SuppressionReason | undefined; + message?: string | undefined; +} + export class SuppressionRepository { constructor(private readonly workspaceId: string) {} /** - * Checks if an email is actively suppressed in the workspace. + * Checks if an email is actively suppressed in the workspace (recipient-level). */ public async isSuppressed(email: string): Promise { if (!email) return false; const cleanEmail = email.toLowerCase().trim(); const count = await SuppressionModel.countDocuments({ workspaceId: this.workspaceId, - email: cleanEmail + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: cleanEmail }, + { email: cleanEmail } + ] + }); + return count > 0; + } + + /** + * Checks if a company is marked Do Not Contact in the workspace. + */ + public async isCompanySuppressed(companyId: string): Promise { + if (!companyId) return false; + const cleanCompanyId = companyId.trim(); + const count = await SuppressionModel.countDocuments({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.COMPANY, + targetId: cleanCompanyId + }); + return count > 0; + } + + /** + * Checks if a domain is suppressed in the workspace. + * Normalizes the domain using canonical normalizeDomain() before lookup. + */ + public async isDomainSuppressed(domainOrEmail: string): Promise { + const normDomain = normalizeDomain(domainOrEmail); + if (!normDomain) return false; + const count = await SuppressionModel.countDocuments({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain }); return count > 0; } + /** + * Evaluates the effective additive suppression state across recipient, company, and domain. + */ + public async evaluateEffectiveSuppression(params: { + email: string; + companyId?: string | null; + }): Promise { + const cleanEmail = (params.email || '').toLowerCase().trim(); + const cleanCompanyId = params.companyId ? params.companyId.trim() : null; + const normDomain = normalizeDomain(cleanEmail); + + const conditions: any[] = []; + if (cleanEmail) { + conditions.push( + { targetType: SuppressionTargetType.RECIPIENT, targetId: cleanEmail }, + { email: cleanEmail } + ); + } + if (cleanCompanyId) { + conditions.push({ targetType: SuppressionTargetType.COMPANY, targetId: cleanCompanyId }); + } + if (normDomain) { + conditions.push({ targetType: SuppressionTargetType.DOMAIN, targetId: normDomain }); + } + + if (conditions.length === 0) { + return { + suppressed: false, + isRecipientSuppressed: false, + isCompanySuppressed: false, + isDomainSuppressed: false, + reasons: [] + }; + } + + const records = await SuppressionModel.find({ + workspaceId: this.workspaceId, + $or: conditions + }); + + if (records.length === 0) { + return { + suppressed: false, + isRecipientSuppressed: false, + isCompanySuppressed: false, + isDomainSuppressed: false, + reasons: [] + }; + } + + let isRecipientSuppressed = false; + let isCompanySuppressed = false; + let isDomainSuppressed = false; + const reasons: EffectiveSuppressionResult['reasons'] = []; + + for (const rec of records) { + const tType = rec.targetType || SuppressionTargetType.RECIPIENT; + if (tType === SuppressionTargetType.RECIPIENT || rec.email === cleanEmail) { + isRecipientSuppressed = true; + reasons.push({ + targetType: SuppressionTargetType.RECIPIENT, + targetId: cleanEmail, + reason: rec.reason, + record: rec + }); + } else if (tType === SuppressionTargetType.COMPANY && cleanCompanyId && rec.targetId === cleanCompanyId) { + isCompanySuppressed = true; + reasons.push({ + targetType: SuppressionTargetType.COMPANY, + targetId: cleanCompanyId, + reason: rec.reason, + record: rec + }); + } else if (tType === SuppressionTargetType.DOMAIN && normDomain && rec.targetId === normDomain) { + isDomainSuppressed = true; + reasons.push({ + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain, + reason: rec.reason, + record: rec + }); + } + } + + // Sort reasons by precedence weight descending to pick primary reason + reasons.sort((a, b) => compareSuppressionPrecedence(b.reason, a.reason)); + const primaryReason = reasons[0]?.reason; + + let message = 'Recipient is suppressed.'; + if (isCompanySuppressed && isDomainSuppressed && isRecipientSuppressed) { + message = `Blocked by recipient suppression, company DNC, and domain suppression (${normDomain}).`; + } else if (isCompanySuppressed) { + message = `Company "${cleanCompanyId}" is marked Do Not Contact in workspace.`; + } else if (isDomainSuppressed) { + message = `Domain "${normDomain}" is suppressed in workspace.`; + } else if (isRecipientSuppressed) { + message = `Recipient "${cleanEmail}" is suppressed in workspace.`; + } + + return { + suppressed: true, + isRecipientSuppressed, + isCompanySuppressed, + isDomainSuppressed, + reasons, + primaryReason, + message + }; + } + /** * Retrieves suppression details for a specific email address. */ @@ -35,7 +195,35 @@ export class SuppressionRepository { const cleanEmail = email.toLowerCase().trim(); return SuppressionModel.findOne({ workspaceId: this.workspaceId, - email: cleanEmail + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: cleanEmail }, + { email: cleanEmail } + ] + }); + } + + /** + * Retrieves suppression details for a specific company. + */ + public async getCompanySuppression(companyId: string): Promise { + if (!companyId) return null; + return SuppressionModel.findOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.COMPANY, + targetId: companyId.trim() + }); + } + + /** + * Retrieves suppression details for a specific domain. + */ + public async getDomainSuppression(domainOrEmail: string): Promise { + const normDomain = normalizeDomain(domainOrEmail); + if (!normDomain) return null; + return SuppressionModel.findOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain }); } @@ -55,7 +243,10 @@ export class SuppressionRepository { const cleanEmail = email.toLowerCase().trim(); const existing = await SuppressionModel.findOne({ workspaceId: this.workspaceId, - email: cleanEmail + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: cleanEmail }, + { email: cleanEmail } + ] }); if (existing) { @@ -70,6 +261,9 @@ export class SuppressionRepository { lastUpdatedAt: new Date().toISOString() }; + existing.targetType = SuppressionTargetType.RECIPIENT; + existing.targetId = cleanEmail; + existing.email = cleanEmail; existing.reason = targetReason; existing.source = comparison >= 0 ? source : existing.source; existing.evidence = mergedEvidence; @@ -84,7 +278,7 @@ export class SuppressionRepository { targetReason, existingReason: existing.reason }, - 'Updated existing suppression record with precedence enforcement' + 'Updated existing email suppression record with precedence enforcement' ); return existing; @@ -92,6 +286,8 @@ export class SuppressionRepository { const created = await SuppressionModel.create({ workspaceId: this.workspaceId, + targetType: SuppressionTargetType.RECIPIENT, + targetId: cleanEmail, email: cleanEmail, reason, source, @@ -114,9 +310,300 @@ export class SuppressionRepository { return created; } + /** + * Idempotently records a company-level Do Not Contact suppression in the workspace. + * Cancels active sequence executions for contacts belonging to this company. + */ + public async suppressCompany( + companyId: string, + reason: SuppressionReason = SuppressionReason.COMPANY_DNC, + source = 'system', + evidence: Record | null = null, + suppressedBy: string | null = null, + notes: string | null = null + ): Promise { + const cleanCompanyId = companyId.trim(); + const existing = await SuppressionModel.findOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.COMPANY, + targetId: cleanCompanyId + }); + + if (existing) { + const comparison = compareSuppressionPrecedence(reason, existing.reason); + const targetReason = comparison > 0 ? reason : existing.reason; + + const mergedEvidence = { + ...(existing.evidence || {}), + ...(evidence || {}), + lastUpdatedReason: reason, + lastUpdatedAt: new Date().toISOString() + }; + + existing.reason = targetReason; + existing.source = comparison >= 0 ? source : existing.source; + existing.evidence = mergedEvidence; + if (notes) existing.notes = notes; + if (suppressedBy) existing.suppressedBy = suppressedBy; + await existing.save(); + + logger.info( + { + workspaceId: this.workspaceId, + companyId: cleanCompanyId, + targetReason, + existingReason: existing.reason + }, + 'Updated existing company DNC record with precedence enforcement' + ); + + return existing; + } + + const created = await SuppressionModel.create({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.COMPANY, + targetId: cleanCompanyId, + companyId: cleanCompanyId, + reason, + source, + evidence, + suppressedAt: new Date(), + suppressedBy, + notes + }); + + // Cascade cancellation to existing active/queued sequence executions for matching company contacts + try { + const matchingContactIds = await ContactModel.find({ + workspaceId: this.workspaceId, + companyId: cleanCompanyId, + deletedAt: null + }).distinct('_id'); + + if (matchingContactIds.length > 0) { + const cancelRes = await SequenceExecutionModel.updateMany( + { + workspaceId: this.workspaceId, + contactId: { $in: matchingContactIds.map(String) }, + status: { $in: ['PENDING', 'RUNNING', 'WAITING'] } + }, + { $set: { status: 'CANCELLED' } } + ); + + logger.info( + { + workspaceId: this.workspaceId, + companyId: cleanCompanyId, + cancelledExecutions: cancelRes.modifiedCount + }, + 'Cancelled queued sequence executions for company DNC cascade' + ); + } + } catch (cancelErr) { + logger.warn( + { cancelErr, workspaceId: this.workspaceId, companyId: cleanCompanyId }, + 'Warning during sequence execution cancellation for company DNC' + ); + } + + logger.info( + { + workspaceId: this.workspaceId, + companyId: cleanCompanyId, + reason, + source + }, + 'Created new company DNC suppression record' + ); + + return created; + } + + /** + * Idempotently records a domain-level suppression in the workspace. + * Cancels active sequence executions for contacts belonging to this domain. + */ + public async suppressDomain( + domainOrEmail: string, + reason: SuppressionReason = SuppressionReason.DOMAIN_SUPPRESSION, + source = 'system', + evidence: Record | null = null, + suppressedBy: string | null = null, + notes: string | null = null + ): Promise { + const normDomain = normalizeDomain(domainOrEmail); + if (!normDomain) { + throw new Error(`Cannot suppress invalid domain: "${domainOrEmail}".`); + } + + const existing = await SuppressionModel.findOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain + }); + + if (existing) { + const comparison = compareSuppressionPrecedence(reason, existing.reason); + const targetReason = comparison > 0 ? reason : existing.reason; + + const mergedEvidence = { + ...(existing.evidence || {}), + ...(evidence || {}), + lastUpdatedReason: reason, + lastUpdatedAt: new Date().toISOString() + }; + + existing.reason = targetReason; + existing.source = comparison >= 0 ? source : existing.source; + existing.evidence = mergedEvidence; + if (notes) existing.notes = notes; + if (suppressedBy) existing.suppressedBy = suppressedBy; + await existing.save(); + + logger.info( + { + workspaceId: this.workspaceId, + domain: normDomain, + targetReason, + existingReason: existing.reason + }, + 'Updated existing domain suppression record with precedence enforcement' + ); + + return existing; + } + + const created = await SuppressionModel.create({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain, + domain: normDomain, + reason, + source, + evidence, + suppressedAt: new Date(), + suppressedBy, + notes + }); + + // Cascade cancellation to existing active/queued sequence executions for matching domain contacts + try { + const escapedDomain = normDomain.replace(/[.*+?^${}()|[\]\\]/g, '\\$&'); + const matchingContactIds = await ContactModel.find({ + workspaceId: this.workspaceId, + $or: [ + { email: { $regex: `@${escapedDomain}$`, $options: 'i' } }, + { 'additionalEmails.email': { $regex: `@${escapedDomain}$`, $options: 'i' } } + ], + deletedAt: null + }).distinct('_id'); + + if (matchingContactIds.length > 0) { + const cancelRes = await SequenceExecutionModel.updateMany( + { + workspaceId: this.workspaceId, + contactId: { $in: matchingContactIds.map(String) }, + status: { $in: ['PENDING', 'RUNNING', 'WAITING'] } + }, + { $set: { status: 'CANCELLED' } } + ); + + logger.info( + { + workspaceId: this.workspaceId, + domain: normDomain, + cancelledExecutions: cancelRes.modifiedCount + }, + 'Cancelled queued sequence executions for domain suppression cascade' + ); + } + } catch (cancelErr) { + logger.warn( + { cancelErr, workspaceId: this.workspaceId, domain: normDomain }, + 'Warning during sequence execution cancellation for domain suppression' + ); + } + + logger.info( + { + workspaceId: this.workspaceId, + domain: normDomain, + reason, + source + }, + 'Created new domain suppression record' + ); + + return created; + } + + /** + * Removes company DNC suppression in the workspace without affecting individual recipient suppressions. + */ + public async unsuppressCompany( + companyId: string, + removedBy?: string + ): Promise<{ unsuppressed: boolean; companyId: string }> { + const cleanCompanyId = companyId.trim(); + const res = await SuppressionModel.deleteOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.COMPANY, + targetId: cleanCompanyId + }); + + const deleted = (res.deletedCount ?? 0) > 0; + logger.info( + { + workspaceId: this.workspaceId, + companyId: cleanCompanyId, + removedBy, + deletedCount: res.deletedCount + }, + 'Unsuppressed company DNC record' + ); + + return { + unsuppressed: deleted, + companyId: cleanCompanyId + }; + } + + /** + * Removes domain suppression in the workspace without affecting individual recipient suppressions. + */ + public async unsuppressDomain( + domainOrEmail: string, + removedBy?: string + ): Promise<{ unsuppressed: boolean; domain: string }> { + const normDomain = normalizeDomain(domainOrEmail); + const res = await SuppressionModel.deleteOne({ + workspaceId: this.workspaceId, + targetType: SuppressionTargetType.DOMAIN, + targetId: normDomain + }); + + const deleted = (res.deletedCount ?? 0) > 0; + logger.info( + { + workspaceId: this.workspaceId, + domain: normDomain, + removedBy, + deletedCount: res.deletedCount + }, + 'Unsuppressed domain suppression record' + ); + + return { + unsuppressed: deleted, + domain: normDomain + }; + } + /** * Removes suppression for an email address (manual unsuppress) and synchronizes * contact eligibility if the contact has no remaining active suppressions (UNSUPPRESS-13). + * Verifies contact is not still blocked by active company DNC or domain suppression. */ public async unsuppress( email: string, @@ -125,7 +612,10 @@ export class SuppressionRepository { const cleanEmail = email.toLowerCase().trim(); const res = await SuppressionModel.deleteOne({ workspaceId: this.workspaceId, - email: cleanEmail + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: cleanEmail }, + { email: cleanEmail } + ] }); const deleted = (res.deletedCount ?? 0) > 0; @@ -159,11 +649,25 @@ export class SuppressionRepository { if (otherEmails.length > 0) { const otherSuppCount = await SuppressionModel.countDocuments({ workspaceId: this.workspaceId, - email: { $in: otherEmails } + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: { $in: otherEmails } }, + { email: { $in: otherEmails } } + ] }); hasOtherSuppression = otherSuppCount > 0; } + // Check if contact is still blocked by company DNC or domain suppression + if (!hasOtherSuppression && contact.companyId) { + hasOtherSuppression = await this.isCompanySuppressed(contact.companyId); + } + if (!hasOtherSuppression && contact.email) { + const domain = normalizeDomain(contact.email); + if (domain) { + hasOtherSuppression = await this.isDomainSuppressed(domain); + } + } + // Narrowest restoration: only restore if no other suppressions exist // and contact was blocked by BOUNCED or INVALID emailStatus if (!hasOtherSuppression) { @@ -224,11 +728,15 @@ export class SuppressionRepository { * Lists suppressions for the workspace. */ public async listSuppressions(filter: { + targetType?: SuppressionTargetType | string | undefined; reason?: string | undefined; limit?: number | undefined; skip?: number | undefined; } = {}): Promise<{ items: SuppressionDocument[]; total: number }> { const query: any = { workspaceId: this.workspaceId }; + if (filter.targetType) { + query.targetType = filter.targetType; + } if (filter.reason) { query.reason = filter.reason; } diff --git a/apps/api/src/routes/suppressions.ts b/apps/api/src/routes/suppressions.ts index 458b737e..79baafe0 100644 --- a/apps/api/src/routes/suppressions.ts +++ b/apps/api/src/routes/suppressions.ts @@ -1,6 +1,6 @@ import { OpenAPIHono } from '@hono/zod-openapi'; import { SuppressionRepository } from '../repositories/suppression/suppression.repository.js'; -import { createSuppressionDtoSchema } from '@leadforge/schema'; +import { createSuppressionDtoSchema, SuppressionTargetType, SuppressionReason } from '@leadforge/schema'; import { successResponse } from '../utils/index.js'; import { getWorkspaceId, getUserId } from './common.js'; import { BadRequestError } from '../errors/index.js'; @@ -11,34 +11,63 @@ export const suppressionsRouter = new OpenAPIHono(); suppressionsRouter.get('/', async (c) => { const wsId = getWorkspaceId(c); const repo = new SuppressionRepository(wsId); + const targetType = c.req.query('targetType') as SuppressionTargetType | undefined; const reason = c.req.query('reason'); const limit = c.req.query('limit') ? parseInt(c.req.query('limit')!, 10) : 50; const skip = c.req.query('skip') ? parseInt(c.req.query('skip')!, 10) : 0; - const result = await repo.listSuppressions({ reason, limit, skip }); + const result = await repo.listSuppressions({ targetType, reason, limit, skip }); return c.json(successResponse(result)); }); -// 2. Check if a specific email is suppressed +// 2. Check if an email, company, or domain is suppressed suppressionsRouter.get('/check', async (c) => { const wsId = getWorkspaceId(c); const email = c.req.query('email'); - if (!email) { - throw new BadRequestError('Query param "email" is required.'); + const companyId = c.req.query('companyId'); + const domain = c.req.query('domain'); + + if (!email && !companyId && !domain) { + throw new BadRequestError('At least one of "email", "companyId", or "domain" query param is required.'); } const repo = new SuppressionRepository(wsId); - const suppression = await repo.getSuppression(email); + + // Exact backward compatibility for single-email query contracts + if (email && !companyId && !domain) { + const suppression = await repo.getSuppression(email); + return c.json( + successResponse({ + email, + suppressed: Boolean(suppression), + suppression: suppression || null + }) + ); + } + + // Multi-target evaluation + const effective = await repo.evaluateEffectiveSuppression({ + email: email || '', + companyId: companyId || null + }); + return c.json( successResponse({ - email, - suppressed: Boolean(suppression), - suppression: suppression || null + email: email || null, + companyId: companyId || null, + domain: domain || null, + suppressed: effective.suppressed, + isRecipientSuppressed: effective.isRecipientSuppressed, + isCompanySuppressed: effective.isCompanySuppressed, + isDomainSuppressed: effective.isDomainSuppressed, + reasons: effective.reasons, + primaryReason: effective.primaryReason || null, + message: effective.message }) ); }); -// 3. Record a suppression +// 3. Record a suppression (recipient, company, or domain) suppressionsRouter.post('/', async (c) => { const wsId = getWorkspaceId(c); const body = await c.req.json().catch(() => ({})); @@ -52,19 +81,66 @@ suppressionsRouter.post('/', async (c) => { } const repo = new SuppressionRepository(wsId); - const record = await repo.suppress( - validated.email, - validated.reason, - validated.source || 'manual', - validated.evidence || null, - userId, - validated.notes || null - ); + const targetType = validated.targetType || SuppressionTargetType.RECIPIENT; + + if (targetType === SuppressionTargetType.COMPANY) { + const compId = (validated.companyId || validated.targetId || '').trim(); + if (!compId) throw new BadRequestError('companyId is required for company suppression.'); + const record = await repo.suppressCompany( + compId, + validated.reason || SuppressionReason.COMPANY_DNC, + validated.source || 'manual', + validated.evidence || null, + userId, + validated.notes || null + ); + return c.json(successResponse(record), 201); + } else if (targetType === SuppressionTargetType.DOMAIN) { + const dom = (validated.domain || validated.targetId || '').trim(); + if (!dom) throw new BadRequestError('domain is required for domain suppression.'); + const record = await repo.suppressDomain( + dom, + validated.reason || SuppressionReason.DOMAIN_SUPPRESSION, + validated.source || 'manual', + validated.evidence || null, + userId, + validated.notes || null + ); + return c.json(successResponse(record), 201); + } else { + const targetEmail = (validated.email || validated.targetId || '').trim(); + if (!targetEmail) throw new BadRequestError('email is required for recipient suppression.'); + const record = await repo.suppress( + targetEmail, + validated.reason || SuppressionReason.MANUAL_SUPPRESSION, + validated.source || 'manual', + validated.evidence || null, + userId, + validated.notes || null + ); + return c.json(successResponse(record), 201); + } +}); - return c.json(successResponse(record), 201); +// 4. Remove company DNC suppression +suppressionsRouter.delete('/company/:companyId', async (c) => { + const wsId = getWorkspaceId(c); + const companyId = decodeURIComponent(c.req.param('companyId')); + const repo = new SuppressionRepository(wsId); + const result = await repo.unsuppressCompany(companyId); + return c.json(successResponse(result)); +}); + +// 5. Remove domain suppression +suppressionsRouter.delete('/domain/:domain', async (c) => { + const wsId = getWorkspaceId(c); + const domain = decodeURIComponent(c.req.param('domain')); + const repo = new SuppressionRepository(wsId); + const result = await repo.unsuppressDomain(domain); + return c.json(successResponse(result)); }); -// 4. Remove suppression (unsuppress) +// 6. Remove recipient suppression (unsuppress) suppressionsRouter.delete('/:email', async (c) => { const wsId = getWorkspaceId(c); const email = decodeURIComponent(c.req.param('email')); diff --git a/apps/api/src/services/automation/automation.service.ts b/apps/api/src/services/automation/automation.service.ts index 248f0b7e..efaca700 100644 --- a/apps/api/src/services/automation/automation.service.ts +++ b/apps/api/src/services/automation/automation.service.ts @@ -1,7 +1,9 @@ import { SequenceModel } from '../../db/models/sequence.model.js'; import { SequenceExecutionModel } from '../../db/models/sequence-execution.model.js'; import { CampaignModel } from '../../db/models/campaign.model.js'; +import { ContactModel } from '../../db/models/contact.model.js'; import { SequenceLogModel } from '../../db/models/sequence-log.model.js'; +import { SuppressionRepository } from '../../repositories/suppression/suppression.repository.js'; import { SequenceStatus, ExecutionStatus } from '@leadforge/schema'; import { ConflictError } from '../../errors/index.js'; @@ -36,7 +38,6 @@ export class AutomationService { _id: id, workspaceId: this.workspaceId } as any); - if (!seq) throw new Error('Sequence not found.'); return seq; } @@ -61,6 +62,27 @@ export class AutomationService { public async createExecution(data: any): Promise { if (data.contactId && data.campaignId) { + // Early policy filtering: check effective workspace suppression (recipient, company DNC, domain suppression) + const contactDoc = await ContactModel.findOne({ + _id: data.contactId, + workspaceId: this.workspaceId, + deletedAt: null + }); + + if (contactDoc) { + const suppressionRepo = new SuppressionRepository(this.workspaceId); + const effectiveSuppression = await suppressionRepo.evaluateEffectiveSuppression({ + email: contactDoc.email || '', + companyId: contactDoc.companyId || data.companyId || null + }); + + if (effectiveSuppression.suppressed) { + throw new ConflictError( + `Cannot enroll contact "${data.contactId}" in campaign: ${effectiveSuppression.message}` + ); + } + } + // Phase 15 (ENROLL-08): Contact cross-campaign active exclusivity check const existingActive = await SequenceExecutionModel.findOne({ workspaceId: this.workspaceId, diff --git a/apps/api/src/services/email/company-dnc-domain-suppression.test.ts b/apps/api/src/services/email/company-dnc-domain-suppression.test.ts new file mode 100644 index 00000000..74319dd2 --- /dev/null +++ b/apps/api/src/services/email/company-dnc-domain-suppression.test.ts @@ -0,0 +1,545 @@ +import { describe, it, expect, vi, beforeEach } from 'vitest'; +import { SuppressionRepository } from '../../repositories/suppression/suppression.repository.js'; +import { EmailService } from './email.service.js'; +import { SuppressionModel } from '../../db/models/suppression.model.js'; +import { ContactModel } from '../../db/models/contact.model.js'; +import { CampaignModel } from '../../db/models/campaign.model.js'; +import { EmailDeliveryModel } from '../../db/models/email-delivery.model.js'; +import { EmailAccountModel } from '../../db/models/email-account.model.js'; +import { SequenceExecutionModel } from '../../db/models/sequence-execution.model.js'; +import { AutomationService } from '../automation/automation.service.js'; +import { DomainPacingService } from '../outreach/domain-pacing.service.js'; +import { + SuppressionReason, + SuppressionTargetType, + normalizeDomain, + ContactStatus, + ContactEmailStatus +} from '@leadforge/schema'; +import { EmailDomainError } from './types.js'; + +vi.mock('../../db/models/suppression.model.js'); +vi.mock('../../db/models/contact.model.js'); +vi.mock('../../db/models/campaign.model.js'); +vi.mock('../../db/models/email-delivery.model.js'); +vi.mock('../../db/models/email-account.model.js'); +vi.mock('../../db/models/sequence-execution.model.js'); + +describe('fix(outreach): enforce company-level DNC and domain suppression cascade (#38)', () => { + const wsA = 'ws_alpha'; + const wsB = 'ws_beta'; + const companyAcme = 'comp_acme_123'; + const companyBeta = 'comp_beta_456'; + + beforeEach(() => { + vi.clearAllMocks(); + }); + + describe('SuppressionRepository — Company & Domain Cascade Semantics', () => { + it('Criterion A: Recipient suppression still works', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.countDocuments as any).mockResolvedValue(1); + + const isSupp = await repo.isSuppressed('jane@example.com'); + expect(isSupp).toBe(true); + expect(SuppressionModel.countDocuments).toHaveBeenCalledWith({ + workspaceId: wsA, + $or: [ + { targetType: SuppressionTargetType.RECIPIENT, targetId: 'jane@example.com' }, + { email: 'jane@example.com' } + ] + }); + }); + + it('Criterion B: Company DNC blocks a linked contact', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.countDocuments as any).mockResolvedValue(1); + + const isCompSupp = await repo.isCompanySuppressed(companyAcme); + expect(isCompSupp).toBe(true); + expect(SuppressionModel.countDocuments).toHaveBeenCalledWith({ + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme + }); + }); + + it('Criterion C: Company DNC blocks multiple contacts under the same company', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.find as any).mockResolvedValue([ + { + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme, + reason: SuppressionReason.COMPANY_DNC + } + ]); + + const evalContact1 = await repo.evaluateEffectiveSuppression({ + email: 'alice@acme.com', + companyId: companyAcme + }); + expect(evalContact1.suppressed).toBe(true); + expect(evalContact1.isCompanySuppressed).toBe(true); + expect(evalContact1.primaryReason).toBe(SuppressionReason.COMPANY_DNC); + + const evalContact2 = await repo.evaluateEffectiveSuppression({ + email: 'charlie@acme.com', + companyId: companyAcme + }); + expect(evalContact2.suppressed).toBe(true); + expect(evalContact2.isCompanySuppressed).toBe(true); + }); + + it('Criterion D: Company DNC blocks contacts across multiple domains sharing the same canonical companyId', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.find as any).mockResolvedValue([ + { + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme, + reason: SuppressionReason.COMPANY_DNC + } + ]); + + // alice@acme.com and bob@acme.co.uk both share companyAcme + const evalAlice = await repo.evaluateEffectiveSuppression({ + email: 'alice@acme.com', + companyId: companyAcme + }); + const evalBobUk = await repo.evaluateEffectiveSuppression({ + email: 'bob@acme.co.uk', + companyId: companyAcme + }); + + expect(evalAlice.suppressed).toBe(true); + expect(evalBobUk.suppressed).toBe(true); + expect(evalAlice.isCompanySuppressed).toBe(true); + expect(evalBobUk.isCompanySuppressed).toBe(true); + }); + + it('Criterion E & F: Domain suppression blocks matching normalized domains with case-insensitivity', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.countDocuments as any).mockResolvedValue(1); + + expect(normalizeDomain('Person@Example.COM')).toBe('example.com'); + expect(normalizeDomain('person@example.com')).toBe('example.com'); + expect(normalizeDomain('USER@EXAMPLE.COM')).toBe('example.com'); + + const isSupp1 = await repo.isDomainSuppressed('Person@Example.COM'); + const isSupp2 = await repo.isDomainSuppressed('example.com'); + expect(isSupp1).toBe(true); + expect(isSupp2).toBe(true); + expect(SuppressionModel.countDocuments).toHaveBeenCalledWith({ + workspaceId: wsA, + targetType: SuppressionTargetType.DOMAIN, + targetId: 'example.com' + }); + }); + + it('Criterion G & T: Workspace isolation — suppression in Workspace A does not affect Workspace B', async () => { + const repoA = new SuppressionRepository(wsA); + const repoB = new SuppressionRepository(wsB); + + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.workspaceId === wsA) return Promise.resolve(1); + return Promise.resolve(0); + }); + + expect(await repoA.isCompanySuppressed(companyAcme)).toBe(true); + expect(await repoB.isCompanySuppressed(companyAcme)).toBe(false); + + expect(await repoA.isDomainSuppressed('acme.com')).toBe(true); + expect(await repoB.isDomainSuppressed('acme.com')).toBe(false); + }); + + it('Criterion H & S: Different company does not inherit company DNC', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetId === companyAcme) return Promise.resolve(1); + return Promise.resolve(0); + }); + + expect(await repo.isCompanySuppressed(companyAcme)).toBe(true); + expect(await repo.isCompanySuppressed(companyBeta)).toBe(false); + }); + + it('Criterion I: Different domain does not inherit domain suppression', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetId === 'blocked.com') return Promise.resolve(1); + return Promise.resolve(0); + }); + + expect(await repo.isDomainSuppressed('blocked.com')).toBe(true); + expect(await repo.isDomainSuppressed('allowed.com')).toBe(false); + expect(await repo.isDomainSuppressed('sub.blocked.com')).toBe(false); + }); + + it('Criterion R: Repeated suppression requests are idempotent', async () => { + const repo = new SuppressionRepository(wsA); + const existingDoc = { + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme, + reason: SuppressionReason.COMPANY_DNC, + source: 'manual', + save: vi.fn().mockResolvedValue(true) + }; + + (SuppressionModel.findOne as any).mockResolvedValue(existingDoc); + + const res = await repo.suppressCompany(companyAcme, SuppressionReason.COMPANY_DNC); + expect(existingDoc.save).toHaveBeenCalled(); + expect(SuppressionModel.create).not.toHaveBeenCalled(); + expect(res.targetId).toBe(companyAcme); + }); + + it('Criterion U, V & W: Independent suppression causes and additive unsuppression', async () => { + const repo = new SuppressionRepository(wsA); + + // Unsuppressing company does NOT delete recipient suppression + (SuppressionModel.deleteOne as any).mockResolvedValue({ deletedCount: 1 }); + const unsuppCompanyRes = await repo.unsuppressCompany(companyAcme); + expect(unsuppCompanyRes.unsuppressed).toBe(true); + expect(SuppressionModel.deleteOne).toHaveBeenCalledWith({ + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme + }); + + // If contact still has active HARD_BOUNCE suppression, evaluateEffectiveSuppression remains suppressed + (SuppressionModel.find as any).mockResolvedValue([ + { + workspaceId: wsA, + targetType: SuppressionTargetType.RECIPIENT, + targetId: 'alice@acme.com', + reason: SuppressionReason.HARD_BOUNCE + } + ]); + + const evalAfterUnsuppCompany = await repo.evaluateEffectiveSuppression({ + email: 'alice@acme.com', + companyId: companyAcme + }); + expect(evalAfterUnsuppCompany.suppressed).toBe(true); + expect(evalAfterUnsuppCompany.isRecipientSuppressed).toBe(true); + expect(evalAfterUnsuppCompany.isCompanySuppressed).toBe(false); + expect(evalAfterUnsuppCompany.primaryReason).toBe(SuppressionReason.HARD_BOUNCE); + }); + }); + + describe('Audience Enrollment & Execution Creation Enforcement', () => { + it('Criterion J: Early policy filtering prevents enrolling company-DNC contacts', async () => { + const autoService = new AutomationService(wsA); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_123', + email: 'blocked@acme.com', + companyId: companyAcme + }); + + (SuppressionModel.find as any).mockResolvedValue([ + { + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme, + reason: SuppressionReason.COMPANY_DNC + } + ]); + + await expect( + autoService.createExecution({ + sequenceId: 'seq_1', + campaignId: 'camp_1', + contactId: 'contact_123' + }) + ).rejects.toThrow(/Cannot enroll contact "contact_123" in campaign/); + + expect(SequenceExecutionModel.prototype.save).not.toHaveBeenCalled(); + }); + + it('Criterion K & L: Suppressing company cascades cancellation to active sequence executions', async () => { + const repo = new SuppressionRepository(wsA); + (SuppressionModel.findOne as any).mockResolvedValue(null); + (SuppressionModel.create as any).mockResolvedValue({ + workspaceId: wsA, + targetType: SuppressionTargetType.COMPANY, + targetId: companyAcme + }); + + (ContactModel.find as any).mockReturnValue({ + distinct: vi.fn().mockResolvedValue(['contact_1', 'contact_2']) + }); + (SequenceExecutionModel.updateMany as any).mockResolvedValue({ modifiedCount: 2 }); + + await repo.suppressCompany(companyAcme); + + expect(SequenceExecutionModel.updateMany).toHaveBeenCalledWith( + { + workspaceId: wsA, + contactId: { $in: ['contact_1', 'contact_2'] }, + status: { $in: ['PENDING', 'RUNNING', 'WAITING'] } + }, + { $set: { status: 'CANCELLED' } } + ); + }); + }); + + describe('Final Server-Authoritative Send Gate (EmailService.send)', () => { + const defaultAccount = { + _id: 'acc_1', + workspaceId: wsA, + email: 'sender@leadforge.ai', + status: 'connected', + provider: 'gmail_oauth', + sendPolicy: { dailyLimit: 100, hourlyLimit: 20 } + }; + + it('Criterion M & N: MANDATORY STALE-CACHE TEST — worker local cache says eligible, Mongo has DNC, API rejects send and provider is never called', async () => { + const emailService = new EmailService(wsA, 'user_1'); + + // Mailbox is active + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + // Contact exists with companyAcme + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_stale', + workspaceId: wsA, + email: 'alice@acme.com', + companyId: companyAcme, + status: ContactStatus.NEW, + emailStatus: ContactEmailStatus.VALID + }); + + // Recipient email itself is not in suppressions + // BUT Company is marked DNC in MongoDB! + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetType === SuppressionTargetType.COMPANY && query.targetId === companyAcme) { + return Promise.resolve(1); + } + return Promise.resolve(0); + }); + + let thrownError: any = null; + try { + await emailService.send({ + accountId: 'acc_1', + to: 'alice@acme.com', + subject: 'Outreach Test', + text: 'Hello Alice', + campaignId: 'camp_1', + contactId: 'contact_stale' + }); + } catch (err) { + thrownError = err; + } + + // Assert local policy rejection + expect(thrownError).toBeInstanceOf(EmailDomainError); + expect(thrownError.code).toBe('COMPANY_DNC'); + expect(thrownError.message).toContain('marked Do Not Contact'); + + // Assert provider is NEVER called + expect(EmailDeliveryModel.create).not.toHaveBeenCalled(); + expect(EmailDeliveryModel.findOneAndUpdate).not.toHaveBeenCalled(); + }); + + it('Criterion O: Company DNC does not create an EmailDelivery record or provider failure', async () => { + const emailService = new EmailService(wsA, 'user_1'); + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_1', + workspaceId: wsA, + email: 'bob@acme.com', + companyId: companyAcme, + status: ContactStatus.NEW, + emailStatus: ContactEmailStatus.VALID + }); + + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetType === SuppressionTargetType.COMPANY && query.targetId === companyAcme) { + return Promise.resolve(1); + } + return Promise.resolve(0); + }); + + await expect( + emailService.send({ + accountId: 'acc_1', + to: 'bob@acme.com', + subject: 'Test', + campaignId: 'camp_1', + contactId: 'contact_1' + }) + ).rejects.toThrow(EmailDomainError); + + expect(EmailDeliveryModel.prototype.save).not.toHaveBeenCalled(); + }); + + it('Criterion P: Company DNC does not trip or increment rejection circuit-breaker counters', async () => { + const emailService = new EmailService(wsA, 'user_1'); + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_1', + workspaceId: wsA, + email: 'bob@acme.com', + companyId: companyAcme + }); + + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetType === SuppressionTargetType.COMPANY) return Promise.resolve(1); + return Promise.resolve(0); + }); + + try { + await emailService.send({ + accountId: 'acc_1', + to: 'bob@acme.com', + subject: 'Test', + campaignId: 'camp_1', + contactId: 'contact_1' + }); + } catch (err: any) { + expect(err.code).toBe('COMPANY_DNC'); + } + + // Verify no failure delivery records were written for the campaign + expect(EmailDeliveryModel.find).not.toHaveBeenCalledWith( + expect.objectContaining({ status: { $in: ['SENT', 'FAILED'] } }) + ); + }); + + it('Criterion Q: Company DNC short-circuits before domain pacing and company cardinality reservation', async () => { + const emailService = new EmailService(wsA, 'user_1'); + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_1', + workspaceId: wsA, + email: 'bob@acme.com', + companyId: companyAcme + }); + + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetType === SuppressionTargetType.COMPANY) return Promise.resolve(1); + return Promise.resolve(0); + }); + + try { + await emailService.send({ + accountId: 'acc_1', + to: 'bob@acme.com', + subject: 'Test', + campaignId: 'camp_1', + contactId: 'contact_1' + }); + } catch (err: any) { + expect(err.code).toBe('COMPANY_DNC'); + } + + // Verify pacing lock was never attempted + expect(EmailDeliveryModel.distinct).not.toHaveBeenCalled(); + }); + + it('Criterion L: Domain suppression blocks matching recipient domain at send gate', async () => { + const emailService = new EmailService(wsA, 'user_1'); + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_domain', + workspaceId: wsA, + email: 'lead@blockeddomain.com', + companyId: null + }); + + (SuppressionModel.countDocuments as any).mockImplementation((query: any) => { + if (query.targetType === SuppressionTargetType.DOMAIN && query.targetId === 'blockeddomain.com') { + return Promise.resolve(1); + } + return Promise.resolve(0); + }); + + await expect( + emailService.send({ + accountId: 'acc_1', + to: 'lead@blockeddomain.com', + subject: 'Test', + contactId: 'contact_domain' + }) + ).rejects.toThrow(/Domain "blockeddomain.com" is suppressed/); + }); + + it('Criterion X: Provider acceptance still behaves normally for a non-suppressed recipient', async () => { + const emailService = new EmailService(wsA, 'user_1'); + (EmailAccountModel.findOne as any).mockResolvedValue(defaultAccount); + + (ContactModel.findOne as any).mockResolvedValue({ + _id: 'contact_valid', + workspaceId: wsA, + email: 'valid.lead@example.com', + companyId: 'comp_valid', + status: ContactStatus.NEW, + emailStatus: ContactEmailStatus.VALID + }); + + // No suppressions active + (SuppressionModel.countDocuments as any).mockResolvedValue(0); + + // Campaign active + (CampaignModel.findOne as any).mockResolvedValue({ + _id: 'camp_active', + workspaceId: wsA, + status: 'ACTIVE' + }); + + // Mock pacing & account limits + vi.spyOn(DomainPacingService.prototype, 'checkAndReservePacing').mockResolvedValue({ + allowed: true, + leaseExpiresAt: new Date(Date.now() + 60000) + } as any); + (emailService as any).accountRepo.resolveEffectiveLimits = vi.fn().mockResolvedValue({ + dailyLimit: 100, + hourlyLimit: 20 + }); + (emailService as any).accountRepo.reserveSendSlot = vi.fn().mockResolvedValue({ success: true }); + (emailService as any).deliveryRepo.reserveDelivery = vi.fn().mockResolvedValue({ + delivery: { _id: 'del_1', status: 'SENDING' }, + isAlreadySent: false + }); + (emailService as any).accounts.buildProvider = vi.fn().mockResolvedValue({ + send: vi.fn().mockResolvedValue({ + messageId: 'gmail_msg_100', + threadId: 'gmail_th_100' + }) + }); + (emailService as any).deliveryRepo.finalizeDelivery = vi.fn().mockResolvedValue({ + _id: 'del_1', + status: 'SENT' + }); + + const res = await emailService.send({ + accountId: 'acc_1', + to: 'valid.lead@example.com', + subject: 'Valid Send', + text: 'Hello', + campaignId: 'camp_active', + contactId: 'contact_valid' + }); + + expect(res.messageId).toBe('gmail_msg_100'); + expect(res.accepted).toContain('valid.lead@example.com'); + }); + + it('Criterion Y: Ambiguous-send behavior remains unchanged and distinct from DNC', () => { + const dncError = new EmailDomainError('COMPANY_DNC', 'Company DNC active'); + const ambiguousError = new EmailDomainError('AMBIGUOUS_SEND_TIMEOUT', 'Google connection timed out', false, false, 'ambiguous'); + + expect(dncError.code).toBe('COMPANY_DNC'); + expect(ambiguousError.code).toBe('AMBIGUOUS_SEND_TIMEOUT'); + expect(ambiguousError.classification).toBe('ambiguous'); + }); + }); +}); diff --git a/apps/api/src/services/email/email.service.ts b/apps/api/src/services/email/email.service.ts index 05ea6156..06364160 100644 --- a/apps/api/src/services/email/email.service.ts +++ b/apps/api/src/services/email/email.service.ts @@ -23,6 +23,7 @@ import { classifyBounce, mapBounceCategoryToFailureCategory, evaluateOutreachEligibility, + normalizeDomain, generateTrackingToken, injectOpenTrackingPixel, rewriteLinksForClickTracking, @@ -208,6 +209,20 @@ export function classifyEmailFailure(err: any): { }; } + // 8b. Company DNC & Domain Suppression policy rejections (local policy, NOT provider failures or hard bounces) + if (code === 'COMPANY_DNC' || code === 'DOMAIN_SUPPRESSED') { + return { + code, + category: EmailFailureCategory.POLICY, + safeHumanMessage: 'Outbound dispatch blocked by company or domain suppression policy.', + technicalMessage: msg, + retryable: false, + ambiguous: false, + bounceCategory: BounceCategory.POLICY_REJECTION, + isHardBounce: false + }; + } + // 9. Specific legacy address-level indicators not caught by numeric status codes if ( code === 'INVALID_RECIPIENT' || @@ -342,30 +357,19 @@ export class EmailService { ); } - // 0a. Pre-flight suppression check: block if recipient is suppressed in workspace (even for direct sends) + const normRecipient = input.to.toLowerCase().trim(); const suppressionRepo = new SuppressionRepository(this.workspaceId); - const isSuppressed = await suppressionRepo.isSuppressed(input.to); + + // 1. Workspace recipient suppression check: block if recipient is suppressed in workspace (even for direct sends) + const isSuppressed = await suppressionRepo.isSuppressed(normRecipient); if (isSuppressed) { throw new EmailDomainError( 'RECIPIENT_SUPPRESSED', - `Recipient "${input.to}" is suppressed in this workspace and cannot receive outreach.` + `Recipient "${normRecipient}" is suppressed in this workspace and cannot receive outreach.` ); } - // 0a. Server-authoritative campaign send authorization check - let campaignDoc: any = null; - if (input.campaignId) { - campaignDoc = await CampaignModel.findOne({ _id: input.campaignId, workspaceId: this.workspaceId }); - if (campaignDoc && campaignDoc.status !== 'ACTIVE') { - throw new EmailDomainError( - 'CAMPAIGN_NOT_ACTIVE', - `Campaign "${input.campaignId}" is in status "${campaignDoc.status}". Sending is not authorized.` - ); - } - } - - // 0b. Server-authoritative contact outreach eligibility check - const normRecipient = input.to.toLowerCase().trim(); + // Resolve contact document to determine canonical company identity let contactDoc: any = null; if (input.contactId && input.contactId !== 'direct-contact') { contactDoc = await ContactModel.findOne({ _id: input.contactId, workspaceId: this.workspaceId }); @@ -381,6 +385,63 @@ export class EmailService { } } + // 2. Company DNC check: block if contact's canonical company is marked Do Not Contact in workspace + const companyId = contactDoc?.companyId || null; + if (companyId) { + const isCompanyDnc = await suppressionRepo.isCompanySuppressed(companyId); + if (isCompanyDnc) { + logger.info( + { + workspaceId: this.workspaceId, + companyId, + contactId: input.contactId, + recipient: normRecipient, + campaignId: input.campaignId + }, + 'Outreach dispatch blocked by company DNC policy' + ); + throw new EmailDomainError( + 'COMPANY_DNC', + `Company "${companyId}" is marked Do Not Contact in this workspace. Outbound outreach to "${normRecipient}" is blocked.` + ); + } + } + + // 3. Domain suppression check: block if recipient domain is suppressed in workspace + const normDomain = normalizeDomain(normRecipient); + if (normDomain) { + const isDomainSuppressed = await suppressionRepo.isDomainSuppressed(normDomain); + if (isDomainSuppressed) { + logger.info( + { + workspaceId: this.workspaceId, + domain: normDomain, + contactId: input.contactId, + recipient: normRecipient, + campaignId: input.campaignId + }, + 'Outreach dispatch blocked by domain suppression policy' + ); + throw new EmailDomainError( + 'DOMAIN_SUPPRESSED', + `Domain "${normDomain}" is suppressed in this workspace. Outbound outreach to "${normRecipient}" is blocked.` + ); + } + } + + // 4. Server-authoritative campaign send authorization check + let campaignDoc: any = null; + if (input.campaignId) { + campaignDoc = await CampaignModel.findOne({ _id: input.campaignId, workspaceId: this.workspaceId }); + if (campaignDoc && campaignDoc.status !== 'ACTIVE') { + throw new EmailDomainError( + 'CAMPAIGN_NOT_ACTIVE', + `Campaign "${input.campaignId}" is in status "${campaignDoc.status}". Sending is not authorized.` + ); + } + } + + // 5. Server-authoritative contact outreach eligibility check if (contactDoc) { const eligibility = evaluateOutreachEligibility({ contact: { diff --git a/apps/api/src/services/email/types.ts b/apps/api/src/services/email/types.ts index c6bb5e7d..f4335fa8 100644 --- a/apps/api/src/services/email/types.ts +++ b/apps/api/src/services/email/types.ts @@ -121,6 +121,8 @@ export interface EmailProviderErrorShape { | 'INVALID_RECIPIENT' | 'INVALID_SUBJECT' | 'RECIPIENT_SUPPRESSED' + | 'COMPANY_DNC' + | 'DOMAIN_SUPPRESSED' | 'AMBIGUOUS_SEND_TIMEOUT' | 'DELIVERY_ALREADY_SENT' | 'DELIVERY_ALREADY_RESERVED' diff --git a/apps/desktop/src/main/workers/plugins/outreach.ts b/apps/desktop/src/main/workers/plugins/outreach.ts index 4c3daa4e..e5d3ca3d 100644 --- a/apps/desktop/src/main/workers/plugins/outreach.ts +++ b/apps/desktop/src/main/workers/plugins/outreach.ts @@ -384,7 +384,23 @@ export async function dispatchOutreach(ctx: JobContext): Promise { sendError.includes('COMPANY_CARDINALITY_EXCEEDED') || sendError.includes('cardinality limit reached'); - if (isAmbiguous) { + const isSuppressedOrDnc = + err.code === 'COMPANY_DNC' || + err.code === 'DOMAIN_SUPPRESSED' || + err.code === 'RECIPIENT_SUPPRESSED' || + sendError.includes('COMPANY_DNC') || + sendError.includes('DOMAIN_SUPPRESSED') || + sendError.includes('RECIPIENT_SUPPRESSED') || + sendError.includes('Do Not Contact') || + sendError.includes('suppressed in this workspace'); + + if (isSuppressedOrDnc) { + skippedCount++; + ctx.emitLog( + `Skipped contact "${contact.email}": blocked by suppression/DNC policy (${err.code || 'POLICY_BLOCKED'}).`, + 'info' + ); + } else if (isAmbiguous) { skippedCount++; ctx.emitLog( `⚠️ Ambiguous delivery outcome for "${contact.email}": send outcome is unconfirmed (pending reconciliation). Blind re-dispatch suppressed to prevent duplicate sending.`, @@ -475,10 +491,11 @@ export async function dispatchOutreach(ctx: JobContext): Promise { // Phase 10: Auto-suppress on hard bounce const isHardBounce = - err.code === 'INVALID_RECIPIENT' || - err.status === 400 || - sendError.includes('INVALID_RECIPIENT') || - sendError.includes('550'); + !isSuppressedOrDnc && + (err.code === 'INVALID_RECIPIENT' || + err.status === 400 || + sendError.includes('INVALID_RECIPIENT') || + sendError.includes('550')); if (isHardBounce) { try { diff --git a/packages/schema/src/entities/suppression.ts b/packages/schema/src/entities/suppression.ts index 3a15c40b..c8032376 100644 --- a/packages/schema/src/entities/suppression.ts +++ b/packages/schema/src/entities/suppression.ts @@ -1,13 +1,18 @@ import { z } from 'zod'; -import { SuppressionReason } from '../enums/index.js'; +import { SuppressionReason, SuppressionTargetType } from '../enums/index.js'; import { entityIdField, emailField } from '../fields/common.js'; export const suppressionReasonSchema = z.nativeEnum(SuppressionReason); +export const suppressionTargetTypeSchema = z.nativeEnum(SuppressionTargetType); export const suppressionRecordSchema = z.object({ id: entityIdField, workspaceId: entityIdField, - email: emailField, + targetType: suppressionTargetTypeSchema.default(SuppressionTargetType.RECIPIENT), + targetId: z.string().min(1), + email: emailField.optional().nullable(), + companyId: z.string().optional().nullable(), + domain: z.string().optional().nullable(), reason: suppressionReasonSchema, source: z.string().default('system'), evidence: z.record(z.any()).optional().nullable(), @@ -20,10 +25,30 @@ export const suppressionRecordSchema = z.object({ export type SuppressionRecord = z.infer; export const createSuppressionDtoSchema = z.object({ - email: emailField, - reason: suppressionReasonSchema, + targetType: suppressionTargetTypeSchema.optional().default(SuppressionTargetType.RECIPIENT), + targetId: z.string().optional(), + email: z.string().optional().nullable(), + companyId: z.string().optional().nullable(), + domain: z.string().optional().nullable(), + reason: suppressionReasonSchema.optional(), source: z.string().optional().default('manual'), - notes: z.string().optional(), - evidence: z.record(z.any()).optional() -}); + notes: z.string().nullable().optional(), + evidence: z.record(z.any()).optional().nullable() +}).refine( + (data) => { + const type = data.targetType || SuppressionTargetType.RECIPIENT; + if (type === SuppressionTargetType.RECIPIENT) { + return Boolean(data.email || data.targetId); + } + if (type === SuppressionTargetType.COMPANY) { + return Boolean(data.companyId || data.targetId); + } + if (type === SuppressionTargetType.DOMAIN) { + return Boolean(data.domain || data.targetId); + } + return false; + }, + { message: 'Must provide an identifier matching targetType (email, companyId, or domain).' } +); export type CreateSuppressionDto = z.infer; + diff --git a/packages/schema/src/enums/index.ts b/packages/schema/src/enums/index.ts index 76f40327..17927807 100644 --- a/packages/schema/src/enums/index.ts +++ b/packages/schema/src/enums/index.ts @@ -250,8 +250,17 @@ export enum EmailQualityStatus { } /** - * Phase 10: Structured reasons for contact / email address suppression. - * Follows strict precedence hierarchy: DO_NOT_CONTACT > UNSUBSCRIBED > SPAM_COMPLAINT > HARD_BOUNCE > MANUAL_SUPPRESSION > POLICY_BLOCK > INVALID_EMAIL. + * Target entity type for workspace suppression policies. + */ +export enum SuppressionTargetType { + RECIPIENT = 'recipient', + COMPANY = 'company', + DOMAIN = 'domain' +} + +/** + * Phase 10: Structured reasons for contact / email address / company / domain suppression. + * Follows strict precedence hierarchy: DO_NOT_CONTACT > COMPANY_DNC > DOMAIN_SUPPRESSION > UNSUBSCRIBED > SPAM_COMPLAINT > HARD_BOUNCE > MANUAL_SUPPRESSION > POLICY_BLOCK > INVALID_EMAIL. */ export enum SuppressionReason { DO_NOT_CONTACT = 'DO_NOT_CONTACT', @@ -260,7 +269,9 @@ export enum SuppressionReason { HARD_BOUNCE = 'HARD_BOUNCE', MANUAL_SUPPRESSION = 'MANUAL_SUPPRESSION', POLICY_BLOCK = 'POLICY_BLOCK', - INVALID_EMAIL = 'INVALID_EMAIL' + INVALID_EMAIL = 'INVALID_EMAIL', + COMPANY_DNC = 'COMPANY_DNC', + DOMAIN_SUPPRESSION = 'DOMAIN_SUPPRESSION' } /** diff --git a/packages/schema/src/utils/email-quality-engine.test.ts b/packages/schema/src/utils/email-quality-engine.test.ts index cc9c7f90..6fe4c1bd 100644 --- a/packages/schema/src/utils/email-quality-engine.test.ts +++ b/packages/schema/src/utils/email-quality-engine.test.ts @@ -96,10 +96,24 @@ describe('Phase 10: Email Quality Decision Engine', () => { }); describe('Suppression Precedence', () => { - it('enforces DO_NOT_CONTACT > UNSUBSCRIBED > HARD_BOUNCE > INVALID', () => { + it('enforces DO_NOT_CONTACT > COMPANY_DNC > DOMAIN_SUPPRESSION > UNSUBSCRIBED > HARD_BOUNCE > INVALID', () => { expect( compareSuppressionPrecedence( SuppressionReason.DO_NOT_CONTACT, + SuppressionReason.COMPANY_DNC + ) + ).toBeGreaterThan(0); + + expect( + compareSuppressionPrecedence( + SuppressionReason.COMPANY_DNC, + SuppressionReason.DOMAIN_SUPPRESSION + ) + ).toBeGreaterThan(0); + + expect( + compareSuppressionPrecedence( + SuppressionReason.DOMAIN_SUPPRESSION, SuppressionReason.UNSUBSCRIBED ) ).toBeGreaterThan(0); diff --git a/packages/schema/src/utils/email-quality-engine.ts b/packages/schema/src/utils/email-quality-engine.ts index b288ddf9..7b3881dd 100644 --- a/packages/schema/src/utils/email-quality-engine.ts +++ b/packages/schema/src/utils/email-quality-engine.ts @@ -20,6 +20,8 @@ import { isKnownRoleAccount, validateEmailStrict, evaluateEmailCandidate } from export const SUPPRESSION_PRECEDENCE_WEIGHTS: Record = { [SuppressionReason.DO_NOT_CONTACT]: 100, + [SuppressionReason.COMPANY_DNC]: 95, + [SuppressionReason.DOMAIN_SUPPRESSION]: 92, [SuppressionReason.UNSUBSCRIBED]: 90, [SuppressionReason.SPAM_COMPLAINT]: 80, [SuppressionReason.HARD_BOUNCE]: 70, diff --git a/packages/schema/src/utils/outreach-eligibility.test.ts b/packages/schema/src/utils/outreach-eligibility.test.ts index 7e3c31fc..4905dcf3 100644 --- a/packages/schema/src/utils/outreach-eligibility.test.ts +++ b/packages/schema/src/utils/outreach-eligibility.test.ts @@ -147,6 +147,26 @@ describe('Contact Outreach Eligibility Policy', () => { expect(res.reason).toBe('EMAIL_SUPPRESSED'); }); + it('rejects contact when company is marked DNC', () => { + const res = evaluateOutreachEligibility({ + contact: { email: 'alice@acme.com', status: ContactStatus.NEW }, + campaign: { status: CampaignStatus.ACTIVE }, + companySuppressed: true + }); + expect(res.eligible).toBe(false); + expect(res.reason).toBe('COMPANY_DNC'); + }); + + it('rejects contact when domain is suppressed', () => { + const res = evaluateOutreachEligibility({ + contact: { email: 'bob@blockeddomain.com', status: ContactStatus.NEW }, + campaign: { status: CampaignStatus.ACTIVE }, + domainSuppressed: true + }); + expect(res.eligible).toBe(false); + expect(res.reason).toBe('DOMAIN_SUPPRESSED'); + }); + it('rejects contact when email domain is a disposable address', () => { const res = evaluateOutreachEligibility({ contact: { email: 'lead@mailinator.com', status: ContactStatus.NEW }, diff --git a/packages/schema/src/utils/outreach-eligibility.ts b/packages/schema/src/utils/outreach-eligibility.ts index df0b806d..24bdef99 100644 --- a/packages/schema/src/utils/outreach-eligibility.ts +++ b/packages/schema/src/utils/outreach-eligibility.ts @@ -60,6 +60,8 @@ export type OutreachIneligibilityReason = | 'CONTACT_ARCHIVED' | 'CONTACT_REPLIED' | 'EMAIL_SUPPRESSED' + | 'COMPANY_DNC' + | 'DOMAIN_SUPPRESSED' | 'EMAIL_DISPOSABLE' | 'EMAIL_INVALID' | 'EMAIL_QUARANTINED' @@ -96,6 +98,8 @@ export interface OutreachEligibilityInput { suppression?: { reason?: string | undefined; } | boolean | null | undefined; + companySuppressed?: boolean | null | undefined; + domainSuppressed?: boolean | null | undefined; campaign?: { id?: string | null | undefined; status?: string | null | undefined; @@ -116,7 +120,7 @@ export interface OutreachEligibilityResult { * Evaluates contact state, email quality status, domain affiliation, and campaign state. */ export function evaluateOutreachEligibility(input: OutreachEligibilityInput): OutreachEligibilityResult { - const { contact, campaign, context, suppression } = input; + const { contact, campaign, context, suppression, companySuppressed, domainSuppressed } = input; const targetEmail = input.recipientEmail || contact.email; // 1. Email existence @@ -124,10 +128,16 @@ export function evaluateOutreachEligibility(input: OutreachEligibilityInput): Ou return { eligible: false, reason: 'CONTACT_MISSING_EMAIL' }; } - // 1a. Explicit suppression check (Dedicated suppression record) + // 1a. Explicit suppression checks (Dedicated suppression records) if (suppression) { return { eligible: false, reason: 'EMAIL_SUPPRESSED' }; } + if (companySuppressed) { + return { eligible: false, reason: 'COMPANY_DNC' }; + } + if (domainSuppressed) { + return { eligible: false, reason: 'DOMAIN_SUPPRESSED' }; + } // 2. Contact CRM status (suppression checks) const contactStatus = (contact.status || '').toUpperCase(); diff --git a/packages/sdk/src/modules/suppressions.ts b/packages/sdk/src/modules/suppressions.ts index cba18403..4bf44137 100644 --- a/packages/sdk/src/modules/suppressions.ts +++ b/packages/sdk/src/modules/suppressions.ts @@ -3,7 +3,11 @@ import { toQueryString } from '../utils/query.js'; export interface SuppressionItem { id?: string; - email: string; + targetType?: 'recipient' | 'company' | 'domain'; + targetId?: string; + email?: string | null; + companyId?: string | null; + domain?: string | null; reason: string; source?: string; evidence?: Record | null; @@ -15,6 +19,7 @@ export class SuppressionsModule { constructor(private client: HttpClient) {} public async list(params?: { + targetType?: string; reason?: string; limit?: number; skip?: number; @@ -24,23 +29,85 @@ export class SuppressionsModule { } public async check( - email: string - ): Promise<{ email: string; suppressed: boolean; suppression: SuppressionItem | null }> { - return this.client.get<{ email: string; suppressed: boolean; suppression: SuppressionItem | null }>( - `/suppressions/check?email=${encodeURIComponent(email)}` - ); + queryOrEmail: string | { email?: string; companyId?: string; domain?: string } + ): Promise<{ + email?: string | null; + companyId?: string | null; + domain?: string | null; + suppressed: boolean; + suppression?: SuppressionItem | null; + isRecipientSuppressed?: boolean; + isCompanySuppressed?: boolean; + isDomainSuppressed?: boolean; + reasons?: any[]; + primaryReason?: string | null; + message?: string; + }> { + if (typeof queryOrEmail === 'string') { + return this.client.get(`/suppressions/check?email=${encodeURIComponent(queryOrEmail)}`); + } + const query = toQueryString(queryOrEmail); + return this.client.get(`/suppressions/check${query}`); } public async create(data: { - email: string; - reason: string; - source?: string; - notes?: string | null; + targetType?: 'recipient' | 'company' | 'domain'; + targetId?: string; + email?: string | null; + companyId?: string | null; + domain?: string | null; + reason?: string; + source?: string | undefined; + notes?: string | null | undefined; evidence?: any; }): Promise { return this.client.post('/suppressions', data); } + public async suppressCompany( + companyId: string, + options?: { reason?: string; source?: string; notes?: string | null; evidence?: any } + ): Promise { + return this.create({ + targetType: 'company', + companyId, + reason: options?.reason || 'COMPANY_DNC', + source: options?.source || 'manual', + notes: options?.notes, + evidence: options?.evidence + }); + } + + public async unsuppressCompany( + companyId: string + ): Promise<{ unsuppressed: boolean; companyId: string }> { + return this.client.delete<{ unsuppressed: boolean; companyId: string }>( + `/suppressions/company/${encodeURIComponent(companyId)}` + ); + } + + public async suppressDomain( + domain: string, + options?: { reason?: string; source?: string; notes?: string | null; evidence?: any } + ): Promise { + return this.create({ + targetType: 'domain', + domain, + reason: options?.reason || 'DOMAIN_SUPPRESSION', + source: options?.source || 'manual', + notes: options?.notes, + evidence: options?.evidence + }); + } + + public async unsuppressDomain( + domain: string + ): Promise<{ unsuppressed: boolean; domain: string }> { + return this.client.delete<{ unsuppressed: boolean; domain: string }>( + `/suppressions/domain/${encodeURIComponent(domain)}` + ); + } + public async delete( email: string ): Promise<{ unsuppressed: boolean; email: string; restoredContactIds?: string[] }> { @@ -49,3 +116,4 @@ export class SuppressionsModule { ); } } +