Compare commits

...
Author SHA1 Message Date
neo773 047bb0a916 nit rename EmailGroupSendType to EmailGroupMessageCategory 2026-05-29 23:55:45 +05:30
neo773 da38aec641 fix(emailing-domain): shorten suppression unique constraint to fit 63-char limit
The constraint name exceeded Postgres' 63-char identifier limit, so it was
truncated in the DB and migrate:generate kept emitting a rename (CI drift).
Shorten it, regenerate the migration off a main baseline, and simplify
suppress() back to an idempotent upsert with an empty-address guard.
2026-05-29 18:54:48 +05:30
neo773 b5f865f0b3 feat(emailing-domain): per-workspace email suppression with SES events
Replace SES account-level suppression with an app-owned, per-workspace list.

- EmailGroupSuppressedRecipient entity keyed (workspaceId, emailAddress, scope);
  scope GLOBAL (hard bounce/complaint) vs CAMPAIGN (unsubscribe), isSuppressed
  flag keeps resubscribe non-destructive.
- Handle SES Email Bounced/Complaint EventBridge events -> suppress recipients.
- Scope-aware getSuppressedAddresses(sendType); email-group sends as CAMPAIGN.
- sendEmail filters to/cc/bcc, returns delivered recipients so persistence and
  contact-creation skip suppressed addresses; throws ALL_RECIPIENTS_SUPPRESSED.
- Drop SES per-workspace contact list from the send path (1-list-per-account
  limit); rely on app-owned suppression.
- Fix AI email tool channel lookup for email_group accounts.
2026-05-29 18:14:02 +05:30
28 changed files with 512 additions and 46 deletions
@@ -0,0 +1,23 @@
import { QueryRunner } from 'typeorm';
import { RegisteredInstanceCommand } from 'src/engine/core-modules/upgrade/decorators/registered-instance-command.decorator';
import { FastInstanceCommand } from 'src/engine/core-modules/upgrade/interfaces/fast-instance-command.interface';
@RegisteredInstanceCommand('2.9.0', 1780060974610)
export class AddEmailGroupSuppressedRecipientFastInstanceCommand implements FastInstanceCommand {
public async up(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query('CREATE TYPE "core"."emailGroupSuppressedRecipient_scope_enum" AS ENUM(\'GLOBAL\', \'CAMPAIGN\')');
await queryRunner.query('CREATE TYPE "core"."emailGroupSuppressedRecipient_reason_enum" AS ENUM(\'HARD_BOUNCE\', \'COMPLAINT\', \'UNSUBSCRIBE\')');
await queryRunner.query('CREATE TYPE "core"."emailGroupSuppressedRecipient_createdbysource_enum" AS ENUM(\'EMAIL\', \'CALENDAR\', \'WORKFLOW\', \'AGENT\', \'API\', \'IMPORT\', \'MANUAL\', \'SYSTEM\', \'WEBHOOK\', \'APPLICATION\')');
await queryRunner.query('CREATE TABLE "core"."emailGroupSuppressedRecipient" ("workspaceId" uuid NOT NULL, "id" uuid NOT NULL DEFAULT uuid_generate_v4(), "emailAddress" character varying NOT NULL, "scope" "core"."emailGroupSuppressedRecipient_scope_enum" NOT NULL, "reason" "core"."emailGroupSuppressedRecipient_reason_enum" NOT NULL, "isSuppressed" boolean NOT NULL DEFAULT true, "providerEventId" character varying, "createdBySource" "core"."emailGroupSuppressedRecipient_createdbysource_enum" NOT NULL, "createdAt" TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT now(), "updatedAt" TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT now(), CONSTRAINT "IDX_EMAIL_GROUP_SUPPRESSED_RECIPIENT_WS_EMAIL_SCOPE_UNIQUE" UNIQUE ("workspaceId", "emailAddress", "scope"), CONSTRAINT "PK_55b0607e539d7941cbaedf1328a" PRIMARY KEY ("id"))');
await queryRunner.query('ALTER TABLE "core"."emailGroupSuppressedRecipient" ADD CONSTRAINT "FK_866066f1e73b748f917fd6fcd80" FOREIGN KEY ("workspaceId") REFERENCES "core"."workspace"("id") ON DELETE CASCADE ON UPDATE NO ACTION');
}
public async down(queryRunner: QueryRunner): Promise<void> {
await queryRunner.query('ALTER TABLE "core"."emailGroupSuppressedRecipient" DROP CONSTRAINT "FK_866066f1e73b748f917fd6fcd80"');
await queryRunner.query('DROP TABLE "core"."emailGroupSuppressedRecipient"');
await queryRunner.query('DROP TYPE "core"."emailGroupSuppressedRecipient_createdbysource_enum"');
await queryRunner.query('DROP TYPE "core"."emailGroupSuppressedRecipient_reason_enum"');
await queryRunner.query('DROP TYPE "core"."emailGroupSuppressedRecipient_scope_enum"');
}
}
@@ -58,6 +58,7 @@ import { DropFieldMetadataIsUniqueColumnFastInstanceCommand } from 'src/database
import { EmailingDomainTenantStatusAndGlobalUniquenessFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-fast-1799000020000-emailing-domain-tenant-status-and-global-uniqueness';
import { EncryptNonSecretApplicationVariableSlowInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-slow-1798400000000-encrypt-non-secret-application-variable';
import { MigrateAiModelPreferencesSlowInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-slow-1799000010000-migrate-ai-model-preferences';
import { AddEmailGroupSuppressedRecipientFastInstanceCommand } from 'src/database/commands/upgrade-version-command/2-9/2-9-instance-command-fast-1780060974610-add-email-group-suppressed-recipient';
export const INSTANCE_COMMANDS = [
AddViewFieldGroupIdIndexOnViewFieldFastInstanceCommand,
@@ -118,4 +119,5 @@ export const INSTANCE_COMMANDS = [
DropFieldMetadataIsUniqueColumnFastInstanceCommand,
MigrateAiModelPreferencesSlowInstanceCommand,
EncryptNonSecretApplicationVariableSlowInstanceCommand,
AddEmailGroupSuppressedRecipientFastInstanceCommand,
];
@@ -0,0 +1,15 @@
import { EmailGroupMessageCategory } from 'src/engine/core-modules/emailing-domain/types/email-group-message-category.type';
import { EmailGroupSuppressionScope } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-scope.type';
export const BLOCKED_SCOPES_BY_MESSAGE_CATEGORY: Record<
EmailGroupMessageCategory,
EmailGroupSuppressionScope[]
> = {
[EmailGroupMessageCategory.TRANSACTIONAL]: [
EmailGroupSuppressionScope.GLOBAL,
],
[EmailGroupMessageCategory.CAMPAIGN]: [
EmailGroupSuppressionScope.GLOBAL,
EmailGroupSuppressionScope.CAMPAIGN,
],
};
@@ -0,0 +1,12 @@
import { EmailGroupSuppressionReason } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-reason.type';
import { EmailGroupSuppressionScope } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-scope.type';
export const SUPPRESSION_SCOPE_BY_REASON: Record<
EmailGroupSuppressionReason,
EmailGroupSuppressionScope
> = {
[EmailGroupSuppressionReason.HARD_BOUNCE]: EmailGroupSuppressionScope.GLOBAL,
[EmailGroupSuppressionReason.COMPLAINT]: EmailGroupSuppressionScope.GLOBAL,
[EmailGroupSuppressionReason.UNSUBSCRIBE]:
EmailGroupSuppressionScope.CAMPAIGN,
};
@@ -21,7 +21,6 @@ describe('AwsSesSendEmailService', () => {
const baseContext = {
tenantName: 'twenty-workspace-ws1',
configurationSetName: 'twenty-workspace-ws1',
contactListName: 'twenty-workspace-ws1',
};
const setUp = () => {
@@ -42,7 +41,7 @@ describe('AwsSesSendEmailService', () => {
return { service, send, handleErrorService };
};
it('should call SendEmail with tenant, config set, and list management options', async () => {
it('should call SendEmail with tenant and config set', async () => {
const { service, send } = setUp();
send.mockResolvedValue({ MessageId: 'msg-1' });
@@ -59,11 +58,8 @@ describe('AwsSesSendEmailService', () => {
Destination: { ToAddresses: ['user@example.com'] },
ConfigurationSetName: 'twenty-workspace-ws1',
TenantName: 'twenty-workspace-ws1',
ListManagementOptions: {
ContactListName: 'twenty-workspace-ws1',
TopicName: 'marketing',
},
});
expect(command.input.ListManagementOptions).toBeUndefined();
expect(command.input.EmailTags).toEqual(
expect.arrayContaining([
{ Name: 'workspace', Value: 'ws1' },
@@ -136,7 +136,6 @@ export class AwsSesDriver implements EmailingDomainDriverInterface {
return this.awsSesSendEmailService.sendEmail(input, {
tenantName: this.buildTenantName(input.workspaceId),
configurationSetName: this.buildConfigurationSetName(input.workspaceId),
contactListName: this.buildContactListName(input.workspaceId),
});
}
@@ -42,7 +42,7 @@ export class AwsSesRegisterDomainService {
ConfigurationSetName: input.configurationSetName,
ReputationOptions: { ReputationMetricsEnabled: true },
SendingOptions: { SendingEnabled: true },
SuppressionOptions: { SuppressedReasons: ['BOUNCE', 'COMPLAINT'] },
SuppressionOptions: { SuppressedReasons: [] },
Tags: [{ Key: 'managed-by', Value: 'twenty' }],
}),
)
@@ -8,7 +8,6 @@ import {
type EmailingDomainSendEmailResult,
} from 'src/engine/core-modules/emailing-domain/drivers/types/send-email';
import { AWS_SES_MARKETING_TOPIC_NAME } from 'src/engine/core-modules/emailing-domain/drivers/aws-ses/constants/aws-ses-marketing-topic-name.constant';
import { AwsSesClientProvider } from 'src/engine/core-modules/emailing-domain/drivers/aws-ses/providers/aws-ses-client.provider';
import { AwsSesHandleErrorService } from 'src/engine/core-modules/emailing-domain/drivers/aws-ses/services/aws-ses-handle-error.service';
import {
@@ -19,7 +18,6 @@ import {
type SendEmailContext = {
tenantName: string;
configurationSetName: string;
contactListName: string;
};
@Injectable()
@@ -75,10 +73,6 @@ export class AwsSesSendEmailService {
},
ConfigurationSetName: context.configurationSetName,
TenantName: context.tenantName,
ListManagementOptions: {
ContactListName: context.contactListName,
TopicName: AWS_SES_MARKETING_TOPIC_NAME,
},
EmailTags: [
{ Name: 'workspace', Value: input.workspaceId },
{ Name: 'domain', Value: input.domain },
@@ -97,7 +91,14 @@ export class AwsSesSendEmailService {
`Sent email ${response.MessageId} from ${input.from} (tenant ${context.tenantName})`,
);
return { messageId: response.MessageId };
return {
messageId: response.MessageId,
deliveredRecipients: {
to: input.to,
cc: input.cc ?? [],
bcc: input.bcc ?? [],
},
};
} catch (error) {
if (error instanceof EmailingDomainDriverException) {
throw error;
@@ -11,6 +11,7 @@ export enum EmailingDomainDriverExceptionCode {
INSUFFICIENT_PERMISSIONS = 'INSUFFICIENT_PERMISSIONS',
CONFIGURATION_ERROR = 'CONFIGURATION_ERROR',
SENDING_SUSPENDED = 'SENDING_SUSPENDED',
ALL_RECIPIENTS_SUPPRESSED = 'ALL_RECIPIENTS_SUPPRESSED',
UNKNOWN = 'UNKNOWN',
}
@@ -26,6 +27,8 @@ const getEmailingDomainDriverExceptionUserFriendlyMessage = (
return msg`Email domain configuration error.`;
case EmailingDomainDriverExceptionCode.SENDING_SUSPENDED:
return msg`Sending is currently suspended for this email domain.`;
case EmailingDomainDriverExceptionCode.ALL_RECIPIENTS_SUPPRESSED:
return msg`All recipients are suppressed for this email domain.`;
case EmailingDomainDriverExceptionCode.TEMPORARY_ERROR:
case EmailingDomainDriverExceptionCode.UNKNOWN:
return STANDARD_ERROR_MESSAGE;
@@ -1,3 +1,5 @@
import { EmailGroupMessageCategory } from 'src/engine/core-modules/emailing-domain/types/email-group-message-category.type';
export type EmailingDomainAttachment = {
filename: string;
content: Buffer;
@@ -14,6 +16,7 @@ export type EmailingDomainEmailContent = {
html?: string;
replyTo?: string[];
attachments?: EmailingDomainAttachment[];
messageCategory?: EmailGroupMessageCategory;
};
export type EmailingDomainSendEmailInput = EmailingDomainEmailContent & {
@@ -23,4 +26,5 @@ export type EmailingDomainSendEmailInput = EmailingDomainEmailContent & {
export type EmailingDomainSendEmailResult = {
messageId: string;
deliveredRecipients: { to: string[]; cc: string[]; bcc: string[] };
};
@@ -0,0 +1,61 @@
import {
Column,
CreateDateColumn,
Entity,
PrimaryGeneratedColumn,
Unique,
UpdateDateColumn,
} from 'typeorm';
import { FieldActorSource } from 'twenty-shared/types';
import { EmailGroupSuppressionReason } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-reason.type';
import { EmailGroupSuppressionScope } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-scope.type';
import { WorkspaceRelatedEntity } from 'src/engine/workspace-manager/types/workspace-related-entity';
@Entity({ name: 'emailGroupSuppressedRecipient', schema: 'core' })
@Unique('IDX_EMAIL_GROUP_SUPPRESSED_RECIPIENT_WS_EMAIL_SCOPE_UNIQUE', [
'workspaceId',
'emailAddress',
'scope',
])
export class EmailGroupSuppressedRecipientEntity extends WorkspaceRelatedEntity {
@PrimaryGeneratedColumn('uuid')
id: string;
@Column({ type: 'varchar', nullable: false })
emailAddress: string;
@Column({
type: 'enum',
enum: Object.values(EmailGroupSuppressionScope),
nullable: false,
})
scope: EmailGroupSuppressionScope;
@Column({
type: 'enum',
enum: Object.values(EmailGroupSuppressionReason),
nullable: false,
})
reason: EmailGroupSuppressionReason;
@Column({ type: 'boolean', default: true, nullable: false })
isSuppressed: boolean;
@Column({ type: 'varchar', nullable: true })
providerEventId: string | null;
@Column({
type: 'enum',
enum: Object.values(FieldActorSource),
nullable: false,
})
createdBySource: FieldActorSource;
@CreateDateColumn({ type: 'timestamptz' })
createdAt: Date;
@UpdateDateColumn({ type: 'timestamptz' })
updatedAt: Date;
}
@@ -8,9 +8,11 @@ import { AwsSesRegisterDomainService } from 'src/engine/core-modules/emailing-do
import { AwsSesHandleErrorService } from 'src/engine/core-modules/emailing-domain/drivers/aws-ses/services/aws-ses-handle-error.service';
import { AwsSesSendEmailService } from 'src/engine/core-modules/emailing-domain/drivers/aws-ses/services/aws-ses-send-email.service';
import { EmailingDomainDriverFactory } from 'src/engine/core-modules/emailing-domain/drivers/emailing-domain-driver.factory';
import { EmailGroupSuppressedRecipientEntity } from 'src/engine/core-modules/emailing-domain/email-group-suppressed-recipient.entity';
import { EmailingDomainEntity } from 'src/engine/core-modules/emailing-domain/emailing-domain.entity';
import { EmailingDomainResolver } from 'src/engine/core-modules/emailing-domain/emailing-domain.resolver';
import { EmailingDomainWorkspaceCleanupJob } from 'src/engine/core-modules/emailing-domain/jobs/emailing-domain-workspace-cleanup.job';
import { EmailGroupSuppressionService } from 'src/engine/core-modules/emailing-domain/services/email-group-suppression.service';
import { EmailingDomainTenantStatusService } from 'src/engine/core-modules/emailing-domain/services/emailing-domain-tenant-status.service';
import { EmailingDomainService } from 'src/engine/core-modules/emailing-domain/services/emailing-domain.service';
import { FeatureFlagModule } from 'src/engine/core-modules/feature-flag/feature-flag.module';
@@ -19,14 +21,22 @@ import { provideWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspac
@Module({
imports: [
TypeORMModule,
NestjsQueryTypeOrmModule.forFeature([EmailingDomainEntity]),
NestjsQueryTypeOrmModule.forFeature([
EmailingDomainEntity,
EmailGroupSuppressedRecipientEntity,
]),
FeatureFlagModule,
PermissionsModule,
],
exports: [EmailingDomainService, EmailingDomainTenantStatusService],
exports: [
EmailingDomainService,
EmailingDomainTenantStatusService,
EmailGroupSuppressionService,
],
providers: [
EmailingDomainService,
EmailingDomainTenantStatusService,
EmailGroupSuppressionService,
EmailingDomainResolver,
EmailingDomainDriverFactory,
EmailingDomainWorkspaceCleanupJob,
@@ -35,6 +45,7 @@ import { provideWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspac
AwsSesRegisterDomainService,
AwsSesSendEmailService,
provideWorkspaceScopedRepository(EmailingDomainEntity),
provideWorkspaceScopedRepository(EmailGroupSuppressedRecipientEntity),
],
})
export class EmailingDomainModule {}
@@ -2,7 +2,9 @@ import { EmailingDomainDriverExceptionCode } from 'src/engine/core-modules/email
import { type EmailingDomainDriverFactory } from 'src/engine/core-modules/emailing-domain/drivers/emailing-domain-driver.factory';
import { EmailingDomainStatus } from 'src/engine/core-modules/emailing-domain/drivers/types/emailing-domain-status.type';
import { EmailingDomainTenantStatus } from 'src/engine/core-modules/emailing-domain/drivers/types/emailing-domain-tenant-status.type';
import { type EmailingDomainEmailContent } from 'src/engine/core-modules/emailing-domain/drivers/types/send-email';
import { type EmailingDomainEntity } from 'src/engine/core-modules/emailing-domain/emailing-domain.entity';
import { type EmailGroupSuppressionService } from 'src/engine/core-modules/emailing-domain/services/email-group-suppression.service';
import { EmailingDomainService } from 'src/engine/core-modules/emailing-domain/services/emailing-domain.service';
import { type WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
@@ -19,14 +21,20 @@ describe('EmailingDomainService.sendEmail', () => {
...overrides,
}) as EmailingDomainEntity;
const buildEmailContent = () => ({
const buildEmailContent = (
overrides: Partial<EmailingDomainEmailContent> = {},
): EmailingDomainEmailContent => ({
from: 'hello@mail.example.com',
to: ['user@example.com'],
subject: 'Hi',
text: 'Body',
...overrides,
});
const setUp = (emailingDomain: EmailingDomainEntity) => {
const setUp = (
emailingDomain: EmailingDomainEntity,
suppressedAddresses: string[] = [],
) => {
const sendEmail = jest.fn().mockResolvedValue({ messageId: 'msg-1' });
const repository = {
findOne: jest.fn().mockResolvedValue(emailingDomain),
@@ -34,7 +42,18 @@ describe('EmailingDomainService.sendEmail', () => {
const factory = {
getCurrentDriver: () => ({ sendEmail }),
} as unknown as EmailingDomainDriverFactory;
const service = new EmailingDomainService(repository, factory);
const suppressionService = {
getSuppressedAddresses: jest
.fn()
.mockResolvedValue(
new Set(suppressedAddresses.map((address) => address.toLowerCase())),
),
} as unknown as EmailGroupSuppressionService;
const service = new EmailingDomainService(
repository,
factory,
suppressionService,
);
return { service, sendEmail };
};
@@ -77,6 +96,43 @@ describe('EmailingDomainService.sendEmail', () => {
},
);
it('removes suppressed recipients but still sends to deliverable ones', async () => {
const { service, sendEmail } = setUp(buildEmailingDomain(), [
'blocked@example.com',
]);
await service.sendEmail(
'ws1',
'domain-1',
buildEmailContent({
to: ['user@example.com', 'Blocked@example.com'],
cc: ['blocked@example.com'],
bcc: ['keep@example.com'],
}),
);
expect(sendEmail).toHaveBeenCalledWith(
expect.objectContaining({
to: ['user@example.com'],
cc: [],
bcc: ['keep@example.com'],
}),
);
});
it('rejects with ALL_RECIPIENTS_SUPPRESSED when every primary recipient is suppressed, without calling the driver', async () => {
const { service, sendEmail } = setUp(buildEmailingDomain(), [
'user@example.com',
]);
await expect(
service.sendEmail('ws1', 'domain-1', buildEmailContent()),
).rejects.toMatchObject({
code: EmailingDomainDriverExceptionCode.ALL_RECIPIENTS_SUPPRESSED,
});
expect(sendEmail).not.toHaveBeenCalled();
});
// Verification is a precondition for the tenant-status check: a domain that
// has not been verified should surface a CONFIGURATION_ERROR rather than
// leaking the tenant pause state to callers who couldn't have used it anyway.
@@ -0,0 +1,88 @@
import { Injectable } from '@nestjs/common';
import { isNonEmptyString } from '@sniptt/guards';
import { FieldActorSource } from 'twenty-shared/types';
import { isNonEmptyArray } from 'twenty-shared/utils';
import { In } from 'typeorm';
import { BLOCKED_SCOPES_BY_MESSAGE_CATEGORY } from 'src/engine/core-modules/emailing-domain/constants/blocked-scopes-by-message-category.constant';
import { SUPPRESSION_SCOPE_BY_REASON } from 'src/engine/core-modules/emailing-domain/constants/suppression-scope-by-reason.constant';
import { EmailGroupSuppressedRecipientEntity } from 'src/engine/core-modules/emailing-domain/email-group-suppressed-recipient.entity';
import { EmailGroupMessageCategory } from 'src/engine/core-modules/emailing-domain/types/email-group-message-category.type';
import { EmailGroupSuppressionReason } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-reason.type';
import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator';
import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
type SuppressRecipientArgs = {
workspaceId: string;
emailAddress: string;
reason: EmailGroupSuppressionReason;
createdBySource: FieldActorSource;
providerEventId?: string | null;
};
@Injectable()
export class EmailGroupSuppressionService {
constructor(
@InjectWorkspaceScopedRepository(EmailGroupSuppressedRecipientEntity)
private readonly suppressedRecipientRepository: WorkspaceScopedRepository<EmailGroupSuppressedRecipientEntity>,
) {}
async getSuppressedAddresses(
workspaceId: string,
emailAddresses: string[],
messageCategory: EmailGroupMessageCategory,
): Promise<Set<string>> {
const normalizedAddresses = [
...new Set(
emailAddresses.map((emailAddress) => emailAddress.trim().toLowerCase()),
),
];
if (!isNonEmptyArray(normalizedAddresses)) {
return new Set();
}
const suppressedRecipients = await this.suppressedRecipientRepository.find(
workspaceId,
{
where: {
emailAddress: In(normalizedAddresses),
scope: In(BLOCKED_SCOPES_BY_MESSAGE_CATEGORY[messageCategory]),
isSuppressed: true,
},
},
);
return new Set(
suppressedRecipients.map((recipient) => recipient.emailAddress),
);
}
async suppress({
workspaceId,
emailAddress,
reason,
createdBySource,
providerEventId = null,
}: SuppressRecipientArgs): Promise<void> {
const normalizedEmailAddress = emailAddress.trim().toLowerCase();
if (!isNonEmptyString(normalizedEmailAddress)) {
return;
}
await this.suppressedRecipientRepository.upsert(
workspaceId,
{
emailAddress: normalizedEmailAddress,
scope: SUPPRESSION_SCOPE_BY_REASON[reason],
reason,
isSuppressed: true,
createdBySource,
providerEventId,
},
['workspaceId', 'emailAddress', 'scope'],
);
}
}
@@ -13,6 +13,8 @@ import {
type EmailingDomainSendEmailResult,
} from 'src/engine/core-modules/emailing-domain/drivers/types/send-email';
import { EmailingDomainEntity } from 'src/engine/core-modules/emailing-domain/emailing-domain.entity';
import { EmailGroupSuppressionService } from 'src/engine/core-modules/emailing-domain/services/email-group-suppression.service';
import { EmailGroupMessageCategory } from 'src/engine/core-modules/emailing-domain/types/email-group-message-category.type';
import { WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity';
import { InjectWorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/inject-workspace-scoped-repository.decorator';
import { WorkspaceScopedRepository } from 'src/engine/twenty-orm/workspace-scoped-repository/workspace-scoped-repository';
@@ -24,6 +26,7 @@ export class EmailingDomainService {
@InjectWorkspaceScopedRepository(EmailingDomainEntity)
private readonly emailingDomainRepository: WorkspaceScopedRepository<EmailingDomainEntity>,
private readonly emailingDomainDriverFactory: EmailingDomainDriverFactory,
private readonly emailGroupSuppressionService: EmailGroupSuppressionService,
) {}
async createEmailingDomain(
@@ -190,8 +193,34 @@ export class EmailingDomainService {
);
}
const suppressedAddresses =
await this.emailGroupSuppressionService.getSuppressedAddresses(
workspaceId,
[
...emailContent.to,
...(emailContent.cc ?? []),
...(emailContent.bcc ?? []),
],
emailContent.messageCategory ?? EmailGroupMessageCategory.TRANSACTIONAL,
);
const isNotSuppressed = (address: string): boolean =>
!suppressedAddresses.has(address.trim().toLowerCase());
const deliverableTo = emailContent.to.filter(isNotSuppressed);
if (deliverableTo.length === 0) {
throw new EmailingDomainDriverException(
`All primary recipients are suppressed for emailing domain ${emailingDomain.domain}`,
EmailingDomainDriverExceptionCode.ALL_RECIPIENTS_SUPPRESSED,
);
}
return this.emailingDomainDriverFactory.getCurrentDriver().sendEmail({
...emailContent,
to: deliverableTo,
cc: emailContent.cc?.filter(isNotSuppressed),
bcc: emailContent.bcc?.filter(isNotSuppressed),
workspaceId,
domain: emailingDomain.domain,
});
@@ -0,0 +1,4 @@
export enum EmailGroupMessageCategory {
TRANSACTIONAL = 'TRANSACTIONAL',
CAMPAIGN = 'CAMPAIGN',
}
@@ -0,0 +1,5 @@
export enum EmailGroupSuppressionReason {
HARD_BOUNCE = 'HARD_BOUNCE',
COMPLAINT = 'COMPLAINT',
UNSUBSCRIBE = 'UNSUBSCRIBE',
}
@@ -0,0 +1,4 @@
export enum EmailGroupSuppressionScope {
GLOBAL = 'GLOBAL',
CAMPAIGN = 'CAMPAIGN',
}
@@ -5,6 +5,7 @@ import { MessagingWebhooksController } from 'src/engine/core-modules/messaging-w
import { SesInboundMailHandlerService } from 'src/engine/core-modules/messaging-webhooks/services/ses-inbound-mail-handler.service';
import { SesInboundWebhookRouterService } from 'src/engine/core-modules/messaging-webhooks/services/ses-inbound-webhook-router.service';
import { SesOutboundSendingStateHandlerService } from 'src/engine/core-modules/messaging-webhooks/services/ses-outbound-sending-state-handler.service';
import { SesOutboundSuppressionHandlerService } from 'src/engine/core-modules/messaging-webhooks/services/ses-outbound-suppression-handler.service';
import { SesOutboundWebhookRouterService } from 'src/engine/core-modules/messaging-webhooks/services/ses-outbound-webhook-router.service';
import { SnsSignatureVerifierService } from 'src/engine/core-modules/messaging-webhooks/services/sns-signature-verifier.service';
import { SnsSubscriptionConfirmerService } from 'src/engine/core-modules/messaging-webhooks/services/sns-subscription-confirmer.service';
@@ -18,6 +19,7 @@ import { TwentyConfigModule } from 'src/engine/core-modules/twenty-config/twenty
SnsSubscriptionConfirmerService,
SesInboundMailHandlerService,
SesOutboundSendingStateHandlerService,
SesOutboundSuppressionHandlerService,
SesInboundWebhookRouterService,
SesOutboundWebhookRouterService,
],
@@ -3,8 +3,8 @@ import { Injectable, Logger } from '@nestjs/common';
import { EmailingDomainTenantStatus } from 'src/engine/core-modules/emailing-domain/drivers/types/emailing-domain-tenant-status.type';
import { EmailingDomainTenantStatusService } from 'src/engine/core-modules/emailing-domain/services/emailing-domain-tenant-status.service';
import { type SesEventBridgeNotification } from 'src/engine/core-modules/messaging-webhooks/types/ses-event-bridge-notification.type';
import { parseWorkspaceIdFromAwsSesResourceArn } from 'src/engine/core-modules/messaging-webhooks/utils/parse-workspace-id-from-aws-ses-resource-arn.util';
import { isDefined, isNonEmptyArray } from 'twenty-shared/utils';
import { resolveWorkspaceIdFromAwsSesResources } from 'src/engine/core-modules/messaging-webhooks/utils/resolve-workspace-id-from-aws-ses-resources.util';
import { isDefined } from 'twenty-shared/utils';
@Injectable()
export class SesOutboundSendingStateHandlerService {
@@ -22,7 +22,7 @@ export class SesOutboundSendingStateHandlerService {
? EmailingDomainTenantStatus.ACTIVE
: EmailingDomainTenantStatus.PAUSED;
const workspaceId = this.resolveWorkspaceIdFromResources(event.resources);
const workspaceId = resolveWorkspaceIdFromAwsSesResources(event.resources);
if (!isDefined(workspaceId)) {
this.logger.warn(
@@ -37,22 +37,4 @@ export class SesOutboundSendingStateHandlerService {
targetStatus,
);
}
private resolveWorkspaceIdFromResources(
resources: string[] | undefined,
): string | null {
if (!isNonEmptyArray(resources)) {
return null;
}
for (const resourceArn of resources) {
const workspaceId = parseWorkspaceIdFromAwsSesResourceArn(resourceArn);
if (isDefined(workspaceId)) {
return workspaceId;
}
}
return null;
}
}
@@ -0,0 +1,111 @@
import { Injectable, Logger } from '@nestjs/common';
import { FieldActorSource } from 'twenty-shared/types';
import { isDefined, isNonEmptyArray } from 'twenty-shared/utils';
import { EmailGroupSuppressionService } from 'src/engine/core-modules/emailing-domain/services/email-group-suppression.service';
import { EmailGroupSuppressionReason } from 'src/engine/core-modules/emailing-domain/types/email-group-suppression-reason.type';
import { type SesEventBridgeNotification } from 'src/engine/core-modules/messaging-webhooks/types/ses-event-bridge-notification.type';
import { resolveWorkspaceIdFromAwsSesResources } from 'src/engine/core-modules/messaging-webhooks/utils/resolve-workspace-id-from-aws-ses-resources.util';
@Injectable()
export class SesOutboundSuppressionHandlerService {
private readonly logger = new Logger(
SesOutboundSuppressionHandlerService.name,
);
constructor(
private readonly emailGroupSuppressionService: EmailGroupSuppressionService,
) {}
async handle(event: SesEventBridgeNotification): Promise<void> {
const suppression = this.resolveSuppression(event);
if (!isDefined(suppression)) {
return;
}
const workspaceId = resolveWorkspaceIdFromAwsSesResources(event.resources);
if (!isDefined(workspaceId)) {
this.logger.warn(
`Could not resolve workspaceId from SES ${event['detail-type']} event resources: ${JSON.stringify(event.resources)}`,
);
return;
}
const results = await Promise.allSettled(
suppression.emailAddresses.map((emailAddress) =>
this.emailGroupSuppressionService.suppress({
workspaceId,
emailAddress,
reason: suppression.reason,
createdBySource: FieldActorSource.WEBHOOK,
providerEventId: suppression.providerEventId,
}),
),
);
if (results.some((result) => result.status === 'rejected')) {
throw new Error(
`Failed to suppress one or more recipients for ${event['detail-type']} event in workspace ${workspaceId}`,
);
}
}
private resolveSuppression(event: SesEventBridgeNotification): {
reason: EmailGroupSuppressionReason;
emailAddresses: string[];
providerEventId: string | null;
} | null {
if (event['detail-type'] === 'Email Bounced') {
const bounce = event.detail?.bounce;
if (bounce?.bounceType !== 'Permanent') {
return null;
}
const emailAddresses = this.extractRecipientAddresses(
bounce.bouncedRecipients,
);
if (!isNonEmptyArray(emailAddresses)) {
return null;
}
return {
reason: EmailGroupSuppressionReason.HARD_BOUNCE,
emailAddresses,
providerEventId: bounce.feedbackId ?? null,
};
}
if (event['detail-type'] === 'Email Complaint Received') {
const complaint = event.detail?.complaint;
const emailAddresses = this.extractRecipientAddresses(
complaint?.complainedRecipients,
);
if (!isNonEmptyArray(emailAddresses)) {
return null;
}
return {
reason: EmailGroupSuppressionReason.COMPLAINT,
emailAddresses,
providerEventId: complaint?.feedbackId ?? null,
};
}
return null;
}
private extractRecipientAddresses = (
recipients: { emailAddress: string }[] | undefined,
): string[] => {
return (recipients ?? [])
.map((recipient) => recipient.emailAddress)
.filter(isDefined);
};
}
@@ -6,6 +6,7 @@ import { isDefined, parseJson } from 'twenty-shared/utils';
import { MessagingWebhookExceptionCode } from 'src/engine/core-modules/messaging-webhooks/messaging-webhook-exception-code.enum';
import { MessagingWebhookException } from 'src/engine/core-modules/messaging-webhooks/messaging-webhook.exception';
import { SesOutboundSendingStateHandlerService } from 'src/engine/core-modules/messaging-webhooks/services/ses-outbound-sending-state-handler.service';
import { SesOutboundSuppressionHandlerService } from 'src/engine/core-modules/messaging-webhooks/services/ses-outbound-suppression-handler.service';
import { SnsSignatureVerifierService } from 'src/engine/core-modules/messaging-webhooks/services/sns-signature-verifier.service';
import { SnsSubscriptionConfirmerService } from 'src/engine/core-modules/messaging-webhooks/services/sns-subscription-confirmer.service';
import { type SesEventBridgeNotification } from 'src/engine/core-modules/messaging-webhooks/types/ses-event-bridge-notification.type';
@@ -18,6 +19,7 @@ export class SesOutboundWebhookRouterService {
private readonly snsSignatureVerifierService: SnsSignatureVerifierService,
private readonly snsSubscriptionConfirmerService: SnsSubscriptionConfirmerService,
private readonly sesOutboundSendingStateHandlerService: SesOutboundSendingStateHandlerService,
private readonly sesOutboundSuppressionHandlerService: SesOutboundSuppressionHandlerService,
) {}
async route(rawBody: Buffer): Promise<void> {
@@ -54,6 +56,15 @@ export class SesOutboundWebhookRouterService {
);
}
if (
event['detail-type'] === 'Email Bounced' ||
event['detail-type'] === 'Email Complaint Received'
) {
await this.sesOutboundSuppressionHandlerService.handle(event);
return;
}
await this.sesOutboundSendingStateHandlerService.handle(event);
}
}
@@ -1,6 +1,16 @@
export type SesEventBridgeDetailType =
| 'Sending Status Enabled'
| 'Sending Status Disabled'
| 'Email Bounced'
| 'Email Complaint Received';
export type SesEventBridgeRecipient = {
emailAddress: string;
};
export type SesEventBridgeNotification = {
source: 'aws.ses';
'detail-type': 'Sending Status Enabled' | 'Sending Status Disabled';
'detail-type': SesEventBridgeDetailType;
resources?: string[];
detail?: {
version?: string;
@@ -11,5 +21,14 @@ export type SesEventBridgeNotification = {
cause?: string;
};
};
bounce?: {
bounceType?: 'Permanent' | 'Transient' | 'Undetermined';
feedbackId?: string;
bouncedRecipients?: SesEventBridgeRecipient[];
};
complaint?: {
feedbackId?: string;
complainedRecipients?: SesEventBridgeRecipient[];
};
};
};
@@ -0,0 +1,21 @@
import { isDefined, isNonEmptyArray } from 'twenty-shared/utils';
import { parseWorkspaceIdFromAwsSesResourceArn } from 'src/engine/core-modules/messaging-webhooks/utils/parse-workspace-id-from-aws-ses-resource-arn.util';
export const resolveWorkspaceIdFromAwsSesResources = (
resources: string[] | undefined,
): string | null => {
if (!isNonEmptyArray(resources)) {
return null;
}
for (const resourceArn of resources) {
const workspaceId = parseWorkspaceIdFromAwsSesResourceArn(resourceArn);
if (isDefined(workspaceId)) {
return workspaceId;
}
}
return null;
};
@@ -362,9 +362,12 @@ export class EmailComposerService {
workspaceId,
);
const messageChannel = connectedAccount.messageChannels.find(
(channel) => channel.handle === connectedAccount.handle,
);
const messageChannel =
connectedAccount.provider === ConnectedAccountProvider.EMAIL_GROUP
? connectedAccount.messageChannels[0]
: connectedAccount.messageChannels.find(
(channel) => channel.handle === connectedAccount.handle,
);
const isSmtpOnlyAccount =
connectedAccount.provider === ConnectedAccountProvider.IMAP_SMTP_CALDAV &&
@@ -6,6 +6,7 @@ import { isDefined } from 'twenty-shared/utils';
import { EmailingDomainStatus } from 'src/engine/core-modules/emailing-domain/drivers/types/emailing-domain-status.type';
import { EmailingDomainEntity } from 'src/engine/core-modules/emailing-domain/emailing-domain.entity';
import { EmailingDomainService } from 'src/engine/core-modules/emailing-domain/services/emailing-domain.service';
import { EmailGroupMessageCategory } from 'src/engine/core-modules/emailing-domain/types/email-group-message-category.type';
import { type ConnectedAccountEntity } from 'src/engine/metadata-modules/connected-account/entities/connected-account.entity';
import {
MessageChannelException,
@@ -53,12 +54,14 @@ export class EmailGroupMessageOutboundService implements MessageOutboundDriver {
from: connectedAccount.handle,
replyTo: [connectedAccount.handle],
attachments: sendMessageInput.attachments,
messageCategory: EmailGroupMessageCategory.CAMPAIGN,
},
);
return {
headerMessageId: result.messageId,
messageExternalId: result.messageId,
deliveredRecipients: result.deliveredRecipients,
};
}
@@ -42,7 +42,7 @@ export class SendEmailService {
sendResult,
subject: data.sanitizedSubject,
body: data.plainTextBody,
recipients: data.recipients,
recipients: sendResult.deliveredRecipients ?? data.recipients,
connectedAccount: data.connectedAccount,
messageChannelId: data.messageChannelId!,
inReplyTo: data.inReplyTo,
@@ -2,4 +2,5 @@ export type SendMessageResult = {
headerMessageId: string;
messageExternalId?: string;
threadExternalId?: string;
deliveredRecipients?: { to: string[]; cc: string[]; bcc: string[] };
};